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 "verify.result": {276 // Failing native checks enter the conversation as user-visible input;277 // passing checks stay out of context (evidence only, v2 §4).278 if (event.payload.ok) break;279 flushAssistant();280 closeToolResults();281 const exit = event.payload.exitCode === null ? "timeout/kill" : `exit ${event.payload.exitCode}`;282 entries.push({283 role: "user",284 blocks: [285 {286 type: "text",287 text:288 `[verify] check "${event.payload.check}" failed (${exit}): ${event.payload.command}\n` +289 event.payload.output,290 },291 ],292 minSeq: event.seq,293 maxSeq: event.seq,294 openToolResults: false,295 });296 break;297 }298 case "context.pruned": {299 flushAssistant();300 applyPrune(event.payload.toolUseIds, event.payload.placeholder);301 break;302 }303 case "context.compacted": {304 flushAssistant();305 closeToolResults();306 applyCompaction(event.payload.checkpointYaml, event.payload.cut);307 break;308 }309 default:310 break; // non-conversation events311 }312 }313 flushAssistant();314315 return entries316 .filter((e) => e.blocks.length > 0)317 .map((e) => ({ role: e.role, content: e.blocks }));318}319320// ───────────────────────── usage totals (real usage only) ─────────────────────────321322export interface ModelUsageTotals extends ModelUsage {323 requests: number;324}325326export interface UsageTotals {327 perModel: Record<string, ModelUsageTotals>;328 totals: ModelUsageTotals;329}330331const EMPTY_USAGE: ModelUsageTotals = {332 inputTokens: 0,333 outputTokens: 0,334 cacheReadTokens: 0,335 cacheWriteTokens: 0,336 requests: 0,337};338339function addUsage(target: ModelUsageTotals, usage: ModelUsage): ModelUsageTotals {340 return {341 inputTokens: target.inputTokens + usage.inputTokens,342 outputTokens: target.outputTokens + usage.outputTokens,343 cacheReadTokens: target.cacheReadTokens + usage.cacheReadTokens,344 cacheWriteTokens: target.cacheWriteTokens + usage.cacheWriteTokens,345 requests: target.requests + 1,346 };347}348349/**350 * Sum real API usage from `ModelResponseCompleted` events only, keyed by the351 * model id recorded by the correlated `ModelRequestStarted` (Absolute Rule #4:352 * no other event may contribute to cost).353 */354export function buildUsageTotals(events: readonly DurableEvent[]): UsageTotals {355 const requestModel = new Map<string, string>();356 const perModel: Record<string, ModelUsageTotals> = {};357 let totals: ModelUsageTotals = { ...EMPTY_USAGE };358359 for (const event of events) {360 if (event.type === "model.request-started") {361 requestModel.set(event.payload.requestId, event.payload.model);362 } else if (event.type === "model.response-completed") {363 const model = requestModel.get(event.payload.requestId) ?? "unknown";364 perModel[model] = addUsage(perModel[model] ?? { ...EMPTY_USAGE }, event.payload.usage);365 totals = addUsage(totals, event.payload.usage);366 }367 }368 return { perModel, totals };369}370371// ───────────────────────── file-change set ─────────────────────────372373export interface FileChange {374 path: string;375 operations: ("write" | "edit")[];376 diffStats: DiffStats; // cumulative377 toolUseIds: string[];378}379380export interface FileReadRecord {381 mtimeMs: number;382 bytes: number;383 toolUseId: string;384}385386export interface FileChangeSet {387 /** Latest recorded baseline (pre-first-edit wins over session-start when both exist). */388 baseline: GitBaseline | null;389 baselineWhen: "session-start" | "pre-first-edit" | null;390 changes: Map<string, FileChange>;391 /** Last read record per path — external-modification detection. */392 reads: Map<string, FileReadRecord>;393}394395/**396 * Fold `BaselineRecorded` / `FileModified` / `FileRead` into the attribution397 * and diff state (EVENT_MODEL.md §6.4). Never erased by compaction — file398 * changes remain first-class even when their conversation is summarized.399 */400export function buildFileChangeSet(events: readonly DurableEvent[]): FileChangeSet {401 const set: FileChangeSet = {402 baseline: null,403 baselineWhen: null,404 changes: new Map(),405 reads: new Map(),406 };407 for (const event of events) {408 switch (event.type) {409 case "git.baseline-recorded": {410 // pre-first-edit refines session-start; a later baseline of the same kind replaces it.411 if (set.baselineWhen !== "pre-first-edit" || event.payload.when === "pre-first-edit") {412 set.baseline = event.payload.baseline;413 set.baselineWhen = event.payload.when;414 }415 break;416 }417 case "file.modified": {418 const existing = set.changes.get(event.payload.path);419 if (existing) {420 existing.operations.push(event.payload.operation);421 existing.diffStats = {422 added: existing.diffStats.added + event.payload.diffStats.added,423 removed: existing.diffStats.removed + event.payload.diffStats.removed,424 };425 existing.toolUseIds.push(event.payload.toolUseId);426 } else {427 set.changes.set(event.payload.path, {428 path: event.payload.path,429 operations: [event.payload.operation],430 diffStats: { ...event.payload.diffStats },431 toolUseIds: [event.payload.toolUseId],432 });433 }434 break;435 }436 case "file.read": {437 set.reads.set(event.payload.path, {438 mtimeMs: event.payload.mtimeMs,439 bytes: event.payload.bytes,440 toolUseId: event.payload.toolUseId,441 });442 break;443 }444 default:445 break;446 }447 }448 return set;449}450451// ───────────────────────── pairing safety (tool_use / tool_result) ─────────────────────────452453/**454 * ToolRequested events lacking a terminal `ToolCompleted`/`ToolFailed`/`ToolCancelled`.455 * Resume recovery closes each with a synthetic `ToolCancelled{reason:"resume-recovery"}`456 * (EVENT_MODEL.md §6.5.2).457 */458export function findDanglingToolUseIds(events: readonly DurableEvent[]): string[] {459 const open = new Map<string, true>();460 for (const event of events) {461 if (event.type === "tool.requested") {462 open.set(event.payload.toolUseId, true);463 } else if (464 event.type === "tool.completed" ||465 event.type === "tool.failed" ||466 event.type === "tool.cancelled"467 ) {468 open.delete(event.payload.toolUseId);469 }470 }471 return [...open.keys()];472}473474/**475 * A compaction cut is pairing-safe iff every `tool_use` at seq ≤ toSeq has its476 * result at seq ≤ toSeq, and no `tool_use` kept outside the cut loses its477 * result to the cut (EVENT_MODEL.md §6.5.3).478 */479export function isPairingSafeCut(480 events: readonly DurableEvent[],481 cut: { fromSeq: number; toSeq: number },482): boolean {483 if (cut.fromSeq > cut.toSeq) return false;484 const requested = new Map<string, number>();485 const closed = new Map<string, number>();486 for (const event of events) {487 if (event.type === "tool.requested") {488 requested.set(event.payload.toolUseId, event.seq);489 } else if (490 event.type === "tool.completed" ||491 event.type === "tool.failed" ||492 event.type === "tool.cancelled"493 ) {494 closed.set(event.payload.toolUseId, event.seq);495 }496 }497 for (const [toolUseId, reqSeq] of requested) {498 const closeSeq = closed.get(toolUseId);499 if (reqSeq <= cut.toSeq && (closeSeq === undefined || closeSeq > cut.toSeq)) return false;500 if (501 reqSeq < cut.fromSeq &&502 closeSeq !== undefined &&503 closeSeq >= cut.fromSeq &&504 closeSeq <= cut.toSeq505 ) {506 return false;507 }508 }509 return true;510}511