From 07642b3fd73d170e83922916fc3471627d817360 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:32:37 -0700 Subject: [PATCH 1/4] Add a token-scoped notification claim to block responses MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CL-7192: the question-response handler needs a way to guarantee at most one send-and-dispatch per question, even under a changed answer or a double-click that beats the UI's disable. Adds a `notifiedAt` / `notificationClaimToken` pair to `block_responses` (migration 0026): `claimBlockResponseNotification` flips `notifiedAt` from null to now in one guarded UPDATE and hands back a fresh token; a second caller for the same key loses. `releaseBlockResponseNotification` only clears the claim when the token it is given still matches — the same token-scoping `turn-claims.ts` uses for CL-7129, applied here as defensive insurance for a future release call site rather than a fix for any reachable double-release today (this store has no TTL or reaper, so nothing can reassign a claim out from under its own holder). Once a question has been claimed, the upsert never resets the claim on a changed answer: re-answering updates the stored payload but never re-notifies, permanently. `block-responses.test.ts` covers the in-memory store's claim/release semantics; `block-responses.drizzle.test.ts` (DB-gated, mirroring `write-claims.drizzle.test.ts`) proves the guarded UPDATE the drizzle store depends on is actually race-safe against a real Postgres connection pool, not just correct-looking single-threaded JS. --- packages/chat/src/block-responses.test.ts | 114 +++++++++++ packages/chat/src/block-responses.ts | 138 ++++++++++--- packages/chat/src/migrations.ts | 8 + packages/chat/src/schema.ts | 8 + .../chat/test/block-responses.drizzle.test.ts | 187 ++++++++++++++++++ packages/chat/test/migrations.test.ts | 1 + 6 files changed, 430 insertions(+), 26 deletions(-) create mode 100644 packages/chat/src/block-responses.test.ts create mode 100644 packages/chat/test/block-responses.drizzle.test.ts 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/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.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", () => { From cdffd25d35e0af695111b7738ec2b296bcd6572f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:33:48 -0700 Subject: [PATCH 2/4] Fix question-response handler: double dispatch, non-atomic write, inverted ordering Three defects in the question-response handler on the *answer* path (`POST /workbenches/:id/messages/:id/blocks/:id/responses`): - Double dispatch: every successful upsert unconditionally sent the answer into the workbench and dispatched a turn, so a changed answer or a double-click that beat the UI's disable ran a second turn. Gated the send behind `claimBlockResponseNotification`: only the submission that wins the claim ever sends and dispatches. - Non-atomic write: the upsert committed before the send, so a thrown send left the row reading "answered" with the agent never notified. A failed send now releases the claim and returns an explicit `notify_failed` 500 telling the caller their answer was saved and to retry, rather than a silent row/timeline mismatch. - Inverted ordering: the turn was dispatched before the machine-readable `block.response` event was posted, so the agent could read the timeline before its own correlation event existed there. The event is now posted (from the upsert's own returned row, not the request's local payload) before the question branch ever gets a chance to dispatch. Tests cover: a changed answer never double-dispatching, a double-click (concurrent identical submissions) never double-dispatching, and a failed notify releasing the claim so a retried submission still reaches the agent. --- packages/chat/src/routes.ts | 126 +++++++----- .../chat/test/block-responses-routes.test.ts | 184 ++++++++++-------- 2 files changed, 186 insertions(+), 124 deletions(-) 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/test/block-responses-routes.test.ts b/packages/chat/test/block-responses-routes.test.ts index 93c31b257..d41b29213 100644 --- a/packages/chat/test/block-responses-routes.test.ts +++ b/packages/chat/test/block-responses-routes.test.ts @@ -9,6 +9,8 @@ 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, @@ -21,6 +23,27 @@ import { 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 +330,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 () => { From e657336623155c559f9547aff52b6f313d696564 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 09:34:13 -0700 Subject: [PATCH 3/4] Update docs: question answers notify at most once (CL-7192) Documents the notification-claim guard on the question-response handler alongside the existing reactions/pins section it sits next to. --- docs/CHAT.md | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) 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. From af2a8e3de4978046a76b102fb6274fad094f88e2 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 16:31:38 -0700 Subject: [PATCH 4/4] Remove unused settleFanout import from block-responses-routes test. --- packages/chat/test/block-responses-routes.test.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/chat/test/block-responses-routes.test.ts b/packages/chat/test/block-responses-routes.test.ts index d41b29213..fa567c208 100644 --- a/packages/chat/test/block-responses-routes.test.ts +++ b/packages/chat/test/block-responses-routes.test.ts @@ -16,7 +16,6 @@ import { createWorkbench, fakePlatform, mountAs, - settleFanout, TENANT, timelineEvents, timelineOf,