summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-06 21:29:52 +0900
committerAdam Malczewski <[email protected]>2026-06-06 21:29:52 +0900
commit2c5bc242a8a99e3b863c247f70b26f5883333677 (patch)
treece938fbef47d8dab44bfe5cce68ebfac7a91ca8a /packages/session-orchestrator/src
parentfedf9c2695476e9ee6f95776b0244acfc37f022f (diff)
downloaddispatch-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.ts117
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts2
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 });
},
};
}