summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-04 23:50:34 +0900
committerAdam Malczewski <[email protected]>2026-06-04 23:50:34 +0900
commit357ad3567480b5483220e2cea266a4f1417d174d (patch)
treef9eba8ebad94d7396546e9e953631109cac43f84 /packages/session-orchestrator/src
parent3390f5ed73674ba12f08ee801869ffa2d5b9b38d (diff)
downloaddispatch-357ad3567480b5483220e2cea266a4f1417d174d.tar.gz
dispatch-357ad3567480b5483220e2cea266a4f1417d174d.zip
feat(core-ext): session-orchestrator + transport-http (parallel); wire into build graph (164 tests)
Diffstat (limited to 'packages/session-orchestrator/src')
-rw-r--r--packages/session-orchestrator/src/extension.ts64
-rw-r--r--packages/session-orchestrator/src/index.ts13
-rw-r--r--packages/session-orchestrator/src/orchestrator.test.ts200
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts60
-rw-r--r--packages/session-orchestrator/src/pure.test.ts69
-rw-r--r--packages/session-orchestrator/src/pure.ts27
6 files changed, 433 insertions, 0 deletions
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts
new file mode 100644
index 0000000..c75bf29
--- /dev/null
+++ b/packages/session-orchestrator/src/extension.ts
@@ -0,0 +1,64 @@
+import { conversationStoreHandle } from "@dispatch/conversation-store";
+import type {
+ Extension,
+ HostAPI,
+ Manifest,
+ ProviderContract,
+ ToolContract,
+} from "@dispatch/kernel";
+import { runTurn } from "@dispatch/kernel";
+import {
+ createSessionOrchestrator,
+ type SessionOrchestrator,
+ sessionOrchestratorHandle,
+} from "./orchestrator.js";
+import { selectFirstProvider } from "./pure.js";
+
+export const manifest: Manifest = {
+ id: "session-orchestrator",
+ name: "Session Orchestrator",
+ version: "0.0.0",
+ apiVersion: "^0.1.0",
+ trust: "bundled",
+ dependsOn: ["conversation-store"],
+ activation: "eager",
+ contributes: {
+ services: ["session-orchestrator/orchestrator"],
+ },
+};
+
+interface ProviderResolvingHostAPI extends HostAPI {
+ readonly getProviders?: () => ReadonlyMap<string, ProviderContract>;
+ readonly getTools?: () => ReadonlyMap<string, ToolContract>;
+}
+
+export function activate(host: HostAPI): void {
+ const conversationStore = host.getService(conversationStoreHandle);
+ const extendedHost = host as ProviderResolvingHostAPI;
+
+ const orchestrator: SessionOrchestrator = createSessionOrchestrator({
+ conversationStore,
+ resolveProvider: () => {
+ if (extendedHost.getProviders !== undefined) {
+ return selectFirstProvider(extendedHost.getProviders());
+ }
+ throw new Error(
+ "HostAPI does not expose getProviders() — change-request: add provider resolution to HostAPI",
+ );
+ },
+ resolveTools: () => {
+ if (extendedHost.getTools !== undefined) {
+ return [...extendedHost.getTools().values()];
+ }
+ return [];
+ },
+ runTurn,
+ });
+
+ host.provideService(sessionOrchestratorHandle, orchestrator);
+}
+
+export const extension: Extension = {
+ manifest,
+ activate,
+};
diff --git a/packages/session-orchestrator/src/index.ts b/packages/session-orchestrator/src/index.ts
new file mode 100644
index 0000000..3270c46
--- /dev/null
+++ b/packages/session-orchestrator/src/index.ts
@@ -0,0 +1,13 @@
+export { extension, manifest } from "./extension.js";
+export {
+ createSessionOrchestrator,
+ type SessionOrchestrator,
+ type SessionOrchestratorDeps,
+ sessionOrchestratorHandle,
+} from "./orchestrator.js";
+export {
+ buildUserMessage,
+ defaultDispatchPolicy,
+ generateTurnId,
+ selectFirstProvider,
+} from "./pure.js";
diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts
new file mode 100644
index 0000000..0d908b2
--- /dev/null
+++ b/packages/session-orchestrator/src/orchestrator.test.ts
@@ -0,0 +1,200 @@
+import type { ConversationStore } from "@dispatch/conversation-store";
+import type { AgentEvent, ChatMessage, ProviderContract, ProviderEvent } from "@dispatch/kernel";
+import { runTurn } from "@dispatch/kernel";
+import { describe, expect, it } from "vitest";
+import { createSessionOrchestrator } from "./orchestrator.js";
+
+function createInMemoryStore(): ConversationStore & {
+ readonly data: Map<string, ChatMessage[]>;
+} {
+ const data = new Map<string, ChatMessage[]>();
+ return {
+ data,
+ async append(conversationId, messages) {
+ const existing = data.get(conversationId) ?? [];
+ data.set(conversationId, [...existing, ...messages]);
+ },
+ async load(conversationId) {
+ return [...(data.get(conversationId) ?? [])];
+ },
+ };
+}
+
+function createFakeProvider(script: ProviderEvent[][]): ProviderContract {
+ let callIndex = 0;
+ return {
+ id: "fake",
+ stream(_messages, _tools) {
+ const events = script[callIndex] ?? [];
+ callIndex++;
+ return (async function* () {
+ for (const event of events) {
+ yield event;
+ }
+ })();
+ },
+ };
+}
+
+function collectEvents(): { events: AgentEvent[]; onEvent: (event: AgentEvent) => void } {
+ const events: AgentEvent[] = [];
+ return { events, onEvent: (event) => events.push(event) };
+}
+
+describe("handleMessage integration", () => {
+ it("loads history, runs turn, emits events, and persists result", async () => {
+ const store = createInMemoryStore();
+ const provider = createFakeProvider([
+ [
+ { type: "text-delta", delta: "Hello" },
+ { type: "text-delta", delta: " there" },
+ { type: "usage", usage: { inputTokens: 5, outputTokens: 3 } },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const orchestrator = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ runTurn,
+ });
+
+ const { events, onEvent } = collectEvents();
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-1",
+ text: "Hi",
+ onEvent,
+ });
+
+ expect(events.length).toBeGreaterThan(0);
+ const textDeltas = events.filter((e) => e.type === "text-delta");
+ expect(textDeltas).toHaveLength(2);
+
+ const stored = store.data.get("conv-1");
+ expect(stored).toBeDefined();
+ expect(stored).toHaveLength(2);
+ expect(stored?.[0]?.role).toBe("user");
+ expect(stored?.[1]?.role).toBe("assistant");
+
+ const userChunks = stored?.[0]?.chunks ?? [];
+ expect(userChunks[0]).toEqual({ type: "text", text: "Hi" });
+
+ const assistantChunks = stored?.[1]?.chunks ?? [];
+ expect(assistantChunks.some((c) => c.type === "text")).toBe(true);
+ });
+
+ it("multi-turn: second call sees first turn in history", async () => {
+ const store = createInMemoryStore();
+ let capturedMessages: ChatMessage[] | undefined;
+
+ let callCount = 0;
+ const provider: ProviderContract = {
+ id: "fake",
+ stream(messages, _tools) {
+ if (callCount === 1) {
+ capturedMessages = [...messages];
+ }
+ callCount++;
+ return (async function* () {
+ yield { type: "text-delta", delta: `Reply ${callCount}` } as ProviderEvent;
+ yield { type: "finish", reason: "stop" } as ProviderEvent;
+ })();
+ },
+ };
+
+ const orchestrator = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ runTurn,
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-multi",
+ text: "First message",
+ onEvent: () => {},
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-multi",
+ text: "Second message",
+ onEvent: () => {},
+ });
+
+ expect(capturedMessages).toBeDefined();
+ expect(capturedMessages?.length).toBeGreaterThanOrEqual(3);
+
+ expect(capturedMessages?.[0]?.role).toBe("user");
+ const firstUserText = capturedMessages?.[0]?.chunks[0];
+ expect(firstUserText).toEqual({ type: "text", text: "First message" });
+
+ expect(capturedMessages?.[1]?.role).toBe("assistant");
+
+ const lastUser = capturedMessages?.findLast((m) => m.role === "user");
+ expect(lastUser).toBeDefined();
+ const lastUserText = lastUser?.chunks[0];
+ 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: () => [],
+ 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([
+ [
+ { type: "text-delta", delta: "ok" },
+ { type: "finish", reason: "stop" },
+ ],
+ ]);
+
+ const orchestrator = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ resolveDispatch: () => ({ maxConcurrent: 4, eager: false }),
+ runTurn,
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-dispatch",
+ text: "test",
+ onEvent: () => {},
+ });
+
+ const stored = store.data.get("conv-dispatch");
+ expect(stored).toBeDefined();
+ expect(stored?.length).toBeGreaterThanOrEqual(1);
+ });
+});
diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts
new file mode 100644
index 0000000..f209d2d
--- /dev/null
+++ b/packages/session-orchestrator/src/orchestrator.ts
@@ -0,0 +1,60 @@
+import type { ConversationStore } from "@dispatch/conversation-store";
+import type {
+ AgentEvent,
+ ChatMessage,
+ ProviderContract,
+ RunTurnInput,
+ RunTurnResult,
+ ToolContract,
+ ToolDispatchPolicy,
+} from "@dispatch/kernel";
+import { defineService } from "@dispatch/kernel";
+import { buildUserMessage, defaultDispatchPolicy, generateTurnId } from "./pure.js";
+
+export interface SessionOrchestrator {
+ handleMessage(input: {
+ conversationId: string;
+ text: string;
+ onEvent: (event: AgentEvent) => void;
+ signal?: AbortSignal;
+ }): Promise<void>;
+}
+
+export const sessionOrchestratorHandle = defineService<SessionOrchestrator>(
+ "session-orchestrator/orchestrator",
+);
+
+export interface SessionOrchestratorDeps {
+ readonly conversationStore: ConversationStore;
+ readonly resolveProvider: () => ProviderContract;
+ readonly resolveTools: () => readonly ToolContract[];
+ readonly resolveDispatch?: () => ToolDispatchPolicy;
+ readonly runTurn: (input: RunTurnInput) => Promise<RunTurnResult>;
+}
+
+export function createSessionOrchestrator(deps: SessionOrchestratorDeps): SessionOrchestrator {
+ return {
+ async handleMessage({ conversationId, text, onEvent, signal }) {
+ const history = await deps.conversationStore.load(conversationId);
+ const userMsg = buildUserMessage(text);
+ const provider = deps.resolveProvider();
+ const tools = deps.resolveTools();
+ const dispatch = deps.resolveDispatch?.() ?? defaultDispatchPolicy();
+ const turnId = generateTurnId();
+
+ const result = await deps.runTurn({
+ provider,
+ messages: [...history, userMsg],
+ tools,
+ dispatch,
+ emit: onEvent,
+ tabId: conversationId,
+ turnId,
+ ...(signal !== undefined ? { signal } : {}),
+ });
+
+ const toPersist: ChatMessage[] = [userMsg, ...result.messages];
+ await deps.conversationStore.append(conversationId, toPersist);
+ },
+ };
+}
diff --git a/packages/session-orchestrator/src/pure.test.ts b/packages/session-orchestrator/src/pure.test.ts
new file mode 100644
index 0000000..e233fca
--- /dev/null
+++ b/packages/session-orchestrator/src/pure.test.ts
@@ -0,0 +1,69 @@
+import type { ProviderContract } from "@dispatch/kernel";
+import { describe, expect, it } from "vitest";
+import {
+ buildUserMessage,
+ defaultDispatchPolicy,
+ generateTurnId,
+ selectFirstProvider,
+} from "./pure.js";
+
+describe("buildUserMessage", () => {
+ it("creates a user message with a single text chunk", () => {
+ const msg = buildUserMessage("hello world");
+ expect(msg.role).toBe("user");
+ expect(msg.chunks).toHaveLength(1);
+ expect(msg.chunks[0]).toEqual({ type: "text", text: "hello world" });
+ });
+
+ it("preserves empty text", () => {
+ const msg = buildUserMessage("");
+ expect(msg.role).toBe("user");
+ expect(msg.chunks[0]).toEqual({ type: "text", text: "" });
+ });
+});
+
+describe("selectFirstProvider", () => {
+ it("returns the first provider from a non-empty map", () => {
+ const provider: ProviderContract = {
+ id: "test-provider",
+ stream: async function* () {},
+ };
+ const providers = new Map<string, ProviderContract>();
+ providers.set("test-provider", provider);
+
+ expect(selectFirstProvider(providers)).toBe(provider);
+ });
+
+ it("throws when the map is empty", () => {
+ const providers = new Map<string, ProviderContract>();
+ expect(() => selectFirstProvider(providers)).toThrow("No providers registered");
+ });
+
+ it("returns the first inserted provider when multiple exist", () => {
+ const first: ProviderContract = { id: "first", stream: async function* () {} };
+ const second: ProviderContract = { id: "second", stream: async function* () {} };
+ const providers = new Map<string, ProviderContract>();
+ providers.set("first", first);
+ providers.set("second", second);
+
+ expect(selectFirstProvider(providers).id).toBe("first");
+ });
+});
+
+describe("defaultDispatchPolicy", () => {
+ it("returns maxConcurrent: 1, eager: true", () => {
+ expect(defaultDispatchPolicy()).toEqual({ maxConcurrent: 1, eager: true });
+ });
+});
+
+describe("generateTurnId", () => {
+ it("returns a string starting with 'turn-'", () => {
+ const id = generateTurnId();
+ expect(id).toMatch(/^turn-/);
+ });
+
+ it("returns unique ids", () => {
+ const ids = new Set(Array.from({ length: 100 }, () => generateTurnId()));
+ expect(ids.size).toBe(100);
+ });
+});
diff --git a/packages/session-orchestrator/src/pure.ts b/packages/session-orchestrator/src/pure.ts
new file mode 100644
index 0000000..46cb79a
--- /dev/null
+++ b/packages/session-orchestrator/src/pure.ts
@@ -0,0 +1,27 @@
+import type { ChatMessage, ProviderContract, ToolDispatchPolicy } from "@dispatch/kernel";
+
+export function buildUserMessage(text: string): ChatMessage {
+ return { role: "user", chunks: [{ type: "text", text }] };
+}
+
+export function selectFirstProvider(
+ providers: ReadonlyMap<string, ProviderContract>,
+): ProviderContract {
+ const first = providers.values().next();
+ if (first.done === true || first.value === undefined) {
+ throw new Error("No providers registered — at least one provider is required to run a turn.");
+ }
+ return first.value;
+}
+
+export function resolveTools(tools: ReadonlyMap<string, unknown>): readonly unknown[] {
+ return [...tools.values()];
+}
+
+export function defaultDispatchPolicy(): ToolDispatchPolicy {
+ return { maxConcurrent: 1, eager: true };
+}
+
+export function generateTurnId(): string {
+ return `turn-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`;
+}