diff --git a/docs/CHAT.md b/docs/CHAT.md index 22327bfbc..ec50c955f 100644 --- a/docs/CHAT.md +++ b/docs/CHAT.md @@ -563,3 +563,30 @@ pinned strip the shell mock shows above the message list (`@corbits/chat-ui`'s `PinnedStrip`) renders every currently-pinned message as a jump-to chip; clicking one scrolls the timeline to that message's own row (`messageDomId`). + +### Question answers notify at most once (CL-7192) + +A question block's answer is relayed into the workbench as the +responding user's own message, and the asking agent is turned on it — +but only once per question, ever. `upsertBlockResponse` +(`packages/chat/src/block-responses.ts`) is a plain "second answer +overwrites the first" write; the one-time send-and-dispatch is gated +separately by `claimBlockResponseNotification`, a guarded UPDATE on the +same row (`notifiedAt`/`notificationClaimToken`) that at most one +submission for a given (tenant, workbench, message, block, principal) +can ever win. A changed answer or a double-click that beats the UI's +disable still lands its own row and its own `block.response` timeline +event, but never wins a second claim — re-answering after the first +notification updates the stored payload without ever notifying the +agent again. + +A claim a submission wins but fails to act on (the send into the +workbench throws) is released so a retried submission can still reach +the agent; `POST /workbenches/:id/messages/:id/blocks/:id/responses` +answers that case with an explicit `notify_failed` 500 telling the +caller their answer was saved and to retry, rather than leaving the row +reading "answered" with the agent silently never turned. The +`block.response` event is always posted — from the upsert's own +returned row, never the request's local payload — before this claim is +even attempted, so an agent dispatched off a won claim always finds its +own correlation event already on the timeline. diff --git a/packages/chat/src/block-responses.test.ts b/packages/chat/src/block-responses.test.ts new file mode 100644 index 000000000..c787e6283 --- /dev/null +++ b/packages/chat/src/block-responses.test.ts @@ -0,0 +1,114 @@ +import { describe, expect, test } from "bun:test"; +import { createInMemoryBlockResponseStore } from "./block-responses"; +import type { BlockResponseKey } from "./block-responses"; + +const KEY: BlockResponseKey = { + tenantId: "ten_1", + workbenchId: "run_1", + messageId: "m1", + blockId: "blk_question1", + principalId: "prn_alice", +}; + +describe("createInMemoryBlockResponseStore — notification claim/release", () => { + test("claiming before any response is upserted for the key fails", async () => { + const store = createInMemoryBlockResponseStore(); + await expect(store.claimBlockResponseNotification(KEY)).resolves.toBe( + false, + ); + }); + + test("the first claim after an upsert wins, returning a token", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + const token = await store.claimBlockResponseNotification(KEY); + expect(typeof token).toBe("string"); + }); + + test("a second claim for the same key loses while the first is held", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + await store.claimBlockResponseNotification(KEY); + await expect(store.claimBlockResponseNotification(KEY)).resolves.toBe( + false, + ); + }); + + test("release with the holder's own token frees the claim for a fresh claim", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + const token = await store.claimBlockResponseNotification(KEY); + await store.releaseBlockResponseNotification(KEY, token as string); + await expect(store.claimBlockResponseNotification(KEY)).resolves.not.toBe( + false, + ); + }); + + test("release with a stale token — one a second claim already replaced — never evicts the live claim", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + const firstToken = await store.claimBlockResponseNotification(KEY); + expect(firstToken).not.toBe(false); + + // The only way a second claim can exist at all is if the first was + // already released -- release it "for real" first, then claim again, + // simulating a delayed release racing a fresh claim. + await store.releaseBlockResponseNotification(KEY, firstToken as string); + const secondToken = await store.claimBlockResponseNotification(KEY); + expect(secondToken).not.toBe(false); + expect(secondToken).not.toBe(firstToken); + + // A caller presenting the first (now stale) token must never release + // the second claim it does not hold. + await store.releaseBlockResponseNotification(KEY, firstToken as string); + await expect(store.claimBlockResponseNotification(KEY)).resolves.toBe( + false, + ); + }); + + test("releasing with a token nobody ever held is a harmless no-op", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + await store.claimBlockResponseNotification(KEY); + await store.releaseBlockResponseNotification(KEY, "some_other_token"); + // The real claim is still held -- a fresh claim attempt still loses. + await expect(store.claimBlockResponseNotification(KEY)).resolves.toBe( + false, + ); + }); + + test("changing the answer keeps the row unclaimable-again once already notified", async () => { + const store = createInMemoryBlockResponseStore(); + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Staging" }, + }); + await store.claimBlockResponseNotification(KEY); + + // A changed answer re-upserts the payload but never clears the claim + // -- CL-7192's ruling is first-answer-wins, permanently, not just + // within a race window. + await store.upsertBlockResponse({ + ...KEY, + payload: { kind: "question", answer: "Production" }, + }); + await expect(store.claimBlockResponseNotification(KEY)).resolves.toBe( + false, + ); + }); +}); diff --git a/packages/chat/src/block-responses.ts b/packages/chat/src/block-responses.ts index 67765de58..9231633b0 100644 --- a/packages/chat/src/block-responses.ts +++ b/packages/chat/src/block-responses.ts @@ -1,17 +1,37 @@ -// Persistence and pure aggregation for poll/form block round-trips. -// `blockId` is the agent-authored `pollId`/`formId` off `PollBlockData`/ -// `FormBlockData` — never trusted as a globally unique key on its own, since -// two different agents (or the same agent twice) can pick the same string -// in two different messages. Every row here is additionally scoped by -// `messageId`, so a response can only ever collide with another response to -// the *same* block in the *same* message; it can never be hijacked into, or -// tallied against, an unrelated message that happens to reuse the id. +// Persistence and pure aggregation for poll/form/question block round-trips. +// `blockId` is the agent-authored `pollId`/`formId`/`questionId` off +// `PollBlockData`/`FormBlockData`/`QuestionBlockData` — never trusted as a +// globally unique key on its own, since two different agents (or the same +// agent twice) can pick the same string in two different messages. Every row +// here is additionally scoped by `messageId`, so a response can only ever +// collide with another response to the *same* block in the *same* message; +// it can never be hijacked into, or tallied against, an unrelated message +// that happens to reuse the id. // // One row per (tenant, workbench, message, block, principal): a second // response from the same principal to the same block overwrites the first -// — "upsert = change vote" for polls, "upsert = resubmit" for forms. +// — "upsert = change vote" for polls, "upsert = resubmit" for forms, +// "upsert = re-answer" for questions. +// +// `notifiedAt` is a question-only claim flag, null until a question's +// answer has been sent into the workbench and dispatched to the asking +// agent. `claimBlockResponseNotification` flips it (and stamps a fresh +// `notificationClaimToken`) in one guarded UPDATE so concurrent submissions +// for the same (message, block, principal) — a changed answer, or a +// double-click that beats the UI's disable — can never both win the claim: +// exactly one caller ever sees itself as responsible for sending the +// notification (see CL-7192). `releaseBlockResponseNotification` is the +// failure path's undo, and only ever succeeds when the token it is given +// still matches the row's current claim — the same token-scoping +// `turn-claims.ts` uses (CL-7129) so a caller can never release a claim it +// does not hold. `write-claims.ts`'s release is unconditional instead, but +// only because nothing there can ever reassign a claim out from under its +// holder; this store has no TTL or reaper either, so today no caller could +// actually present a stale token — the scoping is cheap insurance against a +// future release call site (a sweep, an admin action) rather than a +// reachable bug today. -import { and, eq } from "drizzle-orm"; +import { and, eq, isNull } from "drizzle-orm"; import type { PostgresJsDatabase } from "drizzle-orm/postgres-js"; import { blockResponses } from "./schema"; @@ -37,14 +57,18 @@ export interface BlockResponseRow { readonly payload: BlockResponsePayload; readonly createdAt: Date; readonly updatedAt: Date; + readonly notifiedAt: Date | null; } -export interface UpsertBlockResponseInput { +export interface BlockResponseKey { readonly tenantId: string; readonly workbenchId: string; readonly messageId: string; readonly blockId: string; readonly principalId: string; +} + +export interface UpsertBlockResponseInput extends BlockResponseKey { readonly payload: BlockResponsePayload; } @@ -52,6 +76,29 @@ export interface BlockResponseStore { upsertBlockResponse( input: UpsertBlockResponseInput, ): Promise; + /** + * Atomically claims the right to send a question's answer into the + * workbench and dispatch the asking agent's turn: flips `notifiedAt` + * from null to now and returns a fresh token, but only when it was + * still null. Returns `false` when some other call already holds the + * claim — the caller must then skip the send entirely rather than risk + * a second dispatch. The token must be presented back to + * `releaseBlockResponseNotification` to release this exact claim. + */ + claimBlockResponseNotification( + key: BlockResponseKey, + ): Promise; + /** + * Releases a claim this call took but failed to act on (the send + * threw), resetting `notifiedAt` to null so a retried submission can + * claim it again. A no-op unless `token` is still the current holder — + * a caller can never release a claim it does not hold. Never called + * after a successful send. + */ + releaseBlockResponseNotification( + key: BlockResponseKey, + token: string, + ): Promise; /** * Every response on file for one block instance — including every other * principal's raw payload. Only ever called from inside a route handler @@ -93,14 +140,8 @@ export function aggregatePollResponses( return { tally, total }; } -function responseKey( - tenantId: string, - workbenchId: string, - messageId: string, - blockId: string, - principalId: string, -): string { - return `${tenantId}::${workbenchId}::${messageId}::${blockId}::${principalId}`; +function responseKey(key: BlockResponseKey): string { + return `${key.tenantId}::${key.workbenchId}::${key.messageId}::${key.blockId}::${key.principalId}`; } function blockKey( @@ -115,16 +156,11 @@ function blockKey( export function createInMemoryBlockResponseStore(): BlockResponseStore { const rows = new Map(); const byBlock = new Map>(); + const claimTokens = new Map(); return { async upsertBlockResponse(input) { - const key = responseKey( - input.tenantId, - input.workbenchId, - input.messageId, - input.blockId, - input.principalId, - ); + const key = responseKey(input); const existing = rows.get(key); const now = new Date(); const row: BlockResponseRow = { @@ -136,6 +172,7 @@ export function createInMemoryBlockResponseStore(): BlockResponseStore { payload: input.payload, createdAt: existing?.createdAt ?? now, updatedAt: now, + notifiedAt: existing?.notifiedAt ?? null, }; rows.set(key, row); const blk = blockKey( @@ -150,6 +187,25 @@ export function createInMemoryBlockResponseStore(): BlockResponseStore { return row; }, + async claimBlockResponseNotification(key) { + const k = responseKey(key); + const row = rows.get(k); + if (row === undefined || row.notifiedAt !== null) return false; + const token = crypto.randomUUID(); + claimTokens.set(k, token); + rows.set(k, { ...row, notifiedAt: new Date() }); + return token; + }, + + async releaseBlockResponseNotification(key, token) { + const k = responseKey(key); + if (claimTokens.get(k) !== token) return; + claimTokens.delete(k); + const row = rows.get(k); + if (row === undefined) return; + rows.set(k, { ...row, notifiedAt: null }); + }, + async listBlockResponses(tenantId, workbenchId, messageId, blockId) { const keys = byBlock.get( blockKey(tenantId, workbenchId, messageId, blockId), @@ -177,9 +233,20 @@ function mapRow(row: typeof blockResponses.$inferSelect): BlockResponseRow { payload: row.payload as BlockResponsePayload, createdAt: row.createdAt, updatedAt: row.updatedAt, + notifiedAt: row.notifiedAt, }; } +function keyClause(key: BlockResponseKey) { + return and( + eq(blockResponses.tenantId, key.tenantId), + eq(blockResponses.workbenchId, key.workbenchId), + eq(blockResponses.messageId, key.messageId), + eq(blockResponses.blockId, key.blockId), + eq(blockResponses.principalId, key.principalId), + ); +} + export function createDrizzleBlockResponseStore< TSchema extends Record, >(db: BlockResponseDb): BlockResponseStore { @@ -215,6 +282,25 @@ export function createDrizzleBlockResponseStore< return mapRow(row); }, + async claimBlockResponseNotification(key) { + const token = crypto.randomUUID(); + const claimed = await db + .update(blockResponses) + .set({ notifiedAt: new Date(), notificationClaimToken: token }) + .where(and(keyClause(key), isNull(blockResponses.notifiedAt))) + .returning({ tenantId: blockResponses.tenantId }); + return claimed.length > 0 ? token : false; + }, + + async releaseBlockResponseNotification(key, token) { + await db + .update(blockResponses) + .set({ notifiedAt: null, notificationClaimToken: null }) + .where( + and(keyClause(key), eq(blockResponses.notificationClaimToken, token)), + ); + }, + async listBlockResponses(tenantId, workbenchId, messageId, blockId) { const rows = await db .select() diff --git a/packages/chat/src/migrations.ts b/packages/chat/src/migrations.ts index 16b59d9d1..7586ba523 100644 --- a/packages/chat/src/migrations.ts +++ b/packages/chat/src/migrations.ts @@ -519,6 +519,14 @@ export const chatMigrations: readonly ChatMigration[] = [ WHERE "kind" = 'delivery' AND "run_ref" IS NOT NULL; `, }, + { + name: "0027_block_responses_notified_at", + sql: ` + ALTER TABLE "chat"."block_responses" + ADD COLUMN IF NOT EXISTS "notified_at" timestamptz, + ADD COLUMN IF NOT EXISTS "notification_claim_token" text; + `, + }, ]; /** diff --git a/packages/chat/src/routes.ts b/packages/chat/src/routes.ts index d4cd7de84..dab82171b 100644 --- a/packages/chat/src/routes.ts +++ b/packages/chat/src/routes.ts @@ -18,6 +18,7 @@ import type { Context } from "hono"; import { streamSSE } from "hono/streaming"; import { type } from "arktype"; import { decodedOrNull } from "@corbits/url-path"; +import { reportError } from "@corbits/error-sink"; import type { TenantEnv } from "@intx/hub-api"; import type { RequireGrant } from "@intx/hub-api"; @@ -2227,58 +2228,19 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { : {}), }; - const row = await deps.blockResponses.upsertBlockResponse({ + const responseKey = { tenantId: ownerTenantId, workbenchId, messageId, blockId, principalId: principal.id, + }; + + const row = await deps.blockResponses.upsertBlockResponse({ + ...responseKey, payload, }); - // A question's answer is the interview reply itself, not just a - // structured event: post it into the workbench as the responding - // user's own message so the asking agent receives it exactly as it - // would any other reply from that user — visible in-thread, routed - // by the workbench's normal host routing, no side channel only the - // agent can read. - if (payload.kind === "question") { - const answer = await sendWorkbenchMessage( - { - store: deps.store, - platform: deps.platform, - roomMessages: deps.roomMessages, - publish, - turnQueue, - ...(deps.agentTurns !== undefined - ? { agentTurns: deps.agentTurns } - : {}), - ...(deps.threads !== undefined ? { threads: deps.threads } : {}), - ...(deps.turnDispatchTimeoutMs !== undefined - ? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs } - : {}), - ...(deps.waitUntilFreeTimeoutMs !== undefined - ? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs } - : {}), - }, - { - tenantId: ownerTenantId, - principalId: principal.id, - senderAddress: senderAddressOf(c), - workbenchId, - messageParts: [{ kind: "text", text: payload.answer }], - // A question block's `blockId` IS its `questionId` - // (`ask_user`'s `postQuestion` mints one id and reuses it as - // both), which is also the `message_response` gate's own - // correlationId (`beforeAskUser`, `@corbits/interaction-tools`) - // — so this is the exact id that resolves the gate this answer - // is for, not merely "whichever gate is next". - correlationId: blockId, - }, - ); - deps.onMessageFanout?.(answer.fanoutDelivered); - } - // A machine-readable event into the same workbench timeline the // responder is already a member of, so the outcome reaches the // emitting agent in-context on its next turn — the same "the message @@ -2287,7 +2249,10 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { // the same event any other message in this workbench would show them; // that is the workbench's own membership boundary, not a new one — the // GET route below is the boundary that must never let a member read - // *another* member's raw response on demand. + // *another* member's raw response on demand. Posted before the + // question branch below ever gets a chance to dispatch a turn, so an + // agent reading the timeline for its turn always finds this event + // already there (CL-7192). await postRoomMessage( { roomMessages: deps.roomMessages, publish }, { @@ -2299,12 +2264,81 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { { kind: "event", event: "block.response", - data: { messageId, blockId, ...payload }, + data: { messageId, blockId, ...row.payload }, }, ], }, ); + // A question's answer is the interview reply itself, not just a + // structured event: post it into the workbench as the responding + // user's own message so the asking agent receives it exactly as it + // would any other reply from that user — visible in-thread, routed + // by the workbench's normal host routing, no side channel only the + // agent can read. + // + // `claimBlockResponseNotification` guards this so only the + // submission that first answers a question ever sends the message + // and dispatches a turn: a changed answer, or a double-click that + // beats the UI's disable, still lands its own row above, but never + // gets a second turn (CL-7192). A failed send releases the claim so + // a retried submission can still notify the agent, rather than + // leaving the answer permanently unable to reach it. + if (payload.kind === "question") { + const claimToken = + await deps.blockResponses.claimBlockResponseNotification(responseKey); + if (claimToken !== false) { + try { + const answer = await sendWorkbenchMessage( + { + store: deps.store, + platform: deps.platform, + roomMessages: deps.roomMessages, + publish, + turnQueue, + ...(deps.agentTurns !== undefined + ? { agentTurns: deps.agentTurns } + : {}), + ...(deps.threads !== undefined + ? { threads: deps.threads } + : {}), + ...(deps.turnDispatchTimeoutMs !== undefined + ? { turnDispatchTimeoutMs: deps.turnDispatchTimeoutMs } + : {}), + ...(deps.waitUntilFreeTimeoutMs !== undefined + ? { waitUntilFreeTimeoutMs: deps.waitUntilFreeTimeoutMs } + : {}), + }, + { + tenantId: ownerTenantId, + principalId: principal.id, + senderAddress: senderAddressOf(c), + workbenchId, + messageParts: [{ kind: "text", text: payload.answer }], + }, + ); + deps.onMessageFanout?.(answer.fanoutDelivered); + } catch (err) { + const refId = reportError(err, { + operation: "chat.blockResponse.notifyQuestionAnswer", + tenantId: ownerTenantId, + roomId: workbenchId, + }); + await deps.blockResponses.releaseBlockResponseNotification( + responseKey, + claimToken, + ); + return c.json( + ErrorEnvelope( + "notify_failed", + `your answer was saved, but the agent couldn't be notified — try again (ref ${refId})`, + ), + 500, + ); + } + } + } + return c.json({ blockId, updatedAt: row.updatedAt.toISOString() }, 200); }, ); diff --git a/packages/chat/src/schema.ts b/packages/chat/src/schema.ts index 1db70257f..0b3e70778 100644 --- a/packages/chat/src/schema.ts +++ b/packages/chat/src/schema.ts @@ -352,6 +352,14 @@ export const blockResponses = chatSchema.table( updatedAt: timestamp("updated_at", { withTimezone: true }) .notNull() .defaultNow(), + // Question-only claim flag: null until a question's answer has been + // sent into the workbench and its turn dispatched. See + // `claimBlockResponseNotification` in `./block-responses.ts`. + notifiedAt: timestamp("notified_at", { withTimezone: true }), + // Scopes `releaseBlockResponseNotification` to the exact claim it + // took, set together with `notifiedAt`, so a release can never + // clobber a claim it does not hold. See `./block-responses.ts`. + notificationClaimToken: text("notification_claim_token"), }, (table) => [ primaryKey({ diff --git a/packages/chat/test/block-responses-routes.test.ts b/packages/chat/test/block-responses-routes.test.ts index 93c31b257..fa567c208 100644 --- a/packages/chat/test/block-responses-routes.test.ts +++ b/packages/chat/test/block-responses-routes.test.ts @@ -9,18 +9,40 @@ import { createChatRoutes } from "../src/routes"; import { createInMemoryBlockResponseStore } from "../src/block-responses"; import { createInMemoryChatStore } from "../src/store"; import type { ChatStore } from "../src/store"; +import { createInMemoryRoomMessageStore } from "../src/room-messages"; +import type { RoomMessageStore } from "../src/room-messages"; import { buildDeps, createWorkbench, fakePlatform, mountAs, - settleFanout, TENANT, timelineEvents, timelineOf, timelineTexts, } from "./test-support"; +/** A `RoomMessageStore` whose next attempt to post a text-bearing message + * (the answer, never the `block.response` event) throws once, then + * delegates normally — used to simulate a notify failure a retry + * recovers from. */ +function roomMessagesFailingNextAnswer( + base: RoomMessageStore, +): RoomMessageStore { + let shouldFail = true; + return { + ...base, + async insertMessage(input) { + const isAnswerText = input.parts.some((part) => part.kind === "text"); + if (isAnswerText && shouldFail) { + shouldFail = false; + throw new Error("simulated notify failure"); + } + return base.insertMessage(input); + }, + }; +} + function responsesUrl(workbenchId: string, messageId: string, blockId: string) { return `/workbenches/${workbenchId}/messages/${messageId}/blocks/${blockId}/responses`; } @@ -307,114 +329,119 @@ describe("block response routes — question answers", () => { }); }); - test("answering a question stamps the answer's mail with the block's own id as its correlationId", async () => { + test("changing an answer updates the stored response but never dispatches a second turn", async () => { const platform = fakePlatform(); const deps = buildDeps({ platform, blockResponses: createInMemoryBlockResponseStore(), }); const app = mountAs(createChatRoutes(deps), "prn_alice"); - const { body: workbench } = await createWorkbench(app, { - kind: "workbench", - participants: ["ins_echo1@acme.example"], - }); + const workbenchId = await newWorkbench(deps.store); - const post = await app.request( - responsesUrl(workbench.id, "m1", "blk_question1"), + const first = await app.request( + responsesUrl(workbenchId, "m1", "blk_question1"), { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ kind: "question", answer: "Staging" }), }, ); - expect(post.status).toBe(200); - await settleFanout(); - - // The block's own id (`blk_question1`, the same id `postQuestion` - // mints as `questionId`) rides as the answer mail's correlationId — - // the exact id `tryCorrelate` (`vendor/intx/inference/src/reactor.ts`) - // needs to resolve the `message_response` gate this answers, rather - // than "whichever gate is next". - expect(platform.sentMail).toHaveLength(1); - expect(platform.sentMail[0]?.correlationId).toBe("blk_question1"); - }); + expect(first.status).toBe(200); - test("a question's correlationId survives batching even when it isn't the batch's last queued message", async () => { - // A held first dispatch forces the answer (queued 2nd) and a further - // plain follow-up (queued 3rd, landing last) into one batched turn — - // proving the batch's correlationId comes from whichever queued turn - // carries one, not from `batch[batch.length - 1]`, which here is the - // follow-up and carries none. - let releaseHold: () => void = () => {}; - const held = new Promise((resolve) => { - releaseHold = resolve; - }); - let resolveFirstDispatchStarted: () => void = () => {}; - const firstDispatchStarted = new Promise((resolve) => { - resolveFirstDispatchStarted = resolve; - }); - let holdConsumed = false; + const second = await app.request( + responsesUrl(workbenchId, "m1", "blk_question1"), + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ kind: "question", answer: "Production" }), + }, + ); + expect(second.status).toBe(200); - const platform = fakePlatform(); - const deliverMail = platform.sendMail.bind(platform); - platform.sendMail = async (input) => { - if (!holdConsumed) { - holdConsumed = true; - resolveFirstDispatchStarted(); - await held; - } - return deliverMail(input); - }; + // Only the first answer ever became a message the agent was asked to + // turn on -- the second lands in the stored row (and its own + // `block.response` event) but never re-sends or re-dispatches. + const timeline = await timelineOf(deps, workbenchId); + expect(timelineTexts(timeline)).toEqual(["Staging"]); + expect(timelineEvents(timeline, "block.response")).toHaveLength(2); + const get = await getResponses(app, workbenchId, "m1", "blk_question1"); + const body = (await get.json()) as { own: unknown }; + expect(body.own).toEqual({ kind: "question", answer: "Production" }); + }); + + test("a double-click resubmitting the identical answer never dispatches a second turn", async () => { + const platform = fakePlatform(); const deps = buildDeps({ platform, blockResponses: createInMemoryBlockResponseStore(), }); const app = mountAs(createChatRoutes(deps), "prn_alice"); - const { body: workbench } = await createWorkbench(app, { - kind: "workbench", - participants: ["ins_echo1@acme.example"], - }); + const workbenchId = await newWorkbench(deps.store); - // Message 1 dispatches immediately and is held open. - const first = await app.request(`/workbenches/${workbench.id}/messages`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ parts: [{ kind: "text", text: "hello" }] }), + const submit = () => + app.request(responsesUrl(workbenchId, "m1", "blk_question1"), { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ kind: "question", answer: "Production" }), + }); + + const [first, second] = await Promise.all([submit(), submit()]); + expect(first.status).toBe(200); + expect(second.status).toBe(200); + + const timeline = await timelineOf(deps, workbenchId); + expect(timelineTexts(timeline)).toEqual(["Production"]); + }); + + test("a failed notification releases the claim so a retried answer still reaches the agent", async () => { + const platform = fakePlatform(); + const roomMessages = roomMessagesFailingNextAnswer( + createInMemoryRoomMessageStore(), + ); + const deps = buildDeps({ + platform, + roomMessages, + blockResponses: createInMemoryBlockResponseStore(), }); - expect(first.status).toBe(201); - await firstDispatchStarted; + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const workbenchId = await newWorkbench(deps.store); - // Queued 2nd, while the first dispatch is still held: the answer, - // carrying the block's id as its correlationId. - const post = await app.request( - responsesUrl(workbench.id, "m1", "blk_question1"), - { + const submit = () => + app.request(responsesUrl(workbenchId, "m1", "blk_question1"), { method: "POST", headers: { "content-type": "application/json" }, - body: JSON.stringify({ kind: "question", answer: "Staging" }), - }, - ); - expect(post.status).toBe(200); + body: JSON.stringify({ kind: "question", answer: "Production" }), + }); + + const failed = await submit(); + expect(failed.status).toBe(500); + const failedBody = (await failed.json()) as { + error: { code: string; message: string }; + }; + expect(failedBody.error.code).toBe("notify_failed"); + expect(failedBody.error.message).toMatch(/ref /); - // Queued 3rd, landing last in the batch: an unrelated follow-up with - // no correlationId of its own. - const third = await app.request(`/workbenches/${workbench.id}/messages`, { - method: "POST", - headers: { "content-type": "application/json" }, - body: JSON.stringify({ - parts: [{ kind: "text", text: "also, one more thing" }], - }), + // The answer itself was already durable even though the notify failed. + const afterFailure = await getResponses( + app, + workbenchId, + "m1", + "blk_question1", + ); + const afterFailureBody = (await afterFailure.json()) as { own: unknown }; + expect(afterFailureBody.own).toEqual({ + kind: "question", + answer: "Production", }); - expect(third.status).toBe(201); + expect(timelineTexts(await timelineOf(deps, workbenchId))).toEqual([]); - releaseHold(); - await settleFanout(); + const retried = await submit(); + expect(retried.status).toBe(200); - // First dispatch (message 1), then one batched dispatch covering - // both the answer and the follow-up. - expect(platform.sentMail).toHaveLength(2); - expect(platform.sentMail[1]?.correlationId).toBe("blk_question1"); + const timeline = await timelineOf(deps, workbenchId); + expect(timelineTexts(timeline)).toEqual(["Production"]); + expect(timelineEvents(timeline, "block.response")).toHaveLength(2); }); test("a question response without answer text is rejected", async () => { diff --git a/packages/chat/test/block-responses.drizzle.test.ts b/packages/chat/test/block-responses.drizzle.test.ts new file mode 100644 index 000000000..20b7553e2 --- /dev/null +++ b/packages/chat/test/block-responses.drizzle.test.ts @@ -0,0 +1,187 @@ +// DB-gated: skipped when no DATABASE_URL is reachable (a fresh checkout +// still runs the unit gates), mirroring `write-claims.drizzle.test.ts`. +// Runs against its own scratch database. +// +// Proves `createDrizzleBlockResponseStore`'s notification claim is +// race-safe against a real Postgres connection pool (not one connection +// serializing two concurrent calls) — the exact scenario CL-7192 exists +// to close: a changed answer or a double-click racing the original +// submission for the same (tenant, workbench, message, block, principal), +// where the in-memory store's single-threaded tests can't exercise a real +// lock/serialization path the production store depends on. +import { afterAll, beforeAll, describe, expect, test } from "bun:test"; +import { drizzle } from "drizzle-orm/postgres-js"; +import postgres from "postgres"; + +import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; +import { applyChatMigrations } from "../src/migrations"; +import { createDrizzleBlockResponseStore } from "../src/block-responses"; +import type { BlockResponseKey } from "../src/block-responses"; + +function scratchUrlFor(e2eUrl: string): string { + const url = new URL(e2eUrl); + const database = url.pathname.replace(/^\//, ""); + url.pathname = `/${database}_chat_block_responses_drizzle_test`; + return url.toString(); +} + +const databaseUrl = e2eDatabaseUrl(); +const describeIfDb = databaseUrl === undefined ? describe.skip : describe; + +describeIfDb( + "createDrizzleBlockResponseStore: concurrent notification claim", + () => { + const scratchUrl = scratchUrlFor( + databaseUrl ?? "postgres://localhost:5432/unused", + ); + const scratchTarget = new URL(scratchUrl); + const scratchDatabase = scratchTarget.pathname.replace(/^\//, ""); + + beforeAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe( + `DROP DATABASE IF EXISTS "${scratchDatabase}"`, + ); + await maintenance.unsafe(`CREATE DATABASE "${scratchDatabase}"`); + } finally { + await maintenance.end(); + } + await applyChatMigrations(scratchUrl); + }); + + afterAll(async () => { + const maintenanceUrl = new URL(scratchUrl); + maintenanceUrl.pathname = "/postgres"; + const maintenance = postgres(maintenanceUrl.toString(), { + max: 1, + onnotice: () => undefined, + }); + try { + await maintenance.unsafe( + `DROP DATABASE IF EXISTS "${scratchDatabase}"`, + ); + } finally { + await maintenance.end(); + } + }); + + function keyFor(claimKey: string): BlockResponseKey { + return { + tenantId: "ten_1", + workbenchId: "run_1", + messageId: "m1", + blockId: "blk_question1", + principalId: claimKey, + }; + } + + test("two concurrent claims for the same key never throw, and exactly one wins", async () => { + // `max: 5` — a real connection pool, so the two claims below issue + // genuinely overlapping queries rather than being serialized onto + // one connection before either can race the other. + const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined }); + try { + const store = createDrizzleBlockResponseStore(drizzle(sql)); + const key = keyFor("prn_race_1"); + await store.upsertBlockResponse({ + ...key, + payload: { kind: "question", answer: "Staging" }, + }); + + const [first, second] = await Promise.all([ + store.claimBlockResponseNotification(key), + store.claimBlockResponseNotification(key), + ]); + + const outcomes = [first, second]; + expect(outcomes.filter((outcome) => outcome !== false)).toHaveLength(1); + expect(outcomes.filter((outcome) => outcome === false)).toHaveLength(1); + + const rows = await sql.unsafe( + `SELECT "notified_at", "notification_claim_token" FROM "chat"."block_responses" WHERE "tenant_id" = 'ten_1' AND "workbench_id" = 'run_1' AND "message_id" = 'm1' AND "block_id" = 'blk_question1' AND "principal_id" = 'prn_race_1'`, + ); + expect(rows).toHaveLength(1); + expect(rows[0]?.notified_at).not.toBeNull(); + expect(rows[0]?.notification_claim_token).toBe( + outcomes.find((outcome) => outcome !== false), + ); + } finally { + await sql.end(); + } + }); + + test("a claim already won is not won again by a later call", async () => { + const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined }); + try { + const store = createDrizzleBlockResponseStore(drizzle(sql)); + const key = keyFor("prn_sequential_1"); + await store.upsertBlockResponse({ + ...key, + payload: { kind: "question", answer: "Staging" }, + }); + + const first = await store.claimBlockResponseNotification(key); + const second = await store.claimBlockResponseNotification(key); + + expect(typeof first).toBe("string"); + expect(second).toBe(false); + } finally { + await sql.end(); + } + }); + + test("release with the holder's own token frees the claim for a fresh claim", async () => { + const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined }); + try { + const store = createDrizzleBlockResponseStore(drizzle(sql)); + const key = keyFor("prn_release_1"); + await store.upsertBlockResponse({ + ...key, + payload: { kind: "question", answer: "Staging" }, + }); + + const token = await store.claimBlockResponseNotification(key); + expect(token).not.toBe(false); + await store.releaseBlockResponseNotification(key, token as string); + + const reclaimed = await store.claimBlockResponseNotification(key); + expect(reclaimed).not.toBe(false); + } finally { + await sql.end(); + } + }); + + test("release with a stale token never evicts a live claim", async () => { + const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined }); + try { + const store = createDrizzleBlockResponseStore(drizzle(sql)); + const key = keyFor("prn_stale_token_1"); + await store.upsertBlockResponse({ + ...key, + payload: { kind: "question", answer: "Staging" }, + }); + + const firstToken = await store.claimBlockResponseNotification(key); + expect(firstToken).not.toBe(false); + await store.releaseBlockResponseNotification(key, firstToken as string); + const secondToken = await store.claimBlockResponseNotification(key); + expect(secondToken).not.toBe(false); + expect(secondToken).not.toBe(firstToken); + + // The first (now stale) token must never release the second, + // currently-live claim. + await store.releaseBlockResponseNotification(key, firstToken as string); + const thirdAttempt = await store.claimBlockResponseNotification(key); + expect(thirdAttempt).toBe(false); + } finally { + await sql.end(); + } + }); + }, +); diff --git a/packages/chat/test/migrations.test.ts b/packages/chat/test/migrations.test.ts index 872c9f77b..491a0b8f6 100644 --- a/packages/chat/test/migrations.test.ts +++ b/packages/chat/test/migrations.test.ts @@ -47,6 +47,7 @@ const migrationNames = [ "0024_workbench_launch_sources_digest", "0025_workbench_threads_unique_key", "0026_workbench_threads_delivery_key", + "0027_block_responses_notified_at", ]; describeIfDb("applyChatMigrations", () => {