From 7c068a23410ea3b3c988a5f79ae58783bfdba2ee Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 20 Aug 2026 20:56:32 -0700 Subject: [PATCH 1/4] Add tests for single-run participant reuse and stale turn expiry MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An @mention or workflow command naming a definition already resident in the room must route into the existing participant's run — never mint a sibling — while an explicit re-invite still deliberately places a second instance. A turn row still running past the occurrence timeout must read back failed instead of zombieing (CL-6451). --- packages/chat/src/agent-turns.test.ts | 51 +++++- .../chat/test/agent-turns.drizzle.test.ts | 40 ++++- packages/chat/test/commands.test.ts | 154 +++++++++++++++++- .../chat/test/start-workflow-command.test.ts | 92 +++++++++++ packages/chat/test/workbench-service.test.ts | 43 +++++ 5 files changed, 377 insertions(+), 3 deletions(-) diff --git a/packages/chat/src/agent-turns.test.ts b/packages/chat/src/agent-turns.test.ts index 7fb74188b..c736bfec0 100644 --- a/packages/chat/src/agent-turns.test.ts +++ b/packages/chat/src/agent-turns.test.ts @@ -1,6 +1,9 @@ import { describe, expect, test } from "bun:test"; -import { createInMemoryAgentTurnStore } from "./agent-turns"; +import { + AGENT_TURN_STALE_MS, + createInMemoryAgentTurnStore, +} from "./agent-turns"; const BASE = { tenantId: "ten_1", @@ -106,4 +109,50 @@ describe("createInMemoryAgentTurnStore", () => { }); expect(listed.map((turn) => turn.id)).toEqual([second.id, first.id]); }); + + // CL-6451: a dispatch the supervisor failed (or a hub that died + // mid-turn) never sends the event that closes its row, so a turn + // still `running` past the occurrence timeout is dead by construction + // — it fails visibly instead of showing "typing" forever. + test("a running turn past the stale cutoff reads back failed, never running", async () => { + let clock = 1_000; + const store = createInMemoryAgentTurnStore({ now: () => clock }); + await store.startTurn(BASE); + + clock += AGENT_TURN_STALE_MS + 1; + + expect(await store.findRunningTurn(BASE)).toBeUndefined(); + const listed = await store.listTurns({ + tenantId: BASE.tenantId, + workbenchId: BASE.workbenchId, + }); + expect(listed[0]?.status).toBe("failed"); + expect(listed[0]?.error).not.toBeNull(); + expect(listed[0]?.endedAt).not.toBeNull(); + }); + + test("a running turn within the stale cutoff stays running", async () => { + let clock = 1_000; + const store = createInMemoryAgentTurnStore({ now: () => clock }); + const opened = await store.startTurn(BASE); + + clock += AGENT_TURN_STALE_MS - 1; + + expect((await store.findRunningTurn(BASE))?.id).toBe(opened.id); + }); + + test("starting a new turn expires a stale predecessor rather than leaving two running", async () => { + let clock = 1_000; + const store = createInMemoryAgentTurnStore({ now: () => clock }); + const stale = await store.startTurn(BASE); + + clock += AGENT_TURN_STALE_MS + 1; + const fresh = await store.startTurn(BASE); + + expect((await store.findRunningTurn(BASE))?.id).toBe(fresh.id); + expect( + (await store.getTurn({ tenantId: BASE.tenantId, turnId: stale.id })) + ?.status, + ).toBe("failed"); + }); }); diff --git a/packages/chat/test/agent-turns.drizzle.test.ts b/packages/chat/test/agent-turns.drizzle.test.ts index f84a9ecc2..f468a8d6a 100644 --- a/packages/chat/test/agent-turns.drizzle.test.ts +++ b/packages/chat/test/agent-turns.drizzle.test.ts @@ -14,7 +14,10 @@ import { drizzle } from "drizzle-orm/postgres-js"; import postgres from "postgres"; import { e2eDatabaseUrl } from "../../../scripts/e2e/harness"; -import { createDrizzleAgentTurnStore } from "../src/agent-turns"; +import { + AGENT_TURN_STALE_MS, + createDrizzleAgentTurnStore, +} from "../src/agent-turns"; import { applyChatMigrations } from "../src/migrations"; function scratchUrlFor(e2eUrl: string): string { @@ -160,4 +163,39 @@ describeIfDb("createDrizzleAgentTurnStore", () => { await sql.end(); } }); + + // CL-6451: a dispatch the supervisor failed never sends the event + // that closes its row — a `running` row past the stale cutoff is dead + // by construction and must read back failed. + test("a running turn past the stale cutoff reads back failed", async () => { + const sql = postgres(scratchUrl, { max: 5, onnotice: () => undefined }); + try { + let clock = Date.now(); + const store = createDrizzleAgentTurnStore(drizzle(sql), { + now: () => clock, + }); + const input = { + tenantId: TENANT, + workbenchId: "run_stale", + agentAddress: AGENT, + requestMessageIds: ["msg_stale"], + }; + const opened = await store.startTurn(input); + + // `startedAt` is stamped by the database, not the injected clock, + // so age the clock from a fresh reading with a wide margin. + clock = Date.now() + AGENT_TURN_STALE_MS + 60_000; + + expect(await store.findRunningTurn(input)).toBeUndefined(); + const read = await store.getTurn({ + tenantId: TENANT, + turnId: opened.id, + }); + expect(read?.status).toBe("failed"); + expect(read?.error).not.toBeNull(); + expect(read?.endedAt).not.toBeNull(); + } finally { + await sql.end(); + } + }); }); diff --git a/packages/chat/test/commands.test.ts b/packages/chat/test/commands.test.ts index c9ad595a1..d3c65f021 100644 --- a/packages/chat/test/commands.test.ts +++ b/packages/chat/test/commands.test.ts @@ -5,7 +5,12 @@ // participant's handle, so an ordinary agent mention is untouched. import { describe, expect, test } from "bun:test"; import { createChatRoutes } from "../src/routes"; -import { createCommandRegistry } from "@corbits/commands"; +import { + createCommandRegistry, + createWorkflowCommandPlugin, +} from "@corbits/commands"; +import { createInMemoryAgentTurnStore } from "../src/agent-turns"; +import { startWorkflowCommand } from "../src/workbench-service"; import { buildDeps, createWorkbench, @@ -14,8 +19,41 @@ import { sendText, timelineOf, timelineTexts, + TENANT, } from "./test-support"; +/** + * The full workflow-command wiring the hub composes: the registrar's + * commands are the tenant's invitable definitions, dispatching through + * `startWorkflowCommand` against the same store/platform the routes use. + */ +function buildWorkflowCommandDeps( + platform: ReturnType, +): ReturnType & { + agentTurns: ReturnType; +} { + const registry = createCommandRegistry(); + const agentTurns = createInMemoryAgentTurnStore(); + const deps = buildDeps({ commands: registry, platform, agentTurns }); + registry.registerCommandPlugin( + createWorkflowCommandPlugin({ + listInvitableDefinitions: (tenantId) => + platform.listInvitableDefinitions(tenantId), + startWorkflow: (input) => + startWorkflowCommand( + { + store: deps.store, + platform, + roomMessages: deps.roomMessages, + publish: () => undefined, + }, + input, + ), + }), + ); + return Object.assign(deps, { agentTurns }); +} + describe("workbench command dispatch", () => { test("without an injected registry, a leading slash is posted as an ordinary message", async () => { const deps = buildDeps(); @@ -136,4 +174,118 @@ describe("workbench command dispatch", () => { text: "starting: do the thing", }); }); + + // CL-6451: the participant's mention handle derives from the + // definition's display name ("Myra"), while the workflow command is + // named after the definition's wire name ("assistant") — so the + // known-handle guard alone cannot see that `@assistant` names an agent + // already in the room, and the command path used to mint a SECOND run + // for the same participant. + test("an @name naming an already-resident definition routes to the existing run, never a second one", async () => { + const platform = fakePlatform({ + invitable: [ + { id: "wfd_assistant", name: "assistant", description: "Myra" }, + ], + resolveDefinitionIdByAddress: async (address) => + address === "ins_invited1@acme.example" ? "wfd_assistant" : undefined, + }); + const deps = buildWorkflowCommandDeps(platform); + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const { body: workbench } = await createWorkbench(app, { + kind: "workbench", + }); + + const invite = await app.request(`/workbenches/${workbench.id}/invite`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ definitionId: "wfd_assistant" }), + }); + expect(invite.status).toBe(201); + expect(platform.launchInviteCalls).toHaveLength(1); + + const response = await sendText( + app, + workbench.id, + "@assistant set up a sales workbench", + ); + expect(response.status).toBe(201); + const body = (await response.json()) as Record; + // Not a command dispatch: the message posts as an ordinary message. + expect(body["command"]).toBeUndefined(); + // One participant, one run: no second launch. + expect(platform.launchInviteCalls).toHaveLength(1); + + // The message reached the EXISTING participant's run... + const delivered = platform.sentMail.filter( + (mail) => mail.workbenchId === "ins_invited1", + ); + expect(delivered).toHaveLength(1); + expect(delivered[0]?.content.content).toContain("set up a sales workbench"); + // ...through the ordinary turn pipeline: one turn row, on the same + // occurrence sequence every later turn of this participant rides. + const turns = await deps.agentTurns.listTurns({ + tenantId: TENANT.id, + workbenchId: workbench.id, + }); + expect(turns).toHaveLength(1); + expect(turns[0]).toMatchObject({ + agentAddress: "ins_invited1@acme.example", + childRunId: "turn__0", + status: "running", + }); + // The user's words stay on the timeline, never swallowed by the + // command intercept. + expect(timelineTexts(await timelineOf(deps, workbench.id))).toEqual([ + "@assistant set up a sales workbench", + ]); + + // CL-6453: the next mention rides the SAME run's occurrence + // sequence — turn__0 then turn__1 on one section run — so the + // first exchange lives in the same stepId-keyed history every later + // turn restores. + await sendText(app, workbench.id, "@assistant continue"); + const afterSecond = await deps.agentTurns.listTurns({ + tenantId: TENANT.id, + workbenchId: workbench.id, + }); + expect( + afterSecond + .map((turn) => ({ + agentAddress: turn.agentAddress, + childRunId: turn.childRunId, + })) + .sort((a, b) => a.childRunId.localeCompare(b.childRunId)), + ).toEqual([ + { + agentAddress: "ins_invited1@acme.example", + childRunId: "turn__0", + }, + { + agentAddress: "ins_invited1@acme.example", + childRunId: "turn__1", + }, + ]); + expect(platform.launchInviteCalls).toHaveLength(1); + }); + + test("an @name for a definition NOT in the room still starts it as a command", async () => { + const platform = fakePlatform({ + invitable: [ + { id: "wfd_assistant", name: "assistant", description: "Myra" }, + ], + }); + const deps = buildWorkflowCommandDeps(platform); + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const { body: workbench } = await createWorkbench(app, { + kind: "workbench", + }); + + const response = await sendText(app, workbench.id, "@assistant hello"); + expect(response.status).toBe(201); + const body = (await response.json()) as { + command: { type: string; handle: string }; + }; + expect(body.command.type).toBe("workflow-started"); + expect(platform.launchInviteCalls).toHaveLength(1); + }); }); diff --git a/packages/chat/test/start-workflow-command.test.ts b/packages/chat/test/start-workflow-command.test.ts index 7001ea550..c4135d1ef 100644 --- a/packages/chat/test/start-workflow-command.test.ts +++ b/packages/chat/test/start-workflow-command.test.ts @@ -106,4 +106,96 @@ describe("startWorkflowCommand", () => { ), ).rejects.toThrow(/No workbench "chan_missing"/); }); + + // CL-6451: one room participant = one live run. A workflow command + // naming a definition already resident in the room delivers into the + // existing participant's run instead of minting a sibling. + test("reuses an existing participant launched from the same definition row", async () => { + const store = createInMemoryChatStore(); + const platform = fakePlatform({ + invitable: [{ id: "wfd_echo", name: "echo", description: "Myra" }], + resolveDefinitionIdByAddress: async (address) => + address === "ins_existing@acme.example" ? "wfd_echo" : undefined, + }); + await store.createWorkbenchSettings({ + tenantId: TENANT.id, + workbenchId: "chan_3", + settings: { + "chat/kind": "workbench", + "chat/participants": [ + { address: "ins_existing@acme.example", handle: "myra" }, + ], + }, + updatedBy: "prn_alice", + }); + + const result = await startWorkflowCommand( + { + store, + platform, + roomMessages: createInMemoryRoomMessageStore(), + publish: () => undefined, + }, + { + tenantId: TENANT.id, + principalId: "prn_alice", + workbenchId: "chan_3", + definitionId: "wfd_echo", + args: "again please", + }, + ); + + expect(platform.launchInviteCalls).toHaveLength(0); + expect(result).toEqual({ + handle: "myra", + address: "ins_existing@acme.example", + }); + const opening = platform.sentMail.find( + (mail) => mail.workbenchId === "ins_existing", + ); + expect(opening?.content.content).toBe("again please"); + }); + + test("reuses an existing participant whose recorded definition row differs but projects the same asset", async () => { + const store = createInMemoryChatStore(); + const platform = fakePlatform({ + invitable: [{ id: "wfd_echo_v2", name: "echo" }], + resolveDefinitionIdByAddress: async (address) => + address === "ins_existing@acme.example" ? "wfd_echo_v1" : undefined, + resolveDefinitionAssetId: async (definitionId) => + definitionId === "wfd_echo_v1" || definitionId === "wfd_echo_v2" + ? "ast_echo" + : undefined, + }); + await store.createWorkbenchSettings({ + tenantId: TENANT.id, + workbenchId: "chan_4", + settings: { + "chat/kind": "workbench", + "chat/participants": [ + { address: "ins_existing@acme.example", handle: "echo" }, + ], + }, + updatedBy: "prn_alice", + }); + + const result = await startWorkflowCommand( + { + store, + platform, + roomMessages: createInMemoryRoomMessageStore(), + publish: () => undefined, + }, + { + tenantId: TENANT.id, + principalId: "prn_alice", + workbenchId: "chan_4", + definitionId: "wfd_echo_v2", + args: "hi", + }, + ); + + expect(platform.launchInviteCalls).toHaveLength(0); + expect(result.address).toBe("ins_existing@acme.example"); + }); }); diff --git a/packages/chat/test/workbench-service.test.ts b/packages/chat/test/workbench-service.test.ts index af5c7f059..06192d249 100644 --- a/packages/chat/test/workbench-service.test.ts +++ b/packages/chat/test/workbench-service.test.ts @@ -981,6 +981,49 @@ describe("POST /workbenches/:id/invite", () => { ]); }); + // The explicit invite affordance deliberately CAN place a second + // instance of one definition in a room — proven live by the CL-6329 + // turn-swap proof and the reason handle de-duplication exists. + // CL-6451's anti-sibling rule lives on the command/mention paths, + // never here. + test("an explicit re-invite of the same definition mints a second instance", async () => { + let launches = 0; + const platform = fakePlatform({ + invitable: [{ id: "wfd_echo", name: "echo" }], + launchInvite: async () => { + launches += 1; + return { + instanceId: `ins_invited${launches}`, + address: `ins_invited${launches}@acme.example`, + }; + }, + }); + const deps = buildDeps({ platform }); + const app = mountAs(createChatRoutes(deps), "prn_alice"); + const { body: workbench } = await createWorkbench(app, { + kind: "workbench", + }); + + const invite = () => + app.request(`/workbenches/${workbench.id}/invite`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ definitionId: "wfd_echo" }), + }); + expect((await invite()).status).toBe(201); + expect((await invite()).status).toBe(201); + + expect(platform.launchInviteCalls).toHaveLength(2); + const settingsRow = await deps.store.getWorkbenchSettings( + TENANT.id, + workbench.id, + ); + expect(settingsRow?.settings["chat/participants"]).toEqual([ + { address: "ins_invited1@acme.example", handle: "echo" }, + { address: "ins_invited2@acme.example", handle: "echo-2" }, + ]); + }); + test("appends onto an existing participant list rather than replacing it", async () => { const deps = buildDeps(); const app = mountAs(createChatRoutes(deps), "prn_alice"); From be0c95449b059e89b8e4a0c0261afc65d13cc1ef Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 20 Aug 2026 20:56:32 -0700 Subject: [PATCH 2/4] One room participant, one live run: command paths reuse resident agents; stale turns fail MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An @name resolving to a workflow command is now checked against the definitions the room's agents were launched from — a resident match routes the message into the existing participant's run through the ordinary turn pipeline (queueing behind an in-flight turn) instead of dispatching a launch, and startWorkflowCommand applies the same rule to slash/command dispatch. The explicit invite affordance still always launches: a deliberate second instance of one definition stays possible, which is what handle de-duplication exists for. Turn rows still running past the per-occurrence timeout plus grace are failed on the store's next read or write: a dispatch the supervisor's terminal-or-park backstop killed (or a hub death mid-turn) never sends the closing event, and previously left the turn — and its typing indicator — running forever. Fixes CL-6451 --- packages/chat/src/agent-turns.ts | Bin 11698 -> 14154 bytes packages/chat/src/routes.ts | 53 ++++++++++++++-- packages/chat/src/workbench-service.ts | 82 +++++++++++++++++++++++++ 3 files changed, 130 insertions(+), 5 deletions(-) diff --git a/packages/chat/src/agent-turns.ts b/packages/chat/src/agent-turns.ts index 8010ad216af249fb4d23d0e8a17c9961ffc518ad..04d0a505da1eb566509869672e7128b1d48c0573 100644 GIT binary patch delta 2251 zcmb7FL2nyH6jqv6Gz*1Pm7)kr`(mlE8|+OGAqN{Lm=Zyxwt?6YNFl`I-C29cdS*2< z8z(B3EA_;cGE8$G(|$oTPOr?7>>^oA=)LzVCaxk9L3MKmRg) zA=cTW>4W!JHxW|9eYpK;t+C(O*}k{m*u1;3wbR(YTZcYGQgoreBQVrkkMfjpaIS^>dlyk=0O^~7oN{d86i{b8t(by)`%YN!=xzE({S{(SHQj)CA!8Kqeq!)PQ(h|<6`Bypflei>gDIxLzS|m6y zyg@`*j44bPK}b``jxPm{r66P&>awsjFl2Z@cxdB6ERz+h5D0A2_Dz+K zgzgw&B(h^#ARNSs11mrRu?LfwbfiI25@KiQNv0`Ip|pJvSzS2 z1D#PfLZ89bFFM)7*>nJC2}3!SM#!T)-~3x&hmUdX%ssXMIB5bX6_d2CNlZ;;<$7f} z6fCp@pVSLdEusiZVmbmYP13aD0QR?RxH*bl78b2cdmtj~MQ>CUfJL}9*w0s*Xk5I7 zY848jl`ME>45TcjSkzwTxW$jlr06))$XL@Jtbs*1h@Qk# z2RV*16!HZvoJMfXg)GCRW3NLc zHM^T%_*!zE?4M%cQqIH8(pITKJ{xY=&c}m^}1o&YN$M*TN9r?kfBI ztv9E3vj^Vw2|;gf+|yJj!|ooI%!b}-Io78H;6`YYFHbRFo#;&G@vWsZ_vu{D>CBnaTAck&`>b^e6u}W!rpR%%5qq zwqynKksd*(KrA2wcCHY0E3e_-elmCcuZQi40 m!VWZ}X0nE%>SR4*fyqBjw88WqZJy1%#^Q`%3;a!fGXVewk2g^O diff --git a/packages/chat/src/routes.ts b/packages/chat/src/routes.ts index 8ee67e636..9dcab235d 100644 --- a/packages/chat/src/routes.ts +++ b/packages/chat/src/routes.ts @@ -66,6 +66,7 @@ import { isRecentlyActive } from "./workbench-activity"; import { postRoomMessage, type RoomMessageStore } from "./room-messages"; import type { ConnectGithubBlockData } from "./blocks"; import { + findResidentAgentForDefinition, postCannedGreeting, joinHumanParticipant, launchAndJoinAgent, @@ -777,6 +778,16 @@ async function enrichWithReactionsAndPins( }); } +/** + * What the command intercept decided about an incoming message: + * `command` carries a dispatched command's result; `routeToParticipant` + * says the `@name` named a definition already resident in the room, so + * the message must post normally and its turn must reach that + * participant's existing run (CL-6451) — never a freshly minted one. + */ +type WorkbenchCommandDecision = + { readonly command: CommandResult } | { readonly routeToParticipant: string }; + /** * Decides whether an incoming workbench message opens the command path * at all, and if so, dispatches it. `undefined` — the caller's cue to @@ -784,6 +795,14 @@ async function enrichWithReactionsAndPins( * neither slash- nor `@`-shaped; or an `@name` that names an existing * agent participant's handle rather than a command (mention fan-out * keeps owning that case exactly as before this rollout). + * + * The handle check alone cannot keep "one room participant = one live + * run" (CL-6451): a participant's mention handle derives from its + * definition's display name ("Myra"), while the workflow command is + * named after the definition's wire name ("assistant") — so an `@name` + * that resolves to a command is ALSO checked against the definitions + * the room's agents were launched from, and a resident match routes to + * that participant instead of dispatching a launch. */ async function dispatchWorkbenchCommand( deps: CreateChatRoutesDeps, @@ -793,7 +812,7 @@ async function dispatchWorkbenchCommand( workbenchId: string; text: string; }, -): Promise { +): Promise { if (deps.commands === undefined) return undefined; const ctx = { tenantId: input.tenantId, @@ -802,7 +821,8 @@ async function dispatchWorkbenchCommand( }; if (input.text.startsWith("/")) { - return dispatchSlashCommand(deps.commands, input.text, ctx); + const result = await dispatchSlashCommand(deps.commands, input.text, ctx); + return result === undefined ? undefined : { command: result }; } if (input.text.startsWith("@")) { @@ -824,7 +844,25 @@ async function dispatchWorkbenchCommand( ); if (namesKnownHandle) return undefined; - return dispatchAtCommand(deps.commands, input.text, ctx); + const invitable = await deps.platform.listInvitableDefinitions( + input.tenantId, + ); + const commandDefinition = invitable.find( + (definition) => definition.name === resolved.name, + ); + if (commandDefinition !== undefined) { + const resident = await findResidentAgentForDefinition( + deps.platform, + participants, + commandDefinition.id, + ); + if (resident !== undefined) { + return { routeToParticipant: resident.address }; + } + } + + const result = await dispatchAtCommand(deps.commands, input.text, ctx); + return result === undefined ? undefined : { command: result }; } return undefined; @@ -2088,13 +2126,14 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { // resolving it against the registry only runs once it is // confirmed not to name a known handle, so that mention keeps // its ordinary fan-out behavior exactly as before. - const commandResult = await dispatchWorkbenchCommand(deps, { + const commandDecision = await dispatchWorkbenchCommand(deps, { tenantId: ownerTenantId, principalId: principal.id, workbenchId, text: textOf(messageParts), }); - if (commandResult !== undefined) { + if (commandDecision !== undefined && "command" in commandDecision) { + const commandResult = commandDecision.command; const resultText = textForCommandResult(commandResult); if (resultText !== undefined) { await postRoomMessage( @@ -2136,6 +2175,10 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { ...(parsed.inReplyToMessageId !== undefined ? { inReplyToMessageId: parsed.inReplyToMessageId } : {}), + ...(commandDecision !== undefined && + "routeToParticipant" in commandDecision + ? { forcedRecipientAddress: commandDecision.routeToParticipant } + : {}), }, ); deps.onMessageFanout?.(sent.fanoutDelivered); diff --git a/packages/chat/src/workbench-service.ts b/packages/chat/src/workbench-service.ts index b19e6a0c4..6a1019bb6 100644 --- a/packages/chat/src/workbench-service.ts +++ b/packages/chat/src/workbench-service.ts @@ -174,6 +174,42 @@ export type LaunchAndJoinAgentResult = { readonly joinEventDelivered: Promise; }; +/** + * The agent participant this room already holds for `definitionId`, or + * undefined when none of the room's agents was launched from it. One + * room participant = one live run (CL-6451): every path that could + * start a definition in a room checks residency here first, so a + * mention, a workflow command, or a repeated invite reaches the run the + * room already has instead of minting a sibling. Identity is the + * definition's ASSET when both sides resolve one (a code-sourced deploy + * projects a fresh definition row per wire projection over the same + * asset — see `resolveDefinitionAssetId`), falling back to row-id + * equality when either side's asset is unresolvable. + */ +export async function findResidentAgentForDefinition( + platform: Pick< + WorkbenchLauncher, + "resolveDefinitionIdByAddress" | "resolveDefinitionAssetId" + >, + participants: readonly ParticipantRecord[], + definitionId: string, +): Promise { + const assetId = await platform.resolveDefinitionAssetId(definitionId); + for (const participant of participants) { + if (!isAgentAddress(participant.address)) continue; + const launchedFrom = await platform.resolveDefinitionIdByAddress( + participant.address, + ); + if (launchedFrom === undefined) continue; + if (launchedFrom === definitionId) return participant; + if (assetId === undefined) continue; + const launchedFromAssetId = + await platform.resolveDefinitionAssetId(launchedFrom); + if (launchedFromAssetId === assetId) return participant; + } + return undefined; +} + /** * The invite core: launches the definition's own instance, derives * its friendly mention handle, appends the participant record, posts @@ -181,6 +217,14 @@ export type LaunchAndJoinAgentResult = { * bridge. Shared by `POST .../invite` and chat creation (a chat's * single agent is invited exactly this way, at creation) so the two * paths can never drift. + * + * Always launches: an explicit invite deliberately CAN place a second + * instance of one definition in a room (that is what handle + * de-duplication — "echo", "echo-2" — exists for). "One room + * participant = one live run" (CL-6451) is enforced where the sibling + * was never asked for: the message pipeline's command intercept and + * `startWorkflowCommand` resolve residency via + * `findResidentAgentForDefinition` before ever reaching this launch. */ export async function launchAndJoinAgent( deps: LaunchAndJoinAgentDeps, @@ -696,6 +740,13 @@ export type StartWorkflowCommandResult = { * mailbox. An empty invocation ("/echo" with nothing after it) still * starts the run, mirroring corbits-code's own workflow dispatch: no * args is "Continue.", not "nothing to do". + * + * One room participant = one live run (CL-6451): a command naming a + * definition already resident in the room delivers into the existing + * participant's run — the same anti-sibling rule the message + * pipeline's `@name` intercept enforces — instead of launching again. + * A deliberate second instance stays possible through the explicit + * invite affordance, which always launches. */ export async function startWorkflowCommand( deps: StartWorkflowCommandDeps, @@ -711,6 +762,23 @@ export async function startWorkflowCommand( ); } + const resident = await findResidentAgentForDefinition( + deps.platform, + participantsOf(existing.settings), + input.definitionId, + ); + if (resident !== undefined) { + const text = input.args.trim() !== "" ? input.args.trim() : "Continue."; + await deps.platform.sendMail({ + tenantId: input.tenantId, + workbenchId: localPartOf(resident.address), + principalId: input.principalId, + content: encodeParts([{ kind: "text", text }]), + fromWorkbenchId: input.workbenchId, + }); + return { handle: resident.handle, address: resident.address }; + } + const joined = await launchAndJoinAgent( { store: deps.store, @@ -787,6 +855,17 @@ export type SendWorkbenchMessageInput = { * mentions nobody — the reply gesture is itself an address. */ readonly inReplyToMessageId?: string; + /** + * A participant this message must reach regardless of what its text + * mentions (CL-6451): an `@name` typed as a definition's wire name + * ("assistant") resolves to a participant whose handle derives from + * its display name ("myra"), so mention matching alone would miss it. + * The command intercept resolves that residency and passes the + * participant's address here, so the message rides the ordinary turn + * pipeline — queueing behind an in-flight turn like any mention — + * into the run the room already has. + */ + readonly forcedRecipientAddress?: string; }; export type SendWorkbenchMessageResult = { @@ -951,6 +1030,9 @@ async function routeToRecipients( const recipientSet = new Set( mentionedParticipants(input.messageParts, participants), ); + if (input.forcedRecipientAddress !== undefined) { + recipientSet.add(input.forcedRecipientAddress); + } if (input.inReplyToMessageId !== undefined) { const target = await replyTargetAgent(deps.roomMessages, { tenantId: input.tenantId, From 6eb585a262c9036a2faf023ebad89fa2035db3e0 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 20 Aug 2026 20:56:32 -0700 Subject: [PATCH 3/4] Update docs: single-run command routing and stale-turn failing --- docs/CHAT.md | 23 +++++++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/docs/CHAT.md b/docs/CHAT.md index f07ede96f..1091c2339 100644 --- a/docs/CHAT.md +++ b/docs/CHAT.md @@ -157,6 +157,15 @@ the runtime runs occurrences, older rows can be left `running`. Closing that gap properly needs the runtime to report the occurrence id on the event, not a heuristic here. +**Stale turns fail, never zombie (CL-6451).** A dispatch can die without +any closing event ever reaching the hub — the workflow-host supervisor's +terminal-or-park backstop failing it, or the process dying mid-turn. An +occurrence cannot legitimately outlive the per-turn timeout the section +body enforces, so both turn stores fail any row still `running` past +that timeout plus a settle grace (`AGENT_TURN_STALE_MS`) on their next +read or write: the room shows a failed turn instead of typing forever, +and the reply path can never attribute a later reply to a dead row. + Proved live end to end by `scripts/e2e/cl-6329-turn-swap-proof.ts`: two agents replying in one room under distinct occurrences, three rapid messages serializing into ordered turns, and a sidecar killed @@ -205,6 +214,20 @@ unreadable instance id. Handles are derived from a definition's name at invite time and de-duplicated against every handle already in the workbench (`echo`, `echo-2`, `echo-3`, ...). +**One room participant = one live run (CL-6451).** The `/name` / `@name` +workflow commands resolve residency first +(`findResidentAgentForDefinition`, comparing by the definition's asset so +re-deployed rows still match): a definition already resident reaches the +run the room already has, never a freshly minted sibling that would then +race the original for every message. An `@name` typed as a definition's +wire name (`@assistant`) whose participant answers to a display-name +handle (`@myra`) posts as an ordinary message and rides the normal turn +pipeline into that participant's run, queueing behind an in-flight turn +like any mention. The explicit invite affordance is the one deliberate +way to place a second instance of a definition in a room — it always +launches, which is exactly what handle de-duplication ("echo", "echo-2") +exists for. + A **mention** is `@` followed by a participant's handle at a word boundary, anywhere in a message's text. Mentioning an agent participant triggers **fan-out**: the server sends that agent a single-recipient copy of the From b50c6b6de7873a54cfa1b1b182ca748d6c5091ba Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Thu, 20 Aug 2026 21:06:55 -0700 Subject: [PATCH 4/4] Add CL-6451/CL-6453 live proof: one run per mention, history across turns MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit One real stack (scratch DB, real Ollama): invite 'assistant' (handle 'myra'), send '@assistant' — the message is never intercepted as a command, the room keeps exactly one agent participant, and the reply lands as turn__0 of the invited run. A second '@assistant' answers as turn__1 of the same run and recalls a fact from turn__0, proving the bootstrap exchange lives in the same stepId-keyed durable history. --- scripts/e2e/cl-6451-single-run-proof.ts | 566 ++++++++++++++++++++++++ 1 file changed, 566 insertions(+) create mode 100644 scripts/e2e/cl-6451-single-run-proof.ts diff --git a/scripts/e2e/cl-6451-single-run-proof.ts b/scripts/e2e/cl-6451-single-run-proof.ts new file mode 100644 index 000000000..076a6a87f --- /dev/null +++ b/scripts/e2e/cl-6451-single-run-proof.ts @@ -0,0 +1,566 @@ +// CL-6451's live proof, on one real stack: scratch database, real +// signup, real Ollama, nothing mocked. The bug it proves fixed: a +// participant's mention handle derives from its definition's display +// name ("myra"), while the workflow command registrar names commands +// after the definition's wire name ("assistant") — so `@assistant` +// used to fall past the known-handle guard into the command path and +// mint a SECOND run for an agent already in the room. Messages then +// fanned out to both runs, the original never parked, and the dispatch +// died at the supervisor's terminal-or-park backstop with the turn row +// stuck `running`. +// +// 1. Invite `assistant` into a workbench (run A, handle "myra"), +// then send `@assistant ...` — the message posts as an ordinary +// message (no command result), the room still has exactly ONE +// agent participant, and the reply arrives as occurrence +// `turn__0` of run A. +// 2. (CL-6453) A second `@assistant` message answers from the SAME +// run as `turn__1` and recalls a fact stated in turn 1 — the +// bootstrap exchange lives in the same stepId-keyed durable +// history every later turn restores. +// +// Shares the boot/seed/mint scaffold with `cl-6329-turn-swap-proof.ts`. +// +// Usage: +// E2E_PROVIDER=ollama OLLAMA_BASE_URL=http://localhost:11434 \ +// E2E_OLLAMA_MODEL=llama3.1:8b \ +// DATABASE_URL=postgres://localhost:5432/wb6451proof_e2e \ +// bun run scripts/e2e/cl-6451-single-run-proof.ts +import { expect } from "bun:test"; + +import { mkdtemp } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join as pathJoin } from "node:path"; + +import { resetSchema, setupDatabase } from "../db-setup.ts"; +import { + createGitWorkflowPusher, + createHubAPI, + DEFAULT_WORKFLOWS, + seedTenant, + type ApiCall, +} from "../../packages/hub-client/src/index.ts"; +import { + findPersonalTenant, + testAndPersistCredential, + ensureSeeded, +} from "../../packages/onboarding/src/complete-credential.ts"; +import { + OLLAMA_PLACEHOLDER_SECRET, + ollamaOpenAICompatBaseURL, +} from "../../packages/hub-client/src/credential-test.ts"; +import { + api, + e2eDatabaseUrl, + expectStatus, + freePort, + hop, + provisionSidecar, + startHub, + startSidecar, + type HubHandle, + type SpawnedApp, +} from "./harness.ts"; + +const configuredDatabaseUrl = e2eDatabaseUrl(); +if (configuredDatabaseUrl === undefined) { + throw new Error( + "cl-6451-single-run-proof: DATABASE_URL is not set. This suite proves a " + + "real boot and has nothing honest to assert without one.", + ); +} +const databaseUrl: string = configuredDatabaseUrl; + +const OLLAMA_BASE_URL = process.env["OLLAMA_BASE_URL"]; +if (process.env["E2E_PROVIDER"] !== "ollama" || OLLAMA_BASE_URL === undefined) { + throw new Error( + "cl-6451-single-run-proof: set E2E_PROVIDER=ollama and OLLAMA_BASE_URL. " + + "The proof requires a real completion model actually answering.", + ); +} +const ollamaBaseUrl = OLLAMA_BASE_URL; + +const OLLAMA_MODEL = process.env["E2E_OLLAMA_MODEL"]; +if (OLLAMA_MODEL === undefined || OLLAMA_MODEL === "") { + throw new Error( + "cl-6451-single-run-proof: set E2E_OLLAMA_MODEL to a completion model " + + "the instance at OLLAMA_BASE_URL actually serves (see `ollama list`).", + ); +} +const proofModelSource = { + provider: "openai-compatible", + model: OLLAMA_MODEL, + baseURL: ollamaOpenAICompatBaseURL(ollamaBaseUrl), + apiKey: OLLAMA_PLACEHOLDER_SECRET, +} as const; + +const TURN_TIMEOUT_MS = 300_000; + +const tracked: SpawnedApp[] = []; +const tempDir = (prefix: string) => mkdtemp(pathJoin(tmpdir(), prefix)); +const track = (app: SpawnedApp) => { + tracked.push(app); +}; +process.on("exit", () => { + for (const a of tracked) { + try { + void a.stop(); + } catch { + // Best-effort teardown: a child already gone is fine. + } + } +}); + +function stringField(data: unknown, field: string, what: string): string { + if (typeof data === "object" && data !== null && field in data) { + const value = (data as Record)[field]; + if (typeof value === "string" && value !== "") return value; + } + throw new Error( + `${what}: missing string field "${field}": ${JSON.stringify(data)}`, + ); +} + +function arrayField(data: unknown, field: string, what: string): unknown[] { + if (typeof data === "object" && data !== null && field in data) { + const value = (data as Record)[field]; + if (Array.isArray(value)) return value; + } + throw new Error( + `${what}: missing array field "${field}": ${JSON.stringify(data)}`, + ); +} + +async function signUp( + baseUrl: string, + name: string, +): Promise<{ userId: string; email: string; cookies: string[] }> { + const email = `cl6451-${crypto.randomUUID()}@example.invalid`; + const password = `pw-${crypto.randomUUID()}`; + const res = await api(baseUrl, "POST", "/api/auth/sign-up/email", { + name, + email, + password, + }); + expectStatus(`sign-up for ${name}`, res, 200); + if (res.cookies.length === 0) { + throw new Error(`sign-up for ${name} returned no session cookie`); + } + const userId = stringField( + (res.data as { user: unknown }).user, + "id", + `sign-up user field for ${name}`, + ); + return { userId, email, cookies: res.cookies }; +} + +async function main(): Promise { + const url = databaseUrl; + + await hop("database setup", async () => { + await resetSchema(url); + const report = await setupDatabase(url); + expect(report.action).toBe("migrated"); + }); + + const sidecarId = "cl6451-sidecar"; + const sidecarToken = crypto.randomUUID(); + await provisionSidecar(url, sidecarId, sidecarToken); + + const hub: HubHandle = await hop("hub boot", async () => + startHub({ + databaseUrl: url, + port: freePort(), + sessionSecret: Buffer.from( + crypto.getRandomValues(new Uint8Array(32)), + ).toString("hex"), + dataDir: await tempDir("cl6451-hub-data-"), + }), + ); + track(hub); + const hubPort = Number(new URL(hub.baseUrl).port); + + const sidecar: SpawnedApp = startSidecar({ + hubPort, + sidecarId, + token: sidecarToken, + dataDir: await tempDir("cl6451-sidecar-data-"), + }); + track(sidecar); + + const hubApi: ApiCall = createHubAPI(hub.baseUrl); + const pushWorkflow = createGitWorkflowPusher(); + + const user = await hop("sign up", async () => + signUp(hub.baseUrl, "CL-6451 Proof"), + ); + + const provisioned = await hop("first-login provisioning", async () => { + const res = await api( + hub.baseUrl, + "POST", + "/api/onboarding/provision", + { name: "CL-6451 Proof Bench" }, + user.cookies, + ); + expectStatus("provision", res, 200); + const data = res.data as { kind: string; tenantSlug: string }; + expect(data.kind).toBe("provisioned"); + return data; + }); + + const tenant = await hop("personal bench resolves", async () => { + const found = await findPersonalTenant( + hubApi, + user.cookies, + provisioned.tenantSlug, + ); + if (found === undefined) { + throw new Error( + `findPersonalTenant found nothing for slug ${provisioned.tenantSlug}`, + ); + } + return found; + }); + + const connected = await hop("connect Ollama", async () => { + const result = await testAndPersistCredential({ + api: hubApi, + cookies: user.cookies, + hubUrl: hub.baseUrl, + userId: user.userId, + userEmail: user.email, + provider: "ollama", + apiKey: OLLAMA_PLACEHOLDER_SECRET, + baseURLOverride: ollamaBaseUrl, + pushWorkflow, + log: () => undefined, + }); + if (result.kind !== "connected") { + throw new Error( + `expected the key-path connect to succeed, got: ${JSON.stringify(result)}`, + ); + } + return result; + }); + + await hop("SETUP — the credential's own seed completes", async () => { + const deadline = Date.now() + 180_000; + for (;;) { + try { + await ensureSeeded({ + api: hubApi, + cookies: user.cookies, + hubUrl: hub.baseUrl, + pushWorkflow, + log: () => undefined, + tenant: connected, + provider: "ollama", + apiKey: OLLAMA_PLACEHOLDER_SECRET, + baseURLOverride: ollamaBaseUrl, + }); + return; + } catch (cause) { + if (Date.now() > deadline) throw cause; + await Bun.sleep(1000); + } + } + }); + + await hop( + "SETUP — every default workflow deploys by source-ref", + async () => { + const deadline = Date.now() + 180_000; + for (;;) { + if (sidecar.exited()) { + throw new Error( + `sidecar exited before the seed finished; output:\n${sidecar.output()}`, + ); + } + try { + await seedTenant({ + api: hubApi, + cookies: user.cookies, + hubUrl: hub.baseUrl, + tenant: { + tenantId: tenant.tenantId, + principalId: tenant.principalId, + domain: tenant.tenantDomain, + }, + model: proofModelSource, + pushWorkflow, + log: () => undefined, + workflows: DEFAULT_WORKFLOWS, + confirmDeployments: false, + }); + return; + } catch (cause) { + if (Date.now() > deadline) throw cause; + await Bun.sleep(1000); + } + } + }, + ); + + // No catalog narrowing here, unlike the CL-6329 proof: the seeded + // `assistant` definition pins its own model (the connect flow's + // choice), and disabling the other offerings can disable exactly the + // model that pin names — this proof is about run identity, not model + // selection, so every seeded offering stays launchable. + + const assistant = await hop( + "SETUP — 'assistant' is invitable tenant-wide", + async () => { + const deadline = Date.now() + 60_000; + for (;;) { + const res = await api( + hub.baseUrl, + "GET", + `/api/tenants/${tenant.tenantId}/chat/invitable-definitions`, + undefined, + user.cookies, + ); + if (res.status === 200) { + const items = arrayField(res.data, "items", "invitable") as { + id: string; + name: string; + }[]; + const found = items.find((item) => item.name === "assistant"); + if (found !== undefined) return found; + } + if (Date.now() > deadline) { + throw new Error( + `"assistant" never became invitable: ${JSON.stringify(res.data)}`, + ); + } + await Bun.sleep(1000); + } + }, + ); + + // A WORKBENCH, not a chat: the live bug fired in a room whose agent + // arrived through the invite affordance and was then @named by its + // definition's wire name. + const workbenchId = await hop("SETUP — create a workbench", async () => { + const res = await api( + hub.baseUrl, + "POST", + `/api/tenants/${tenant.tenantId}/chat/workbenches`, + { kind: "workbench", name: "CL-6451 proof room" }, + user.cookies, + ); + expectStatus("create workbench", res, 201); + return stringField(res.data, "id", "create workbench"); + }); + + async function listAgents(): Promise<{ address: string; handle: string }[]> { + const res = await api( + hub.baseUrl, + "GET", + `/api/tenants/${tenant.tenantId}/chat/workbenches/${workbenchId}/agents`, + undefined, + user.cookies, + ); + expectStatus("read the room's agents", res, 200); + return arrayField(res.data, "items", "room participants") as { + address: string; + handle: string; + }[]; + } + + const agentAddress = await hop( + "SETUP — invite `assistant` (run A)", + async () => { + const deadline = Date.now() + 120_000; + for (;;) { + const res = await api( + hub.baseUrl, + "POST", + `/api/tenants/${tenant.tenantId}/chat/workbenches/${workbenchId}/invite`, + { definitionId: assistant.id }, + user.cookies, + ); + if (res.status === 201) { + return stringField(res.data, "address", "invite"); + } + if (Date.now() > deadline) { + throw new Error(`invite never landed: ${JSON.stringify(res.data)}`); + } + await Bun.sleep(1000); + } + }, + ); + const [agentRunId] = agentAddress.split("@"); + if (agentRunId === undefined) { + throw new Error(`agent address is not a run address: ${agentAddress}`); + } + + // The precondition that made the bug reachable: the participant's + // handle is NOT the definition's wire name, so the known-handle guard + // alone cannot recognize `@assistant` as this participant. + const handle = (await listAgents()).find( + (p) => p.address === agentAddress, + )?.handle; + if (handle === undefined) { + throw new Error(`the invited agent is not on the room's participant list`); + } + if (handle === assistant.name) { + throw new Error( + `precondition broken: the handle (${handle}) equals the definition's ` + + `wire name — this stack cannot reproduce the guard miss`, + ); + } + console.log( + ` TRANSCRIPT — invited ${agentAddress} as @${handle}; ` + + `the command name is "${assistant.name}"`, + ); + + type Turn = { + id: string; + agentAddress: string; + childRunId: string; + status: string; + replyMessageId: string | null; + error: string | null; + }; + + async function listTurns(): Promise { + const res = await api( + hub.baseUrl, + "GET", + `/api/tenants/${tenant.tenantId}/chat/workbenches/${workbenchId}/turns`, + undefined, + user.cookies, + ); + expectStatus("list the room's turns", res, 200); + return arrayField(res.data, "items", "list turns") as Turn[]; + } + + async function listRoomMessages(): Promise< + { id: string; address: string; text: string }[] + > { + const res = await api( + hub.baseUrl, + "GET", + `/api/tenants/${tenant.tenantId}/chat/workbenches/${workbenchId}/messages`, + undefined, + user.cookies, + ); + expectStatus("list room messages", res, 200); + const items = arrayField(res.data, "items", "list room messages") as { + id: string; + sender: { address: string }; + parts: { kind: string; text?: string }[]; + }[]; + return items.map((i) => ({ + id: i.id, + address: i.sender.address, + text: i.parts.map((p) => p.text ?? "").join(""), + })); + } + + async function sendAndAwaitTurn( + text: string, + expectedChildRunId: string, + ): Promise<{ replyText: string }> { + const res = await api( + hub.baseUrl, + "POST", + `/api/tenants/${tenant.tenantId}/chat/workbenches/${workbenchId}/messages`, + { parts: [{ kind: "text", text }] }, + user.cookies, + ); + expectStatus(`send "${text.slice(0, 40)}"`, res, 201); + // The heart of CL-6451: the `@assistant` message must NOT have been + // intercepted as a workflow command. + if ( + typeof res.data === "object" && + res.data !== null && + "command" in res.data + ) { + throw new Error( + `the @${assistant.name} message was dispatched as a command — a ` + + `second run was started: ${JSON.stringify(res.data)}`, + ); + } + + const deadline = Date.now() + TURN_TIMEOUT_MS; + for (;;) { + const turns = await listTurns(); + const settled = turns.find( + (t) => + t.childRunId === expectedChildRunId && + t.status === "completed" && + t.replyMessageId !== null, + ); + if (settled !== undefined) { + expect(settled.agentAddress).toBe(agentAddress); + const reply = (await listRoomMessages()).find( + (m) => m.id === settled.replyMessageId, + ); + const replyText = reply?.text ?? ""; + console.log( + ` TRANSCRIPT — ${settled.childRunId}: ${replyText.slice(0, 160)}`, + ); + return { replyText }; + } + const failed = turns.find((t) => t.status === "failed"); + if (failed !== undefined) { + throw new Error( + `a turn failed instead of completing: ${JSON.stringify(failed)}\n` + + `sidecar output:\n${sidecar.output()}`, + ); + } + if (Date.now() > deadline) { + throw new Error( + `no completed ${expectedChildRunId} turn within the budget: ` + + `${JSON.stringify(turns)}\nsidecar output:\n${sidecar.output()}`, + ); + } + await Bun.sleep(2000); + } + } + + await hop( + "PROOF 1 — @assistant routes into run A, never a second run", + async () => { + await sendAndAwaitTurn( + `@${assistant.name} My favorite color is teal. ` + + `Acknowledge in one short sentence.`, + "turn__0", + ); + const agents = await listAgents(); + if (agents.length !== 1) { + throw new Error( + `the room grew a second agent participant: ` + JSON.stringify(agents), + ); + } + expect(agents[0]?.address).toBe(agentAddress); + }, + ); + + await hop( + "PROOF 2 (CL-6453) — turn__1 rides the same run and remembers turn__0", + async () => { + const { replyText } = await sendAndAwaitTurn( + `@${assistant.name} What is my favorite color? ` + + `Answer with just the color name.`, + "turn__1", + ); + if (!/teal/i.test(replyText)) { + throw new Error( + `turn__1's reply does not recall turn__0's fact ("teal"): ` + + `"${replyText}" — the bootstrap exchange is missing from the ` + + `durable history`, + ); + } + const agents = await listAgents(); + expect(agents).toHaveLength(1); + }, + ); + + console.log("\nBoth CL-6451/CL-6453 proofs passed."); +} + +await main(); +process.exit(0);