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 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