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
11 changes: 11 additions & 0 deletions apps/server/src/db.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,17 @@ export class Store {
);
return result.rows.map((row) => row.data as T);
}
async listByStatus<T = Record<string, unknown>>(
owner: string,
kind: string,
status: string,
): Promise<T[]> {
const result = await this.db.query(
"SELECT data FROM records WHERE owner=$1 AND kind=$2 AND data->>'status'=$3 ORDER BY updated_at DESC,id",
[owner, kind, status],
);
return result.rows.map((row) => row.data as T);
}
async put<T extends { id: string }>(owner: string, kind: string, value: T): Promise<T> {
await this.db.query(
"INSERT INTO records(owner,kind,id,data) VALUES($1,$2,$3,$4::jsonb) ON CONFLICT(owner,kind,id) DO UPDATE SET data=excluded.data,updated_at=now()",
Expand Down
12 changes: 6 additions & 6 deletions apps/server/src/engine/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -524,17 +524,17 @@ export class AgentService {
w.mail.filter((mail) => /^Sent\b/i.test(mail.label)).map((mail) => mail.id),
);
const completedSources = new Set(
(await this.db.list<AgentTask>(owner, "tasks"))
.filter((task) => task.status === "succeeded" && typeof task.input.messageId === "string")
(await this.db.listByStatus<AgentTask>(owner, "tasks", "succeeded"))
.filter((task) => typeof task.input.messageId === "string")
.map((task) => `${task.kind}:${task.input.messageId}`),
);
const obsolete = (kind: AgentTask["kind"], messageId: unknown) =>
typeof messageId === "string" &&
(sentIds.has(messageId) || completedSources.has(`${kind}:${messageId}`));
// Retire earlier suggestions as well as preventing new duplicates. A concurrent
// acceptance wins its own compare-and-swap and is never overwritten here.
for (const idea of await this.db.list<Idea>(owner, "ideas"))
if (idea.status === "new" && obsolete(idea.kind, idea.input.messageId))
for (const idea of await this.db.listByStatus<Idea>(owner, "ideas", "new"))
if (obsolete(idea.kind, idea.input.messageId))
await this.db.compareAndSwap(
owner,
"ideas",
Expand Down Expand Up @@ -583,8 +583,8 @@ export class AgentService {
createdAt: date(),
} satisfies Idea);
}
for (const goal of await this.db.list<Goal>(owner, "goals"))
if (goal.status === "active" && !goal.milestones.length) {
for (const goal of await this.db.listByStatus<Goal>(owner, "goals", "active"))
if (!goal.milestones.length) {
const id = hash(`goal:${goal.id}:${goal.description}`);
await this.db.insertIfAbsent(owner, "ideas", {
id,
Expand Down
30 changes: 30 additions & 0 deletions tests/persistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,3 +31,33 @@ test("idle Postgres client errors are logged instead of crashing the process", a
await pool.end();
}
});
test("listByStatus matches list plus status filtering", async () => {
const store = await createStore();
try {
for (const [id, status] of [
["i0", "new"],
["i1", "dismissed"],
["i2", "new"],
["i3", "accepted"],
] as const)
await store.put("owner", "ideas", { id, status });
await store.put("owner", "ideas", { id: "i-none" });

const scoped = await store.listByStatus<{ id: string; status?: string }>(
"owner",
"ideas",
"new",
);
const expected = (await store.list<{ id: string; status?: string }>("owner", "ideas")).filter(
(item) => item.status === "new",
);
assert.deepEqual(scoped, expected);

await store.put("other", "ideas", { id: "x", status: "new" });
await store.put("owner", "goals", { id: "g", status: "new" });
assert.equal((await store.listByStatus("owner", "ideas", "new")).length, 2);
assert.equal((await store.listByStatus("nobody", "ideas", "new")).length, 0);
} finally {
await store.close();
}
});
Loading