From 484f9c16446755a65820a01e03d1ef39f6c0e664 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 01:12:04 -0700 Subject: [PATCH 1/2] Add a test for two overlapping turns in two benches resolveMemberWorkbenches returns workbench ids plural by design: an agent can be mid-turn in two benches at once. Prove the orchestrator doesn't let one turn's inference/tool events, reply, or bracket-close observe or consume the other's in-flight state. --- packages/chat/test/chat-orchestrator.test.ts | 161 +++++++++++++++++++ 1 file changed, 161 insertions(+) diff --git a/packages/chat/test/chat-orchestrator.test.ts b/packages/chat/test/chat-orchestrator.test.ts index d89da3cf3..b816f068d 100644 --- a/packages/chat/test/chat-orchestrator.test.ts +++ b/packages/chat/test/chat-orchestrator.test.ts @@ -521,6 +521,167 @@ describe("createChatOrchestrator", () => { orchestrator.dispose(); }); + // CL-7196: `resolveMemberWorkbenches` deliberately returns workbench ids + // plural — one address can be a member of several benches, and each + // bench's turn runs its own sidecar connection with its own + // `sessionId`. Two such turns for the SAME agent, interleaved on the + // shared `agent.event` stream, must never share in-flight state: room + // A's inference/tool events, reply, and bracket-close must never + // observe or consume anything room B's turn accumulated, and vice + // versa — regardless of which turn's bracket closes first. + test("two overlapping turns for one agent in two benches do not observe each other's state", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + const agentAddress = "ins_echo1@ten1.workbench.test"; + await agentTurns.startTurn({ + tenantId: "ten_1", + workbenchId: "ins_room_a", + agentAddress, + requestMessageIds: ["msg_a"], + }); + await agentTurns.startTurn({ + tenantId: "ten_1", + workbenchId: "ins_room_b", + agentAddress, + requestMessageIds: ["msg_b"], + }); + const events = createSidecarEmitter(); + const orchestrator = createChatOrchestrator({ + db: createFakeDb({ id: "ins_echo1", tenantId: "ten_1" }) as never, + store: { + listWorkbenchSettings: async () => [ + workbenchRow("ins_room_a", [agentAddress]), + workbenchRow("ins_room_b", [agentAddress]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + events, + agentTurns, + claims: fakeClaims(), + approvals: { findByCorrelationId: async () => null }, + }); + + // Room B's turn starts thinking first: a text block, then a tool + // call — accumulating in its own turn's bucket. + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_b", + event: { + type: "inference.done", + data: { + turn: { + content: [ + { type: "text", text: "Let me check that." }, + { + type: "tool_call", + id: "call_1", + name: "web_search", + arguments: {}, + }, + ], + }, + }, + }, + }); + // Room A's turn, on a wholly different session, produces its own + // text while room B's turn is still open. + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_a", + event: { + type: "inference.done", + data: { turn: { content: [{ type: "text", text: "Sure, one sec." }] } }, + }, + }); + // Room B's tool call resolves. + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_b", + event: { + type: "tool.done", + data: { result: { callId: "call_1", content: "3 results found" } }, + }, + }); + // Room A replies and closes first. + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_a", + event: { + type: "connector.reply", + data: { content: "Sure, one sec.", fromWorkbenchId: "ins_room_a" }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_a", + event: { + type: "message.run.ended", + data: { status: "completed", fromWorkbenchId: "ins_room_a" }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + + // Room B replies and closes after room A. + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_b", + event: { + type: "connector.reply", + data: { + content: "Here's what I found.", + fromWorkbenchId: "ins_room_b", + }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + events.emit("agent.event", { + agentAddress, + sessionId: "ses_room_b", + event: { + type: "message.run.ended", + data: { status: "completed", fromWorkbenchId: "ins_room_b" }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + + // Room A's own reply carries only its own text — never room B's + // "Let me check that." text or its tool-trace. + const postedInA = room.posted.filter( + (message) => message.workbenchId === "ins_room_a", + ); + expect(postedInA).toHaveLength(1); + expect(postedInA[0]?.parts).toEqual([ + { kind: "text", text: "Sure, one sec." }, + ]); + + // Room B's own reply keeps its full accumulated [text, tool-trace] — + // never swallowed by room A's bracket close, and never replaced by + // just the flattened `content` string. + const postedInB = room.posted.filter( + (message) => message.workbenchId === "ins_room_b", + ); + expect(postedInB).toHaveLength(1); + expect(postedInB[0]?.parts).toEqual([ + { kind: "text", text: "Let me check that." }, + { + kind: "tool-trace", + name: "web_search", + input: {}, + status: "success", + output: "3 results found", + }, + ]); + + // Neither bench got a spurious "I didn't manage to answer that one" + // drop notice from the other turn's flag being stolen. + expect(room.posted).toHaveLength(2); + + orchestrator.dispose(); + }); + // CL-6378: a turn's `inference.done` events already split the model's // output into prose and tool calls (see `event-collector.ts`'s // `handleInferenceDone`), and `tool.done` resolves each call's From 8069b47b438c8ed8f94892908b6a7b6713d7f725 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 04:45:55 -0700 Subject: [PATCH 2/2] Key in-flight turn state by turn, not agent address resolveMemberWorkbenches returns workbench ids plural: one agent address can be mid-turn in several benches at once, each on its own sidecar session. repliedAddresses, notifiedDropAddresses, and the reply-parts accumulator's partsByAddress all keyed on the address alone, so two concurrent turns for the same agent shared one bucket and could consume or drop each other's state. Key all three on agent address + sidecar sessionId instead. --- packages/chat/src/chat-orchestrator.ts | 128 ++++++++++++++----------- 1 file changed, 72 insertions(+), 56 deletions(-) diff --git a/packages/chat/src/chat-orchestrator.ts b/packages/chat/src/chat-orchestrator.ts index a67dcb35a..6deb22877 100644 --- a/packages/chat/src/chat-orchestrator.ts +++ b/packages/chat/src/chat-orchestrator.ts @@ -196,6 +196,21 @@ function gateBlockedCorrelationId(event: unknown): string | undefined { : undefined; } +/** + * Identifies one turn on the shared `agent.event` stream: an agent + * address alone is not enough, since `resolveMemberWorkbenches` (above) + * deliberately returns workbench ids plural — one address can be a member + * of several benches, each with its own turn running at once. Every + * `agent.event` frame carries a `sessionId` regardless of which inner + * event type it wraps (`vendor/intx/hub-sessions/src/ws/sidecar-events.ts`'s + * `SidecarEventMap["agent.event"]`), scoped to the one sidecar connection + * that turn's occurrence runs on — the correlator this module's in-flight + * turn state (below) keys on instead of the address alone. + */ +function turnKeyFor(agentAddress: string, sessionId: string): string { + return `${agentAddress}:${sessionId}`; +} + /** * Turn-scoped reply assembly (CL-6378). A turn's `inference.done` events * already carry the model's output pre-split into prose and tool calls @@ -205,14 +220,13 @@ function gateBlockedCorrelationId(event: unknown): string | undefined { * `Part[]` a reply message posts as, so a tool call the model made * becomes a `ToolTracePart` the UI renders as `ToolBlock`, never JSON * folded into a `TextPart`'s prose. One accumulator per process, keyed by - * agent address (this stream carries no turnId to key on more precisely, - * matching `repliedAddresses`' own per-address bookkeeping above); reset - * the moment a turn's `connector.reply` or turn-drop notice consumes it. + * turn (see `turnKeyFor` below): reset the moment a turn's + * `connector.reply` or turn-drop notice consumes it. */ export function createReplyPartsAccumulator(): { - onInferenceDone(agentAddress: string, blocks: ReplyContentBlock[]): void; + onInferenceDone(turnKey: string, blocks: ReplyContentBlock[]): void; onToolDone( - agentAddress: string, + turnKey: string, result: { callId: string; content: unknown; @@ -220,18 +234,18 @@ export function createReplyPartsAccumulator(): { detail?: unknown; }, ): void; - /** Returns and clears the address's accumulated parts, or undefined if - * nothing was ever accumulated for it this turn. */ - take(agentAddress: string): Part[] | undefined; + /** Returns and clears the turn's accumulated parts, or undefined if + * nothing was ever accumulated for it. */ + take(turnKey: string): Part[] | undefined; } { - const partsByAddress = new Map(); - const toolTraceIndexByAddress = new Map>(); + const partsByTurn = new Map(); + const toolTraceIndexByTurn = new Map>(); return { - onInferenceDone(agentAddress, blocks) { - const parts = partsByAddress.get(agentAddress) ?? []; + onInferenceDone(turnKey, blocks) { + const parts = partsByTurn.get(turnKey) ?? []; const toolTraceIndex = - toolTraceIndexByAddress.get(agentAddress) ?? new Map(); + toolTraceIndexByTurn.get(turnKey) ?? new Map(); for (const block of blocks) { if (block.kind === "text") { parts.push({ kind: "text", text: block.text }); @@ -245,14 +259,12 @@ export function createReplyPartsAccumulator(): { }); } } - partsByAddress.set(agentAddress, parts); - toolTraceIndexByAddress.set(agentAddress, toolTraceIndex); + partsByTurn.set(turnKey, parts); + toolTraceIndexByTurn.set(turnKey, toolTraceIndex); }, - onToolDone(agentAddress, result) { - const parts = partsByAddress.get(agentAddress); - const index = toolTraceIndexByAddress - .get(agentAddress) - ?.get(result.callId); + onToolDone(turnKey, result) { + const parts = partsByTurn.get(turnKey); + const index = toolTraceIndexByTurn.get(turnKey)?.get(result.callId); if (parts === undefined || index === undefined) return; const existing = parts[index]; if (existing === undefined || existing.kind !== "tool-trace") return; @@ -292,10 +304,10 @@ export function createReplyPartsAccumulator(): { }); } }, - take(agentAddress) { - const parts = partsByAddress.get(agentAddress); - partsByAddress.delete(agentAddress); - toolTraceIndexByAddress.delete(agentAddress); + take(turnKey) { + const parts = partsByTurn.get(turnKey); + partsByTurn.delete(turnKey); + toolTraceIndexByTurn.delete(turnKey); return parts; }, }; @@ -997,29 +1009,31 @@ export function createChatOrchestrator( // doc comment for what this does and doesn't cover. const postedApprovalIds = new Set(); - // Every address with a `connector.reply` pending delivery for its - // current turn — added the moment reply content is seen, cleared the - // moment that turn's own `message.run.ended` bracket closes (see - // below), keyed per address rather than per message — this stream - // carries no messageId to correlate on more precisely — though chat - // only needs a presence bit here, never the reply text itself. - const repliedAddresses = new Set(); + // Every turn with a `connector.reply` pending delivery — added the + // moment reply content is seen, cleared the moment that same turn's own + // `message.run.ended` bracket closes (see below). Keyed by `turnKeyFor` + // (agent address + sidecar session id), never by address alone: one + // address can be mid-turn in several benches at once + // (`resolveMemberWorkbenches` returns workbench ids plural by design), + // and each such turn runs its own sidecar connection with its own + // `sessionId` — the correlator every `agent.event` frame carries + // regardless of which inner event type it wraps. + const repliedTurns = new Set(); // Process-lifetime idempotency guard for the turn-drop notice below, - // keyed the same coarse way as `repliedAddresses` (this stream carries - // no turnId to key a redelivered `message.run.ended` on more - // precisely): set the moment a notice is posted for an address's - // silent turn, cleared the moment that address's next `connector.reply` - // arrives OR the next `message.run.started` opens so a genuinely new - // turn is never suppressed by a stale entry. The `message.run.started` - // clear matters independently of a reply landing: two turns in a row - // can each end with zero `connector.reply` (an inference turn that + // keyed the same way as `repliedTurns`: set the moment a notice is + // posted for a turn that ended silently, cleared the moment that same + // turn's next `connector.reply` arrives OR its `message.run.started` + // reopens (a session can host more than one turn across its lifetime, + // and a genuinely new turn on that session must never be suppressed by + // a stale entry the previous turn on it left behind). Two turns in a + // row can each end with zero `connector.reply` (an inference turn that // produced no text is not rare under load — see the notice's own // wording below), and without re-arming on every turn's OPEN, the // second silent turn's own notice would be swallowed by the guard // still set from the first, leaving the user staring at a thread that // looks permanently dead with no trace anywhere it ever tried again. - const notifiedDropAddresses = new Set(); + const notifiedDropTurns = new Set(); // See `PendingDelegationThread`'s own doc comment above `postReply`. const pendingDelegationThreads = new Map(); @@ -1029,15 +1043,17 @@ export function createChatOrchestrator( const unsubscribe = deps.events.on( "agent.event", - ({ agentAddress, event }) => { + ({ agentAddress, sessionId, event }) => { // Any event at all counts as activity, not just `connector.reply` // below — an agent mid-inference must never be undeployed out // from under itself by the idle sweep just because it hasn't // replied yet. deps.recordActivity?.(agentAddress); + const turnKey = turnKeyFor(agentAddress, sessionId); + if (messageRunStarted(event)) { - notifiedDropAddresses.delete(agentAddress); + notifiedDropTurns.delete(turnKey); return; } @@ -1048,17 +1064,17 @@ export function createChatOrchestrator( // an inference step or a tool settling is not itself a reply. const blocks = inferenceDoneBlocks(event); if (blocks !== undefined) { - replyParts.onInferenceDone(agentAddress, blocks); + replyParts.onInferenceDone(turnKey, blocks); } const toolResult = toolDoneResult(event); if (toolResult !== undefined) { - replyParts.onToolDone(agentAddress, toolResult); + replyParts.onToolDone(turnKey, toolResult); } const content = connectorReplyContent(event); if (content !== undefined) { - repliedAddresses.add(agentAddress); - notifiedDropAddresses.delete(agentAddress); + repliedTurns.add(turnKey); + notifiedDropTurns.delete(turnKey); // Prefer this turn's accumulated [text, tool-trace, text, ...] // parts over the flattened `content` string — the string is a // fallback for the rare case nothing was ever accumulated (e.g. @@ -1067,7 +1083,7 @@ export function createChatOrchestrator( // `connector.reply` case). Never wrap `content` verbatim when // structured parts exist, or a leaked JSON-shaped tool call // sitting in that string would still show up as prose. - const accumulated = replyParts.take(agentAddress); + const accumulated = replyParts.take(turnKey); const parts: Part[] = accumulated !== undefined && accumulated.length > 0 ? accumulated @@ -1111,22 +1127,22 @@ export function createChatOrchestrator( // came back out." Beyond the error log, an honest notice now goes // into the workbench itself — a human staring at a stalled thread // during saturated inference (stress round 3) must never see - // nothing at all — guarded by `notifiedDropAddresses` so a - // redelivered `message.run.ended` (sidecar reconnect, wire-layer - // replay) posts the notice once, not once per delivery. + // nothing at all — guarded by `notifiedDropTurns` so a redelivered + // `message.run.ended` (sidecar reconnect, wire-layer replay) posts + // the notice once, not once per delivery. const ended = messageRunEnded(event); if (ended !== undefined) { - // A turn's bracket closed either way — nothing accumulated for - // it (if any) belongs to a future turn, not this one. - replyParts.take(agentAddress); + // This turn's bracket closed either way — nothing accumulated + // for it (if any) belongs to a future turn, not this one. + replyParts.take(turnKey); - const hadReply = repliedAddresses.delete(agentAddress); + const hadReply = repliedTurns.delete(turnKey); if (!hadReply) { const errorMessage = ended.errorMessage ?? "no error reported"; log.error`chat orchestrator: agent ${agentAddress}'s turn ended (${ended.status}) with no reply ever posted to any workbench: ${errorMessage}`; - if (!notifiedDropAddresses.has(agentAddress)) { - notifiedDropAddresses.add(agentAddress); + if (!notifiedDropTurns.has(turnKey)) { + notifiedDropTurns.add(turnKey); const noticeContent = ended.status === "failed" ? errorMessage