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%
15.5 KB · 488 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 "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