diff options
| author | Adam Malczewski <[email protected]> | 2026-06-11 14:48:30 +0900 |
|---|---|---|
| committer | Adam Malczewski <[email protected]> | 2026-06-11 14:48:30 +0900 |
| commit | bfbad3af79cab23f52be0f6388311a5798b7fd04 (patch) | |
| tree | 7bda9843aca8d7fad589ba0b7b5579e2627a3585 /packages/session-orchestrator/src | |
| parent | 58e2ad559cccc8b35c513818e253b04e60af69b8 (diff) | |
| download | dispatch-bfbad3af79cab23f52be0f6388311a5798b7fd04.tar.gz dispatch-bfbad3af79cab23f52be0f6388311a5798b7fd04.zip | |
feat(cache-warming): CR-3 — manual warm resets timer + nextWarmAt/lastWarmAt surface
FE CR-3 (backend-handoff-cache-warming-timer.md). The inversion: session-orchestrator's
warm() (the single chokepoint for manual /chat/warm AND the automatic timer) emits a
warmCompleted bus event; cache-warming subscribes and does ALL post-warm handling. So a
manual warm now re-arms the timer + refreshes the surface with NO transport-http change
(core can't depend on the standard cache-warming ext).
- session-orchestrator: warmCompleted event hook + emit from warm() on success
- cache-warming: warmCompleted subscriber unifies result handling (manual + automatic);
adds nextWarmAt/lastWarmAt state + a custom 'cache-warming-timer' surface field
- fix: createWarmService was missing the emit dep (deps.emit?. silently no-oped) →
wired it + made emit REQUIRED so it can't regress
Live-verified vs claude haiku: manual POST /chat/warm now logs cache-warming 'warm
complete' ~2s after the turn (not the 4-min timer) → manual warm reaches the warmer.
800 vitest + 109 bun green; tsc -b 0; biome clean.
Diffstat (limited to 'packages/session-orchestrator/src')
| -rw-r--r-- | packages/session-orchestrator/src/extension.ts | 7 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/index.ts | 3 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.test.ts | 108 | ||||
| -rw-r--r-- | packages/session-orchestrator/src/orchestrator.ts | 21 |
4 files changed, 136 insertions, 3 deletions
diff --git a/packages/session-orchestrator/src/extension.ts b/packages/session-orchestrator/src/extension.ts index 12d387c..4175b0c 100644 --- a/packages/session-orchestrator/src/extension.ts +++ b/packages/session-orchestrator/src/extension.ts @@ -21,7 +21,11 @@ export const manifest: Manifest = { activation: "eager", contributes: { services: ["session-orchestrator/orchestrator", "session-orchestrator/warm"], - hooks: ["session-orchestrator/turn-started", "session-orchestrator/turn-settled"], + hooks: [ + "session-orchestrator/turn-started", + "session-orchestrator/turn-settled", + "session-orchestrator/warm-completed", + ], }, }; @@ -64,6 +68,7 @@ export function activate(host: HostAPI): void { runTurn, logger: host.logger, now: () => Date.now(), + emit: (hook, payload) => host.emit(hook, payload), }, activeConversations, ); diff --git a/packages/session-orchestrator/src/index.ts b/packages/session-orchestrator/src/index.ts index 37ae5ce..2daf278 100644 --- a/packages/session-orchestrator/src/index.ts +++ b/packages/session-orchestrator/src/index.ts @@ -10,8 +10,11 @@ export { type TurnLifecyclePayload, turnSettled, turnStarted, + type WarmCompletedPayload, type WarmResult, type WarmService, + type WarmServiceDeps, + warmCompleted, } from "./orchestrator.js"; export { buildUserMessage, diff --git a/packages/session-orchestrator/src/orchestrator.test.ts b/packages/session-orchestrator/src/orchestrator.test.ts index 5d512ea..ba4912a 100644 --- a/packages/session-orchestrator/src/orchestrator.test.ts +++ b/packages/session-orchestrator/src/orchestrator.test.ts @@ -17,6 +17,7 @@ import { createSessionOrchestrator, createWarmService, type TurnLifecyclePayload, + type WarmCompletedPayload, } from "./orchestrator.js"; import type { ToolAssembly } from "./tools-filter.js"; @@ -1060,6 +1061,7 @@ describe("warm service", () => { resolveTools: () => [toolA], applyToolsFilter: identityApplyToolsFilter, runTurn, + emit: () => {}, }; const { activeConversations } = createSessionOrchestrator(deps); @@ -1115,6 +1117,7 @@ describe("warm service", () => { resolveTools: () => [], applyToolsFilter: identityApplyToolsFilter, runTurn: blockingRunTurn, + emit: () => {}, }; const { orchestrator, activeConversations } = createSessionOrchestrator(deps); @@ -1158,6 +1161,7 @@ describe("warm service", () => { resolveTools: () => [], applyToolsFilter: identityApplyToolsFilter, runTurn, + emit: () => {}, }; const { activeConversations } = createSessionOrchestrator(deps); @@ -1201,6 +1205,7 @@ describe("warm service", () => { resolveTools: () => [], applyToolsFilter: identityApplyToolsFilter, runTurn, + emit: () => {}, }; const { activeConversations } = createSessionOrchestrator(deps); @@ -1215,4 +1220,107 @@ describe("warm service", () => { cacheWriteTokens: 100, }); }); + + it("warm emits warmCompleted with the usage on success", async () => { + const store = createInMemoryStore(); + const existingMsg: ChatMessage = { + role: "user", + chunks: [{ type: "text", text: "existing" }], + }; + await store.append("conv-warm-emit", [existingMsg]); + + const provider: ProviderContract = { + id: "p", + stream: async function* () { + yield { + type: "usage", + usage: { inputTokens: 200, outputTokens: 10, cacheReadTokens: 150, cacheWriteTokens: 50 }, + } as ProviderEvent; + yield { type: "finish", reason: "stop" } as ProviderEvent; + }, + }; + + const emitted: Array<{ hook: string; payload: WarmCompletedPayload }> = []; + const fakeEmit = <TPayload>(hook: EventHookDescriptor<TPayload>, payload: TPayload): void => { + emitted.push({ hook: hook.id, payload: payload as WarmCompletedPayload }); + }; + + const deps = { + conversationStore: store, + resolveProvider: () => provider, + resolveTools: () => [], + applyToolsFilter: identityApplyToolsFilter, + runTurn, + emit: fakeEmit, + }; + + const { activeConversations } = createSessionOrchestrator(deps); + const warmService = createWarmService(deps, activeConversations); + + const result = await warmService.warm("conv-warm-emit"); + + if (!("inputTokens" in result)) throw new Error("expected success"); + + expect(emitted).toHaveLength(1); + expect(emitted[0]?.hook).toBe("session-orchestrator/warm-completed"); + expect(emitted[0]?.payload.conversationId).toBe("conv-warm-emit"); + expect(emitted[0]?.payload.usage).toEqual(result); + }); + + it("warm does NOT emit warmCompleted when it refuses (conversation generating / no history)", async () => { + const store = createInMemoryStore(); + + const provider: ProviderContract = { + id: "p", + stream: async function* () { + yield { type: "text-delta", delta: "slow" } as ProviderEvent; + yield { type: "finish", reason: "stop" } as ProviderEvent; + }, + }; + + const emitted: Array<{ hook: string }> = []; + const fakeEmit = <TPayload>(hook: EventHookDescriptor<TPayload>, _payload: TPayload): void => { + emitted.push({ hook: hook.id }); + }; + + const blockingRunTurn = async (_input: RunTurnInput): Promise<RunTurnResult> => { + await new Promise<void>((resolve) => setTimeout(resolve, 50)); + return { + messages: [{ role: "assistant", chunks: [{ type: "text", text: "done" }] }], + usage: { inputTokens: 1, outputTokens: 1 }, + finishReason: "stop", + }; + }; + + const deps = { + conversationStore: store, + resolveProvider: () => provider, + resolveTools: () => [], + applyToolsFilter: identityApplyToolsFilter, + runTurn: blockingRunTurn, + emit: fakeEmit, + }; + + const { orchestrator, activeConversations } = createSessionOrchestrator(deps); + const warmService = createWarmService(deps, activeConversations); + + // Refuse because conversation is generating + const turnPromise = orchestrator.handleMessage({ + conversationId: "conv-refuse-gen", + text: "test", + onEvent: () => {}, + }); + + const genResult = await warmService.warm("conv-refuse-gen"); + expect(genResult).toEqual({ error: "conversation is generating" }); + + await turnPromise; + + // Refuse because no history + const noHistResult = await warmService.warm("conv-refuse-empty"); + expect(noHistResult).toEqual({ error: "no history" }); + + const warmEmits = emitted.filter((e) => e.hook === "session-orchestrator/warm-completed"); + expect(warmEmits).toHaveLength(0); + }); }); diff --git a/packages/session-orchestrator/src/orchestrator.ts b/packages/session-orchestrator/src/orchestrator.ts index c39bc06..6df92c8 100644 --- a/packages/session-orchestrator/src/orchestrator.ts +++ b/packages/session-orchestrator/src/orchestrator.ts @@ -35,6 +35,16 @@ export const turnStarted: EventHookDescriptor<TurnLifecyclePayload> = export const turnSettled: EventHookDescriptor<TurnLifecyclePayload> = defineEventHook<TurnLifecyclePayload>("session-orchestrator/turn-settled"); +/** Payload for the warmCompleted bus event. */ +export interface WarmCompletedPayload { + readonly conversationId: string; + readonly usage: WarmResult; +} + +/** Fired when a warm probe succeeds (both automatic and manual paths). */ +export const warmCompleted: EventHookDescriptor<WarmCompletedPayload> = + defineEventHook<WarmCompletedPayload>("session-orchestrator/warm-completed"); + // --- Warm service --- export interface WarmResult { @@ -89,6 +99,11 @@ export interface SessionOrchestratorDeps { readonly emit?: <TPayload>(hook: EventHookDescriptor<TPayload>, payload: TPayload) => void; } +/** Deps for the warm service — emit is REQUIRED so warmCompleted is never silently dropped. */ +export type WarmServiceDeps = SessionOrchestratorDeps & { + readonly emit: <TPayload>(hook: EventHookDescriptor<TPayload>, payload: TPayload) => void; +}; + export interface SessionOrchestratorBundle { readonly orchestrator: SessionOrchestrator; /** The shared active-conversations set, for use by createWarmService. */ @@ -187,7 +202,7 @@ export function createSessionOrchestrator( } export function createWarmService( - deps: SessionOrchestratorDeps, + deps: WarmServiceDeps, activeConversations: ReadonlySet<string>, ): WarmService { return { @@ -247,7 +262,9 @@ export function createWarmService( } } - return { inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens }; + const result: WarmResult = { inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens }; + deps.emit(warmCompleted, { conversationId, usage: result }); + return result; }, }; } |
