diff options
| author | Adam Malczewski <[email protected]> | 2026-06-06 21:29:52 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-06 21:29:52 +0900 |
| commit | 2c5bc242a8a99e3b863c247f70b26f5883333677 (patch) | |
| tree | ce938fbef47d8dab44bfe5cce68ebfac7a91ca8a /packages/session-orchestrator/src | |
| parent | fedf9c2695476e9ee6f95776b0244acfc37f022f (diff) | |
| download | dispatch-2c5bc242a8a99e3b863c247f70b26f5883333677.tar.gz dispatch-2c5bc242a8a99e3b863c247f70b26f5883333677.zip | |
feat(kernel-runtime,session-orchestrator): emit turn lifecycle events
Close a gap found live: neither transport emitted turn-start/done/turn-sealed
(the wire defined them; nothing fired them). turn-sealed is the FE's
cache-commit signal (frontend-design §6.3); done ends the stream.
- kernel-runtime: runTurn emits turn-start first and done (with finishReason)
last, on every exit path (stop/tool-calls/max-steps/error/aborted).
- session-orchestrator: emits turn-sealed after conversationStore.append
succeeds (the kernel touches no DB, so the post-persist seal is the
orchestrator's). Not emitted if append throws.
No contract change (all three wire types already existed). Verified live: HTTP
/chat and WS chat both stream turn-start … done turn-sealed.
typecheck clean, 494 vitest + 80 bun, biome clean.
Diffstat (limited to 'packages/session-orchestrator/src')
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.test.ts | 117 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.ts | 2 |
2 files changed, 119 insertions, 0 deletions
diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts index 39c95d5..b648d42 100644 --- a/packages/session-orchestrator/src/orchestrator.test.ts +++ b/packages/session-orchestrator/src/orchestrator.test.ts @@ -364,3 +364,120 @@ describe("handleMessage model resolution", () => { expect(captured[1]?.cwd).toBeUndefined(); }); }); + +describe("turn-sealed event", () => { + it("emits turn-sealed after persisting the turn", async () => { + const store = createInMemoryStore(); + const provider = createFakeProvider([ + [ + { type: "text-delta", delta: "ok" }, + { type: "finish", reason: "stop" }, + ], + ]); + + const orchestrator = createSessionOrchestrator({ + conversationStore: store, + resolveProvider: () => provider, + resolveTools: () => [], + runTurn, + }); + + const { events, onEvent } = collectEvents(); + + await orchestrator.handleMessage({ + conversationId: "conv-seal", + text: "test", + onEvent, + }); + + const sealedEvents = events.filter((e) => e.type === "turn-sealed"); + expect(sealedEvents).toHaveLength(1); + const sealed = sealedEvents[0] as AgentEvent & { type: "turn-sealed" }; + expect(sealed.conversationId).toBe("conv-seal"); + expect(sealed.turnId).toMatch(/^turn-/); + }); + + it("turn-sealed is emitted after the store append", async () => { + const store = createInMemoryStore(); + const provider = createFakeProvider([ + [ + { type: "text-delta", delta: "ok" }, + { type: "finish", reason: "stop" }, + ], + ]); + + const ordering: string[] = []; + const wrappedStore: ConversationStore = { + async append(conversationId, messages) { + await store.append(conversationId, messages); + ordering.push("append"); + }, + async load(conversationId) { + return store.load(conversationId); + }, + async loadSince(conversationId, sinceSeq) { + return store.loadSince(conversationId, sinceSeq); + }, + }; + + const orchestrator = createSessionOrchestrator({ + conversationStore: wrappedStore, + resolveProvider: () => provider, + resolveTools: () => [], + runTurn, + }); + + await orchestrator.handleMessage({ + conversationId: "conv-order", + text: "test", + onEvent: (event) => { + if (event.type === "turn-sealed") { + ordering.push("turn-sealed"); + } + }, + }); + + expect(ordering).toEqual(["append", "turn-sealed"]); + }); + + it("does not emit turn-sealed when append throws", async () => { + const provider = createFakeProvider([ + [ + { type: "text-delta", delta: "ok" }, + { type: "finish", reason: "stop" }, + ], + ]); + + const failingStore: ConversationStore = { + async append() { + throw new Error("storage failure"); + }, + async load() { + return []; + }, + async loadSince() { + return []; + }, + }; + + const orchestrator = createSessionOrchestrator({ + conversationStore: failingStore, + resolveProvider: () => provider, + resolveTools: () => [], + runTurn, + }); + + const { events, onEvent } = collectEvents(); + + await expect( + orchestrator.handleMessage({ + conversationId: "conv-fail", + text: "test", + onEvent, + }), + ).rejects.toThrow("storage failure"); + + const sealedEvents = events.filter((e) => e.type === "turn-sealed"); + expect(sealedEvents).toHaveLength(0); + }); +}); diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts index a9ff4ff..311b620 100644 --- a/packages/session-orchestrator/src/orchestrator.ts +++ b/packages/session-orchestrator/src/orchestrator.ts @@ -92,6 +92,8 @@ export function createSessionOrchestrator(deps: SessionOrchestratorDeps): Sessio const toPersist: ChatMessage[] = [userMsg, ...result.messages]; await deps.conversationStore.append(conversationId, toPersist); + + onEvent({ type: "turn-sealed", conversationId, turnId }); }, }; } |
