diff --git a/packages/chat/src/chat-orchestrator.ts b/packages/chat/src/chat-orchestrator.ts index 7c5163725..a67dcb35a 100644 --- a/packages/chat/src/chat-orchestrator.ts +++ b/packages/chat/src/chat-orchestrator.ts @@ -29,6 +29,7 @@ import { toolDoneResult, type ReplyContentBlock, } from "@corbits/agent-events"; +import { reportError } from "@corbits/error-sink"; import type { Memory } from "@corbits/memory"; import { persistedArtifactsForFinalizedTurn, @@ -313,11 +314,9 @@ function flattenReplyText(parts: readonly Part[]): string { /** * Resolves an agent address on the event stream to every chat workbench it is - * a member of, per the durable `workbench_settings` store. Shared by every - * poster below rather than assuming "exactly one workbench": an agent is - * invited to exactly one workbench today by construction, but this resolves - * defensively so a store that ever showed more than one still gets every - * member workbench, not a guess at "the" one. + * a member of, per the durable `workbench_settings` store. An agent can be + * invited to more than one workbench; callers that post a turn's reply must + * still pick the originating room rather than spraying every membership. */ async function resolveMemberWorkbenches( deps: ChatOrchestratorDeps, @@ -451,17 +450,126 @@ type TurnOutcome = | { readonly status: "completed" } | { readonly status: "failed"; readonly error: string }; +type ReplyMember = { + readonly workbenchId: string; + readonly turnId: string | undefined; + readonly childRunId: string | undefined; +}; + +function stringField(value: unknown): string | undefined { + return typeof value === "string" && value !== "" ? value : undefined; +} + +/** + * Hints already on the inbound event for the originating room: the + * dispatch mail's `fromWorkbenchId` / From local-part / `replyTo`, or a + * turn id. The sidecar `agent.event` stream often carries none of these + * (address + payload only); callers must still refuse to spray when + * correlation is missing and more than one running turn exists. + */ +function originatingHintsFromEvent(event: unknown): readonly string[] { + if (typeof event !== "object" || event === null) return []; + const record = event as Record; + const data = + typeof record["data"] === "object" && record["data"] !== null + ? (record["data"] as Record) + : undefined; + const hints: string[] = []; + const take = (value: unknown) => { + const field = stringField(value); + if (field !== undefined) hints.push(field); + }; + take(data?.["workbenchId"]); + take(data?.["fromWorkbenchId"]); + take(data?.["replyTo"]); + take(data?.["from"]); + take(data?.["turnId"]); + take(record["workbenchId"]); + take(record["fromWorkbenchId"]); + take(record["replyTo"]); + take(record["from"]); + take(record["turnId"]); + return hints; +} + +function correlateOriginatingWorkbench( + event: unknown, + members: readonly ReplyMember[], +): string | undefined { + const memberIds = new Set(members.map((member) => member.workbenchId)); + for (const hint of originatingHintsFromEvent(event)) { + if (memberIds.has(hint)) return hint; + const local = localPartOf(hint); + if (memberIds.has(local)) return local; + const byTurnId = members.filter((member) => member.turnId === hint); + if (byTurnId.length === 1) { + const match = byTurnId[0]; + if (match !== undefined) return match.workbenchId; + } + const byChildRunId = members.filter((member) => member.childRunId === hint); + if (byChildRunId.length === 1) { + const match = byChildRunId[0]; + if (match !== undefined) return match.workbenchId; + } + } + return undefined; +} + async function postReply( deps: ChatOrchestratorDeps, pendingDelegationThreads: Map, agentAddress: string, parts: readonly Part[], outcome: TurnOutcome, + event?: unknown, ): Promise { const resolved = await resolveMemberWorkbenches(deps, agentAddress); if (resolved === undefined) return; + const members: ReplyMember[] = []; for (const workbenchId of resolved.workbenchIds) { + const turn = await deps.agentTurns?.findRunningTurn({ + tenantId: resolved.tenantId, + workbenchId, + agentAddress: resolved.roomAddress, + }); + members.push({ + workbenchId, + turnId: turn?.id, + childRunId: turn?.childRunId, + }); + } + const runningIds = members + .filter((member) => member.turnId !== undefined) + .map((member) => member.workbenchId); + const correlated = correlateOriginatingWorkbench(event, members); + const targetIds = + correlated !== undefined + ? [correlated] + : runningIds.length === 1 + ? runningIds + : runningIds.length === 0 && resolved.workbenchIds.length === 1 + ? resolved.workbenchIds + : []; + + if (targetIds.length === 0) { + reportError( + new Error(`dropping ${agentAddress}'s reply: no originating workbench`), + { + operation: "chat.postReply", + tenantId: resolved.tenantId, + agentId: agentAddress, + extra: { + workbenchIds: resolved.workbenchIds, + runningWorkbenchIds: runningIds, + }, + }, + ); + log.error`chat orchestrator: dropping ${agentAddress}'s reply — no originating workbench among ${String(resolved.workbenchIds.length)} membership(s) (${String(runningIds.length)} running)`; + return; + } + + for (const workbenchId of targetIds) { const turn = await deps.agentTurns?.findRunningTurn({ tenantId: resolved.tenantId, workbenchId, @@ -964,9 +1072,16 @@ export function createChatOrchestrator( accumulated !== undefined && accumulated.length > 0 ? accumulated : [{ kind: "text", text: content }]; - void postReply(deps, pendingDelegationThreads, agentAddress, parts, { - status: "completed", - }).catch((cause: unknown) => { + void postReply( + deps, + pendingDelegationThreads, + agentAddress, + parts, + { + status: "completed", + }, + event, + ).catch((cause: unknown) => { log.error`chat orchestrator: failed to post ${agentAddress}'s reply: ${ cause instanceof Error ? cause.message : String(cause) }`; @@ -1022,6 +1137,7 @@ export function createChatOrchestrator( agentAddress, [{ kind: "text", text: noticeContent }], { status: "failed", error: errorMessage }, + event, ).catch((cause: unknown) => { log.error`chat orchestrator: failed to post ${agentAddress}'s turn-drop notice: ${ cause instanceof Error ? cause.message : String(cause) diff --git a/packages/chat/src/workbench-service.ts b/packages/chat/src/workbench-service.ts index ee38afc31..c487b8e05 100644 --- a/packages/chat/src/workbench-service.ts +++ b/packages/chat/src/workbench-service.ts @@ -1139,9 +1139,9 @@ async function routeToRecipients( if (target !== undefined) recipientSet.add(target); } // No mention and no agent reply target: the default-routing case, - // where the host receives every such message unconditionally — the - // same standing relationship a chat's one agent has always had, so - // it needs no re-situating context either. + // where the host receives every such message unconditionally. Assemble + // this-room turn context the same way mention fan-out does, so a host + // shared across rooms is not asked with another room's rows. const isDefaultRouting = recipientSet.size === 0; if (isDefaultRouting) { const host = participants.find((participant) => @@ -1171,7 +1171,7 @@ async function routeToRecipients( ); const contextText = - !isDefaultRouting && recipients.length > 0 + recipients.length > 0 ? await assembleTurnContext({ roomMessages: deps.roomMessages, tenantId: input.tenantId, diff --git a/packages/chat/test/chat-orchestrator.test.ts b/packages/chat/test/chat-orchestrator.test.ts index 711f8c078..d89da3cf3 100644 --- a/packages/chat/test/chat-orchestrator.test.ts +++ b/packages/chat/test/chat-orchestrator.test.ts @@ -1,10 +1,10 @@ // Proves the orchestrator's own wiring: a `connector.reply` on the // shared `agent.event` stream resolves the replying address to its -// folded run, finds every workbench whose participants carry that -// address (defensively — more than one, if the store ever shows it), -// and posts the reply onto each one's own timeline via -// `postRoomMessage`, sent by the replying agent's address and carrying -// its run id. Non-reply events are ignored for posting but still bump +// folded run, finds the originating workbench (the room with a running +// turn, a unique membership, or an explicit from/replyTo/turn hint on +// the event), and posts the reply onto that timeline via +// `postRoomMessage`. A host in two rooms never gets one reply sprayed +// into both. Non-reply events are ignored for posting but still bump // activity; an address the store never produced (no folded run) is // ignored outright; `dispose` stops the subscription. // @@ -274,6 +274,253 @@ describe("createChatOrchestrator", () => { orchestrator.dispose(); }); + test("a reply for room B does not land in room A when the same agent is in both", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + await agentTurns.startTurn({ + tenantId: "ten_1", + workbenchId: "ins_room_b", + agentAddress: "ins_echo1@ten1.workbench.test", + 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", ["ins_echo1@ten1.workbench.test"]), + workbenchRow("ins_room_b", ["ins_echo1@ten1.workbench.test"]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + events, + agentTurns, + claims: fakeClaims(), + approvals: { findByCorrelationId: async () => null }, + }); + + events.emit("agent.event", { + agentAddress: "ins_echo1@ten1.workbench.test", + sessionId: "ses_1", + event: { type: "connector.reply", data: { content: "reply for room B" } }, + }); + + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(room.posted.map((message) => message.workbenchId)).toEqual([ + "ins_room_b", + ]); + expect(room.posted[0]).toMatchObject({ + workbenchId: "ins_room_b", + parts: [{ kind: "text", text: "reply for room B" }], + }); + + orchestrator.dispose(); + }); + + // CL-7172: waitUntilFree is per workbench, so the same host can have a + // running turn in room A and room B at once. Collecting every running + // row would post one connector.reply into both rooms and finishTurn + // both — then the second reply would see no running turns and two + // memberships and silently drop. Originating workbench is singular. + test("a connector.reply does not post into two rooms that both have a running turn for the same agent", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + const agentAddress = "ins_echo1@ten1.workbench.test"; + const turnA = await agentTurns.startTurn({ + tenantId: "ten_1", + workbenchId: "ins_room_a", + agentAddress, + requestMessageIds: ["msg_a"], + }); + const turnB = 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 }, + }); + + events.emit("agent.event", { + agentAddress, + sessionId: "ses_1", + event: { type: "connector.reply", data: { content: "one reply" } }, + }); + + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(room.posted.map((message) => message.workbenchId)).toEqual([]); + expect( + (await agentTurns.getTurn({ tenantId: "ten_1", turnId: turnA.id })) + ?.status, + ).toBe("running"); + expect( + (await agentTurns.getTurn({ tenantId: "ten_1", turnId: turnB.id })) + ?.status, + ).toBe("running"); + + orchestrator.dispose(); + }); + + test("a connector.reply posts nowhere when the same agent is in two rooms and neither has a running turn", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + const events = createSidecarEmitter(); + const orchestrator = createChatOrchestrator({ + db: createFakeDb({ id: "ins_echo1", tenantId: "ten_1" }) as never, + store: { + listWorkbenchSettings: async () => [ + workbenchRow("ins_room_a", ["ins_echo1@ten1.workbench.test"]), + workbenchRow("ins_room_b", ["ins_echo1@ten1.workbench.test"]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + events, + agentTurns, + claims: fakeClaims(), + approvals: { findByCorrelationId: async () => null }, + }); + + events.emit("agent.event", { + agentAddress: "ins_echo1@ten1.workbench.test", + sessionId: "ses_1", + event: { type: "connector.reply", data: { content: "orphan reply" } }, + }); + + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(room.posted.map((message) => message.workbenchId)).toEqual([]); + + orchestrator.dispose(); + }); + + test("a connector.reply still posts when the agent has a single membership and no running turn", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + const events = createSidecarEmitter(); + const orchestrator = createChatOrchestrator({ + db: createFakeDb({ id: "ins_echo1", tenantId: "ten_1" }) as never, + store: { + listWorkbenchSettings: async () => [ + workbenchRow("ins_workbench1", ["ins_echo1@ten1.workbench.test"]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + events, + agentTurns, + claims: fakeClaims(), + approvals: { findByCorrelationId: async () => null }, + }); + + events.emit("agent.event", { + agentAddress: "ins_echo1@ten1.workbench.test", + sessionId: "ses_1", + event: { type: "connector.reply", data: { content: "solo fallback" } }, + }); + + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(room.posted.map((message) => message.workbenchId)).toEqual([ + "ins_workbench1", + ]); + expect(room.posted[0]).toMatchObject({ + workbenchId: "ins_workbench1", + parts: [{ kind: "text", text: "solo fallback" }], + }); + + orchestrator.dispose(); + }); + + test("a connector.reply correlated to room B does not finish room A's running turn", async () => { + const room = fakeRoom(); + const agentTurns = createInMemoryAgentTurnStore(); + const agentAddress = "ins_echo1@ten1.workbench.test"; + const turnA = 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 }, + }); + + events.emit("agent.event", { + agentAddress, + sessionId: "ses_1", + event: { + type: "connector.reply", + data: { + content: "reply for room B", + fromWorkbenchId: "ins_room_b", + }, + }, + }); + + await new Promise((resolve) => setTimeout(resolve, 0)); + + expect(room.posted.map((message) => message.workbenchId)).toEqual([ + "ins_room_b", + ]); + expect( + (await agentTurns.getTurn({ tenantId: "ten_1", turnId: turnA.id })) + ?.status, + ).toBe("running"); + expect( + ( + await agentTurns.findRunningTurn({ + tenantId: "ten_1", + workbenchId: "ins_room_b", + agentAddress, + }) + )?.status, + ).toBeUndefined(); + + 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 @@ -565,41 +812,6 @@ describe("createChatOrchestrator", () => { orchestrator.dispose(); }); - test("posts into every workbench when the store shows the address in more than one", async () => { - const room = fakeRoom(); - const events = createSidecarEmitter(); - const orchestrator = createChatOrchestrator({ - db: createFakeDb({ id: "ins_echo1", tenantId: "ten_1" }) as never, - store: { - listWorkbenchSettings: async () => [ - workbenchRow("ins_workbench1", ["ins_echo1@ten1.workbench.test"]), - workbenchRow("ins_workbench2", ["ins_echo1@ten1.workbench.test"]), - ], - }, - roomMessages: room.roomMessages, - publish: room.publish, - platform: fakeMail().platform, - events, - claims: fakeClaims(), - approvals: { findByCorrelationId: async () => null }, - }); - - events.emit("agent.event", { - agentAddress: "ins_echo1@ten1.workbench.test", - sessionId: "ses_1", - event: { type: "connector.reply", data: { content: "hi" } }, - }); - - await new Promise((resolve) => setTimeout(resolve, 0)); - - expect(room.posted.map((message) => message.workbenchId).sort()).toEqual([ - "ins_workbench1", - "ins_workbench2", - ]); - - orchestrator.dispose(); - }); - // CL-6137 / turn-drop notice (stress round 3): a turn that never // emits `connector.reply` content is no longer invisible — the // bracket close posts an honest in-workbench notice — but a reply this diff --git a/packages/chat/test/workbench-service.test.ts b/packages/chat/test/workbench-service.test.ts index 70fa21472..fc0de3606 100644 --- a/packages/chat/test/workbench-service.test.ts +++ b/packages/chat/test/workbench-service.test.ts @@ -614,7 +614,65 @@ describe("message fan-out", () => { }); }); - test("a chat's fan-out carries no context block, even with a full history", async () => { + test("a default-route turn in room B does not include room A's rows", async () => { + const deps = buildDeps(); + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const { body: roomA } = await createWorkbench(app, { + kind: "workbench", + name: "room A", + participants: ["ins_echo1@acme.example"], + }); + const { body: roomB } = await createWorkbench(app, { + kind: "workbench", + name: "room B", + participants: ["ins_echo1@acme.example"], + }); + + await app.request(`/workbenches/${roomA.id}/messages`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + parts: [{ kind: "text", text: "room A secret: the launch is on Mars" }], + }), + }); + await settleFanout(); + await nextTimelineMoment(); + + await app.request(`/workbenches/${roomB.id}/messages`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + parts: [{ kind: "text", text: "room B note: we meet on Tuesday" }], + }), + }); + await settleFanout(); + await nextTimelineMoment(); + + const response = await app.request(`/workbenches/${roomB.id}/messages`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + parts: [{ kind: "text", text: "hello, no mention here" }], + }), + }); + expect(response.status).toBe(201); + await settleFanout(); + + const platform = deps.platform as ReturnType; + const roomBMail = platform.sentMail.filter( + (mail) => mail.fromWorkbenchId === roomB.id, + ); + const copy = roomBMail[roomBMail.length - 1]; + const [contextPart] = decodeParts(copy?.content ?? { content: "" }); + const contextText = contextPart?.kind === "text" ? contextPart.text : ""; + + expect(contextText).toContain("room B note: we meet on Tuesday"); + expect(contextText).toContain("hello, no mention here"); + expect(contextText).not.toContain("room A secret"); + expect(contextText).not.toContain("the launch is on Mars"); + }); + + test("a chat's default-route fan-out carries this-room context", async () => { const deps = buildDeps({ platform: fakePlatform({ invitable: [{ id: "wfd_echo", name: "Echo" }] }), }); @@ -640,13 +698,14 @@ describe("message fan-out", () => { }, ); expect(response.status).toBe(201); + await settleFanout(); const platform = deps.platform as ReturnType; const fanned = platform.sentMail[platform.sentMail.length - 1]; - expect(fanned?.content).toEqual({ - content: "hello, no mention here", - replyTo: workbench.id, - }); + const [contextPart] = decodeParts(fanned?.content ?? { content: "" }); + const contextText = contextPart?.kind === "text" ? contextPart.text : ""; + expect(contextText).toContain("earlier turn"); + expect(contextText).toContain("hello, no mention here"); }); test("a mention fan-out under the context window carries no dropped-history recap", async () => {