/** * KHAELOR * File: src/session/projections.ts * Description: Pure folds over the durable event stream — conversation, usage totals, file-change set, pairing safety. * * Author: Simon-Pierre Boucher * Contact: contact@spboucher.ai */ import type { DiffStats, DurableEvent, GitBaseline, ModelUsage, } from "./events.js"; // ───────────────────────── conversation (LlmHistory-shaped) ───────────────────────── export type ContentBlock = | { type: "text"; text: string } | { type: "thinking"; thinking: string; signature: string } | { type: "tool_use"; id: string; name: string; input: unknown } | { type: "tool_result"; tool_use_id: string; content: string; is_error?: boolean }; export interface ConversationMessage { role: "user" | "assistant"; content: ContentBlock[]; } /** Fixed template wrapper for compaction checkpoints (EVENT_MODEL.md §6.2.6). */ export const COMPACTION_MESSAGE_PREFIX = "[Context checkpoint — earlier conversation was summarized. Structured checkpoint follows.]\n\n"; interface Entry { role: "user" | "assistant"; blocks: ContentBlock[]; minSeq: number; maxSeq: number; /** Open tool-result user messages accept further tool_result blocks. */ openToolResults: boolean; } interface PendingAssistant { requestId: string; blocks: { blockIndex: number; block: ContentBlock }[]; minSeq: number; maxSeq: number; } /** * Build the conversation state (Anthropic-`messages[]`-shaped) from the durable * stream, in `seq` order. Deterministic by construction: `ContextPruned` and * `ContextCompacted` are re-applied from their recorded payloads on every * rebuild — same events, same bytes (EVENT_MODEL.md §6.2). */ export function buildConversation(events: readonly DurableEvent[]): ConversationMessage[] { const entries: Entry[] = []; const steeringTexts = new Map(); // SteeringQueued id → text let pending: PendingAssistant | null = null; const flushAssistant = (extendToSeq?: number): void => { if (!pending) return; const sorted = [...pending.blocks].sort((a, b) => a.blockIndex - b.blockIndex); entries.push({ role: "assistant", blocks: sorted.map((b) => b.block), minSeq: pending.minSeq, maxSeq: extendToSeq !== undefined ? Math.max(pending.maxSeq, extendToSeq) : pending.maxSeq, openToolResults: false, }); pending = null; }; const closeToolResults = (): void => { const last = entries[entries.length - 1]; if (last) last.openToolResults = false; }; const addAssistantBlock = ( requestId: string, blockIndex: number, block: ContentBlock, seq: number, ): void => { if (pending && pending.requestId !== requestId) flushAssistant(); closeToolResults(); if (!pending) { pending = { requestId, blocks: [], minSeq: seq, maxSeq: seq }; } pending.blocks.push({ blockIndex, block }); pending.minSeq = Math.min(pending.minSeq, seq); pending.maxSeq = Math.max(pending.maxSeq, seq); }; const addToolResult = (block: ContentBlock, seq: number): void => { flushAssistant(); const last = entries[entries.length - 1]; if (last && last.role === "user" && last.openToolResults) { last.blocks.push(block); last.maxSeq = Math.max(last.maxSeq, seq); return; } entries.push({ role: "user", blocks: [block], minSeq: seq, maxSeq: seq, openToolResults: true }); }; const applyPrune = (toolUseIds: readonly string[], placeholder: string): void => { const ids = new Set(toolUseIds); for (const entry of entries) { for (let i = 0; i < entry.blocks.length; i++) { const block = entry.blocks[i] as ContentBlock; if (block.type === "tool_result" && ids.has(block.tool_use_id)) { entry.blocks[i] = { ...block, content: placeholder }; } } } }; const applyCompaction = ( checkpointYaml: string, cut: { fromSeq: number; toSeq: number }, ): void => { const synthetic: Entry = { role: "user", blocks: [{ type: "text", text: COMPACTION_MESSAGE_PREFIX + checkpointYaml }], minSeq: cut.fromSeq, maxSeq: cut.toSeq, openToolResults: false, }; let insertAt = -1; const kept: Entry[] = []; for (const entry of entries) { if (entry.minSeq >= cut.fromSeq && entry.maxSeq <= cut.toSeq) { if (insertAt === -1) insertAt = kept.length; continue; // consumed by the checkpoint (may include earlier checkpoints) } kept.push(entry); } if (insertAt === -1) { // Nothing matched (empty cut): insert after the last entry that ends before the cut. insertAt = kept.findIndex((e) => e.minSeq > cut.toSeq); if (insertAt === -1) insertAt = kept.length; } kept.splice(insertAt, 0, synthetic); entries.length = 0; entries.push(...kept); }; for (const event of events) { switch (event.type) { case "user.message-created": { flushAssistant(); closeToolResults(); entries.push({ role: "user", blocks: [{ type: "text", text: event.payload.text }], minSeq: event.seq, maxSeq: event.seq, openToolResults: false, }); break; } case "user.steering-queued": { steeringTexts.set(event.id, event.payload.text); break; } case "user.steering-injected": { const text = steeringTexts.get(event.payload.queuedEventId); if (text === undefined) break; // unknown reference — skip, never fabricate const block: ContentBlock = { type: "text", text }; let target: Entry | undefined; if (event.payload.seam === "post-tool-batch") { target = entries.find( (e) => e.role === "user" && e.minSeq <= event.payload.afterSeq && event.payload.afterSeq <= e.maxSeq, ); } if (!target) { for (let i = entries.length - 1; i >= 0; i--) { const candidate = entries[i] as Entry; if (candidate.role === "user") { target = candidate; break; } } } if (target) { target.blocks.push(block); target.maxSeq = Math.max(target.maxSeq, event.seq); } else { entries.push({ role: "user", blocks: [block], minSeq: event.seq, maxSeq: event.seq, openToolResults: false, }); } break; } case "model.text-block-completed": { addAssistantBlock( event.payload.requestId, event.payload.blockIndex, { type: "text", text: event.payload.text }, event.seq, ); break; } case "model.thinking-block-completed": { addAssistantBlock( event.payload.requestId, event.payload.blockIndex, { type: "thinking", thinking: event.payload.thinking, signature: event.payload.signature, }, event.seq, ); break; } case "tool.requested": { addAssistantBlock( event.payload.requestId, event.payload.blockIndex, { type: "tool_use", id: event.payload.toolUseId, name: event.payload.toolName, input: event.payload.input, }, event.seq, ); break; } case "model.response-completed": { if (pending && (pending as PendingAssistant).requestId === event.payload.requestId) { flushAssistant(event.seq); } break; } case "model.request-failed": { // Settled blocks of a failed request remain durable history. flushAssistant(event.seq); break; } case "tool.completed": { addToolResult( { type: "tool_result", tool_use_id: event.payload.toolUseId, content: event.payload.modelText }, event.seq, ); break; } case "tool.failed": { addToolResult( { type: "tool_result", tool_use_id: event.payload.toolUseId, content: event.payload.modelText, is_error: true, }, event.seq, ); break; } case "tool.cancelled": { addToolResult( { type: "tool_result", tool_use_id: event.payload.toolUseId, content: event.payload.modelText }, event.seq, ); break; } case "verify.result": { // Failing native checks enter the conversation as user-visible input; // passing checks stay out of context (evidence only, v2 §4). if (event.payload.ok) break; flushAssistant(); closeToolResults(); const exit = event.payload.exitCode === null ? "timeout/kill" : `exit ${event.payload.exitCode}`; entries.push({ role: "user", blocks: [ { type: "text", text: `[verify] check "${event.payload.check}" failed (${exit}): ${event.payload.command}\n` + event.payload.output, }, ], minSeq: event.seq, maxSeq: event.seq, openToolResults: false, }); break; } case "context.pruned": { flushAssistant(); applyPrune(event.payload.toolUseIds, event.payload.placeholder); break; } case "context.compacted": { flushAssistant(); closeToolResults(); applyCompaction(event.payload.checkpointYaml, event.payload.cut); break; } default: break; // non-conversation events } } flushAssistant(); return entries .filter((e) => e.blocks.length > 0) .map((e) => ({ role: e.role, content: e.blocks })); } // ───────────────────────── usage totals (real usage only) ───────────────────────── export interface ModelUsageTotals extends ModelUsage { requests: number; } export interface UsageTotals { perModel: Record; totals: ModelUsageTotals; } const EMPTY_USAGE: ModelUsageTotals = { inputTokens: 0, outputTokens: 0, cacheReadTokens: 0, cacheWriteTokens: 0, requests: 0, }; function addUsage(target: ModelUsageTotals, usage: ModelUsage): ModelUsageTotals { return { inputTokens: target.inputTokens + usage.inputTokens, outputTokens: target.outputTokens + usage.outputTokens, cacheReadTokens: target.cacheReadTokens + usage.cacheReadTokens, cacheWriteTokens: target.cacheWriteTokens + usage.cacheWriteTokens, requests: target.requests + 1, }; } /** * Sum real API usage from `ModelResponseCompleted` events only, keyed by the * model id recorded by the correlated `ModelRequestStarted` (Absolute Rule #4: * no other event may contribute to cost). */ export function buildUsageTotals(events: readonly DurableEvent[]): UsageTotals { const requestModel = new Map(); const perModel: Record = {}; let totals: ModelUsageTotals = { ...EMPTY_USAGE }; for (const event of events) { if (event.type === "model.request-started") { requestModel.set(event.payload.requestId, event.payload.model); } else if (event.type === "model.response-completed") { const model = requestModel.get(event.payload.requestId) ?? "unknown"; perModel[model] = addUsage(perModel[model] ?? { ...EMPTY_USAGE }, event.payload.usage); totals = addUsage(totals, event.payload.usage); } } return { perModel, totals }; } // ───────────────────────── file-change set ───────────────────────── export interface FileChange { path: string; operations: ("write" | "edit")[]; diffStats: DiffStats; // cumulative toolUseIds: string[]; } export interface FileReadRecord { mtimeMs: number; bytes: number; toolUseId: string; } export interface FileChangeSet { /** Latest recorded baseline (pre-first-edit wins over session-start when both exist). */ baseline: GitBaseline | null; baselineWhen: "session-start" | "pre-first-edit" | null; changes: Map; /** Last read record per path — external-modification detection. */ reads: Map; } /** * Fold `BaselineRecorded` / `FileModified` / `FileRead` into the attribution * and diff state (EVENT_MODEL.md §6.4). Never erased by compaction — file * changes remain first-class even when their conversation is summarized. */ export function buildFileChangeSet(events: readonly DurableEvent[]): FileChangeSet { const set: FileChangeSet = { baseline: null, baselineWhen: null, changes: new Map(), reads: new Map(), }; for (const event of events) { switch (event.type) { case "git.baseline-recorded": { // pre-first-edit refines session-start; a later baseline of the same kind replaces it. if (set.baselineWhen !== "pre-first-edit" || event.payload.when === "pre-first-edit") { set.baseline = event.payload.baseline; set.baselineWhen = event.payload.when; } break; } case "file.modified": { const existing = set.changes.get(event.payload.path); if (existing) { existing.operations.push(event.payload.operation); existing.diffStats = { added: existing.diffStats.added + event.payload.diffStats.added, removed: existing.diffStats.removed + event.payload.diffStats.removed, }; existing.toolUseIds.push(event.payload.toolUseId); } else { set.changes.set(event.payload.path, { path: event.payload.path, operations: [event.payload.operation], diffStats: { ...event.payload.diffStats }, toolUseIds: [event.payload.toolUseId], }); } break; } case "file.read": { set.reads.set(event.payload.path, { mtimeMs: event.payload.mtimeMs, bytes: event.payload.bytes, toolUseId: event.payload.toolUseId, }); break; } default: break; } } return set; } // ───────────────────────── pairing safety (tool_use / tool_result) ───────────────────────── /** * ToolRequested events lacking a terminal `ToolCompleted`/`ToolFailed`/`ToolCancelled`. * Resume recovery closes each with a synthetic `ToolCancelled{reason:"resume-recovery"}` * (EVENT_MODEL.md §6.5.2). */ export function findDanglingToolUseIds(events: readonly DurableEvent[]): string[] { const open = new Map(); for (const event of events) { if (event.type === "tool.requested") { open.set(event.payload.toolUseId, true); } else if ( event.type === "tool.completed" || event.type === "tool.failed" || event.type === "tool.cancelled" ) { open.delete(event.payload.toolUseId); } } return [...open.keys()]; } /** * A compaction cut is pairing-safe iff every `tool_use` at seq ≤ toSeq has its * result at seq ≤ toSeq, and no `tool_use` kept outside the cut loses its * result to the cut (EVENT_MODEL.md §6.5.3). */ export function isPairingSafeCut( events: readonly DurableEvent[], cut: { fromSeq: number; toSeq: number }, ): boolean { if (cut.fromSeq > cut.toSeq) return false; const requested = new Map(); const closed = new Map(); for (const event of events) { if (event.type === "tool.requested") { requested.set(event.payload.toolUseId, event.seq); } else if ( event.type === "tool.completed" || event.type === "tool.failed" || event.type === "tool.cancelled" ) { closed.set(event.payload.toolUseId, event.seq); } } for (const [toolUseId, reqSeq] of requested) { const closeSeq = closed.get(toolUseId); if (reqSeq <= cut.toSeq && (closeSeq === undefined || closeSeq > cut.toSeq)) return false; if ( reqSeq < cut.fromSeq && closeSeq !== undefined && closeSeq >= cut.fromSeq && closeSeq <= cut.toSeq ) { return false; } } return true; }