SPB Git

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%
16.3 KB · 511 lines typescript
Raw Blame History
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