feat(team-mode): add team tasklist claim with tests
This commit is contained in:
@@ -0,0 +1,99 @@
|
|||||||
|
/// <reference types="bun-types" />
|
||||||
|
|
||||||
|
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()
|
||||||
|
}
|
||||||
|
})
|
||||||
@@ -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<boolean> {
|
||||||
|
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<Task> {
|
||||||
|
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 })
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user