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 @@ -79,6 +79,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;
}
// Update only the selected checkbox against the current row, preserving concurrent
// additions, renames, reordering and other milestones' completion state.
async setMilestoneDone<T>(owner: string, goalId: string, milestoneId: string, done: boolean) {
Expand Down
56 changes: 53 additions & 3 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,34 @@ const hash = (text: string) => createHash("sha256").update(text).digest("hex");
type MonitorPage = { id: string; hash?: string; timelessHash?: string; lines: string[] };
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;
readonly search: SearchService;
Expand Down Expand Up @@ -87,8 +115,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,7 @@ export class AgentService {
}
private async publishOutcome(owner: string, saved: AgentTask) {
const task = await this.getTask(owner, saved.id);
let outcomeReconciled = true;
if (task.status === "succeeded") {
await this.notify(
owner,
Expand All @@ -909,9 +942,13 @@ export class AgentService {
`task-done:${task.id}`,
);
if (task.goalId && !task.milestoneId) {
let milestoneReconciled = false;
for (let attempt = 0; attempt < 8; attempt++) {
const goal = await this.db.get<Goal>(owner, "goals", task.goalId);
if (!goal || goal.milestones.some((m) => m.id === task.id)) break;
if (!goal || goal.milestones.some((m) => m.id === task.id)) {
milestoneReconciled = true;
break;
}
if (
await this.db.compareAndSwap(
owner,
Expand All @@ -922,9 +959,12 @@ export class AgentService {
milestones: [...goal.milestones, { id: task.id, title: task.title, done: true }],
},
)
)
) {
milestoneReconciled = true;
break;
}
}
outcomeReconciled = milestoneReconciled;
}
} else if (task.status === "failed") {
await this.notify(
Expand Down Expand Up @@ -965,6 +1005,16 @@ export class AgentService {
{ status: "active", error: task.error },
{ status: "paused" },
);
const key = outcomeKey(task);
if (key && outcomeReconciled)
await this.db.compareAndSetJsonPath<AgentTask>(
Comment thread
kvnloo marked this conversation as resolved.
owner,
"tasks",
task.id,
outcomeMarkerExpected(task),
["state", "publishedOutcome"],
key,
);
}
private async document(
owner: string,
Expand Down
240 changes: 240 additions & 0 deletions tests/maintain-outcome.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,240 @@
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, Goal } 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();
}
});

test("outcome marker waits for goal milestone reconciliation after contention", async () => {
const db = await createStore();
const originalCompareAndSwap = db.compareAndSwap.bind(db);
try {
const goal: Goal = {
id: "goal1",
title: "Goal",
description: "",
category: "Personal",
status: "active",
milestones: [],
createdAt: new Date().toISOString(),
};
await db.put("owner", "goals", goal);
await db.put("owner", "tasks", task("task1", "succeeded", { result: "done", goalId: goal.id }));
let rejectedGoalWrites = 0;
db.compareAndSwap = (async (
owner: string,
kind: string,
id: string,
expected: Record<string, unknown>,
patch: Record<string, unknown>,
) => {
if (kind === "goals" && rejectedGoalWrites < 8) {
rejectedGoalWrites++;
return null;
}
return originalCompareAndSwap(owner, kind, id, expected, patch);
}) as Store["compareAndSwap"];
const service = serviceWith(db);

await service.maintain();
assert.equal(rejectedGoalWrites, 8);
assert.equal((await db.get<Goal>("owner", "goals", goal.id))?.milestones.length, 0);
assert.equal(
(await db.get<AgentTask>("owner", "tasks", "task1"))?.state.publishedOutcome,
undefined,
);
assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 1);

await service.maintain();
assert.deepEqual((await db.get<Goal>("owner", "goals", goal.id))?.milestones, [
{ id: "task1", title: "Task task1", done: true },
]);
assert.equal(
(await db.get<AgentTask>("owner", "tasks", "task1"))?.state.publishedOutcome,
"task-done:task1",
);
assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 1);
} finally {
db.compareAndSwap = originalCompareAndSwap as Store["compareAndSwap"];
await db.close();
}
});

test("linked milestone completion stays manual while the outcome is marked published", async () => {
const db = await createStore();
try {
const goal: Goal = {
id: "goal1",
title: "Goal",
description: "",
category: "Personal",
status: "active",
milestones: [{ id: "manual", title: "Manual milestone", done: false }],
createdAt: new Date().toISOString(),
};
await db.put("owner", "goals", goal);
await db.put(
"owner",
"tasks",
task("task1", "succeeded", {
result: "done",
goalId: goal.id,
milestoneId: "manual",
}),
);
const service = serviceWith(db);

await service.maintain();

assert.deepEqual((await db.get<Goal>("owner", "goals", goal.id))?.milestones, [
{ id: "manual", title: "Manual milestone", done: false },
]);
assert.equal(
(await db.get<AgentTask>("owner", "tasks", "task1"))?.state.publishedOutcome,
"task-done:task1",
);
assert.equal((await notifications(db)).filter((n) => n.taskId === "task1").length, 1);
} finally {
await db.close();
}
});
Loading