diff --git a/apps/server/src/storage/contract.ts b/apps/server/src/storage/contract.ts index 8bdafe2a..ccd5d2b6 100644 --- a/apps/server/src/storage/contract.ts +++ b/apps/server/src/storage/contract.ts @@ -1,3 +1,5 @@ +import { researchWorkspace } from "./research-inline-contract-fixtures"; +import { researchInlineContract } from "./research-inline-contract"; import { describe, expect, it } from "bun:test"; import { @@ -11,33 +13,12 @@ import { StorageError } from "./errors"; import { BACKGROUND_JOB_PROGRESS_LIMIT } from "./model"; import type { StorageFactory as Factory } from "./contract-support"; -import type { CreateResearchWorkspace, JsonValue, Lease } from "./model"; +import type { JsonValue } from "./model"; function attempt(action: () => Promise): Promise { return Promise.resolve().then(action); } -function researchWorkspace( - channelId: string, - createdBy: string, - lease: Lease, - overrides: Partial = {}, -): CreateResearchWorkspace { - return { - id: id("research-workspace"), - channelId, - title: "API compatibility research", - proposedQuestion: "Which API contracts changed?", - origin: "sidebar", - createdBy, - idempotencyKey: id("create-research"), - fingerprint: id("research-fingerprint"), - now: new Date("2026-01-07T03:04:05.000Z"), - lease, - ...overrides, - }; -} - /** The behavioral gate every built-in storage adapter must pass. */ export function storageContract(name: string, factory: Factory): void { describe(`${name} storage`, () => { @@ -1932,6 +1913,7 @@ export function storageContract(name: string, factory: Factory): void { }); researchPublicationContract(factory); + researchInlineContract(factory); it("rejects new archived research mutations but preserves exact replays", async () => { let storage = await opened(factory); diff --git a/apps/server/src/storage/memory/adapter.test.ts b/apps/server/src/storage/memory/adapter.test.ts index b093fe48..f08db6b4 100644 --- a/apps/server/src/storage/memory/adapter.test.ts +++ b/apps/server/src/storage/memory/adapter.test.ts @@ -1,4 +1,78 @@ +import { expect, it } from "bun:test"; import { storageContract } from "../contract"; +import { backgroundJob, contractId as id, userAndChannel } from "../contract-support"; import { MemoryStorage } from "./adapter"; storageContract("memory", () => new MemoryStorage()); + +it("commits all research projections before a linked job can become terminal", async () => { + let storage = new MemoryStorage(); + await storage.migrate(); + try { + let { channelId, userId, lease } = await userAndChannel(storage); + let now = new Date("2026-01-07T03:04:05.000Z"); + let references: { id: string; jobId: string }[] = []; + for (let index = 0; index < 2; index++) { + let workspaceId = id(`projection-workspace-${index}`); + let started = await storage.research.start({ + id: workspaceId, + channelId, + title: "Inline research", + question: `Question ${index}`, + origin: "inline", + createdBy: userId, + turnId: id(`projection-turn-${index}`), + messageId: id(`projection-message-${index}`), + requestId: id(`projection-request-${index}`), + idempotencyKey: id(`projection-start-${index}`), + fingerprint: id(`projection-fingerprint-${index}`), + now, + lease, + }); + let job = await storage.jobs.enqueue(backgroundJob(channelId, lease, { + type: "research-answer", + targetKey: `research-answer:workspace:${workspaceId}:turn:${started.turn.id}:answer`, + availableAt: now, + now, + })); + await storage.research.linkJob({ + channelId, + workspaceId, + turnId: started.turn.id, + role: "answer", + jobId: job.job.id, + now, + lease, + }); + references.push({ id: workspaceId, jobId: job.job.id }); + } + + let commit = storage.collaboration.commit({ + channelId, + lease, + expectedRevision: 0, + operationId: id("multi-projection-commit"), + epoch: "projection-epoch", + update: new Uint8Array([1]), + events: [], + now, + researchProjections: references.map(({ id }) => ({ id, action: "add" })), + }); + let revisionBeforeTerminalization = Promise.resolve().then(async () => { + let state = await storage.collaboration.load(channelId, now); + await storage.jobs.cancel({ + channelId, + jobId: references[0]!.jobId, + now, + lease, + }); + return state?.channel.revision; + }); + + let [result, observedRevision] = await Promise.all([commit, revisionBeforeTerminalization]); + expect(observedRevision).toBe(result.revision); + expect(result.revision).toBe(1); + } finally { + await storage.close(); + } +}); diff --git a/apps/server/src/storage/memory/adapter.ts b/apps/server/src/storage/memory/adapter.ts index a541d255..1cec35e9 100644 --- a/apps/server/src/storage/memory/adapter.ts +++ b/apps/server/src/storage/memory/adapter.ts @@ -3,6 +3,7 @@ import { documentSlug, documentSlugCandidate } from "../../channels/slug"; import { availableChannelTitle } from "../../channels/title"; import { MemoryBackgroundJobStore } from "./jobs"; import { MemoryResearchWorkspaceStore } from "./research"; +import { researchProjectionAllowed, ResearchProjectionConflict } from "../model"; import type { AddUserProject, @@ -742,14 +743,27 @@ export class MemoryStorage implements StorageAdapter { }); } - #commit(input: CommitChannel): Promise { + async #commit(input: CommitChannel): Promise { + this.#assertLease(input.lease); + let replay = this.#operations.get(input.channelId)?.get(input.operationId); + if (replay) return { ...replay, repeated: true }; + for (let change of input.researchProjections ?? []) { + // Keep every projection check and the commit in one synchronous turn. + let detail = this.#research.getCurrent(input.channelId, change.id); + let initial = detail?.turns.find(turn => turn.kind === "initial"); + let jobId = initial?.answerJobId ?? initial?.evidenceJobId; + let job = jobId ? this.#jobs.detail(input.channelId, jobId)?.job : undefined; + if (!researchProjectionAllowed(input.channelId, change, detail?.workspace, initial, job)) { + throw new ResearchProjectionConflict(change.id); + } + } this.#assertLease(input.lease); let found = this.#channels.get(input.channelId); if (!found) throw missing(`channel ${input.channelId} does not exist`); let operations = this.#operations.get(input.channelId) ?? new Map(); this.#operations.set(input.channelId, operations); let repeated = operations.get(input.operationId); - if (repeated) return Promise.resolve({ ...repeated, repeated: true }); + if (repeated) return { ...repeated, repeated: true }; if (found.archivedAt && !input.allowArchived) { throw conflict(`channel ${input.channelId} is archived`); } diff --git a/apps/server/src/storage/memory/research.ts b/apps/server/src/storage/memory/research.ts index f27ad6e5..2711b715 100644 --- a/apps/server/src/storage/memory/research.ts +++ b/apps/server/src/storage/memory/research.ts @@ -27,6 +27,7 @@ import type { PublishInitialResearchReportResult, ResearchJobRole, ResearchMessage, + ResearchTerminalRecovery, ResearchTurn, ResearchWorkspace, ResearchWorkspaceDetail, @@ -189,6 +190,9 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { if (input.origin !== "inline" && input.origin !== "planner") { throw conflict("research workspace start origin is invalid"); } + if (input.inlineReference && input.origin !== "planner") { + throw conflict("inline reference requires Planner origin"); + } if ( input.origin === "inline" && input.originMessageId !== undefined || input.origin === "planner" && !input.originMessageId @@ -208,6 +212,9 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { `research workspace idempotency key ${input.idempotencyKey} was reused`, ); } + if (input.inlineReference && !repeated.inlineReference) { + repeated.inlineReference = "pending"; + } return { workspace: workspace(repeated), turn: turn(repeatedTurn), @@ -238,6 +245,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { confirmedQuery: input.question, origin: input.origin, originMessageId: input.originMessageId, + ...(input.inlineReference ? { inlineReference: input.inlineReference } : {}), createdBy: input.createdBy, confirmedBy: input.createdBy, revision: 0, @@ -280,6 +288,70 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { }; }); + readonly markReferencePlaced = (input: { + channelId: string; + workspaceId: string; + lease: Lease; + }): Promise => + this.#mutate(async () => { + this.#assertChannel(input.channelId, input.lease); + let found = this.#requireWorkspace(input.channelId, input.workspaceId); + if (!found.inlineReference) throw conflict("research reference placement is not required"); + found.inlineReference = "placed"; + }); + + readonly listReferenceRecovery = async ( + limit: number, + afterId?: string, + channelId?: string, + ): Promise => { + let count = Math.min(100, Math.max(1, limit)); + return [...this.#workspaces.values()] + .filter(value => !channelId || value.channelId === channelId) + .filter(value => + value.inlineReference === "pending" + || value.inlineReference === "placed" + && !(this.#turns.get(value.id) ?? [])[0]?.evidenceJobId + ) + .filter(value => !afterId || compareId(value.id, afterId) > 0) + .sort((left, right) => compareId(left.id, right.id)) + .slice(0, count) + .map(workspace); + }; + + readonly listTerminalRecovery = async ( + limit: number, + afterId?: string, + channelId?: string, + ): Promise => { + let count = Math.min(100, Math.max(1, limit)); + let candidates = [...this.#workspaces.values()] + .filter(value => value.origin === "planner" && value.inlineReference === "placed") + .filter(value => !channelId || value.channelId === channelId) + .filter(value => !afterId || compareId(value.id, afterId) > 0) + .sort((left, right) => compareId(left.id, right.id)); + let page: ResearchTerminalRecovery[] = []; + for (let value of candidates) { + let initial = (this.#turns.get(value.id) ?? []).find(turn => turn.kind === "initial"); + if (!initial) continue; + let answer = initial.answerJobId + ? await this.#options.job(value.channelId, initial.answerJobId) + : undefined; + let evidence = initial.evidenceJobId + ? await this.#options.job(value.channelId, initial.evidenceJobId) + : undefined; + let jobId: string | undefined; + if ( + answer?.job.state === "failed" + || answer?.job.state === "completed" && value.publishedChannelId + ) jobId = answer.job.id; + else if (evidence?.job.state === "failed") jobId = evidence.job.id; + if (jobId) page.push({ id: value.id, channelId: value.channelId, jobId }); + if (page.length === count) break; + } + return page; + }; + readonly confirm = ( input: ConfirmResearchWorkspace, ): Promise => @@ -682,10 +754,14 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { readonly list = async ( channelId: string, limit: number, + includePlanner = true, ): Promise => { let count = Math.min(100, Math.max(1, limit)); return [...this.#workspaces.values()] - .filter(value => value.channelId === channelId) + .filter(value => + value.channelId === channelId + && (includePlanner || value.origin !== "planner") + ) .sort((left, right) => right.updatedAt.getTime() - left.updatedAt.getTime() || compareId(left.id, right.id) ) @@ -697,6 +773,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { repositoryId: string, limit: number, includeArchived = false, + includePlanner = true, ): Promise => { let count = Math.min(RESEARCH_REPOSITORY_WORKSPACE_LIMIT, Math.max(1, limit)); let orderedChannels = this.#options.channels(repositoryId) @@ -706,6 +783,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { ); let workspacesByChannel = new Map(); for (let saved of this.#workspaces.values()) { + if (!includePlanner && saved.origin === "planner") continue; let values = workspacesByChannel.get(saved.channelId) ?? []; values.push(saved); workspacesByChannel.set(saved.channelId, values); @@ -744,10 +822,10 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { return { channels: groups, truncated }; }; - readonly get = async ( + getCurrent( channelId: string, workspaceId: string, - ): Promise => { + ): ResearchWorkspaceDetail | undefined { let found = this.#workspaces.get(workspaceId); if (!found || found.channelId !== channelId) return undefined; return { @@ -755,6 +833,19 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore { turns: (this.#turns.get(found.id) ?? []).map(turn), messages: (this.#messages.get(found.id) ?? []).map(message), }; + } + + readonly get = async ( + channelId: string, + workspaceId: string, + ): Promise => this.getCurrent(channelId, workspaceId); + + readonly findByIdempotencyKey = async ( + channelId: string, + idempotencyKey: string, + ): Promise => { + let workspaceId = this.#idempotency.get(key(channelId, idempotencyKey)); + return workspaceId ? this.get(channelId, workspaceId) : undefined; }; readonly findTurnByJob = async ( diff --git a/apps/server/src/storage/model.ts b/apps/server/src/storage/model.ts index 8a1cdb10..869ac60b 100644 --- a/apps/server/src/storage/model.ts +++ b/apps/server/src/storage/model.ts @@ -1,3 +1,5 @@ +import { StorageError } from "./errors"; + /** Values storage adapters may persist without knowing their domain schema. */ export type JsonValue = | null @@ -236,10 +238,22 @@ export type CommitChannel = { sidecar?: JsonValue; events: EventInput[]; now: Date; + /** Browser Research reference changes that must remain authorized at commit time. */ + researchProjections?: ResearchProjectionChange[]; /** Lifecycle and shutdown maintenance may persist after document archival. */ allowArchived?: boolean; }; +export type ResearchProjectionChange = { id: string; action: "add" | "remove" }; + +/** A rejected browser projection is an invalid batch, not a storage outage. */ +export class ResearchProjectionConflict extends StorageError { + constructor(id: string) { + super("conflict", `research projection ${id} is no longer authorized`); + this.name = "ResearchProjectionConflict"; + } +} + export type CommitResult = { revision: number; sequence: number; @@ -484,6 +498,8 @@ export type ResearchWorkspace = { confirmedQuery: string | undefined; origin: ResearchWorkspaceOrigin; originMessageId: string | undefined; + /** Present only for Planner requests that must be placed in the parent document. */ + inlineReference?: "pending" | "placed"; createdBy: string; confirmedBy: string | undefined; revision: number; @@ -495,6 +511,12 @@ export type ResearchWorkspace = { export type ResearchWorkspaceSummary = ResearchWorkspace; +export type ResearchTerminalRecovery = { + id: string; + channelId: string; + jobId: string; +}; + export const RESEARCH_REPOSITORY_CHANNEL_LIMIT = 1_000; export const RESEARCH_REPOSITORY_WORKSPACE_LIMIT = 500; export const RESEARCH_REPOSITORY_CHANNEL_WORKSPACE_LIMIT = 100; @@ -531,6 +553,38 @@ export type ResearchTurn = { updatedAt: Date; }; +export function researchProjectionAllowed( + channelId: string, + change: ResearchProjectionChange, + workspace: + | Pick + | undefined, + initial: + | Pick< + ResearchTurn, + "id" | "workspaceId" | "kind" | "evidenceJobId" | "answerJobId" + > + | undefined, + job: Pick | undefined, +): boolean { + if (!workspace || workspace.id !== change.id || workspace.channelId !== channelId) return false; + if (!initial || initial.kind !== "initial" || initial.workspaceId !== change.id) return false; + if (change.action === "add" && workspace.origin !== "inline") return false; + let role = initial.answerJobId ? "answer" : "evidence"; + let jobId = initial.answerJobId ?? initial.evidenceJobId; + if ( + jobId && ( + !job || job.channelId !== channelId || job.type !== `research-${role}` + || job.targetKey !== `research-${role}:workspace:${change.id}:turn:${initial.id}:${role}` + ) + ) return false; + let terminal = job && ["failed", "cancelled", "superseded"].includes(job.state); + return change.action === "add" + ? !!job && !workspace.publishedChannelId + && ["pending", "paused", "running", "completed"].includes(job.state) + : !!workspace.publishedChannelId || !!terminal || !jobId; +} + export type ResearchMessageAuthorKind = "member" | "agent" | "system"; export type ResearchMessage = { @@ -578,6 +632,7 @@ export type StartResearchWorkspace = { question: string; origin: Extract; originMessageId?: string; + inlineReference?: "pending"; createdBy: string; createdByHandle?: string; turnId: string; diff --git a/apps/server/src/storage/port.ts b/apps/server/src/storage/port.ts index 43773fd8..84843697 100644 --- a/apps/server/src/storage/port.ts +++ b/apps/server/src/storage/port.ts @@ -46,7 +46,9 @@ import type { RenewBackgroundJob, ReplaceChannel, RequeueBackgroundJob, + ResearchTerminalRecovery, ResearchTurn, + ResearchWorkspace, ResearchWorkspaceDetail, ResearchWorkspaceRepositoryList, ResearchWorkspaceSummary, @@ -177,6 +179,21 @@ export interface BackgroundJobStore { export interface ResearchWorkspaceStore { create(input: CreateResearchWorkspace): Promise; start(input: StartResearchWorkspace): Promise; + markReferencePlaced(input: { + channelId: string; + workspaceId: string; + lease: Lease; + }): Promise; + listReferenceRecovery( + limit: number, + afterId?: string, + channelId?: string, + ): Promise; + listTerminalRecovery( + limit: number, + afterId?: string, + channelId?: string, + ): Promise; confirm(input: ConfirmResearchWorkspace): Promise; appendTurn(input: AppendResearchTurn): Promise; linkJob(input: LinkResearchTurnJob): Promise; @@ -189,13 +206,22 @@ export interface ResearchWorkspaceStore { publishInitialReport( input: PublishInitialResearchReport, ): Promise; - list(channelId: string, limit: number): Promise; + list( + channelId: string, + limit: number, + includePlanner?: boolean, + ): Promise; listRepository( repositoryId: string, limit: number, includeArchived?: boolean, + includePlanner?: boolean, ): Promise; get(channelId: string, workspaceId: string): Promise; + findByIdempotencyKey( + channelId: string, + idempotencyKey: string, + ): Promise; findTurnByJob(channelId: string, jobId: string): Promise; } diff --git a/apps/server/src/storage/postgres/adapter.test.ts b/apps/server/src/storage/postgres/adapter.test.ts index 2c502b5f..a2dcfdfb 100644 --- a/apps/server/src/storage/postgres/adapter.test.ts +++ b/apps/server/src/storage/postgres/adapter.test.ts @@ -5,6 +5,8 @@ import { join } from "node:path"; import { documentSlug } from "../../channels/slug"; import { storageContract } from "../contract"; import { StorageError } from "../errors"; +import { ResearchProjectionConflict } from "../model"; +import { contractId, userAndChannel } from "../contract-support"; import { PostgresStorage } from "./adapter"; import { migrate, verifyMigrations } from "./migrations"; import { backfillDocumentSlugs } from "./migrations/002_document_slugs"; @@ -366,6 +368,136 @@ if (url) { } }); + it("serializes a Research projection commit behind a concurrent job failure", async () => { + let storage = new PostgresStorage(url); + let sql = new SQL(url); + let locked = Promise.withResolvers(); + let release = Promise.withResolvers(); + let blocker: Promise | undefined; + try { + await storage.migrate(); + let { channelId, userId, lease } = await userAndChannel(storage); + let jobLease = await storage.leases.acquire( + contractId("research-job-writer"), + contractId("research-worker"), + 60_000, + ); + if (!jobLease) throw new Error("test could not acquire a job storage lease"); + let now = new Date(); + let workspaceId = contractId("inline-research"); + let started = await storage.research.start({ + id: workspaceId, + channelId, + title: "Inline research", + question: "What changed?", + origin: "inline", + createdBy: userId, + turnId: contractId("initial-turn"), + messageId: contractId("initial-message"), + requestId: contractId("initial-request"), + idempotencyKey: contractId("research-start"), + fingerprint: contractId("research-fingerprint"), + now, + lease, + }); + let job = await storage.jobs.enqueue({ + id: contractId("research-answer"), + channelId, + type: "research-answer", + version: 1, + origin: "user", + targetKey: `research-answer:workspace:${workspaceId}:turn:${started.turn.id}:answer`, + idempotencyKey: contractId("research-enqueue"), + fingerprint: contractId("research-job-fingerprint"), + input: { question: "What changed?" }, + availableAt: now, + now, + lease: jobLease, + }); + await storage.research.linkJob({ + channelId, + workspaceId, + turnId: started.turn.id, + role: "answer", + jobId: job.job.id, + now, + lease: jobLease, + }); + let [claimed] = await storage.jobs.claim({ + channelId, + claimOwner: contractId("worker"), + count: 1, + ttlMs: 60_000, + now: new Date(now.getTime() + 1), + lease: jobLease, + }); + if (!claimed) throw new Error("Research answer job was not claimable"); + let initial = await storage.collaboration.load(channelId, now); + if (!initial) throw new Error("Research parent channel could not be loaded"); + + blocker = sql.begin(async transaction => { + await transaction`SELECT id FROM background_jobs WHERE id = ${job.job.id} FOR UPDATE`; + locked.resolve(); + await release.promise; + }); + await locked.promise; + + let failed = storage.jobs.fail({ + channelId, + jobId: job.job.id, + claimOwner: claimed.claimOwner!, + claimGeneration: claimed.claimGeneration, + reason: "test failure", + now: new Date(now.getTime() + 2), + lease: jobLease, + }); + let waitForLockWaiters = async (expected: number) => { + let deadline = Date.now() + 3_000; + while (Date.now() < deadline) { + let [row] = await sql<{ count: number }[]>` + SELECT count(*)::int AS count + FROM pg_stat_activity + WHERE datname = current_database() + AND wait_event_type = 'Lock' + AND query ILIKE '%FROM background_jobs%' + AND query ILIKE '%FOR UPDATE%' + `; + if (row?.count === expected) return; + await Bun.sleep(10); + } + throw new Error(`expected ${expected} Research row-lock waiters`); + }; + await waitForLockWaiters(1); + + let committing = storage.collaboration.commit({ + channelId, + lease, + expectedRevision: initial.channel.revision, + operationId: contractId("research-projection-race"), + epoch: "research-race-epoch", + update: new Uint8Array([1]), + sidecar: { revision: 1 }, + events: [], + now, + researchProjections: [{ id: workspaceId, action: "add" }], + }); + await waitForLockWaiters(2); + release.resolve(); + await blocker; + expect((await failed).state).toBe("failed"); + await expect(committing).rejects.toBeInstanceOf(ResearchProjectionConflict); + + let saved = await storage.collaboration.load(channelId, now); + expect(saved?.channel.revision).toBe(initial.channel.revision); + expect(saved?.updates).toHaveLength(0); + } finally { + release.resolve(); + await blocker?.catch(() => {}); + await storage.close(); + await sql.close(); + } + }); + it("backfills unique readable slugs for channels created before the slug migration", async () => { let sql = new SQL(url); let schema = `slug_migration_${crypto.randomUUID().replaceAll("-", "")}`; @@ -485,6 +617,7 @@ if (url) { "012_child_channels", "013_research_child_publication", "014_inline_research", + "015_planner_inline_reference", ]); expect(await sql<{ table: string | null }[]>`SELECT to_regclass('channel_slugs') AS table`) .toEqual([ diff --git a/apps/server/src/storage/postgres/adapter.ts b/apps/server/src/storage/postgres/adapter.ts index 74d45a38..caf53f9e 100644 --- a/apps/server/src/storage/postgres/adapter.ts +++ b/apps/server/src/storage/postgres/adapter.ts @@ -7,6 +7,7 @@ import { migrate, verifyMigrations } from "./migrations"; import { PostgresNavigationStore } from "./navigation"; import { PostgresBackgroundJobStore } from "./jobs"; import { PostgresResearchWorkspaceStore } from "./research"; +import { researchProjectionAllowed, ResearchProjectionConflict } from "../model"; import type { TransactionSQL } from "bun"; import type { @@ -1274,6 +1275,84 @@ export class PostgresStorage implements StorageAdapter { `channel ${input.channelId} is at revision ${current}, expected ${input.expectedRevision}`, ); } + for ( + let change of [...(input.researchProjections ?? [])].sort((a, b) => + a.id.localeCompare(b.id) + ) + ) { + // Publication locks channel, workspace, initial turn, then job in this order. + let [workspace] = await transaction< + Array<{ + id: string; + channelId: string; + origin: "inline" | "sidebar" | "planner"; + publishedChannelId: string | null; + }> + >` + SELECT id, channel_id AS "channelId", origin, + published_channel_id AS "publishedChannelId" + FROM research_workspaces + WHERE id = ${change.id} AND channel_id = ${input.channelId} + FOR UPDATE + `; + let [initial] = workspace + ? await transaction< + Array<{ + id: string; + workspaceId: string; + kind: "initial" | "follow-up" | "search-more"; + evidenceJobId: string | null; + answerJobId: string | null; + }> + >` + SELECT id, workspace_id AS "workspaceId", kind, + evidence_job_id AS "evidenceJobId", answer_job_id AS "answerJobId" + FROM research_turns + WHERE workspace_id = ${change.id} AND ordinal = 1 + FOR UPDATE + ` + : []; + let jobId = initial?.answerJobId ?? initial?.evidenceJobId; + let [job] = jobId + ? await transaction< + Array<{ + channelId: string; + type: string; + targetKey: string; + state: + | "pending" + | "paused" + | "running" + | "completed" + | "failed" + | "cancelled" + | "superseded"; + }> + >` + SELECT channel_id AS "channelId", type, target_key AS "targetKey", state + FROM background_jobs + WHERE id = ${jobId} AND channel_id = ${input.channelId} + FOR UPDATE + ` + : []; + if ( + !researchProjectionAllowed( + input.channelId, + change, + workspace + ? { ...workspace, publishedChannelId: workspace.publishedChannelId ?? undefined } + : undefined, + initial + ? { + ...initial, + evidenceJobId: initial.evidenceJobId ?? undefined, + answerJobId: initial.answerJobId ?? undefined, + } + : undefined, + job, + ) + ) throw new ResearchProjectionConflict(change.id); + } let sequence = integer(locked.nextSequence, "channel sequence"); let revision = current + 1; await transaction` diff --git a/apps/server/src/storage/postgres/migrations.ts b/apps/server/src/storage/postgres/migrations.ts index a389b627..60a001c2 100644 --- a/apps/server/src/storage/postgres/migrations.ts +++ b/apps/server/src/storage/postgres/migrations.ts @@ -57,6 +57,9 @@ const MIGRATIONS = [{ }, { id: "014_inline_research", path: join(import.meta.dir, "migrations/014_inline_research.sql"), +}, { + id: "015_planner_inline_reference", + path: join(import.meta.dir, "migrations/015_planner_inline_reference.sql"), }] satisfies Migration[]; /** Navigation shipped as 002 before document slugs claimed that number on main. */ diff --git a/apps/server/src/storage/postgres/migrations/015_planner_inline_reference.sql b/apps/server/src/storage/postgres/migrations/015_planner_inline_reference.sql new file mode 100644 index 00000000..038b7346 --- /dev/null +++ b/apps/server/src/storage/postgres/migrations/015_planner_inline_reference.sql @@ -0,0 +1,7 @@ +ALTER TABLE research_workspaces + ADD COLUMN inline_reference text + CHECK (inline_reference IS NULL OR inline_reference IN ('pending', 'placed')); + +CREATE INDEX research_workspaces_inline_recovery + ON research_workspaces (id) + WHERE inline_reference IS NOT NULL; diff --git a/apps/server/src/storage/postgres/research.ts b/apps/server/src/storage/postgres/research.ts index c4fbbc37..cef46d9c 100644 --- a/apps/server/src/storage/postgres/research.ts +++ b/apps/server/src/storage/postgres/research.ts @@ -28,6 +28,7 @@ import type { PublishInitialResearchReportResult, ResearchMessage, ResearchMessageAuthorKind, + ResearchTerminalRecovery, ResearchTurn, ResearchTurnKind, ResearchWorkspace, @@ -62,6 +63,7 @@ type WorkspaceRow = { confirmedQuery: unknown; origin: unknown; originMessageId: unknown; + inlineReference: unknown; createdBy: unknown; confirmedBy: unknown; revision: unknown; @@ -172,6 +174,7 @@ const WORKSPACE_COLUMNS = ` confirmed_query AS "confirmedQuery", origin, origin_message_id AS "originMessageId", + inline_reference AS "inlineReference", created_by AS "createdBy", confirmed_by AS "confirmedBy", revision, @@ -309,6 +312,15 @@ function workspace(row: WorkspaceRow): ResearchWorkspace { throw corrupt(`research workspace ${id} has an invalid origin`); } let revision = integer(row.revision, "research workspace revision"); + let inlineReference = row.inlineReference === null ? undefined : row.inlineReference; + if ( + inlineReference !== undefined && inlineReference !== "pending" && inlineReference !== "placed" + ) { + throw corrupt(`research workspace ${id} has an invalid inline reference state`); + } + if (inlineReference && (row.origin !== "planner" || !row.originMessageId)) { + throw corrupt(`research workspace ${id} has an invalid inline reference origin`); + } let createdAt = date(row.createdAt, "research workspace creation time"); let updatedAt = date(row.updatedAt, "research workspace update time"); if (updatedAt < createdAt) throw corrupt(`research workspace ${id} has invalid timestamps`); @@ -324,6 +336,7 @@ function workspace(row: WorkspaceRow): ResearchWorkspace { confirmedQuery, origin: row.origin as ResearchWorkspaceOrigin, originMessageId: optionalText(row.originMessageId, "research workspace origin message id"), + ...(inlineReference ? { inlineReference } : {}), createdBy: text(row.createdBy, "research workspace creating member"), confirmedBy, revision, @@ -518,6 +531,9 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { if (input.origin !== "inline" && input.origin !== "planner") { throw conflict("research workspace start origin is invalid"); } + if (input.inlineReference && input.origin !== "planner") { + throw conflict("inline reference requires Planner origin"); + } let originMessageId = optionalInput( input.originMessageId, "workspace origin message id", @@ -552,6 +568,15 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { `research workspace idempotency key ${idempotencyKey} was reused`, ); } + if (input.inlineReference && !repeated.inlineReference) { + let [upgraded] = await transaction` + UPDATE research_workspaces SET inline_reference = 'pending' + WHERE id = ${repeated.id} + RETURNING ${transaction.unsafe(WORKSPACE_COLUMNS)} + `; + if (!upgraded) throw corrupt("research reference upgrade returned no record"); + repeated = workspace(upgraded); + } return { workspace: repeated, turn: repeatedTurn, @@ -574,11 +599,12 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { INSERT INTO research_workspaces ( id, channel_id, title, proposed_question, confirmed_query, origin, origin_message_id, + inline_reference, created_by, confirmed_by, revision, next_turn_ordinal, next_message_sequence, idempotency_key, fingerprint, created_at, updated_at ) VALUES ( ${workspaceId}, ${channelId}, ${title}, ${question}, ${question}, - ${input.origin}, ${originMessageId ?? null}, + ${input.origin}, ${originMessageId ?? null}, ${input.inlineReference ?? null}, ${createdBy}, ${createdBy}, 0, 2, 2, ${idempotencyKey}, ${fingerprint}, ${now}, ${now} ) @@ -615,6 +641,81 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { }; })); + readonly markReferencePlaced = (input: { + channelId: string; + workspaceId: string; + lease: Lease; + }): Promise => + this.#run("mark research reference placed", () => + this.#sql.begin(async transaction => { + await this.#fence(transaction, input.lease); + let [placed] = await transaction<{ id: string }[]>` + UPDATE research_workspaces SET inline_reference = 'placed' + WHERE channel_id = ${input.channelId} AND id = ${input.workspaceId} + AND inline_reference IS NOT NULL + RETURNING id + `; + if (!placed) throw conflict("research reference placement is not required"); + })); + + readonly listReferenceRecovery = ( + limit: number, + afterId?: string, + channelId?: string, + ): Promise => + this.#run("list research reference recovery", async () => { + let rows = await this.#sql` + SELECT ${this.#sql.unsafe(WORKSPACE_COLUMNS)} + FROM research_workspaces + WHERE (${afterId ?? ""} = '' OR id > ${afterId ?? ""}) + AND (${channelId ?? ""} = '' OR channel_id = ${channelId ?? ""}) + AND (inline_reference = 'pending' OR inline_reference = 'placed' + AND NOT EXISTS ( + SELECT 1 FROM research_turns + WHERE workspace_id = research_workspaces.id AND ordinal = 1 + AND evidence_job_id IS NOT NULL + )) + ORDER BY id ASC + LIMIT ${Math.min(100, Math.max(1, limit))} + `; + return rows.map(workspace); + }); + + readonly listTerminalRecovery = ( + limit: number, + afterId?: string, + channelId?: string, + ): Promise => + this.#run("list research terminal recovery", async () => { + let rows = await this.#sql<{ id: string; channelId: string; jobId: string }[]>` + SELECT research_workspaces.id, + research_workspaces.channel_id AS "channelId", + CASE + WHEN answer.state = 'failed' OR + (answer.state = 'completed' AND published_channel_id IS NOT NULL) + THEN answer.id + ELSE evidence.id + END AS "jobId" + FROM research_workspaces + JOIN research_turns initial ON initial.workspace_id = research_workspaces.id + AND initial.kind = 'initial' + LEFT JOIN background_jobs answer ON answer.id = initial.answer_job_id + AND answer.channel_id = research_workspaces.channel_id + LEFT JOIN background_jobs evidence ON evidence.id = initial.evidence_job_id + AND evidence.channel_id = research_workspaces.channel_id + WHERE research_workspaces.origin = 'planner' + AND research_workspaces.inline_reference = 'placed' + AND (${afterId ?? ""} = '' OR research_workspaces.id > ${afterId ?? ""}) + AND (${channelId ?? ""} = '' OR research_workspaces.channel_id = ${channelId ?? ""}) + AND (answer.state = 'failed' + OR answer.state = 'completed' AND published_channel_id IS NOT NULL + OR evidence.state = 'failed') + ORDER BY research_workspaces.id ASC + LIMIT ${Math.min(100, Math.max(1, limit))} + `; + return rows; + }); + readonly confirm = ( input: ConfirmResearchWorkspace, ): Promise => @@ -1115,6 +1216,7 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { readonly list = ( channelId: string, limit: number, + includePlanner = true, ): Promise => this.#run("list research workspaces", async () => { let count = Math.min(100, Math.max(1, limit)); @@ -1122,6 +1224,7 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { SELECT ${this.#sql.unsafe(WORKSPACE_COLUMNS)} FROM research_workspaces WHERE channel_id = ${channelId} + AND (${includePlanner} OR origin <> 'planner') ORDER BY updated_at DESC, id ASC LIMIT ${count} `; @@ -1132,6 +1235,7 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { repositoryId: string, limit: number, includeArchived = false, + includePlanner = true, ): Promise => this.#run( "list repository research workspaces", @@ -1174,6 +1278,7 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { SELECT ${transaction.unsafe(WORKSPACE_COLUMNS)} FROM research_workspaces WHERE research_workspaces.channel_id = repository_channels.id + AND (${includePlanner} OR research_workspaces.origin <> 'planner') ORDER BY research_workspaces.updated_at DESC, research_workspaces.id ASC LIMIT ${RESEARCH_REPOSITORY_CHANNEL_WORKSPACE_LIMIT} ) AS workspaces @@ -1250,6 +1355,18 @@ export class PostgresResearchWorkspaceStore implements ResearchWorkspaceStore { }), ); + readonly findByIdempotencyKey = ( + channelId: string, + idempotencyKey: string, + ): Promise => + this.#run("find research workspace by request key", async () => { + let [row] = await this.#sql<{ id: unknown }[]>` + SELECT id FROM research_workspaces + WHERE channel_id = ${channelId} AND idempotency_key = ${idempotencyKey} + `; + return row ? this.get(channelId, text(row.id, "research workspace id")) : undefined; + }); + readonly findTurnByJob = (channelId: string, jobId: string): Promise => this.#run("find research turn by job", async () => { let rows = await this.#sql< diff --git a/apps/server/src/storage/research-inline-contract-1.ts b/apps/server/src/storage/research-inline-contract-1.ts new file mode 100644 index 00000000..e8bb19ce --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-1.ts @@ -0,0 +1,124 @@ +import { expect, it } from "bun:test"; +import { + backgroundJob, + contractId as id, + openedStorage as opened, + userAndChannel, +} from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract1(factory: Factory): void { + for (let terminal of ["failed", "cancelled", "published"] as const) { + it(`conditions Research projection commits on ${terminal} request state`, async () => { + let storage = await opened(factory); + try { + let { channelId, userId, lease } = await userAndChannel(storage); + let now = new Date("2026-01-07T03:04:05.000Z"); + let workspaceId = id("inline-research"); + let started = await storage.research.start({ + id: workspaceId, + channelId, + title: "Inline research", + question: "What changed?", + origin: "inline", + createdBy: userId, + turnId: id("initial-turn"), + messageId: id("initial-message"), + requestId: id("initial-request"), + idempotencyKey: id("inline-start"), + fingerprint: id("inline-fingerprint"), + now, + lease, + }); + let job = await storage.jobs.enqueue(backgroundJob(channelId, lease, { + type: "research-answer", + targetKey: `research-answer:workspace:${workspaceId}:turn:${started.turn.id}:answer`, + availableAt: now, + now, + })); + await storage.research.linkJob({ + channelId, + workspaceId, + turnId: started.turn.id, + role: "answer", + jobId: job.job.id, + now, + lease, + }); + let revision = 0; + let commit = async (action: "add" | "remove") => { + let result = await storage.collaboration.commit({ + channelId, + lease, + expectedRevision: revision, + operationId: id(`research-${action}`), + epoch: "research-epoch", + update: new Uint8Array([revision + 1]), + sidecar: { revision: revision + 1 }, + events: [], + now, + researchProjections: [{ id: workspaceId, action }], + }); + revision = result.revision; + }; + await commit("add"); + await expect(commit("remove")).rejects.toMatchObject({ failure: "conflict" }); + let [claimed] = await storage.jobs.claim({ + channelId, + claimOwner: id("worker"), + count: 1, + ttlMs: 60_000, + now, + lease, + }); + if (!claimed) throw new Error("research job was not claimable"); + if (terminal === "failed") { + await storage.jobs.fail({ + channelId, + jobId: job.job.id, + claimOwner: claimed.claimOwner!, + claimGeneration: claimed.claimGeneration, + reason: "test failure", + now, + lease, + }); + } else if (terminal === "cancelled") { + await storage.jobs.cancel({ channelId, jobId: job.job.id, now, lease }); + } else { + await storage.jobs.settle({ + channelId, + jobId: job.job.id, + claimOwner: claimed.claimOwner!, + claimGeneration: claimed.claimGeneration, + artifact: { source: "# Report\n" }, + now, + lease, + }); + await commit("add"); + await storage.research.publishInitialReport({ + channelId, + workspaceId, + answerJobId: job.job.id, + title: "Report", + initial: { + generation: id("generation"), + epoch: "child-epoch", + source: "# Report\n", + sourceHash: "sha256:report", + document: new Uint8Array([1]), + sidecar: { revision: 0 }, + }, + now, + lease, + }); + } + await expect(commit("add")).rejects.toMatchObject({ failure: "conflict" }); + expect((await storage.collaboration.load(channelId, now))?.channel.revision) + .toBe(revision); + await commit("remove"); + } finally { + await storage.close(); + } + }); + } +} diff --git a/apps/server/src/storage/research-inline-contract-2.ts b/apps/server/src/storage/research-inline-contract-2.ts new file mode 100644 index 00000000..f6b03458 --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-2.ts @@ -0,0 +1,51 @@ +import { expect, it } from "bun:test"; +import { contractId as id, openedStorage as opened, userAndChannel } from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract2(factory: Factory): void { + it("allows removing a legacy unlinked Research reference", async () => { + let storage = await opened(factory); + try { + let { channelId, userId, lease } = await userAndChannel(storage); + let now = new Date("2026-01-07T03:04:05.000Z"); + let workspaceId = id("legacy-inline-research"); + await storage.research.start({ + id: workspaceId, + channelId, + title: "Legacy inline research", + question: "What changed?", + origin: "inline", + createdBy: userId, + turnId: id("legacy-initial-turn"), + messageId: id("legacy-initial-message"), + requestId: id("legacy-initial-request"), + idempotencyKey: id("legacy-inline-start"), + fingerprint: id("legacy-inline-fingerprint"), + now, + lease, + }); + let base = { + channelId, + lease, + expectedRevision: 0, + operationId: id("legacy-removal"), + epoch: "legacy-epoch", + update: new Uint8Array([1]), + events: [], + now, + }; + await expect(storage.collaboration.commit({ + ...base, + researchProjections: [{ id: workspaceId, action: "add" }], + })).rejects.toMatchObject({ failure: "conflict" }); + expect( + await storage.collaboration.commit({ + ...base, + researchProjections: [{ id: workspaceId, action: "remove" }], + }), + ).toMatchObject({ revision: 1 }); + } finally { + await storage.close(); + } + }); +} diff --git a/apps/server/src/storage/research-inline-contract-3.ts b/apps/server/src/storage/research-inline-contract-3.ts new file mode 100644 index 00000000..43c62fbe --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-3.ts @@ -0,0 +1,44 @@ +import { expect, it } from "bun:test"; +import { contractId as id, openedStorage as opened, userAndChannel } from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract3(factory: Factory): void { + it("retains Planner inline placement recovery through exact retries", async () => { + let storage = await opened(factory); + try { + let { userId, channelId, lease } = await userAndChannel(storage); + let input = { + id: id("planner-inline-workspace"), + channelId, + title: "Planner inline research", + question: "Which API contracts changed?", + origin: "planner" as const, + originMessageId: id("planner-inline-origin"), + inlineReference: "pending" as const, + createdBy: userId, + turnId: id("planner-inline-turn"), + messageId: id("planner-inline-message"), + requestId: id("planner-inline-request"), + idempotencyKey: id("planner-inline-start"), + fingerprint: "sha256:planner-inline-start", + now: new Date("2026-01-07T03:04:05.000Z"), + lease, + }; + let started = await storage.research.start(input); + expect(started.workspace.inlineReference).toBe("pending"); + expect((await storage.research.listReferenceRecovery(100)).map(value => value.id)) + .toContain(input.id); + await storage.research.markReferencePlaced({ channelId, workspaceId: input.id, lease }); + await storage.research.markReferencePlaced({ channelId, workspaceId: input.id, lease }); + expect((await storage.research.get(channelId, input.id))?.workspace.inlineReference) + .toBe("placed"); + expect((await storage.research.listReferenceRecovery(100)).map(value => value.id)) + .toContain(input.id); + let repeated = await storage.research.start(input); + expect(repeated.workspace.inlineReference).toBe("placed"); + expect(repeated.repeated).toBe(true); + } finally { + await storage.close(); + } + }); +} diff --git a/apps/server/src/storage/research-inline-contract-4.ts b/apps/server/src/storage/research-inline-contract-4.ts new file mode 100644 index 00000000..2d68b30f --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-4.ts @@ -0,0 +1,53 @@ +import { expect, it } from "bun:test"; +import { contractId as id, openedStorage as opened, userAndChannel } from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract4(factory: Factory): void { + it("finds a Planner request by its channel-scoped durable key", async () => { + let storage = await opened(factory); + try { + let { userId, channelId, repositoryId, lease } = await userAndChannel(storage); + let input = { + id: id("lookup-workspace"), + channelId, + title: "Lookup research", + question: "Which API contracts changed?", + origin: "planner" as const, + originMessageId: id("lookup-origin"), + inlineReference: "pending" as const, + createdBy: userId, + createdByHandle: "octocat", + turnId: id("lookup-turn"), + messageId: id("lookup-message"), + requestId: id("lookup-request"), + idempotencyKey: id("lookup-key"), + fingerprint: "sha256:lookup-fingerprint", + now: new Date("2026-01-07T03:04:05.000Z"), + lease, + }; + await storage.research.start(input); + let other = await storage.channels.create({ + id: id("other-channel"), + repositoryId, + repositoryOwner: "octo-org", + repositoryName: "score", + title: "Other document", + createdBy: userId, + now: input.now, + }); + expect(await storage.research.findByIdempotencyKey(channelId, input.idempotencyKey)) + .toEqual(await storage.research.get(channelId, input.id)); + expect(await storage.research.findByIdempotencyKey(channelId, "missing")) + .toBeUndefined(); + expect(await storage.research.findByIdempotencyKey(other.id, input.idempotencyKey)) + .toBeUndefined(); + await storage.research.markReferencePlaced({ channelId, workspaceId: input.id, lease }); + expect( + (await storage.research.findByIdempotencyKey(channelId, input.idempotencyKey)) + ?.workspace.inlineReference, + ).toBe("placed"); + } finally { + await storage.close(); + } + }); +} diff --git a/apps/server/src/storage/research-inline-contract-5.ts b/apps/server/src/storage/research-inline-contract-5.ts new file mode 100644 index 00000000..00dbda7f --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-5.ts @@ -0,0 +1,87 @@ +import { expect, it } from "bun:test"; +import { + backgroundJob, + contractId as id, + openedStorage as opened, + userAndChannel, +} from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract5(factory: Factory): void { + it("pages only terminal Planner inline requests by workspace ID", async () => { + let storage = await opened(factory); + try { + let { userId, channelId, lease } = await userAndChannel(storage); + let expected: string[] = []; + for (let index = 0; index < 4; index++) { + let workspaceId = `terminal-${index}-${crypto.randomUUID()}`; + let started = await storage.research.start({ + id: workspaceId, + channelId, + title: `Research ${index}`, + question: `Question ${index}`, + origin: "planner", + originMessageId: id(`origin-${index}`), + ...(index < 3 ? { inlineReference: "pending" as const } : {}), + createdBy: userId, + turnId: id(`turn-${index}`), + messageId: id(`message-${index}`), + requestId: id(`request-${index}`), + idempotencyKey: id(`start-${index}`), + fingerprint: id(`fingerprint-${index}`), + now: new Date("2026-01-07T03:04:05.000Z"), + lease, + }); + if (index < 3) { + await storage.research.markReferencePlaced({ channelId, workspaceId, lease }); + } + let job = await storage.jobs.enqueue(backgroundJob(channelId, lease, { + type: "research-evidence", + targetKey: `research-evidence:workspace:${workspaceId}:turn:${started.turn.id}:evidence`, + availableAt: new Date("2026-01-07T03:04:05.000Z"), + now: new Date("2026-01-07T03:04:05.000Z"), + })); + await storage.research.linkJob({ + channelId, + workspaceId, + turnId: started.turn.id, + role: "evidence", + jobId: job.job.id, + now: new Date("2026-01-07T03:04:06.000Z"), + lease, + }); + let [claimed] = await storage.jobs.claim({ + channelId, + claimOwner: `worker-${index}`, + count: 1, + ttlMs: 60_000, + now: new Date("2026-01-07T03:04:07.000Z"), + lease, + }); + await storage.jobs.fail({ + channelId, + jobId: claimed!.id, + claimOwner: `worker-${index}`, + claimGeneration: claimed!.claimGeneration, + reason: "failed", + now: new Date("2026-01-07T03:04:08.000Z"), + lease, + }); + if (index < 3) expected.push(workspaceId); + } + let recovered: string[] = []; + let after: string | undefined; + for (let index = 0; index < 3; index++) { + let page = await storage.research.listTerminalRecovery(1, after, channelId); + expect(page).toHaveLength(1); + recovered.push(page[0]!.id); + after = page[0]!.id; + } + expect(recovered).toEqual(expected); + expect(await storage.research.listTerminalRecovery(1, after, channelId)) + .toEqual([]); + } finally { + await storage.close(); + } + }); +} diff --git a/apps/server/src/storage/research-inline-contract-6.ts b/apps/server/src/storage/research-inline-contract-6.ts new file mode 100644 index 00000000..9175df4c --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-6.ts @@ -0,0 +1,39 @@ +import { expect, it } from "bun:test"; +import { contractId as id, openedStorage as opened, userAndChannel } from "./contract-support"; +import type { StorageFactory as Factory } from "./contract-support"; +import { researchWorkspace } from "./research-inline-contract-fixtures"; + +export function researchInlineContract6(factory: Factory): void { + it("filters hidden Planner requests before workspace listing limits", async () => { + let storage = await opened(factory); + try { + let { userId, channelId, repositoryId, lease } = await userAndChannel(storage); + let sidebar = await storage.research.create(researchWorkspace(channelId, userId, lease)); + let later = new Date(sidebar.workspace.updatedAt.getTime() + 1000); + await storage.research.start({ + id: id("hidden-planner-workspace"), + channelId, + title: "Hidden Planner request", + question: "Research a release", + origin: "planner", + originMessageId: id("hidden-planner-origin"), + inlineReference: "pending", + createdBy: userId, + turnId: id("hidden-planner-turn"), + messageId: id("hidden-planner-message"), + requestId: id("hidden-planner-request"), + idempotencyKey: id("hidden-planner-start"), + fingerprint: "sha256:hidden-planner-start", + now: later, + lease, + }); + expect((await storage.research.list(channelId, 1, false)).map(value => value.id)) + .toEqual([sidebar.workspace.id]); + let repository = await storage.research.listRepository(repositoryId, 1, false, false); + expect(repository.channels.flatMap(group => group.workspaces.map(value => value.id))) + .toContain(sidebar.workspace.id); + } finally { + await storage.close(); + } + }); +} diff --git a/apps/server/src/storage/research-inline-contract-fixtures.ts b/apps/server/src/storage/research-inline-contract-fixtures.ts new file mode 100644 index 00000000..ce80ab55 --- /dev/null +++ b/apps/server/src/storage/research-inline-contract-fixtures.ts @@ -0,0 +1,23 @@ +import { contractId as id } from "./contract-support"; +import type { CreateResearchWorkspace, Lease } from "./model"; + +export function researchWorkspace( + channelId: string, + createdBy: string, + lease: Lease, + overrides: Partial = {}, +): CreateResearchWorkspace { + return { + id: id("research-workspace"), + channelId, + title: "API compatibility research", + proposedQuestion: "Which API contracts changed?", + origin: "sidebar", + createdBy, + idempotencyKey: id("create-research"), + fingerprint: id("research-fingerprint"), + now: new Date("2026-01-07T03:04:05.000Z"), + lease, + ...overrides, + }; +} diff --git a/apps/server/src/storage/research-inline-contract.ts b/apps/server/src/storage/research-inline-contract.ts new file mode 100644 index 00000000..bb95da05 --- /dev/null +++ b/apps/server/src/storage/research-inline-contract.ts @@ -0,0 +1,16 @@ +import { researchInlineContract1 } from "./research-inline-contract-1"; +import { researchInlineContract2 } from "./research-inline-contract-2"; +import { researchInlineContract3 } from "./research-inline-contract-3"; +import { researchInlineContract4 } from "./research-inline-contract-4"; +import { researchInlineContract5 } from "./research-inline-contract-5"; +import { researchInlineContract6 } from "./research-inline-contract-6"; +import type { StorageFactory as Factory } from "./contract-support"; + +export function researchInlineContract(factory: Factory): void { + researchInlineContract1(factory); + researchInlineContract2(factory); + researchInlineContract3(factory); + researchInlineContract4(factory); + researchInlineContract5(factory); + researchInlineContract6(factory); +}