From 4270af3ac9810931b9d056a9896d6c63ae42efa4 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Tue, 28 Apr 2026 10:46:08 +0900 Subject: [PATCH] feat(team-mode): add team tasklist claim with tests --- .../team-mode/team-tasklist/claim.test.ts | 99 +++++++++++++++++++ src/features/team-mode/team-tasklist/claim.ts | 98 ++++++++++++++++++ 2 files changed, 197 insertions(+) create mode 100644 src/features/team-mode/team-tasklist/claim.test.ts create mode 100644 src/features/team-mode/team-tasklist/claim.ts diff --git a/src/features/team-mode/team-tasklist/claim.test.ts b/src/features/team-mode/team-tasklist/claim.test.ts new file mode 100644 index 000000000..9b41abff8 --- /dev/null +++ b/src/features/team-mode/team-tasklist/claim.test.ts @@ -0,0 +1,99 @@ +/// + +import { expect, test } from "bun:test" +import { writeFile } from "node:fs/promises" +import path from "node:path" + +import { getTasksDir, resolveBaseDir } from "../team-registry" +import { claimTask, AlreadyClaimedError, BlockedByError } from "./claim" +import { createTask } from "./store" +import { createTaskInput, createTasklistFixture } from "./test-support" +import { updateTaskStatus } from "./update" + +test("claimTask allows exactly one concurrent claimant", async () => { + // given + const fixture = await createTasklistFixture() + + try { + const task = await createTask(fixture.teamRunId, createTaskInput(), fixture.config) + + // when + const claimResults = await Promise.allSettled([ + claimTask(fixture.teamRunId, task.id, "member-a", fixture.config), + claimTask(fixture.teamRunId, task.id, "member-b", fixture.config), + ]) + + const successfulClaims = claimResults.filter((result) => result.status === "fulfilled") + const failedClaims = claimResults.filter((result) => result.status === "rejected") + + // then + expect(successfulClaims).toHaveLength(1) + expect(failedClaims).toHaveLength(1) + expect(failedClaims[0]?.status).toBe("rejected") + if (failedClaims[0]?.status === "rejected") { + expect(failedClaims[0].reason).toBeInstanceOf(AlreadyClaimedError) + } + } finally { + await fixture.cleanup() + } +}) + +test("claimTask rejects blocked tasks until blockers complete", async () => { + // given + const fixture = await createTasklistFixture() + + try { + const blockerTask = await createTask(fixture.teamRunId, createTaskInput({ subject: "blocker" }), fixture.config) + const blockedTask = await createTask( + fixture.teamRunId, + createTaskInput({ subject: "blocked", blockedBy: [blockerTask.id] }), + fixture.config, + ) + + // when + let blockedError: unknown = null + try { + await claimTask(fixture.teamRunId, blockedTask.id, "member-a", fixture.config) + } catch (error) { + blockedError = error + } + + // then + expect(blockedError).toBeInstanceOf(BlockedByError) + + // given + await claimTask(fixture.teamRunId, blockerTask.id, "member-b", fixture.config) + await updateTaskStatus(fixture.teamRunId, blockerTask.id, "in_progress", "member-b", fixture.config) + await updateTaskStatus(fixture.teamRunId, blockerTask.id, "completed", "member-b", fixture.config) + + // when + const claimedTask = await claimTask(fixture.teamRunId, blockedTask.id, "member-a", fixture.config) + + // then + expect(claimedTask.status).toBe("claimed") + expect(claimedTask.owner).toBe("member-a") + } finally { + await fixture.cleanup() + } +}) + +test("claimTask reaps a stale claim lock before claiming", async () => { + // given + const fixture = await createTasklistFixture() + + try { + const task = await createTask(fixture.teamRunId, createTaskInput(), fixture.config) + const tasksDirectory = getTasksDir(resolveBaseDir(fixture.config), fixture.teamRunId) + const staleLockPath = path.join(tasksDirectory, "claims", `${task.id}.lock`) + await writeFile(staleLockPath, `member-z\n999999\n${Date.now() - 600_000}\n`) + + // when + const claimedTask = await claimTask(fixture.teamRunId, task.id, "member-a", fixture.config) + + // then + expect(claimedTask.status).toBe("claimed") + expect(claimedTask.owner).toBe("member-a") + } finally { + await fixture.cleanup() + } +}) diff --git a/src/features/team-mode/team-tasklist/claim.ts b/src/features/team-mode/team-tasklist/claim.ts new file mode 100644 index 000000000..986b83916 --- /dev/null +++ b/src/features/team-mode/team-tasklist/claim.ts @@ -0,0 +1,98 @@ +import { access, mkdir } from "node:fs/promises" +import path from "node:path" + +import type { TeamModeConfig } from "../../../config/schema/team-mode" +import { getTasksDir, resolveBaseDir } from "../team-registry" +import { atomicWrite, detectStaleLock, reapStaleLock, withLock } from "../team-state-store/locks" +import { TaskSchema } from "../types" +import type { Task } from "../types" +import { canClaim } from "./dependencies" +import { getTask } from "./get" +import { listTasks } from "./list" + +const CLAIM_STALE_AFTER_MS = 300_000 + +async function lockExists(lockPath: string): Promise { + try { + await access(lockPath) + return true + } catch { + return false + } +} + +function getBlockingTaskIds(task: Task, allTasks: Task[]): string[] { + return task.blockedBy.filter((blockerId) => { + const blockerTask = allTasks.find((candidateTask) => candidateTask.id === blockerId) + return blockerTask !== undefined && blockerTask.status !== "completed" + }) +} + +export class AlreadyClaimedError extends Error { + constructor(message = "already_claimed") { + super(message) + this.name = "AlreadyClaimedError" + } +} + +export class BlockedByError extends Error { + constructor(public readonly blockers: string[]) { + super(`blocked by ${blockers.join(",")}`) + this.name = "BlockedByError" + } +} + +export async function claimTask( + teamRunId: string, + taskId: string, + memberName: string, + config: TeamModeConfig, +): Promise { + const baseDirectory = resolveBaseDir(config) + const tasksDirectory = getTasksDir(baseDirectory, teamRunId) + const claimsDirectory = path.join(tasksDirectory, "claims") + const taskPath = path.join(tasksDirectory, `${taskId}.json`) + const claimLockPath = path.join(claimsDirectory, `${taskId}.lock`) + + await mkdir(claimsDirectory, { recursive: true, mode: 0o700 }) + + const task = await getTask(teamRunId, taskId, config) + if (task.status !== "pending") { + throw new AlreadyClaimedError() + } + + const allTasks = await listTasks(teamRunId, config) + if (!canClaim(task, allTasks)) { + throw new BlockedByError(getBlockingTaskIds(task, allTasks)) + } + + if (await detectStaleLock(claimLockPath, CLAIM_STALE_AFTER_MS)) { + await reapStaleLock(claimLockPath) + } else if (await lockExists(claimLockPath)) { + throw new AlreadyClaimedError() + } + + return withLock(claimLockPath, async () => { + const refreshedTask = await getTask(teamRunId, taskId, config) + if (refreshedTask.status !== "pending") { + throw new AlreadyClaimedError() + } + + const refreshedTasks = await listTasks(teamRunId, config) + if (!canClaim(refreshedTask, refreshedTasks)) { + throw new BlockedByError(getBlockingTaskIds(refreshedTask, refreshedTasks)) + } + + const now = Date.now() + const updatedTask = TaskSchema.parse({ + ...refreshedTask, + status: "claimed", + owner: memberName, + claimedAt: now, + updatedAt: now, + }) + + await atomicWrite(taskPath, `${JSON.stringify(updatedTask, null, 2)}\n`) + return updatedTask + }, { ownerTag: memberName, staleAfterMs: CLAIM_STALE_AFTER_MS }) +}