diff --git a/apps/server/src/db.ts b/apps/server/src/db.ts index 29075322..98f37efc 100644 --- a/apps/server/src/db.ts +++ b/apps/server/src/db.ts @@ -57,6 +57,21 @@ export class Store { ); return (result.rows[0]?.data as T | undefined) ?? null; } + async compareAndSetJsonPath( + owner: string, + kind: string, + id: string, + expected: Record, + path: string[], + value: unknown, + ): Promise { + if (path.length === 0) throw new Error("compareAndSetJsonPath requires a non-empty path"); + const result = await this.db.query( + "UPDATE records SET data=jsonb_set(data,$5::text[],$6::jsonb,true),updated_at=now() WHERE owner=$1 AND kind=$2 AND id=$3 AND data @> $4::jsonb RETURNING data", + [owner, kind, id, JSON.stringify(expected), path, JSON.stringify(value)], + ); + return (result.rows[0]?.data as T | undefined) ?? null; + } async insertIfAbsent( owner: string, kind: string, diff --git a/apps/server/src/engine/service.ts b/apps/server/src/engine/service.ts index 3f2af061..9b1493db 100644 --- a/apps/server/src/engine/service.ts +++ b/apps/server/src/engine/service.ts @@ -39,6 +39,34 @@ import { LostLeaseError, type TaskContext, TaskWorker } from "./worker.ts"; const hash = (text: string) => createHash("sha256").update(text).digest("hex"); const date = () => new Date().toISOString(); const terminal = new Set(["succeeded", "failed", "cancelled"]); +const outcomeKey = (task: AgentTask): string | undefined => { + if (task.status === "succeeded") return `task-done:${task.id}`; + if (task.status === "failed") return `task-error:${task.id}:${task.attempts}`; + if (task.status === "waiting_input") return `input:${task.id}:${hash(task.question ?? "")}`; + if (task.status === "waiting_approval") return `review:${task.actionId}`; + if (task.status === "cancelled") return `cancelled:${task.id}`; + const notice = z + .object({ title: z.string(), body: z.string(), key: z.string() }) + .safeParse(task.state.notice); + if ((task.status === "scheduled" || (task.status === "paused" && task.error)) && notice.success) + return notice.data.key; + return undefined; +}; + +const outcomeMarkerExpected = (task: AgentTask): Record => { + const expected: Record = { status: task.status }; + if (task.status === "failed") expected.attempts = task.attempts; + if (task.status === "waiting_input" && task.question !== undefined) + expected.question = task.question; + if (task.status === "waiting_approval" && task.actionId !== undefined) + expected.actionId = task.actionId; + const notice = z + .object({ title: z.string(), body: z.string(), key: z.string() }) + .safeParse(task.state.notice); + if ((task.status === "scheduled" || task.status === "paused") && notice.success) + expected.state = { notice: notice.data }; + return expected; +}; export class AgentService { readonly worker: TaskWorker; private maintenance?: ReturnType; @@ -75,8 +103,12 @@ export class AgentService { this.refreshing = true; try { // Recover publications if the process exited after committing an outcome. - for (const { owner, value } of await this.db.scan("tasks")) + // A matching durable marker means this exact outcome was already published. + for (const { owner, value } of await this.db.scan("tasks")) { + const key = outcomeKey(value); + if (key && value.state.publishedOutcome === key) continue; await this.publishOutcome(owner, value); + } for (const { owner, value } of await this.db.scan("monitors")) await this.activateMonitor(owner, value); for (const { owner, value } of await this.db.scan("ideas")) @@ -900,6 +932,16 @@ export class AgentService { { status: "active", error: task.error }, { status: "paused" }, ); + const key = outcomeKey(task); + if (key) + await this.db.compareAndSetJsonPath( + owner, + "tasks", + task.id, + outcomeMarkerExpected(task), + ["state", "publishedOutcome"], + key, + ); } private async document( owner: string, diff --git a/tests/maintain-outcome.test.ts b/tests/maintain-outcome.test.ts new file mode 100644 index 00000000..f0d01119 --- /dev/null +++ b/tests/maintain-outcome.test.ts @@ -0,0 +1,146 @@ +import assert from "node:assert/strict"; +import { createHash } from "node:crypto"; +import { test } from "node:test"; +import { createStore, type Store } from "../apps/server/src/db.ts"; +import { AgentService } from "../apps/server/src/engine/service.ts"; +import type { AgentNotification, AgentTask } from "../packages/domain/src/agent.ts"; + +function task(id: string, status: AgentTask["status"], extra: Partial = {}): AgentTask { + return { + id, + title: `Task ${id}`, + prompt: `Task ${id}`, + kind: "agent", + status, + plan: [], + evidence: [], + input: {}, + state: {}, + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + attempts: status === "failed" ? 3 : 1, + error: status === "failed" ? "simulated failure" : null, + leaseId: null, + leaseUntil: null, + artifactIds: [], + ...extra, + }; +} + +function serviceWith(db: Store) { + const service = new AgentService( + db, + {} as never, + {} as never, + {} as never, + {} as never, + {} as never, + ); + return service as unknown as { maintain(): Promise }; +} + +async function notifications(db: Store): Promise { + return (await db.scan("notifications")).map(({ value }) => value); +} + +test("maintenance recovers a lost publication once, then skips the marked row", async () => { + const db = await createStore(); + try { + await db.put("owner", "tasks", task("task1", "succeeded", { result: "done" })); + const service = serviceWith(db); + + await service.maintain(); + assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 1); + assert.equal( + (await db.get("owner", "tasks", "task1"))?.state.publishedOutcome, + "task-done:task1", + ); + + let taskGets = 0; + const originalGet = db.get.bind(db); + db.get = (async (owner: string, kind: string, id: string) => { + if (kind === "tasks") taskGets++; + return originalGet(owner, kind, id); + }) as Store["get"]; + await service.maintain(); + db.get = originalGet; + assert.equal(taskGets, 0); + assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 1); + } finally { + await db.close(); + } +}); + +test("outcome marker preserves a concurrent sibling state write", async () => { + const db = await createStore(); + try { + await db.put( + "owner", + "tasks", + task("task1", "succeeded", { result: "done", state: { existing: "value" } }), + ); + const service = serviceWith(db); + const originalInsert = db.insertIfAbsent.bind(db); + let injected = false; + db.insertIfAbsent = (async (owner: string, kind: string, value: { id: string }) => { + if (kind === "notifications" && !injected) { + injected = true; + const latest = await db.get("owner", "tasks", "task1"); + assert.ok(latest); + await db.compareAndSwap( + "owner", + "tasks", + "task1", + { status: "succeeded" }, + { state: { ...latest.state, concurrent: "kept" } }, + ); + } + return originalInsert(owner, kind, value); + }) as Store["insertIfAbsent"]; + + await service.maintain(); + const saved = await db.get("owner", "tasks", "task1"); + assert.equal(saved?.state.existing, "value"); + assert.equal(saved?.state.concurrent, "kept"); + assert.equal(saved?.state.publishedOutcome, "task-done:task1"); + } finally { + await db.close(); + } +}); + +test("a same-status outcome change is not hidden by a stale marker", async () => { + const db = await createStore(); + try { + await db.put("owner", "tasks", task("task1", "waiting_input", { question: "Old question?" })); + const service = serviceWith(db); + const originalInsert = db.insertIfAbsent.bind(db); + let changed = false; + db.insertIfAbsent = (async (owner: string, kind: string, value: { id: string }) => { + if (kind === "notifications" && !changed) { + changed = true; + await db.compareAndSwap( + "owner", + "tasks", + "task1", + { status: "waiting_input", question: "Old question?" }, + { question: "New question?" }, + ); + } + return originalInsert(owner, kind, value); + }) as Store["insertIfAbsent"]; + + await service.maintain(); + let saved = await db.get("owner", "tasks", "task1"); + assert.equal(saved?.question, "New question?"); + assert.equal(saved?.state.publishedOutcome, undefined); + + db.insertIfAbsent = originalInsert as Store["insertIfAbsent"]; + await service.maintain(); + saved = await db.get("owner", "tasks", "task1"); + const expected = `input:task1:${createHash("sha256").update("New question?").digest("hex")}`; + assert.equal(saved?.state.publishedOutcome, expected); + assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 2); + } finally { + await db.close(); + } +});