feat(FN-1253): implement task checkout leasing end-to-end
- Add checkout lease types and conflict error exports, plus DB schema v20 migration for checkedOutBy/checkedOutAt - Persist checkout lease fields in TaskStore and add AgentStore checkout/release/force-release/get-holder operations - Add dashboard checkout API routes for acquire/release/force-release/status with explicit 409 conflict and 403 holder enforcement - Enforce checkout ownership in heartbeat execution with graceful checkout_conflict exits when another agent holds the lease - Expand core and dashboard test coverage for schema, store behavior, API routes, and leasing workflows, and document leasing behavior in AGENTS.md
This commit is contained in:
37
AGENTS.md
37
AGENTS.md
@@ -706,6 +706,43 @@ interface WakeContext {
|
|||||||
- `InProcessRuntime.stop()` stops the trigger scheduler before stopping the HeartbeatMonitor
|
- `InProcessRuntime.stop()` stops the trigger scheduler before stopping the HeartbeatMonitor
|
||||||
- `InProcessRuntime.getTriggerScheduler()` — Returns the scheduler instance for testing access
|
- `InProcessRuntime.getTriggerScheduler()` — Returns the scheduler instance for testing access
|
||||||
|
|
||||||
|
## Checkout Leasing
|
||||||
|
|
||||||
|
Task ownership now supports explicit checkout leases modeled after Paperclip's checkout/release flow.
|
||||||
|
|
||||||
|
### Pattern
|
||||||
|
|
||||||
|
- Acquire ownership with `POST /api/tasks/:id/checkout` using `{ agentId }`
|
||||||
|
- Release ownership with `POST /api/tasks/:id/release` using `{ agentId }`
|
||||||
|
- Admin override with `POST /api/tasks/:id/force-release`
|
||||||
|
- Read current lease state with `GET /api/tasks/:id/checkout`
|
||||||
|
|
||||||
|
`AgentStore` exposes matching methods:
|
||||||
|
- `checkoutTask(agentId, taskId)`
|
||||||
|
- `releaseTask(agentId, taskId)`
|
||||||
|
- `forceReleaseTask(taskId)`
|
||||||
|
- `getCheckedOutBy(taskId)`
|
||||||
|
|
||||||
|
### Conflict Semantics
|
||||||
|
|
||||||
|
- Checkout conflicts return **409 Conflict** when another agent already holds the lease
|
||||||
|
- Response shape: `{ error: "Task is already checked out", currentHolder, taskId }`
|
||||||
|
- Clients **must not retry 409 automatically** — this is ownership contention, not a transient failure
|
||||||
|
|
||||||
|
### Heartbeat Enforcement
|
||||||
|
|
||||||
|
`HeartbeatMonitor.executeHeartbeat()` validates checkout before work begins:
|
||||||
|
- If `task.checkedOutBy` is set to another agent, the run exits gracefully with `reason: "checkout_conflict"`
|
||||||
|
- Heartbeat execution is **read-only with respect to lease ownership** — it does not auto-checkout
|
||||||
|
- Scheduler/API callers are responsible for obtaining checkout before starting work
|
||||||
|
|
||||||
|
### API Reference
|
||||||
|
|
||||||
|
- `POST /api/tasks/:id/checkout` — Acquire lease (`{ agentId }`)
|
||||||
|
- `POST /api/tasks/:id/release` — Release lease (`{ agentId }`)
|
||||||
|
- `POST /api/tasks/:id/force-release` — Force release lease
|
||||||
|
- `GET /api/tasks/:id/checkout` — Read lease state (`{ checkedOutBy, checkedOutAt }`)
|
||||||
|
|
||||||
## Dashboard Task Creation
|
## Dashboard Task Creation
|
||||||
|
|
||||||
The dashboard provides two UI surfaces for creating tasks:
|
The dashboard provides two UI surfaces for creating tasks:
|
||||||
|
|||||||
@@ -51,7 +51,7 @@ describe("TaskStore task documents", () => {
|
|||||||
|
|
||||||
expect(tableNames.has("task_documents")).toBe(true);
|
expect(tableNames.has("task_documents")).toBe(true);
|
||||||
expect(tableNames.has("task_document_revisions")).toBe(true);
|
expect(tableNames.has("task_document_revisions")).toBe(true);
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
const index = db
|
const index = db
|
||||||
.prepare(
|
.prepare(
|
||||||
|
|||||||
@@ -12,12 +12,13 @@
|
|||||||
*/
|
*/
|
||||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
||||||
import { AgentStore } from "./agent-store.js";
|
import { AgentStore } from "./agent-store.js";
|
||||||
|
import { TaskStore } from "./store.js";
|
||||||
import { rm } from "node:fs/promises";
|
import { rm } from "node:fs/promises";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
import { mkdtempSync, existsSync, writeFileSync, readFileSync } from "node:fs";
|
import { mkdtempSync, existsSync, writeFileSync, readFileSync } from "node:fs";
|
||||||
import { tmpdir } from "node:os";
|
import { tmpdir } from "node:os";
|
||||||
import { createHash } from "node:crypto";
|
import { createHash } from "node:crypto";
|
||||||
import type { AgentCapability, AgentState } from "./types.js";
|
import { CheckoutConflictError, type AgentCapability, type AgentState } from "./types.js";
|
||||||
|
|
||||||
function makeTmpDir(): string {
|
function makeTmpDir(): string {
|
||||||
return mkdtempSync(join(tmpdir(), "kb-agent-store-test-"));
|
return mkdtempSync(join(tmpdir(), "kb-agent-store-test-"));
|
||||||
@@ -1385,6 +1386,111 @@ describe("AgentStore", () => {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("checkout leasing", () => {
|
||||||
|
let taskStore: TaskStore;
|
||||||
|
let holderId: string;
|
||||||
|
let otherAgentId: string;
|
||||||
|
let taskId: string;
|
||||||
|
|
||||||
|
beforeEach(async () => {
|
||||||
|
taskStore = new TaskStore(rootDir);
|
||||||
|
await taskStore.init();
|
||||||
|
|
||||||
|
store = new AgentStore({ rootDir, taskStore });
|
||||||
|
await store.init();
|
||||||
|
|
||||||
|
const holder = await store.createAgent({ name: "Checkout Holder", role: "executor" });
|
||||||
|
const other = await store.createAgent({ name: "Checkout Other", role: "executor" });
|
||||||
|
const task = await taskStore.createTask({ description: "Task for checkout leasing tests" });
|
||||||
|
|
||||||
|
holderId = holder.id;
|
||||||
|
otherAgentId = other.id;
|
||||||
|
taskId = task.id;
|
||||||
|
});
|
||||||
|
|
||||||
|
it("checkoutTask acquires a lease and stamps checkedOutAt", async () => {
|
||||||
|
const updated = await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
expect(updated.checkedOutBy).toBe(holderId);
|
||||||
|
expect(updated.checkedOutAt).toBeDefined();
|
||||||
|
|
||||||
|
const persisted = await taskStore.getTask(taskId);
|
||||||
|
expect(persisted?.checkedOutBy).toBe(holderId);
|
||||||
|
expect(persisted?.checkedOutAt).toBeDefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("checkoutTask is idempotent when the same agent re-checks out", async () => {
|
||||||
|
const first = await store.checkoutTask(holderId, taskId);
|
||||||
|
const second = await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
expect(second.checkedOutBy).toBe(holderId);
|
||||||
|
expect(second.checkedOutAt).toBe(first.checkedOutAt);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("checkoutTask throws CheckoutConflictError when already held by another agent", async () => {
|
||||||
|
await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
try {
|
||||||
|
await store.checkoutTask(otherAgentId, taskId);
|
||||||
|
throw new Error("Expected checkout conflict");
|
||||||
|
} catch (error) {
|
||||||
|
expect(error).toBeInstanceOf(CheckoutConflictError);
|
||||||
|
const conflict = error as CheckoutConflictError;
|
||||||
|
expect(conflict.taskId).toBe(taskId);
|
||||||
|
expect(conflict.currentHolderId).toBe(holderId);
|
||||||
|
expect(conflict.requestedById).toBe(otherAgentId);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
it("checkoutTask throws when agent is missing", async () => {
|
||||||
|
await expect(store.checkoutTask("agent-missing", taskId)).rejects.toThrow("Agent agent-missing not found");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("checkoutTask throws when task is missing", async () => {
|
||||||
|
await expect(store.checkoutTask(holderId, "FN-404")).rejects.toThrow("Task FN-404 not found");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("releaseTask clears checkedOutBy and checkedOutAt for the holder", async () => {
|
||||||
|
await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
const released = await store.releaseTask(holderId, taskId);
|
||||||
|
expect(released.checkedOutBy).toBeUndefined();
|
||||||
|
expect(released.checkedOutAt).toBeUndefined();
|
||||||
|
|
||||||
|
const persisted = await taskStore.getTask(taskId);
|
||||||
|
expect(persisted?.checkedOutBy).toBeUndefined();
|
||||||
|
expect(persisted?.checkedOutAt).toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("releaseTask throws for a non-holder agent", async () => {
|
||||||
|
await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
await expect(store.releaseTask(otherAgentId, taskId)).rejects.toThrow("Cannot release: not the checkout holder");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("releaseTask is idempotent when task is already released", async () => {
|
||||||
|
const released = await store.releaseTask(holderId, taskId);
|
||||||
|
|
||||||
|
expect(released.checkedOutBy).toBeUndefined();
|
||||||
|
expect(released.checkedOutAt).toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("forceReleaseTask clears checkout regardless of holder", async () => {
|
||||||
|
await store.checkoutTask(holderId, taskId);
|
||||||
|
|
||||||
|
const released = await store.forceReleaseTask(taskId);
|
||||||
|
expect(released.checkedOutBy).toBeUndefined();
|
||||||
|
expect(released.checkedOutAt).toBeUndefined();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("getCheckedOutBy returns holder ID when checked out and undefined otherwise", async () => {
|
||||||
|
expect(await store.getCheckedOutBy(taskId)).toBeUndefined();
|
||||||
|
|
||||||
|
await store.checkoutTask(holderId, taskId);
|
||||||
|
expect(await store.getCheckedOutBy(taskId)).toBe(holderId);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
// ── resetAgent ────────────────────────────────────────────────────
|
// ── resetAgent ────────────────────────────────────────────────────
|
||||||
|
|
||||||
describe("resetAgent", () => {
|
describe("resetAgent", () => {
|
||||||
|
|||||||
@@ -40,8 +40,10 @@ import type {
|
|||||||
AgentRating,
|
AgentRating,
|
||||||
AgentRatingSummary,
|
AgentRatingSummary,
|
||||||
AgentRatingInput,
|
AgentRatingInput,
|
||||||
|
Task,
|
||||||
} from "./types.js";
|
} from "./types.js";
|
||||||
import { AGENT_VALID_TRANSITIONS, agentToConfigSnapshot, diffConfigSnapshots } from "./types.js";
|
import { AGENT_VALID_TRANSITIONS, agentToConfigSnapshot, diffConfigSnapshots, CheckoutConflictError } from "./types.js";
|
||||||
|
import type { TaskStore } from "./store.js";
|
||||||
import { computeAccessState } from "./agent-permissions.js";
|
import { computeAccessState } from "./agent-permissions.js";
|
||||||
import { Database } from "./db.js";
|
import { Database } from "./db.js";
|
||||||
|
|
||||||
@@ -76,8 +78,10 @@ type TypedEventEmitter<Events extends Record<string, unknown[]>> = {
|
|||||||
|
|
||||||
/** Options for AgentStore constructor */
|
/** Options for AgentStore constructor */
|
||||||
export interface AgentStoreOptions {
|
export interface AgentStoreOptions {
|
||||||
/** Root directory for fn data (default: .fusion) */
|
/** Root directory for kb data (default: .fusion) */
|
||||||
rootDir?: string;
|
rootDir?: string;
|
||||||
|
/** Optional TaskStore for checkout/release operations */
|
||||||
|
taskStore?: TaskStore;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Agent data as stored on disk */
|
/** Agent data as stored on disk */
|
||||||
@@ -119,11 +123,13 @@ export class AgentStore extends EventEmitter {
|
|||||||
private agentsDir: string;
|
private agentsDir: string;
|
||||||
private locks: Map<string, AgentLock> = new Map();
|
private locks: Map<string, AgentLock> = new Map();
|
||||||
private _db: Database | null = null;
|
private _db: Database | null = null;
|
||||||
|
private taskStore?: TaskStore;
|
||||||
|
|
||||||
constructor(options: AgentStoreOptions = {}) {
|
constructor(options: AgentStoreOptions = {}) {
|
||||||
super();
|
super();
|
||||||
this.rootDir = options.rootDir ?? ".fusion";
|
this.rootDir = options.rootDir ?? ".fusion";
|
||||||
this.agentsDir = join(this.rootDir, "agents");
|
this.agentsDir = join(this.rootDir, "agents");
|
||||||
|
this.taskStore = options.taskStore;
|
||||||
}
|
}
|
||||||
|
|
||||||
private get db(): Database {
|
private get db(): Database {
|
||||||
@@ -813,6 +819,93 @@ export class AgentStore extends EventEmitter {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Acquire a checkout lease for a task.
|
||||||
|
* Throws CheckoutConflictError when another agent already holds the lease.
|
||||||
|
*/
|
||||||
|
async checkoutTask(agentId: string, taskId: string): Promise<Task> {
|
||||||
|
if (!this.taskStore) {
|
||||||
|
throw new Error("TaskStore not configured for checkout operations");
|
||||||
|
}
|
||||||
|
|
||||||
|
const agent = await this.getAgent(agentId);
|
||||||
|
if (!agent) {
|
||||||
|
throw new Error(`Agent ${agentId} not found`);
|
||||||
|
}
|
||||||
|
|
||||||
|
const task = await this.taskStore.getTask(taskId);
|
||||||
|
if (!task) {
|
||||||
|
throw new Error(`Task ${taskId} not found`);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (task.checkedOutBy && task.checkedOutBy !== agentId) {
|
||||||
|
throw new CheckoutConflictError(taskId, task.checkedOutBy, agentId);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (task.checkedOutBy === agentId) {
|
||||||
|
return task;
|
||||||
|
}
|
||||||
|
|
||||||
|
const updated = await this.taskStore.updateTask(taskId, { checkedOutBy: agentId });
|
||||||
|
await this.taskStore.logEntry(taskId, `Checked out by agent ${agentId}`);
|
||||||
|
return updated;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Release a checkout lease for a task.
|
||||||
|
*/
|
||||||
|
async releaseTask(agentId: string, taskId: string): Promise<Task> {
|
||||||
|
if (!this.taskStore) {
|
||||||
|
throw new Error("TaskStore not configured for checkout operations");
|
||||||
|
}
|
||||||
|
|
||||||
|
const task = await this.taskStore.getTask(taskId);
|
||||||
|
if (!task) {
|
||||||
|
throw new Error(`Task ${taskId} not found`);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (task.checkedOutBy && task.checkedOutBy !== agentId) {
|
||||||
|
throw new Error("Cannot release: not the checkout holder");
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!task.checkedOutBy) {
|
||||||
|
return task;
|
||||||
|
}
|
||||||
|
|
||||||
|
const updated = await this.taskStore.updateTask(taskId, { checkedOutBy: null });
|
||||||
|
await this.taskStore.logEntry(taskId, `Released by agent ${agentId}`);
|
||||||
|
return updated;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Force release a task checkout lease regardless of holder.
|
||||||
|
*/
|
||||||
|
async forceReleaseTask(taskId: string): Promise<Task> {
|
||||||
|
if (!this.taskStore) {
|
||||||
|
throw new Error("TaskStore not configured for checkout operations");
|
||||||
|
}
|
||||||
|
|
||||||
|
const updated = await this.taskStore.updateTask(taskId, { checkedOutBy: null });
|
||||||
|
await this.taskStore.logEntry(taskId, "Checkout force-released");
|
||||||
|
return updated;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Get the current checkout lease holder for a task.
|
||||||
|
*/
|
||||||
|
async getCheckedOutBy(taskId: string): Promise<string | undefined> {
|
||||||
|
if (!this.taskStore) {
|
||||||
|
throw new Error("TaskStore not configured for checkout operations");
|
||||||
|
}
|
||||||
|
|
||||||
|
const task = await this.taskStore.getTask(taskId);
|
||||||
|
if (!task) {
|
||||||
|
throw new Error(`Task ${taskId} not found`);
|
||||||
|
}
|
||||||
|
|
||||||
|
return task.checkedOutBy;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Reset budget token usage counters for an agent.
|
* Reset budget token usage counters for an agent.
|
||||||
* @param agentId - The agent ID
|
* @param agentId - The agent ID
|
||||||
|
|||||||
@@ -106,7 +106,7 @@ describe("Database", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("seeds schema version", () => {
|
it("seeds schema version", () => {
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("seeds lastModified", () => {
|
it("seeds lastModified", () => {
|
||||||
@@ -129,7 +129,7 @@ describe("Database", () => {
|
|||||||
|
|
||||||
it("is idempotent - calling init() twice does not fail", () => {
|
it("is idempotent - calling init() twice does not fail", () => {
|
||||||
expect(() => db.init()).not.toThrow();
|
expect(() => db.init()).not.toThrow();
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
});
|
});
|
||||||
|
|
||||||
it("does not overwrite existing config on re-init", () => {
|
it("does not overwrite existing config on re-init", () => {
|
||||||
@@ -736,7 +736,7 @@ describe("schema migrations", () => {
|
|||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
// Verify version bumped to 5 (includes v1→v2, v2→v3, v3→v4, and v4→v5 migrations)
|
// Verify version bumped to 5 (includes v1→v2, v2→v3, v3→v4, and v4→v5 migrations)
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
// Verify new columns exist and existing data is intact
|
// Verify new columns exist and existing data is intact
|
||||||
const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>;
|
const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>;
|
||||||
@@ -761,11 +761,11 @@ describe("schema migrations", () => {
|
|||||||
const db = new Database(kbDir);
|
const db = new Database(kbDir);
|
||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
// Re-init should not fail
|
// Re-init should not fail
|
||||||
db.init();
|
db.init();
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
db.close();
|
db.close();
|
||||||
});
|
});
|
||||||
@@ -781,7 +781,7 @@ describe("schema migrations", () => {
|
|||||||
|
|
||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'agentRatings'").all() as Array<{ name: string }>;
|
const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'agentRatings'").all() as Array<{ name: string }>;
|
||||||
expect(tables).toEqual([{ name: "agentRatings" }]);
|
expect(tables).toEqual([{ name: "agentRatings" }]);
|
||||||
@@ -805,7 +805,7 @@ describe("schema migrations", () => {
|
|||||||
|
|
||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'mission_events'").all() as Array<{ name: string }>;
|
const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'mission_events'").all() as Array<{ name: string }>;
|
||||||
expect(tables).toEqual([{ name: "mission_events" }]);
|
expect(tables).toEqual([{ name: "mission_events" }]);
|
||||||
@@ -909,7 +909,7 @@ describe("schema migrations", () => {
|
|||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
// Verify version bumped to 5
|
// Verify version bumped to 5
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
|
|
||||||
// Verify new columns exist and existing data is intact
|
// Verify new columns exist and existing data is intact
|
||||||
const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>;
|
const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>;
|
||||||
@@ -1119,7 +1119,7 @@ describe("createDatabase factory", () => {
|
|||||||
const db = createDatabase(kbDir);
|
const db = createDatabase(kbDir);
|
||||||
db.init();
|
db.init();
|
||||||
|
|
||||||
expect(db.getSchemaVersion()).toBe(19);
|
expect(db.getSchemaVersion()).toBe(20);
|
||||||
expect(db.getLastModified()).toBeGreaterThan(0);
|
expect(db.getLastModified()).toBeGreaterThan(0);
|
||||||
|
|
||||||
db.close();
|
db.close();
|
||||||
|
|||||||
@@ -59,7 +59,7 @@ export function fromJson<T>(json: string | null | undefined): T | undefined {
|
|||||||
|
|
||||||
// ── Schema Definition ────────────────────────────────────────────────
|
// ── Schema Definition ────────────────────────────────────────────────
|
||||||
|
|
||||||
const SCHEMA_VERSION = 19;
|
const SCHEMA_VERSION = 20;
|
||||||
|
|
||||||
function normalizeTaskComments(
|
function normalizeTaskComments(
|
||||||
steeringComments: SteeringComment[] | undefined,
|
steeringComments: SteeringComment[] | undefined,
|
||||||
@@ -742,6 +742,13 @@ export class Database {
|
|||||||
this.db.exec("CREATE INDEX IF NOT EXISTS idxAiSessionsLock ON ai_sessions(lockedByTab)");
|
this.db.exec("CREATE INDEX IF NOT EXISTS idxAiSessionsLock ON ai_sessions(lockedByTab)");
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (version < 20) {
|
||||||
|
this.applyMigration(20, () => {
|
||||||
|
this.addColumnIfMissing("tasks", "checkedOutBy", "TEXT");
|
||||||
|
this.addColumnIfMissing("tasks", "checkedOutAt", "TEXT");
|
||||||
|
});
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -1,5 +1,5 @@
|
|||||||
export { COLUMNS, COLUMN_LABELS, COLUMN_DESCRIPTIONS, VALID_TRANSITIONS, DEFAULT_SETTINGS, DEFAULT_GLOBAL_SETTINGS, DEFAULT_PROJECT_SETTINGS, GLOBAL_SETTINGS_KEYS, PROJECT_SETTINGS_KEYS, THINKING_LEVELS, THEME_MODES, COLOR_THEMES, WORKFLOW_STEP_TEMPLATES, AGENT_PERMISSIONS, agentToConfigSnapshot, diffConfigSnapshots } from "./types.js";
|
export { COLUMNS, COLUMN_LABELS, COLUMN_DESCRIPTIONS, VALID_TRANSITIONS, DEFAULT_SETTINGS, DEFAULT_GLOBAL_SETTINGS, DEFAULT_PROJECT_SETTINGS, GLOBAL_SETTINGS_KEYS, PROJECT_SETTINGS_KEYS, THINKING_LEVELS, THEME_MODES, COLOR_THEMES, WORKFLOW_STEP_TEMPLATES, AGENT_PERMISSIONS, agentToConfigSnapshot, diffConfigSnapshots, CheckoutConflictError } from "./types.js";
|
||||||
export type { Column, IssueInfo, IssueState, PrInfo, PrStatus, Task, TaskAttachment, TaskComment, TaskCommentInput, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskCreateInput, TaskDetail, AgentLogEntry, AgentLogType, AgentRole, BoardConfig, MergeDetails, MergeResult, Settings, GlobalSettings, ProjectSettings, SettingsScope, TaskStep, StepStatus, TaskLogEntry, ActivityLogEntry, ActivityEventType, ThinkingLevel, ThemeMode, ColorTheme, PlanningQuestion, PlanningSummary, PlanningResponse, PlanningQuestionType, ArchivedTaskEntry, BatchStatusRequest, BatchStatusResponse, BatchStatusEntry, BatchStatusResult, ModelPreset, WorkflowStep, WorkflowStepMode, WorkflowStepPhase, WorkflowStepInput, WorkflowStepResult, WorkflowStepTemplate, Agent, OrgTreeNode, AgentState, AgentDetail, AgentCreateInput, AgentUpdateInput, AgentApiKey, AgentApiKeyCreateResult, AgentCapability, AgentPromptTemplate, AgentPromptsConfig, AgentPermission, TaskAssignSource, AgentAccessState, AgentHeartbeatConfig, AgentBudgetConfig, AgentBudgetStatus, InstructionsBundleConfig, MessageResponseMode, AgentHeartbeatEvent, AgentHeartbeatRun, HeartbeatInvocationSource, AgentTaskSession, AgentRating, AgentRatingSummary, AgentRatingInput, AgentConfigSnapshot, RevisionFieldDiff, AgentConfigRevision, AgentStats, ReflectionTrigger, ReflectionMetrics, AgentReflection, AgentPerformanceSummary, NtfyNotificationEvent, SteeringComment, ParticipantType, MessageType, Message, MessageCreateInput, MessageFilter, Mailbox } from "./types.js";
|
export type { Column, IssueInfo, IssueState, PrInfo, PrStatus, Task, TaskAttachment, TaskComment, TaskCommentInput, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskCreateInput, TaskDetail, AgentLogEntry, AgentLogType, AgentRole, BoardConfig, MergeDetails, MergeResult, Settings, GlobalSettings, ProjectSettings, SettingsScope, TaskStep, StepStatus, TaskLogEntry, ActivityLogEntry, ActivityEventType, ThinkingLevel, ThemeMode, ColorTheme, PlanningQuestion, PlanningSummary, PlanningResponse, PlanningQuestionType, ArchivedTaskEntry, BatchStatusRequest, BatchStatusResponse, BatchStatusEntry, BatchStatusResult, ModelPreset, WorkflowStep, WorkflowStepMode, WorkflowStepPhase, WorkflowStepInput, WorkflowStepResult, WorkflowStepTemplate, Agent, OrgTreeNode, AgentState, AgentDetail, AgentCreateInput, AgentUpdateInput, AgentApiKey, AgentApiKeyCreateResult, AgentCapability, AgentPromptTemplate, AgentPromptsConfig, AgentPermission, TaskAssignSource, AgentAccessState, AgentHeartbeatConfig, AgentBudgetConfig, AgentBudgetStatus, InstructionsBundleConfig, MessageResponseMode, AgentHeartbeatEvent, AgentHeartbeatRun, HeartbeatInvocationSource, AgentTaskSession, AgentRating, AgentRatingSummary, AgentRatingInput, AgentConfigSnapshot, RevisionFieldDiff, AgentConfigRevision, AgentStats, ReflectionTrigger, ReflectionMetrics, AgentReflection, AgentPerformanceSummary, NtfyNotificationEvent, SteeringComment, ParticipantType, MessageType, Message, MessageCreateInput, MessageFilter, Mailbox, CheckoutLease } from "./types.js";
|
||||||
export { AGENT_VALID_TRANSITIONS } from "./types.js";
|
export { AGENT_VALID_TRANSITIONS } from "./types.js";
|
||||||
export {
|
export {
|
||||||
BUILTIN_AGENT_PROMPTS,
|
BUILTIN_AGENT_PROMPTS,
|
||||||
|
|||||||
@@ -230,6 +230,8 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
missionId: row.missionId || undefined,
|
missionId: row.missionId || undefined,
|
||||||
sliceId: row.sliceId || undefined,
|
sliceId: row.sliceId || undefined,
|
||||||
assignedAgentId: row.assignedAgentId || undefined,
|
assignedAgentId: row.assignedAgentId || undefined,
|
||||||
|
checkedOutBy: row.checkedOutBy || undefined,
|
||||||
|
checkedOutAt: row.checkedOutAt || undefined,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -279,10 +281,10 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
summary, thinkingLevel, createdAt, updatedAt, columnMovedAt,
|
summary, thinkingLevel, createdAt, updatedAt, columnMovedAt,
|
||||||
dependencies, steps, log, attachments, steeringComments,
|
dependencies, steps, log, attachments, steeringComments,
|
||||||
comments, workflowStepResults, prInfo, issueInfo, mergeDetails,
|
comments, workflowStepResults, prInfo, issueInfo, mergeDetails,
|
||||||
breakIntoSubtasks, enabledWorkflowSteps, modifiedFiles, missionId, sliceId, assignedAgentId
|
breakIntoSubtasks, enabledWorkflowSteps, modifiedFiles, missionId, sliceId, assignedAgentId, checkedOutBy, checkedOutAt
|
||||||
) VALUES (
|
) VALUES (
|
||||||
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?,
|
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?,
|
||||||
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
|
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
|
||||||
)
|
)
|
||||||
`).run(
|
`).run(
|
||||||
task.id,
|
task.id,
|
||||||
@@ -332,6 +334,8 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
task.missionId ?? null,
|
task.missionId ?? null,
|
||||||
task.sliceId ?? null,
|
task.sliceId ?? null,
|
||||||
task.assignedAgentId ?? null,
|
task.assignedAgentId ?? null,
|
||||||
|
task.checkedOutBy ?? null,
|
||||||
|
task.checkedOutAt ?? null,
|
||||||
);
|
);
|
||||||
this.db.bumpLastModified();
|
this.db.bumpLastModified();
|
||||||
}
|
}
|
||||||
@@ -1247,7 +1251,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
|
|
||||||
async updateTask(
|
async updateTask(
|
||||||
id: string,
|
id: string,
|
||||||
updates: { title?: string; description?: string; prompt?: string; worktree?: string | null; status?: string | null; dependencies?: string[]; blockedBy?: string | null; assignedAgentId?: string | null; paused?: boolean; baseBranch?: string | null; branch?: string | null; baseCommitSha?: string | null; size?: "S" | "M" | "L"; reviewLevel?: number; mergeRetries?: number; stuckKillCount?: number | null; recoveryRetryCount?: number | null; nextRecoveryAt?: string | null; enabledWorkflowSteps?: string[]; modelProvider?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningModelId?: string | null; thinkingLevel?: string | null; error?: string | null; summary?: string | null; sessionFile?: string | null; workflowStepResults?: import("./types.js").WorkflowStepResult[] | null; mergeDetails?: import("./types.js").MergeDetails | null; modifiedFiles?: string[] | null; missionId?: string | null; sliceId?: string | null },
|
updates: { title?: string; description?: string; prompt?: string; worktree?: string | null; status?: string | null; dependencies?: string[]; blockedBy?: string | null; assignedAgentId?: string | null; checkedOutBy?: string | null; checkedOutAt?: string | null; paused?: boolean; baseBranch?: string | null; branch?: string | null; baseCommitSha?: string | null; size?: "S" | "M" | "L"; reviewLevel?: number; mergeRetries?: number; stuckKillCount?: number | null; recoveryRetryCount?: number | null; nextRecoveryAt?: string | null; enabledWorkflowSteps?: string[]; modelProvider?: string | null; modelId?: string | null; validatorModelProvider?: string | null; validatorModelId?: string | null; planningModelProvider?: string | null; planningModelId?: string | null; thinkingLevel?: string | null; error?: string | null; summary?: string | null; sessionFile?: string | null; workflowStepResults?: import("./types.js").WorkflowStepResult[] | null; mergeDetails?: import("./types.js").MergeDetails | null; modifiedFiles?: string[] | null; missionId?: string | null; sliceId?: string | null },
|
||||||
): Promise<Task> {
|
): Promise<Task> {
|
||||||
return this.withTaskLock(id, async () => {
|
return this.withTaskLock(id, async () => {
|
||||||
// Validate that task doesn't depend on itself
|
// Validate that task doesn't depend on itself
|
||||||
@@ -1303,6 +1307,14 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
|
|||||||
} else if (updates.assignedAgentId !== undefined) {
|
} else if (updates.assignedAgentId !== undefined) {
|
||||||
task.assignedAgentId = updates.assignedAgentId;
|
task.assignedAgentId = updates.assignedAgentId;
|
||||||
}
|
}
|
||||||
|
if (updates.checkedOutBy === null) {
|
||||||
|
task.checkedOutBy = undefined;
|
||||||
|
task.checkedOutAt = undefined;
|
||||||
|
} else if (updates.checkedOutBy !== undefined) {
|
||||||
|
task.checkedOutBy = updates.checkedOutBy;
|
||||||
|
// Auto-set checkedOutAt when acquiring a lease (use provided value or generate timestamp)
|
||||||
|
task.checkedOutAt = updates.checkedOutAt ?? new Date().toISOString();
|
||||||
|
}
|
||||||
if (updates.paused !== undefined) task.paused = updates.paused || undefined;
|
if (updates.paused !== undefined) task.paused = updates.paused || undefined;
|
||||||
if (updates.baseBranch === null) {
|
if (updates.baseBranch === null) {
|
||||||
task.baseBranch = undefined;
|
task.baseBranch = undefined;
|
||||||
|
|||||||
@@ -540,6 +540,26 @@ export interface MergeDetails {
|
|||||||
autoResolvedCount?: number;
|
autoResolvedCount?: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Represents an agent's checkout lease on a task. */
|
||||||
|
export interface CheckoutLease {
|
||||||
|
/** The agent ID that holds the lease */
|
||||||
|
agentId: string;
|
||||||
|
/** ISO-8601 timestamp when the lease was acquired */
|
||||||
|
checkedOutAt: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Thrown when a checkout is attempted on a task already checked out by another agent. */
|
||||||
|
export class CheckoutConflictError extends Error {
|
||||||
|
constructor(
|
||||||
|
public readonly taskId: string,
|
||||||
|
public readonly currentHolderId: string,
|
||||||
|
public readonly requestedById: string,
|
||||||
|
) {
|
||||||
|
super(`Task ${taskId} is already checked out by agent ${currentHolderId}`);
|
||||||
|
this.name = "CheckoutConflictError";
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
export interface Task {
|
export interface Task {
|
||||||
id: string;
|
id: string;
|
||||||
title?: string;
|
title?: string;
|
||||||
@@ -639,6 +659,10 @@ export interface Task {
|
|||||||
thinkingLevel?: ThinkingLevel;
|
thinkingLevel?: ThinkingLevel;
|
||||||
/** Explicitly assigned agent ID for task-agent linking. Distinct from Agent.taskId active execution state. */
|
/** Explicitly assigned agent ID for task-agent linking. Distinct from Agent.taskId active execution state. */
|
||||||
assignedAgentId?: string;
|
assignedAgentId?: string;
|
||||||
|
/** Agent ID currently holding the checkout lease for this task. Undefined when no active lease. */
|
||||||
|
checkedOutBy?: string;
|
||||||
|
/** ISO-8601 timestamp when the checkout lease was acquired. */
|
||||||
|
checkedOutAt?: string;
|
||||||
/** Path to the persisted agent session file, enabling pause/resume without
|
/** Path to the persisted agent session file, enabling pause/resume without
|
||||||
* losing conversation context. Set when execution starts; cleared on
|
* losing conversation context. Set when execution starts; cleared on
|
||||||
* completion or terminal failure. */
|
* completion or terminal failure. */
|
||||||
|
|||||||
@@ -2202,6 +2202,230 @@ describe("PATCH /tasks/:id/assign and GET /agents/:id/tasks", () => {
|
|||||||
}, 30_000);
|
}, 30_000);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("Task checkout routes", () => {
|
||||||
|
let tempDir: string;
|
||||||
|
let fusionDir: string;
|
||||||
|
let store: TaskStore;
|
||||||
|
let agentAId: string;
|
||||||
|
let agentBId: string;
|
||||||
|
let taskState: TaskDetail;
|
||||||
|
|
||||||
|
beforeEach(async () => {
|
||||||
|
tempDir = mkdtempSync(join(tmpdir(), "kb-routes-task-checkout-"));
|
||||||
|
fusionDir = join(tempDir, ".fusion");
|
||||||
|
mkdirSync(fusionDir, { recursive: true });
|
||||||
|
|
||||||
|
const { AgentStore } = await import("@fusion/core");
|
||||||
|
const agentStore = new AgentStore({ rootDir: fusionDir });
|
||||||
|
await agentStore.init();
|
||||||
|
|
||||||
|
const agentA = await agentStore.createAgent({
|
||||||
|
name: "Checkout Agent A",
|
||||||
|
role: "executor",
|
||||||
|
});
|
||||||
|
const agentB = await agentStore.createAgent({
|
||||||
|
name: "Checkout Agent B",
|
||||||
|
role: "executor",
|
||||||
|
});
|
||||||
|
|
||||||
|
agentAId = agentA.id;
|
||||||
|
agentBId = agentB.id;
|
||||||
|
|
||||||
|
taskState = {
|
||||||
|
...FAKE_TASK_DETAIL,
|
||||||
|
id: "FN-300",
|
||||||
|
checkedOutBy: undefined,
|
||||||
|
checkedOutAt: undefined,
|
||||||
|
};
|
||||||
|
|
||||||
|
store = createMockStore({
|
||||||
|
getFusionDir: vi.fn().mockReturnValue(fusionDir),
|
||||||
|
getTask: vi.fn().mockImplementation(async (id: string) => {
|
||||||
|
if (id !== taskState.id) {
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
return { ...taskState };
|
||||||
|
}),
|
||||||
|
updateTask: vi.fn().mockImplementation(async (id: string, updates: { checkedOutBy?: string | null; checkedOutAt?: string | null }) => {
|
||||||
|
if (id !== taskState.id) {
|
||||||
|
throw new Error(`Task ${id} not found`);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (updates.checkedOutBy === null) {
|
||||||
|
taskState.checkedOutBy = undefined;
|
||||||
|
taskState.checkedOutAt = undefined;
|
||||||
|
} else if (updates.checkedOutBy !== undefined) {
|
||||||
|
taskState.checkedOutBy = updates.checkedOutBy;
|
||||||
|
taskState.checkedOutAt = updates.checkedOutAt ?? new Date().toISOString();
|
||||||
|
}
|
||||||
|
|
||||||
|
return { ...taskState };
|
||||||
|
}),
|
||||||
|
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||||
|
} as any);
|
||||||
|
}, 30_000);
|
||||||
|
|
||||||
|
afterEach(() => {
|
||||||
|
rmSync(tempDir, { recursive: true, force: true });
|
||||||
|
});
|
||||||
|
|
||||||
|
function buildApp() {
|
||||||
|
const app = express();
|
||||||
|
app.use(express.json());
|
||||||
|
app.use("/api", createApiRoutes(store));
|
||||||
|
return app;
|
||||||
|
}
|
||||||
|
|
||||||
|
it("POST /tasks/:id/checkout — returns 200 on success", async () => {
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
expect(res.body.checkedOutBy).toBe(agentAId);
|
||||||
|
expect(res.body.checkedOutAt).toBeTruthy();
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/checkout — returns 409 on conflict", async () => {
|
||||||
|
await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentBId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(409);
|
||||||
|
expect(res.body.error).toBe("Task is already checked out");
|
||||||
|
expect(res.body.currentHolder).toBe(agentAId);
|
||||||
|
expect(res.body.taskId).toBe(taskState.id);
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/checkout — returns 400 when agentId is missing", async () => {
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({}),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(400);
|
||||||
|
expect(res.body.error).toBe("agentId is required");
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/release — returns 200 on success", async () => {
|
||||||
|
await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/release`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
expect(res.body.checkedOutBy).toBeUndefined();
|
||||||
|
|
||||||
|
const statusRes = await GET(buildApp(), `/api/tasks/${taskState.id}/checkout`);
|
||||||
|
expect(statusRes.status).toBe(200);
|
||||||
|
expect(statusRes.body.checkedOutBy).toBeNull();
|
||||||
|
expect(statusRes.body.checkedOutAt).toBeNull();
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/release — returns 403 for wrong holder", async () => {
|
||||||
|
await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/release`,
|
||||||
|
JSON.stringify({ agentId: agentBId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(403);
|
||||||
|
expect(res.body.error).toBe("Not the checkout holder");
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/release — returns 400 when agentId is missing", async () => {
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/release`,
|
||||||
|
JSON.stringify({}),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(400);
|
||||||
|
expect(res.body.error).toBe("agentId is required");
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("GET /tasks/:id/checkout — returns checkout status", async () => {
|
||||||
|
await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const res = await GET(buildApp(), `/api/tasks/${taskState.id}/checkout`);
|
||||||
|
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
expect(res.body.checkedOutBy).toBe(agentAId);
|
||||||
|
expect(res.body.checkedOutAt).toBeTruthy();
|
||||||
|
}, 20_000);
|
||||||
|
|
||||||
|
it("POST /tasks/:id/force-release — clears active checkout", async () => {
|
||||||
|
await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/checkout`,
|
||||||
|
JSON.stringify({ agentId: agentAId }),
|
||||||
|
{ "Content-Type": "application/json" },
|
||||||
|
);
|
||||||
|
|
||||||
|
const res = await REQUEST(
|
||||||
|
buildApp(),
|
||||||
|
"POST",
|
||||||
|
`/api/tasks/${taskState.id}/force-release`,
|
||||||
|
);
|
||||||
|
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
|
||||||
|
const statusRes = await GET(buildApp(), `/api/tasks/${taskState.id}/checkout`);
|
||||||
|
expect(statusRes.status).toBe(200);
|
||||||
|
expect(statusRes.body.checkedOutBy).toBeNull();
|
||||||
|
expect(statusRes.body.checkedOutAt).toBeNull();
|
||||||
|
}, 20_000);
|
||||||
|
});
|
||||||
|
|
||||||
describe("Attachment routes", () => {
|
describe("Attachment routes", () => {
|
||||||
const FAKE_ATTACHMENT: TaskAttachment = {
|
const FAKE_ATTACHMENT: TaskAttachment = {
|
||||||
filename: "1234-screenshot.png",
|
filename: "1234-screenshot.png",
|
||||||
|
|||||||
@@ -2891,6 +2891,120 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
// Acquire checkout lease for a task
|
||||||
|
router.post("/tasks/:id/checkout", async (req, res) => {
|
||||||
|
try {
|
||||||
|
const { agentId } = req.body ?? {};
|
||||||
|
if (typeof agentId !== "string" || agentId.trim().length === 0) {
|
||||||
|
throw badRequest("agentId is required");
|
||||||
|
}
|
||||||
|
|
||||||
|
const scopedStore = await getScopedStore(req);
|
||||||
|
const { AgentStore } = await import("@fusion/core");
|
||||||
|
const agentStore = new AgentStore({
|
||||||
|
rootDir: scopedStore.getFusionDir(),
|
||||||
|
taskStore: scopedStore,
|
||||||
|
});
|
||||||
|
await agentStore.init();
|
||||||
|
|
||||||
|
const task = await agentStore.checkoutTask(agentId, req.params.id);
|
||||||
|
res.json(task);
|
||||||
|
} catch (err: any) {
|
||||||
|
if (err instanceof ApiError) {
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
|
if (err?.name === "CheckoutConflictError") {
|
||||||
|
res.status(409).json({
|
||||||
|
error: "Task is already checked out",
|
||||||
|
currentHolder: err.currentHolderId,
|
||||||
|
taskId: err.taskId,
|
||||||
|
});
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (err?.message?.includes("not found")) {
|
||||||
|
throw notFound(err.message);
|
||||||
|
}
|
||||||
|
rethrowAsApiError(err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Release checkout lease for a task
|
||||||
|
router.post("/tasks/:id/release", async (req, res) => {
|
||||||
|
try {
|
||||||
|
const { agentId } = req.body ?? {};
|
||||||
|
if (typeof agentId !== "string" || agentId.trim().length === 0) {
|
||||||
|
throw badRequest("agentId is required");
|
||||||
|
}
|
||||||
|
|
||||||
|
const scopedStore = await getScopedStore(req);
|
||||||
|
const { AgentStore } = await import("@fusion/core");
|
||||||
|
const agentStore = new AgentStore({
|
||||||
|
rootDir: scopedStore.getFusionDir(),
|
||||||
|
taskStore: scopedStore,
|
||||||
|
});
|
||||||
|
await agentStore.init();
|
||||||
|
|
||||||
|
const task = await agentStore.releaseTask(agentId, req.params.id);
|
||||||
|
res.json(task);
|
||||||
|
} catch (err: any) {
|
||||||
|
if (err instanceof ApiError) {
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
|
if (err?.message?.includes("not the checkout holder")) {
|
||||||
|
throw new ApiError(403, "Not the checkout holder");
|
||||||
|
}
|
||||||
|
if (err?.message?.includes("not found")) {
|
||||||
|
throw notFound(err.message);
|
||||||
|
}
|
||||||
|
rethrowAsApiError(err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Force release checkout lease for a task
|
||||||
|
router.post("/tasks/:id/force-release", async (req, res) => {
|
||||||
|
try {
|
||||||
|
const scopedStore = await getScopedStore(req);
|
||||||
|
const { AgentStore } = await import("@fusion/core");
|
||||||
|
const agentStore = new AgentStore({
|
||||||
|
rootDir: scopedStore.getFusionDir(),
|
||||||
|
taskStore: scopedStore,
|
||||||
|
});
|
||||||
|
await agentStore.init();
|
||||||
|
|
||||||
|
const task = await agentStore.forceReleaseTask(req.params.id);
|
||||||
|
res.json(task);
|
||||||
|
} catch (err: any) {
|
||||||
|
if (err instanceof ApiError) {
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
|
if (err?.message?.includes("not found")) {
|
||||||
|
throw notFound(err.message);
|
||||||
|
}
|
||||||
|
rethrowAsApiError(err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// Get checkout lease state for a task
|
||||||
|
router.get("/tasks/:id/checkout", async (req, res) => {
|
||||||
|
try {
|
||||||
|
const scopedStore = await getScopedStore(req);
|
||||||
|
const task = await scopedStore.getTask(req.params.id);
|
||||||
|
if (!task) {
|
||||||
|
throw notFound("Task not found");
|
||||||
|
}
|
||||||
|
|
||||||
|
res.json({
|
||||||
|
checkedOutBy: task.checkedOutBy ?? null,
|
||||||
|
checkedOutAt: task.checkedOutAt ?? null,
|
||||||
|
});
|
||||||
|
} catch (err: any) {
|
||||||
|
if (err instanceof ApiError) {
|
||||||
|
throw err;
|
||||||
|
}
|
||||||
|
rethrowAsApiError(err);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
// Delete task
|
// Delete task
|
||||||
router.delete("/tasks/:id", async (req, res) => {
|
router.delete("/tasks/:id", async (req, res) => {
|
||||||
try {
|
try {
|
||||||
|
|||||||
@@ -703,6 +703,25 @@ export class HeartbeatMonitor {
|
|||||||
return (await this.store.getRunDetail(agentId, run.id))!;
|
return (await this.store.getRunDetail(agentId, run.id))!;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Checkout enforcement: agent must hold the lease to work on this task.
|
||||||
|
// The heartbeat only validates existing checkout state — it does NOT attempt
|
||||||
|
// to acquire a checkout itself. The calling system (scheduler, API trigger)
|
||||||
|
// is responsible for checking out the task before the heartbeat starts.
|
||||||
|
if (taskDetail.checkedOutBy && taskDetail.checkedOutBy !== agentId) {
|
||||||
|
heartbeatLog.warn(
|
||||||
|
`Agent ${agentId} does not hold checkout for ${taskId} (held by ${taskDetail.checkedOutBy}) — graceful exit`
|
||||||
|
);
|
||||||
|
await this.completeRun(agentId, run.id, {
|
||||||
|
status: "completed",
|
||||||
|
resultJson: {
|
||||||
|
reason: "checkout_conflict",
|
||||||
|
taskId,
|
||||||
|
checkedOutBy: taskDetail.checkedOutBy,
|
||||||
|
},
|
||||||
|
});
|
||||||
|
return (await this.store.getRunDetail(agentId, run.id))!;
|
||||||
|
}
|
||||||
|
|
||||||
// Track usage via callbacks
|
// Track usage via callbacks
|
||||||
const STDOUT_EXCERPT_LIMIT = 4000;
|
const STDOUT_EXCERPT_LIMIT = 4000;
|
||||||
let outputLength = 0;
|
let outputLength = 0;
|
||||||
|
|||||||
Reference in New Issue
Block a user