import { join } from "node:path"; import { type Api, type AssistantMessage, type AssistantMessageEventStream, calculateCost, createAssistantMessageEventStream, getCurrentSystemPrompt, getCurrentTools, type Model, type SimpleStreamOptions, type TranscriptContext, } from "@earendil-works/pi-ai"; import type { KlausQueryHandle } from "./agent-sdk.js"; import { FileSessionStore, KlausCacheCorruptionError, type KlausSessionStore, loadCheckpoints, MemorySessionStore, StagingSessionStore, saveCheckpoint, } from "./cache.js"; import { dbg, diagnosticKind } from "./debug.js"; import { contextFingerprint, createAssistant, isKlausModelId, type KlausContentEvent, type KlausContext, type KlausQueryCallbacks, type KlausRequest, klausSelector, modelFacingContent, requiredConstrainedTool, SDK_VERSION, toKlausTools, } from "./protocol.js"; import { replayImages, replayPrompt } from "./replay.js"; import { rewriteSubscriptionPrompt } from "./subscription-prompt.js"; import type { PiToolResult } from "./tool-bridge.js"; interface ActiveQuery { id: symbol; sessionId?: string; handle?: KlausQueryHandle; sink?: PiSink; toolCallIds: Set; tools: KlausRequest["tools"]; boundaryFingerprint?: string; boundaryMessageCount?: number; closed: boolean; closePromise?: Promise; timer?: ReturnType; timeoutMs?: number; idleLease?: ReturnType; /** Set while this child is idle between turns and may take the next one. */ reuseLease?: ReturnType; sdkSessionId?: string; position?: string; messageCount?: number; /** Lineage of the turns this child already holds, computed without tools so * a tool-list change can be answered with a swap instead of a replay. */ lineageFingerprint?: string; abort?: () => void; signal?: AbortSignal; secondary: boolean; committable: boolean; modelId: string; cwd: string; thinking: SimpleStreamOptions["reasoning"]; thinkingBudget?: number; maxTokens?: number; systemPrompt: string; headers: SimpleStreamOptions["headers"]; env: SimpleStreamOptions["env"]; metadata?: Record; emittedContent: boolean; context: KlausContext; store: KlausSessionStore; staging: StagingSessionStore; } interface CompletedQuery { sdkSessionId: string; fingerprint: string; position?: string; piLeafId?: string; messageCount: number; protocol: 1; store: KlausSessionStore; } interface SessionScope { id: string; persisted: boolean; dir: string; } /** How long a finished child stays alive waiting for the next turn. Tests may * shorten it, because proving expiry must not cost two minutes per run. */ function reuseLeaseMs(): number { const override = process.env.VITEST ? Number(process.env.KLAUS_REUSE_LEASE_MS) : Number.NaN; return Number.isSafeInteger(override) && override > 0 ? override : 2 * 60_000; } type ResolveOAuth = (modelId: string) => Promise; type CurrentCwd = () => string; type CurrentSession = () => SessionScope | undefined; type Warn = (message: string) => void; class PiSink { readonly stream = createAssistantMessageEventStream(); readonly output: AssistantMessage; private indices = new Map(); private toolJson = new Map(); private finished = false; constructor(private readonly model: Model) { dbg?.("coordinator.sink.create"); this.output = createAssistant(model); this.stream.push({ type: "start", partial: this.output }); } content(event: KlausContentEvent): void { dbg?.("coordinator.sink.content", { index: event.index, length: "delta" in event ? event.delta.length : "text" in event ? event.text.length : "thinking" in event ? event.thinking.length : undefined, }); if (this.finished) return; switch (event.type) { case "text-start": { const index = this.output.content.length; this.indices.set(event.index, index); this.output.content.push({ type: "text", text: "" }); this.stream.push({ type: "text_start", contentIndex: index, partial: this.output, }); break; } case "text-delta": { const index = this.indices.get(event.index); const block = index === undefined ? undefined : this.output.content[index]; if (index === undefined || block?.type !== "text") return; block.text += event.delta; this.stream.push({ type: "text_delta", contentIndex: index, delta: event.delta, partial: this.output, }); break; } case "text-end": { const index = this.indices.get(event.index); if (index === undefined) return; this.stream.push({ type: "text_end", contentIndex: index, content: event.text, partial: this.output, }); break; } case "thinking-start": { const index = this.output.content.length; this.indices.set(event.index, index); this.output.content.push({ type: "thinking", thinking: "" }); this.stream.push({ type: "thinking_start", contentIndex: index, partial: this.output, }); break; } case "thinking-delta": { const index = this.indices.get(event.index); const block = index === undefined ? undefined : this.output.content[index]; if (index === undefined || block?.type !== "thinking") return; block.thinking += event.delta; this.stream.push({ type: "thinking_delta", contentIndex: index, delta: event.delta, partial: this.output, }); break; } case "thinking-end": { const index = this.indices.get(event.index); const block = index === undefined ? undefined : this.output.content[index]; if (index === undefined || block?.type !== "thinking") return; block.thinkingSignature = event.signature; block.redacted = event.redacted; this.stream.push({ type: "thinking_end", contentIndex: index, content: event.thinking, partial: this.output, }); break; } case "tool-start": { const index = this.output.content.length; this.indices.set(event.index, index); this.toolJson.set(event.index, ""); this.output.content.push({ type: "toolCall", id: event.id, name: event.name, arguments: {}, }); this.stream.push({ type: "toolcall_start", contentIndex: index, partial: this.output, }); break; } case "tool-delta": { const index = this.indices.get(event.index); if (index === undefined) return; this.toolJson.set( event.index, `${this.toolJson.get(event.index) ?? ""}${event.delta}`, ); this.stream.push({ type: "toolcall_delta", contentIndex: index, delta: event.delta, partial: this.output, }); break; } case "tool-end": { const index = this.indices.get(event.index); const block = index === undefined ? undefined : this.output.content[index]; if (index === undefined || block?.type !== "toolCall") return; block.arguments = event.arguments; this.stream.push({ type: "toolcall_end", contentIndex: index, toolCall: block, partial: this.output, }); break; } } } toolUse( usage: { input: number; output: number; cacheRead: number; cacheWrite: number; }, responseId?: string, responseModel?: string, ): void { dbg?.("coordinator.sink.toolUse", { finished: this.finished }); if (this.finished) return; this.finished = true; this.applyUsage(usage, responseId, responseModel); this.output.stopReason = "toolUse"; this.stream.push({ type: "done", reason: "toolUse", message: this.output }); this.stream.end(); } done( usage: { input: number; output: number; cacheRead: number; cacheWrite: number; }, responseId?: string, responseModel?: string, stopReason: "stop" | "length" = "stop", ): void { dbg?.("coordinator.sink.done", { finished: this.finished }); if (this.finished) return; this.finished = true; this.applyUsage(usage, responseId, responseModel); this.output.stopReason = stopReason; this.stream.push({ type: "done", reason: stopReason, message: this.output, }); this.stream.end(); } private applyUsage( usage: { input: number; output: number; cacheRead: number; cacheWrite: number; }, responseId?: string, responseModel?: string, ): void { Object.assign(this.output.usage, usage); this.output.usage.totalTokens = usage.input + usage.output + usage.cacheRead + usage.cacheWrite; const fallbackCost = this.model.compat && "allowedFallbackModels" in this.model.compat && responseModel !== this.model.id ? this.model.compat.allowedFallbackModels?.find( (fallback) => fallback.provider === "anthropic" && fallback.model === responseModel, )?.cost : undefined; calculateCost( fallbackCost && responseModel ? { ...this.model, id: responseModel, cost: fallbackCost } : this.model, this.output.usage, ); this.output.responseId = responseId; this.output.responseModel = responseModel; } error(error: unknown, aborted = false): void { dbg?.("coordinator.sink.error", { finished: this.finished, aborted, kind: diagnosticKind(error), }); if (this.finished) return; this.finished = true; this.output.stopReason = aborted ? "aborted" : "error"; this.output.errorMessage = error instanceof Error ? error.message : String(error); this.stream.push({ type: "error", reason: this.output.stopReason, error: this.output, }); this.stream.end(); } } export class QueryCoordinator { private readonly bySession = new Map>(); private readonly anonymous = new Set(); private readonly completed = new Map(); private readonly memoryStores = new Map(); private readonly diskStores = new Map(); private readonly pendingCheckpoints = new Map(); private cacheReady: Promise = Promise.resolve(); constructor( private readonly resolveOAuth: ResolveOAuth, private readonly currentCwd: CurrentCwd, private readonly currentSession: CurrentSession = () => undefined, private readonly warn: Warn = () => undefined, ) { dbg?.("coordinator.create"); } stream( model: Model, context: TranscriptContext, options: SimpleStreamOptions = {}, ): AssistantMessageEventStream { const currentContext: KlausContext = { messages: context.messages.filter( (message): message is KlausContext["messages"][number] => message.role !== "system", ), systemPrompt: getCurrentSystemPrompt(context.messages), tools: getCurrentTools(context.messages), }; dbg?.("coordinator.stream", { messageCount: currentContext.messages.length, toolCount: currentContext.tools?.length ?? 0, }); const sink = new PiSink(model); void this.invoke(sink, model, currentContext, options); return sink.stream; } preparePrimaryCache(scope: SessionScope): Promise { dbg?.("coordinator.preparePrimaryCache", { persisted: scope.persisted }); if (!scope.persisted) { this.cacheReady = Promise.resolve(); return this.cacheReady; } const store = new FileSessionStore( join(scope.dir, ".klaus", SDK_VERSION, "sessions"), ); this.diskStores.set(scope.id, store); this.cacheReady = store.prepare(); return this.cacheReady; } async persistCheckpoint(sessionId: string, piLeafId: string): Promise { dbg?.("coordinator.persistCheckpoint.start"); const checkpoint = this.pendingCheckpoints.get(sessionId); const root = this.checkpointRoot(sessionId); if (!checkpoint || !root) { dbg?.("coordinator.persistCheckpoint.skip", { hasCheckpoint: Boolean(checkpoint), hasRoot: Boolean(root), }); return; } checkpoint.piLeafId = piLeafId; dbg?.("coordinator.persistCheckpoint.position", { messageCount: checkpoint.messageCount, }); await saveCheckpoint(root, sessionId, checkpoint); this.pendingCheckpoints.delete(sessionId); dbg?.("coordinator.persistCheckpoint.end"); } async closeAll(reason: string): Promise { dbg?.("coordinator.closeAll.start"); const states = [ ...this.anonymous, ...[...this.bySession.values()].flatMap((states) => [...states]), ]; await Promise.all(states.map((state) => this.closeState(state, reason))); this.completed.clear(); this.memoryStores.clear(); this.diskStores.clear(); this.pendingCheckpoints.clear(); dbg?.("coordinator.closeAll.end", { stateCount: states.length }); } private async invoke( sink: PiSink, model: Model, context: KlausContext, options: SimpleStreamOptions, ): Promise { dbg?.("coordinator.invoke.start", { messageCount: context.messages.length, }); try { if (!isKlausModelId(model.id)) throw new Error(`Unsupported Klaus model ${model.id}.`); const required = requiredConstrainedTool(context.tools); if (required) throw new Error( `Tool ${required} requires constrained sampling, which Klaus cannot provide.`, ); const batch = recentToolResults(context); dbg?.("coordinator.invoke.batch", { resultCount: batch.results.length }); const continuation = this.findContinuation( options.sessionId, batch.results, ); const tools = toKlausTools(context.tools); dbg?.("coordinator.invoke.continuation", { found: Boolean(continuation), hasHandle: Boolean(continuation?.handle), }); if (continuation?.handle) { const resultCount = batch.results.length; const boundaryMessageCount = context.messages.length - resultCount; const boundaryFingerprint = contextFingerprint( { systemPrompt: context.systemPrompt, messages: context.messages.slice(0, boundaryMessageCount), tools: context.tools, }, `${model.id}\0${this.currentCwd()}`, ); const incompatible = batch.userAfter || continuation.boundaryMessageCount !== boundaryMessageCount || continuation.boundaryFingerprint !== boundaryFingerprint || JSON.stringify(continuation.tools) !== JSON.stringify(tools) || continuation.modelId !== model.id || continuation.cwd !== this.currentCwd() || continuation.systemPrompt !== (context.systemPrompt ?? "") || continuation.thinking !== options.reasoning || continuation.thinkingBudget !== thinkingBudget(options) || continuation.maxTokens !== options.maxTokens || !sameStringRecord(continuation.headers, options.headers) || !sameStringRecord(continuation.env, options.env); dbg?.("coordinator.invoke.compatibility", { incompatible }); if (incompatible) { await this.closeState( continuation, "Pi request changed across a tool boundary; replaying canonically.", ); } else { continuation.sink = sink; continuation.context = context; if (continuation.idleLease) clearTimeout(continuation.idleLease); continuation.idleLease = undefined; this.replaceAbort(continuation, options.signal); continuation.timeoutMs = options.timeoutMs; this.armTimeout(continuation, continuation.timeoutMs); if (options.signal?.aborted) { await this.closeState( continuation, "Pi aborted the Klaus query.", true, ); return; } if (options.onResponse && continuation.metadata) { await options.onResponse( { status: 200, headers: continuation.metadata }, model, ); } if (options.signal?.aborted) { await this.closeState( continuation, "Pi aborted the Klaus query.", true, ); return; } for (const result of batch.results) { await continuation.handle.bridge.waitForPending( [result.id], options.signal, ); if (options.signal?.aborted) { await this.closeState( continuation, "Pi aborted the Klaus query.", true, ); return; } if (!continuation.handle.bridge.settle(result)) { throw new Error(`Unknown Klaus tool result ${result.id}.`); } } return; } } await this.start( sink, model, context, options, tools, false, batch.results.length === 0 ? this.findReusable(options.sessionId) : undefined, ); } catch (error) { dbg?.("coordinator.invoke.error", { kind: diagnosticKind(error) }); sink.error(error, options.signal?.aborted); } } /** `reusable` is an idle child from this session's previous turn. It is used * only when this turn resolves to exactly the delta the child already * expects, and it is closed on any other outcome. */ private async start( sink: PiSink, model: Model, context: KlausContext, options: SimpleStreamOptions, tools: KlausRequest["tools"], forceCold = false, reusable?: ActiveQuery, ): Promise { dbg?.("coordinator.start.begin", { messageCount: context.messages.length }); await this.cacheReady; dbg?.("coordinator.start.cacheReady"); const store = this.storeFor(options.sessionId); const cwd = this.currentCwd(); const identity = `${model.id}\0${cwd}`; const candidates: CompletedQuery[] = []; const memory = !forceCold && options.sessionId ? this.completed.get(options.sessionId) : undefined; if (memory?.position) candidates.push(memory); dbg?.("coordinator.start.memoryCandidate", { found: Boolean(memory) }); const checkpointRoot = this.checkpointRoot(options.sessionId); if (!forceCold && options.sessionId && checkpointRoot) { try { for (const checkpoint of await loadCheckpoints( checkpointRoot, options.sessionId, )) { if (checkpoint.position) candidates.push({ ...checkpoint, store }); } } catch (error) { dbg?.("coordinator.start.checkpointLoadError", { kind: diagnosticKind(error), }); candidates.length = memory?.position ? 1 : 0; } } dbg?.("coordinator.start.candidates", { count: candidates.length }); const prefixFingerprints = new Map(); let completed: CompletedQuery | undefined; for (const checkpoint of candidates) { if (checkpoint.messageCount >= context.messages.length) continue; let fingerprint = prefixFingerprints.get(checkpoint.messageCount); if (!fingerprint) { fingerprint = contextFingerprint( { systemPrompt: context.systemPrompt, messages: context.messages.slice(0, checkpoint.messageCount), tools: context.tools, }, identity, ); prefixFingerprints.set(checkpoint.messageCount, fingerprint); } if ( checkpoint.fingerprint === fingerprint && (!completed || checkpoint.messageCount > completed.messageCount) ) { completed = checkpoint; } } dbg?.("coordinator.start.completed", { found: Boolean(completed), messageCount: completed?.messageCount, }); const offered = reusable && !forceCold && reusable.handle?.isIdle() === true ? this.reuseFit(reusable, context, model, options, tools, identity) : undefined; dbg?.("coordinator.start.reuseFit"); const reuse = offered ? reusable : undefined; if (reusable && !reuse) { await this.closeState( reusable, "Klaus cannot reuse the idle child for this turn.", ); } const reuseBoundary = reuse?.messageCount; const deltaContext: KlausContext = { systemPrompt: context.systemPrompt, messages: reuseBoundary !== undefined ? context.messages.slice(reuseBoundary) : completed ? context.messages.slice(completed.messageCount) : context.messages, tools: context.tools, }; const linearResume = (reuse !== undefined || completed !== undefined) && deltaContext.messages.length === 1 && isTextUser(deltaContext.messages[0]); const canResume = reuse === undefined && completed !== undefined; dbg?.("coordinator.start.resumePlan", { linearResume }); const staging = new StagingSessionStore(completed?.store ?? store); const request: KlausRequest = { modelId: model.id, selector: klausSelector(model.id), systemPrompt: context.systemPrompt ?? "", prompt: linearResume ? finalUserPrompt(deltaContext) : replayPrompt(deltaContext), images: replayImages(deltaContext), tools, thinking: options.reasoning, thinkingBudget: thinkingBudget(options), maxTokens: options.maxTokens, cwd, sessionId: options.sessionId, headers: options.headers ?? {}, env: options.env ?? {}, timeoutMs: options.timeoutMs, sessionStore: staging, resume: canResume ? completed?.sdkSessionId : undefined, resumeAt: canResume ? completed?.position : undefined, forkSession: canResume, }; dbg?.("coordinator.start.request", { imageCount: request.images?.length ?? 0, }); const payloadRequest: KlausRequest = { ...request }; delete payloadRequest.sessionStore; delete payloadRequest.resume; delete payloadRequest.resumeAt; delete payloadRequest.forkSession; dbg?.("coordinator.start.payloadHook", { enabled: options.onPayload !== undefined, }); const hookPayload = options.onPayload ? await options.onPayload(payloadRequest, model) : undefined; const transformedPayload = hookPayload === undefined ? payloadRequest : hookPayload; const transformed = isRequest(transformedPayload) ? { ...transformedPayload, sessionStore: request.sessionStore, resume: request.resume, resumeAt: request.resumeAt, forkSession: request.forkSession, } : transformedPayload; dbg?.("coordinator.start.payloadTransformed", { transformed: hookPayload !== undefined, }); if (!isRequest(transformed)) throw new Error("Klaus onPayload returned an invalid request."); if ( transformed.resume && (transformed.modelId !== request.modelId || transformed.selector !== request.selector || transformed.systemPrompt !== request.systemPrompt || transformed.cwd !== request.cwd || JSON.stringify(transformed.tools) !== JSON.stringify(request.tools)) ) { await this.start(sink, model, context, options, tools, true); return; } const scope = this.currentSession(); dbg?.("coordinator.start.scope", { persisted: scope?.persisted }); if ( reuse && (transformed.prompt !== request.prompt || transformed.systemPrompt !== request.systemPrompt || transformed.modelId !== request.modelId || transformed.selector !== request.selector || transformed.cwd !== request.cwd || JSON.stringify(transformed.tools) !== JSON.stringify(request.tools)) ) { dbg?.("coordinator.start.reuseRejectedByHook"); await this.closeState( reuse, "Klaus onPayload changed a reused turn; replaying canonically.", ); await this.start(sink, model, context, options, tools, true); return; } const state: ActiveQuery = reuse ?? { id: Symbol("klaus-query"), sessionId: options.sessionId, sink, toolCallIds: new Set(), tools, closed: false, timeoutMs: transformed.timeoutMs, secondary: Boolean(options.sessionId && options.sessionId !== scope?.id), committable: true, modelId: model.id, cwd, thinking: options.reasoning, thinkingBudget: thinkingBudget(options), maxTokens: options.maxTokens, systemPrompt: context.systemPrompt ?? "", headers: options.headers, env: options.env, emittedContent: false, context, store: completed?.store ?? store, staging, }; if (reuse) { if (state.reuseLease) clearTimeout(state.reuseLease); state.reuseLease = undefined; state.sink = sink; state.context = context; state.tools = tools; state.toolCallIds = new Set(); state.emittedContent = false; state.timeoutMs = transformed.timeoutMs; state.boundaryFingerprint = undefined; state.boundaryMessageCount = undefined; } else { this.addState(state); } dbg?.("coordinator.start.stateAdded"); this.armTimeout(state, transformed.timeoutMs); this.replaceAbort(state, options.signal); if (options.signal?.aborted) { await this.closeState(state, "Pi aborted the Klaus query.", true); return; } const callbacks: KlausQueryCallbacks = { onReady: async (metadata) => { dbg?.("coordinator.callback.ready"); if (state.closed) return; state.metadata = metadata; await options.onResponse?.({ status: 200, headers: metadata }, model); }, onActivity: () => { dbg?.("coordinator.callback.activity", { closed: state.closed }); this.touchState(state, state.timeoutMs); }, onNotice: (message) => { dbg?.("coordinator.callback.notice"); this.warn(message); }, onContent: (event) => { dbg?.("coordinator.callback.content", { index: event.index }); if (state.closed) return; state.emittedContent = true; if (event.type === "tool-start") state.toolCallIds.add(event.id); state.sink?.content(event); }, onToolBoundary: (result) => { dbg?.("coordinator.callback.toolBoundary", { closed: state.closed }); if (state.closed) return; if (!state.sessionId) { state.sink?.error( new Error("Klaus cannot continue a sessionless tool call."), ); void this.closeState( state, "Sessionless tool continuation is unsupported.", ); return; } if (state.timer) clearTimeout(state.timer); state.timer = undefined; state.sink?.toolUse( result.usage, result.responseId, result.responseModel, ); if (state.sink) { state.boundaryMessageCount = state.context.messages.length + 1; state.boundaryFingerprint = contextFingerprint( { systemPrompt: state.context.systemPrompt, messages: [...state.context.messages, state.sink.output], tools: state.context.tools, }, `${state.modelId}\0${state.cwd}`, ); } state.sink = undefined; if (state.secondary) { state.idleLease = setTimeout( () => void this.closeState( state, "Klaus secondary query idle lease expired.", ), 15 * 60_000, ); } }, onResult: (result) => { dbg?.("coordinator.callback.result", { closed: state.closed }); if (state.closed) return; const resultSink = state.sink; void (async () => { if ( state.committable && state.sessionId && result.sdkSessionId && resultSink ) { try { await state.staging.commit(); if (!state.closed && state.committable) { const completed: CompletedQuery = { sdkSessionId: result.sdkSessionId, position: result.position, messageCount: state.context.messages.length + 1, protocol: 1, fingerprint: contextFingerprint( { systemPrompt: state.context.systemPrompt, messages: [...state.context.messages, resultSink.output], tools: state.context.tools, }, `${state.modelId}\0${state.cwd}`, ), store: state.store, }; dbg?.("coordinator.checkpoint.published", { messageCount: completed.messageCount, }); this.completed.set(state.sessionId, completed); state.sdkSessionId = completed.sdkSessionId; state.position = completed.position; state.messageCount = completed.messageCount; state.lineageFingerprint = contextFingerprint( { systemPrompt: state.context.systemPrompt, messages: [...state.context.messages, resultSink.output], }, `${state.modelId}\0${state.cwd}`, ); if (this.checkpointRoot(state.sessionId)) { this.pendingCheckpoints.set(state.sessionId, completed); } } } catch (error) { this.warn( `Klaus could not cache the completed session: ${error instanceof Error ? error.message : String(error)}`, ); } } if (!state.closed) { resultSink?.done( result.usage, result.responseId, result.responseModel, result.stopReason, ); } if (this.keepForReuse(state)) return; await this.closeState(state, "Klaus query completed."); })().catch((error) => { resultSink?.error(error); void this.closeState(state, "Klaus checkpoint failed."); }); }, onError: (error) => { dbg?.("coordinator.callback.error", { kind: diagnosticKind(error), closed: state.closed, }); if (state.closed) return; if ( transformed.resume && !state.emittedContent && state.sink && isReplayableResumeError(error) ) { const retrySink = state.sink; state.sink = undefined; if (state.sessionId) { this.completed.delete(state.sessionId); this.pendingCheckpoints.delete(state.sessionId); } void (async () => { await this.closeState( state, "Klaus resume failed; replaying canonically.", ); await this.start(retrySink, model, context, options, tools, true); })().catch((retryError) => retrySink.error(retryError)); return; } state.sink?.error(error); void this.closeState(state, error.message); }, }; if (reuse) { try { if (offered?.fit === "swap") { dbg?.("coordinator.reuse.swapTools"); await state.handle?.setTools(tools); } if (state.metadata) { await options.onResponse?.( { status: 200, headers: state.metadata }, model, ); } if (state.closed) return; dbg?.("coordinator.reuse.send"); state.handle?.send( { prompt: transformed.prompt, images: transformed.images }, callbacks, ); } catch (error) { dbg?.("coordinator.reuse.error", { kind: diagnosticKind(error) }); state.sink = undefined; await this.closeState(state, "Klaus could not reuse the idle child."); await this.start(sink, model, context, options, tools, forceCold); } return; } try { dbg?.("coordinator.oauth.resolve.start"); const oauth = await this.resolveOAuth(model.id); dbg?.("coordinator.oauth.resolve.end"); if (state.closed) return; dbg?.("coordinator.sdk.start"); const outgoing = rewriteSubscriptionPrompt(transformed, this.warn); // Claude Agent SDK costs about 200 ms to parse and is needed only once // Klaus handles a request, not while every Pi process starts. const { startSdkQuery } = await import("./agent-sdk.js"); const handle = await startSdkQuery(outgoing, oauth, callbacks); dbg?.("coordinator.sdk.started"); if (state.closed) { await handle.close("Klaus query closed during startup."); return; } state.handle = handle; dbg?.("coordinator.handle.attached"); } catch (error) { dbg?.("coordinator.start.error", { kind: diagnosticKind(error) }); state.sink?.error(error, options.signal?.aborted); await this.closeState(state, "Klaus query startup failed."); } } /** Decides whether an idle child can take this turn, and whether it needs a * tool swap first. Lineage is compared without tools so a Pi tool-list change * alone does not force a canonical replay. */ private reuseFit( state: ActiveQuery, context: KlausContext, model: Model, options: SimpleStreamOptions, tools: KlausRequest["tools"], identity: string, ): { fit: "reuse" | "swap" } | undefined { const boundary = state.messageCount; if (boundary === undefined || state.closed) return undefined; const delta = context.messages.slice(boundary); if (delta.length !== 1 || !isTextUser(delta[0])) return undefined; if ( state.modelId !== model.id || state.cwd !== this.currentCwd() || state.systemPrompt !== (context.systemPrompt ?? "") || state.thinking !== options.reasoning || state.thinkingBudget !== thinkingBudget(options) || state.maxTokens !== options.maxTokens || !sameStringRecord(state.headers, options.headers) || !sameStringRecord(state.env, options.env) ) { return undefined; } const lineage = contextFingerprint( { systemPrompt: context.systemPrompt, messages: context.messages.slice(0, boundary), }, identity, ); if (state.lineageFingerprint !== lineage) return undefined; return { fit: JSON.stringify(state.tools) === JSON.stringify(tools) ? "reuse" : "swap", }; } /** A finished child keeps running for a short window so the next turn in the * same session skips the child spawn and initialize round trip. */ private keepForReuse(state: ActiveQuery): boolean { const reusable = !state.closed && state.committable && Boolean(state.sessionId) && Boolean(state.sdkSessionId) && Boolean(state.position) && state.handle !== undefined; dbg?.("coordinator.keepForReuse"); if (!reusable) return false; state.sink = undefined; if (state.timer) clearTimeout(state.timer); state.timer = undefined; if (state.reuseLease) clearTimeout(state.reuseLease); state.reuseLease = setTimeout( () => void this.closeState(state, "Klaus idle child lease expired."), reuseLeaseMs(), ); return true; } private findReusable(sessionId: string | undefined): ActiveQuery | undefined { if (!sessionId) return undefined; const found = [...(this.bySession.get(sessionId) ?? [])].find( (state) => !state.closed && state.reuseLease !== undefined, ); dbg?.("coordinator.findReusable", { found: Boolean(found) }); return found; } private findContinuation( sessionId: string | undefined, results: PiToolResult[], ): ActiveQuery | undefined { dbg?.("coordinator.findContinuation"); if (!sessionId || results.length === 0) return undefined; return [...(this.bySession.get(sessionId) ?? [])].find( (state) => !state.closed && /** An idle child between turns holds no parked call, so its tool ids * from the finished turn must never capture a continuation. */ state.reuseLease === undefined && results.every((result) => state.toolCallIds.has(result.id)), ); } private addState(state: ActiveQuery): void { dbg?.("coordinator.addState"); if (!state.sessionId) { this.anonymous.add(state); return; } const states = this.bySession.get(state.sessionId) ?? new Set(); if ([...states].some((existing) => !existing.closed)) { state.committable = false; for (const existing of states) existing.committable = false; } states.add(state); this.bySession.set(state.sessionId, states); } private replaceAbort( state: ActiveQuery, signal: AbortSignal | undefined, ): void { dbg?.("coordinator.replaceAbort", { hadSignal: Boolean(state.signal), aborted: signal?.aborted, }); if (state.signal && state.abort) { state.signal.removeEventListener("abort", state.abort); } state.signal = signal; state.abort = signal ? () => void this.closeState(state, "Pi aborted the Klaus query.", true) : undefined; if (signal && state.abort) { signal.addEventListener("abort", state.abort, { once: true }); } } private touchState(state: ActiveQuery, timeoutMs: number | undefined): void { dbg?.("coordinator.touchState", { timeoutMs }); this.armTimeout(state, timeoutMs); if (!state.idleLease) return; clearTimeout(state.idleLease); state.idleLease = setTimeout( () => void this.closeState( state, "Klaus secondary query idle lease expired.", ), 15 * 60_000, ); } private checkpointRoot(sessionId: string | undefined): string | undefined { const scope = this.currentSession(); dbg?.("coordinator.checkpointRoot", { persisted: scope?.persisted }); return sessionId && scope?.id === sessionId && scope.persisted ? join(scope.dir, ".klaus", SDK_VERSION, "checkpoints") : undefined; } private storeFor(sessionId: string | undefined): KlausSessionStore { dbg?.("coordinator.storeFor"); if (!sessionId) return new MemorySessionStore(); const scope = this.currentSession(); if (scope?.id === sessionId && scope.persisted) { const existing = this.diskStores.get(sessionId); if (existing) return existing; const store = new FileSessionStore( join(scope.dir, ".klaus", SDK_VERSION, "sessions"), ); this.diskStores.set(sessionId, store); return store; } const existing = this.memoryStores.get(sessionId); if (existing) return existing; const store = new MemorySessionStore(); this.memoryStores.set(sessionId, store); return store; } private armTimeout(state: ActiveQuery, timeoutMs: number | undefined): void { dbg?.("coordinator.armTimeout", { timeoutMs, closed: state.closed }); if (state.timer) clearTimeout(state.timer); if (!timeoutMs || state.closed || !state.sink) return; state.timer = setTimeout(() => { state.sink?.error( new Error("Klaus timed out waiting for Claude activity."), ); void this.closeState(state, "Klaus activity timeout."); }, timeoutMs); } private async closeState( state: ActiveQuery, reason: string, aborted = false, ): Promise { dbg?.("coordinator.closeState.start", { aborted, hasHandle: Boolean(state.handle), }); if (state.closePromise) return state.closePromise; state.closed = true; state.closePromise = (async () => { if (state.timer) clearTimeout(state.timer); if (state.idleLease) clearTimeout(state.idleLease); if (state.reuseLease) clearTimeout(state.reuseLease); if (state.signal && state.abort) { state.signal.removeEventListener("abort", state.abort); } state.sink?.error(new Error(reason), aborted); if (state.handle) await state.handle.close(reason); if (state.sessionId) { const states = this.bySession.get(state.sessionId); states?.delete(state); if (states?.size === 0) this.bySession.delete(state.sessionId); } else { this.anonymous.delete(state); } dbg?.("coordinator.closeState.end"); })(); return state.closePromise; } } function sameStringRecord( left: Record | undefined, right: Record | undefined, ): boolean { const entries = (value: Record | undefined) => Object.entries(value ?? {}).sort(([leftKey], [rightKey]) => leftKey < rightKey ? -1 : leftKey > rightKey ? 1 : 0, ); return JSON.stringify(entries(left)) === JSON.stringify(entries(right)); } export function isReplayableResumeError(error: Error): boolean { return ( error instanceof KlausCacheCorruptionError || /^Klaus cache corruption:/i.test(error.message) || /^Resume rejected by --resume-drops-turn:/i.test(error.message) || /^No message found with message\.uuid of:/i.test(error.message) ); } function thinkingBudget(options: SimpleStreamOptions): number | undefined { dbg?.("coordinator.thinkingBudget"); switch (options.reasoning) { case "minimal": return options.thinkingBudgets?.minimal; case "low": return options.thinkingBudgets?.low; case "medium": return options.thinkingBudgets?.medium; case "high": case "xhigh": case "max": return options.thinkingBudgets?.high; default: return undefined; } } function isTextUser( message: KlausContext["messages"][number] | undefined, ): boolean { return ( message?.role === "user" && (typeof message.content === "string" || message.content.every((item) => item.type === "text")) ); } function finalUserPrompt(context: KlausContext): string { dbg?.("coordinator.finalUserPrompt", { messageCount: context.messages.length, }); const message = context.messages.at(-1); if (message?.role !== "user") return replayPrompt(context); return typeof message.content === "string" ? message.content : message.content .filter((item) => item.type === "text") .map((item) => item.text) .join("\n"); } function recentToolResults(context: KlausContext): { results: PiToolResult[]; userAfter: boolean; } { dbg?.("coordinator.recentToolResults.start", { messageCount: context.messages.length, }); const results: PiToolResult[] = []; let userAfter = false; for (let index = context.messages.length - 1; index >= 0; index -= 1) { const message = context.messages[index]; if (message.role === "assistant") break; if (message.role === "user") { userAfter = true; continue; } results.unshift({ id: message.toolCallId, content: modelFacingContent(message.content), isError: message.isError, }); } dbg?.("coordinator.recentToolResults.end"); return { results, userAfter }; } function isRequest(value: unknown): value is KlausRequest { dbg?.("coordinator.isRequest"); if (typeof value !== "object" || value === null || Array.isArray(value)) return false; const request = value as Partial; return ( typeof request.modelId === "string" && isKlausModelId(request.modelId) && typeof request.selector === "string" && typeof request.systemPrompt === "string" && typeof request.prompt === "string" && (request.images === undefined || (Array.isArray(request.images) && request.images.every( (image) => typeof image === "object" && image !== null && typeof image.data === "string" && typeof image.mimeType === "string", ))) && (request.thinking === undefined || ["minimal", "low", "medium", "high", "xhigh", "max"].includes( request.thinking, )) && (request.thinkingBudget === undefined || (Number.isInteger(request.thinkingBudget) && request.thinkingBudget > 0)) && (request.maxTokens === undefined || (Number.isSafeInteger(request.maxTokens) && request.maxTokens > 0)) && (request.timeoutMs === undefined || (Number.isFinite(request.timeoutMs) && request.timeoutMs >= 0)) && (request.sessionId === undefined || typeof request.sessionId === "string") && Array.isArray(request.tools) && request.tools.every( (tool) => typeof tool === "object" && tool !== null && typeof tool.name === "string" && typeof tool.description === "string" && typeof tool.inputSchema === "object" && tool.inputSchema !== null && !Array.isArray(tool.inputSchema), ) && typeof request.cwd === "string" && typeof request.headers === "object" && request.headers !== null && Object.values(request.headers).every( (value) => typeof value === "string" || value === null, ) && (request.env === undefined || (typeof request.env === "object" && request.env !== null && Object.values(request.env).every((value) => typeof value === "string"))) ); }