Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 27 additions & 0 deletions docs/CHAT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
114 changes: 114 additions & 0 deletions packages/chat/src/block-responses.test.ts
Original file line number Diff line number Diff line change
@@ -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,
);
});
});
138 changes: 112 additions & 26 deletions packages/chat/src/block-responses.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -37,21 +57,48 @@ 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;
}

export interface BlockResponseStore {
upsertBlockResponse(
input: UpsertBlockResponseInput,
): Promise<BlockResponseRow>;
/**
* 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<string | false>;
/**
* 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<void>;
/**
* Every response on file for one block instance — including every other
* principal's raw payload. Only ever called from inside a route handler
Expand Down Expand Up @@ -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(
Expand All @@ -115,16 +156,11 @@ function blockKey(
export function createInMemoryBlockResponseStore(): BlockResponseStore {
const rows = new Map<string, BlockResponseRow>();
const byBlock = new Map<string, Set<string>>();
const claimTokens = new Map<string, string>();

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 = {
Expand All @@ -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(
Expand All @@ -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),
Expand Down Expand Up @@ -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<string, unknown>,
>(db: BlockResponseDb<TSchema>): BlockResponseStore {
Expand Down Expand Up @@ -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()
Expand Down
8 changes: 8 additions & 0 deletions packages/chat/src/migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
`,
},
];

/**
Expand Down
Loading
Loading