From 7107c6d260a25ef1903c7c3277e9152d09e6267b Mon Sep 17 00:00:00 2001 From: Maggie Appleton <5599295+MaggieAppleton@users.noreply.github.com> Date: Sat, 3 Oct 2026 08:02:38 +0100 Subject: [PATCH] Place and recover inline research requests --- ...accepted-research-storage.test-fixtures.ts | 55 +++ apps/server/src/main.ts | 136 ++++++-- apps/server/src/research/inline-1.test.ts | 136 ++++++++ apps/server/src/research/inline-2.test.ts | 135 ++++++++ apps/server/src/research/inline-3.test.ts | 139 ++++++++ apps/server/src/research/inline-4.test.ts | 64 ++++ .../server/src/research/inline-memory.test.ts | 90 +++++ .../src/research/inline.test-fixtures.ts | 77 +++++ .../research/placement-close.test-fixtures.ts | 164 +++++++++ .../src/research/placement-close.test.ts | 22 ++ apps/server/src/research/placement.test.ts | 57 ++++ apps/server/src/research/placement.ts | 30 ++ apps/server/src/research/service.test.ts | 33 +- apps/server/src/research/service.ts | 284 +++++++++++++--- apps/server/src/research/start-identity.ts | 49 +++ .../src/storage/postgres/startup.test.ts | 312 +++++++++++++++++- 16 files changed, 1682 insertions(+), 101 deletions(-) create mode 100644 apps/server/src/conversation-plan/accepted-research-storage.test-fixtures.ts create mode 100644 apps/server/src/research/inline-1.test.ts create mode 100644 apps/server/src/research/inline-2.test.ts create mode 100644 apps/server/src/research/inline-3.test.ts create mode 100644 apps/server/src/research/inline-4.test.ts create mode 100644 apps/server/src/research/inline-memory.test.ts create mode 100644 apps/server/src/research/inline.test-fixtures.ts create mode 100644 apps/server/src/research/placement-close.test-fixtures.ts create mode 100644 apps/server/src/research/placement-close.test.ts create mode 100644 apps/server/src/research/placement.test.ts create mode 100644 apps/server/src/research/placement.ts create mode 100644 apps/server/src/research/start-identity.ts diff --git a/apps/server/src/conversation-plan/accepted-research-storage.test-fixtures.ts b/apps/server/src/conversation-plan/accepted-research-storage.test-fixtures.ts new file mode 100644 index 00000000..3e106ae3 --- /dev/null +++ b/apps/server/src/conversation-plan/accepted-research-storage.test-fixtures.ts @@ -0,0 +1,55 @@ +import { createHash } from "node:crypto"; +import { broadcast } from "../wire"; +import { JobRegistry } from "../jobs/registry"; +import { JobService } from "../jobs/service"; +import { researchAnswerDefinition, researchEvidenceDefinition } from "../jobs/research-workspace"; +import { ResearchWorkspaceService } from "../research/service"; +import { REPORT } from "../research/test-support"; +import * as Plan from "../plan/service"; +import type { headingMemory } from "../chat/job-heading-memory.test-fixtures"; + +async function evidence() { + return { findings: [], sources: [] }; +} +async function privateEvidence() { + return { findings: [] }; +} +async function report() { + return REPORT; +} +async function answer() { + return { text: "Offline answer", sourceUrls: [] }; +} + +export function researchStorage(h: Awaited>) { + let jobs = new JobService({ + storage: h.opened.storage, + registry: new JobRegistry([ + researchEvidenceDefinition({ config: { agent: true, model: "offline" }, engine: evidence }), + researchAnswerDefinition({ + config: { agent: true, model: "offline" }, + engines: { private: privateEvidence, synthesize: report, answer }, + }), + ]), + lease: () => h.opened.lease, + }); + let research = new ResearchWorkspaceService({ + storage: h.opened.storage, + jobs, + lease: () => h.opened.lease, + current: async () => ({ + channelId: h.plan.id, + revision: h.plan.revision, + source: Plan.source(h.plan), + sourceHash: `sha256:${createHash("sha256").update(Plan.source(h.plan)).digest("hex")}`, + }), + publish: (channelId, workspaceId, revision) => + broadcast(h.opened.server, channelId, { + kind: "research:changed", + ts: 0, + workspaceId, + revision, + }), + }); + return { research, jobs }; +} diff --git a/apps/server/src/main.ts b/apps/server/src/main.ts index f57b0bc9..e83433e7 100644 --- a/apps/server/src/main.ts +++ b/apps/server/src/main.ts @@ -37,7 +37,8 @@ import * as Inject from "./questions/inject"; import * as Marks from "./comments/inject"; import * as Questions from "./questions/service"; import { registerResearchWorkspaceRoutes } from "./research/routes"; -import { ResearchWorkspaceService } from "./research/service"; +import { ResearchWorkspaceError, ResearchWorkspaceService } from "./research/service"; +import { placeResearchReference as placeResearch } from "./research/placement"; import * as Rooms from "./rooms"; import { admit } from "./socket/admission"; import { StorageError } from "./storage/errors"; @@ -70,6 +71,7 @@ const LEASE_RENEW_MS = 10_000; const LEASE_SAFETY_MS = 5_000; const SESSION_CLEANUP_MS = 5 * 60_000; const ACCESS_RECHECK_MS = 60_000; +const RESEARCH_RECOVERY_RETRY_MS = 10_000; let server: Server; let heldLease: Lease | undefined; @@ -81,6 +83,8 @@ let cleaningSessions: Promise | undefined; let ownerBindings: ActiveOwnerBindings | undefined; let jobRunner: JobRunner | undefined; let researchService: ResearchWorkspaceService | undefined; +let researchRecoveryTimer: ReturnType | undefined; +let recoveringResearch: Promise | undefined; let referenceService: ReferenceService | undefined; let summaryCoordinator: DocumentSummaryCoordinator | undefined; let descriptionProjector: DocumentDescriptionProjector | undefined; @@ -119,17 +123,8 @@ function presence(server: Server, room: Rooms.Room): void { }); } -/** - * Attach the document to a room, once. - * - * Two clients opening at the same moment must not build two documents, so the - * first stores its promise and the second waits on it. - */ -async function plan(room: Rooms.Room, server: Server): Promise { - if (room.closing) await room.closing; - if (deletingChannels.has(room.id)) throw new Error("document is unavailable"); - if (room.plan) return room.plan; - let backend: Service.Backend = { +function documentBackend(): Service.Backend { + return { storage, lease: () => { if (!heldLease) throw new Error("storage writer lease is unavailable"); @@ -141,6 +136,56 @@ async function plan(room: Rooms.Room, server: Server): Promise summaryCoordinator?.schedule(target), }; +} + +function placeResearchReference( + channelId: string, + workspaceId: string, +): Promise<"placed" | "deferred"> { + return placeResearch(channelId, workspaceId, { + get: Rooms.get, + exclusive: withDocumentLock, + detached: async (id, action) => { + let detached = await Service.open(id, documentBackend(), server); + try { + return await action(detached); + } finally { + await Service.close(detached); + } + }, + }); +} + +function scheduleResearchRecovery(deferred: number): void { + if (deferred === 0 || researchRecoveryTimer || draining) return; + researchRecoveryTimer = setTimeout(() => { + researchRecoveryTimer = undefined; + if (draining) return; + recoveringResearch = researchService!.recoverPendingPlannerInline(placeResearchReference).then( + result => { + scheduleResearchRecovery(result.deferred); + }, + err => { + console.error("chopin: research reference recovery failed -", err); + signal(); + }, + ).finally(() => { + recoveringResearch = undefined; + }); + }, RESEARCH_RECOVERY_RETRY_MS); +} + +/** + * Attach the document to a room, once. + * + * Two clients opening at the same moment must not build two documents, so the + * first stores its promise and the second waits on it. + */ +async function plan(room: Rooms.Room, server: Server): Promise { + if (room.closing) await room.closing; + if (deletingChannels.has(room.id)) throw new Error("document is unavailable"); + if (room.plan) return room.plan; + let backend = documentBackend(); let opening = room.opening ??= withDocumentLock(room.id, async () => { if (deletingChannels.has(room.id)) throw new Error("document is unavailable"); if (room.plan) return room.plan; @@ -186,14 +231,23 @@ function chat(room: Rooms.Room, ws: Socket): Chat.Room { ? async request => { let service = researchService; if (!service) throw new Error("research workspaces are unavailable"); - let created = await service.startPlanner({ - channelId: room.id, - question: request.question, - originMessageId: request.entryId, - requestedBy: request.userId, - requestedByHandle: request.handle, - beforeStart: () => jobRunner?.ownerAvailable(room.id) ?? Promise.resolve(), - }); + let created; + try { + created = await service.startPlannerInline({ + channelId: room.id, + question: request.question, + originMessageId: request.entryId, + requestedBy: request.userId, + requestedByHandle: request.handle, + beforeStart: () => jobRunner?.ownerAvailable(room.id) ?? Promise.resolve(), + placeReference: id => placeResearchReference(room.id, id), + }); + } catch (err) { + if (err instanceof ResearchWorkspaceError && err.code === "not-ready") { + scheduleResearchRecovery(1); + } + throw err; + } return { workspaceId: created.request.id, state: created.request.state, @@ -644,6 +698,8 @@ function drain(): Promise { } }; await attempt(() => server.stop(true)); + if (researchRecoveryTimer) clearTimeout(researchRecoveryTimer); + if (recoveringResearch) await attempt(() => recoveringResearch!); if (sessionCleanup) clearInterval(sessionCleanup); for (let result of await Promise.allSettled([cleaningSessions])) { if (result.status === "rejected") record(result.reason); @@ -808,6 +864,11 @@ async function restoreChannelLocked(channelId: string, now: Date) { channelId, () => storage.channels.restore({ id: channelId, now }), ); + let recovery = await researchService?.recoverPendingPlannerInline( + placeResearchReference, + channelId, + ); + scheduleResearchRecovery(recovery?.deferred ?? 0); summaryCoordinator?.resume(channelId); announceChannel(result.channel); if (summaryCoordinator) void summaryCoordinator.ensure(channelId).catch(() => {}); @@ -907,6 +968,32 @@ function announceResearchChanged(channelId: string, workspaceId: string, revisio }); } +function announceResearchTerminal(channelId: string, id: string, text: string): Promise { + return withDocumentLock(channelId, async () => { + let active = Rooms.get(channelId)?.plan; + if (active) { + let existing = active.chat.entries.find(entry => entry.id === id); + await Chat.noticeOnce( + { chat: active.chat, plan: active, server, room: channelId }, + id, + existing?.text ?? text, + ); + return; + } + let detached = await Service.open(channelId, documentBackend(), server); + try { + let existing = detached.chat.entries.find(entry => entry.id === id); + await Chat.noticeOnce( + { chat: detached.chat, plan: detached, server, room: channelId }, + id, + existing?.text ?? text, + ); + } finally { + await Service.close(detached); + } + }); +} + async function currentDocumentTarget( channelId: string, ): Promise { @@ -1047,6 +1134,7 @@ researchService = new ResearchWorkspaceService({ }, current: currentDocumentTarget, publish: announceResearchChanged, + terminalNotice: announceResearchTerminal, }); referenceService = new ReferenceService({ storage, @@ -1185,10 +1273,16 @@ leaseRenewal = setInterval(renewLease, LEASE_RENEW_MS); cleanSessions(); sessionCleanup = setInterval(cleanSessions, SESSION_CLEANUP_MS); +let listening: Server | undefined; try { - server = listen(); + server = listening = listen(); + let recovery = await researchService.recoverPendingPlannerInline(placeResearchReference); + scheduleResearchRecovery(recovery.deferred); + await researchService.recoverTerminalPlannerInline(); jobRunner.start(); } catch (err) { + if (listening) await listening.stop(true).catch(() => {}); + if (researchRecoveryTimer) clearTimeout(researchRecoveryTimer); if (sessionCleanup) clearInterval(sessionCleanup); if (leaseRenewal) clearInterval(leaseRenewal); if (leaseWatchdog) clearTimeout(leaseWatchdog); diff --git a/apps/server/src/research/inline-1.test.ts b/apps/server/src/research/inline-1.test.ts new file mode 100644 index 00000000..be7c7137 --- /dev/null +++ b/apps/server/src/research/inline-1.test.ts @@ -0,0 +1,136 @@ +import { expect, it } from "bun:test"; +import { answerArtifact, evidenceArtifact, settle, setup } from "./test-support"; +import { acceptedOffer, terminalNotices } from "./inline.test-fixtures"; + +it("reads a pending, pre-job, then linked accepted request without starting work", async () => { + let context = await setup(); + let offer = acceptedOffer(context.userId); + expect(await context.service.acceptedOfferLink(context.channelId, offer)) + .toEqual({ status: "pending" }); + let placed = false; + let input = { + channelId: context.channelId, + question: offer.brief, + originMessageId: offer.source.messageId, + requestedBy: offer.action!.principalId, + requestedByHandle: offer.action!.actor.handle, + placeReference: async () => { + if (!placed) throw new Error("placement interrupted"); + return "placed" as const; + }, + }; + await expect(context.service.startPlannerInline(input)).rejects.toThrow( + "placement interrupted", + ); + let beforeJobs = await context.jobs.list(context.channelId, 100); + let beforePublications = context.publications.length; + let preJob = await context.restart().acceptedOfferLink(context.channelId, offer); + expect(preJob.status).toBe("unlinked"); + expect(typeof preJob.researchRequestId).toBe("string"); + if (!preJob.researchRequestId) throw new Error("pre-job request ID was not found"); + expect(await context.jobs.list(context.channelId, 100)).toEqual(beforeJobs); + expect(context.publications).toHaveLength(beforePublications); + placed = true; + let started = await context.restart().startPlannerInline(input); + expect(started.request.id).toBe(preJob.researchRequestId); + expect(await context.restart().acceptedOfferLink(context.channelId, offer)) + .toEqual({ status: "linked", researchRequestId: started.request.id }); + for (let index = 0; index < 105; index++) { + await context.storage.research.start({ + id: `other-${index}`, + channelId: context.channelId, + title: `Other ${index}`, + question: `Other question ${index}`, + origin: "planner", + originMessageId: `other-message-${index}`, + createdBy: context.userId, + turnId: `other-turn-${index}`, + messageId: `other-message-record-${index}`, + requestId: `other-request-${index}`, + idempotencyKey: `other-key-${index}`, + fingerprint: `other-fingerprint-${index}`, + now: context.advance(), + lease: context.lease, + }); + } + expect((await context.storage.research.list(context.channelId, 100, true)) + .some(item => item.id === started.request.id)).toBe(false); + expect(await context.restart().acceptedOfferLink(context.channelId, offer)) + .toEqual({ status: "linked", researchRequestId: started.request.id }); + expect(await context.service.acceptedOfferLink("other-channel", offer)) + .toEqual({ status: "pending" }); +}); + +it("rejects colliding or altered original consent identities", async () => { + let context = await setup(); + let offer = acceptedOffer(context.userId); + await context.service.startPlannerInline({ + channelId: context.channelId, + question: offer.brief, + originMessageId: offer.source.messageId, + requestedBy: context.userId, + requestedByHandle: "octocat", + placeReference: async () => "placed", + }); + for ( + let changed of [ + { ...offer, brief: "Different brief" }, + { ...offer, action: { ...offer.action!, principalId: "another-user" } }, + { + ...offer, + action: { + ...offer.action!, + actor: { kind: "member" as const, handle: "someone-else" }, + }, + }, + ] + ) { + await expect(context.service.acceptedOfferLink(context.channelId, changed)) + .rejects.toMatchObject({ code: "invalid-state" }); + } + await expect(context.service.acceptedOfferLink(context.channelId, { + ...offer, + status: "dismissed", + })).rejects.toMatchObject({ code: "invalid-request" }); + let ordinary = await setup(); + let ordinaryOffer = acceptedOffer(ordinary.userId); + await ordinary.service.startPlanner({ + channelId: ordinary.channelId, + question: ordinaryOffer.brief, + originMessageId: ordinaryOffer.source.messageId, + requestedBy: ordinary.userId, + requestedByHandle: "octocat", + }); + await expect(ordinary.service.acceptedOfferLink(ordinary.channelId, ordinaryOffer)) + .rejects.toMatchObject({ code: "invalid-state" }); +}); + +it("emits a stable ready notice after publishing a Planner inline child", async () => { + let context = await setup(); + let started = await context.service.startPlannerInline({ + channelId: context.channelId, + question: "Check the release", + originMessageId: "message-ready", + requestedBy: context.userId, + placeReference: async () => "placed", + }); + let evidence = await settle(context, "research-evidence", evidenceArtifact); + await context.service.jobChanged(evidence.job); + let answer = await settle(context, "research-answer", answerArtifact); + let { service, notices } = terminalNotices(context); + await service.jobChanged(answer.job); + + let request = await service.request(context.channelId, started.request.id); + if (request?.stage !== "ready") throw new Error("published research child is not ready"); + expect(notices).toEqual([{ + channelId: context.channelId, + id: `research-ready:${started.request.id}:${answer.job.id}`, + text: `Research is ready. [Open the research document](` + + `/documents/octo-org/score/${context.channel.slug}/children/${request.child.slug}).`, + }]); + await service.jobChanged(answer.job); + expect(notices.map(value => value.id)).toEqual([ + `research-ready:${started.request.id}:${answer.job.id}`, + `research-ready:${started.request.id}:${answer.job.id}`, + ]); +}); diff --git a/apps/server/src/research/inline-2.test.ts b/apps/server/src/research/inline-2.test.ts new file mode 100644 index 00000000..f87f8a3b --- /dev/null +++ b/apps/server/src/research/inline-2.test.ts @@ -0,0 +1,135 @@ +import { expect, it } from "bun:test"; +import { ResearchWorkspaceService } from "./service"; +import { + answerArtifact, + evidenceArtifact, + REPOSITORY_ID, + requestId, + settle, + setup, +} from "./test-support"; +import { failInitialEvidence, terminalNotices } from "./inline.test-fixtures"; + +it("emits a stable failed notice without exposing a request ID", async () => { + let context = await setup(); + let started = await context.service.startPlannerInline({ + channelId: context.channelId, + question: "Check the release", + originMessageId: "message-failed", + requestedBy: context.userId, + placeReference: async () => "placed", + }); + let failedJobId = await failInitialEvidence(context); + let failed = await context.jobs.get(context.channelId, failedJobId); + if (!failed) throw new Error("failed research job is missing"); + let { service, notices } = terminalNotices(context); + await service.jobChanged(failed.job); + expect((await service.request(context.channelId, started.request.id))?.stage).toBe("failed"); + expect(notices).toEqual([{ + channelId: context.channelId, + id: `research-failed:${started.request.id}:${failedJobId}`, + text: "Research could not be completed. You can retry it from the research card.", + }]); + expect(notices[0]!.text).not.toContain(started.request.id); + await service.jobChanged(failed.job); + expect(notices[1]).toEqual(notices[0]); +}); + +it("recovers missed ready and failed notices across repeated startup scans", async () => { + let context = await setup(); + let ready = await context.service.startPlannerInline({ + channelId: context.channelId, + question: "Check the release", + originMessageId: "missed-ready", + requestedBy: context.userId, + placeReference: async () => "placed", + }); + let evidence = await settle(context, "research-evidence", evidenceArtifact); + await context.service.jobChanged(evidence.job); + let answer = await settle(context, "research-answer", answerArtifact); + await context.service.jobChanged(answer.job); + let failed = await context.service.startPlannerInline({ + channelId: context.channelId, + question: "Check the failure", + originMessageId: "missed-failed", + requestedBy: context.userId, + placeReference: async () => "placed", + }); + let failedJobId = await failInitialEvidence(context); + let notices = new Map(); + let recovered = new ResearchWorkspaceService({ + storage: context.storage, + jobs: context.jobs, + lease: () => context.lease, + current: async () => undefined, + publish: () => {}, + terminalNotice: async (_channelId, id, text) => { + let previous = notices.get(id); + if (previous && previous !== text) throw new Error("conflicting notice"); + notices.set(id, text); + }, + }); + await recovered.recoverTerminalPlannerInline(); + expect([...notices.keys()]).toEqual([ + `research-ready:${ready.request.id}:${answer.job.id}`, + `research-failed:${failed.request.id}:${failedJobId}`, + ]); + await recovered.recoverTerminalPlannerInline(); + expect(notices.size).toBe(2); +}); + +it("does not announce terminal legacy research without an inline card", async () => { + let context = await setup(); + await context.service.start({ + channelId: context.channelId, + question: "Old research", + requestId: requestId(75), + requestedBy: context.userId, + }); + let failedJobId = await failInitialEvidence(context); + let failed = await context.jobs.get(context.channelId, failedJobId); + if (!failed) throw new Error("failed research job is missing"); + let { service, notices } = terminalNotices(context); + await service.jobChanged(failed.job); + expect(notices).toEqual([]); +}); + +it("places one inline Planner request before queueing or publishing it", async () => { + let context = await setup(); + let references = new Set(); + let placements = 0; + let input = { + channelId: context.channelId, + question: "Which public evidence supports version 3?", + originMessageId: "01K39QZG000000000000000003", + requestedBy: context.userId, + placeReference: async (id: string) => { + placements++; + if (placements === 1) { + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + expect(context.publications).toEqual([]); + expect( + (await context.storage.research.get(context.channelId, id))?.workspace + .inlineReference, + ).toBe("pending"); + } + references.add(id); + return "placed" as const; + }, + }; + let started = await context.service.startPlannerInline(input); + expect(references.has(started.request.id)).toBe(true); + expect(started.request).toMatchObject({ state: "pending", stage: "queued" }); + expect((await context.storage.research.get(context.channelId, started.request.id))?.workspace) + .toMatchObject({ origin: "planner", originMessageId: input.originMessageId }); + expect(context.publications).toHaveLength(1); + expect(await context.service.request(context.channelId, started.request.id)) + .toMatchObject({ id: started.request.id, stage: "queued" }); + expect(await context.service.get(context.channelId, started.request.id)).toBeUndefined(); + expect(await context.service.list(context.channelId)).toEqual([]); + expect((await context.service.listRepository(REPOSITORY_ID)).channels).toEqual([]); + let replay = await context.service.startPlannerInline(input); + expect(replay).toEqual({ ...started, repeated: true }); + expect(placements).toBe(2); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(1); +}); diff --git a/apps/server/src/research/inline-3.test.ts b/apps/server/src/research/inline-3.test.ts new file mode 100644 index 00000000..f21e4c9b --- /dev/null +++ b/apps/server/src/research/inline-3.test.ts @@ -0,0 +1,139 @@ +import { expect, it } from "bun:test"; +import { setup } from "./test-support"; + +it("retries a failed inline placement without publishing or enqueuing early", async () => { + let context = await setup(); + let placed = false; + let input = { + channelId: context.channelId, + question: "Which public evidence supports version 3?", + originMessageId: "01K39QZG000000000000000004", + requestedBy: context.userId, + placeReference: async () => { + if (!placed) throw new Error("document commit failed"); + return "placed" as const; + }, + }; + await expect(context.service.startPlannerInline(input)).rejects.toThrow( + "document commit failed", + ); + expect(context.publications).toEqual([]); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + let [durable] = await context.storage.research.list(context.channelId, 100); + expect(durable?.origin).toBe("planner"); + placed = true; + await context.restart().recoverPendingPlannerInline(async (channelId, id) => { + expect(channelId).toBe(context.channelId); + expect(id).toBe(durable?.id); + return "placed"; + }); + expect( + (await context.storage.research.get(context.channelId, durable!.id))?.workspace + .inlineReference, + ).toBe("placed"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(1); + let retry = await context.restart().startPlannerInline(input); + expect(retry.repeated).toBe(true); + expect(retry.request.stage).toBe("queued"); + expect(retry.request.id).toBe(durable?.id); +}); + +it("recovers after the card commit but before the placement marker advances", async () => { + let context = await setup(); + let references = new Set(); + let placed = context.storage.research.markReferencePlaced; + let interrupted = true; + context.storage.research.markReferencePlaced = async input => { + if (interrupted) throw new Error("crash after document commit"); + return placed(input); + }; + let input = { + channelId: context.channelId, + question: "Which API contracts changed?", + originMessageId: "01K39QZG000000000000000007", + requestedBy: context.userId, + placeReference: async (id: string) => { + references.add(id); + return "placed" as const; + }, + }; + await expect(context.service.startPlannerInline(input)) + .rejects.toThrow("crash after document commit"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + interrupted = false; + await context.restart().recoverPendingPlannerInline(async (_channelId, id) => { + references.add(id); + return "placed"; + }); + expect(references.size).toBe(1); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(1); +}); + +it("defers an implementation-locked card and recovers other requests", async () => { + let context = await setup(); + let ids: string[] = []; + for (let suffix of ["a", "b"]) { + await expect(context.service.startPlannerInline({ + channelId: context.channelId, + question: `Research brief ${suffix}`, + originMessageId: `01K39QZG00000000000000000${suffix}`, + requestedBy: context.userId, + placeReference: async id => { + ids.push(id); + return "deferred"; + }, + })).rejects.toMatchObject({ code: "not-ready" }); + } + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + let locked = ids[0]!; + let result = await context.restart().recoverPendingPlannerInline(async (_channelId, id) => + id === locked ? "deferred" : "placed" + ); + expect(result).toEqual({ deferred: 1 }); + expect( + (await context.storage.research.get(context.channelId, locked))?.workspace + .inlineReference, + ).toBe("pending"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(1); + let resumed = await context.restart().recoverPendingPlannerInline(async () => "placed"); + expect(resumed).toEqual({ deferred: 0 }); + expect( + (await context.storage.research.get(context.channelId, locked))?.workspace + .inlineReference, + ).toBe("placed"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(2); +}); + +it("places pending cards when research execution is temporarily disabled", async () => { + let context = await setup({ answer: false }); + let saved = await context.storage.research.start({ + id: "disabled-research", + channelId: context.channelId, + title: "Disabled research", + question: "Which API contracts changed?", + origin: "planner", + originMessageId: "01K39QZG000000000000000008", + inlineReference: "pending", + createdBy: context.userId, + turnId: "disabled-turn", + messageId: "disabled-message", + requestId: "disabled-request", + idempotencyKey: "disabled-research", + fingerprint: "disabled-research", + now: context.advance(), + lease: context.lease, + }); + let placed: string[] = []; + await context.service.recoverPendingPlannerInline(async (_channelId, id) => { + placed.push(id); + return "placed"; + }, context.channelId); + expect(placed).toEqual([saved.workspace.id]); + expect( + (await context.storage.research.get(context.channelId, saved.workspace.id)) + ?.workspace.inlineReference, + ).toBe("placed"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + expect((await context.storage.research.listReferenceRecovery(100)).map(value => value.id)) + .toContain(saved.workspace.id); +}); diff --git a/apps/server/src/research/inline-4.test.ts b/apps/server/src/research/inline-4.test.ts new file mode 100644 index 00000000..73836958 --- /dev/null +++ b/apps/server/src/research/inline-4.test.ts @@ -0,0 +1,64 @@ +import { expect, it } from "bun:test"; +import { ResearchWorkspaceError } from "./service"; +import { + answerArtifact, + createChild, + evidenceArtifact, + reconciledRequest, + REPORT, + REPOSITORY_ID, + settle, + setup, +} from "./test-support"; +import { childFixture, failInitialEvidence } from "./inline.test-fixtures"; + +it("refuses inline Planner research for an absent or child channel", async () => { + let context = await setup(); + let child = await createChild(childFixture(context)); + let placements = 0; + for (let channelId of ["missing-channel", child.id]) { + await expect(context.service.startPlannerInline({ + channelId, + question: "Create a child report", + originMessageId: "01K39QZG000000000000000005", + requestedBy: context.userId, + placeReference: async () => { + placements++; + return "placed"; + }, + })).rejects.toBeInstanceOf(ResearchWorkspaceError); + expect(await context.storage.research.list(channelId, 100)).toEqual([]); + } + expect(placements).toBe(0); +}); + +it("keeps a Planner card linked through failure, retry, and child publication", async () => { + let context = await setup(); + let references = new Set(); + let input = { + channelId: context.channelId, + question: "Which API contracts changed?", + originMessageId: "01K39QZG000000000000000006", + requestedBy: context.userId, + placeReference: async (id: string) => { + references.add(id); + return "placed" as const; + }, + }; + let started = await context.service.startPlannerInline(input); + await failInitialEvidence(context); + expect(await context.service.request(context.channelId, started.request.id)) + .toMatchObject({ stage: "failed" }); + await context.service.retryRequest({ + channelId: context.channelId, + workspaceId: started.request.id, + }); + await settle(context, "research-evidence", evidenceArtifact); + await reconciledRequest(context.service, context.channelId, started.request.id); + await settle(context, "research-answer", answerArtifact); + let ready = await reconciledRequest(context.service, context.channelId, started.request.id); + expect(ready).toMatchObject({ stage: "ready", child: { title: REPORT.title } }); + expect(references).toEqual(new Set([started.request.id])); + expect((await context.storage.channels.list(REPOSITORY_ID, 100)).channels + .filter(channel => channel.parentChannelId === context.channelId)).toHaveLength(1); +}); diff --git a/apps/server/src/research/inline-memory.test.ts b/apps/server/src/research/inline-memory.test.ts new file mode 100644 index 00000000..a0ac0b74 --- /dev/null +++ b/apps/server/src/research/inline-memory.test.ts @@ -0,0 +1,90 @@ +import { expect, it, spyOn } from "bun:test"; + +import { setup } from "./test-support"; + +it("recovers one durably committed reference after its placement marker fails", async () => { + let context = await setup(); + let placements: Array<{ repeated: boolean }> = []; + let placeReference = async (workspaceId: string) => { + let saved = await context.storage.collaboration.load(context.channelId, context.advance()); + let committed = await context.storage.collaboration.commit({ + channelId: context.channelId, + lease: context.lease, + expectedRevision: saved!.channel.revision, + operationId: `research-reference:${workspaceId}`, + epoch: "research-test", + sidecar: { researchReference: workspaceId }, + events: [], + now: context.advance(), + }); + placements.push(committed); + return "placed" as const; + }; + let mark = spyOn(context.storage.research, "markReferencePlaced").mockRejectedValueOnce( + new Error("placement marker unavailable"), + ); + try { + await expect(context.service.startPlannerInline({ + channelId: context.channelId, + question: "Which APIs changed?", + originMessageId: "message-memory-recovery", + requestedBy: context.userId, + placeReference, + })).rejects.toThrow("placement marker unavailable"); + } finally { + mark.mockRestore(); + } + let [pending] = await context.storage.research.listReferenceRecovery(100); + expect(pending?.inlineReference).toBe("pending"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + expect(context.publications).toEqual([]); + expect((await context.storage.collaboration.load(context.channelId, context.advance()))?.sidecar) + .toEqual({ researchReference: pending!.id }); + let restarted = context.restart(); + await restarted.recoverPendingPlannerInline((_channelId, id) => placeReference(id)); + await restarted.recoverPendingPlannerInline((_channelId, id) => placeReference(id)); + expect(placements.map(value => value.repeated)).toEqual([false, true]); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toHaveLength(1); + expect((await context.storage.research.get(context.channelId, pending!.id))?.workspace) + .toMatchObject({ inlineReference: "placed" }); + expect( + (await context.storage.collaboration.load(context.channelId, context.advance())) + ?.channel.revision, + ).toBe(1); +}); + +it("does not enqueue when live authorization fails after reference placement", async () => { + let context = await setup(); + await expect(context.service.startPlannerInline({ + channelId: context.channelId, + question: "Which APIs changed?", + originMessageId: "message-owner-denied", + requestedBy: context.userId, + placeReference: async () => "placed", + beforeStart: () => { + throw new Error("active owner revoked"); + }, + })).rejects.toThrow("active owner revoked"); + let [pending] = await context.storage.research.listReferenceRecovery(100); + expect(pending?.inlineReference).toBe("placed"); + expect((await context.jobs.list(context.channelId, 100))?.jobs).toEqual([]); + expect(context.publications).toEqual([]); +}); + +it("rejects a new inline request after its parent is archived", async () => { + let context = await setup(); + await context.storage.channels.archive({ id: context.channelId, now: context.advance() }); + let placed = false; + await expect(context.service.startPlannerInline({ + channelId: context.channelId, + question: "Which APIs changed?", + originMessageId: "message-archived-parent", + requestedBy: context.userId, + placeReference: async () => { + placed = true; + return "placed"; + }, + })).rejects.toMatchObject({ code: "invalid-request" }); + expect(placed).toBe(false); + expect(await context.storage.research.list(context.channelId, 100)).toEqual([]); +}); diff --git a/apps/server/src/research/inline.test-fixtures.ts b/apps/server/src/research/inline.test-fixtures.ts new file mode 100644 index 00000000..049a2d10 --- /dev/null +++ b/apps/server/src/research/inline.test-fixtures.ts @@ -0,0 +1,77 @@ +import type { ConversationPlan } from "@chopin/protocol"; +import { ResearchWorkspaceService } from "./service"; +import { setup } from "./test-support"; + +export function acceptedOffer( + principalId: string, + brief = "Which public evidence supports version 3?", +): ConversationPlan.ResearchOffer { + return { + id: "offer-lookup", + needId: "evidence", + contextId: "version-3", + source: { + messageId: "message-lookup", + author: { kind: "member", handle: "octocat" }, + quote: brief, + start: 0, + end: brief.length, + }, + brief, + status: "accepted", + action: { + id: "action-lookup", + kind: "research", + actor: { kind: "member", handle: "octocat" }, + principalId, + at: 1, + }, + }; +} + +export async function failInitialEvidence(context: Awaited>) { + let [claimed] = await context.storage.jobs.claim({ + channelId: context.channelId, + claimOwner: "failing-worker", + count: 1, + ttlMs: 30_000, + now: context.advance(), + lease: context.lease, + }); + if (!claimed) throw new Error("initial evidence job was not claimable"); + await context.storage.jobs.fail({ + channelId: context.channelId, + jobId: claimed.id, + claimOwner: "failing-worker", + claimGeneration: claimed.claimGeneration, + reason: "internal-provider-detail", + now: context.advance(), + lease: context.lease, + }); + return claimed.id; +} + +export function terminalNotices(context: Awaited>) { + let notices: Array<{ channelId: string; id: string; text: string }> = []; + let service = new ResearchWorkspaceService({ + storage: context.storage, + jobs: context.jobs, + lease: () => context.lease, + current: async () => undefined, + publish: () => {}, + terminalNotice: async (channelId, id, text) => { + notices.push({ channelId, id, text }); + }, + }); + return { service, notices }; +} + +export function childFixture(context: Awaited>) { + return { + jobs: context.jobs, + lease: context.lease, + now: context.advance, + parent: context.channel, + storage: context.storage, + }; +} diff --git a/apps/server/src/research/placement-close.test-fixtures.ts b/apps/server/src/research/placement-close.test-fixtures.ts new file mode 100644 index 00000000..7440aabe --- /dev/null +++ b/apps/server/src/research/placement-close.test-fixtures.ts @@ -0,0 +1,164 @@ +import assert from "node:assert/strict"; +import { placeResearchReference } from "./placement"; +import { readFileSync } from "node:fs"; +import { parse } from "@babel/parser"; +import * as Plan from "../plan/service"; +import * as Chat from "../chat/service"; +import * as Rooms from "../rooms"; +import { openPlan } from "../testing/plan"; +import { authenticatedMemory } from "../chat/job-authenticated-memory.test-fixtures"; +import { headingHarness } from "../chat/job-heading-harness.test-fixtures"; +import { createPlannerAgent } from "../harness/agents"; +import { openPlannerSession } from "../harness/session"; +import { createConversationRuntime } from "../conversation-plan/runtime"; +import { researchStorage } from "../conversation-plan/accepted-research-storage.test-fixtures"; +// Isolate a failed lock cycle so it cannot strand Chat in the parent test process. +let timeout = setTimeout(() => process.exit(2), 5_000); +let opened = await openPlan(); +let plan = opened.plan; +let identity = await authenticatedMemory(opened); +let context = identity.context; +let runtime = createConversationRuntime({ + config: { agent: true }, + server: () => opened.server, + unavailable: () => false, +}); +let ws = { + data: { room: plan.id, client: "client", handle: "test", principalId: "U_test" }, + send() {}, +} as never; +let room = Rooms.join(ws); +room.plan = plan; +let main = readFileSync(new URL("../main.ts", import.meta.url), "utf8"); +let p = parse(main, { sourceType: "module", plugins: ["typescript"] }).program; +let names = ["withDocumentLock", "closeRoom"]; +let source = names.map(name => { + let n = p.body.find(n => n.type === "FunctionDeclaration" && n.id?.name === name)!; + return main.slice(n.start!, n.end!); +}).join("\n"); +let transpiled = new Bun.Transpiler({ loader: "ts" }).transformSync(source); +let { closeRoom, withDocumentLock } = new Function( + "documentLocks", + "Rooms", + "Service", + "conversationRuntime", + "documentBackend", + "server", + transpiled + "\nreturn {closeRoom,withDocumentLock};", +)(new Map(), Rooms, Plan, runtime, () => opened.backend, opened.server); +let { research, jobs } = researchStorage({ opened, plan } as Parameters[0]); +let entered = Promise.withResolvers(); +let release = Promise.withResolvers(); +let placementRequested = false; +let workspaceId: string | undefined; +let resultReturned = false; +context.createResearch = async request => { + let created = await research.startPlannerInline({ + channelId: plan.id, + question: request.question, + originMessageId: request.entryId, + requestedBy: request.userId, + requestedByHandle: request.handle, + placeReference: async id => { + workspaceId = id; + entered.resolve(); + await release.promise; + placementRequested = true; + return placeResearchReference(plan.id, id, { + get: Rooms.get, + exclusive: withDocumentLock, + detached: async (channelId, action) => { + let detached = await Plan.open(channelId, opened.backend, opened.server); + try { + return await action(detached); + } finally { + await Plan.close(detached); + } + }, + }); + }, + }); + return { + workspaceId: created.request.id, + state: created.request.state, + stage: created.request.stage, + }; +}; +let driver = headingHarness(async () => {}, async () => { + resultReturned = true; + driver.finish(); +}, { + name: "create_research_workspace", + input: { question: "Research the exact API compatibility." }, +}); +let agent = createPlannerAgent(driver.fake); +context.openPlannerSession = (owner, channel) => + openPlannerSession(owner, channel, { + agent, + githubTools: async () => ({ ok: true, value: {} }), + createSandbox: async () => + ({ + defaultWorkingDirectory: "/tmp", + async run() { + return { exitCode: 0, stdout: "", stderr: "" }; + }, + async destroy() {}, + }) as never, + registerCredential: () => () => {}, + }); +await Chat.send(context, ws, { + kind: "chat:send", + ts: 0, + rid: "send", + requestId: crypto.randomUUID(), + to: "planner", + text: "Create research on exact API compatibility.", +}); +await entered.promise; +let closed = false; +let closeError: unknown; +let closing = closeRoom(room, true).then( + () => closed = true, + (error: unknown) => closeError = error, +); +await Bun.sleep(10); +console.log("beforeRelease", { + closed, + closing: !!room.closing, + chatClosed: plan.chat.closed, + chatRunning: !!plan.chat.running, +}); +release.resolve(); +await Promise.race([closing, Bun.sleep(300)]); +console.log("afterRelease", { + closed, + closeError: String(closeError ?? ""), + placementRequested, + resultReturned, + closing: !!room.closing, + chatClosed: plan.chat.closed, + chatRunning: !!plan.chat.running, +}); +if (!closed || closeError) { + console.log( + "REPRODUCED actual Chat/current Harness tool waits for documentLock held by closeRoom, which waits for Chat.running.", + ); + process.exit(1); +} +assert.equal(closeError, undefined); +assert.equal(closed, true); +assert.equal(room.closing, undefined); +assert.equal(plan.chat.running, undefined); +assert.equal(resultReturned, true); +assert.ok(workspaceId); +let stored = await opened.storage.research.get(plan.id, workspaceId); +assert.equal(stored?.workspace.id, workspaceId); +assert.equal(stored?.workspace.inlineReference, "pending"); +assert.equal(stored?.turns[0]?.evidenceJobId, undefined); +assert.deepEqual((await jobs.list(plan.id, 100))?.jobs, []); +console.log("Deferred request retains its identity, pending reference, and no queued work."); +Rooms.forget(room); +identity.revokeAll(); +clearTimeout(timeout); +console.log("No close deadlock in current Harness execution."); +process.exit(0); diff --git a/apps/server/src/research/placement-close.test.ts b/apps/server/src/research/placement-close.test.ts new file mode 100644 index 00000000..94bd3e29 --- /dev/null +++ b/apps/server/src/research/placement-close.test.ts @@ -0,0 +1,22 @@ +import { expect, test } from "bun:test"; + +test("actual research host tool defers placement while document close drains Chat", async () => { + let child = Bun.spawn([ + process.execPath, + new URL("./placement-close.test-fixtures.ts", import.meta.url).pathname, + ], { + stdout: "pipe", + stderr: "pipe", + }); + let [exitCode, stdout, stderr] = await Promise.all([ + child.exited, + new Response(child.stdout).text(), + new Response(child.stderr).text(), + ]); + expect(exitCode, `${stdout}\n${stderr}`).toBe(0); + expect(stderr).toContain("Research card placement is deferred"); + expect(stdout).toContain("No close deadlock in current Harness execution."); + expect(stdout).toContain( + "Deferred request retains its identity, pending reference, and no queued work.", + ); +}, 10_000); diff --git a/apps/server/src/research/placement.test.ts b/apps/server/src/research/placement.test.ts new file mode 100644 index 00000000..1ed4f068 --- /dev/null +++ b/apps/server/src/research/placement.test.ts @@ -0,0 +1,57 @@ +import { expect, test } from "bun:test"; +import { openPlan } from "../testing/plan"; +import * as Plan from "../plan/service"; +import { type PlacementDependencies, placeResearchReference } from "./placement"; + +test("closing room defers before requesting its document lock", async () => { + let calls: string[] = []; + let deps: PlacementDependencies = { + get: () => ({ plan: undefined, closing: Promise.resolve() }), + exclusive: async (_id, action) => { + calls.push("lock"); + return action(); + }, + detached: async () => { + throw new Error("must not open detached plan"); + }, + }; + expect(await placeResearchReference("channel", "workspace", deps)).toBe("deferred"); + expect(calls).toEqual([]); +}); + +test("closing that begins while waiting for the document lock also defers", async () => { + let closing: Promise | undefined; + let deps: PlacementDependencies = { + get: () => ({ plan: undefined, closing }), + exclusive: async (_id, action) => { + closing = Promise.resolve(); + return action(); + }, + detached: async () => { + throw new Error("must not open detached plan"); + }, + }; + expect(await placeResearchReference("channel", "workspace", deps)).toBe("deferred"); +}); + +test("closing Plan persistence defers live and detached placement without writes", async () => { + let opened = await openPlan(); + let plan = opened.plan; + let revision = plan.persistence.revision; + let live = true; + let deps: PlacementDependencies = { + get: () => live ? { plan, closing: undefined } : undefined, + exclusive: async (_id, action) => action(), + detached: async (_id, action) => action(plan), + }; + try { + plan.persistence.closing = true; + expect(await placeResearchReference(plan.id, "workspace", deps)).toBe("deferred"); + live = false; + expect(await placeResearchReference(plan.id, "workspace", deps)).toBe("deferred"); + expect(plan.persistence.revision).toBe(revision); + } finally { + plan.persistence.closing = false; + await Plan.close(plan); + } +}); diff --git a/apps/server/src/research/placement.ts b/apps/server/src/research/placement.ts new file mode 100644 index 00000000..45ccb567 --- /dev/null +++ b/apps/server/src/research/placement.ts @@ -0,0 +1,30 @@ +import * as Plan from "../plan/service"; +import type { Room } from "../rooms"; + +export type PlacementDependencies = { + get(channelId: string): Pick | undefined; + exclusive(channelId: string, action: () => Promise): Promise; + detached(channelId: string, action: (plan: Plan.Plan) => Promise): Promise; +}; + +export function placeResearchReference( + channelId: string, + workspaceId: string, + deps: PlacementDependencies, +): Promise<"placed" | "deferred"> { + // Closing waits for Chat while holding this lock; a host tool must not wait behind it. + if (deps.get(channelId)?.closing) return Promise.resolve("deferred"); + return deps.exclusive(channelId, async () => { + let room = deps.get(channelId); + if (room?.closing) return "deferred"; + let active = room?.plan; + if (active) { + if (active.persistence.closing) return "deferred"; + return Plan.placeResearchReference(active, workspaceId); + } + return deps.detached(channelId, plan => { + if (plan.persistence.closing) return Promise.resolve("deferred"); + return Plan.placeResearchReference(plan, workspaceId); + }); + }); +} diff --git a/apps/server/src/research/service.test.ts b/apps/server/src/research/service.test.ts index 86251d6f..fee67a3a 100644 --- a/apps/server/src/research/service.test.ts +++ b/apps/server/src/research/service.test.ts @@ -1,3 +1,4 @@ +import { childFixture, failInitialEvidence } from "./inline.test-fixtures"; import { describe, expect, it, spyOn } from "bun:test"; import { parse } from "@chopin/dialect"; @@ -30,38 +31,6 @@ import { import type { ResearchReport } from "../jobs/research-workspace"; import type { JsonValue } from "../storage/model"; -async function failInitialEvidence(context: Awaited>) { - let [claimed] = await context.storage.jobs.claim({ - channelId: context.channelId, - claimOwner: "failing-worker", - count: 1, - ttlMs: 30_000, - now: context.advance(), - lease: context.lease, - }); - if (!claimed) throw new Error("initial evidence job was not claimable"); - await context.storage.jobs.fail({ - channelId: context.channelId, - jobId: claimed.id, - claimOwner: "failing-worker", - claimGeneration: claimed.claimGeneration, - reason: "internal-provider-detail", - now: context.advance(), - lease: context.lease, - }); - return claimed.id; -} - -function childFixture(context: Awaited>) { - return { - jobs: context.jobs, - lease: context.lease, - now: context.advance, - parent: context.channel, - storage: context.storage, - }; -} - async function legacyChildEvidence(context: Awaited>) { let fixture = childFixture(context); let child = await createChild(fixture); diff --git a/apps/server/src/research/service.ts b/apps/server/src/research/service.ts index 6a8b78a9..b3734b95 100644 --- a/apps/server/src/research/service.ts +++ b/apps/server/src/research/service.ts @@ -1,4 +1,6 @@ -import { createHash } from "node:crypto"; +import { digest, fingerprint, startIdentity } from "./start-identity"; +import type { ValidatedStartResearchRequest } from "./start-identity"; +import { childDocumentPath } from "@chopin/protocol/document-url"; import { jobDetail as serializedJobDetail } from "../jobs/browser"; import { @@ -7,11 +9,12 @@ import { parseResearchEvidenceArtifact, } from "../jobs/research-workspace"; import { RESEARCH_REPOSITORY_WORKSPACE_LIMIT, researchAttemptDisposition } from "../storage/model"; +import { StorageError } from "../storage/errors"; import { MAX_TITLE_LENGTH } from "../channels/title"; import { publishInitialResearchChild } from "./publication"; import { projectRequestView } from "./request-view"; -import type { Job, Research } from "@chopin/protocol"; +import type { ConversationPlan, Job, Research } from "@chopin/protocol"; import type { ResearchAnswerInput, ResearchEvidence, @@ -60,6 +63,7 @@ export type ResearchWorkspaceServiceOptions = { workspaceId: string, revision: number, ) => void | Promise; + terminalNotice?: (channelId: string, id: string, text: string) => Promise; clock?: () => Date; id?: () => string; }; @@ -101,15 +105,9 @@ export type StartPlannerResearchRequest = { beforeStart?: () => void | Promise; }; -type ValidatedStartResearchRequest = { - channelId: string; - question: string; - scope: string; - origin: "inline" | "planner"; - originMessageId?: string; - requestedBy: string; - requestedByHandle?: string; - beforeStart?: () => void | Promise; +export type StartPlannerInlineResearchRequest = StartPlannerResearchRequest & { + /** Must durably place the canonical card before research work can be enqueued. */ + placeReference: (workspaceId: string) => Promise<"placed" | "deferred">; }; export type ConfirmResearchDraft = { @@ -186,14 +184,6 @@ const SAFE_ID = /^[A-Za-z0-9][A-Za-z0-9._:-]*$/; const HANDLE = /^[A-Za-z0-9](?:[A-Za-z0-9-]{0,38})$/; const ACTIVE_JOB_STATES = new Set(["pending", "paused", "running"]); -function digest(value: string): string { - return createHash("sha256").update(value).digest("hex"); -} - -function fingerprint(kind: string, value: JsonValue): string { - return digest(`${kind}\0${JSON.stringify(value)}`); -} - function safeId(value: unknown, field: string, maximum = MAX_ID): string { if ( typeof value !== "string" || value.length < 1 || value.length > maximum @@ -364,6 +354,7 @@ export class ResearchWorkspaceService { #lease: () => Lease; #current: (channelId: string) => Promise; #publish: ResearchWorkspaceServiceOptions["publish"]; + #terminalNotice: ResearchWorkspaceServiceOptions["terminalNotice"]; #clock: () => Date; #id: () => string; #tails = new Map>(); @@ -374,6 +365,7 @@ export class ResearchWorkspaceService { this.#lease = options.lease; this.#current = options.current; this.#publish = options.publish; + this.#terminalNotice = options.terminalNotice; this.#clock = options.clock ?? (() => new Date()); this.#id = options.id ?? (() => crypto.randomUUID()); } @@ -417,41 +409,114 @@ export class ResearchWorkspaceService { }); } + async startPlannerInline( + input: StartPlannerInlineResearchRequest, + ): Promise { + let channelId = safeId(input.channelId, "Channel id"); + let question = researchBrief(input.question); + let originMessageId = safeId( + input.originMessageId, + "Origin message id", + MAX_ORIGIN_MESSAGE_ID, + ); + let requestedBy = opaqueId(input.requestedBy, "Requesting member id"); + let requestedByHandle = handle(input.requestedByHandle); + return this.#start({ + channelId, + question, + scope: originMessageId, + origin: "planner-inline", + originMessageId, + requestedBy, + ...(requestedByHandle ? { requestedByHandle } : {}), + ...(input.beforeStart ? { beforeStart: input.beforeStart } : {}), + placeReference: input.placeReference, + }); + } + + /** Read only: accepted consent may precede request creation or job linkage. */ + async acceptedOfferLink( + channelId: string, + offer: ConversationPlan.ResearchOffer, + ): Promise<{ status: "pending" | "unlinked" | "linked"; researchRequestId?: string }> { + let verifiedChannelId = safeId(channelId, "Channel id"); + if (offer.status !== "accepted" || offer.action?.kind !== "research") { + throw new ResearchWorkspaceError("invalid-request", "Research offer is not accepted."); + } + let question = researchBrief(offer.brief); + let originMessageId = safeId( + offer.source.messageId, + "Origin message id", + MAX_ORIGIN_MESSAGE_ID, + ); + let requestedBy = opaqueId(offer.action.principalId, "Requesting member id"); + let requestedByHandle = handle(offer.action.actor.handle); + if (!requestedByHandle) { + throw new ResearchWorkspaceError("invalid-state", "Research offer has no consent actor."); + } + let identity = startIdentity({ + channelId: verifiedChannelId, + question, + scope: originMessageId, + origin: "planner-inline", + originMessageId, + requestedBy, + requestedByHandle, + }); + let detail = await this.#storage.research.findByIdempotencyKey( + verifiedChannelId, + identity.idempotencyKey, + ); + if (!detail) return { status: "pending" }; + let { workspace, turns, messages } = detail; + let initial = turns.find(turn => turn.ordinal === 1); + let firstMessage = initial + && messages.find(message => message.turnId === initial.id && message.authorKind === "member"); + if ( + workspace.channelId !== verifiedChannelId + || workspace.idempotencyKey !== identity.idempotencyKey + || workspace.fingerprint !== identity.fingerprint + || workspace.origin !== "planner" || workspace.originMessageId !== originMessageId + || workspace.proposedQuestion !== question || workspace.confirmedQuery !== question + || workspace.createdBy !== requestedBy + || !workspace.inlineReference + || !initial || initial.kind !== "initial" || initial.workspaceId !== workspace.id + || initial.requestId !== originMessageId || initial.fingerprint !== identity.fingerprint + || initial.question !== question || initial.requestedBy !== requestedBy + || !firstMessage || firstMessage.userId !== requestedBy + || firstMessage.userHandle !== requestedByHandle || firstMessage.text !== question + ) { + throw new ResearchWorkspaceError("invalid-state", "Research offer request identity differs."); + } + return { + status: workspace.inlineReference === "placed" && initial.evidenceJobId + ? "linked" + : "unlinked", + researchRequestId: workspace.id, + }; + } + async #start(input: ValidatedStartResearchRequest): Promise { await this.#requireTopLevelChannel(input.channelId); this.#requireDefinition("research-evidence"); this.#requireDefinition("research-answer"); - let requestFingerprint = fingerprint( - `research-${input.origin}`, - input.origin === "inline" - ? { - channelId: input.channelId, - question: input.question, - requestedBy: input.requestedBy, - requestedByHandle: input.requestedByHandle ?? null, - } - : { - channelId: input.channelId, - question: input.question, - originMessageId: input.originMessageId!, - requestedBy: input.requestedBy, - requestedByHandle: input.requestedByHandle ?? null, - }, - ); + let identity = startIdentity(input); + let durableOrigin = input.origin === "planner-inline" ? "planner" : input.origin; let stored = await this.#storage.research.start({ id: this.#newId("Workspace id"), channelId: input.channelId, title: researchWorkspaceTitle(input.question), question: input.question, - origin: input.origin, + origin: durableOrigin, + ...(input.origin === "planner-inline" ? { inlineReference: "pending" as const } : {}), ...(input.originMessageId ? { originMessageId: input.originMessageId } : {}), createdBy: input.requestedBy, ...(input.requestedByHandle ? { createdByHandle: input.requestedByHandle } : {}), turnId: this.#newId("Turn id"), messageId: this.#newId("Message id"), requestId: input.scope, - idempotencyKey: `research-${input.origin}:${digest(input.scope).slice(0, 48)}`, - fingerprint: requestFingerprint, + idempotencyKey: identity.idempotencyKey, + fingerprint: identity.fingerprint, now: this.#time(), lease: this.#lease(), }); @@ -461,6 +526,28 @@ export class ResearchWorkspaceService { if (!initial) { throw new ResearchWorkspaceError("invalid-state", "Research request has no initial work."); } + if (current.workspace.inlineReference && !input.placeReference) { + throw new ResearchWorkspaceError( + "invalid-state", + "Inline Planner research requires a document placement callback.", + ); + } + if (input.placeReference) { + let placement = await input.placeReference(current.workspace.id); + if (placement === "deferred") { + throw new ResearchWorkspaceError( + "not-ready", + "Research card placement is deferred while implementation is active.", + ); + } + if (current.workspace.inlineReference === "pending") { + await this.#storage.research.markReferencePlaced({ + channelId: input.channelId, + workspaceId: current.workspace.id, + lease: this.#lease(), + }); + } + } if (initial.evidenceJobId === undefined) await input.beforeStart?.(); await this.#ensureEvidence(current.workspace, initial); }); @@ -469,6 +556,78 @@ export class ResearchWorkspaceService { return { request, repeated: stored.repeated }; } + /** Recover a committed Planner request even when its initiating turn will never replay. */ + async recoverPendingPlannerInline( + placeReference: ( + channelId: string, + workspaceId: string, + ) => Promise<"placed" | "deferred">, + channelId?: string, + ): Promise<{ deferred: number }> { + if (channelId !== undefined) safeId(channelId, "Channel id"); + let afterId: string | undefined; + let deferred = 0; + while (true) { + let page = await this.#storage.research.listReferenceRecovery(100, afterId, channelId); + for (let workspace of page) { + let outcome = await this.#exclusive(workspace.channelId, workspace.id, async () => { + let detail = await this.#stored(workspace.channelId, workspace.id); + if (!detail.workspace.inlineReference) return "skipped"; + let channel = await this.#storage.channels.get(workspace.channelId); + if (!channel || channel.archivedAt) return "skipped"; + let placement = await placeReference(workspace.channelId, workspace.id); + if (placement === "deferred") return "deferred"; + if (detail.workspace.inlineReference === "pending") { + await this.#storage.research.markReferencePlaced({ + channelId: workspace.channelId, + workspaceId: workspace.id, + lease: this.#lease(), + }); + } + let initial = detail.turns.find(value => value.kind === "initial"); + if (!initial) { + throw new ResearchWorkspaceError( + "invalid-state", + "Research request has no initial work.", + ); + } + if (this.#hasDefinition("research-evidence") && this.#hasDefinition("research-answer")) { + await this.#ensureEvidence(detail.workspace, initial); + } + return "placed"; + }); + if (outcome === "deferred") deferred++; + } + if (page.length < 100) return { deferred }; + afterId = page.at(-1)!.id; + } + } + + /** Revisit terminal Planner cards whose Chat notice may have been missed before shutdown. */ + async recoverTerminalPlannerInline(channelId?: string): Promise { + if (channelId !== undefined) safeId(channelId, "Channel id"); + let afterId: string | undefined; + while (true) { + let page = await this.#storage.research.listTerminalRecovery(100, afterId, channelId); + for (let candidate of page) { + try { + let channel = await this.#storage.channels.get(candidate.channelId); + if (!channel || channel.parentChannelId || channel.archivedAt) continue; + let job = await this.#jobs.get(candidate.channelId, candidate.jobId); + if (job) await this.jobChanged(job.job); + } catch (error) { + if ( + !(error instanceof ResearchWorkspaceError && error.code === "invalid-state") + && !(error instanceof StorageError && error.failure === "corrupt") + ) throw error; + console.warn(`[research] skipped invalid terminal request ${candidate.id}:`, error); + } + } + if (page.length < 100) return; + afterId = page.at(-1)!.id; + } + } + async createDraft(input: CreateResearchDraft): Promise { let channelId = safeId(input.channelId, "Channel id"); let createdBy = opaqueId(input.createdBy, "Creating member id"); @@ -647,7 +806,9 @@ export class ResearchWorkspaceService { if (!this.#validId(channelId) || !this.#validId(workspaceId)) return undefined; return this.#exclusive(channelId, workspaceId, async () => { let detail = await this.#storage.research.get(channelId, workspaceId); - return detail && this.#browserView(detail); + return detail?.workspace.origin === "planner" + ? undefined + : detail && this.#browserView(detail); }); } @@ -682,7 +843,7 @@ export class ResearchWorkspaceService { if (!Number.isSafeInteger(limit) || limit < 1 || limit > 100) { throw new ResearchWorkspaceError("invalid-request", "Research workspace limit is invalid."); } - return (await this.#storage.research.list(channelId, limit)).map(summary); + return (await this.#storage.research.list(channelId, limit, false)).map(summary); } async listRepository( @@ -701,6 +862,7 @@ export class ResearchWorkspaceService { repositoryId, limit, includeArchived, + false, ); return { channels: listed.channels.map(group => ({ @@ -837,7 +999,7 @@ export class ResearchWorkspaceService { async jobChanged(job: JobView): Promise { if ( - job.state !== "completed" + (job.state !== "completed" && job.state !== "failed") || (job.type !== "research-evidence" && job.type !== "research-answer") || !this.#validId(job.channelId) ) return; @@ -845,7 +1007,38 @@ export class ResearchWorkspaceService { let linked = await this.#storage.research.findTurnByJob(job.channelId, job.id); if (!linked) return; await this.#exclusive(job.channelId, linked.workspaceId, async () => { - await this.#reconciled(job.channelId, linked.workspaceId); + let detail = await this.#reconciled(job.channelId, linked.workspaceId); + if (!detail || linked.kind !== "initial" || detail.workspace.inlineReference !== "placed") { + return; + } + let request = await this.#requestView(detail); + if (request.stage !== "ready" && request.stage !== "failed") return; + let activeTurn = detail.turns.find(value => value.id === linked.id); + if ( + request.stage === "ready" + && (job.state !== "completed" || job.id !== activeTurn?.answerJobId) + || request.stage === "failed" + && (job.state !== "failed" + || job.id !== activeTurn?.answerJobId && job.id !== activeTurn?.evidenceJobId) + ) { + return; + } + let id = `research-${request.stage}:${detail.workspace.id}:${job.id}`; + let text = "Research could not be completed. You can retry it from the research card."; + if (request.stage === "ready") { + let parent = await this.#storage.channels.get(job.channelId); + if (!parent) { + throw new ResearchWorkspaceError("invalid-state", "Parent document is missing."); + } + let path = childDocumentPath( + parent.repositoryOwner, + parent.repositoryName, + parent.slug, + request.child.slug, + ); + text = `Research is ready. [Open the research document](${path}).`; + } + await this.#terminalNotice?.(job.channelId, id, text); }); } @@ -1366,10 +1559,11 @@ export class ResearchWorkspaceService { } async #requireTopLevelChannel(channelId: string): Promise { - if (!await this.#isTopLevelChannel(channelId)) { + let channel = await this.#storage.channels.get(channelId); + if (!channel || channel.parentChannelId || channel.archivedAt) { throw new ResearchWorkspaceError( "invalid-request", - "Child documents cannot start research.", + "Channel cannot start research.", ); } } diff --git a/apps/server/src/research/start-identity.ts b/apps/server/src/research/start-identity.ts new file mode 100644 index 00000000..7258fa22 --- /dev/null +++ b/apps/server/src/research/start-identity.ts @@ -0,0 +1,49 @@ +import { createHash } from "node:crypto"; +import type { JsonValue } from "../storage/model"; + +export type ValidatedStartResearchRequest = { + channelId: string; + question: string; + scope: string; + origin: "inline" | "planner" | "planner-inline"; + originMessageId?: string; + requestedBy: string; + requestedByHandle?: string; + beforeStart?: () => void | Promise; + placeReference?: (workspaceId: string) => Promise<"placed" | "deferred">; +}; + +export function startIdentity(input: ValidatedStartResearchRequest): { + idempotencyKey: string; + fingerprint: string; +} { + let durableOrigin = input.origin === "planner-inline" ? "planner" : input.origin; + return { + idempotencyKey: `research-${durableOrigin}:${digest(input.scope).slice(0, 48)}`, + fingerprint: fingerprint( + `research-${durableOrigin}`, + input.origin === "inline" + ? { + channelId: input.channelId, + question: input.question, + requestedBy: input.requestedBy, + requestedByHandle: input.requestedByHandle ?? null, + } + : { + channelId: input.channelId, + question: input.question, + originMessageId: input.originMessageId!, + requestedBy: input.requestedBy, + requestedByHandle: input.requestedByHandle ?? null, + }, + ), + }; +} + +export function digest(value: string): string { + return createHash("sha256").update(value).digest("hex"); +} + +export function fingerprint(kind: string, value: JsonValue): string { + return digest(`${kind}\0${JSON.stringify(value)}`); +} diff --git a/apps/server/src/storage/postgres/startup.test.ts b/apps/server/src/storage/postgres/startup.test.ts index 6a9e991c..b0f935be 100644 --- a/apps/server/src/storage/postgres/startup.test.ts +++ b/apps/server/src/storage/postgres/startup.test.ts @@ -1,9 +1,14 @@ +import { SQL } from "bun"; import { afterEach, describe, expect, it } from "bun:test"; import { join } from "node:path"; +import * as Service from "../../plan/service"; +import { backgroundJob } from "../contract-support"; import { PostgresStorage } from "./adapter"; import type { Subprocess } from "bun"; +import type { SocketData } from "../../wire"; +import type { Server } from "bun"; let database = process.env.TEST_DATABASE_URL; let running: Subprocess[] = []; @@ -18,6 +23,7 @@ function spawn( port: number, stderr: "ignore" | "pipe" = "ignore", extra: Record = {}, + stdout: "ignore" | "pipe" = "ignore", ): Subprocess { let child = Bun.spawn(["bun", join(import.meta.dir, "../../main.ts")], { env: { @@ -36,13 +42,28 @@ function spawn( STORAGE_DRIVER: "postgres", ...extra, }, - stdout: "ignore", + stdout, stderr, }); running.push(child); return child; } +async function started(child: Subprocess): Promise { + let output = child.stdout as ReadableStream; + let reader = output.getReader(); + let content = ""; + try { + while (!content.includes("chopin ยท")) { + let next = await reader.read(); + if (next.done) throw new Error("server exited before startup completed"); + content += new TextDecoder().decode(next.value); + } + } finally { + reader.releaseLock(); + } +} + async function ready(port: number): Promise { for (let attempt = 0; attempt < 200; attempt++) { try { @@ -61,7 +82,8 @@ if (database) { let setup = new PostgresStorage(database); await setup.migrate(); await setup.close(); - let first = spawn(9071); + let first = spawn(9071, "ignore", {}, "pipe"); + await started(first); await ready(9071); let session = await fetch("http://127.0.0.1:9071/api/session"); expect(await session.json()).toEqual({ user: null, agent: false }); @@ -108,7 +130,8 @@ if (database) { first.kill("SIGTERM"); expect(await first.exited).toBe(0); - let replacement = spawn(9072); + let replacement = spawn(9072, "ignore", {}, "pipe"); + await started(replacement); await ready(9072); expect(await storage.sessions.get(sessionId, now)).toBeUndefined(); expect((await storage.collaboration.load(channelId, now))!.agent).toMatchObject({ @@ -122,6 +145,289 @@ if (database) { expect(await replacement.exited).toBe(0); await storage.close(); }, 20_000); + + it("recovers an inline Research card without an open room, once across restarts", async () => { + let storage = new PostgresStorage(database); + await storage.migrate(); + let now = new Date(); + let userId = `user-${crypto.randomUUID()}`; + let channelId = crypto.randomUUID(); + let workspaceId = crypto.randomUUID(); + await storage.users.put({ id: userId, login: "mona", avatarUrl: "", now }); + await storage.channels.create({ + id: channelId, + repositoryId: `repository-${crypto.randomUUID()}`, + repositoryOwner: "octo-org", + repositoryName: "score", + title: "Research recovery", + createdBy: userId, + now, + }); + let lease = await storage.leases.acquire("chopin:writer", crypto.randomUUID(), 30_000); + if (!lease) throw new Error("could not acquire setup lease"); + await storage.research.start({ + id: workspaceId, + channelId, + title: "Research request", + question: "What changed?", + origin: "planner", + originMessageId: crypto.randomUUID(), + inlineReference: "pending", + createdBy: userId, + turnId: crypto.randomUUID(), + messageId: crypto.randomUUID(), + requestId: crypto.randomUUID(), + idempotencyKey: `research-recovery-${workspaceId}`, + fingerprint: `research-recovery-${workspaceId}`, + now, + lease, + }); + await storage.leases.release(lease); + + let countCards = async () => { + let stored = await storage.collaboration.load(channelId, new Date()); + if (!stored) throw new Error("document is missing"); + let source = (await Service.readStored(stored)).source; + return source.split(``).length - 1; + }; + for (let port of [9073, 9074]) { + let child = spawn(port, "ignore", {}, "pipe"); + await started(child); + await ready(port); + expect((await storage.research.get(channelId, workspaceId))?.workspace.inlineReference) + .toBe("placed"); + expect(await countCards()).toBe(1); + child.kill("SIGTERM"); + expect(await child.exited).toBe(0); + } + await storage.close(); + }); + + it("recovers one missed terminal Chat notice across repeated server starts", async () => { + let storage = new PostgresStorage(database); + await storage.migrate(); + let now = new Date(); + let userId = `user-${crypto.randomUUID()}`; + let channelId = crypto.randomUUID(); + let workspaceId = crypto.randomUUID(); + await storage.users.put({ id: userId, login: "mona", avatarUrl: "", now }); + await storage.channels.create({ + id: channelId, + repositoryId: `repository-${crypto.randomUUID()}`, + repositoryOwner: "octo-org", + repositoryName: "score", + title: "Terminal notice recovery", + createdBy: userId, + now, + }); + let lease = await storage.leases.acquire("chopin:writer", crypto.randomUUID(), 30_000); + if (!lease) throw new Error("could not acquire setup lease"); + let request = await storage.research.start({ + id: workspaceId, + channelId, + title: "Research request", + question: "What failed?", + origin: "planner", + originMessageId: crypto.randomUUID(), + inlineReference: "pending", + createdBy: userId, + turnId: crypto.randomUUID(), + messageId: crypto.randomUUID(), + requestId: crypto.randomUUID(), + idempotencyKey: `terminal-${workspaceId}`, + fingerprint: `terminal-${workspaceId}`, + now, + lease, + }); + await storage.research.markReferencePlaced({ channelId, workspaceId, lease }); + let queued = await storage.jobs.enqueue(backgroundJob(channelId, lease, { + type: "research-evidence", + targetKey: `research-evidence:workspace:${workspaceId}:turn:${request.turn.id}:evidence`, + availableAt: now, + now, + })); + await storage.research.linkJob({ + channelId, + workspaceId, + turnId: request.turn.id, + role: "evidence", + jobId: queued.job.id, + now, + lease, + }); + let [claimed] = await storage.jobs.claim({ + channelId, + claimOwner: "failing-worker", + count: 1, + ttlMs: 30_000, + now, + lease, + }); + await storage.jobs.fail({ + channelId, + jobId: claimed!.id, + claimOwner: "failing-worker", + claimGeneration: claimed!.claimGeneration, + reason: "failure", + now, + lease, + }); + await storage.leases.release(lease); + + let noticeId = `research-failed:${workspaceId}:${queued.job.id}`; + for (let port of [9076, 9077]) { + let child = spawn(port, "ignore", {}, "pipe"); + await started(child); + await ready(port); + child.kill("SIGTERM"); + expect(await child.exited).toBe(0); + let durable = await storage.collaboration.load(channelId, new Date()); + expect(JSON.stringify(durable?.sidecar).split(noticeId).length - 1).toBe(1); + } + await storage.close(); + }); + + it("starts while placement is locked and recovers one card and job after release", async () => { + let storage = new PostgresStorage(database); + await storage.migrate(); + let now = new Date(); + let userId = `user-${crypto.randomUUID()}`; + let channelId = crypto.randomUUID(); + let workspaceId = crypto.randomUUID(); + await storage.users.put({ id: userId, login: "mona", avatarUrl: "", now }); + await storage.channels.create({ + id: channelId, + repositoryId: `repository-${crypto.randomUUID()}`, + repositoryOwner: "octo-org", + repositoryName: "score", + title: "Locked research recovery", + createdBy: userId, + now, + }); + let lease = await storage.leases.acquire("chopin:writer", crypto.randomUUID(), 30_000); + if (!lease) throw new Error("could not acquire setup lease"); + let backend: Service.Backend = { + storage, + lease: () => lease, + fatal: error => { + throw error; + }, + }; + let server = { + publish() { + return 0; + }, + } as unknown as Server; + let plan = await Service.open(channelId, backend, server); + plan.graph = { + versions: [{ + number: 1, + revision: 1, + planRevision: plan.revision, + state: "locked", + definition: { + tasks: [{ + id: "task", + title: "Implement", + context: "Research is pending.", + goal: "Complete the work.", + acceptance: ["Implementation is complete.", "Verification is recorded."], + dependsOn: [], + }], + }, + }], + }; + plan.execution = { + id: crypto.randomUUID(), + user: "mona", + client: { name: "Codex", version: "1" }, + session: crypto.randomUUID(), + planRevision: plan.revision, + graphVersion: 1, + graphRevision: 1, + repository: "octo-org/score", + branch: "test/locked-research", + commit: "deadbeef", + startedAt: now.toISOString(), + }; + await Service.persist(plan); + await Service.close(plan); + await storage.research.start({ + id: workspaceId, + channelId, + title: "Research request", + question: "What changed?", + origin: "planner", + originMessageId: crypto.randomUUID(), + inlineReference: "pending", + createdBy: userId, + turnId: crypto.randomUUID(), + messageId: crypto.randomUUID(), + requestId: crypto.randomUUID(), + idempotencyKey: `research-recovery-${workspaceId}`, + fingerprint: `research-recovery-${workspaceId}`, + now, + lease, + }); + await storage.leases.release(lease); + + let countCards = async () => { + let stored = await storage.collaboration.load(channelId, new Date()); + if (!stored) throw new Error("document is missing"); + let source = (await Service.readStored(stored)).source; + return source.split(``).length - 1; + }; + let child = spawn(9075, "pipe", { AGENT: "on" }, "pipe"); + await started(child); + await ready(9075); + expect(await countCards()).toBe(0); + expect((await storage.research.get(channelId, workspaceId))?.workspace.inlineReference) + .toBe("pending"); + expect((await storage.research.get(channelId, workspaceId))?.turns[0]?.evidenceJobId) + .toBeUndefined(); + + // Simulate a durable lifecycle release while this server has no open room. + let sql = new SQL(database); + try { + await sql` + UPDATE channel_state + SET sidecar = to_jsonb(jsonb_set( + (sidecar #>> '{}')::jsonb - 'execution', + '{graph,versions,0,state}', + '"approved"'::jsonb + )::text) + WHERE channel_id = ${channelId} + `; + } finally { + await sql.close(); + } + let placed = false; + for (let attempt = 0; attempt < 150; attempt++) { + let detail = await storage.research.get(channelId, workspaceId); + if (detail?.workspace.inlineReference === "placed" && detail.turns[0]?.evidenceJobId) { + placed = true; + break; + } + await Bun.sleep(100); + } + expect(placed).toBe(true); + expect(await countCards()).toBe(1); + let researchJobs = async () => + (await storage.jobs.list(channelId, 100))?.jobs + .filter(job => job.type === "research-evidence") ?? []; + expect(await researchJobs()).toHaveLength(1); + child.kill("SIGTERM"); + expect(await child.exited).toBe(0); + + let replacement = spawn(9076, "pipe", { AGENT: "on" }, "pipe"); + await started(replacement); + await ready(9076); + expect(await countCards()).toBe(1); + expect(await researchJobs()).toHaveLength(1); + replacement.kill("SIGTERM"); + expect(await replacement.exited).toBe(0); + await storage.close(); + }, 30_000); }); } else { describe("postgres server lifecycle", () => {