diff --git a/apps/server/src/db.ts b/apps/server/src/db.ts index 29075322..6fd23d2c 100644 --- a/apps/server/src/db.ts +++ b/apps/server/src/db.ts @@ -30,6 +30,17 @@ export class Store { ); return result.rows.map((row) => row.data as T); } + async listByStatus>( + owner: string, + kind: string, + status: string, + ): Promise { + 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(owner: string, kind: string, value: T): Promise { 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()", diff --git a/apps/server/src/engine/service.ts b/apps/server/src/engine/service.ts index 3f2af061..8aeb276e 100644 --- a/apps/server/src/engine/service.ts +++ b/apps/server/src/engine/service.ts @@ -524,8 +524,8 @@ 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(owner, "tasks")) - .filter((task) => task.status === "succeeded" && typeof task.input.messageId === "string") + (await this.db.listByStatus(owner, "tasks", "succeeded")) + .filter((task) => typeof task.input.messageId === "string") .map((task) => `${task.kind}:${task.input.messageId}`), ); const obsolete = (kind: AgentTask["kind"], messageId: unknown) => @@ -533,8 +533,8 @@ export class AgentService { (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(owner, "ideas")) - if (idea.status === "new" && obsolete(idea.kind, idea.input.messageId)) + for (const idea of await this.db.listByStatus(owner, "ideas", "new")) + if (obsolete(idea.kind, idea.input.messageId)) await this.db.compareAndSwap( owner, "ideas", @@ -583,8 +583,8 @@ export class AgentService { createdAt: date(), } satisfies Idea); } - for (const goal of await this.db.list(owner, "goals")) - if (goal.status === "active" && !goal.milestones.length) { + for (const goal of await this.db.listByStatus(owner, "goals", "active")) + if (!goal.milestones.length) { const id = hash(`goal:${goal.id}:${goal.description}`); await this.db.insertIfAbsent(owner, "ideas", { id, diff --git a/tests/persistence.test.ts b/tests/persistence.test.ts index e371655c..5535c3aa 100644 --- a/tests/persistence.test.ts +++ b/tests/persistence.test.ts @@ -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(); + } +});