Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions apps/server/src/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,21 @@ export class Store {
);
return (result.rows[0]?.data as T | undefined) ?? null;
}
async compareAndSetJsonPath<T>(
owner: string,
kind: string,
id: string,
expected: Record<string, unknown>,
path: string[],
value: unknown,
): Promise<T | null> {
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<T extends { id: string }>(
owner: string,
kind: string,
Expand Down
44 changes: 43 additions & 1 deletion apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> => {
const expected: Record<string, unknown> = { 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<typeof setInterval>;
Expand Down Expand Up @@ -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<AgentTask>("tasks"))
// A matching durable marker means this exact outcome was already published.
for (const { owner, value } of await this.db.scan<AgentTask>("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<Monitor>("monitors"))
await this.activateMonitor(owner, value);
for (const { owner, value } of await this.db.scan<Idea>("ideas"))
Expand Down Expand Up @@ -900,6 +932,16 @@ export class AgentService {
{ status: "active", error: task.error },
{ status: "paused" },
);
const key = outcomeKey(task);
if (key)
await this.db.compareAndSetJsonPath<AgentTask>(
owner,
"tasks",
task.id,
outcomeMarkerExpected(task),
["state", "publishedOutcome"],
key,
);
}
private async document(
owner: string,
Expand Down
146 changes: 146 additions & 0 deletions tests/maintain-outcome.test.ts
Original file line number Diff line number Diff line change
@@ -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> = {}): 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<void> };
}

async function notifications(db: Store): Promise<AgentNotification[]> {
return (await db.scan<AgentNotification>("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<AgentTask>("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<AgentTask>("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<AgentTask>("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<AgentTask>("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<AgentTask>("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();
}
});
Loading