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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 23 additions & 0 deletions docs/CHAT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
51 changes: 50 additions & 1 deletion packages/chat/src/agent-turns.test.ts
Original file line number Diff line number Diff line change
@@ -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",
Expand Down Expand Up @@ -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");
});
});
Binary file modified packages/chat/src/agent-turns.ts
Binary file not shown.
53 changes: 48 additions & 5 deletions packages/chat/src/routes.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -777,13 +778,31 @@ async function enrichWithReactionsAndPins<T extends WireMessageItem>(
});
}

/**
* 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
* post the message normally — for: no registry injected; text that is
* 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,
Expand All @@ -793,7 +812,7 @@ async function dispatchWorkbenchCommand(
workbenchId: string;
text: string;
},
): Promise<CommandResult | undefined> {
): Promise<WorkbenchCommandDecision | undefined> {
if (deps.commands === undefined) return undefined;
const ctx = {
tenantId: input.tenantId,
Expand All @@ -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("@")) {
Expand All @@ -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;
Expand Down Expand Up @@ -2088,13 +2126,14 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono<TenantEnv> {
// 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(
Expand Down Expand Up @@ -2136,6 +2175,10 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono<TenantEnv> {
...(parsed.inReplyToMessageId !== undefined
? { inReplyToMessageId: parsed.inReplyToMessageId }
: {}),
...(commandDecision !== undefined &&
"routeToParticipant" in commandDecision
? { forcedRecipientAddress: commandDecision.routeToParticipant }
: {}),
},
);
deps.onMessageFanout?.(sent.fanoutDelivered);
Expand Down
82 changes: 82 additions & 0 deletions packages/chat/src/workbench-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -174,13 +174,57 @@ export type LaunchAndJoinAgentResult = {
readonly joinEventDelivered: Promise<void>;
};

/**
* 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<ParticipantRecord | undefined> {
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
* the join event onto the workbench's timeline, and arms the reply
* 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,
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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 = {
Expand Down Expand Up @@ -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,
Expand Down
40 changes: 39 additions & 1 deletion packages/chat/test/agent-turns.drizzle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
}
});
});
Loading
Loading