spb/khaelor Public
KHAELOR — a terminal-native autonomous engineering agent powered by Anthropic.
TypeScript 82.9%
HTML 14.9%
CSS 1.1%
JavaScript 0.7%
1/**2 * KHAELOR3 * File: src/session/projections.ts4 * Description: Pure folds over the durable event stream — conversation, usage totals, file-change set, pairing safety.5 *6 * Author: Simon-Pierre Boucher7 * Contact: contact@spboucher.ai8 */910import type {11 DiffStats,12 DurableEvent,13 GitBaseline,14 ModelUsage,15} from "./events.js";1617// ───────────────────────── conversation (LlmHistory-shaped) ─────────────────────────1819export type ContentBlock =20 | { type: "text"; text: string }21 | { type: "thinking"; thinking: string; signature: string }22 | { type: "tool_use"; id: string; name: string; input: unknown }23 | { type: "tool_result"; tool_use_id: string; content: string; is_error?: boolean };2425export interface ConversationMessage {26 role: "user" | "assistant";27 content: ContentBlock[];28}2930/** Fixed template wrapper for compaction checkpoints (EVENT_MODEL.md §6.2.6). */31export const COMPACTION_MESSAGE_PREFIX =32 "[Context checkpoint — earlier conversation was summarized. Structured checkpoint follows.]\n\n";3334interface Entry {35 role: "user" | "assistant";36 blocks: ContentBlock[];37 minSeq: number;38 maxSeq: number;39 /** Open tool-result user messages accept further tool_result blocks. */40 openToolResults: boolean;41}4243interface PendingAssistant {44 requestId: string;45 blocks: { blockIndex: number; block: ContentBlock }[];46 minSeq: number;47 maxSeq: number;48}4950/**51 * Build the conversation state (Anthropic-`messages[]`-shaped) from the durable52 * stream, in `seq` order. Deterministic by construction: `ContextPruned` and53 * `ContextCompacted` are re-applied from their recorded payloads on every54 * rebuild — same events, same bytes (EVENT_MODEL.md §6.2).55 */56export function buildConversation(events: readonly DurableEvent[]): ConversationMessage[] {57 const entries: Entry[] = [];58 const steeringTexts = new Map<string, string>(); // SteeringQueued id → text59 let pending: PendingAssistant | null = null;6061 const flushAssistant = (extendToSeq?: number): void => {62 if (!pending) return;63 const sorted = [...pending.blocks].sort((a, b) => a.blockIndex - b.blockIndex);64 entries.push({65 role: "assistant",66 blocks: sorted.map((b) => b.block),67 minSeq: pending.minSeq,68 maxSeq: extendToSeq !== undefined ? Math.max(pending.maxSeq, extendToSeq) : pending.maxSeq,69 openToolResults: false,70 });71 pending = null;72 };7374 const closeToolResults = (): void => {75 const last = entries[entries.length - 1];76 if (last) last.openToolResults = false;77 };7879 const addAssistantBlock = (80 requestId: string,81 blockIndex: number,82 block: ContentBlock,83 seq: number,84 ): void => {85 if (pending && pending.requestId !== requestId) flushAssistant();86 closeToolResults();87 if (!pending) {88 pending = { requestId, blocks: [], minSeq: seq, maxSeq: seq };89 }90 pending.blocks.push({ blockIndex, block });91 pending.minSeq = Math.min(pending.minSeq, seq);92 pending.maxSeq = Math.max(pending.maxSeq, seq);93 };9495 const addToolResult = (block: ContentBlock, seq: number): void => {96 flushAssistant();97 const last = entries[entries.length - 1];98 if (last && last.role === "user" && last.openToolResults) {99 last.blocks.push(block);100 last.maxSeq = Math.max(last.maxSeq, seq);101 return;102 }103 entries.push({ role: "user", blocks: [block], minSeq: seq, maxSeq: seq, openToolResults: true });104 };105106 const applyPrune = (toolUseIds: readonly string[], placeholder: string): void => {107 const ids = new Set(toolUseIds);108 for (const entry of entries) {109 for (let i = 0; i < entry.blocks.length; i++) {110 const block = entry.blocks[i] as ContentBlock;111 if (block.type === "tool_result" && ids.has(block.tool_use_id)) {112 entry.blocks[i] = { ...block, content: placeholder };113 }114 }115 }116 };117118 const applyCompaction = (119 checkpointYaml: string,120 cut: { fromSeq: number; toSeq: number },121 ): void => {122 const synthetic: Entry = {123 role: "user",124 blocks: [{ type: "text", text: COMPACTION_MESSAGE_PREFIX + checkpointYaml }],125 minSeq: cut.fromSeq,126 maxSeq: cut.toSeq,127 openToolResults: false,128 };129 let insertAt = -1;130 const kept: Entry[] = [];131 for (const entry of entries) {132 if (entry.minSeq >= cut.fromSeq && entry.maxSeq <= cut.toSeq) {133 if (insertAt === -1) insertAt = kept.length;134 continue; // consumed by the checkpoint (may include earlier checkpoints)135 }136 kept.push(entry);137 }138 if (insertAt === -1) {139 // Nothing matched (empty cut): insert after the last entry that ends before the cut.140 insertAt = kept.findIndex((e) => e.minSeq > cut.toSeq);141 if (insertAt === -1) insertAt = kept.length;142 }143 kept.splice(insertAt, 0, synthetic);144 entries.length = 0;145 entries.push(...kept);146 };147148 for (const event of events) {149 switch (event.type) {150 case "user.message-created": {151 flushAssistant();152 closeToolResults();153 entries.push({154 role: "user",155 blocks: [{ type: "text", text: event.payload.text }],156 minSeq: event.seq,157 maxSeq: event.seq,158 openToolResults: false,159 });160 break;161 }162 case "user.steering-queued": {163 steeringTexts.set(event.id, event.payload.text);164 break;165 }166 case "user.steering-injected": {167 const text = steeringTexts.get(event.payload.queuedEventId);168 if (text === undefined) break; // unknown reference — skip, never fabricate169 const block: ContentBlock = { type: "text", text };170 let target: Entry | undefined;171 if (event.payload.seam === "post-tool-batch") {172 target = entries.find(173 (e) =>174 e.role === "user" &&175 e.minSeq <= event.payload.afterSeq &&176 event.payload.afterSeq <= e.maxSeq,177 );178 }179 if (!target) {180 for (let i = entries.length - 1; i >= 0; i--) {181 const candidate = entries[i] as Entry;182 if (candidate.role === "user") {183 target = candidate;184 break;185 }186 }187 }188 if (target) {189 target.blocks.push(block);190 target.maxSeq = Math.max(target.maxSeq, event.seq);191 } else {192 entries.push({193 role: "user",194 blocks: [block],195 minSeq: event.seq,196 maxSeq: event.seq,197 openToolResults: false,198 });199 }200 break;201 }202 case "model.text-block-completed": {203 addAssistantBlock(204 event.payload.requestId,205 event.payload.blockIndex,206 { type: "text", text: event.payload.text },207 event.seq,208 );209 break;210 }211 case "model.thinking-block-completed": {212 addAssistantBlock(213 event.payload.requestId,214 event.payload.blockIndex,215 {216 type: "thinking",217 thinking: event.payload.thinking,218 signature: event.payload.signature,219 },220 event.seq,221 );222 break;223 }224 case "tool.requested": {225 addAssistantBlock(226 event.payload.requestId,227 event.payload.blockIndex,228 {229 type: "tool_use",230 id: event.payload.toolUseId,231 name: event.payload.toolName,232 input: event.payload.input,233 },234 event.seq,235 );236 break;237 }238 case "model.response-completed": {239 if (pending && (pending as PendingAssistant).requestId === event.payload.requestId) {240 flushAssistant(event.seq);241 }242 break;243 }244 case "model.request-failed": {245 // Settled blocks of a failed request remain durable history.246 flushAssistant(event.seq);247 break;248 }249 case "tool.completed": {250 addToolResult(251 { type: "tool_result", tool_use_id: event.payload.toolUseId, content: event.payload.modelText },252 event.seq,253 );254 break;255 }256 case "tool.failed": {257 addToolResult(258 {259 type: "tool_result",260 tool_use_id: event.payload.toolUseId,261 content: event.payload.modelText,262 is_error: true,263 },264 event.seq,265 );266 break;267 }268 case "tool.cancelled": {269 addToolResult(270 { type: "tool_result", tool_use_id: event.payload.toolUseId, content: event.payload.modelText },271 event.seq,272 );273 break;274 }275 case "context.pruned": {276 flushAssistant();277 applyPrune(event.payload.toolUseIds, event.payload.placeholder);278 break;279 }280 case "context.compacted": {281 flushAssistant();282 closeToolResults();283 applyCompaction(event.payload.checkpointYaml, event.payload.cut);284 break;285 }286 default:287 break; // non-conversation events288 }289 }290 flushAssistant();291292 return entries293 .filter((e) => e.blocks.length > 0)294 .map((e) => ({ role: e.role, content: e.blocks }));295}296297// ───────────────────────── usage totals (real usage only) ─────────────────────────298299export interface ModelUsageTotals extends ModelUsage {300 requests: number;301}302303export interface UsageTotals {304 perModel: Record<string, ModelUsageTotals>;305 totals: ModelUsageTotals;306}307308const EMPTY_USAGE: ModelUsageTotals = {309 inputTokens: 0,310 outputTokens: 0,311 cacheReadTokens: 0,312 cacheWriteTokens: 0,313 requests: 0,314};315316function addUsage(target: ModelUsageTotals, usage: ModelUsage): ModelUsageTotals {317 return {318 inputTokens: target.inputTokens + usage.inputTokens,319 outputTokens: target.outputTokens + usage.outputTokens,320 cacheReadTokens: target.cacheReadTokens + usage.cacheReadTokens,321 cacheWriteTokens: target.cacheWriteTokens + usage.cacheWriteTokens,322 requests: target.requests + 1,323 };324}325326/**327 * Sum real API usage from `ModelResponseCompleted` events only, keyed by the328 * model id recorded by the correlated `ModelRequestStarted` (Absolute Rule #4:329 * no other event may contribute to cost).330 */331export function buildUsageTotals(events: readonly DurableEvent[]): UsageTotals {332 const requestModel = new Map<string, string>();333 const perModel: Record<string, ModelUsageTotals> = {};334 let totals: ModelUsageTotals = { ...EMPTY_USAGE };335336 for (const event of events) {337 if (event.type === "model.request-started") {338 requestModel.set(event.payload.requestId, event.payload.model);339 } else if (event.type === "model.response-completed") {340 const model = requestModel.get(event.payload.requestId) ?? "unknown";341 perModel[model] = addUsage(perModel[model] ?? { ...EMPTY_USAGE }, event.payload.usage);342 totals = addUsage(totals, event.payload.usage);343 }344 }345 return { perModel, totals };346}347348// ───────────────────────── file-change set ─────────────────────────349350export interface FileChange {351 path: string;352 operations: ("write" | "edit")[];353 diffStats: DiffStats; // cumulative354 toolUseIds: string[];355}356357export interface FileReadRecord {358 mtimeMs: number;359 bytes: number;360 toolUseId: string;361}362363export interface FileChangeSet {364 /** Latest recorded baseline (pre-first-edit wins over session-start when both exist). */365 baseline: GitBaseline | null;366 baselineWhen: "session-start" | "pre-first-edit" | null;367 changes: Map<string, FileChange>;368 /** Last read record per path — external-modification detection. */369 reads: Map<string, FileReadRecord>;370}371372/**373 * Fold `BaselineRecorded` / `FileModified` / `FileRead` into the attribution374 * and diff state (EVENT_MODEL.md §6.4). Never erased by compaction — file375 * changes remain first-class even when their conversation is summarized.376 */377export function buildFileChangeSet(events: readonly DurableEvent[]): FileChangeSet {378 const set: FileChangeSet = {379 baseline: null,380 baselineWhen: null,381 changes: new Map(),382 reads: new Map(),383 };384 for (const event of events) {385 switch (event.type) {386 case "git.baseline-recorded": {387 // pre-first-edit refines session-start; a later baseline of the same kind replaces it.388 if (set.baselineWhen !== "pre-first-edit" || event.payload.when === "pre-first-edit") {389 set.baseline = event.payload.baseline;390 set.baselineWhen = event.payload.when;391 }392 break;393 }394 case "file.modified": {395 const existing = set.changes.get(event.payload.path);396 if (existing) {397 existing.operations.push(event.payload.operation);398 existing.diffStats = {399 added: existing.diffStats.added + event.payload.diffStats.added,400 removed: existing.diffStats.removed + event.payload.diffStats.removed,401 };402 existing.toolUseIds.push(event.payload.toolUseId);403 } else {404 set.changes.set(event.payload.path, {405 path: event.payload.path,406 operations: [event.payload.operation],407 diffStats: { ...event.payload.diffStats },408 toolUseIds: [event.payload.toolUseId],409 });410 }411 break;412 }413 case "file.read": {414 set.reads.set(event.payload.path, {415 mtimeMs: event.payload.mtimeMs,416 bytes: event.payload.bytes,417 toolUseId: event.payload.toolUseId,418 });419 break;420 }421 default:422 break;423 }424 }425 return set;426}427428// ───────────────────────── pairing safety (tool_use / tool_result) ─────────────────────────429430/**431 * ToolRequested events lacking a terminal `ToolCompleted`/`ToolFailed`/`ToolCancelled`.432 * Resume recovery closes each with a synthetic `ToolCancelled{reason:"resume-recovery"}`433 * (EVENT_MODEL.md §6.5.2).434 */435export function findDanglingToolUseIds(events: readonly DurableEvent[]): string[] {436 const open = new Map<string, true>();437 for (const event of events) {438 if (event.type === "tool.requested") {439 open.set(event.payload.toolUseId, true);440 } else if (441 event.type === "tool.completed" ||442 event.type === "tool.failed" ||443 event.type === "tool.cancelled"444 ) {445 open.delete(event.payload.toolUseId);446 }447 }448 return [...open.keys()];449}450451/**452 * A compaction cut is pairing-safe iff every `tool_use` at seq ≤ toSeq has its453 * result at seq ≤ toSeq, and no `tool_use` kept outside the cut loses its454 * result to the cut (EVENT_MODEL.md §6.5.3).455 */456export function isPairingSafeCut(457 events: readonly DurableEvent[],458 cut: { fromSeq: number; toSeq: number },459): boolean {460 if (cut.fromSeq > cut.toSeq) return false;461 const requested = new Map<string, number>();462 const closed = new Map<string, number>();463 for (const event of events) {464 if (event.type === "tool.requested") {465 requested.set(event.payload.toolUseId, event.seq);466 } else if (467 event.type === "tool.completed" ||468 event.type === "tool.failed" ||469 event.type === "tool.cancelled"470 ) {471 closed.set(event.payload.toolUseId, event.seq);472 }473 }474 for (const [toolUseId, reqSeq] of requested) {475 const closeSeq = closed.get(toolUseId);476 if (reqSeq <= cut.toSeq && (closeSeq === undefined || closeSeq > cut.toSeq)) return false;477 if (478 reqSeq < cut.fromSeq &&479 closeSeq !== undefined &&480 closeSeq >= cut.fromSeq &&481 closeSeq <= cut.toSeq482 ) {483 return false;484 }485 }486 return true;487}488