diff options
| author | Adam Malczewski <[email protected]> | 2026-06-12 20:13:55 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-12 20:13:55 +0900 |
| commit | 020e051040001320955a70d6dcaab2d833013196 (patch) | |
| tree | 1a0921487ae3c89befdbccc1754cd399c07ce1b9 /packages/session-orchestrator/src | |
| parent | 35197ed933044d322d0a653c4e88a5f3e475fe76 (diff) | |
| download | dispatch-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.ts | 1 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.test.ts | 153 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.ts | 54 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/pure.test.ts | 33 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/pure.ts | 19 |
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 { |
