summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-12 02:25:57 +0900
committerAdam Malczewski <[email protected]>2026-06-12 02:25:57 +0900
commit86b5137c4f7f2bcc08f0580f1edaa05d14015e63 (patch)
tree20215a5dccf1b76cf9cf95fceec1eca450aa8559 /packages/session-orchestrator/src
parent839f5c02676a0b7def27ace7125fbec3aa08bda5 (diff)
downloaddispatch-86b5137c4f7f2bcc08f0580f1edaa05d14015e63.tar.gz
dispatch-86b5137c4f7f2bcc08f0580f1edaa05d14015e63.zip
feat(turns): detached turns + multi-client live view
A turn no longer dies when its WebSocket connection closes. The turn-broadcast hub moves into the core (session-orchestrator): turns run detached, persist at seal regardless of clients, and fan out AgentEvents to N subscribers per conversation with in-flight buffer replay for late-joiners. transport-ws stops aborting turns on socket close and gains chat.subscribe/chat.unsubscribe so a second device (or a reloaded browser) can watch a running turn. - @dispatch/transport-contract 0.6.0->0.7.0: chat.subscribe/chat.unsubscribe WS ops - session-orchestrator: startTurn/subscribe/isActive; persistent subscribers + per-turn buffer (two-map model); handleMessage = convenience wrapper (no signal) - transport-ws: per-connection chat-subscription fan-out; no turn-abort-on-close - transport-http: test fakes updated for the widened interface (runtime unchanged) - design notes/turn-continuity-design.md; FE courier frontend-turn-continuity-handoff.md Live-verified vs flash (2-client WS): sender disconnect mid-turn -> other client streams to done + turn persists; late-join replays turn from turn-start. 891 vitest + transport bun green; tsc -b EXIT 0; biome clean.
Diffstat (limited to 'packages/session-orchestrator/src')
-rw-r--r--packages/session-orchestrator/src/index.ts3
-rw-r--r--packages/session-orchestrator/src/orchestrator.test.ts484
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts173
3 files changed, 591 insertions, 69 deletions
diff --git a/packages/session-orchestrator/src/index.ts b/packages/session-orchestrator/src/index.ts
index 2daf278..8bc99d2 100644
--- a/packages/session-orchestrator/src/index.ts
+++ b/packages/session-orchestrator/src/index.ts
@@ -6,7 +6,10 @@ export {
type SessionOrchestrator,
type SessionOrchestratorBundle,
type SessionOrchestratorDeps,
+ type StartTurnInput,
+ type StartTurnResult,
sessionOrchestratorHandle,
+ type TurnEventListener,
type TurnLifecyclePayload,
turnSettled,
turnStarted,
diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts
index b33bdcc..efa5d4e 100644
--- a/packages/session-orchestrator/src/orchestrator.test.ts
+++ b/packages/session-orchestrator/src/orchestrator.test.ts
@@ -208,39 +208,6 @@ describe("handleMessage integration", () => {
expect(lastUserText).toEqual({ type: "text", text: "Second message" });
});
- it("passes abort signal through to runTurn", async () => {
- const store = createInMemoryStore();
- const ac = new AbortController();
- ac.abort();
-
- const provider = createFakeProvider([
- [
- { type: "text-delta", delta: "should not appear" },
- { type: "finish", reason: "stop" },
- ],
- ]);
-
- const { orchestrator } = createSessionOrchestrator({
- conversationStore: store,
- resolveProvider: () => provider,
- resolveTools: () => [],
- applyToolsFilter: identityApplyToolsFilter,
- runTurn,
- });
-
- await orchestrator.handleMessage({
- conversationId: "conv-abort",
- text: "test",
- onEvent: () => {},
- signal: ac.signal,
- });
-
- const stored = store.data.get("conv-abort");
- expect(stored).toBeDefined();
- expect(stored).toHaveLength(1);
- expect(stored?.[0]?.role).toBe("user");
- });
-
it("uses custom dispatch policy when resolveDispatch is provided", async () => {
const store = createInMemoryStore();
const provider = createFakeProvider([
@@ -558,7 +525,7 @@ describe("turn-sealed event", () => {
expect(ordering).toEqual(["append", "appendMetrics", "turn-sealed"]);
});
- it("does not emit turn-sealed when append throws", async () => {
+ it("does not emit turn-sealed when append throws — emits error event instead", async () => {
const provider = createFakeProvider([
[
{ type: "text-delta", delta: "ok" },
@@ -598,16 +565,18 @@ describe("turn-sealed event", () => {
const { events, onEvent } = collectEvents();
- await expect(
- orchestrator.handleMessage({
- conversationId: "conv-fail",
- text: "test",
- onEvent,
- }),
- ).rejects.toThrow("storage failure");
+ await orchestrator.handleMessage({
+ conversationId: "conv-fail",
+ text: "test",
+ onEvent,
+ });
const sealedEvents = events.filter((e) => e.type === "turn-sealed");
expect(sealedEvents).toHaveLength(0);
+
+ const errorEvents = events.filter((e) => e.type === "error");
+ expect(errorEvents).toHaveLength(1);
+ expect((errorEvents[0] as AgentEvent & { type: "error" }).message).toBe("storage failure");
});
});
@@ -936,17 +905,18 @@ describe("turn metrics persistence", () => {
const { events, onEvent } = collectEvents();
- await expect(
- orchestrator.handleMessage({
- conversationId: "conv-fail-metrics",
- text: "test",
- onEvent,
- }),
- ).rejects.toThrow("storage failure");
+ await orchestrator.handleMessage({
+ conversationId: "conv-fail-metrics",
+ text: "test",
+ onEvent,
+ });
const sealedEvents = events.filter((e) => e.type === "turn-sealed");
expect(sealedEvents).toHaveLength(0);
expect(metricsAppended).toBe(false);
+
+ const errorEvents = events.filter((e) => e.type === "error");
+ expect(errorEvents).toHaveLength(1);
});
});
@@ -1628,3 +1598,421 @@ describe("cwd persistence", () => {
expect(captured[0]?.cwd).toBeUndefined();
});
});
+
+describe("detached turn hub", () => {
+ function waitForEvent(
+ orchestrator: ReturnType<typeof createSessionOrchestrator>["orchestrator"],
+ conversationId: string,
+ eventType: string,
+ ): Promise<AgentEvent> {
+ return new Promise((resolve) => {
+ const unsub = orchestrator.subscribe(conversationId, (event) => {
+ if (event.type === eventType) {
+ unsub();
+ resolve(event);
+ }
+ });
+ });
+ }
+
+ it("subscribe-BEFORE-startTurn delivers — listener receives full ordered event sequence", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Hello" },
+ { type: "text-delta", delta: " world" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ const events: AgentEvent[] = [];
+ const unsub = orchestrator.subscribe("conv-pre-sub", (e) => events.push(e));
+
+ orchestrator.startTurn({ conversationId: "conv-pre-sub", text: "Hi" });
+
+ const sealed = waitForEvent(orchestrator, "conv-pre-sub", "turn-sealed");
+ await sealed;
+
+ unsub();
+
+ expect(events.length).toBeGreaterThan(0);
+ const types = events.map((e) => e.type);
+ expect(types[0]).toBe("turn-start");
+ expect(types).toContain("text-delta");
+ expect(types[types.length - 1]).toBe("turn-sealed");
+
+ const textDeltas = events.filter((e) => e.type === "text-delta");
+ expect(textDeltas).toHaveLength(2);
+ });
+
+ it("multi-subscriber fan-out (subscribed before start) — two listeners receive identical ordered events", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Hello" },
+ { type: "text-delta", delta: " world" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ const eventsA: AgentEvent[] = [];
+ const eventsB: AgentEvent[] = [];
+ const unsubA = orchestrator.subscribe("conv-fanout-pre", (e) => eventsA.push(e));
+ const unsubB = orchestrator.subscribe("conv-fanout-pre", (e) => eventsB.push(e));
+
+ orchestrator.startTurn({ conversationId: "conv-fanout-pre", text: "Hi" });
+
+ const sealed = waitForEvent(orchestrator, "conv-fanout-pre", "turn-sealed");
+ await sealed;
+
+ unsubA();
+ unsubB();
+
+ expect(eventsA.length).toBeGreaterThan(0);
+ expect(eventsA).toEqual(eventsB);
+
+ const types = eventsA.map((e) => e.type);
+ expect(types[0]).toBe("turn-start");
+ expect(types[types.length - 1]).toBe("turn-sealed");
+ });
+
+ it("late-join replay — subscriber added mid-turn receives buffered events then live events, no gap/dup", async () => {
+ const store = createInMemoryStore();
+ let emitBarrierResolve: (() => void) | undefined;
+ const emitBarrier = new Promise<void>((resolve) => {
+ emitBarrierResolve = resolve;
+ });
+
+ let callIndex = 0;
+ const provider: ProviderContract = {
+ id: "fake",
+ stream() {
+ const idx = callIndex++;
+ return (async function* () {
+ if (idx === 0) {
+ yield { type: "text-delta", delta: "Hello" } as ProviderEvent;
+ yield { type: "text-delta", delta: " world" } as ProviderEvent;
+ await emitBarrier;
+ yield { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } } as ProviderEvent;
+ yield { type: "finish", reason: "stop" } as ProviderEvent;
+ }
+ })();
+ },
+ };
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ orchestrator.startTurn({ conversationId: "conv-latejoin", text: "Hi" });
+
+ const earlyEvents: AgentEvent[] = [];
+ const unsubEarly = orchestrator.subscribe("conv-latejoin", (e) => earlyEvents.push(e));
+
+ await new Promise<void>((resolve) => setTimeout(resolve, 10));
+
+ const lateEvents: AgentEvent[] = [];
+ const unsubLate = orchestrator.subscribe("conv-latejoin", (e) => lateEvents.push(e));
+
+ const earlySnapshot = [...earlyEvents];
+ expect(earlySnapshot.length).toBeGreaterThanOrEqual(2);
+ expect(earlySnapshot.some((e) => e.type === "turn-start")).toBe(true);
+ expect(earlySnapshot.some((e) => e.type === "text-delta")).toBe(true);
+
+ expect(lateEvents.length).toBe(earlySnapshot.length);
+ expect(lateEvents).toEqual(earlySnapshot);
+
+ emitBarrierResolve?.();
+
+ const sealed = waitForEvent(orchestrator, "conv-latejoin", "turn-sealed");
+ await sealed;
+
+ unsubEarly();
+ unsubLate();
+
+ expect(earlyEvents.length).toBeGreaterThan(earlySnapshot.length);
+ expect(lateEvents.length).toBe(earlyEvents.length);
+ expect(lateEvents).toEqual(earlyEvents);
+ });
+
+ it("subscriber persists across turns — one subscriber receives events from two sequential turns", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Turn1" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ [
+ { type: "text-delta", delta: "Turn2" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ const allEvents: AgentEvent[] = [];
+ const unsub = orchestrator.subscribe("conv-persist", (e) => allEvents.push(e));
+
+ // First turn
+ orchestrator.startTurn({ conversationId: "conv-persist", text: "First" });
+ const sealed1 = waitForEvent(orchestrator, "conv-persist", "turn-sealed");
+ await sealed1;
+
+ // Second turn
+ orchestrator.startTurn({ conversationId: "conv-persist", text: "Second" });
+ const sealed2 = waitForEvent(orchestrator, "conv-persist", "turn-sealed");
+ await sealed2;
+
+ unsub();
+
+ const turnStarts = allEvents.filter((e) => e.type === "turn-start");
+ expect(turnStarts).toHaveLength(2);
+
+ const turnSealeds = allEvents.filter((e) => e.type === "turn-sealed");
+ expect(turnSealeds).toHaveLength(2);
+
+ const textDeltas = allEvents.filter((e) => e.type === "text-delta");
+ expect(textDeltas).toHaveLength(2);
+ expect((textDeltas[0] as AgentEvent & { type: "text-delta" }).delta).toBe("Turn1");
+ expect((textDeltas[1] as AgentEvent & { type: "text-delta" }).delta).toBe("Turn2");
+ });
+
+ it("detached completion — turn runs to completion with zero subscribers and persists", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Hello" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ const result = orchestrator.startTurn({ conversationId: "conv-detached", text: "Hi" });
+ expect(result.started).toBe(true);
+ const sealed = waitForEvent(orchestrator, "conv-detached", "turn-sealed");
+
+ await sealed;
+
+ const stored = store.data.get("conv-detached");
+ expect(stored).toBeDefined();
+ expect(stored).toHaveLength(2);
+ expect(stored?.[0]?.role).toBe("user");
+ expect(stored?.[1]?.role).toBe("assistant");
+ });
+
+ it("single-flight reject — startTurn while active returns already-active, no second turn", async () => {
+ const store = createInMemoryStore();
+ let resolveRunTurn: (() => void) | undefined;
+ const runTurnBlocker = new Promise<void>((resolve) => {
+ resolveRunTurn = resolve;
+ });
+
+ const provider: ProviderContract = {
+ id: "fake",
+ stream: async function* () {
+ yield { type: "text-delta", delta: "slow" } as ProviderEvent;
+ yield { type: "finish", reason: "stop" } as ProviderEvent;
+ },
+ };
+
+ const blockingRunTurn = async (_input: RunTurnInput): Promise<RunTurnResult> => {
+ await runTurnBlocker;
+ return {
+ messages: [{ role: "assistant", chunks: [{ type: "text", text: "done" }] }],
+ usage: { inputTokens: 1, outputTokens: 1 },
+ finishReason: "stop",
+ };
+ };
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: blockingRunTurn,
+ });
+
+ const first = orchestrator.startTurn({ conversationId: "conv-singleflight", text: "first" });
+ expect(first.started).toBe(true);
+ const sealed = waitForEvent(orchestrator, "conv-singleflight", "turn-sealed");
+
+ const second = orchestrator.startTurn({ conversationId: "conv-singleflight", text: "second" });
+ expect(second.started).toBe(false);
+ if (!second.started) {
+ expect(second.reason).toBe("already-active");
+ }
+
+ resolveRunTurn?.();
+ await sealed;
+
+ expect(store.data.get("conv-singleflight")?.length).toBe(2);
+ });
+
+ it("isActive false after seal — subscribe replays nothing after turn-sealed", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "ok" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ expect(orchestrator.isActive("conv-cleared")).toBe(false);
+
+ orchestrator.startTurn({ conversationId: "conv-cleared", text: "test" });
+ const sealed = waitForEvent(orchestrator, "conv-cleared", "turn-sealed");
+
+ await sealed;
+
+ expect(orchestrator.isActive("conv-cleared")).toBe(false);
+
+ const lateEvents: AgentEvent[] = [];
+ const unsub = orchestrator.subscribe("conv-cleared", (e) => lateEvents.push(e));
+ unsub();
+
+ expect(lateEvents).toHaveLength(0);
+ });
+
+ it("handleMessage convenience — drives turn end-to-end via onEvent and resolves on seal", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Hello" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn,
+ });
+
+ const events: AgentEvent[] = [];
+ await orchestrator.handleMessage({
+ conversationId: "conv-hm",
+ text: "Hi",
+ onEvent: (e) => events.push(e),
+ });
+
+ const types = events.map((e) => e.type);
+ expect(types[0]).toBe("turn-start");
+ expect(types).toContain("text-delta");
+ expect(types[types.length - 1]).toBe("turn-sealed");
+
+ const stored = store.data.get("conv-hm");
+ expect(stored).toHaveLength(2);
+ });
+
+ it("handleMessage already-active emits error event and resolves without hanging", async () => {
+ const store = createInMemoryStore();
+ let resolveRunTurn: (() => void) | undefined;
+ const runTurnBlocker = new Promise<void>((resolve) => {
+ resolveRunTurn = resolve;
+ });
+
+ const provider: ProviderContract = {
+ id: "fake",
+ stream: async function* () {
+ yield { type: "text-delta", delta: "slow" } as ProviderEvent;
+ yield { type: "finish", reason: "stop" } as ProviderEvent;
+ },
+ };
+
+ const blockingRunTurn = async (_input: RunTurnInput): Promise<RunTurnResult> => {
+ await runTurnBlocker;
+ return {
+ messages: [{ role: "assistant", chunks: [{ type: "text", text: "done" }] }],
+ usage: { inputTokens: 1, outputTokens: 1 },
+ finishReason: "stop",
+ };
+ };
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: blockingRunTurn,
+ });
+
+ const firstEvents: AgentEvent[] = [];
+ const firstPromise = orchestrator.handleMessage({
+ conversationId: "conv-hm-active",
+ text: "first",
+ onEvent: (e) => firstEvents.push(e),
+ });
+
+ await new Promise<void>((resolve) => setTimeout(resolve, 10));
+
+ const secondEvents: AgentEvent[] = [];
+ const secondPromise = orchestrator.handleMessage({
+ conversationId: "conv-hm-active",
+ text: "second",
+ onEvent: (e) => secondEvents.push(e),
+ });
+
+ await secondPromise;
+
+ expect(secondEvents).toHaveLength(1);
+ expect(secondEvents[0]?.type).toBe("error");
+ expect((secondEvents[0] as AgentEvent & { type: "error" }).message).toBe(
+ "turn already active for this conversation",
+ );
+
+ resolveRunTurn?.();
+ await firstPromise;
+
+ expect(firstEvents.some((e) => e.type === "turn-sealed")).toBe(true);
+ });
+});
diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts
index e86729c..0421108 100644
--- a/packages/session-orchestrator/src/orchestrator.ts
+++ b/packages/session-orchestrator/src/orchestrator.ts
@@ -18,6 +18,26 @@ import { createMetricsAccumulator } from "./metrics.js";
import { buildUserMessage, defaultDispatchPolicy, generateTurnId } from "./pure.js";
import type { ToolAssembly } from "./tools-filter.js";
+// --- Broadcast hub types ---
+
+export interface StartTurnInput {
+ readonly conversationId: string;
+ readonly text: string;
+ readonly modelName?: string;
+ readonly cwd?: string;
+}
+
+export type StartTurnResult =
+ | { readonly started: true; readonly turnId: string }
+ | { readonly started: false; readonly reason: "already-active" };
+
+export type TurnEventListener = (event: AgentEvent) => void;
+
+interface ActiveTurn {
+ buffer: AgentEvent[];
+ turnId: string;
+}
+
// --- Lifecycle event hooks ---
/** Context carried on turn-lifecycle events, enough to replicate the turn's request prefix. */
@@ -66,11 +86,13 @@ export const cacheWarmHandle: ServiceHandle<WarmService> = defineService<WarmSer
);
export interface SessionOrchestrator {
+ startTurn(input: StartTurnInput): StartTurnResult;
+ subscribe(conversationId: string, listener: TurnEventListener): () => void;
+ isActive(conversationId: string): boolean;
handleMessage(input: {
conversationId: string;
text: string;
onEvent: (event: AgentEvent) => void;
- signal?: AbortSignal;
modelName?: string;
cwd?: string;
}): Promise<void>;
@@ -114,31 +136,57 @@ export function createSessionOrchestrator(
deps: SessionOrchestratorDeps,
): SessionOrchestratorBundle {
const activeConversations = new Set<string>();
+ const subscribers = new Map<string, Set<TurnEventListener>>();
+ const activeTurns = new Map<string, ActiveTurn>();
+
+ function emitToHub(conversationId: string, event: AgentEvent): void {
+ const turn = activeTurns.get(conversationId);
+ if (turn !== undefined) {
+ turn.buffer.push(event);
+ }
+ const listeners = subscribers.get(conversationId);
+ if (listeners !== undefined) {
+ for (const listener of listeners) {
+ listener(event);
+ }
+ }
+ }
- const orchestrator: SessionOrchestrator = {
- async handleMessage({ conversationId, text, onEvent, signal, modelName, cwd }) {
- activeConversations.add(conversationId);
-
- const effectiveCwd =
- cwd !== undefined
- ? cwd
- : ((await deps.conversationStore.getCwd(conversationId)) ?? undefined);
-
- const payload: TurnLifecyclePayload = {
- conversationId,
- ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
- ...(modelName !== undefined ? { modelName } : {}),
- };
+ function runTurnDetached(
+ conversationId: string,
+ text: string,
+ modelName: string | undefined,
+ cwd: string | undefined,
+ ): void {
+ const turnId = generateTurnId();
+ activeTurns.set(conversationId, { buffer: [], turnId });
+ activeConversations.add(conversationId);
+
+ const effectiveCwdPromise =
+ cwd !== undefined
+ ? Promise.resolve(cwd)
+ : deps.conversationStore.getCwd(conversationId).then((c) => c ?? undefined);
+
+ const payloadPromise = effectiveCwdPromise.then((effectiveCwd) => ({
+ conversationId,
+ ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
+ ...(modelName !== undefined ? { modelName } : {}),
+ }));
+
+ payloadPromise.then((payload) => {
deps.emit?.(turnStarted, payload);
+ });
+ void (async () => {
try {
+ const effectiveCwd = await effectiveCwdPromise;
+
if (cwd !== undefined) {
await deps.conversationStore.setCwd(conversationId, cwd);
}
const history = await deps.conversationStore.load(conversationId);
const userMsg = buildUserMessage(text);
- const turnId = generateTurnId();
let provider: ProviderContract;
let modelOverride: string | undefined;
@@ -146,7 +194,7 @@ export function createSessionOrchestrator(
if (modelName !== undefined && deps.resolveModel !== undefined) {
const resolved = deps.resolveModel(modelName);
if (resolved === undefined) {
- onEvent({
+ emitToHub(conversationId, {
type: "error",
conversationId,
turnId,
@@ -172,7 +220,7 @@ export function createSessionOrchestrator(
const emitAndAccumulate = (event: AgentEvent): void => {
metrics.ingest(event);
- onEvent(event);
+ emitToHub(conversationId, event);
};
const opts: RunTurnInput = {
@@ -187,7 +235,6 @@ export function createSessionOrchestrator(
? { providerOpts: { model: modelOverride } satisfies ProviderStreamOptions }
: {}),
...(turnLogger !== undefined ? { logger: turnLogger } : {}),
- ...(signal !== undefined ? { signal } : {}),
...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
...(deps.now !== undefined ? { now: deps.now } : {}),
};
@@ -200,11 +247,95 @@ export function createSessionOrchestrator(
const turnMetrics = metrics.build(turnId);
await deps.conversationStore.appendMetrics(conversationId, turnMetrics);
- onEvent({ type: "turn-sealed", conversationId, turnId });
+ emitToHub(conversationId, { type: "turn-sealed", conversationId, turnId });
+ } catch (err) {
+ const message = err instanceof Error ? err.message : String(err);
+ emitToHub(conversationId, {
+ type: "error",
+ conversationId,
+ turnId,
+ message,
+ });
} finally {
+ activeTurns.delete(conversationId);
activeConversations.delete(conversationId);
- deps.emit?.(turnSettled, payload);
+ void payloadPromise.then((payload) => {
+ deps.emit?.(turnSettled, payload);
+ });
+ }
+ })();
+ }
+
+ const orchestrator: SessionOrchestrator = {
+ startTurn({ conversationId, text, modelName, cwd }) {
+ if (activeTurns.has(conversationId)) {
+ return { started: false, reason: "already-active" };
+ }
+ runTurnDetached(conversationId, text, modelName, cwd);
+ const turn = activeTurns.get(conversationId);
+ const turnId = turn !== undefined ? turn.turnId : "";
+ return { started: true, turnId };
+ },
+
+ subscribe(conversationId, listener) {
+ let listeners = subscribers.get(conversationId);
+ if (listeners === undefined) {
+ listeners = new Set();
+ subscribers.set(conversationId, listeners);
+ }
+ const turn = activeTurns.get(conversationId);
+ if (turn !== undefined) {
+ const snapshot = [...turn.buffer];
+ listeners.add(listener);
+ for (const event of snapshot) {
+ listener(event);
+ }
+ } else {
+ listeners.add(listener);
+ }
+ return () => {
+ const set = subscribers.get(conversationId);
+ if (set !== undefined) {
+ set.delete(listener);
+ if (set.size === 0) {
+ subscribers.delete(conversationId);
+ }
+ }
+ };
+ },
+
+ isActive(conversationId) {
+ return activeTurns.has(conversationId);
+ },
+
+ async handleMessage({ conversationId, text, onEvent, modelName, cwd }) {
+ const turnInput: StartTurnInput = {
+ conversationId,
+ text,
+ ...(modelName !== undefined ? { modelName } : {}),
+ ...(cwd !== undefined ? { cwd } : {}),
+ };
+ const result = orchestrator.startTurn(turnInput);
+ if (!result.started) {
+ const errorTurnId = generateTurnId();
+ onEvent({
+ type: "error",
+ conversationId,
+ turnId: errorTurnId,
+ message: "turn already active for this conversation",
+ });
+ return;
}
+
+ await new Promise<void>((resolve) => {
+ const unsubscribe = orchestrator.subscribe(conversationId, (event) => {
+ onEvent(event);
+ if (event.type === "turn-sealed" || event.type === "error") {
+ unsubscribe();
+ resolve();
+ }
+ });
+ });
},
};