From 68046b201151d6246bf7300a8a7585a4771eaaaf Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 23 Jul 2026 01:47:11 -0400 Subject: [PATCH] test(gateway): consolidate agent event fixtures (#112878) --- .../server-chat.agent-events.test-helpers.ts | 58 + src/gateway/server-chat.agent-events.test.ts | 2716 ++++++----------- 2 files changed, 993 insertions(+), 1781 deletions(-) create mode 100644 src/gateway/server-chat.agent-events.test-helpers.ts diff --git a/src/gateway/server-chat.agent-events.test-helpers.ts b/src/gateway/server-chat.agent-events.test-helpers.ts new file mode 100644 index 000000000000..e1cb96cf9560 --- /dev/null +++ b/src/gateway/server-chat.agent-events.test-helpers.ts @@ -0,0 +1,58 @@ +import type { AgentEventPayload, AgentEventStream } from "../infra/agent-events.js"; +import type { ChatRunRegistration, ChatRunState } from "./server-chat-state.js"; + +type AgentEventHandler = (event: AgentEventPayload) => void; + +type AgentEventOverrideKey = + | "agentId" + | "lifecycleGeneration" + | "seq" + | "sessionId" + | "sessionKey" + | "ts"; +type AgentEventOverrides = { + [Key in AgentEventOverrideKey]?: AgentEventPayload[Key] | undefined; +}; +type AgentEventCase = readonly [ + stream: AgentEventStream, + data: Record, + overrides?: AgentEventOverrides, +]; + +export function emitAgentEvent( + handler: AgentEventHandler, + runId: string, + stream: AgentEventStream, + data: Record, + overrides: AgentEventOverrides = {}, +) { + handler({ runId, seq: 1, stream, ts: Date.now(), data, ...overrides }); +} + +export function emitAgentEvents( + handler: AgentEventHandler, + runId: string, + events: readonly AgentEventCase[], +) { + events.forEach(([stream, data, overrides], index) => + emitAgentEvent(handler, runId, stream, data, { seq: index + 1, ...overrides }), + ); +} + +export function registerChatRun( + state: ChatRunState, + runId: string, + sessionKey: string, + clientRunId: string, + overrides: Omit = {}, +) { + state.registry.add(runId, { sessionKey, clientRunId, ...overrides }); +} + +export function registerNamedChatRun( + state: ChatRunState, + name: string, + overrides: Omit = {}, +) { + registerChatRun(state, `run-${name}`, `session-${name}`, `client-${name}`, overrides); +} diff --git a/src/gateway/server-chat.agent-events.test.ts b/src/gateway/server-chat.agent-events.test.ts index a26b7faee70c..62725459f51d 100644 --- a/src/gateway/server-chat.agent-events.test.ts +++ b/src/gateway/server-chat.agent-events.test.ts @@ -63,6 +63,12 @@ vi.mock("./session-utils.js", () => { import { getRuntimeConfig } from "../config/io.js"; import { resolveHeartbeatVisibility } from "../infra/heartbeat-visibility.js"; +import { + emitAgentEvent, + emitAgentEvents, + registerChatRun, + registerNamedChatRun, +} from "./server-chat.agent-events.test-helpers.js"; import { createAgentEventHandler, createChatRunState, @@ -182,20 +188,26 @@ describe("agent event handler", () => { text: string, field: "text" | "delta" = "text", ): ReturnType { - harness.chatRunState.registry.add("run-1", { - sessionKey: "session-1", - clientRunId: "client-1", - }); - harness.handler({ - runId: "run-1", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { [field]: text }, - }); + registerChatRun(harness.chatRunState, "run-1", "session-1", "client-1"); + emitAgentEvent(harness.handler, "run-1", "assistant", { [field]: text }); return harness; } + function mockSessionEntry( + entry: ReturnType["entry"], + canonicalKey = "session-1", + ) { + vi.mocked(loadSessionEntry).mockReturnValue({ + cfg: {}, + storePath: "/tmp/sessions.json", + store: {}, + entry, + canonicalKey, + storeKeys: [canonicalKey], + legacyKey: undefined, + }); + } + function chatBroadcastCalls(broadcast: ReturnType) { return broadcast.mock.calls.filter(([event]) => event === "chat"); } @@ -211,23 +223,19 @@ describe("agent event handler", () => { it("carries prepared validation diagnostics into active run state", () => { const updateRunToolErrorSummary = vi.fn(); const { chatRunState, handler } = createHarness({ updateRunToolErrorSummary }); - chatRunState.registry.add("provider-run", { - sessionKey: "session-1", - clientRunId: "client-run", - }); - - handler({ - runId: "provider-run", - seq: 1, - stream: "tool", - ts: 1_000, - data: { + registerChatRun(chatRunState, "provider-run", "session-1", "client-run"); + emitAgentEvent( + handler, + "provider-run", + "tool", + { phase: "result", name: "edit", isError: true, toolErrorSummary: "edit tool validation failed: edits: must be an array", }, - }); + { ts: 1_000 }, + ); expect(updateRunToolErrorSummary).toHaveBeenCalledWith({ runId: "provider-run", @@ -238,22 +246,18 @@ describe("agent event handler", () => { it("records, replaces, dismisses, and clears normalized plan snapshots", () => { const { chatRunState, handler } = createHarness(); - chatRunState.registry.add("provider-run", { - sessionKey: "session-1", - clientRunId: "client-run", - }); - - handler({ - runId: "provider-run", - seq: 1, - stream: "plan", - ts: 1_000, - data: { + registerChatRun(chatRunState, "provider-run", "session-1", "client-run"); + emitAgentEvent( + handler, + "provider-run", + "plan", + { phase: "update", explanation: " Initial plan ", steps: ["Legacy step", { step: "Active step", status: "in_progress" }], }, - }); + { ts: 1_000 }, + ); expect(chatRunState.planSnapshots.get("client-run")).toEqual({ explanation: "Initial plan", steps: [ @@ -262,27 +266,27 @@ describe("agent event handler", () => { ], }); - handler({ - runId: "provider-run", - seq: 2, - stream: "plan", - ts: 1_100, - data: { - phase: "update", - steps: [{ step: "Replacement", status: "completed" }], - }, - }); + emitAgentEvent( + handler, + "provider-run", + "plan", + { phase: "update", steps: [{ step: "Replacement", status: "completed" }] }, + { seq: 2, ts: 1_100 }, + ); expect(chatRunState.planSnapshots.get("client-run")).toEqual({ steps: [{ step: "Replacement", status: "completed" }], }); - handler({ - runId: "provider-run", - seq: 3, - stream: "plan", - ts: 1_200, - data: { phase: "update", steps: [] }, - }); + emitAgentEvent( + handler, + "provider-run", + "plan", + { phase: "update", steps: [] }, + { + seq: 3, + ts: 1_200, + }, + ); expect(chatRunState.planSnapshots.get("client-run")).toEqual({ steps: [] }); chatRunState.planSnapshots.set("client-run", { @@ -298,29 +302,22 @@ describe("agent event handler", () => { ] as const)("clears stale validation diagnostics on $stream progress", (progressEvent) => { const updateRunToolErrorSummary = vi.fn(); const { chatRunState, handler } = createHarness({ updateRunToolErrorSummary }); - chatRunState.registry.add("provider-run", { - sessionKey: "session-1", - clientRunId: "client-run", - }); - - handler({ - runId: "provider-run", - seq: 1, - stream: "tool", - ts: 1_000, - data: { + registerChatRun(chatRunState, "provider-run", "session-1", "client-run"); + emitAgentEvent( + handler, + "provider-run", + "tool", + { phase: "result", name: "edit", isError: true, toolErrorSummary: "edit tool validation failed: invalid arguments", }, - }); - handler({ - runId: "provider-run", + { ts: 1_000 }, + ); + emitAgentEvent(handler, "provider-run", progressEvent.stream, progressEvent.data, { seq: 2, - stream: progressEvent.stream, ts: 1_100, - data: progressEvent.data, }); expect(updateRunToolErrorSummary).toHaveBeenLastCalledWith({ @@ -397,18 +394,31 @@ describe("agent event handler", () => { activeModel: "moonshotai/Kimi-K2.5", } as const; + const SESSION_OWNERSHIP = { + spawnedBy: "agent:main:main", + spawnedWorkspaceDir: "/tmp/subagent", + forkedFromParent: true, + spawnDepth: 2, + subagentRole: "orchestrator", + subagentControlScope: "children", + lastThreadId: 42, + fastMode: true, + verboseLevel: "on", + } as const; + + const OWNED_SESSION_ROW = { + key: "session-1", + kind: "direct" as const, + ...SESSION_OWNERSHIP, + updatedAt: 1_200, + }; + function emitLifecycleEnd( handler: ReturnType["handler"], runId: string, seq = 2, ) { - handler({ - runId, - seq, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvent(handler, runId, "lifecycle", { phase: "end" }, { seq }); } function emitFallbackLifecycle(params: { @@ -417,14 +427,16 @@ describe("agent event handler", () => { seq?: number; sessionKey?: string; }) { - params.handler({ - runId: params.runId, - seq: params.seq ?? 1, - stream: "lifecycle", - ts: Date.now(), - ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), - data: { ...FALLBACK_LIFECYCLE_DATA }, - }); + emitAgentEvent( + params.handler, + params.runId, + "lifecycle", + { ...FALLBACK_LIFECYCLE_DATA }, + { + seq: params.seq ?? 1, + ...(params.sessionKey ? { sessionKey: params.sessionKey } : {}), + }, + ); } function expectSingleAgentBroadcastPayload(broadcast: ReturnType) { @@ -451,116 +463,61 @@ describe("agent event handler", () => { it("injects isHeartbeat into agent broadcast payloads when present in run context", () => { const harness = createHarness(); - registerAgentRunContext("run-heartbeat-true", { sessionKey: "session-1", isHeartbeat: true }); - registerAgentRunContext("run-heartbeat-false", { sessionKey: "session-2", isHeartbeat: false }); + for (const { runId, sessionKey, isHeartbeat, ts, needsUnlinkedEvent } of [ + { + runId: "run-heartbeat-true", + sessionKey: "session-1", + isHeartbeat: true, + ts: 100, + needsUnlinkedEvent: true, + }, + { + runId: "run-heartbeat-false", + sessionKey: "session-2", + isHeartbeat: false, + ts: 101, + needsUnlinkedEvent: false, + }, + { + runId: "run-normal", + sessionKey: "session-3", + isHeartbeat: undefined, + ts: 102, + needsUnlinkedEvent: false, + }, + ] as const) { + if (isHeartbeat !== undefined) { + registerAgentRunContext(runId, { sessionKey, isHeartbeat }); + } + if (needsUnlinkedEvent) { + emitAgentEvent(harness.handler, runId, "assistant", { text: "hello" }, { ts }); + } + registerChatRun(harness.chatRunState, runId, sessionKey, runId); + emitAgentEvent( + harness.handler, + runId, + "assistant", + { text: "hello" }, + { + seq: needsUnlinkedEvent ? 2 : 1, + ts, + }, + ); - // 1. isHeartbeat: true - harness.handler({ - runId: "run-heartbeat-true", - seq: 1, - stream: "assistant", - ts: 100, - data: { text: "hello" }, - }); - - const agentPayload1 = requireRecord( - requireCall( + for (const payload of [ harness.broadcast.mock.calls.find(([event]) => event === "agent")?.[1], - "agent broadcast payload", - ), - "agent broadcast payload", - ); - expect(agentPayload1.isHeartbeat).toBe(true); - - // sessionKey is required for nodeSendToSession to be called - harness.chatRunState.registry.add("run-heartbeat-true", { - sessionKey: "session-1", - clientRunId: "run-heartbeat-true", - }); - harness.handler({ - runId: "run-heartbeat-true", - seq: 2, - stream: "assistant", - ts: 100, - data: { text: "hello" }, - }); - - const nodeSendPayload1 = requireRecord( - requireCall( harness.nodeSendToSession.mock.calls.find(([, event]) => event === "agent")?.[2], - "agent node-send payload", - ), - "agent node-send payload", - ); - expect(nodeSendPayload1.isHeartbeat).toBe(true); - - harness.broadcast.mockClear(); - harness.nodeSendToSession.mockClear(); - - // 2. isHeartbeat: false - harness.chatRunState.registry.add("run-heartbeat-false", { - sessionKey: "session-2", - clientRunId: "run-heartbeat-false", - }); - harness.handler({ - runId: "run-heartbeat-false", - seq: 1, - stream: "assistant", - ts: 101, - data: { text: "hello" }, - }); - - const agentPayload2 = requireRecord( - requireCall( - harness.broadcast.mock.calls.find(([event]) => event === "agent")?.[1], - "agent broadcast payload", - ), - "agent broadcast payload", - ); - expect(agentPayload2.isHeartbeat).toBe(false); - - const nodeSendPayload2 = requireRecord( - requireCall( - harness.nodeSendToSession.mock.calls.find(([, event]) => event === "agent")?.[2], - "agent node-send payload", - ), - "agent node-send payload", - ); - expect(nodeSendPayload2.isHeartbeat).toBe(false); - - harness.broadcast.mockClear(); - harness.nodeSendToSession.mockClear(); - - // 3. isHeartbeat: undefined (absent) - harness.chatRunState.registry.add("run-normal", { - sessionKey: "session-3", - clientRunId: "run-normal", - }); - harness.handler({ - runId: "run-normal", - seq: 1, - stream: "assistant", - ts: 102, - data: { text: "hello" }, - }); - - const normalBroadcast = requireRecord( - requireCall( - harness.broadcast.mock.calls.find(([event]) => event === "agent")?.[1], - "normal agent broadcast payload", - ), - "normal agent broadcast payload", - ); - expect("isHeartbeat" in normalBroadcast).toBe(false); - - const normalNodeSend = requireRecord( - requireCall( - harness.nodeSendToSession.mock.calls.find(([, event]) => event === "agent")?.[2], - "normal agent node-send payload", - ), - "normal agent node-send payload", - ); - expect("isHeartbeat" in normalNodeSend).toBe(false); + ]) { + const record = requireRecord(requireCall(payload, "agent payload"), "agent payload"); + if (isHeartbeat === undefined) { + expect(record).not.toHaveProperty("isHeartbeat"); + } else { + expect(record.isHeartbeat).toBe(isHeartbeat); + } + } + harness.broadcast.mockClear(); + harness.nodeSendToSession.mockClear(); + } }); it.each(["text", "delta"] as const)("emits chat delta for assistant %s-only events", (field) => { @@ -585,10 +542,7 @@ describe("agent event handler", () => { it("keeps internal context private when it spans delta-only events", () => { const { broadcast, chatRunState, handler, nowSpy } = createHarness({ now: 1_000 }); - chatRunState.registry.add("run-split-context", { - sessionKey: "session-split-context", - clientRunId: "client-split-context", - }); + registerNamedChatRun(chatRunState, "split-context"); const deltas = [ `Visible\n${INTERNAL_RUNTIME_CONTEXT_BEGIN}\n`, @@ -596,13 +550,7 @@ describe("agent event handler", () => { `${INTERNAL_RUNTIME_CONTEXT_END}\nAfter`, ]; deltas.forEach((delta, index) => { - handler({ - runId: "run-split-context", - seq: index + 1, - stream: "assistant", - ts: Date.now(), - data: { delta }, - }); + emitAgentEvent(handler, "run-split-context", "assistant", { delta }, { seq: index + 1 }); }); emitLifecycleEnd(handler, "run-split-context", 4); @@ -619,9 +567,7 @@ describe("agent event handler", () => { it("emits the first assistant chat.send timing event to the originating Control UI", () => { const { broadcastToConnIds, chatRunState, handler, nowSpy } = createHarness({ now: 1_000 }); - chatRunState.registry.add("run-1", { - sessionKey: "session-1", - clientRunId: "client-1", + registerChatRun(chatRunState, "run-1", "session-1", "client-1", { chatSendTiming: { ackedAtMs: 0, connId: "conn-control-ui", @@ -630,20 +576,10 @@ describe("agent event handler", () => { }, }); - handler({ - runId: "run-1", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); - handler({ - runId: "run-1", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world again" }, - }); + emitAgentEvents(handler, "run-1", [ + ["assistant", { text: "Hello world" }], + ["assistant", { text: "Hello world again" }], + ]); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); expect(broadcastToConnIds).toHaveBeenCalledWith( @@ -664,9 +600,7 @@ describe("agent event handler", () => { it("emits first assistant chat.send timing when text first flushes on final", () => { const { broadcastToConnIds, chatRunState, handler, nowSpy } = createHarness({ now: 1_000 }); - chatRunState.registry.add("run-final-only", { - sessionKey: "session-1", - clientRunId: "client-final", + registerChatRun(chatRunState, "run-final-only", "session-1", "client-final", { chatSendTiming: { ackedAtMs: 0, connId: "conn-control-ui", @@ -676,13 +610,7 @@ describe("agent event handler", () => { }); chatRunState.buffers.set("client-final", "Final only reply"); - handler({ - runId: "run-final-only", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvent(handler, "run-final-only", "lifecycle", { phase: "end" }); expect(broadcastToConnIds).toHaveBeenCalledWith( "chat.send_timing", @@ -702,19 +630,11 @@ describe("agent event handler", () => { it("projects typed run startup status onto the active chat stream", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-startup", { - sessionKey: "session-startup", - clientRunId: "client-startup", + registerNamedChatRun(chatRunState, "startup", { agentId: "main", }); - handler({ - runId: "run-startup", - seq: 1, - stream: "run_status", - ts: Date.now(), - data: { phase: "preparing_context" }, - }); + emitAgentEvent(handler, "run-startup", "run_status", { phase: "preparing_context" }); expect(chatBroadcastCalls(broadcast)).toEqual([ [ @@ -738,20 +658,17 @@ describe("agent event handler", () => { let now = 10_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-throttle", { - sessionKey: "session-agent-throttle", - clientRunId: "client-agent-throttle", - }); + registerNamedChatRun(chatRunState, "agent-throttle"); for (let i = 0; i < 5; i += 1) { now = 10_000 + i * 20; - handler({ - runId: "run-agent-throttle", - seq: i + 1, - stream: "assistant", - ts: Date.now(), - data: { text: "x".repeat(i + 1), delta: "x" }, - }); + emitAgentEvent( + handler, + "run-agent-throttle", + "assistant", + { text: "x".repeat(i + 1), delta: "x" }, + { seq: i + 1 }, + ); } const agentCalls = agentBroadcastCalls(broadcast); @@ -773,41 +690,26 @@ describe("agent event handler", () => { let now = 20_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-flush", { - sessionKey: "session-agent-flush", - clientRunId: "client-agent-flush", - }); + registerNamedChatRun(chatRunState, "agent-flush"); - handler({ - runId: "run-agent-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello", delta: "Hello" }, - }); + emitAgentEvent(handler, "run-agent-flush", "assistant", { text: "Hello", delta: "Hello" }); now = 20_050; - handler({ - runId: "run-agent-flush", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world", delta: " world" }, - }); + emitAgentEvent( + handler, + "run-agent-flush", + "assistant", + { text: "Hello world", delta: " world" }, + { seq: 2 }, + ); now = 20_090; - handler({ - runId: "run-agent-flush", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world!", delta: "!" }, - }); - handler({ - runId: "run-agent-flush", - seq: 4, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvent( + handler, + "run-agent-flush", + "assistant", + { text: "Hello world!", delta: "!" }, + { seq: 3 }, + ); + emitAgentEvent(handler, "run-agent-flush", "lifecycle", { phase: "end" }, { seq: 4 }); const agentCalls = agentBroadcastCalls(broadcast); expect(agentCalls).toHaveLength(3); @@ -858,34 +760,25 @@ describe("agent event handler", () => { let now = 22_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-window", { - sessionKey: "session-agent-window", - clientRunId: "client-agent-window", - }); + registerNamedChatRun(chatRunState, "agent-window"); - handler({ - runId: "run-agent-window", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hel", delta: "Hel" }, - }); + emitAgentEvent(handler, "run-agent-window", "assistant", { text: "Hel", delta: "Hel" }); now = 22_050; - handler({ - runId: "run-agent-window", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello", delta: "lo" }, - }); + emitAgentEvent( + handler, + "run-agent-window", + "assistant", + { text: "Hello", delta: "lo" }, + { seq: 2 }, + ); now = 22_200; - handler({ - runId: "run-agent-window", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello!", delta: "!" }, - }); + emitAgentEvent( + handler, + "run-agent-window", + "assistant", + { text: "Hello!", delta: "!" }, + { seq: 3 }, + ); const agentCalls = agentBroadcastCalls(broadcast); expect(agentCalls).toHaveLength(3); @@ -924,34 +817,28 @@ describe("agent event handler", () => { let now = 23_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-cross-stream", { - sessionKey: "session-agent-cross-stream", - clientRunId: "client-agent-cross-stream", - }); + registerNamedChatRun(chatRunState, "agent-cross-stream"); - handler({ - runId: "run-agent-cross-stream", - seq: 1, - stream: "thinking", - ts: Date.now(), - data: { text: "Think", delta: "Think" }, + emitAgentEvent(handler, "run-agent-cross-stream", "thinking", { + text: "Think", + delta: "Think", }); now = 23_050; - handler({ - runId: "run-agent-cross-stream", - seq: 2, - stream: "thinking", - ts: Date.now(), - data: { text: "Thinking", delta: "ing" }, - }); + emitAgentEvent( + handler, + "run-agent-cross-stream", + "thinking", + { text: "Thinking", delta: "ing" }, + { seq: 2 }, + ); now = 23_080; - handler({ - runId: "run-agent-cross-stream", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { text: "Answer", delta: "Answer" }, - }); + emitAgentEvent( + handler, + "run-agent-cross-stream", + "assistant", + { text: "Answer", delta: "Answer" }, + { seq: 3 }, + ); const agentCalls = agentBroadcastCalls(broadcast); expect(agentCalls.map(([, payload]) => (payload as { seq?: number }).seq)).toEqual([1, 2, 3]); @@ -975,26 +862,17 @@ describe("agent event handler", () => { let now = 25_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-start", { - sessionKey: "session-agent-start", - clientRunId: "client-agent-start", - }); + registerNamedChatRun(chatRunState, "agent-start"); - handler({ - runId: "run-agent-start", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); + emitAgentEvent(handler, "run-agent-start", "lifecycle", { phase: "start" }); now = 25_050; - handler({ - runId: "run-agent-start", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello", delta: "Hello" }, - }); + emitAgentEvent( + handler, + "run-agent-start", + "assistant", + { text: "Hello", delta: "Hello" }, + { seq: 2 }, + ); const agentCalls = agentBroadcastCalls(broadcast); expect(agentCalls).toHaveLength(2); @@ -1017,20 +895,17 @@ describe("agent event handler", () => { let now = 27_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-thinking", { - sessionKey: "session-agent-thinking", - clientRunId: "client-agent-thinking", - }); + registerNamedChatRun(chatRunState, "agent-thinking"); for (let i = 0; i < 5; i += 1) { now = 27_000 + i * 20; - handler({ - runId: "run-agent-thinking", - seq: i + 1, - stream: "thinking", - ts: Date.now(), - data: { text: "t".repeat(i + 1), delta: "t" }, - }); + emitAgentEvent( + handler, + "run-agent-thinking", + "thinking", + { text: "t".repeat(i + 1), delta: "t" }, + { seq: i + 1 }, + ); } const agentCalls = agentBroadcastCalls(broadcast); @@ -1054,49 +929,34 @@ describe("agent event handler", () => { let now = 30_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-agent-media", { - sessionKey: "session-agent-media", - clientRunId: "client-agent-media", - }); + registerNamedChatRun(chatRunState, "agent-media"); - handler({ - runId: "run-agent-media", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Look", delta: "Look" }, - }); + emitAgentEvent(handler, "run-agent-media", "assistant", { text: "Look", delta: "Look" }); now = 30_050; - handler({ - runId: "run-agent-media", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Look", delta: "", mediaUrls: ["https://example.test/image.png"] }, - }); + emitAgentEvent( + handler, + "run-agent-media", + "assistant", + { text: "Look", delta: "", mediaUrls: ["https://example.test/image.png"] }, + { seq: 2 }, + ); now = 30_070; - handler({ - runId: "run-agent-media", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { text: "Look elsewhere", delta: "", replace: true }, - }); + emitAgentEvent( + handler, + "run-agent-media", + "assistant", + { text: "Look elsewhere", delta: "", replace: true }, + { seq: 3 }, + ); now = 30_090; - handler({ - runId: "run-agent-media", - seq: 4, - stream: "assistant", - ts: Date.now(), - data: { text: "Look elsewhere now", delta: " now" }, - }); - handler({ - runId: "run-agent-media", - seq: 5, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvent( + handler, + "run-agent-media", + "assistant", + { text: "Look elsewhere now", delta: " now" }, + { seq: 4 }, + ); + emitAgentEvent(handler, "run-agent-media", "lifecycle", { phase: "end" }, { seq: 5 }); const agentCalls = agentBroadcastCalls(broadcast); expect(agentCalls).toHaveLength(5); @@ -1187,15 +1047,9 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_000, }); - chatRunState.registry.add("run-2", { sessionKey: "session-2", clientRunId: "client-2" }); + registerNamedChatRun(chatRunState, "2"); - handler({ - runId: "run-2", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: replyText }, - }); + emitAgentEvent(handler, "run-2", "assistant", { text: replyText }); emitLifecycleEnd(handler, "run-2"); const payload = expectSingleFinalChatPayload(broadcast) as { message?: unknown }; @@ -1209,16 +1063,10 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_100, }); - chatRunState.registry.add("run-3", { sessionKey: "session-3", clientRunId: "client-3" }); + registerNamedChatRun(chatRunState, "3"); for (const text of ["NO", "NO_", "NO_RE", "NO_REPLY"]) { - handler({ - runId: "run-3", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text }, - }); + emitAgentEvent(handler, "run-3", "assistant", { text }); } emitLifecycleEnd(handler, "run-3"); @@ -1237,19 +1085,10 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_150, }); - chatRunState.registry.add("run-control", { - sessionKey: "session-control", - clientRunId: "client-control", - }); + registerNamedChatRun(chatRunState, "control"); for (const text of fragments) { - handler({ - runId: "run-control", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text }, - }); + emitAgentEvent(handler, "run-control", "assistant", { text }); } emitLifecycleEnd(handler, "run-control"); @@ -1264,15 +1103,9 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_200, }); - chatRunState.registry.add("run-4", { sessionKey: "session-4", clientRunId: "client-4" }); + registerNamedChatRun(chatRunState, "4"); - handler({ - runId: "run-4", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "No" }, - }); + emitAgentEvent(handler, "run-4", "assistant", { text: "No" }); emitLifecycleEnd(handler, "run-4"); const payload = expectSingleFinalChatPayload(broadcast) as { @@ -1287,22 +1120,12 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_250, }); - chatRunState.registry.add("run-4b", { sessionKey: "session-4b", clientRunId: "client-4b" }); + registerNamedChatRun(chatRunState, "4b"); - handler({ - runId: "run-4b", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "NO_REPLYThe user" }, - }); - handler({ - runId: "run-4b", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "NO_REPLYThe user is saying hello" }, - }); + emitAgentEvents(handler, "run-4b", [ + ["assistant", { text: "NO_REPLYThe user" }], + ["assistant", { text: "NO_REPLYThe user is saying hello" }], + ]); emitLifecycleEnd(handler, "run-4b"); const chatCalls = chatBroadcastCalls(broadcast); @@ -1327,27 +1150,12 @@ describe("agent event handler", () => { let now = 10_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-flush", { - sessionKey: "session-flush", - clientRunId: "client-flush", - }); + registerNamedChatRun(chatRunState, "flush"); - handler({ - runId: "run-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello" }, - }); + emitAgentEvent(handler, "run-flush", "assistant", { text: "Hello" }); now = 10_100; - handler({ - runId: "run-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent(handler, "run-flush", "assistant", { text: "Hello world" }); emitLifecycleEnd(handler, "run-flush"); @@ -1374,27 +1182,21 @@ describe("agent event handler", () => { let now = 10_500; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-segmented", { - sessionKey: "session-segmented", - clientRunId: "client-segmented", - }); + registerNamedChatRun(chatRunState, "segmented"); - handler({ - runId: "run-segmented", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Before tool call", delta: "Before tool call" }, + emitAgentEvent(handler, "run-segmented", "assistant", { + text: "Before tool call", + delta: "Before tool call", }); now = 10_700; - handler({ - runId: "run-segmented", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "After tool call", delta: "\nAfter tool call" }, - }); + emitAgentEvent( + handler, + "run-segmented", + "assistant", + { text: "After tool call", delta: "\nAfter tool call" }, + { seq: 2 }, + ); emitLifecycleEnd(handler, "run-segmented", 3); @@ -1422,27 +1224,21 @@ describe("agent event handler", () => { let now = 10_800; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-segmented-flush", { - sessionKey: "session-segmented-flush", - clientRunId: "client-segmented-flush", - }); + registerNamedChatRun(chatRunState, "segmented-flush"); - handler({ - runId: "run-segmented-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Before tool call", delta: "Before tool call" }, + emitAgentEvent(handler, "run-segmented-flush", "assistant", { + text: "Before tool call", + delta: "Before tool call", }); now = 10_860; - handler({ - runId: "run-segmented-flush", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "After tool call", delta: "\nAfter tool call" }, - }); + emitAgentEvent( + handler, + "run-segmented-flush", + "assistant", + { text: "After tool call", delta: "\nAfter tool call" }, + { seq: 2 }, + ); emitLifecycleEnd(handler, "run-segmented-flush", 3); @@ -1470,27 +1266,12 @@ describe("agent event handler", () => { let now = 11_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-no-dup-flush", { - sessionKey: "session-no-dup-flush", - clientRunId: "client-no-dup-flush", - }); + registerNamedChatRun(chatRunState, "no-dup-flush"); - handler({ - runId: "run-no-dup-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello" }, - }); + emitAgentEvent(handler, "run-no-dup-flush", "assistant", { text: "Hello" }); now = 11_200; - handler({ - runId: "run-no-dup-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent(handler, "run-no-dup-flush", "assistant", { text: "Hello world" }); emitLifecycleEnd(handler, "run-no-dup-flush"); @@ -1509,27 +1290,18 @@ describe("agent event handler", () => { let now = 11_250; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-unchanged-snapshot", { - sessionKey: "session-unchanged-snapshot", - clientRunId: "client-unchanged-snapshot", - }); + registerNamedChatRun(chatRunState, "unchanged-snapshot"); - handler({ - runId: "run-unchanged-snapshot", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent(handler, "run-unchanged-snapshot", "assistant", { text: "Hello world" }); now = 11_450; - handler({ - runId: "run-unchanged-snapshot", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent( + handler, + "run-unchanged-snapshot", + "assistant", + { text: "Hello world" }, + { seq: 2 }, + ); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls).toHaveLength(1); @@ -1543,27 +1315,12 @@ describe("agent event handler", () => { let now = 11_300; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-replacement", { - sessionKey: "session-replacement", - clientRunId: "client-replacement", - }); + registerNamedChatRun(chatRunState, "replacement"); - handler({ - runId: "run-replacement", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent(handler, "run-replacement", "assistant", { text: "Hello world" }); now = 11_500; - handler({ - runId: "run-replacement", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Goodbye world" }, - }); + emitAgentEvent(handler, "run-replacement", "assistant", { text: "Goodbye world" }, { seq: 2 }); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls).toHaveLength(2); @@ -1585,27 +1342,12 @@ describe("agent event handler", () => { let now = 11_700; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-short-replacement-flush", { - sessionKey: "session-short-replacement-flush", - clientRunId: "client-short-replacement-flush", - }); + registerNamedChatRun(chatRunState, "short-replacement-flush"); - handler({ - runId: "run-short-replacement-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Hello world" }, - }); + emitAgentEvent(handler, "run-short-replacement-flush", "assistant", { text: "Hello world" }); now = 11_760; - handler({ - runId: "run-short-replacement-flush", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Hi" }, - }); + emitAgentEvent(handler, "run-short-replacement-flush", "assistant", { text: "Hi" }, { seq: 2 }); emitLifecycleEnd(handler, "run-short-replacement-flush", 3); @@ -1630,27 +1372,12 @@ describe("agent event handler", () => { it("cleans up agent run sequence tracking when lifecycle completes", () => { const { agentRunSeq, chatRunState, handler, nowSpy } = createHarness({ now: 2_500 }); - chatRunState.registry.add("run-cleanup", { - sessionKey: "session-cleanup", - clientRunId: "client-cleanup", - }); + registerNamedChatRun(chatRunState, "cleanup"); - handler({ - runId: "run-cleanup", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "done" }, - }); + emitAgentEvent(handler, "run-cleanup", "assistant", { text: "done" }); expect(agentRunSeq.get("run-cleanup")).toBe(1); - handler({ - runId: "run-cleanup", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvent(handler, "run-cleanup", "lifecycle", { phase: "end" }, { seq: 2 }); expect(agentRunSeq.has("run-cleanup")).toBe(false); expect(agentRunSeq.has("client-cleanup")).toBe(false); @@ -1661,18 +1388,9 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler, nowSpy } = createHarness({ now: 2_500, }); - chatRunState.registry.add("run-stale-tail", { - sessionKey: "session-stale-tail", - clientRunId: "client-stale-tail", - }); + registerNamedChatRun(chatRunState, "stale-tail"); - handler({ - runId: "run-stale-tail", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "done" }, - }); + emitAgentEvent(handler, "run-stale-tail", "assistant", { text: "done" }); emitLifecycleEnd(handler, "run-stale-tail"); const errorCallsBeforeStaleEvent = broadcast.mock.calls.filter( ([event, payload]) => @@ -1680,13 +1398,7 @@ describe("agent event handler", () => { ).length; const sessionChatCallsBeforeStaleEvent = sessionChatCalls(nodeSendToSession).length; - handler({ - runId: "run-stale-tail", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { text: "late tail" }, - }); + emitAgentEvent(handler, "run-stale-tail", "assistant", { text: "late tail" }, { seq: 3 }); const errorCalls = broadcast.mock.calls.filter( ([event, payload]) => @@ -1711,41 +1423,32 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-tool-flush", }); - chatRunState.registry.add("run-tool-flush", { - sessionKey: "session-tool-flush", - clientRunId: "client-tool-flush", - }); + registerNamedChatRun(chatRunState, "tool-flush"); registerAgentRunContext("run-tool-flush", { sessionKey: "session-tool-flush", verboseLevel: "off", }); toolEventRecipients.add("run-tool-flush", "conn-1"); - handler({ - runId: "run-tool-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Before tool" }, - }); + emitAgentEvent(handler, "run-tool-flush", "assistant", { text: "Before tool" }); // Throttled assistant update (within 150ms window). now = 12_050; - handler({ - runId: "run-tool-flush", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "Before tool expanded" }, - }); + emitAgentEvent( + handler, + "run-tool-flush", + "assistant", + { text: "Before tool expanded" }, + { seq: 2 }, + ); - handler({ - runId: "run-tool-flush", - seq: 3, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "read", toolCallId: "tool-flush-1" }, - }); + emitAgentEvent( + handler, + "run-tool-flush", + "tool", + { phase: "start", name: "read", toolCallId: "tool-flush-1" }, + { seq: 3 }, + ); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls).toHaveLength(2); @@ -1774,13 +1477,7 @@ describe("agent event handler", () => { registerAgentRunContext("run-tool", { sessionKey: "session-1", verboseLevel: "on" }); toolEventRecipients.add("run-tool", "conn-1"); - handler({ - runId: "run-tool", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "read", toolCallId: "t1" }, - }); + emitAgentEvent(handler, "run-tool", "tool", { phase: "start", name: "read", toolCallId: "t1" }); expect(broadcast).not.toHaveBeenCalled(); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); @@ -1794,12 +1491,10 @@ describe("agent event handler", () => { registerAgentRunContext("run-tool-off", { sessionKey: "session-1", verboseLevel: "off" }); toolEventRecipients.add("run-tool-off", "conn-1"); - handler({ - runId: "run-tool-off", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "read", toolCallId: "t2" }, + emitAgentEvent(handler, "run-tool-off", "tool", { + phase: "start", + name: "read", + toolCallId: "t2", }); // Tool events always broadcast to registered WS recipients @@ -1814,27 +1509,17 @@ describe("agent event handler", () => { now: 1_000, resolveSessionKeyForRun: () => "session-1", }); - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { sessionId: "session-1", verboseLevel: "on", updatedAt: 1_500 }, - canonicalKey: "session-1", - storeKeys: ["session-1"], - legacyKey: undefined, - }); + mockSessionEntry({ sessionId: "session-1", verboseLevel: "on", updatedAt: 1_500 }); registerAgentRunContext("run-tool-toggle", { sessionKey: "session-1", verboseLevel: "off", }); - handler({ - runId: "run-tool-toggle", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "read", toolCallId: "t-toggle" }, + emitAgentEvent(handler, "run-tool-toggle", "tool", { + phase: "start", + name: "read", + toolCallId: "t-toggle", }); const nodeToolCalls = nodeSendToSession.mock.calls.filter(([, event]) => event === "agent"); @@ -1852,27 +1537,17 @@ describe("agent event handler", () => { now: 2_000, resolveSessionKeyForRun: () => "session-1", }); - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { sessionId: "session-1", verboseLevel: "off", updatedAt: 1_500 }, - canonicalKey: "session-1", - storeKeys: ["session-1"], - legacyKey: undefined, - }); + mockSessionEntry({ sessionId: "session-1", verboseLevel: "off", updatedAt: 1_500 }); registerAgentRunContext("run-tool-inline", { sessionKey: "session-1", verboseLevel: "on", }); - handler({ - runId: "run-tool-inline", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "read", toolCallId: "t-inline" }, + emitAgentEvent(handler, "run-tool-inline", "tool", { + phase: "start", + name: "read", + toolCallId: "t-inline", }); const nodeToolCalls = nodeSendToSession.mock.calls.filter(([, event]) => event === "agent"); @@ -1884,36 +1559,23 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-1", }); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "session-1", - kind: "direct", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", - updatedAt: 1_200, - }); + vi.mocked(loadGatewaySessionRow).mockReturnValue(OWNED_SESSION_ROW); registerAgentRunContext("run-session-tool", { sessionKey: "session-1", verboseLevel: "off" }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "run-session-tool", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-session-tool", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-session-1", args: { command: "echo hi" }, }, - }); + { ts: 1_234 }, + ); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); expect(requireMockArg(broadcastToConnIds, 0, 0, "session tool event")).toBe("session.tool"); @@ -1921,15 +1583,7 @@ describe("agent event handler", () => { expectRecordFields(sessionToolPayload, { runId: "run-session-tool", sessionKey: "session-1", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", + ...SESSION_OWNERSHIP, stream: "tool", ts: 1_234, }); @@ -1967,20 +1621,18 @@ describe("agent event handler", () => { status: "running", updatedAt: 1_200, }); - chatRunState.registry.add("run-global-tool", { - sessionKey: "global", + registerChatRun(chatRunState, "run-global-tool", "global", "client-global-tool", { agentId: "work", - clientRunId: "client-global-tool", }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "run-global-tool", - seq: 1, - stream: "tool", - ts: 1_234, - data: { phase: "start", name: "exec", toolCallId: "tool-global-1" }, - }); + emitAgentEvent( + handler, + "run-global-tool", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-global-1" }, + { ts: 1_234 }, + ); expect(loadGatewaySessionRow).toHaveBeenCalledWith("global", { agentId: "work" }); expect(requireMockArg(broadcastToConnIds, 0, 0, "session tool event")).toBe("session.tool"); @@ -2013,18 +1665,18 @@ describe("agent event handler", () => { sessionEventSubscribers.subscribe("conn-overlap"); sessionEventSubscribers.subscribe("conn-session-only"); - handler({ - runId: "run-session-dedupe-tool", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-session-dedupe-tool", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-session-dedupe-1", args: { command: "echo hi" }, }, - }); + { ts: 1_234 }, + ); expect(broadcastToConnIds).toHaveBeenCalledTimes(2); expect(requireMockArg(broadcastToConnIds, 0, 0, "run tool event")).toBe("agent"); @@ -2056,18 +1708,18 @@ describe("agent event handler", () => { toolEventRecipients.add("run-heartbeat-tool", "conn-run"); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "run-heartbeat-tool", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-heartbeat-tool", + "tool", + { phase: "start", name: "read", toolCallId: "tool-heartbeat-1", args: { path: "HEARTBEAT.md" }, }, - }); + { ts: 1_234 }, + ); expect(broadcastToConnIds).not.toHaveBeenCalled(); const nodeToolCalls = nodeSendToSession.mock.calls.filter(([, event]) => event === "agent"); @@ -2079,36 +1731,23 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-1", }); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "session-1", - kind: "direct", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", - updatedAt: 1_200, - }); + vi.mocked(loadGatewaySessionRow).mockReturnValue(OWNED_SESSION_ROW); registerAgentRunContext("run-tool-owner", { sessionKey: "session-1", verboseLevel: "off" }); toolEventRecipients.add("run-tool-owner", "conn-run"); - handler({ - runId: "run-tool-owner", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-tool-owner", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-run-1", args: { command: "echo hi" }, }, - }); + { ts: 1_234 }, + ); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); expect(requireMockArg(broadcastToConnIds, 0, 0, "run tool event")).toBe("agent"); @@ -2116,15 +1755,7 @@ describe("agent event handler", () => { expectRecordFields(runToolPayload, { runId: "run-tool-owner", sessionKey: "session-1", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", + ...SESSION_OWNERSHIP, stream: "tool", ts: 1_234, }); @@ -2152,12 +1783,11 @@ describe("agent event handler", () => { verboseLevel: "on", }); - handler({ - runId: "run-tool-search-node", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-tool-search-node", + "tool", + { phase: "start", name: "tool_search_code", toolCallId: "tool-search-node-1", @@ -2165,7 +1795,8 @@ describe("agent event handler", () => { code: 'return await openclaw.tools.call("openclaw:core:exec", { command: "echo hi" });', }, }, - }); + { ts: 1_234 }, + ); const payload = requireMockArg(nodeSendToSession, 0, 2, "node tool-search payload") as { stream?: string; @@ -2201,35 +1832,22 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-1", }); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "session-1", - kind: "direct", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", - updatedAt: 1_200, - }); + vi.mocked(loadGatewaySessionRow).mockReturnValue(OWNED_SESSION_ROW); registerAgentRunContext("run-tool-node", { sessionKey: "session-1", verboseLevel: "on" }); - handler({ - runId: "run-tool-node", - seq: 1, - stream: "tool", - ts: 1_234, - data: { + emitAgentEvent( + handler, + "run-tool-node", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-node-1", args: { command: "echo hi" }, }, - }); + { ts: 1_234 }, + ); expect(requireMockArg(nodeSendToSession, 0, 0, "node tool session")).toBe("session-1"); expect(requireMockArg(nodeSendToSession, 0, 1, "node tool event")).toBe("agent"); @@ -2237,15 +1855,7 @@ describe("agent event handler", () => { expectRecordFields(nodeToolPayload, { runId: "run-tool-node", sessionKey: "session-1", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - lastThreadId: 42, - fastMode: true, - verboseLevel: "on", + ...SESSION_OWNERSHIP, stream: "tool", ts: 1_234, }); @@ -2283,27 +1893,27 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-finished", - seq: 1, - stream: "lifecycle", - ts: 1_000, - data: { + emitAgentEvent( + handler, + "run-finished", + "lifecycle", + { phase: "start", startedAt: 900, }, - }); - handler({ - runId: "run-finished", - seq: 2, - stream: "lifecycle", - ts: 1_800, - data: { + { ts: 1_000 }, + ); + emitAgentEvent( + handler, + "run-finished", + "lifecycle", + { phase: "end", startedAt: 900, endedAt: 1_700, }, - }); + { seq: 2, ts: 1_800 }, + ); await waitForFast(() => { expect( @@ -2368,31 +1978,27 @@ describe("agent event handler", () => { }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "old-run", - seq: 1, - stream: "lifecycle", - sessionKey: "session-reset", - sessionId: "old-session", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "old-run", + "lifecycle", + { phase: "start", startedAt: 2_100, }, - }); - handler({ - runId: "old-run", - seq: 2, - stream: "lifecycle", - sessionKey: "session-reset", - sessionId: "old-session", - ts: 2_200, - data: { + { sessionKey: "session-reset", sessionId: "old-session", ts: 2_100 }, + ); + emitAgentEvent( + handler, + "old-run", + "lifecycle", + { phase: "error", endedAt: 2_200, error: "old run failed", }, - }); + { seq: 2, sessionKey: "session-reset", sessionId: "old-session", ts: 2_200 }, + ); await waitForFast(() => { expect( @@ -2425,11 +2031,8 @@ describe("agent event handler", () => { }); it("suppresses late interrupted pre-restart lifecycle events from live projections", () => { - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", @@ -2441,10 +2044,8 @@ describe("agent event handler", () => { }, ], }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); const { broadcast, broadcastToConnIds, @@ -2458,26 +2059,26 @@ describe("agent event handler", () => { lifecycleErrorRetryGraceMs: 0, }); sessionEventSubscribers.subscribe("conn-session"); - chatRunState.registry.add("interrupted-run", { - sessionKey: "session-recovery", - clientRunId: "interrupted-run", - }); + registerChatRun(chatRunState, "interrupted-run", "session-recovery", "interrupted-run"); - handler({ - runId: "interrupted-run", - lifecycleGeneration: "pre-restart", - seq: 2, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "interrupted-run", + "lifecycle", + { phase: "end", aborted: true, stopReason: "restart", endedAt: 2_100, }, - }); + { + seq: 2, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); expect( @@ -2494,11 +2095,8 @@ describe("agent event handler", () => { }); it("projects successful completion when a restart marker was persisted before abort", async () => { - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", @@ -2510,10 +2108,8 @@ describe("agent event handler", () => { }, ], }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); const markTrackedRunTerminalPersisted = vi.fn(); const trackTrackedRunTerminalPersistence = vi.fn(); const { @@ -2531,24 +2127,29 @@ describe("agent event handler", () => { trackTrackedRunTerminalPersistence, }); sessionEventSubscribers.subscribe("conn-session"); - chatRunState.registry.add("completed-during-marker-write", { - sessionKey: "session-recovery", - clientRunId: "completed-during-marker-write", - }); + registerChatRun( + chatRunState, + "completed-during-marker-write", + "session-recovery", + "completed-during-marker-write", + ); - handler({ - runId: "completed-during-marker-write", - lifecycleGeneration: "pre-restart", - seq: 2, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "completed-during-marker-write", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { + seq: 2, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); expect(chatBroadcastCalls(broadcast)).toHaveLength(1); expect(persistGatewaySessionLifecycleEventMock).toHaveBeenCalledTimes(1); @@ -2590,21 +2191,16 @@ describe("agent event handler", () => { lifecycleGeneration: "pre-restart", }, ]; - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", abortedLastRun: true, restartRecoveryRuns, }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); vi.mocked(loadGatewaySessionRow).mockReturnValue({ key: "session-recovery", kind: "direct", @@ -2619,19 +2215,22 @@ describe("agent event handler", () => { }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "completed-run", - lifecycleGeneration: "pre-restart", - seq: 2, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "completed-run", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { + seq: 2, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); await waitForFast(() => { expect( @@ -2661,21 +2260,16 @@ describe("agent event handler", () => { lifecycleGeneration: "pre-restart-b", }, ]; - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", abortedLastRun: true, restartRecoveryRuns, }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); let currentRow = { key: "session-recovery", kind: "direct" as const, @@ -2753,19 +2347,14 @@ describe("agent event handler", () => { }); it("reloads canonical state when a restart marker races terminal persistence", async () => { - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); let currentRow = { key: "session-recovery", kind: "direct" as const, @@ -2788,19 +2377,22 @@ describe("agent event handler", () => { }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "completed-run", - lifecycleGeneration: "pre-restart", - seq: 2, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "completed-run", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { + seq: 2, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); expect( broadcastToConnIds.mock.calls.filter(([event]) => event === "sessions.changed"), @@ -2845,18 +2437,16 @@ describe("agent event handler", () => { }); sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "run-failed-write", - seq: 2, - stream: "lifecycle", - sessionKey: "session-failed-write", - sessionId: "session-failed-write", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "run-failed-write", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { seq: 2, sessionKey: "session-failed-write", sessionId: "session-failed-write", ts: 2_100 }, + ); await waitForFast(() => { expect( @@ -2876,11 +2466,8 @@ describe("agent event handler", () => { }); it("does not clear a same-id retry when an old restart terminal arrives", () => { - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", @@ -2891,10 +2478,8 @@ describe("agent event handler", () => { }, ], }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); registerAgentRunContext("shared-run", { sessionKey: "session-recovery", sessionId: "session-recovery", @@ -2907,25 +2492,25 @@ describe("agent event handler", () => { resolveActiveLifecycleGenerationForRun: () => "post-restart", }); agentRunSeq.set("shared-run", 4); - chatRunState.registry.add("shared-run", { - sessionKey: "session-recovery", - clientRunId: "shared-run", - }); + registerChatRun(chatRunState, "shared-run", "session-recovery", "shared-run"); chatRunState.buffers.set("shared-run", "new retry output"); - handler({ - runId: "shared-run", - lifecycleGeneration: "pre-restart", - seq: 3, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "shared-run", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { + seq: 3, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); expect(chatRunState.registry.peek("shared-run")).toBeDefined(); expect(chatRunState.buffers.get("shared-run")).toBe("new retry output"); @@ -2948,32 +2533,28 @@ describe("agent event handler", () => { lifecycleErrorRetryGraceMs: 100, resolveActiveLifecycleGenerationForRun: () => activeLifecycleGeneration, }); - chatRunState.registry.add("shared-run", { - sessionKey: "session-recovery", - clientRunId: "shared-run", - }); + registerChatRun(chatRunState, "shared-run", "session-recovery", "shared-run"); - handler({ - runId: "shared-run", - lifecycleGeneration: "pre-restart", - seq: 1, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_000, - data: { + emitAgentEvent( + handler, + "shared-run", + "lifecycle", + { phase: "error", error: "retryable provider failure", endedAt: 2_000, }, - }); + { + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_000, + }, + ); expect(vi.getTimerCount()).toBe(1); - vi.mocked(loadSessionEntry).mockReturnValue({ - cfg: {}, - storePath: "/tmp/sessions.json", - store: {}, - entry: { + mockSessionEntry( + { sessionId: "session-recovery", updatedAt: 2_000, status: "running", @@ -2984,10 +2565,8 @@ describe("agent event handler", () => { }, ], }, - canonicalKey: "session-recovery", - storeKeys: ["session-recovery"], - legacyKey: undefined, - }); + "session-recovery", + ); activeLifecycleGeneration = "post-restart"; registerAgentRunContext("shared-run", { sessionKey: "session-recovery", @@ -2995,19 +2574,22 @@ describe("agent event handler", () => { lifecycleGeneration: activeLifecycleGeneration, }); - handler({ - runId: "shared-run", - lifecycleGeneration: "pre-restart", - seq: 2, - stream: "lifecycle", - sessionKey: "session-recovery", - sessionId: "session-recovery", - ts: 2_100, - data: { + emitAgentEvent( + handler, + "shared-run", + "lifecycle", + { phase: "end", endedAt: 2_100, }, - }); + { + seq: 2, + lifecycleGeneration: "pre-restart", + sessionKey: "session-recovery", + sessionId: "session-recovery", + ts: 2_100, + }, + ); expect(vi.getTimerCount()).toBe(0); vi.advanceTimersByTime(100); @@ -3022,14 +2604,13 @@ describe("agent event handler", () => { lifecycleErrorRetryGraceMs: 100, }); - handler({ - runId: "run-dispose", - seq: 1, - stream: "lifecycle", - sessionKey: "session-dispose", - ts: 2_000, - data: { phase: "error", error: "retryable provider failure" }, - }); + emitAgentEvent( + handler, + "run-dispose", + "lifecycle", + { phase: "error", error: "retryable provider failure" }, + { sessionKey: "session-dispose", ts: 2_000 }, + ); expect(vi.getTimerCount()).toBe(1); handler.dispose(); @@ -3061,22 +2642,19 @@ describe("agent event handler", () => { handler, } = createHarness(); sessionEventSubscribers.subscribe("conn-session"); - chatRunState.registry.add("provider-run", { - sessionKey: "session-finished", - clientRunId: "client-run", - }); + registerChatRun(chatRunState, "provider-run", "session-finished", "client-run"); - handler({ - runId: "provider-run", - seq: 2, - stream: "lifecycle", - ts: 1_800, - data: { + emitAgentEvent( + handler, + "provider-run", + "lifecycle", + { phase: "end", startedAt: 900, endedAt: 1_700, }, - }); + { seq: 2, ts: 1_800 }, + ); expect(clearTrackedActiveRun).toHaveBeenCalledWith({ runId: "provider-run", @@ -3114,22 +2692,19 @@ describe("agent event handler", () => { } }, }); - chatRunState.registry.add("provider-run", { - sessionKey: "session-finished", - clientRunId: "client-run", - }); + registerChatRun(chatRunState, "provider-run", "session-finished", "client-run"); - handler({ - runId: "provider-run", - seq: 2, - stream: "lifecycle", - ts: 1_800, - data: { + emitAgentEvent( + handler, + "provider-run", + "lifecycle", + { phase: "end", startedAt: 900, endedAt: 1_700, }, - }); + { seq: 2, ts: 1_800 }, + ); const providerGuard = trackedActiveRuns.get("provider-run"); const retryGuard = trackedActiveRuns.get("client-run"); @@ -3141,19 +2716,16 @@ describe("agent event handler", () => { it("keeps aborted chat run markers through terminal lifecycle cleanup", () => { const { broadcast, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-aborted", { - sessionKey: "session-aborted", - clientRunId: "client-aborted", - }); + registerNamedChatRun(chatRunState, "aborted"); chatRunState.abortedRuns.set("client-aborted", createChatAbortMarker()); - handler({ - runId: "run-aborted", - seq: 2, - stream: "lifecycle", - ts: 1_500, - data: { phase: "end", aborted: true, stopReason: "rpc" }, - }); + emitAgentEvent( + handler, + "run-aborted", + "lifecycle", + { phase: "end", aborted: true, stopReason: "rpc" }, + { seq: 2, ts: 1_500 }, + ); expect(chatRunState.abortedRuns.has("client-aborted")).toBe(true); expect(chatRunState.registry.peek("run-aborted")).toBeUndefined(); @@ -3162,22 +2734,24 @@ describe("agent event handler", () => { it("projects lifecycle self-aborts with their validation diagnostic", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("provider-validation-loop", { - sessionKey: "session-validation-loop", - clientRunId: "client-validation-loop", - }); + registerChatRun( + chatRunState, + "provider-validation-loop", + "session-validation-loop", + "client-validation-loop", + ); - handler({ - runId: "provider-validation-loop", - seq: 2, - stream: "lifecycle", - ts: 1_500, - data: { + emitAgentEvent( + handler, + "provider-validation-loop", + "lifecycle", + { phase: "end", aborted: true, toolErrorSummary: "edit tool validation failed: edits: must be an array", }, - }); + { seq: 2, ts: 1_500 }, + ); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls).toHaveLength(1); @@ -3198,18 +2772,20 @@ describe("agent event handler", () => { { stopReason: "timeout", expectedState: "error" }, ])("preserves $stopReason lifecycle abort classification", ({ stopReason, expectedState }) => { const { broadcast, chatRunState, handler } = createHarness(); - chatRunState.registry.add(`provider-${stopReason}`, { - sessionKey: `session-${stopReason}`, - clientRunId: `client-${stopReason}`, - }); + registerChatRun( + chatRunState, + `provider-${stopReason}`, + `session-${stopReason}`, + `client-${stopReason}`, + ); - handler({ - runId: `provider-${stopReason}`, - seq: 2, - stream: "lifecycle", - ts: 1_500, - data: { phase: "end", aborted: true, stopReason }, - }); + emitAgentEvent( + handler, + `provider-${stopReason}`, + "lifecycle", + { phase: "end", aborted: true, stopReason }, + { seq: 2, ts: 1_500 }, + ); expect( expectDefined( @@ -3225,23 +2801,25 @@ describe("agent event handler", () => { it("does not forward unsafe lifecycle abort diagnostics", () => { const { broadcast, chatRunState, handler } = createHarness(); - chatRunState.registry.add("provider-unsafe-abort", { - sessionKey: "session-unsafe-abort", - clientRunId: "client-unsafe-abort", - }); + registerChatRun( + chatRunState, + "provider-unsafe-abort", + "session-unsafe-abort", + "client-unsafe-abort", + ); - handler({ - runId: "provider-unsafe-abort", - seq: 2, - stream: "lifecycle", - ts: 1_500, - data: { + emitAgentEvent( + handler, + "provider-unsafe-abort", + "lifecycle", + { phase: "end", aborted: true, stopReason: "aborted", toolErrorSummary: "browser failed\nsecret output", }, - }); + { seq: 2, ts: 1_500 }, + ); const payload = expectDefined( chatBroadcastCalls(broadcast)[0], @@ -3253,17 +2831,13 @@ describe("agent event handler", () => { it("preserves timeout terminal precedence for abort-marked lifecycle events", () => { const { broadcast, chatRunState, handler } = createHarness(); - chatRunState.registry.add("provider-timeout", { - sessionKey: "session-timeout", - clientRunId: "client-timeout", - }); + registerChatRun(chatRunState, "provider-timeout", "session-timeout", "client-timeout"); - handler({ - runId: "provider-timeout", - seq: 2, - stream: "lifecycle", - ts: 1_500, - data: { + emitAgentEvent( + handler, + "provider-timeout", + "lifecycle", + { phase: "end", aborted: true, stopReason: "timeout", @@ -3271,7 +2845,8 @@ describe("agent event handler", () => { providerStarted: true, error: "agent provider timeout", }, - }); + { seq: 2, ts: 1_500 }, + ); const payload = expectDefined( chatBroadcastCalls(broadcast)[0], @@ -3300,25 +2875,22 @@ describe("agent event handler", () => { ({ marker }) => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness({ now: 2_000 }); chatRunState.abortedRuns.set("client-stale-abort", marker()); - chatRunState.registry.add("run-stale-abort", { - sessionKey: "session-stale-abort", - clientRunId: "client-stale-abort", - }); + registerNamedChatRun(chatRunState, "stale-abort"); - handler({ - runId: "run-stale-abort", - seq: 1, - stream: "assistant", - ts: 2_100, - data: { text: "Fresh output", delta: "Fresh output" }, - }); - handler({ - runId: "run-stale-abort", - seq: 2, - stream: "lifecycle", - ts: 2_200, - data: { phase: "end" }, - }); + emitAgentEvent( + handler, + "run-stale-abort", + "assistant", + { text: "Fresh output", delta: "Fresh output" }, + { ts: 2_100 }, + ); + emitAgentEvent( + handler, + "run-stale-abort", + "lifecycle", + { phase: "end" }, + { seq: 2, ts: 2_200 }, + ); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls).toHaveLength(2); @@ -3334,26 +2906,23 @@ describe("agent event handler", () => { it("honors same-millisecond abort markers from the current same-key run", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness({ now: 3_000 }); - chatRunState.registry.add("run-current-abort", { - sessionKey: "session-current-abort", - clientRunId: "client-current-abort", - }); + registerNamedChatRun(chatRunState, "current-abort"); chatRunState.abortedRuns.set("client-current-abort", createChatAbortMarker()); - handler({ - runId: "run-current-abort", - seq: 1, - stream: "assistant", - ts: 3_100, - data: { text: "Suppressed output", delta: "Suppressed output" }, - }); - handler({ - runId: "run-current-abort", - seq: 2, - stream: "lifecycle", - ts: 3_200, - data: { phase: "end", aborted: true, stopReason: "rpc" }, - }); + emitAgentEvent( + handler, + "run-current-abort", + "assistant", + { text: "Suppressed output", delta: "Suppressed output" }, + { ts: 3_100 }, + ); + emitAgentEvent( + handler, + "run-current-abort", + "lifecycle", + { phase: "end", aborted: true, stopReason: "rpc" }, + { seq: 2, ts: 3_200 }, + ); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); expect(sessionChatCalls(nodeSendToSession)).toHaveLength(0); @@ -3366,15 +2935,8 @@ describe("agent event handler", () => { key: "session-finished", kind: "direct", updatedAt: 1_650, - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - fastMode: true, + ...SESSION_OWNERSHIP, sendPolicy: "deny", - verboseLevel: "on", responseUsage: "full", totalTokens: 42, totalTokensFresh: true, @@ -3397,17 +2959,17 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-finished", - seq: 2, - stream: "lifecycle", - ts: 1_800, - data: { + emitAgentEvent( + handler, + "run-finished", + "lifecycle", + { phase: "end", startedAt: 900, endedAt: 1_700, }, - }); + { seq: 2, ts: 1_800 }, + ); await waitForFast(() => { expect( @@ -3420,15 +2982,8 @@ describe("agent event handler", () => { expectPayloadFields(requireMockArg(broadcastToConnIds, 0, 1, "sessions changed payload"), { sessionKey: "session-finished", phase: "end", - spawnedBy: "agent:main:main", - spawnedWorkspaceDir: "/tmp/subagent", - forkedFromParent: true, - spawnDepth: 2, - subagentRole: "orchestrator", - subagentControlScope: "children", - fastMode: true, + ...SESSION_OWNERSHIP, sendPolicy: "deny", - verboseLevel: "on", responseUsage: "full", totalTokens: 42, totalTokensFresh: true, @@ -3469,13 +3024,13 @@ describe("agent event handler", () => { sessionEventSubscribers.subscribe("conn-session"); - handler({ - runId: "run-global", - seq: 2, - stream: "lifecycle", - ts: 1_800, - data: { phase: "end", endedAt: 1_700 }, - }); + emitAgentEvent( + handler, + "run-global", + "lifecycle", + { phase: "end", endedAt: 1_700 }, + { seq: 2, ts: 1_800 }, + ); await waitForFast(() => { expect( @@ -3490,63 +3045,42 @@ describe("agent event handler", () => { expect(requireRecord(payload.session, "nested session")).not.toHaveProperty("goal"); }); - it("keeps tool output for Control UI recipients when verbose is on", () => { + it.each([ + { + name: "keeps tool output for Control UI recipients when verbose is on", + runId: "run-tool-on", + toolCallId: "t3", + verboseLevel: "on", + partialResult: { content: [{ type: "text", text: "partial" }] }, + }, + { + name: "keeps tool output when verbose is full", + runId: "run-tool-full", + toolCallId: "t4", + verboseLevel: "full", + partialResult: undefined, + }, + ] as const)("$name", ({ runId, toolCallId, verboseLevel, partialResult }) => { const { broadcastToConnIds, toolEventRecipients, handler } = createHarness({ resolveSessionKeyForRun: () => "session-1", }); - - registerAgentRunContext("run-tool-on", { sessionKey: "session-1", verboseLevel: "on" }); - toolEventRecipients.add("run-tool-on", "conn-1"); - - handler({ - runId: "run-tool-on", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { - phase: "result", - name: "exec", - toolCallId: "t3", - result: { content: [{ type: "text", text: "secret" }] }, - partialResult: { content: [{ type: "text", text: "partial" }] }, - }, + const result = { content: [{ type: "text", text: "secret" }] }; + registerAgentRunContext(runId, { sessionKey: "session-1", verboseLevel }); + toolEventRecipients.add(runId, "conn-1"); + emitAgentEvent(handler, runId, "tool", { + phase: "result", + name: "exec", + toolCallId, + result, + ...(partialResult ? { partialResult } : {}), }); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); const payload = requireMockArg(broadcastToConnIds, 0, 1, "tool output payload") as { data?: Record; }; - expect(payload.data?.result).toEqual({ content: [{ type: "text", text: "secret" }] }); - expect(payload.data?.partialResult).toEqual({ content: [{ type: "text", text: "partial" }] }); - }); - - it("keeps tool output when verbose is full", () => { - const { broadcastToConnIds, toolEventRecipients, handler } = createHarness({ - resolveSessionKeyForRun: () => "session-1", - }); - - registerAgentRunContext("run-tool-full", { sessionKey: "session-1", verboseLevel: "full" }); - toolEventRecipients.add("run-tool-full", "conn-1"); - - const result = { content: [{ type: "text", text: "secret" }] }; - handler({ - runId: "run-tool-full", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { - phase: "result", - name: "exec", - toolCallId: "t4", - result, - }, - }); - - expect(broadcastToConnIds).toHaveBeenCalledTimes(1); - const payload = requireMockArg(broadcastToConnIds, 0, 1, "full tool output payload") as { - data?: Record; - }; expect(payload.data?.result).toEqual(result); + expect(payload.data?.partialResult).toEqual(partialResult); }); it("broadcasts fallback events to agent subscribers and node session", () => { @@ -3571,10 +3105,12 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness({ resolveSessionKeyForRun: () => "session-fallback", }); - chatRunState.registry.add("run-fallback-internal", { - sessionKey: "session-fallback", - clientRunId: "run-fallback-client", - }); + registerChatRun( + chatRunState, + "run-fallback-internal", + "session-fallback", + "run-fallback-client", + ); emitFallbackLifecycle({ handler, runId: "run-fallback-internal" }); @@ -3591,19 +3127,11 @@ describe("agent event handler", () => { it("keeps selected-agent global chat events scoped to the linked agent", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness(); - chatRunState.registry.add("run-global-main", { - sessionKey: "global", + registerChatRun(chatRunState, "run-global-main", "global", "client-global-main", { agentId: "main", - clientRunId: "client-global-main", }); - handler({ - runId: "run-global-main", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "main global reply" }, - }); + emitAgentEvent(handler, "run-global-main", "assistant", { text: "main global reply" }); const chatPayload = chatBroadcastCalls(broadcast)[0]?.[1] as { agentId?: string; @@ -3623,19 +3151,11 @@ describe("agent event handler", () => { it("persists selected-agent global lifecycle state with the linked agent", () => { const { broadcastToConnIds, chatRunState, handler, sessionEventSubscribers } = createHarness(); sessionEventSubscribers.subscribe("conn-1"); - chatRunState.registry.add("run-global-work", { - sessionKey: "global", + registerChatRun(chatRunState, "run-global-work", "global", "client-global-work", { agentId: "work", - clientRunId: "client-global-work", }); - handler({ - runId: "run-global-work", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); + emitAgentEvent(handler, "run-global-work", "lifecycle", { phase: "start" }); expect(persistGatewaySessionLifecycleEventMock).toHaveBeenCalledWith( expect.objectContaining({ @@ -3658,20 +3178,12 @@ describe("agent event handler", () => { it("logs when start session persistence fails", async () => { const { chatRunState, handler, sessionEventSubscribers } = createHarness(); sessionEventSubscribers.subscribe("conn-1"); - chatRunState.registry.add("run-global-work", { - sessionKey: "global", + registerChatRun(chatRunState, "run-global-work", "global", "client-global-work", { agentId: "work", - clientRunId: "client-global-work", }); persistGatewaySessionLifecycleEventMock.mockRejectedValueOnce(new Error("start disk full")); - handler({ - runId: "run-global-work", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); + emitAgentEvent(handler, "run-global-work", "lifecycle", { phase: "start" }); await waitForFast(() => { expect(logErrorMock).toHaveBeenCalledTimes(1); @@ -3686,23 +3198,15 @@ describe("agent event handler", () => { createHarness(); sessionMessageSubscribers.subscribe("conn-main", "agent:main:global"); sessionMessageSubscribers.subscribe("conn-work", "agent:work:global"); - chatRunState.registry.add("run-hidden-main", { - sessionKey: "global", + registerChatRun(chatRunState, "run-hidden-main", "global", "client-hidden-main", { agentId: "main", - clientRunId: "client-hidden-main", }); registerAgentRunContext("run-hidden-main", { sessionKey: "global", isControlUiVisible: false, }); - handler({ - runId: "run-hidden-main", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "hidden main global reply" }, - }); + emitAgentEvent(handler, "run-hidden-main", "assistant", { text: "hidden main global reply" }); const chatCall = broadcastToConnIds.mock.calls.find(([event]) => event === "chat"); expect(chatCall?.[2]).toEqual(new Set(["conn-main"])); @@ -3722,21 +3226,14 @@ describe("agent event handler", () => { createHarness(); sessionMessageSubscribers.subscribe("conn-main", "agent:main:global"); sessionMessageSubscribers.subscribe("conn-ops", "agent:ops:global"); - chatRunState.registry.add("run-hidden-default", { - sessionKey: "global", - clientRunId: "client-hidden-default", - }); + registerChatRun(chatRunState, "run-hidden-default", "global", "client-hidden-default"); registerAgentRunContext("run-hidden-default", { sessionKey: "global", isControlUiVisible: false, }); - handler({ - runId: "run-hidden-default", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "hidden default global reply" }, + emitAgentEvent(handler, "run-hidden-default", "assistant", { + text: "hidden default global reply", }); const chatCall = broadcastToConnIds.mock.calls.find(([event]) => event === "chat"); @@ -3754,25 +3251,12 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-fallback", lifecycleErrorRetryGraceMs: 100, }); - chatRunState.registry.add("run-fallback-retry", { - sessionKey: "session-fallback", - clientRunId: "run-fallback-client", - }); + registerChatRun(chatRunState, "run-fallback-retry", "session-fallback", "run-fallback-client"); - handler({ - runId: "run-fallback-retry", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "draft" }, - }); - handler({ - runId: "run-fallback-retry", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "provider failed" }, - }); + emitAgentEvents(handler, "run-fallback-retry", [ + ["assistant", { text: "draft" }], + ["lifecycle", { phase: "error", error: "provider failed" }], + ]); expect(chatRunState.registry.peek("run-fallback-retry")).toMatchObject({ sessionKey: "session-fallback", @@ -3834,20 +3318,10 @@ describe("agent event handler", () => { }); registerAgentRunContext("run-terminal-error", { sessionKey: "session-terminal-error" }); - handler({ - runId: "run-terminal-error", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "partial" }, - }); - handler({ - runId: "run-terminal-error", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "still broken" }, - }); + emitAgentEvents(handler, "run-terminal-error", [ + ["assistant", { text: "partial" }], + ["lifecycle", { phase: "error", error: "still broken" }], + ]); expect(clearAgentRunContext).not.toHaveBeenCalled(); expect(agentRunSeq.get("run-terminal-error")).toBe(2); @@ -3879,16 +3353,10 @@ describe("agent event handler", () => { sessionKey: "session-terminal-error", }); - handler({ - runId: "run-terminal-final-failure", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { - phase: "error", - error: "LLM request failed: network connection error.", - fallbackExhaustedFailure: true, - }, + emitAgentEvent(handler, "run-terminal-final-failure", "lifecycle", { + phase: "error", + error: "LLM request failed: network connection error.", + fallbackExhaustedFailure: true, }); const finalPayload = chatBroadcastCalls(broadcast).at(-1)?.[1] as { @@ -3920,27 +3388,11 @@ describe("agent event handler", () => { sessionKey: "session-terminal-error", }); - handler({ - runId: "run-terminal-late-tool", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); - handler({ - runId: "run-terminal-late-tool", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "request timed out" }, - }); - handler({ - runId: "run-terminal-late-tool", - seq: 3, - stream: "tool", - ts: Date.now(), - data: { phase: "result", name: "exec" }, - }); + emitAgentEvents(handler, "run-terminal-late-tool", [ + ["lifecycle", { phase: "start" }], + ["lifecycle", { phase: "error", error: "request timed out" }], + ["tool", { phase: "result", name: "exec" }], + ]); vi.advanceTimersByTime(99); @@ -3982,27 +3434,11 @@ describe("agent event handler", () => { sessionKey: "session-terminal-error", }); - handler({ - runId: "run-terminal-late-lifecycle", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); - handler({ - runId: "run-terminal-late-lifecycle", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "request timed out" }, - }); - handler({ - runId: "run-terminal-late-lifecycle", - seq: 3, - stream: "lifecycle", - ts: Date.now(), - data: { msg: "status update" }, - }); + emitAgentEvents(handler, "run-terminal-late-lifecycle", [ + ["lifecycle", { phase: "start" }], + ["lifecycle", { phase: "error", error: "request timed out" }], + ["lifecycle", { msg: "status update" }], + ]); vi.advanceTimersByTime(100); @@ -4028,27 +3464,11 @@ describe("agent event handler", () => { sessionKey: "session-terminal-retry", }); - handler({ - runId: "run-terminal-retry", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); - handler({ - runId: "run-terminal-retry", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "attempt failed" }, - }); - handler({ - runId: "run-terminal-retry", - seq: 3, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); + emitAgentEvents(handler, "run-terminal-retry", [ + ["lifecycle", { phase: "start" }], + ["lifecycle", { phase: "error", error: "attempt failed" }], + ["lifecycle", { phase: "start" }], + ]); vi.advanceTimersByTime(100); @@ -4149,15 +3569,9 @@ describe("agent event handler", () => { }); registerAgentRunContext("run-detected-error", { sessionKey: "session-detected-error" }); - handler({ - runId: "run-detected-error", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { - phase: "error", - error: Object.assign(new Error("Too many requests"), { code: 429 }), - }, + emitAgentEvent(handler, "run-detected-error", "lifecycle", { + phase: "error", + error: Object.assign(new Error("Too many requests"), { code: 429 }), }); const payload = chatBroadcastCalls(broadcast).at(-1)?.[1] as { @@ -4186,20 +3600,10 @@ describe("agent event handler", () => { }); registerAgentRunContext("run-chat-send", { sessionKey: "session-chat-send" }); - handler({ - runId: "run-chat-send", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "partial" }, - }); - handler({ - runId: "run-chat-send", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "chat.send failed" }, - }); + emitAgentEvents(handler, "run-chat-send", [ + ["assistant", { text: "partial" }], + ["lifecycle", { phase: "error", error: "chat.send failed" }], + ]); vi.advanceTimersByTime(100); @@ -4219,18 +3623,12 @@ describe("agent event handler", () => { lifecycleErrorRetryGraceMs: 100, isChatSendRunActive: (runId) => runId === "run-chat-send", }); - chatRunState.registry.add("run-chat-send", { - sessionKey: "session-chat-send", - clientRunId: "run-chat-send", - }); + registerChatRun(chatRunState, "run-chat-send", "session-chat-send", "run-chat-send"); registerAgentRunContext("run-chat-send", { sessionKey: "session-chat-send" }); - handler({ - runId: "run-chat-send", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "chat.send failed" }, + emitAgentEvent(handler, "run-chat-send", "lifecycle", { + phase: "error", + error: "chat.send failed", }); vi.advanceTimersByTime(100); @@ -4262,13 +3660,7 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-hidden", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Reply from quietchat" }, - }); + emitAgentEvent(handler, "run-hidden", "assistant", { text: "Reply from quietchat" }); emitLifecycleEnd(handler, "run-hidden", 2); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); @@ -4299,12 +3691,8 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-hidden", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "visible only to the selected session" }, + emitAgentEvent(handler, "run-hidden", "assistant", { + text: "visible only to the selected session", }); emitLifecycleEnd(handler, "run-hidden", 2); @@ -4387,60 +3775,27 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-hidden-commentary", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { - text: "I will inspect the files first.", - delta: "I will inspect the files first.", - phase: "commentary", - }, - }); - handler({ - runId: "run-hidden-commentary", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { - text: "Untagged text frame must not mirror.", - delta: "Untagged text frame must not mirror.", - }, - }); - handler({ - runId: "run-hidden-commentary", - seq: 3, - stream: "assistant", - ts: Date.now(), - data: { - delta: "Untagged delta-only stream must not mirror.", - }, - }); - handler({ - runId: "run-hidden-commentary", - seq: 4, - stream: "assistant", - ts: Date.now(), - data: { text: "Terminal echo without delta" }, - }); - handler({ - runId: "run-hidden-commentary", - seq: 5, - stream: "assistant", - ts: Date.now(), - data: { text: "Final answer", delta: "Final answer", phase: "final_answer" }, - }); - handler({ - runId: "run-hidden-commentary", - seq: 6, - stream: "assistant", - ts: Date.now(), - data: { - delta: "Streaming commentary delta.", - phase: "commentary", - }, - }); + emitAgentEvents(handler, "run-hidden-commentary", [ + [ + "assistant", + { + text: "I will inspect the files first.", + delta: "I will inspect the files first.", + phase: "commentary", + }, + ], + [ + "assistant", + { + text: "Untagged text frame must not mirror.", + delta: "Untagged text frame must not mirror.", + }, + ], + ["assistant", { delta: "Untagged delta-only stream must not mirror." }], + ["assistant", { text: "Terminal echo without delta" }], + ["assistant", { text: "Final answer", delta: "Final answer", phase: "final_answer" }], + ["assistant", { delta: "Streaming commentary delta.", phase: "commentary" }], + ]); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); expect(agentBroadcastCalls(broadcast)).toHaveLength(0); @@ -4502,16 +3857,10 @@ describe("agent event handler", () => { }); chatRunState.abortedRuns.set("run-hidden-commentary-aborted", 1_000); - handler({ - runId: "run-hidden-commentary-aborted", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { - text: "This aborted commentary must not be mirrored.", - delta: "This aborted commentary must not be mirrored.", - phase: "commentary", - }, + emitAgentEvent(handler, "run-hidden-commentary-aborted", "assistant", { + text: "This aborted commentary must not be mirrored.", + delta: "This aborted commentary must not be mirrored.", + phase: "commentary", }); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); @@ -4534,17 +3883,11 @@ describe("agent event handler", () => { verboseLevel: "off", }); - handler({ - runId: "run-hidden", - seq: 1, - stream: "item", - ts: Date.now(), - data: { - kind: "status", - title: "Fast", - phase: "update", - summary: "💨Fast: auto-off(8s>=5s)", - }, + emitAgentEvent(handler, "run-hidden", "item", { + kind: "status", + title: "Fast", + phase: "update", + summary: "💨Fast: auto-off(8s>=5s)", }); expect(agentBroadcastCalls(broadcast)).toHaveLength(0); @@ -4586,27 +3929,18 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "session-tool-remap", }); - chatRunState.registry.add("run-tool-internal", { - sessionKey: "session-tool-remap", - clientRunId: "run-tool-client", - }); + registerChatRun(chatRunState, "run-tool-internal", "session-tool-remap", "run-tool-client"); registerAgentRunContext("run-tool-internal", { sessionKey: "session-tool-remap", verboseLevel: "on", }); toolEventRecipients.add("run-tool-internal", "conn-1"); - handler({ - runId: "run-tool-internal", - seq: 1, - stream: "tool", - ts: Date.now(), - data: { - phase: "start", - name: "exec", - toolCallId: "tool-remap-1", - args: { command: "true" }, - }, + emitAgentEvent(handler, "run-tool-internal", "tool", { + phase: "start", + name: "exec", + toolCallId: "tool-remap-1", + args: { command: "true" }, }); expect(broadcastToConnIds).toHaveBeenCalledTimes(1); @@ -4620,24 +3954,15 @@ describe("agent event handler", () => { const { broadcast, nodeSendToSession, chatRunState, handler } = createHarness({ now: 2_000, }); - chatRunState.registry.add("run-heartbeat", { - sessionKey: "session-heartbeat", - clientRunId: "client-heartbeat", - }); + registerNamedChatRun(chatRunState, "heartbeat"); registerAgentRunContext("run-heartbeat", { sessionKey: "session-heartbeat", isHeartbeat: true, verboseLevel: "off", }); - handler({ - runId: "run-heartbeat", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { - text: "HEARTBEAT_OK Read HEARTBEAT.md if it exists (workspace context). Follow it strictly.", - }, + emitAgentEvent(handler, "run-heartbeat", "assistant", { + text: "HEARTBEAT_OK Read HEARTBEAT.md if it exists (workspace context). Follow it strictly.", }); expect(chatBroadcastCalls(broadcast)).toHaveLength(0); @@ -4656,10 +3981,7 @@ describe("agent event handler", () => { }); const { broadcast, chatRunState, handler } = createHarness({ now: 3_000 }); - chatRunState.registry.add("run-heartbeat-alert", { - sessionKey: "session-heartbeat-alert", - clientRunId: "client-heartbeat-alert", - }); + registerNamedChatRun(chatRunState, "heartbeat-alert"); registerAgentRunContext("run-heartbeat-alert", { sessionKey: "session-heartbeat-alert", isHeartbeat: true, @@ -4667,14 +3989,8 @@ describe("agent event handler", () => { }); const alert = `Disk usage crossed 95 percent on /data. ${"Cleanup required. ".repeat(20)}`; - handler({ - runId: "run-heartbeat-alert", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { - text: `HEARTBEAT_OK ${alert}`, - }, + emitAgentEvent(handler, "run-heartbeat-alert", "assistant", { + text: `HEARTBEAT_OK ${alert}`, }); emitLifecycleEnd(handler, "run-heartbeat-alert"); @@ -4686,90 +4002,58 @@ describe("agent event handler", () => { }); describe("spawnedBy enrichment in chat and agent broadcasts", () => { - it("includes spawnedBy in chat delta broadcasts for subagent sessions", () => { + function mockSessionLineage(key: string, spawnedBy?: string) { vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:abc", + key, kind: "direct", updatedAt: null, - spawnedBy: "agent:conductor:task:parent-1", + ...(spawnedBy ? { spawnedBy } : {}), }); + } + + it.each([ + { + name: "includes spawnedBy in chat delta broadcasts for subagent sessions", + runId: "run-sub-1", + clientRunId: "client-sub-1", + events: [["assistant", { text: "hello from subagent" }]], + state: "delta", + expectsNodeCopy: true, + }, + { + name: "includes spawnedBy in chat final broadcasts for subagent sessions", + runId: "run-sub-final", + clientRunId: "client-sub-final", + events: [ + ["assistant", { text: "done" }], + ["lifecycle", { phase: "end" }], + ], + state: "final", + expectsNodeCopy: false, + }, + ] as const)("$name", ({ runId, clientRunId, events, state, expectsNodeCopy }) => { + mockSessionLineage("agent:coder:subagent:abc", "agent:conductor:task:parent-1"); const { broadcast, nodeSendToSession, handler, chatRunState } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:abc", }); + registerChatRun(chatRunState, runId, "agent:coder:subagent:abc", clientRunId); + emitAgentEvents(handler, runId, events); - chatRunState.registry.add("run-sub-1", { - sessionKey: "agent:coder:subagent:abc", - clientRunId: "client-sub-1", - }); - - handler({ - runId: "run-sub-1", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "hello from subagent" }, - }); - - const chatCalls = chatBroadcastCalls(broadcast); - expect(chatCalls.length).toBeGreaterThanOrEqual(1); - const [, payload] = expectDefined(chatCalls[0], "chatCalls[0] test invariant"); + const payload = requireCall( + chatBroadcastCalls(broadcast).find(([, candidate]) => candidate.state === state)?.[1], + `${state} chat payload`, + ); expectPayloadFields(payload, { sessionKey: "agent:coder:subagent:abc", spawnedBy: "agent:conductor:task:parent-1", - state: "delta", - }); - - const nodeCalls = sessionChatCalls(nodeSendToSession); - expect(nodeCalls.length).toBeGreaterThanOrEqual(1); - expectPayloadFields(nodeCalls[0]?.[2], { - spawnedBy: "agent:conductor:task:parent-1", - }); - }); - - it("includes spawnedBy in chat final broadcasts for subagent sessions", () => { - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:abc", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-1", - }); - - const { broadcast, handler, chatRunState } = createHarness({ - resolveSessionKeyForRun: () => "agent:coder:subagent:abc", - }); - - chatRunState.registry.add("run-sub-final", { - sessionKey: "agent:coder:subagent:abc", - clientRunId: "client-sub-final", - }); - - handler({ - runId: "run-sub-final", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "done" }, - }); - - handler({ - runId: "run-sub-final", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); - - const chatCalls = chatBroadcastCalls(broadcast); - const finalCall = requireCall( - chatCalls.find(([, p]) => p.state === "final"), - "final chat call", - ); - expectPayloadFields(finalCall[1], { - sessionKey: "agent:coder:subagent:abc", - spawnedBy: "agent:conductor:task:parent-1", - state: "final", + state, }); + if (expectsNodeCopy) { + expectPayloadFields(sessionChatCalls(nodeSendToSession)[0]?.[2], { + spawnedBy: "agent:conductor:task:parent-1", + }); + } }); it("marks a yielded final as waiting instead of parent-task completion", () => { @@ -4777,29 +4061,22 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "agent:main:main", }); - chatRunState.registry.add("run-yielded", { - sessionKey: "agent:main:main", - clientRunId: "client-yielded", + registerChatRun(chatRunState, "run-yielded", "agent:main:main", "client-yielded"); + emitAgentEvent(handler, "run-yielded", "assistant", { + text: "Waiting for registered continuation work.", }); - handler({ - runId: "run-yielded", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "Waiting for registered continuation work." }, - }); - handler({ - runId: "run-yielded", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { + emitAgentEvent( + handler, + "run-yielded", + "lifecycle", + { phase: "end", yielded: true, livenessState: "paused", stopReason: "end_turn", }, - }); + { seq: 2 }, + ); const finalCall = requireCall( chatBroadcastCalls(broadcast).find(([, payload]) => payload.state === "final"), @@ -4819,22 +4096,13 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "agent:main:main", }); - chatRunState.registry.add("run-aborted", { - sessionKey: "agent:main:main", - clientRunId: "client-aborted", - }); - handler({ - runId: "run-aborted", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { - phase: "end", - aborted: true, - yielded: true, - livenessState: "paused", - stopReason: "end_turn", - }, + registerChatRun(chatRunState, "run-aborted", "agent:main:main", "client-aborted"); + emitAgentEvent(handler, "run-aborted", "lifecycle", { + phase: "end", + aborted: true, + yielded: true, + livenessState: "paused", + stopReason: "end_turn", }); const finalCall = requireCall( @@ -4851,28 +4119,15 @@ describe("agent event handler", () => { }); it("omits spawnedBy from chat broadcasts for non-subagent sessions", () => { - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:main:main", - kind: "direct", - updatedAt: null, - }); + mockSessionLineage("agent:main:main"); const { broadcast, handler, chatRunState } = createHarness({ resolveSessionKeyForRun: () => "agent:main:main", }); - chatRunState.registry.add("run-main", { - sessionKey: "agent:main:main", - clientRunId: "client-main", - }); + registerChatRun(chatRunState, "run-main", "agent:main:main", "client-main"); - handler({ - runId: "run-main", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "hello from main" }, - }); + emitAgentEvent(handler, "run-main", "assistant", { text: "hello from main" }); const chatCalls = chatBroadcastCalls(broadcast); expect(chatCalls.length).toBeGreaterThanOrEqual(1); @@ -4886,19 +4141,16 @@ describe("agent event handler", () => { resolveSessionKeyForRun: () => "agent:main:main", }); - chatRunState.registry.add("run-no-lineage", { - sessionKey: "agent:main:main", - clientRunId: "client-no-lineage", - }); + registerChatRun(chatRunState, "run-no-lineage", "agent:main:main", "client-no-lineage"); for (let seq = 1; seq <= 5; seq++) { - handler({ - runId: "run-no-lineage", - seq, - stream: "assistant", - ts: Date.now() + seq * 200, - data: { text: `message ${seq}` }, - }); + emitAgentEvent( + handler, + "run-no-lineage", + "assistant", + { text: `message ${seq}` }, + { seq, ts: Date.now() + seq * 200 }, + ); } // The chat delta path invokes resolveSpawnedBy only. Non-subagent, @@ -4915,12 +4167,7 @@ describe("agent event handler", () => { }); it("includes spawnedBy in non-tool agent event broadcasts for subagent sessions", () => { - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:xyz", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-2", - }); + mockSessionLineage("agent:coder:subagent:xyz", "agent:conductor:task:parent-2"); const { broadcast, handler } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:xyz", @@ -4928,13 +4175,7 @@ describe("agent event handler", () => { registerAgentRunContext("run-agent-sub", { sessionKey: "agent:coder:subagent:xyz" }); - handler({ - runId: "run-agent-sub", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); + emitAgentEvent(handler, "run-agent-sub", "lifecycle", { phase: "start" }); const agentCalls = broadcast.mock.calls.filter(([event]) => event === "agent"); expect(agentCalls.length).toBeGreaterThanOrEqual(1); @@ -4945,38 +4186,19 @@ describe("agent event handler", () => { }); it("includes spawnedBy in chat error final broadcasts for subagent sessions", () => { - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:err", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-err", - }); + mockSessionLineage("agent:coder:subagent:err", "agent:conductor:task:parent-err"); const { broadcast, handler, chatRunState } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:err", lifecycleErrorRetryGraceMs: 0, }); - chatRunState.registry.add("run-sub-err", { - sessionKey: "agent:coder:subagent:err", - clientRunId: "client-sub-err", - }); + registerChatRun(chatRunState, "run-sub-err", "agent:coder:subagent:err", "client-sub-err"); - handler({ - runId: "run-sub-err", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "partial" }, - }); - - handler({ - runId: "run-sub-err", - seq: 2, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "error", error: "provider failed" }, - }); + emitAgentEvents(handler, "run-sub-err", [ + ["assistant", { text: "partial" }], + ["lifecycle", { phase: "error", error: "provider failed" }], + ]); const chatCalls = chatBroadcastCalls(broadcast); const errorCall = requireCall( @@ -4994,51 +4216,42 @@ describe("agent event handler", () => { let now = 20_000; const nowSpy = vi.spyOn(Date, "now").mockImplementation(() => now); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:flush", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-flush", - }); + mockSessionLineage("agent:coder:subagent:flush", "agent:conductor:task:parent-flush"); const { broadcast, chatRunState, toolEventRecipients, handler } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:flush", }); - chatRunState.registry.add("run-sub-flush", { - sessionKey: "agent:coder:subagent:flush", - clientRunId: "client-sub-flush", - }); + registerChatRun( + chatRunState, + "run-sub-flush", + "agent:coder:subagent:flush", + "client-sub-flush", + ); registerAgentRunContext("run-sub-flush", { sessionKey: "agent:coder:subagent:flush", verboseLevel: "off", }); toolEventRecipients.add("run-sub-flush", "conn-flush"); - handler({ - runId: "run-sub-flush", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "before tool" }, - }); + emitAgentEvent(handler, "run-sub-flush", "assistant", { text: "before tool" }); now = 20_050; - handler({ - runId: "run-sub-flush", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "before tool expanded" }, - }); + emitAgentEvent( + handler, + "run-sub-flush", + "assistant", + { text: "before tool expanded" }, + { seq: 2 }, + ); - handler({ - runId: "run-sub-flush", - seq: 3, - stream: "tool", - ts: Date.now(), - data: { phase: "start", name: "exec", toolCallId: "tool-flush-sub" }, - }); + emitAgentEvent( + handler, + "run-sub-flush", + "tool", + { phase: "start", name: "exec", toolCallId: "tool-flush-sub" }, + { seq: 3 }, + ); const chatCalls = chatBroadcastCalls(broadcast); const flushedDelta = requireCall( @@ -5056,12 +4269,7 @@ describe("agent event handler", () => { }); it("includes spawnedBy in seq gap error broadcasts for subagent sessions", () => { - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:gap", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-gap", - }); + mockSessionLineage("agent:coder:subagent:gap", "agent:conductor:task:parent-gap"); const { broadcast, handler } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:gap", @@ -5069,21 +4277,10 @@ describe("agent event handler", () => { registerAgentRunContext("run-sub-gap", { sessionKey: "agent:coder:subagent:gap" }); - handler({ - runId: "run-sub-gap", - seq: 1, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "start" }, - }); - - handler({ - runId: "run-sub-gap", - seq: 5, - stream: "assistant", - ts: Date.now(), - data: { text: "skipped seq" }, - }); + emitAgentEvents(handler, "run-sub-gap", [ + ["lifecycle", { phase: "start" }], + ["assistant", { text: "skipped seq" }, { seq: 5 }], + ]); const agentCalls = broadcast.mock.calls.filter(([event]) => event === "agent"); const gapError = requireCall( @@ -5099,44 +4296,19 @@ describe("agent event handler", () => { it("caches spawnedBy lookup so repeated events for the same subagent session only load the row once", () => { vi.mocked(loadGatewaySessionRow).mockClear(); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:cache-test", - kind: "direct", - updatedAt: null, - spawnedBy: "agent:conductor:task:parent-cache", - }); + mockSessionLineage("agent:coder:subagent:cache-test", "agent:conductor:task:parent-cache"); const { broadcast, handler, chatRunState } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:cache-test", }); - chatRunState.registry.add("run-cache", { - sessionKey: "agent:coder:subagent:cache-test", - clientRunId: "client-cache", - }); + registerChatRun(chatRunState, "run-cache", "agent:coder:subagent:cache-test", "client-cache"); - // Fire multiple events for the same session - handler({ - runId: "run-cache", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "chunk 1" }, - }); - handler({ - runId: "run-cache", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "chunk 2" }, - }); - handler({ - runId: "run-cache", - seq: 3, - stream: "lifecycle", - ts: Date.now(), - data: { phase: "end" }, - }); + emitAgentEvents(handler, "run-cache", [ + ["assistant", { text: "chunk 1" }], + ["assistant", { text: "chunk 2" }], + ["lifecycle", { phase: "end" }], + ]); // Key assertion: loadGatewaySessionRow called exactly once despite 3 events expect(loadGatewaySessionRow).toHaveBeenCalledTimes(1); @@ -5153,36 +4325,18 @@ describe("agent event handler", () => { it("caches null spawnedBy for eligible subagent sessions that lack a spawnedBy value", () => { vi.mocked(loadGatewaySessionRow).mockClear(); - vi.mocked(loadGatewaySessionRow).mockReturnValue({ - key: "agent:coder:subagent:no-lineage", - kind: "direct", - updatedAt: null, - // no spawnedBy field - }); + mockSessionLineage("agent:coder:subagent:no-lineage"); const { broadcast, handler, chatRunState } = createHarness({ resolveSessionKeyForRun: () => "agent:coder:subagent:no-lineage", }); - chatRunState.registry.add("run-null", { - sessionKey: "agent:coder:subagent:no-lineage", - clientRunId: "client-null", - }); + registerChatRun(chatRunState, "run-null", "agent:coder:subagent:no-lineage", "client-null"); - handler({ - runId: "run-null", - seq: 1, - stream: "assistant", - ts: Date.now(), - data: { text: "chunk 1" }, - }); - handler({ - runId: "run-null", - seq: 2, - stream: "assistant", - ts: Date.now(), - data: { text: "chunk 2" }, - }); + emitAgentEvents(handler, "run-null", [ + ["assistant", { text: "chunk 1" }], + ["assistant", { text: "chunk 2" }], + ]); // null result is cached — only one DB call despite two events expect(loadGatewaySessionRow).toHaveBeenCalledTimes(1);