summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-12 20:13:55 +0900
committerAdam Malczewski <[email protected]>2026-06-12 20:13:55 +0900
commit020e051040001320955a70d6dcaab2d833013196 (patch)
tree1a0921487ae3c89befdbccc1754cd399c07ce1b9 /packages/session-orchestrator/src
parent35197ed933044d322d0a653c4e88a5f3e475fe76 (diff)
downloaddispatch-020e051040001320955a70d6dcaab2d833013196.tar.gz
dispatch-020e051040001320955a70d6dcaab2d833013196.zip
feat(reasoning-effort): persisted per-conversation + per-turn override, threaded to providers
- conversation-store: get/setReasoningEffort (own key space, mirrors cwd) - session-orchestrator: resolveReasoningEffort (override -> stored -> 'high'), StartTurnInput.reasoningEffort, warm() parity (cache-safe) - transport-http: /chat validation (400 on bad level) + GET/PUT /conversations/:id/reasoning-effort - transport-ws: chat.send threading + validation - cli: --effort <low|medium|high|xhigh|max> 993 vitest + 189 bun tests green; typecheck + biome clean.
Diffstat (limited to 'packages/session-orchestrator/src')
-rw-r--r--packages/session-orchestrator/src/index.ts1
-rw-r--r--packages/session-orchestrator/src/orchestrator.test.ts153
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts54
-rw-r--r--packages/session-orchestrator/src/pure.test.ts33
-rw-r--r--packages/session-orchestrator/src/pure.ts19
5 files changed, 244 insertions, 16 deletions
diff --git a/packages/session-orchestrator/src/index.ts b/packages/session-orchestrator/src/index.ts
index b99c15e..711fa5a 100644
--- a/packages/session-orchestrator/src/index.ts
+++ b/packages/session-orchestrator/src/index.ts
@@ -25,6 +25,7 @@ export {
buildUserMessage,
defaultDispatchPolicy,
generateTurnId,
+ resolveReasoningEffort,
selectFirstProvider,
} from "./pure.js";
export { type ToolAssembly, toolsFilter } from "./tools-filter.js";
diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts
index 799fef5..39996b0 100644
--- a/packages/session-orchestrator/src/orchestrator.test.ts
+++ b/packages/session-orchestrator/src/orchestrator.test.ts
@@ -7,6 +7,7 @@ import type {
ProviderContract,
ProviderEvent,
ProviderStreamOptions,
+ ReasoningEffort,
RunTurnInput,
RunTurnResult,
StoredChunk,
@@ -27,14 +28,17 @@ function createInMemoryStore(): ConversationStore & {
readonly data: Map<string, ChatMessage[]>;
readonly metricsData: Map<string, TurnMetrics[]>;
readonly cwdData: Map<string, string>;
+ readonly effortData: Map<string, ReasoningEffort>;
} {
const data = new Map<string, ChatMessage[]>();
const metricsData = new Map<string, TurnMetrics[]>();
const cwdData = new Map<string, string>();
+ const effortData = new Map<string, ReasoningEffort>();
return {
data,
metricsData,
cwdData,
+ effortData,
async append(conversationId, messages) {
const existing = data.get(conversationId) ?? [];
data.set(conversationId, [...existing, ...messages]);
@@ -69,6 +73,12 @@ function createInMemoryStore(): ConversationStore & {
async setCwd(conversationId, cwd) {
cwdData.set(conversationId, cwd);
},
+ async getReasoningEffort(conversationId) {
+ return effortData.get(conversationId) ?? null;
+ },
+ async setReasoningEffort(conversationId, effort) {
+ effortData.set(conversationId, effort);
+ },
};
}
@@ -288,7 +298,7 @@ describe("handleMessage model resolution", () => {
expect(captured).toHaveLength(1);
expect(captured[0]?.provider).toBe(resolvedProvider);
- expect(captured[0]?.providerOpts).toEqual({ model: "gpt-4" });
+ expect(captured[0]?.providerOpts).toEqual({ reasoningEffort: "high", model: "gpt-4" });
expect(captured[0]?.cwd).toBe("/work/dir");
});
@@ -349,7 +359,7 @@ describe("handleMessage model resolution", () => {
expect(captured).toHaveLength(1);
expect(captured[0]?.provider).toBe(fallbackProvider);
- expect(captured[0]?.providerOpts).toBeUndefined();
+ expect(captured[0]?.providerOpts).toEqual({ reasoningEffort: "high" });
});
it("cwd is forwarded to RunTurnInput.cwd and absent when not provided", async () => {
@@ -502,6 +512,12 @@ describe("turn-sealed event", () => {
async setCwd(conversationId, cwd) {
await store.setCwd(conversationId, cwd);
},
+ async getReasoningEffort(conversationId) {
+ return store.getReasoningEffort(conversationId);
+ },
+ async setReasoningEffort(conversationId, effort) {
+ await store.setReasoningEffort(conversationId, effort);
+ },
};
const { orchestrator } = createSessionOrchestrator({
@@ -553,6 +569,10 @@ describe("turn-sealed event", () => {
return null;
},
async setCwd() {},
+ async getReasoningEffort() {
+ return null;
+ },
+ async setReasoningEffort() {},
};
const { orchestrator } = createSessionOrchestrator({
@@ -893,6 +913,10 @@ describe("turn metrics persistence", () => {
return null;
},
async setCwd() {},
+ async getReasoningEffort() {
+ return null;
+ },
+ async setReasoningEffort() {},
};
const { orchestrator } = createSessionOrchestrator({
@@ -2256,3 +2280,128 @@ describe("closeConversation (CR-4c)", () => {
expect(orchestrator.closeConversation("conv-never-seen").abortedTurn).toBe(false);
});
});
+
+describe("reasoning effort resolution", () => {
+ it("override wins over stored → provider receives the override level", async () => {
+ const store = createInMemoryStore();
+ await store.setReasoningEffort("conv-effort-override", "low");
+ const provider: ProviderContract = { id: "p", stream: async function* () {} };
+ const { captured, captureRunTurn } = createCapturingRunTurn();
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: captureRunTurn,
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-effort-override",
+ text: "hi",
+ onEvent: () => {},
+ reasoningEffort: "max",
+ });
+
+ expect(captured).toHaveLength(1);
+ expect(captured[0]?.providerOpts?.reasoningEffort).toBe("max");
+ });
+
+ it("no override, store has a value → provider receives the stored value", async () => {
+ const store = createInMemoryStore();
+ await store.setReasoningEffort("conv-effort-stored", "xhigh");
+ const provider: ProviderContract = { id: "p", stream: async function* () {} };
+ const { captured, captureRunTurn } = createCapturingRunTurn();
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: captureRunTurn,
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-effort-stored",
+ text: "hi",
+ onEvent: () => {},
+ });
+
+ expect(captured).toHaveLength(1);
+ expect(captured[0]?.providerOpts?.reasoningEffort).toBe("xhigh");
+ });
+
+ it("no override, store empty → provider receives 'high' (default)", async () => {
+ const store = createInMemoryStore();
+ const provider: ProviderContract = { id: "p", stream: async function* () {} };
+ const { captured, captureRunTurn } = createCapturingRunTurn();
+
+ const { orchestrator } = createSessionOrchestrator({
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: captureRunTurn,
+ });
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-effort-default",
+ text: "hi",
+ onEvent: () => {},
+ });
+
+ expect(captured).toHaveLength(1);
+ expect(captured[0]?.providerOpts?.reasoningEffort).toBe("high");
+ });
+
+ it("warm receives the same resolved effort as a real turn for the same conversation", async () => {
+ const store = createInMemoryStore();
+ await store.append("conv-warm-effort", [
+ { role: "user", chunks: [{ type: "text", text: "hi" }] },
+ ]);
+ await store.setReasoningEffort("conv-warm-effort", "medium");
+
+ let warmOpts: ProviderStreamOptions | undefined;
+
+ const provider: ProviderContract = {
+ id: "p",
+ stream(_messages, _tools, opts) {
+ warmOpts = opts;
+ return (async function* () {
+ yield {
+ type: "usage",
+ usage: { inputTokens: 1, outputTokens: 1, cacheReadTokens: 0, cacheWriteTokens: 0 },
+ } as ProviderEvent;
+ yield { type: "finish", reason: "stop" } as ProviderEvent;
+ })();
+ },
+ };
+
+ const { captured, captureRunTurn } = createCapturingRunTurn();
+
+ const deps = {
+ conversationStore: store,
+ resolveProvider: () => provider,
+ resolveTools: () => [],
+ applyToolsFilter: identityApplyToolsFilter,
+ runTurn: captureRunTurn,
+ emit: () => {},
+ };
+
+ const { orchestrator, activeConversations } = createSessionOrchestrator(deps);
+ const warmService = createWarmService(deps, activeConversations);
+
+ await warmService.warm("conv-warm-effort");
+
+ await orchestrator.handleMessage({
+ conversationId: "conv-warm-effort",
+ text: "hi",
+ onEvent: () => {},
+ });
+
+ expect(warmOpts?.reasoningEffort).toBe("medium");
+ expect(captured).toHaveLength(1);
+ expect(captured[0]?.providerOpts?.reasoningEffort).toBe("medium");
+ expect(warmOpts?.reasoningEffort).toBe(captured[0]?.providerOpts?.reasoningEffort);
+ });
+});
diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts
index b0a1083..5b2f264 100644
--- a/packages/session-orchestrator/src/orchestrator.ts
+++ b/packages/session-orchestrator/src/orchestrator.ts
@@ -7,6 +7,7 @@ import type {
ProviderContract,
ProviderEvent,
ProviderStreamOptions,
+ ReasoningEffort,
RunTurnInput,
RunTurnResult,
ToolContract,
@@ -15,7 +16,12 @@ import type {
} from "@dispatch/kernel";
import { defineEventHook, defineService, type ServiceHandle } from "@dispatch/kernel";
import { createMetricsAccumulator } from "./metrics.js";
-import { buildUserMessage, defaultDispatchPolicy, generateTurnId } from "./pure.js";
+import {
+ buildUserMessage,
+ defaultDispatchPolicy,
+ generateTurnId,
+ resolveReasoningEffort,
+} from "./pure.js";
import type { ToolAssembly } from "./tools-filter.js";
// --- Broadcast hub types ---
@@ -25,6 +31,7 @@ export interface StartTurnInput {
readonly text: string;
readonly modelName?: string;
readonly cwd?: string;
+ readonly reasoningEffort?: ReasoningEffort;
}
export type StartTurnResult =
@@ -119,6 +126,7 @@ export interface SessionOrchestrator {
onEvent: (event: AgentEvent) => void;
modelName?: string;
cwd?: string;
+ reasoningEffort?: ReasoningEffort;
}): Promise<void>;
}
@@ -181,6 +189,7 @@ export function createSessionOrchestrator(
text: string,
modelName: string | undefined,
cwd: string | undefined,
+ reasoningEffortOverride: ReasoningEffort | undefined,
): void {
const turnId = generateTurnId();
const controller = new AbortController();
@@ -194,11 +203,15 @@ export function createSessionOrchestrator(
? Promise.resolve(cwd)
: deps.conversationStore.getCwd(conversationId).then((c) => c ?? undefined);
- const payloadPromise = effectiveCwdPromise.then((effectiveCwd) => ({
- conversationId,
- ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
- ...(modelName !== undefined ? { modelName } : {}),
- }));
+ const storedEffortPromise = deps.conversationStore.getReasoningEffort(conversationId);
+
+ const payloadPromise = Promise.all([effectiveCwdPromise, storedEffortPromise]).then(
+ ([effectiveCwd]) => ({
+ conversationId,
+ ...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
+ ...(modelName !== undefined ? { modelName } : {}),
+ }),
+ );
payloadPromise.then((payload) => {
deps.emit?.(turnStarted, payload);
@@ -206,12 +219,17 @@ export function createSessionOrchestrator(
void (async () => {
try {
- const effectiveCwd = await effectiveCwdPromise;
+ const [effectiveCwd, storedEffort] = await Promise.all([
+ effectiveCwdPromise,
+ storedEffortPromise,
+ ]);
if (cwd !== undefined) {
await deps.conversationStore.setCwd(conversationId, cwd);
}
+ const resolvedEffort = resolveReasoningEffort(reasoningEffortOverride, storedEffort);
+
const history = await deps.conversationStore.load(conversationId);
const userMsg = buildUserMessage(text);
@@ -250,6 +268,11 @@ export function createSessionOrchestrator(
emitToHub(conversationId, event);
};
+ const providerOpts: ProviderStreamOptions = {
+ reasoningEffort: resolvedEffort,
+ ...(modelOverride !== undefined ? { model: modelOverride } : {}),
+ };
+
const opts: RunTurnInput = {
provider,
messages: [...history, userMsg],
@@ -259,9 +282,7 @@ export function createSessionOrchestrator(
conversationId,
turnId,
signal: controller.signal,
- ...(modelOverride !== undefined
- ? { providerOpts: { model: modelOverride } satisfies ProviderStreamOptions }
- : {}),
+ providerOpts,
...(turnLogger !== undefined ? { logger: turnLogger } : {}),
...(effectiveCwd !== undefined ? { cwd: effectiveCwd } : {}),
...(deps.now !== undefined ? { now: deps.now } : {}),
@@ -295,11 +316,11 @@ export function createSessionOrchestrator(
}
const orchestrator: SessionOrchestrator = {
- startTurn({ conversationId, text, modelName, cwd }) {
+ startTurn({ conversationId, text, modelName, cwd, reasoningEffort }) {
if (activeTurns.has(conversationId)) {
return { started: false, reason: "already-active" };
}
- runTurnDetached(conversationId, text, modelName, cwd);
+ runTurnDetached(conversationId, text, modelName, cwd, reasoningEffort);
const turn = activeTurns.get(conversationId);
const turnId = turn !== undefined ? turn.turnId : "";
return { started: true, turnId };
@@ -346,12 +367,13 @@ export function createSessionOrchestrator(
return { abortedTurn };
},
- async handleMessage({ conversationId, text, onEvent, modelName, cwd }) {
+ async handleMessage({ conversationId, text, onEvent, modelName, cwd, reasoningEffort }) {
const turnInput: StartTurnInput = {
conversationId,
text,
...(modelName !== undefined ? { modelName } : {}),
...(cwd !== undefined ? { cwd } : {}),
+ ...(reasoningEffort !== undefined ? { reasoningEffort } : {}),
};
const result = orchestrator.startTurn(turnInput);
if (!result.started) {
@@ -424,6 +446,11 @@ export function createWarmService(
...(cwd !== undefined ? { cwd } : {}),
});
+ // Resolve reasoning effort the SAME way the real turn does (stored → "high";
+ // no per-turn override on warm). A mismatch here silently busts the prompt cache.
+ const storedEffort = await deps.conversationStore.getReasoningEffort(conversationId);
+ const resolvedEffort = resolveReasoningEffort(undefined, storedEffort);
+
const probeMsg: ChatMessage = {
role: "user",
chunks: [{ type: "text", text: "reply with just a ." }],
@@ -439,6 +466,7 @@ export function createWarmService(
const warmLogger = deps.logger?.child({ conversationId, attrs: { warm: true } });
const providerOpts: ProviderStreamOptions = {
maxTokens: 1,
+ reasoningEffort: resolvedEffort,
...(modelOverride !== undefined ? { model: modelOverride } : {}),
...(warmLogger !== undefined ? { logger: warmLogger } : {}),
};
diff --git a/packages/session-orchestrator/src/pure.test.ts b/packages/session-orchestrator/src/pure.test.ts
index e233fca..9e5d3c4 100644
--- a/packages/session-orchestrator/src/pure.test.ts
+++ b/packages/session-orchestrator/src/pure.test.ts
@@ -4,6 +4,7 @@ import {
buildUserMessage,
defaultDispatchPolicy,
generateTurnId,
+ resolveReasoningEffort,
selectFirstProvider,
} from "./pure.js";
@@ -67,3 +68,35 @@ describe("generateTurnId", () => {
expect(ids.size).toBe(100);
});
});
+
+describe("resolveReasoningEffort", () => {
+ it("override wins over stored", () => {
+ expect(resolveReasoningEffort("low", "high")).toBe("low");
+ expect(resolveReasoningEffort("max", "medium")).toBe("max");
+ });
+
+ it("stored wins over default", () => {
+ expect(resolveReasoningEffort(undefined, "medium")).toBe("medium");
+ expect(resolveReasoningEffort(undefined, "xhigh")).toBe("xhigh");
+ });
+
+ it("default is 'high' when both are absent", () => {
+ expect(resolveReasoningEffort(undefined, null)).toBe("high");
+ });
+
+ it("all 5 levels pass through as override", () => {
+ expect(resolveReasoningEffort("low", null)).toBe("low");
+ expect(resolveReasoningEffort("medium", null)).toBe("medium");
+ expect(resolveReasoningEffort("high", null)).toBe("high");
+ expect(resolveReasoningEffort("xhigh", null)).toBe("xhigh");
+ expect(resolveReasoningEffort("max", null)).toBe("max");
+ });
+
+ it("all 5 levels pass through as stored", () => {
+ expect(resolveReasoningEffort(undefined, "low")).toBe("low");
+ expect(resolveReasoningEffort(undefined, "medium")).toBe("medium");
+ expect(resolveReasoningEffort(undefined, "high")).toBe("high");
+ expect(resolveReasoningEffort(undefined, "xhigh")).toBe("xhigh");
+ expect(resolveReasoningEffort(undefined, "max")).toBe("max");
+ });
+});
diff --git a/packages/session-orchestrator/src/pure.ts b/packages/session-orchestrator/src/pure.ts
index 46cb79a..85edd14 100644
--- a/packages/session-orchestrator/src/pure.ts
+++ b/packages/session-orchestrator/src/pure.ts
@@ -1,9 +1,26 @@
-import type { ChatMessage, ProviderContract, ToolDispatchPolicy } from "@dispatch/kernel";
+import type {
+ ChatMessage,
+ ProviderContract,
+ ReasoningEffort,
+ ToolDispatchPolicy,
+} from "@dispatch/kernel";
export function buildUserMessage(text: string): ChatMessage {
return { role: "user", chunks: [{ type: "text", text }] };
}
+/**
+ * Resolve the reasoning-effort level for a turn:
+ * per-turn override → persisted per-conversation value → default `"high"`.
+ * Pure — no I/O, no ambient state.
+ */
+export function resolveReasoningEffort(
+ override: ReasoningEffort | undefined,
+ stored: ReasoningEffort | null,
+): ReasoningEffort {
+ return override ?? stored ?? "high";
+}
+
export function selectFirstProvider(
providers: ReadonlyMap<string, ProviderContract>,
): ProviderContract {