summaryrefslogtreecommitdiffhomepage
path: root/packages/session-orchestrator/src
diff options
context:
space:
mode:
authorAdam Malczewski <[email protected]>2026-06-11 14:48:30 +0900
committerAdam Malczewski <[email protected]>2026-06-11 14:48:30 +0900
commitbfbad3af79cab23f52be0f6388311a5798b7fd04 (patch)
tree7bda9843aca8d7fad589ba0b7b5579e2627a3585 /packages/session-orchestrator/src
parent58e2ad559cccc8b35c513818e253b04e60af69b8 (diff)
downloaddispatch-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.ts7
-rw-r--r--packages/session-orchestrator/src/index.ts3
-rw-r--r--packages/session-orchestrator/src/orchestrator.test.ts108
-rw-r--r--packages/session-orchestrator/src/orchestrator.ts21
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;
},
};
}