SLOPSHOPPER

lcm

Lossless context management — DAG-based summarization that preserves every message

newguardpromptmodelprocessnetwork
★ 45v0.16.0MITupdated 2026-10-03lossless-claude/lcm
A shopper browsing a rack in a slop shop
README

<strong>lcm</strong><br> Shared memory infrastructure for coding agents

DAG-based summarization, SQLite-backed message persistence, promoted long-term memory, MCP retrieval tools

<a href="https://www.npmjs.com/package/@lossless-claude/lcm"><img src="https://img.shields.io/npm/v/@lossless-claude/lcm" alt="npm"></a> <a href="LICENSE"><img src="https://img.shields.io/github/license/lossless-claude/lcm" alt="License: MIT"></a> <a href="package.json"><img src="https://img.shields.io/node/v/@lossless-claude/lcm" alt="Node"></a> <a href="https://github.com/anthropics/claude-code"><img src="https://img.shields.io/badge/Claude_Code-hooks%20%2B%20MCP-7c3aed" alt="Claude Code"></a>

<a href="https://lossless-claude.com">Website</a> &bull; <a href="#runtime-model">Runtime Model</a> &bull; <a href="#installation">Installation</a> &bull; <a href="#mcp-tools">MCP Tools</a> &bull; <a href="#development">Development</a>


lcm replaces sliding-window forgetfulness with a persistent memory runtime for both humans and agents.

  • Every message is stored in a project SQLite database.
  • Older context is compacted into a DAG of summaries instead of being dropped.
  • Durable decisions and findings are promoted into cross-session memory.
  • Claude Code already has end-to-end hook integration, while VS Code, Codex, and Oh My Pi use connector-based workflows on the same backend today.

Humans and agents use the same backend. The integration surface differs by client, but the memory model is shared.

This repo started as a fork of lossless-claw by Martian Engineering, adapted for Claude Code. The LCM model and DAG architecture originate from the Voltropy paper.

Runtime Model

flowchart LR
  subgraph Clients["Clients"]
    CC["Claude Code<br/>hooks + MCP"]
  end

  CC --> D["lcm daemon"]

  D --> DB[("project SQLite DAG")]
  D --> PM[("promoted memory FTS5")]
  D --> TOOLS["MCP tools<br/>search / grep / expand / describe / store / stats / doctor"]

Capabilities by integration path

PathRestorePrompt hintsTurn writebackAutomatic compactionNotes
Claude CodeYesYesYes, via transcript/hooksYesPrimary hook-based integration
GitHub Copilot (VS Code)NoYes, via skill/rulesNoNoRepo-local skill can teach Copilot to call lcm, but there is no automatic restore or turn capture yet
CodexYesYesYes, via native lifecycle hooksLCM memory compacts on PreCompact; native compaction continueslcm connectors install codex installs the hooks; see docs/vscode-codex.md. MCP config in .codex/config.toml is still manual
Oh My PiYesYesYes, via native lifecycle hooksYeslcm connectors install omp installs the hooks and --type mcp registers the MCP server; see docs/omp.md.

LCM Model

PhaseWhat happens
PersistRaw messages are stored in SQLite per conversation
SummarizeOlder messages are grouped into leaf summaries
CondenseSummaries roll up into higher-level DAG nodes
PromoteDurable insights are copied into cross-session memory
RestoreNew sessions recover context from summaries and promoted memory
RecallAgents query, expand, and inspect memory on demand

Nothing is dropped. Raw messages remain in the database. Summaries point back to their sources. Promoted memory remains searchable across sessions.

flowchart TD
  A["conversation / tool output"] --> B["persist raw messages"]
  B --> C["compact into leaf summaries"]
  C --> D["condense into deeper DAG nodes"]
  C --> E["promote durable insights"]
  D --> F["restore future context"]
  E --> F
  F --> G["search / grep / describe / expand / store"]

Installation

Prerequisites

  • Node.js 22+
  • Claude Code if you want hook-based automation
  • GitHub Copilot in VS Code if you want VS Code integration
  • Codex CLI if you want Codex connector installation, summarization, or transcript import
  • Oh My Pi if you want Oh My Pi connector installation or transcript import

Claude Code

Install the lcm binary first:

npm install -g @lossless-claude/lcm  # provides the `lcm` command
claude plugin marketplace add lossless-claude/lcm
claude plugin install lcm@lossless-claude
lcm install

lcm install writes config, registers MCP, installs the /memory skill and lcm.md, and sets up the Claude Code, Codex, and Oh My Pi integrations. It reports one outcome per harness and exits non-zero on any failure; --dry-run writes nothing. Run from the Claude Code plugin (/memory install) it skips the connector installs and leaves the MCP entry to the plugin manifest; use the npm CLI for all harnesses.

VS Code (GitHub Copilot)

Install the lcm binary first:

npm install -g @lossless-claude/lcm

Then install the repo-local Copilot connector:

lcm connectors install github-copilot
lcm connectors doctor github-copilot

This creates a workspace skill under .agents/skills/lcm-memory/SKILL.md. Both Codex and Copilot read that directory, so one installed file serves either host.

Codex

Install the lcm binary first:

npm install -g @lossless-claude/lcm

Then install the Codex connector:

lcm connectors install codex
lcm connectors doctor codex

The default connector installs native hooks for automatic restore, prompt recall, incremental turn capture, and compaction continuity. Review and trust them in Codex /hooks; connector diagnostics distinguish configuration from activation. See Codex setup. Paginated rollout rewrites resume capture only for a proven subagent tail match; other paginated recovery mismatches remain blocked and are reported by lcm doctor. Stored history is preserved.

Import older Codex or Oh My Pi sessions, or replay all supported history:

lcm import --codex
lcm import --omp
lcm import --replay

If you also want MCP inside Codex, run lcm connectors install codex --type mcp. Today that prints the TOML block you must add manually to .codex/config.toml.

See docs/vscode-codex.md and docs/omp.md for the current connector setup paths and known shortcomings.

Hooks

The plugin registers seven hooks. Every hook fails open (exit 0) and, before running, removes stale copies of itself left in settings.json by older installers so nothing fires twice.

HookCommandPurpose
PreCompactlcm compact --hookWrites a DAG summary before compaction
SessionStartlcm restoreRestores project context, recent summaries, and promoted memory
SessionEndlcm session-endIngests the completed Claude transcript
UserPromptSubmitlcm user-promptSearches memory and injects prompt-time hints
Stoplcm session-snapshotRolling transcript ingest, throttled
PostToolUselcm post-toolPassive learning: records decisions, plans, files, commands
PostToolUseFailurelcm post-toolPassive learning: records tool errors
flowchart LR
  SS["SessionStart"] --> CONV["Conversation"]
  CONV --> UP["UserPromptSubmit<br/>(each prompt)"]
  UP --> CONV
  CONV --> PC["PreCompact<br/>(if context fills)"]
  PC --> CONV
  CONV --> SE["SessionEnd"]

MCP Tools

ToolPurpose
lcm_searchSearch across episodic memory (messages and summaries) and promoted memory
lcm_grepRegex or full-text search across raw messages and summaries
lcm_expandDecompress a summary node into its source content by traversing the DAG
lcm_describeInspect session or summary metadata, lineage and explicit commit references
lcm_storePersist durable memory manually with optional tags
lcm_statsShow token savings, compression ratios, and usage statistics
lcm_doctorDiagnose setup and report stale stores, record-less stores and orphan summaries

CLI

# Setup & diagnostics
lcm install                # setup wizard
lcm uninstall              # remove hooks, MCP, and config
lcm doctor                 # diagnostics, bounded store lists (stale, record-less) and orphan-summary ids
lcm doctor --verbose       # complete store lists, orphan ids and event details
lcm doctor --repair-manual-attribution # preview manual memory session attribution; explicit --apply requires an offline hold and backs up each changed store
lcm diagnose               # scan recent sessions for hook failures
lcm status                 # daemon + summarizer mode
lcm -V                     # version

# Memory inspection
lcm search "query"        # search episodic and promoted memory
lcm grep "pattern"        # search messages and summaries
lcm describe <nodeId>      # inspect session or summary metadata and commit references
lcm expand <nodeId>        # expand a summary node into source detail
lcm store "content"       # persist a durable memory entry
lcm stats                  # memory and compression overview
lcm stats -v               # per-conversation breakdown
lcm stats --warning-backtest # offline environment-warning backtest for the current project; warnings stay off
lcm stats --pool           # connection pool statistics
lcm stats --pool --json    # connection pool statistics as JSON

# Compaction & promotion
lcm compact                # compact the current project
lcm compact --all          # compact all tracked projects
lcm compact --verbose      # per-session token detail
lcm compact --replay       # compact sequentially with threaded context (resumable)
lcm compact --replay --restart  # discard recorded progress and start from scratch
lcm promote                # promote durable insights to long-term memory
lcm promote --all          # promote across all tracked projects

# Import / export
lcm import                 # import Claude Code, Codex and OMP sessions for the current project
lcm import --all           # import all projects
lcm import --replay        # import and compact with threaded context (resumable)
lcm import --replay --restart # discard recorded progress and start from scratch
lcm import --provider codex   # import Codex sessions (--codex is the short form)
lcm import --provider omp     # import Oh My Pi sessions (--omp is the short form)
lcm export                 # export promoted knowledge to JSON on stdout
lcm export --all --output <f> # every project, written to files; --tags, --since filter
lcm import-knowledge <f>   # import a knowledge JSON file

# Connectors (wire lcm into other AI agents)
lcm connectors list        # list available agents and installed connectors
lcm connectors install <a> # install a connector for an agent (--type rules|mcp|skill|hooks)
lcm connectors remove <a>  # remove a connector for an agent
lcm connectors doctor      # check connector health
lcm connectors install <a> --global  # in the agent's user-level config, not this repo

# Sensitive data
lcm sensitive add <pat>    # add a redaction pattern (project-scoped)
lcm sensitive add --global # add a global redaction pattern
lcm sensitive list         # list all active patterns
lcm sensitive test <str>   # test what gets redacted
lcm sensitive purge --yes  # remove all stored data for the current project

# Daemon
lcm daemon start --detach  # start daemon in background
lcm daemon restart         # pick up changed LCM_* values
lcm daemon stop --hold     # keep it down so hooks cannot respawn it (--minutes <n>, --reason <text>)

# Hook handlers (internal — called by Claude Code and Codex hooks)
lcm compact --hook         # PreCompact hook (Claude Code)
lcm restore                # SessionStart hook (Claude Code)
lcm session-end            # SessionEnd hook (Claude Code)
lcm user-prompt            # UserPromptSubmit hook (Claude Code)
lcm post-tool              # PostToolUse + PostToolUseFailure hooks (Claude Code, passive learning)
lcm session-snapshot       # Stop hook (Claude Code, rolling ingest)
lcm codex-hook             # native Codex lifecycle hook — see docs/vscode-codex.md

# MCP server
lcm mcp                    # start MCP server

lcm help [command] prints this same reference from the CLI. lcm bench is a development tool: it benchmarks search against a local corpus, which the npm package does not include; see docs/search.md.

Configuration

All environment variables are optional. The default summarizer mode is auto. The daemon reads the tuning values when it starts: after changing one, run lcm daemon restart.

VariableDefaultDescription
LCM_SUMMARY_PROVIDERautoauto, claude-process, codex-process, copilot-process, omp-process, anthropic, openai, disabled, or session
ANTHROPIC_API_KEYunsetRead only when llm.provider is anthropic (or session falling back to it) and llm.apiKey is unset
LCM_HOME~/.lossless-claudeWhere the daemon, databases, sidecars and logs live
LCM_ENABLEDtrueSet to false to make every Claude Code and Codex command hook a no-op while keeping the plugin registered
LCM_CONTEXT_THRESHOLD0.75Context fill ratio that triggers compaction
LCM_FRESH_TAIL_COUNT8Most recent raw messages protected from compaction
LCM_LEAF_MIN_FANOUT3Minimum raw messages outside the fresh tail before a leaf pass runs
LCM_CONDENSED_MIN_FANOUT2Minimum same-depth summaries before they are condensed
LCM_CONDENSED_MIN_FANOUT_HARD1The same minimum during a hard-trigger sweep
LCM_LEAF_CHUNK_TOKENS20000Maximum source tokens per leaf compaction pass
LCM_CONDENSED_TARGET_TOKENS900Target size for condensed summaries

auto resolves per caller:

  • Claude caller -> claude-process
  • Codex caller -> codex-process
  • Copilot caller -> copilot-process
  • OMP caller -> omp-process
  • explicit config or LCM_SUMMARY_PROVIDER override always takes precedence

See docs/configuration.md for tuning notes and deeper operational guidance.

Development

npm install
npm run build
npx vitest
npx tsc --noEmit

Build before testing: the suite asserts dist/ was built from the current sources, and fails with the rebuild command instead of silently testing a stale binary. LCM_SKIP_CACHE_SYNC=1 npm run build skips the plugin-cache sync when only dist/ matters.

To score a candidate summarizer model against the real compaction engine, see docs/summarizer-bench.md. The bench is opt-in — it is skipped unless LCM_EVAL_MODEL and LCM_EVAL_CORPUS_DIR are set, so npx vitest never calls a paid API.

Repository layout

bin/
  lcm.ts                      CLI entry point (binary: lcm)
src/
  compaction.ts               DAG compaction engine
  connectors/                 client integration adapters
  daemon/                     HTTP daemon, lifecycle, config, routes
  db/                         SQLite schema + promoted memory
  hooks/                      Claude hook handlers + auto-heal
  llm/                        summarizer backends
  mcp/                        MCP server + tool definitions
  store/                      conversation and summary persistence
installer/
  install.ts                  setup wizard
  uninstall.ts                cleanup
test/
  bench/                      summarizer eval bench (opt-in, see docs/summarizer-bench.md)
  ...                         Vitest suites

Privacy

All conversation data is stored locally in ~/.lossless-claude/. Nothing is sent to any lossless-claude server.

If you configure an external summarizer (claude-process, anthropic, openai, etc.), messages are sent to that provider for summarization — after built-in secret redaction. lcm scrubs common secret patterns (API keys, tokens, passwords) from message content before writing to SQLite and before sending to the summarizer.

Add project-specific patterns with lcm sensitive add "MY_PATTERN". See docs/privacy.md for full details.

Technical Notes

  • Claude Code integration is hook-first.
  • The daemon is shared; the memory backend is client-agnostic.
  • The repo carries the original lossless-claw lineage; the current runtime is Claude Code oriented.

Acknowledgments

lcm stands on the shoulders of lossless-claw, the original implementation by Martian Engineering. The DAG-based compaction architecture, the LCM memory model, and the foundational design decisions all originate there.

The underlying theory comes from the LCM paper by Voltropy.

License

MIT

Source 9 files
hooks/lcm-hooks.ts 1030 lines
1// hooks/lcm-hooks.ts — lcm's function-hooks module (Claude Code early access).
2//
3// Claude Code loads it whenever its mods (function hooks) are on; the command hooks remain
4// for when they are not. It claims the session in a temp-dir file, rewritten in every
5// classic event a claim-reading command hook listens to, before that command hook runs,
6// and withdrawn at session.end. While the claim stands the SessionStart, PostToolUse,
7// PostToolUseFailure, UserPromptSubmit and Stop command hooks stay silent
8// (functionHooksOwnSession in src/hooks/session-claim.ts) and this module does their work
9// through the daemon:
10//   tool.call       → POST /tool-event    (the daemon writes the passive-learning rows)
11//   prompt.submit   → POST /prompt-search (memory hits ride as hidden context on the prompt)
12//   prompt.section  → the learning instruction is appended once to the system prompt's
13//                     `memory` section, instead of to every prompt.
14//   prompt.context  → POST /restore (the session's memory, as one named context block)
15//   turn.complete   → POST /ingest (the daemon reads the transcript delta) and
16//                     POST /promote-events, at most once a minute, replacing the Stop hook.
17// The module has no Node and no SQLite, so the daemon does every write.
18//
19// Types: Claude Code writes them to .claude-plugin/types/ when it loads this plugin from a
20// folder you own (`claude --plugin-dir .`); then `import type { Register } from "claude-code"`.
21// They are written from the running build and not committed, so the check is on demand:
22// `npm run typecheck:hooks`. Run it after a Claude Code update too — the API is early
23// access, and a green run is the answer to whether the release moved anything under this.
24// `claude plugin validate` reads this file statically: `$` may only be passed to a function
25// declared at the top level, and calls on it must be spelled `$.noun.method(...)`.
26import type { Register, EngineInterface } from "claude-code";
27import { endSharedSessionBudget, sharedSessionOutputBudget, type SessionOutputBudget } from "./model-budget.js";
28import { endShadowSessionBudget } from "./shadow-budget.js";
29import type { ShadowDeadline } from "./shadow-deadline.js";
30import type { ShadowAppend } from "./shadow-boundaries.js";
31import { ShadowSessionState, runCompactionShadow, type ShadowTransport, type ShadowEngine } from "./compaction-shadow.js";
32
33/** Same set the PostToolUse matcher in plugin.json names; `mcp__*` is matched by prefix. */
34const CAPTURED_TOOLS = new Set([
35  "Agent", "AskUserQuestion", "Bash", "EnterPlanMode", "ExitPlanMode",
36  "Read", "Edit", "Write", "Glob", "Grep", "TaskCreate", "TaskUpdate", "Skill",
37]);
38const DEFAULT_PORT = 3737;
39/** Fallback when the host command cannot report TMPDIR; matches Node os.tmpdir() on POSIX. */
40const DEFAULT_TMP_DIR = "/tmp";
41/** Same default as hooks.snapshotIntervalSec for the Stop command hook: one ingest a minute. */
42const INGEST_INTERVAL_MS = 60_000;
43let lastIngestAt = 0;
44
45// Verbatim copy of LEARNING_INSTRUCTION in src/guidance.ts (the module cannot import from src/);
46// test/hooks/learning-instruction.test.ts fails when the two drift.
47const LEARNING_INSTRUCTION = `<learning-instruction>
48When you recognize a durable insight, call lcm_store immediately:
49- decision: an architectural or design choice, with the trade-off that settled it
50- preference: how the user wants things done
51- root-cause: a bug cause that took effort to uncover
52- pattern: a codebase convention documented nowhere else
53- gotcha: a non-obvious pitfall
54- solution: a non-trivial fix worth remembering
55- workflow: a multi-step process that works
56
57Tag each store with type: plus one of project: or scope:; add source: when the origin matters for trust, and priority: rarely.
58Usage: lcm_store(text: "concise insight with why", tags: ["type:decision", "project:<repo>"])
59
60When you act on a surfaced memory (use it to inform a decision, avoid a known pitfall, or reference it in your work), emit:
61lcm_store(text: "Acted on memory <id> — <one-line how>", tags: ["signal:memory_used", "memory_id:<id>"])
62
63When you check a surfaced memory against current evidence, vote on it (reason is required both ways):
64lcm_store(text: "<what confirmed it, e.g. a file, test, or command output>", tags: ["signal:memory_vote", "vote:+1", "memory_id:<id>"])
65lcm_store(text: "<what contradicts it>", tags: ["signal:memory_vote", "vote:-1", "memory_id:<id>"])
66"Not relevant here" is not a -1 — only a real contradiction is.
67</learning-instruction>`;
68
69/** What the module cannot read for itself: everything outside $.fs's project-and-temp reach. */
70type HostEnv = { port: number; token: string | null; tmpDir: string };
71let hostEnv: Promise<HostEnv> | null = null;
72/** Routes the running daemon answered 404 for: an older lcm build. Logged once each, not per call. */
73const missingRoutes = new Set<string>();
74
75type HookObservation = {
76  hook: string;
77  operation: string;
78  kind: "delivery" | "execution";
79  status: string;
80  reason: string;
81  count: number;
82};
83type HookSnapshot = {
84  counts: Map<string, HookObservation>;
85  failures: { hook: string; operation: string; code: string; at: number }[];
86  seq: number;
87  truncated: boolean;
88  writing: boolean;
89};
90const hookSnapshots = new Map<string, HookSnapshot>();
91const MAX_HOOK_OBSERVATIONS = 128;
92const MAX_ACTIVE_HOOK_SESSIONS = 32;
93const hookSnapshotGeneration = Date.now();
94const failedHookSnapshotWrites = new Set<string>();
95
96function sessionFileId(sessionId: string): string {
97  return encodeURIComponent(sessionId).replace(/[_.!~*'()]/g,
98    (char) => `%${char.charCodeAt(0).toString(16).toUpperCase().padStart(2, "0")}`);
99}
100
101function noteHook(
102  sessionId: string, hook: string, operation: string,
103  kind: HookObservation["kind"], status: string, reason = "",
104): void {
105  if (!sessionId?.trim()) return;
106  const snapshot = hookSnapshots.get(sessionId) ?? {
107    counts: new Map<string, HookObservation>(), failures: [], seq: 0, truncated: false, writing: false,
108  };
109  if (!hookSnapshots.has(sessionId) && hookSnapshots.size >= MAX_ACTIVE_HOOK_SESSIONS) {
110    hookSnapshots.delete(hookSnapshots.keys().next().value!);
111  }
112  hookSnapshots.set(sessionId, snapshot);
113  noteObservation(snapshot, { hook, operation, kind, status, reason, count: 1 });
114}
115function noteObservation(snapshot: HookSnapshot, observation: HookObservation): void {
116  const key = JSON.stringify([observation.hook, observation.operation, observation.kind, observation.status, observation.reason]);
117  const existing = snapshot.counts.get(key);
118  if (existing) existing.count++;
119  else addObservation(snapshot, key, observation);
120  if (observation.status !== "failed" && observation.status !== "rejected") return;
121  snapshot.failures.push({ hook: observation.hook, operation: observation.operation, code: observation.reason || observation.status, at: Date.now() });
122  if (snapshot.failures.length <= 32) return;
123  snapshot.failures.shift(); snapshot.truncated = true;
124}
125function addObservation(snapshot: HookSnapshot, key: string, observation: HookObservation): void {
126  if (snapshot.counts.size >= MAX_HOOK_OBSERVATIONS) {
127    snapshot.counts.delete(snapshot.counts.keys().next().value!); snapshot.truncated = true;
128  }
129  snapshot.counts.set(key, observation);
130}
131
132async function flushHookObservations($: EngineInterface, sessionId: string): Promise<void> {
133  if (!sessionId?.trim()) return;
134  const snapshot = hookSnapshots.get(sessionId);
135  if (!snapshot || snapshot.writing) return;
136  snapshot.writing = true;
137  const pending = writeHookSnapshot($, { snapshot, sessionId }).catch(() => undefined)
138    .finally(() => { snapshot.writing = false; });
139  try {
140    await Promise.race([
141      pending,
142      $.clock.sleep(250),
143    ]);
144  } catch {
145    return; // Diagnostic host failure leaves the turn usable.
146  }
147}
148
149async function writeHookSnapshot($: EngineInterface, { snapshot, sessionId }: { snapshot: HookSnapshot; sessionId: string }): Promise<void> {
150    const { tmpDir } = await readHostEnv($);
151    const safeId = sessionFileId(sessionId);
152    const cwd = await $.session.cwd().catch(() => "");
153    const seq = ++snapshot.seq;
154    const path = `${tmpDir}/lcm-hook-observe-${safeId}-${seq % 2}.json`;
155    const content = JSON.stringify({
156      version: 1, harness: "claude-function", sessionId, cwd, seq, generation: hookSnapshotGeneration,
157      updatedAt: Date.now(), truncated: snapshot.truncated,
158      observations: [...snapshot.counts.values()], failures: snapshot.failures,
159    });
160    try {
161      await $.fs.write(path, content);
162    } catch {
163      warnSnapshotWrite($, sessionId);
164    }
165
166}
167
168function warnSnapshotWrite($: EngineInterface, sessionId: string): void {
169  if (failedHookSnapshotWrites.has(sessionId)) return;
170  failedHookSnapshotWrites.add(sessionId);
171  try { $.ui.log("[lcm] hook observation snapshot could not be written"); }
172  catch { return; }
173}
174
175/** No config file yet, or one being rewritten, both mean the compiled-in port. */
176function parsePort(configJson: string): number {
177  try {
178    const port = JSON.parse(configJson).daemon?.port;
179    return typeof port === "number" ? port : DEFAULT_PORT;
180  } catch {
181    return DEFAULT_PORT;
182  }
183}
184
185// The daemon's port and bearer token live under the lcm home, and TMPDIR is an
186// environment value; neither is reachable through $.fs (project and temp dir only), so
187// one host command reads all three once per module load. LCM_HOME moves that home, which
188// is what lets a sandbox run this module without touching the host's own lcm state.
189function readHostEnv($: EngineInterface): Promise<HostEnv> {
190  hostEnv ??= $.process.run(["sh", "-c",
191    // Trimmed before the emptiness test, so a blank LCM_HOME falls back here exactly as
192    // it does in lcmHome(); otherwise the two sides disagree about the root.
193    'H=$(printf %s "${LCM_HOME:-}" | sed "s/^[[:space:]]*//; s/[[:space:]]*$//"); '
194    + '[ -n "$H" ] || H="$HOME/.lossless-claude"; '
195    + 'cat "$H/daemon.token" 2>/dev/null; echo; echo "__CONFIG__"; '
196    + 'cat "$H/config.json" 2>/dev/null; echo; echo "__TMPDIR__"; '
197    + 'printf %s "${TMPDIR:-/tmp}"',
198  ]).then(({ stdout }) => {
199    const [tokenPart, afterToken = ""] = stdout.split("__CONFIG__");
200    const [configPart, tmpPart = ""] = afterToken.split("__TMPDIR__");
201    return {
202      port: parsePort(configPart),
203      token: tokenPart.trim() || null,
204      tmpDir: tmpPart.trim().replace(/\/+$/, "") || DEFAULT_TMP_DIR,
205    };
206  }, () => ({ port: DEFAULT_PORT, token: null, tmpDir: DEFAULT_TMP_DIR }));
207  return hostEnv;
208}
209
210/**
211 * Tell the command hooks this module is live for this session, so they stay silent
212 * instead of doing the same work again. Read by functionHooksOwnSession in
213 * src/hooks/session-claim.ts, which names the same file from Node's os.tmpdir() and
214 * counts it only while `ts` is fresh and no `ended` is set. No claim means "not mine",
215 * costing a duplicate row rather than a lost one, so nothing here is worth failing an
216 * event over.
217 */
218async function writeClaim($: EngineInterface, sessionId: string, ended?: string): Promise<void> {
219  if (!sessionId) return;
220  const { tmpDir } = await readHostEnv($);
221  const file = `${tmpDir}/lcm-claim-${sessionFileId(sessionId)}.json`;
222  await $.fs.write(file, JSON.stringify({ sessionId, ts: Date.now(), ...(ended ? { ended } : {}) }));
223}
224
225/** Sessions whose claim could not be written, so the failure is logged once, not per event. */
226const unclaimedSessions = new Set<string>();
227
228async function claimSession($: EngineInterface, sessionId: string, hook: string): Promise<void> {
229  let claimed = true;
230  await writeClaim($, sessionId).catch((error: unknown) => {
231    claimed = false;
232    if (unclaimedSessions.has(sessionId)) return;
233    unclaimedSessions.add(sessionId);
234    $.ui.log(`[lcm] could not claim the session, command hooks stay active: ${String(error)}`);
235  });
236  noteHook(sessionId, hook, "claim", "execution", claimed ? "completed" : "failed",
237    claimed ? "" : "write-failed");
238}
239
240/**
241 * The claim is rewritten in each classic event a claim-reading command hook listens to.
242 * Modules run before the command hooks of the same event, which run inside `next(e)`,
243 * so the claim is fresh when they read it. That also covers a session id that begins
244 * with /clear, /resume or /branch, where no session.start fires: its classic
245 * SessionStart claims it before the SessionStart command hook can restore.
246 */
247function registerSessionClaim(on: On, shadow?: ShadowSessionState): void {
248  on("classic.SessionStart", async ($, e, next) => {
249    if (["clear", "resume", "fork"].includes(e.source)) shadow?.reset(e.session_id);
250    await claimSession($, e.session_id, "classic.SessionStart");
251    return next(e);
252  });
253  on("classic.UserPromptSubmit", async ($, e, next) => {
254    await claimSession($, e.session_id, "classic.UserPromptSubmit");
255    return next(e);
256  });
257  on("classic.PostToolUse", async ($, e, next) => {
258    await claimSession($, e.session_id, "classic.PostToolUse");
259    return next(e);
260  });
261  on("classic.PostToolUseFailure", async ($, e, next) => {
262    await claimSession($, e.session_id, "classic.PostToolUseFailure");
263    return next(e);
264  });
265  on("classic.Stop", async ($, e, next) => {
266    await claimSession($, e.session_id, "classic.Stop");
267    return next(e);
268  });
269}
270
271/**
272 * session.end fires on exit, /clear, /resume and /branch. The ending id's claim is
273 * withdrawn so a later run of that id without this module records again; a crash skips
274 * this, and its claim lapses with its timestamp instead.
275 */
276function registerSessionEnd(on: On, shadow?: ShadowSessionState): void {
277  on("session.end", async ($, e, next) => {
278    shadow?.reset(e.sessionId);
279    endSharedSessionBudget(e.sessionId); endShadowSessionBudget(e.sessionId);
280    await writeClaim($, e.sessionId, e.reason).catch(() => { /* the claim lapses on its own */ });
281    return next(e);
282  });
283}
284
285/** Last time this module asked the host to start the daemon; one attempt per cooldown window. */
286let lastDaemonStartAt = 0;
287const DAEMON_START_COOLDOWN_MS = 60_000;
288const DAEMON_START_TIMEOUT_MS = 15_000;
289const HEALTH_PROBE_TIMEOUT_MS = 500;
290const DAEMON_POST_TIMEOUT_MS = 5_000;
291const RESTORE_TIMEOUT_MS = 10_000;
292/** The daemon's summary-job long poll holds for up to 25 seconds. */
293const SUMMARY_POLL_TIMEOUT_MS = 30_000;
294let warnedBusyDaemon = false;
295/** POSIX exit code for "command not found": no `lcm` binary on PATH. */
296const EXIT_COMMAND_NOT_FOUND = 127;
297/** `daemon start --automatic` reports an active hold with EX_TEMPFAIL. */
298const EXIT_DAEMON_HELD = 75;
299let warnedNoLcmBinary = false;
300
301/**
302 * The daemon exits when idle and the command hooks used to bring it back (ensureDaemon).
303 * While this module owns the events, it must do the same, or every tool call and prompt
304 * between the idle exit and the next SessionStart is lost. `lcm daemon start --detach`
305 * spawns the daemon and waits until /health answers.
306 */
307async function startDaemon($: EngineInterface): Promise<boolean> {
308  const now = Date.now();
309  if (now - lastDaemonStartAt < DAEMON_START_COOLDOWN_MS) return false;
310  lastDaemonStartAt = now;
311  const run = await $.process.run(
312    ["sh", "-c", 'command -v lcm >/dev/null 2>&1 || exit 127; exec lcm daemon start --detach --automatic'],
313    { timeoutMs: DAEMON_START_TIMEOUT_MS },
314  ).catch(() => null);
315  if (run?.exitCode === EXIT_DAEMON_HELD) {
316    if (lastDaemonStartAt === now) lastDaemonStartAt = 0;
317    return false;
318  }
319  return reportDaemonLaunch($, run);
320}
321function reportDaemonLaunch($: EngineInterface, run: Awaited<ReturnType<EngineInterface["process"]["run"]>> | null): boolean {
322  if (run?.exitCode === 0) return true;
323  if (run?.exitCode === EXIT_COMMAND_NOT_FOUND && !warnedNoLcmBinary) {
324    warnedNoLcmBinary = true;
325    $.ui.log("[lcm] daemon is down and no `lcm` binary is on PATH to start it; events are lost until a command hook restarts it");
326  } else if (run && run.exitCode !== EXIT_COMMAND_NOT_FOUND) {
327    $.ui.log(`[lcm] daemon start failed (exit ${run.exitCode}): ${run.stderr.trim().split("\n")[0] ?? ""}`);
328  }
329  return false;
330}
331
332type PostOutcome = { body: Record<string, unknown> | null; connectionFailed: boolean; httpStatus?: number };
333
334/** A refused connection fails at once; a listener that is busy fails only after a wait. */
335const REFUSED_WITHIN_MS = 1_000;
336
337/**
338 * Whether a call started at `startedAt` failed for lack of a listener, which a start and a
339 * retry can fix. The error comes from the host's realm, so its fields are read rather than
340 * trusted to `instanceof`; one that names no refusal is judged by how long it took.
341 */
342function isDaemonUnreachableError(error: unknown, startedAt: number): boolean {
343  const fields = (value: unknown) =>
344    (typeof value === "object" && value !== null ? value : {}) as { code?: unknown; message?: unknown; cause?: unknown; name?: unknown };
345  const failure = fields(error);
346  if (failure.name === "TimeoutError" || failure.name === "AbortError") return false;
347  const text = [failure.code, fields(failure.cause).code, failure.message].map(String).join(" ");
348  return /ECONNREFUSED|connection refused/i.test(text) || Date.now() - startedAt < REFUSED_WITHIN_MS;
349}
350
351/** Bound our wait even on hosts that do not put a deadline on HTTP fetches. */
352async function fetchDaemon(
353  $: EngineInterface, url: string, init: Parameters<EngineInterface["http"]["fetch"]>[1], timeoutMs: number,
354): ReturnType<EngineInterface["http"]["fetch"]> {
355  return Promise.race([
356    Promise.resolve().then(() => $.http.fetch(url, init)),
357    Promise.resolve().then(() => $.clock.sleep(timeoutMs)).then(() => {
358      const error = new Error("daemon did not answer within its deadline");
359      error.name = "TimeoutError";
360      throw error;
361    }),
362  ]);
363}
364
365function reportBusyDaemon($: EngineInterface): void {
366  if (warnedBusyDaemon) return;
367  warnedBusyDaemon = true;
368  $.ui.log("[lcm] daemon busy, did not answer; delivery is unconfirmed and the next capture will retry");
369}
370
371function delivery(outcome: PostOutcome): ["accepted" | "rejected" | "unconfirmed", string] {
372  return outcome.body ? ["accepted", ""]
373    : outcome.httpStatus !== undefined ? ["rejected", `http-${outcome.httpStatus}`]
374      : ["unconfirmed", "unconfirmed"];
375}
376
377/** A 404 means an older lcm build is listening. Say so once per route, not per call. */
378function logMissingRoute($: EngineInterface, route: string, consequence: string): void {
379  if (missingRoutes.has(route)) return;
380  missingRoutes.add(route);
381  $.ui.log(`[lcm] daemon has no ${route} route (older lcm build); ${consequence}`);
382}
383
384async function postOnce($: EngineInterface, route: string, body: unknown): Promise<PostOutcome> {
385  let startedAt = Date.now();
386  let res: Awaited<ReturnType<EngineInterface["http"]["fetch"]>>;
387  try {
388    const { port, token } = await readHostEnv($);
389    startedAt = Date.now();
390    res = await fetchDaemon($, `http://127.0.0.1:${port}${route}`, {
391      method: "POST",
392      headers: { "content-type": "application/json", ...(token ? { authorization: `Bearer ${token}` } : {}) },
393      body: JSON.stringify(body),
394    }, route === "/restore" ? RESTORE_TIMEOUT_MS : DAEMON_POST_TIMEOUT_MS);
395  } catch (error) {
396    if (isDaemonUnreachableError(error, startedAt)) return { body: null, connectionFailed: true };
397    // A timed-out or dropped call may still complete. The next Stop snapshot or scan
398    // captures the turn, so starting another daemon or retrying now adds no safety.
399    reportBusyDaemon($);
400    return { body: null, connectionFailed: false };
401  }
402  if (res.status === 404) {
403    // Not "the command hooks still record": they stand down while this module holds the
404    // session, so an older daemon means the call is simply lost.
405    logMissingRoute($, route, "this call is dropped until the daemon is upgraded");
406    return { body: null, connectionFailed: false, httpStatus: res.status };
407  }
408  if (!res.ok) {
409    $.ui.log(`[lcm] ${route}: daemon answered ${res.status}`);
410    return { body: null, connectionFailed: false, httpStatus: res.status };
411  }
412  try {
413    return { body: JSON.parse(res.text) as Record<string, unknown>, connectionFailed: false };
414  } catch {
415    // A listener answered, so neither a start nor a retry applies.
416    $.ui.log(`[lcm] ${route}: daemon answered ${res.status} with a body that is not JSON`);
417    return { body: null, connectionFailed: false };
418  }
419}
420
421/**
422 * POST `body` to the daemon; resolves to the parsed JSON, or null when the daemon lacks the
423 * route or cannot be reached. Only a refused connection starts the daemon and retries once.
424 */
425async function postDaemonOutcome($: EngineInterface, route: string, body: unknown): Promise<PostOutcome> {
426  const first = await postOnce($, route, body);
427  if (!first.connectionFailed) return first;
428  if (!(await startDaemon($))) return first;
429  const second = await postOnce($, route, body);
430  if (second.connectionFailed) $.ui.log(`[lcm] ${route}: daemon still unreachable after start`);
431  return second;
432}
433
434async function postDaemon($: EngineInterface, route: string, body: unknown): Promise<Record<string, unknown> | null> {
435  return (await postDaemonOutcome($, route, body)).body;
436}
437
438type SummaryJob = {
439  id: string;
440  session_id: string;
441  kind: "leaf" | "condensed";
442  system: string;
443  prompt: string;
444  maxTokens: number;
445  pool?: true;
446};
447type SummaryAnswer = {
448  text: string;
449  providerId: "session:haiku" | "session:fork" | "session-pool:haiku" | "session-pool:sonnet";
450  usage: { input_tokens: number; output_tokens: number; estimated: boolean };
451  priorUsage?: UsageAttempt[];
452};
453type UsageAttempt = {
454  providerId: "session:haiku" | "session:fork" | "session-pool:haiku" | "session-pool:sonnet";
455  usage: { input_tokens: number; output_tokens: number; estimated: boolean };
456  /** Whether this call failed to provide a usable answer. */
457  failed?: boolean;
458};
459type SummaryFailure = Error & { usageAttempts?: UsageAttempt[]; usageUnknown?: boolean };
460const DEFAULT_SUMMARY_OUTPUT_CAP = 50_000;
461/** The engine's own estimate ratio; used only when the host reports no usage. */
462const CHARS_PER_TOKEN = 4;
463/** The host capped the long poll below the daemon's hold, so poll short and pause between. */
464const SHORT_POLL_PAUSE_MS = 2_000;
465/** The daemon has no summarizer route; a later respawn may bring it back. */
466const MISSING_ROUTE_RETRY_MS = 60_000;
467/** The daemon answered but not with a job: back off before asking again. */
468const POLL_BACKOFF_MS = 5_000;
469const WORKER_ANSWER_ATTEMPTS = 3;
470let summaryPollerStarted = false;
471
472function summaryDelay($: EngineInterface, ms: number): Promise<void> {
473  return new Promise((resolve) => { $.clock.after(ms, resolve); });
474}
475
476type WorkerModel = "haiku" | "sonnet";
477
478async function completeSummary($: EngineInterface, job: SummaryJob, remainingTokens: number, workerModel?: WorkerModel): Promise<SummaryAnswer> {
479  const providerId = workerModel ? `session-pool:${workerModel}` as const : "session:haiku";
480  const result = await $.model.complete({
481    model: workerModel ?? "haiku", system: job.system, prompt: job.prompt, maxTokens: Math.min(job.maxTokens, remainingTokens),
482  });
483  // Earlier engine builds returned the text directly; current builds return a
484  // discriminated result and do not reject when the provider cannot answer.
485  if (typeof result !== "string" && !result.isAnswered) {
486    const failure = new Error(result.reason) as SummaryFailure;
487    failure.usageAttempts = [{ providerId, usage: {
488      input_tokens: result.usage.input_tokens,
489      output_tokens: result.usage.output_tokens, estimated: false,
490    }, failed: true }];
491    throw failure;
492  }
493  const text = typeof result === "string" ? result : result.text;
494  // The engine can mark a completion answered without its text being a string (the
495  // exact shape is not confirmed — the module has no committed types for it). Treat
496  // that the same as an unanswered job instead of throwing out of `.trim()`.
497  if (typeof text !== "string") {
498    const failure = new Error("model.complete: answer text was not a string") as SummaryFailure;
499    if (typeof result !== "string") {
500      failure.usageAttempts = [{ providerId, usage: {
501        input_tokens: result.usage.input_tokens,
502        output_tokens: result.usage.output_tokens, estimated: false,
503      }, failed: true }];
504    }
505    throw failure;
506  }
507  const trimmed = text.trim();
508  const usage = typeof result === "string"
509    ? { input_tokens: Math.ceil((job.system.length + job.prompt.length) / CHARS_PER_TOKEN),
510        output_tokens: Math.ceil(trimmed.length / CHARS_PER_TOKEN), estimated: true }
511    : { input_tokens: result.usage.input_tokens,
512        output_tokens: result.usage.output_tokens, estimated: false };
513  return {
514    text: trimmed, providerId,
515    usage,
516  };
517}
518
519async function answerSummary($: EngineInterface, job: SummaryJob, remainingTokens: number): Promise<SummaryAnswer> {
520  if (job.kind !== "condensed") return completeSummary($, job, remainingTokens);
521  const fork = await $.model.fork({ prompt: `${job.system}\n\n${job.prompt}` }).catch(() => null);
522  if (fork && usableSummaryFork(fork)) {
523    return { text: fork.text.trim(), providerId: "session:fork",
524      usage: { input_tokens: fork.usage.input_tokens, output_tokens: fork.usage.output_tokens, estimated: false } };
525  }
526  const forkTextMalformed = Boolean(fork && "text" in fork && "usage" in fork);
527  const forkAttempt: UsageAttempt | undefined = fork && "usage" in fork ? { providerId: "session:fork",
528    usage: { input_tokens: fork.usage.input_tokens, output_tokens: fork.usage.output_tokens, estimated: false }, failed: true } : undefined;
529  const forkSpent = forkAttempt?.usage.output_tokens ?? 0;
530  if (forkSpent >= remainingTokens) {
531    const failure = new Error("spend cap") as SummaryFailure;
532    failure.usageAttempts = forkAttempt ? [forkAttempt] : []; throw failure;
533  }
534  return fallbackSummary($, { job, remainingTokens: remainingTokens - forkSpent, forkAttempt, forkTextMalformed });
535}
536function usableSummaryFork(fork: Awaited<ReturnType<EngineInterface["model"]["fork"]>>): fork is Extract<typeof fork, { text: string }> {
537  return "text" in fork && "usage" in fork && typeof fork.text === "string";
538}
539async function fallbackSummary($: EngineInterface, context: { job: SummaryJob; remainingTokens: number; forkAttempt?: UsageAttempt; forkTextMalformed: boolean }): Promise<SummaryAnswer> {
540  const { job, remainingTokens, forkAttempt, forkTextMalformed } = context;
541  let answer: SummaryAnswer;
542  try { answer = await completeSummary($, job, remainingTokens); }
543  catch (error) { throw failedSummaryFallback(error, forkAttempt, forkTextMalformed); }
544  return forkAttempt ? { ...answer, priorUsage: [forkAttempt] } : answer;
545}
546function failedSummaryFallback(error: unknown, attempt: UsageAttempt | undefined, malformed: boolean): unknown {
547  if (!attempt) return error;
548  const failure = (error instanceof Error ? error : new Error(String(error))) as SummaryFailure;
549  if (malformed) failure.message = `model.fork: answer text was not a string; fallback: ${failure.message}`;
550  failure.usageUnknown ||= !failure.usageAttempts?.length;
551  failure.usageAttempts = [attempt, ...(failure.usageAttempts ?? [])];
552  return failure;
553}
554
555/**
556 * One poll of `/summarize-jobs/next`. A `job` is one to run; `wait` means the daemon
557 * answered something other than a job and the caller should back off for that long.
558 */
559type PollOutcome =
560  | { job: SummaryJob }
561  | { wait: number; shortPoll?: true };
562
563async function fetchNextJob($: EngineInterface, { port, token }: HostEnv, sessionId: string, shortPoll: boolean, worker = false) {
564  const binding = worker ? `&caller_session_id=${encodeURIComponent(sessionId)}&cwd=${encodeURIComponent(await $.session.cwd())}&client=claude&transport=hook` : "";
565  return fetchDaemon($,
566    `http://127.0.0.1:${port}/summarize-jobs/next?${worker ? "worker_id" : "session_id"}=${encodeURIComponent(sessionId)}${binding}${shortPoll ? "&wait_ms=0" : ""}`,
567    { headers: token ? { authorization: `Bearer ${token}` } : {} }, SUMMARY_POLL_TIMEOUT_MS,
568  );
569}
570
571/** Turns a non-job response into how long to wait before asking again. */
572function classifyPollResponse($: EngineInterface, status: number): { wait: number } {
573  if (status === 204) return { wait: 0 };
574  if (status === 404) {
575    // An older build answered, or the daemon was swapped mid-session. A later
576    // respawn may bring the route back, so keep checking, slowly.
577    logMissingRoute($, "/summarize-jobs/next", "retrying every minute");
578    return { wait: MISSING_ROUTE_RETRY_MS };
579  }
580  if (status === 401) hostEnv = null;
581  $.ui.log(`[lcm] /summarize-jobs/next: daemon answered ${status}`);
582  return { wait: POLL_BACKOFF_MS };
583}
584
585async function nextSummaryJob(
586  $: EngineInterface, sessionId: string, shortPoll: boolean, worker = false,
587): Promise<PollOutcome> {
588  let job: SummaryJob | undefined;
589  let startedAt = Date.now();
590  let response: Awaited<ReturnType<typeof fetchNextJob>>;
591  try {
592    const host = await readHostEnv($);
593    // Timed from the request alone: a slow host-environment read is not the daemon's wait.
594    startedAt = Date.now();
595    response = await fetchNextJob($, host, sessionId, shortPoll, worker);
596  } catch (error) {
597    // Some hosts cap HTTP request duration below the daemon's 25-second hold.
598    if (isDaemonUnreachableError(error, startedAt)) {
599      await startDaemon($);
600      hostEnv = null; // A restarted daemon may have a new bearer token.
601    }
602    return { wait: POLL_BACKOFF_MS, shortPoll: true };
603  }
604  if (!response.ok || response.status === 204) return classifyPollResponse($, response.status);
605  try {
606    job = JSON.parse(response.text).job as SummaryJob;
607  } catch {
608    // A malformed 200 backs off like a transport failure instead of stopping the poller.
609    return { wait: POLL_BACKOFF_MS, shortPoll: true };
610  }
611  return validatePolledJob($, job, { sessionId, worker });
612}
613function validatePolledJob($: EngineInterface, job: SummaryJob | undefined, { sessionId, worker }: { sessionId: string; worker: boolean }): PollOutcome {
614  if (worker && !job) return { wait: 0 };
615  if (!job || !summaryJobMatches(job, sessionId, worker)) {
616    $.ui.log("[lcm] discarded summary job for a different session"); return { wait: POLL_BACKOFF_MS };
617  }
618  return { job };
619}
620function summaryJobMatches(job: SummaryJob, sessionId: string, worker: boolean): boolean {
621  return worker ? job.pool === true : job.pool !== true && job.session_id === sessionId;
622}
623
624async function postSummaryAnswer($: EngineInterface, job: SummaryJob, route: string, body: unknown): Promise<void> {
625  const sessionId = job.pool ? await $.session.id() : job.session_id;
626  try {
627    const boundBody = job.pool ? { ...(body as Record<string, unknown>),
628      worker_id: sessionId, caller_session_id: sessionId, cwd: await $.session.cwd(), client: "claude", transport: "hook",
629      providerId: (body as Record<string, unknown>).providerId ?? `session-pool:${await $.env.get("LCM_SUMMARIZE_WORKER_MODEL") || "haiku"}`,
630    } : body;
631    for (let attempt = 0; attempt < WORKER_ANSWER_ATTEMPTS; attempt++) {
632      const outcome = await postDaemonOutcome($, route, boundBody);
633      noteHook(sessionId, "session.start", "summary-answer", "delivery", ...delivery(outcome));
634      // Submissions are idempotent: a delivered answer with a lost response is
635      // discarded on retry. Reuse the answer rather than spending on completion again.
636      if (!retrySummaryDelivery(job, outcome)) return;
637      if (attempt + 1 === WORKER_ANSWER_ATTEMPTS) return;
638      if (outcome.httpStatus === 401) hostEnv = null;
639      await summaryDelay($, POLL_BACKOFF_MS);
640    }
641  } catch (error) {
642    noteHook(sessionId, "session.start", "summary-answer", "delivery", "unconfirmed", "transport");
643    throw error;
644  } finally {
645    await flushHookObservations($, sessionId);
646  }
647}
648
649function retrySummaryDelivery(job: SummaryJob, outcome: PostOutcome): boolean {
650  if (outcome.body || !job.pool) return false;
651  if (outcome.httpStatus === undefined) return true;
652  return outcome.httpStatus >= 500 || outcome.httpStatus === 401 || outcome.httpStatus === 429;
653}
654
655/** Answers one job. Returns the output tokens it spent, or null when the cap was hit. */
656async function serveSummaryJob(
657  $: EngineInterface, job: SummaryJob, budget: SessionOutputBudget, workerModel?: WorkerModel,
658): Promise<number | null> {
659  const route = `/summarize-jobs/${encodeURIComponent(job.id)}`;
660  const lease = await budget.reserveOrdinary(job.maxTokens, !workerModel && job.kind === "condensed");
661  if (!lease) {
662    await postSummaryAnswer($, job, route, { error: "spend cap" });
663    return null;
664  }
665  let answer: SummaryAnswer;
666  try {
667    answer = workerModel ? await completeSummary($, job, lease.allowance, workerModel)
668      : await answerSummary($, job, lease.allowance);
669  } catch (error) {
670    const { attempts, used, unknown } = failedSummaryAccounting(error);
671    lease.settle(used, unknown);
672    const overshot = budget.snapshot().overshoot > 0;
673    await postSummaryAnswer($, job, route, { error: overshot ? "spend cap" : error instanceof Error ? error.message : String(error),
674      ...(attempts.length ? { usageAttempts: attempts } : {}) });
675    return overshot ? null : used;
676  }
677  const attempts: UsageAttempt[] = [...(answer.priorUsage ?? []),
678    { providerId: answer.providerId, usage: answer.usage, failed: false }];
679  const totalOutput = attempts.reduce((sum, attempt) => sum + attempt.usage.output_tokens, 0);
680  lease.settle(totalOutput);
681  if (budget.snapshot().overshoot > 0) {
682    await postSummaryAnswer($, job, route, { error: "spend cap", usageAttempts: attempts });
683    return null;
684  }
685  if (!answer.text) {
686    await postSummaryAnswer($, job, route, { error: "empty summary", usageAttempts: [
687      ...(answer.priorUsage ?? []), { providerId: answer.providerId, usage: answer.usage, failed: true },
688    ] });
689    return totalOutput;
690  }
691  const { priorUsage, ...body } = answer;
692  await postSummaryAnswer($, job, route, priorUsage ? { ...body, usageAttempts: priorUsage } : body);
693  return totalOutput;
694}
695function failedSummaryAccounting(error: unknown): { attempts: UsageAttempt[]; used: number; unknown: boolean } {
696  const failure = error as SummaryFailure;
697  const attempts = failure?.usageAttempts ?? [];
698  const known = attempts.filter(attempt => Number.isSafeInteger(attempt.usage.output_tokens) && attempt.usage.output_tokens >= 0);
699  return { attempts, used: known.reduce((sum, attempt) => sum + attempt.usage.output_tokens, 0),
700    unknown: Boolean(failure?.usageUnknown) || attempts.length === 0 || known.length !== attempts.length };
701}
702
703/** One request at a time also serializes jobs from concurrent daemon compactions. */
704type ServingConfig = { cap: number; worker: boolean; workerModel?: WorkerModel };
705async function workerDeclared($: EngineInterface): Promise<boolean> {
706  try { return await $.env.get("LCM_SUMMARIZE_WORKER") === "1"; }
707  catch { return false; }
708}
709async function summaryServingConfig($: EngineInterface, configuredCap: number): Promise<ServingConfig | null> {
710  const worker = await workerDeclared($);
711  if (!worker) return { cap: configuredCap, worker };
712  const model = await $.env.get("LCM_SUMMARIZE_WORKER_MODEL") ?? "haiku";
713  if (model !== "haiku" && model !== "sonnet") {
714    $.ui.log("[lcm] worker model must be haiku or sonnet; worker stopped"); return null;
715  }
716  const setting = await $.env.get("LCM_SUMMARIZE_WORKER_MAX_OUTPUT_TOKENS");
717  const cap = setting === undefined ? configuredCap : Number(setting);
718  if (!Number.isSafeInteger(cap) || cap < 0) {
719    $.ui.log("[lcm] worker output cap must be a non-negative integer; worker stopped"); return null;
720  }
721  return { cap, worker, workerModel: model };
722}
723async function pollSummaries($: EngineInterface, configuredCap: number): Promise<void> {
724  const config = await summaryServingConfig($, configuredCap);
725  if (!config || config.cap === 0) return;
726  let shortPoll = false;
727  while (true) {
728    if (shortPoll) await summaryDelay($, SHORT_POLL_PAUSE_MS);
729    const sessionId = await $.session.id();
730    const outcome = await nextSummaryJob($, sessionId, shortPoll, config.worker);
731    if ("wait" in outcome) {
732      shortPoll = outcome.shortPoll ?? shortPoll;
733      await waitSummaryPoll($, outcome.wait); continue;
734    }
735    if (!await continueSummaryServing($, outcome.job, { sessionId, config })) return;
736  }
737}
738async function continueSummaryServing($: EngineInterface, job: SummaryJob, { sessionId, config }: { sessionId: string; config: ServingConfig }): Promise<boolean> {
739  const budget = sharedSessionOutputBudget(sessionId, config.cap);
740  const spent = await serveSummaryJob($, job, budget, config.workerModel);
741  return spent !== null && (!config.worker || budget.snapshot().available > 0);
742}
743async function waitSummaryPoll($: EngineInterface, delay: number): Promise<void> {
744  if (delay > 0) await summaryDelay($, delay);
745}
746
747type On = Parameters<Register>[0];
748
749function startSummaryPoller($: EngineInterface, summaryCap: number): void {
750  if (summaryPollerStarted) return;
751  summaryPollerStarted = true;
752  void pollSummaries($, summaryCap).catch((error) => {
753    $.ui.log(`[lcm] session summarizer stopped: ${String(error)}`);
754  });
755}
756
757/** Keep checking late command-hook registration without holding session.start open. */
758async function retryWorkerEnrollment($: EngineInterface, sessionId: string, cwd: string, summaryCap: number): Promise<void> {
759  while (!summaryPollerStarted) {
760    await summaryDelay($, POLL_BACKOFF_MS);
761    if (await $.session.id() !== sessionId) return;
762    const enrollment = await postDaemon($, "/worker-session", {
763      session_id: sessionId, cwd, client: "claude", declared: true, action: "check",
764    }).catch(() => null);
765    if (enrollment?.enrolled === true && typeof enrollment.warning === "string") {
766      $.ui.log(`[lcm] ${enrollment.warning}`);
767      startSummaryPoller($, summaryCap);
768      return;
769    }
770    $.ui.log(`[lcm] worker mode refused: ${enrollment?.reason ?? "command-hook enrollment could not be confirmed"}; retrying`);
771  }
772}
773
774/** Read the daemon's address and make sure it is listening before the first prompt. */
775async function awaitWorkerEnrollment($: EngineInterface, sessionId: string, cwd: string): Promise<Record<string, unknown> | null> {
776  let enrollment: Record<string, unknown> | null = null;
777  let waited = 0, delay = 250;
778  do {
779    enrollment = await postDaemon($, "/worker-session", { session_id: sessionId, cwd, client: "claude", declared: true, action: "check" }).catch(() => null);
780    if (workerConfirmed(enrollment) || waited >= 30_000) break;
781    const pause = Math.min(delay, 30_000 - waited);
782    await summaryDelay($, pause); waited += pause; delay = Math.min(delay * 2, 2_000);
783  } while (true);
784  return enrollment;
785}
786function workerConfirmed(enrollment: Record<string, unknown> | null): boolean {
787  return enrollment?.enrolled === true && typeof enrollment.warning === "string";
788}
789async function startupWorker($: EngineInterface, sessionId: string, summaryCap: number): Promise<boolean> {
790  if (!await workerDeclared($)) return true;
791  const cwd = await $.session.cwd(), enrollment = await awaitWorkerEnrollment($, sessionId, cwd);
792  if (workerConfirmed(enrollment)) { $.ui.log(`[lcm] ${enrollment!.warning}`); return true; }
793  $.ui.log(`[lcm] worker mode refused: ${enrollment?.reason ?? "command-hook enrollment could not be confirmed"}; retrying`);
794  void retryWorkerEnrollment($, sessionId, cwd, summaryCap).catch(error => { $.ui.log(`[lcm] worker enrollment retry stopped: ${String(error)}`); });
795  return false;
796}
797function probeSessionDaemon($: EngineInterface): void {
798  void readHostEnv($).then(({ port }) => {
799    const startedAt = Date.now();
800    return fetchDaemon($, `http://127.0.0.1:${port}/health`, undefined, HEALTH_PROBE_TIMEOUT_MS).then(() => undefined,
801      async (error: unknown) => {
802        if (isDaemonUnreachableError(error, startedAt)) await startDaemon($);
803        else reportBusyDaemon($);
804      });
805  });
806}
807function registerSessionStart(on: On, summaryCap: number): void {
808  on("session.start", async ($, e, next) => {
809    const sessionId = await $.session.id();
810    const mayServe = await startupWorker($, sessionId, summaryCap);
811    await claimSession($, sessionId, "session.start");
812    await flushHookObservations($, sessionId);
813    probeSessionDaemon($);
814    // The housekeeping the SessionStart command hook awaited. Nothing reads its result,
815    // and the session has no reason to wait for a prune.
816    void $.session.cwd().then(cwd => postDaemon($, "/session-scavenge", { cwd }));
817    // Catch-up sweep for conversations of the same project a prior session left
818    // uncompacted (it ended without SessionEnd). The daemon selects, caps and
819    // fires the actual compaction requests; this call only triggers it.
820    void $.session.cwd().then(cwd => postDaemon($, "/session-start-compact", { cwd, session_id: sessionId }));
821    if (mayServe) startSummaryPoller($, summaryCap);
822    return next(e);
823  });
824}
825
826type RestoreResponse = {
827  context?: string;
828  insights?: { content: string; confidence?: number; tags: string[] }[];
829};
830
831/** The text the SessionStart command hook printed: the daemon's context, plus its insights. */
832function restoreBlockText(body: Record<string, unknown> | null): string {
833  const result = (body ?? {}) as RestoreResponse;
834  const context = result.context ?? "";
835  const insights = result.insights ?? [];
836  if (insights.length === 0) return context;
837  const seen = new Set<string>();
838  const lines = insights.filter((i) => !seen.has(i.content) && seen.add(i.content))
839    .map((i) => `- ${i.content}${typeof i.confidence === "number" ? ` (confidence: ${i.confidence})` : ""}`).join("\n");
840  return `${context}\n<learned-insights source="passive-capture">\n`
841    + `Recent learnings from your previous sessions:\n${lines}\n</learned-insights>`;
842}
843
844/**
845 * `prompt.context` replaces the SessionStart command hook's stdout. It fires once per
846 * conversation and again after compaction and `/clear` — the three moments the command
847 * hook ran — but carries no reason for firing, so the daemon decides which content to
848 * return from the mark `/compact` left for this session.
849 */
850function registerRestoreContext(on: On): void {
851  on("prompt.context", async ($, e, next) => {
852    const { blocks } = await next(e);
853    const [session_id, cwd] = await Promise.all([$.session.id(), $.session.cwd()]);
854    const outcome = await postDaemonOutcome($, "/restore", { session_id, cwd });
855    const restored = outcome.body;
856    const text = restoreBlockText(restored);
857    noteHook(session_id, "prompt.context", "restore", "delivery", ...delivery(outcome));
858    if (restored) noteHook(session_id, "prompt.context", "restore", "execution", "completed",
859      text.trim() ? "context" : "no-context");
860    if (!text.trim()) return { blocks };
861    return { blocks: [...blocks, { name: "lcm", text }] };
862  });
863}
864
865/**
866 * The learning instruction goes into the system prompt once, cached for the session,
867 * instead of riding on every prompt as the command hook's stdout did.
868 */
869function registerLearningInstruction(on: On): void {
870  on("prompt.section", { name: "memory" }, async ($, e, next) => {
871    const section = await next(e);
872    noteHook(await $.session.id(), "prompt.section", "instruction", "execution", "completed");
873    return { text: `${section.text ?? ""}\n\n${LEARNING_INSTRUCTION}` };
874  });
875}
876
877/** Memory hits ride as hidden context on the prompt; the user never sees them. */
878function registerPromptSearch(on: On): void {
879  on("prompt.submit", async ($, e, next) => {
880    if (!e.text.trim()) {
881      noteHook(await $.session.id(), "prompt.submit", "search", "execution", "skipped", "empty-prompt");
882      return next(e);
883    }
884    const [result, search] = await Promise.all([
885      next(e),
886      Promise.all([$.session.id(), $.session.cwd()]).then(([session_id, cwd]) =>
887        postDaemonOutcome($, "/prompt-search", {
888          query: e.text, cwd, session_id,
889          learningInstructionBytes: 0, // the instruction lives in prompt.section now, not in this budget
890          recordEvents: true,
891          format: "context",
892        })),
893    ]);
894    const sessionId = await $.session.id();
895    noteHook(sessionId, "prompt.submit", "search", "delivery", ...delivery(search));
896    if (search.body) noteHook(sessionId, "prompt.submit", "search", "execution", "completed",
897      typeof search.body.context === "string" ? "context" : "no-context");
898    if (result.drop !== undefined) return result;
899    const context = typeof search.body?.context === "string" ? search.body.context : null;
900    if (!context) return result;
901    return { ...result, context: [...(result.context ?? []), context] };
902  });
903}
904
905/**
906 * The session's transcript is ingested incrementally as turns end, so memory does not wait
907 * for SessionEnd. The daemon derives the transcript path from session id and cwd.
908 */
909function registerTurnIngest(on: On): void {
910  on("turn.complete", async ($, e, next) => {
911    const result = await next(e);
912    const session_id = await $.session.id();
913    const now = Date.now();
914    if (now - lastIngestAt < INGEST_INTERVAL_MS) {
915      noteHook(session_id, "turn.complete", "capture", "execution", "deferred", "throttled");
916      await flushHookObservations($, session_id);
917      return result;
918    }
919    lastIngestAt = now;
920    const cwd = await $.session.cwd();
921    const captured = await postDaemonOutcome($, "/ingest", { session_id, cwd });
922    noteHook(session_id, "turn.complete", "capture", "delivery", ...delivery(captured));
923    if (captured.body) noteHook(session_id, "turn.complete", "capture", "execution", "completed");
924    const promoted = await postDaemonOutcome($, "/promote-events", { cwd });
925    noteHook(session_id, "turn.complete", "promote-events", "delivery", ...delivery(promoted));
926    await flushHookObservations($, session_id);
927    return result;
928  });
929}
930
931/** The same body the PostToolUse command hook reads on stdin, so both paths write one row shape. */
932function toolEventPayload(
933  event: Record<string, unknown>, result: { isError?: boolean; result?: unknown; text?: string },
934  session_id: string, cwd: string,
935): Record<string, unknown> {
936  const { tool, tool_use_id, ...tool_input } = event as { tool: string; tool_use_id: string } & Record<string, unknown>;
937  const failed = result.isError === true;
938  return {
939    session_id, cwd, tool_use_id, harness: "claude-function",
940    tool_name: tool,
941    tool_input,
942    tool_response: result.result,
943    tool_output: failed ? { isError: true } : undefined,
944    hook_event_name: failed ? "PostToolUseFailure" : "PostToolUse",
945    error: failed ? result.text : undefined,
946  };
947}
948
949/** One hook replaces both PostToolUse and PostToolUseFailure; the daemon writes the rows. */
950function registerToolCapture(on: On): void {
951  on("tool.call", async ($, e, next) => {
952    const captured = CAPTURED_TOOLS.has(e.tool) || e.tool.startsWith("mcp__");
953    if (!captured) {
954      const result = await next(e);
955      noteHook(await $.session.id(), "tool.call", "tool-capture", "execution", "skipped", "unmatched-tool");
956      return result;
957    }
958
959    const result = await next(e);
960    if (result.deny !== undefined) {
961      noteHook(await $.session.id(), "tool.call", "tool-capture", "execution", "skipped", "denied");
962      return result;
963    }
964
965    const [session_id, cwd] = await Promise.all([$.session.id(), $.session.cwd()]);
966    const recorded = await postDaemonOutcome($, "/tool-event", toolEventPayload(e, result, session_id, cwd));
967    noteHook(session_id, "tool.call", "tool-capture", "delivery", ...delivery(recorded));
968    if (recorded.body) noteHook(session_id, "tool.call", "tool-capture", "execution", "completed",
969      typeof recorded.body.recorded === "number" && recorded.body.recorded > 0 ? "events" : "no-match");
970    return result;
971  });
972}
973
974export const register: Register = (on, options) => {
975  const summaryCap = normalizedSummaryCap(options.sessionSummarizerMaxOutputTokens);
976  const shadow = options.compactionShadow === true ? shadowSessions : undefined;
977  registerSessionStart(on, summaryCap);
978  registerSessionClaim(on, shadow);
979  registerSessionEnd(on, shadow);
980  registerRestoreContext(on);
981  registerLearningInstruction(on);
982  registerPromptSearch(on);
983  registerTurnIngest(on);
984  registerToolCapture(on);
985  if (options.compactionShadow === true) registerCompactionShadow(on, summaryCap);
986};
987const shadowSessions = new ShadowSessionState();
988function shadowTransport($: EngineInterface): ShadowTransport {
989  return { post: (route, body, deadline) => postShadow($, { route, body }, deadline), observe: (status, reason, sessionId) => noteHook(sessionId ?? "unknown", "session.compact", "shadow", "execution", status, reason) };
990}
991async function postShadow($: EngineInterface, { route, body }: { route: string; body: unknown }, deadline?: ShadowDeadline): Promise<PostOutcome> {
992  if (!deadline) return postDaemonOutcome($, route, body);
993  deadline.check();
994  const { port, token } = await readHostEnv($);
995  deadline.check();
996  const response = await $.http.fetch(`http://127.0.0.1:${port}${route}`, {
997    method: "POST", headers: { "content-type": "application/json", ...(token ? { authorization: `Bearer ${token}` } : {}) }, body: JSON.stringify(body),
998  });
999  deadline.check();
1000  return { body: response.ok ? JSON.parse(response.text) as Record<string, unknown> : null, connectionFailed: false, httpStatus: response.status };
1001}
1002function shadowEngine($: EngineInterface): ShadowEngine {
1003  return { clock: { after: (ms, callback) => $.clock.after(ms, callback) }, env: { get: () => $.env.get("LCM_SUMMARIZE_WORKER") }, session: { id: () => $.session.id(), cwd: () => $.session.cwd(), model: () => $.session.model() },
1004    model: { fork: request => $.model.fork(request), complete: (request, options) => $.model.complete(request, options) } };
1005}
1006function registerCompactionShadow(on: On, cap: number): void {
1007  on("session.append", async ($, event, next) => {
1008    let sessionId: string;
1009    try { sessionId = await $.session.id(); }
1010    catch { noteHook("unknown", "session.append", "shadow", "execution", "unavailable", "identity"); return next(event); }
1011    const owner = event.agentId === undefined ? shadowSessions.beginAppend(sessionId, event.door) : undefined;
1012    const result = await next(event);
1013    settledShadowAppend(owner, result);
1014    return result;
1015  });
1016  on("session.compact", async ($, event, next) => {
1017    return runCompactionShadow(shadowEngine($), { event, next, state: shadowSessions, cap, signal: next.signal }, shadowTransport($));
1018  });
1019}
1020function settledShadowAppend(owner: ShadowAppend | undefined, result: import("claude-code").EventResult<"session.append">): void {
1021  if (!owner || result.deny !== undefined) return;
1022  if (result.message.role === undefined) return;
1023  const text = result.message.content.filter(block => block.type === "text").map(block => block.text).join("\n");
1024  shadowSessions.storedAppend(owner, { uuid: result.uuid, text });
1025}
1026function normalizedSummaryCap(value: unknown): number {
1027  if (typeof value !== "number" || !Number.isFinite(value)) return DEFAULT_SUMMARY_OUTPUT_CAP;
1028  return Math.min(Number.MAX_SAFE_INTEGER, Math.max(0, Math.floor(value)));
1029}
1030
hooks/model-budget.ts 120 lines
1export type BudgetLease = { maxTokens?: number; allowance: number; settle(outputTokens: number | null, usageUnknown?: boolean): void };
2export class SessionOutputBudget {
3  private spent = 0;
4  private reserved = 0;
5  private unbounded = false;
6  private usageUnknown = false;
7  private waiters: (() => void)[] = [];
8  /** Reservations still awaiting this budget; one of them may take a lease the moment it wakes. */
9  private pending = 0;
10  private release?: () => void;
11  constructor(readonly cap: number) {
12    if (!Number.isSafeInteger(cap) || cap < 0) throw new Error("Invalid session output budget");
13  }
14  async reserveComplete(maxTokens: number, count = 1): Promise<BudgetLease | null> {
15    positiveInteger(maxTokens); positiveInteger(count);
16    return this.awaiting(async () => {
17      while (this.unbounded) await this.changed();
18      const allowance = Math.min(maxTokens, Math.floor(this.snapshot().available / count));
19      if (!allowance) return null;
20      this.reserved += allowance * count;
21      return this.lease(allowance * count, allowance);
22    });
23  }
24  /** A cut-time fork must start now or produce an unavailable outcome. */
25  reserveFork(): BudgetLease | null {
26    if (this.unbounded || this.reserved || !this.snapshot().available) return null;
27    this.unbounded = true;
28    return this.lease(this.snapshot().available);
29  }
30  /** Ordinary condensed jobs may wait for a preceding bounded request. */
31  async reserveExclusive(): Promise<BudgetLease | null> {
32    return this.awaiting(async () => {
33      while (this.unbounded || this.reserved) await this.changed();
34      return this.reserveFork();
35    });
36  }
37  async reserveOrdinary(maxTokens: number, fork: boolean): Promise<BudgetLease | null> {
38    return this.awaiting(async () => {
39      while (this.unbounded || this.reserved) await this.changed();
40      return fork ? this.reserveFork() : this.reserveComplete(maxTokens);
41    });
42  }
43  /**
44   * The session ended: `release` runs once, as soon as no lease or waiting reservation is
45   * outstanding. Whoever still holds this object keeps using it; only the registry forgets it.
46   */
47  retire(release: () => void): void { this.release = release; this.releaseWhenIdle(); }
48  /** The id is in use again, so the pending release no longer applies. */
49  revive(): void { this.release = undefined; }
50  private async awaiting<T>(reserve: () => Promise<T>): Promise<T> {
51    this.pending++;
52    try { return await reserve(); }
53    finally { this.pending--; this.releaseWhenIdle(); }
54  }
55  private releaseWhenIdle(): void {
56    const idle = !this.reserved && !this.unbounded && !this.pending;
57    if (!this.release || !idle) return;
58    const release = this.release; this.release = undefined; release();
59  }
60  snapshot() {
61    return { spent: this.spent, reserved: this.reserved, unbounded: this.unbounded, usageUnknown: this.usageUnknown,
62      available: Math.max(0, this.cap - this.spent - this.reserved), overshoot: Math.max(0, this.spent - this.cap) };
63  }
64  private changed(): Promise<void> { return new Promise(resolve => this.waiters.push(resolve)); }
65  private lease(allowance: number, maxTokens?: number): BudgetLease {
66    let settled = false;
67    return { allowance, ...(maxTokens !== undefined ? { maxTokens } : {}), settle: (output, usageUnknown = false) => {
68      if (settled) throw new Error("Budget reservation already settled");
69      settled = true;
70      const reported = validOutputUsage(output) ? output : null;
71      this.finish({ allowance, maxTokens }, { output: reported, usageUnknown: usageUnknown || reported === null });
72    } };
73  }
74  private finish({ allowance, maxTokens }: { allowance: number; maxTokens?: number }, { output, usageUnknown }: { output: number | null; usageUnknown: boolean }): void {
75    if (maxTokens === undefined) this.unbounded = false;
76    else this.reserved -= allowance;
77    const unknown = usageUnknown || output === null;
78    this.usageUnknown ||= unknown;
79    this.spent += unknown ? Math.max(allowance, output ?? 0) : output!;
80    const waiters = this.waiters; this.waiters = [];
81    waiters.forEach(resolve => resolve());
82    this.releaseWhenIdle();
83  }
84}
85function validOutputUsage(value: unknown): value is number {
86  return typeof value === "number" && Number.isSafeInteger(value) && value >= 0;
87}
88const owners = new WeakMap<object, Map<string, SessionOutputBudget>>();
89const moduleOwner = {};
90/** Dispatch facades need not have stable object identity; this owner does. */
91export function sharedSessionOutputBudget(sessionId: string, cap: number): SessionOutputBudget {
92  return sessionOutputBudget(moduleOwner, sessionId, cap);
93}
94export function sessionOutputBudget(owner: object, sessionId: string, cap: number): SessionOutputBudget {
95  const sessions = owners.get(owner) ?? new Map<string, SessionOutputBudget>();
96  const prior = sessions.get(sessionId);
97  if (prior) {
98    if (prior.cap !== cap) throw new Error("Session output budget changed");
99    prior.revive();
100    return prior;
101  }
102  const budget = new SessionOutputBudget(cap);
103  sessions.set(sessionId, budget); owners.set(owner, sessions);
104  return budget;
105}
106/** Forget a session's budget once nothing is still spending against it. */
107export function retireSessionBudget(sessions: Map<string, SessionOutputBudget>, sessionId: string): void {
108  sessions.get(sessionId)?.retire(() => sessions.delete(sessionId));
109}
110export function endSharedSessionBudget(sessionId: string): void {
111  const sessions = owners.get(moduleOwner);
112  if (sessions) retireSessionBudget(sessions, sessionId);
113}
114export function _sharedSessionBudgetIdsForTesting(): string[] {
115  return [...owners.get(moduleOwner)?.keys() ?? []];
116}
117function positiveInteger(value: number): void {
118  if (!Number.isSafeInteger(value) || value <= 0) throw new Error("Invalid output reservation");
119}
120
hooks/shadow-budget.ts 28 lines
1import { retireSessionBudget, SessionOutputBudget, type BudgetLease } from "./model-budget.js";
2
3/** Shadow reservations and charges never reduce or block ordinary summaries. */
4class ShadowOutputBudget extends SessionOutputBudget {
5  override reserveFork(): BudgetLease | null {
6    return this.snapshot().usageUnknown ? null : super.reserveFork();
7  }
8  override async reserveComplete(maxTokens: number, count = 1): Promise<BudgetLease | null> {
9    if (this.snapshot().usageUnknown) return null;
10    const lease = await super.reserveComplete(maxTokens, count);
11    if (this.snapshot().usageUnknown) { lease?.settle(0); return null; }
12    return lease;
13  }
14}
15const sessions = new Map<string, ShadowOutputBudget>();
16export function shadowSessionOutputBudget(sessionId: string, cap: number): SessionOutputBudget {
17  const existing = sessions.get(sessionId);
18  if (existing) {
19    if (existing.cap !== cap) throw new Error("Session shadow output budget changed");
20    existing.revive();
21    return existing;
22  }
23  const budget = new ShadowOutputBudget(cap);
24  sessions.set(sessionId, budget); return budget;
25}
26export function endShadowSessionBudget(sessionId: string): void { retireSessionBudget(sessions, sessionId); }
27export function _shadowSessionBudgetIdsForTesting(): string[] { return [...sessions.keys()]; }
28
hooks/shadow-deadline.ts 39 lines
1import type { EngineInterface, Timer } from "claude-code";
2
3export const SHADOW_WAIT_MS = 2000;
4export class ShadowInterrupted extends Error {
5  constructor(readonly outcome: "unavailable" | "cancelled") { super(outcome); }
6}
7
8/** One operation-wide deadline; late prerequisites cannot start another stage. */
9export class ShadowDeadline {
10  private closed = false;
11  private interruption?: ShadowInterrupted;
12  private reject!: (error: ShadowInterrupted) => void;
13  private readonly stopped = new Promise<never>((_resolve, reject) => { this.reject = reject; });
14  private readonly expiresAt = Date.now() + SHADOW_WAIT_MS;
15  private readonly timer: Timer;
16  private readonly abort = () => this.stop("cancelled");
17  constructor(clock: Pick<EngineInterface["clock"], "after">, private readonly signal?: AbortSignal) {
18    void this.stopped.catch(() => {});
19    this.timer = clock.after(SHADOW_WAIT_MS, () => this.stop("unavailable"));
20    signal?.addEventListener("abort", this.abort, { once: true });
21    if (signal?.aborted) this.abort();
22  }
23  check(): void {
24    if (Date.now() >= this.expiresAt) this.stop("unavailable");
25    if (this.interruption) throw this.interruption;
26    if (this.closed) throw new ShadowInterrupted("unavailable");
27  }
28  wait<T>(work: Promise<T>): Promise<T> { return Promise.race([work, this.stopped]); }
29  close(): void {
30    this.closed = true;
31    this.signal?.removeEventListener("abort", this.abort);
32    try { this.timer.cancel(); } catch { /* Cleanup never replaces native. */ }
33  }
34  private stop(outcome: ShadowInterrupted["outcome"]): void {
35    if (this.closed || this.interruption) return;
36    this.interruption = new ShadowInterrupted(outcome); this.reject(this.interruption);
37  }
38}
39
hooks/shadow-boundaries.ts 15 lines
1export type ShadowAppend = { sessionId: string; epoch: number; order: number; door: string };
2
3/** Pending or denied appends never replace the last stored model-visible row. */
4export class ShadowBoundaries {
5  private order = 0;
6  private stored = new Map<string, { uuid: string; order: number }>();
7  begin(sessionId: string, epoch: number, door: string): ShadowAppend { return { sessionId, epoch, order: ++this.order, door }; }
8  complete(owner: ShadowAppend, uuid: string): void {
9    if (owner.order > (this.stored.get(owner.sessionId)?.order ?? 0)) this.stored.set(owner.sessionId, { uuid, order: owner.order });
10  }
11  boundary(sessionId: string): string | undefined { return this.stored.get(sessionId)?.uuid; }
12  entries(): [string, string][] { return [...this.stored].map(([sessionId, row]) => [sessionId, row.uuid]); }
13  reset(sessionId: string): void { this.stored.delete(sessionId); }
14}
15
hooks/compaction-shadow.ts 253 lines
1import type { EngineInterface, SessionCompactInput, SessionCompactResult, SessionMessage } from "claude-code";
2import { captureHeaderModel, executeHeaderFork, executeHeaderPair, type HeaderCall, type HeaderOutcome, type SessionModelAtCut } from "./compaction-header.js";
3import { shadowSessionOutputBudget } from "./shadow-budget.js";
4import { ShadowDeadline, ShadowInterrupted } from "./shadow-deadline.js";
5import { ShadowBoundaries, type ShadowAppend } from "./shadow-boundaries.js";
6
7export type ShadowEngine = { clock: Pick<EngineInterface["clock"], "after">; session: Pick<EngineInterface["session"], "id" | "cwd" | "model">; model: Pick<EngineInterface["model"], "fork" | "complete">; env: Pick<EngineInterface["env"], "get"> };
8
9type Message = { role: "user" | "assistant"; text: string; handle?: string };
10export type ShadowTransport = { post(route: string, body: unknown, deadline?: ShadowDeadline): Promise<{ body: Record<string, unknown> | null; httpStatus?: number }>; observe(status: string, reason: string, sessionId?: string): void };
11type Deferred<T> = { promise: Promise<T>; resolve(value: T): void; reject(error: unknown): void };
12function deferred<T>(): Deferred<T> {
13  let resolve!: (value: T) => void, reject!: (error: unknown) => void;
14  const promise = new Promise<T>((yes, no) => { resolve = yes; reject = no; });
15  void promise.catch(() => {});
16  return { promise, resolve, reject };
17}
18type Cut = {
19  epoch: number; shadowAdmission?: "cancelled" | "unavailable"; sessionId: string; cwd: string; cutId: string; snapshotHash: string; sourceHash: string; model: SessionModelAtCut; fork?: HeaderCall;
20  ready: Deferred<HeaderCall>; paired: Deferred<void>; cancel: Deferred<void>; messages: readonly SessionMessage[]; appended: { uuid: string; text: string }[];
21  clock: ShadowEngine["clock"]; signal?: AbortSignal; cancelled: boolean; complete?: HeaderCall; startedAt: number; setupMs: number; nativeMs: number; pairingMs: number; hookMs: number;
22};
23type BoundaryEpoch = { uuid: string; epoch: number };
24export class ShadowSessionState {
25  private epochs = new Map<string, number>();
26  private boundaries = new ShadowBoundaries();
27  private cuts = new Map<string, Set<Cut>>();
28  private tasks = new Set<Promise<void>>();
29  boundary(sessionId: string): string | undefined { return this.boundaries.boundary(sessionId); }
30  epoch(sessionId: string): number { return this.epochs.get(sessionId) ?? 0; }
31  freezeBoundaries(): ReadonlyMap<string, BoundaryEpoch> {
32    return new Map(this.boundaries.entries().map(([sessionId, uuid]) => [sessionId, { uuid, epoch: this.epoch(sessionId) }]));
33  }
34  beginAppend(sessionId: string, door: string): ShadowAppend { return this.boundaries.begin(sessionId, this.epoch(sessionId), door); }
35  storedAppend(owner: ShadowAppend, { uuid, text }: { uuid: string; text: string }): void {
36    if (owner.epoch !== this.epoch(owner.sessionId)) return;
37    this.boundaries.complete(owner, uuid);
38    if (owner.door === "compaction") this.cuts.get(owner.sessionId)?.forEach(cut => cut.appended.push({ uuid, text }));
39  }
40  reset(sessionId: string): void {
41    this.epochs.set(sessionId, this.epoch(sessionId) + 1);
42    this.boundaries.reset(sessionId);
43    this.cuts.get(sessionId)?.forEach(cancelCut);
44  }
45  add(cut: Cut): void { if (this.epoch(cut.sessionId) !== cut.epoch) cancelCut(cut); const cuts = this.cuts.get(cut.sessionId) ?? new Set<Cut>(); cuts.add(cut); this.cuts.set(cut.sessionId, cuts); }
46  remove(cut: Cut): void { this.cuts.get(cut.sessionId)?.delete(cut); }
47  own(task: Promise<void>, transport: ShadowTransport): void {
48    const handled = task.catch(() => { transport.observe("unconfirmed", "background"); }).finally(() => this.tasks.delete(handled));
49    this.tasks.add(handled);
50  }
51}
52function cancelCut(cut: Cut): void {
53  cut.cancelled = true; cut.cancel.resolve(); cut.ready.reject(new Error("dispatch ended"));
54}
55const object = (value: unknown): value is Record<string, unknown> => value !== null && typeof value === "object" && !Array.isArray(value);
56function descriptors(messages: readonly SessionMessage[]): Message[] { return messages.map(row => ({ role: row.role, text: row.text, ...(row.handle ? { handle: row.handle } : {}) })); }
57function jobCall(value: unknown, kind: "fork" | "complete", sourceHash: string): HeaderCall {
58  if (!object(value)) throw new Error("header job unavailable");
59  const prompt = kind === "fork" ? value.forkPrompt : value.completePrompt;
60  const promptHash = kind === "fork" ? value.forkPromptHash : value.promptHash;
61  if (typeof prompt !== "string" || typeof promptHash !== "string" || typeof value.inputHash !== "string") throw new Error("invalid header job");
62  return { prompt, promptHash, inputHash: kind === "fork" ? sourceHash : value.inputHash, evidence: value.evidence as HeaderCall["evidence"] };
63}
64type Invocation = { event: SessionCompactInput; next(event: SessionCompactInput): Promise<SessionCompactResult>; state: ShadowSessionState; cap: number; signal?: AbortSignal };
65type ShadowDispatch = { engine: ShadowEngine; call: Invocation; transport: ShadowTransport };
66export async function runCompactionShadow(engine: ShadowEngine, call: Invocation, transport: ShadowTransport): Promise<SessionCompactResult> {
67  if (call.event.agentId !== undefined || call.event.trigger === "precompute") return call.next(call.event);
68  const startedAt = performance.now();
69  const dispatch = { engine, call, transport };
70  const cut = await captureCut(dispatch);
71  if (!cut) return call.next(call.event);
72  cut.startedAt = startedAt; cut.setupMs = performance.now() - startedAt;
73  if (call.state.epoch(cut.sessionId) !== cut.epoch) cancelCut(cut);
74  const boundTransport: ShadowTransport = { post: transport.post, observe: (status, reason) => transport.observe(status, reason, cut!.sessionId) };
75  const abort = () => cancelCut(cut);
76  call.signal?.addEventListener("abort", abort, { once: true });
77  if (call.signal?.aborted) abort();
78  startOrCancel({ ...dispatch, transport: boundTransport, cut });
79  try { return await nativeAtCut(call, cut, boundTransport); }
80  finally { call.signal?.removeEventListener("abort", abort); }
81}
82function startOrCancel(context: ShadowDispatch & { cut: Cut }): void {
83  const { engine, cut, call, transport } = context;
84  try {
85    if (cut.cancelled) {
86      cut.shadowAdmission = "cancelled"; transport.observe("cancelled", "admission"); call.state.remove(cut);
87    }
88    else startArms(engine, { cut, state: call.state, cap: call.cap }, transport);
89  } catch {
90    transport.observe("unavailable", "executor");
91    deliverUnstarted({ cut, state: call.state, outcome: "unavailable" }, transport);
92  }
93}
94async function captureCut(context: ShadowDispatch): Promise<Cut | undefined> {
95  const { engine, call, transport } = context;
96  let deadline: ShadowDeadline | undefined;
97  const attempt: { cut?: Cut } = {};
98  let accepted = false;
99  try {
100    deadline = new ShadowDeadline(engine.clock, call.signal); deadline.check();
101    const result = await deadline.wait(captureWithinDeadline({ context, deadline, attempt }));
102    deadline.check(); accepted = result !== undefined; return result;
103  } catch (error) {
104    transport.observe(error instanceof ShadowInterrupted ? error.outcome : "unavailable", "admission"); return undefined;
105  } finally {
106    deadline?.close();
107    if (!accepted && attempt.cut) { cancelCut(attempt.cut); call.state.remove(attempt.cut); }
108  }
109}
110async function captureWithinDeadline(input: { context: ShadowDispatch; deadline: ShadowDeadline; attempt: { cut?: Cut } }): Promise<Cut | undefined> {
111  const { context: { engine, call, transport }, deadline, attempt } = input;
112  const frozen = structuredClone(call.event), boundaries = call.state.freezeBoundaries();
113  const [identity, worker] = await Promise.all([
114    Promise.all([engine.session.id(), engine.session.cwd(), captureHeaderModel(engine)]),
115    engine.env.get("LCM_SUMMARIZE_WORKER").catch(() => undefined),
116  ]);
117  deadline.check();
118  if (worker === "1") return undefined;
119  return admit({ identity, frozen, boundaries, state: call.state, clock: engine.clock, signal: call.signal, attempt }, { transport, deadline });
120}
121async function nativeAtCut(call: Invocation, cut: Cut, transport: ShadowTransport): Promise<SessionCompactResult> {
122  let result: SessionCompactResult;
123  const nativeAt = performance.now();
124  try { result = await call.next(call.event); }
125  catch (error) { cut.nativeMs = performance.now() - nativeAt; await pairNative(cut, null, transport); throw error; }
126  cut.nativeMs = performance.now() - nativeAt; await pairNative(cut, result, transport);
127  return result;
128}
129function deliverUnstarted({ cut, state, outcome }: { cut: Cut; state: ShadowSessionState; outcome: "unavailable" | "aborted" }, transport: ShadowTransport): void {
130  state.own(cut.paired.promise.then(async () => {
131    await Promise.all((["A", "B", "C"] as const).map(arm => postArm(cut, outcome === "aborted" ? cancelledArm(cut, arm) : unavailableArm(cut, arm), transport)));
132  }).finally(() => state.remove(cut)), transport);
133}
134function startArms(engine: ShadowEngine, { cut, state, cap }: { cut: Cut; state: ShadowSessionState; cap: number }, transport: ShadowTransport): void {
135  const budget = shadowSessionOutputBudget(cut.sessionId, cap);
136  const fork = cut.fork ? executeHeaderFork(engine, cut.fork, { model: cut.model, budget }).catch(() => unavailableArm(cut, "A")) : Promise.resolve(unavailableArm(cut, "A"));
137  const forkDelivery = deliverFork(fork, cut, transport);
138  state.own(forkDelivery, transport);
139  const input = cut.fork ? cut.ready.promise : cut.ready.promise.then(() => { throw new Error("admitted input unavailable"); });
140  let pairStarted = false;
141  const activeEngine: ShadowEngine = { ...engine, model: { ...engine.model, complete: (request, options) => {
142    if (cut.cancelled) return Promise.reject(new Error("session ended"));
143    pairStarted = true; return engine.model.complete(request, options);
144  } } };
145  const pair = executeHeaderPair(activeEngine, input, { model: cut.model, budget, canStart: () => !cut.cancelled }).catch(() => ({ B: unavailableArm(cut, "B"), C: unavailableArm(cut, "C") }));
146  const pairDelivery = Promise.race([pair, cut.cancel.promise.then(() => pairStarted ? pair : ({ B: cancelledArm(cut, "B"), C: cancelledArm(cut, "C") }))]).then(async outcomes => {
147    await cut.paired.promise;
148    await Promise.all([postArm(cut, outcomes.B, transport), postArm(cut, outcomes.C, transport)]);
149  });
150  state.own(pairDelivery, transport);
151  state.own(Promise.all([forkDelivery, pairDelivery]).then(() => state.remove(cut), () => state.remove(cut)), transport);
152}
153type AdmissionInput = { attempt: { cut?: Cut }; clock: ShadowEngine["clock"]; signal?: AbortSignal; identity: [string, string, SessionModelAtCut]; frozen: SessionCompactInput; boundaries: ReadonlyMap<string, BoundaryEpoch>; state: ShadowSessionState };
154async function admit({ identity: [sessionId, cwd, model], frozen, boundaries, state, clock, signal, attempt }: AdmissionInput, { transport, deadline }: { transport: ShadowTransport; deadline: ShadowDeadline }): Promise<Cut | undefined> {
155  const { messages, instructions, trigger } = frozen;
156  const boundary = boundaries.get(sessionId);
157  if (!boundary) { transport.observe("unavailable", "boundary", sessionId); return undefined; }
158  const cutId = crypto.randomUUID();
159  const cut = pendingCut({ sessionId, cwd, model, cutId, messages, clock, signal }, boundary.epoch);
160  attempt.cut = cut; state.add(cut);
161  if (cut.cancelled) { transport.observe("cancelled", "admission", sessionId); return undefined; }
162  const response = await transport.post("/compaction-shadow/start", { cwd, session_id: sessionId, cut_id: cutId, model: model.id, trigger, instructions,
163    boundary_uuid: boundary.uuid, engine_messages: descriptors(messages), prepare_header: true }, deadline);
164  deadline.check();
165  return admittedCut(response.body, cut, transport);
166}
167type AdmissionIdentity = Pick<Cut, "sessionId" | "cwd" | "model" | "cutId" | "messages" | "clock" | "signal">;
168function admittedCut(body: unknown, pending: Cut, transport: ShadowTransport): Cut | undefined {
169  if (!object(body) || body.admitted !== true) { transport.observe("unavailable", refusalReason(body)); return undefined; }
170  const { cut, snapshot } = admittedBinding(body, pending.cutId);
171  let fork: HeaderCall | undefined;
172  try { fork = jobCall(body.job, "fork", snapshot.sourceHash as string); }
173  catch { transport.observe("unavailable", "header-input"); }
174  pending.snapshotHash = cut.snapshotHash as string; pending.sourceHash = snapshot.sourceHash as string; pending.fork = fork;
175  return pending;
176}
177function pendingCut(identity: AdmissionIdentity, epoch: number): Cut {
178  return { ...identity, epoch, snapshotHash: "", sourceHash: "",
179    ready: deferred<HeaderCall>(), paired: deferred<void>(), cancel: deferred<void>(), appended: [], cancelled: false,
180    startedAt: 0, setupMs: 0, nativeMs: 0, pairingMs: 0, hookMs: 0 };
181}
182function refusalReason(body: unknown): string { return object(body) && typeof body.reason === "string" ? body.reason : "admission"; }
183function admittedBinding(body: Record<string, unknown>, cutId: string) {
184  if (!object(body.cut) || !object(body.snapshot)) throw new Error("cut binding unavailable");
185  if (body.cut.cutId !== cutId || typeof body.cut.snapshotHash !== "string" || typeof body.snapshot.sourceHash !== "string") throw new Error("cut identity mismatch");
186  return { cut: body.cut, snapshot: body.snapshot };
187}
188function unavailableArm(cut: Cut, arm: "A" | "B" | "C"): HeaderOutcome {
189  return { arm, outcome: "unavailable", header: null, text: "", usage: null, requestedModel: arm === "C" ? "sonnet" : cut.model.id,
190    inputHash: null, promptHash: null, durationMs: 0, queueMs: 0 };
191}
192
193function binding(cut: Cut) { return { cwd: cut.cwd, session_id: cut.sessionId, cut_id: cut.cutId, snapshot_hash: cut.snapshotHash }; }
194async function deliverFork(fork: Promise<HeaderOutcome>, cut: Cut, transport: ShadowTransport): Promise<void> {
195  const result = await fork;
196  await cut.paired.promise; await postArm(cut, result, transport);
197}
198function cancelledArm(cut: Cut, arm: "A" | "B" | "C"): HeaderOutcome {
199  const input = arm === "A" ? cut.fork : cut.complete;
200  return { arm, outcome: "aborted", header: null, text: "", usage: null, requestedModel: arm === "C" ? "sonnet" : cut.model.id,
201    inputHash: input?.inputHash ?? null, promptHash: input?.promptHash ?? null, durationMs: 0, queueMs: 0 };
202}
203async function postArm(cut: Cut, result: HeaderOutcome, transport: ShadowTransport): Promise<void> {
204  const timings = { setupMs: cut.setupMs, nativeMs: cut.nativeMs, pairingMs: cut.pairingMs, hookMs: cut.hookMs };
205  const response = await transport.post("/compaction-shadow/arm", { ...binding(cut), arm: result.arm, attempt_id: "first", record: { ...result, timings, costUsd: null, usageAttempts: [] } });
206  if (!response.body?.stored) transport.observe("unconfirmed", "arm-delivery");
207}
208function extractedNative(cut: Cut, result: SessionCompactResult | null) {
209  if (!result || result.skip !== undefined) return { text: "", outcome: result ? "skipped" : "aborted", fidelity: result ? "skipped" : "aborted", tail: [], observedMessages: [], candidateIndices: [] };
210  const original = new Map(cut.messages.map((message, index) => [message.handle, { message, index }]));
211  const candidates = result.messages.flatMap((message, index) => !message.handle || !original.has(message.handle) ? [index] : []);
212  const observed = { observedMessages: descriptors(result.messages), candidateIndices: candidates, usage: result.usage, tokensBefore: result.tokensBefore, tokensAfter: result.tokensAfter };
213  if (original.size !== cut.messages.length || original.has(undefined)) return { text: "", outcome: "unavailable", fidelity: "native-tail-unverified", tail: [], ...observed };
214  if (!uniqueSummary(candidates)) return { text: "", outcome: "unavailable", fidelity: "native-summary-unverified", tail: [], ...observed };
215  const tail = result.messages.slice(1);
216  if (!tailMatches(tail, original)) return { text: "", outcome: "unavailable", fidelity: "native-tail-unverified", tail: [], ...observed };
217  const summary = result.messages[0];
218  return { text: summary.text, outcome: "answered", fidelity: "verified", tail: descriptors(tail), ...observed,
219    ...nativeAppendIdentity(cut, summary.text) };
220}
221function nativeAppendIdentity(cut: Cut, text: string): { summaryUuid?: string } {
222  const rows = cut.appended.filter(row => row.text === text);
223  return rows.length === 1 ? { summaryUuid: rows[0].uuid } : {};
224}
225function uniqueSummary(indices: readonly number[]): boolean { return indices.length === 1 && indices[0] === 0; }
226
227function tailMatches(tail: readonly SessionMessage[], original: ReadonlyMap<string | undefined, { message: SessionMessage; index: number }>): boolean {
228  let prior = -1;
229  return tail.every(message => {
230    const source = original.get(message.handle);
231    if (!source || source.index <= prior) return false;
232    if (!sameMessage(source.message, message)) return false;
233    prior = source.index; return true;
234  });
235}
236function sameMessage(before: SessionMessage, after: SessionMessage): boolean { return before.role === after.role && before.text === after.text; }
237
238async function pairNative(cut: Cut, result: SessionCompactResult | null, transport: ShadowTransport): Promise<void> {
239  const pairingAt = performance.now();
240  let deadline: ShadowDeadline | undefined;
241  try {
242    deadline = new ShadowDeadline(cut.clock, cut.signal); deadline.check();
243    const record = { ...extractedNative(cut, result), durationMs: cut.nativeMs, ...(cut.shadowAdmission ? { shadowAdmission: cut.shadowAdmission } : {}) };
244    const response = await deadline.wait(transport.post("/compaction-shadow/native", { ...binding(cut), record, prepare_header: !cut.shadowAdmission }, deadline));
245    deadline.check();
246    if (cut.shadowAdmission) return;
247    if (record.fidelity !== "verified" || !response.body?.stored || !response.body.job) throw new Error("native input unavailable");
248    cut.complete = jobCall(response.body.job, "complete", cut.sourceHash);
249    cut.ready.resolve(cut.complete);
250  } catch { cut.ready.reject(new Error("native pairing unavailable")); transport.observe("unavailable", "native-pairing"); }
251  finally { deadline?.close(); cut.pairingMs = performance.now() - pairingAt; cut.hookMs = performance.now() - cut.startedAt; cut.paired.resolve(); }
252}
253
hooks/compaction-header.ts 138 lines
1import type { EngineInterface } from "claude-code";
2import type { BudgetLease, SessionOutputBudget } from "./model-budget.js";
3import { validCompactionHeader, validModelName, type CompactionHeader } from "./compaction-header-schema.js";
4import { freezeCitationEvidence, resolveHeaderCitations, type CitationEvidence, type HeaderCitations } from "./header-citations.js";
5
6type HeaderEngine = { session: Pick<EngineInterface["session"], "model">; model: Pick<EngineInterface["model"], "fork" | "complete"> };
7export type HeaderCall = { prompt: string; inputHash: string; promptHash: string; maxTokens?: number; evidence?: CitationEvidence };
8const CAPTURED_MODEL = Symbol("session model at cut");
9export type SessionModelAtCut = { readonly id: string; readonly [CAPTURED_MODEL]: true };
10export type HeaderUsage = { input_tokens: number; output_tokens: number; cache_read_input_tokens: number; cache_creation_input_tokens: number };
11type Arm = "A" | "B" | "C";
12type HeaderFailure = "api-error" | "empty-reply" | "aborted" | "nothing-to-fork" | "invalid-output" | "unavailable" | "spend-cap" | "unconfirmed";
13type OutcomeMetadata = {
14  arm: Arm; requestedModel: string; usage: HeaderUsage | null; inputHash: string | null; promptHash: string | null;
15  durationMs: number; queueMs: number; options?: { maxTokens: number }; status?: number | null; errorKind?: string;
16  budget?: ReturnType<SessionOutputBudget["snapshot"]>;
17  citations?: HeaderCitations; usageUnknown?: boolean; refusalReason?: "usageUnknown" | "spendCap";
18};
19export type HeaderOutcome = OutcomeMetadata & ({ outcome: "answered"; text: string; header: CompactionHeader } | { outcome: HeaderFailure; text: string; header: null });
20type ExecutionContext = { model: SessionModelAtCut; budget: SessionOutputBudget; canStart?: () => boolean };
21export interface HeaderExecutor {
22  captureModel($: HeaderEngine): Promise<SessionModelAtCut>;
23  fork($: HeaderEngine, job: HeaderCall, context: ExecutionContext): Promise<HeaderOutcome>;
24  pair($: HeaderEngine, job: Promise<HeaderCall>, context: ExecutionContext): Promise<{ B: HeaderOutcome; C: HeaderOutcome }>;
25}
26const DEFAULT_HEADER_MAX_TOKENS = 4096;
27const MAX_COMPLETION_TOKENS = 64_000;
28const CLEAN_HEADER_ARMS = 2;
29const object = (value: unknown): value is Record<string, unknown> => value !== null && typeof value === "object" && !Array.isArray(value);
30export async function captureHeaderModel($: HeaderEngine): Promise<SessionModelAtCut> {
31  const id = await $.session.model();
32  if (!validModelName(id)) throw new Error("Invalid session model identifier");
33  return Object.freeze({ id, [CAPTURED_MODEL]: true as const });
34}
35type CallContext = { arm: Arm; model: string; job?: HeaderCall; queuedAt: number; startedAt: number; maxTokens?: number };
36function metadata(call: CallContext, usage: HeaderUsage | null): OutcomeMetadata {
37  return { arm: call.arm, requestedModel: call.model, inputHash: call.job?.inputHash ?? null, promptHash: call.job?.promptHash ?? null, usage,
38    durationMs: Date.now() - call.startedAt, queueMs: call.startedAt - call.queuedAt,
39    ...(call.maxTokens !== undefined ? { options: { maxTokens: call.maxTokens } } : {}) };
40}
41function failure(call: CallContext, outcome: HeaderFailure, usage: HeaderUsage | null = null): HeaderOutcome {
42  return { ...metadata(call, usage), outcome, header: null, text: "" };
43}
44function readUsage(value: unknown): HeaderUsage | null {
45  if (!object(value)) return null;
46  const keys = ["input_tokens", "output_tokens", "cache_read_input_tokens", "cache_creation_input_tokens"] as const;
47  if (!keys.every(key => Number.isSafeInteger(value[key]) && (value[key] as number) >= 0)) return null;
48  return { input_tokens: value.input_tokens as number, output_tokens: value.output_tokens as number,
49    cache_read_input_tokens: value.cache_read_input_tokens as number, cache_creation_input_tokens: value.cache_creation_input_tokens as number };
50}
51function answered(call: CallContext, text: string, usage: HeaderUsage): HeaderOutcome {
52  if (!text.trim()) return failure(call, "empty-reply", usage);
53  let header: unknown;
54  try { header = JSON.parse(text); }
55  catch { return { ...failure(call, "invalid-output", usage), text }; }
56  if (!validCompactionHeader(header)) return { ...failure(call, "invalid-output", usage), text };
57  return { ...metadata(call, usage), outcome: "answered", text, header, citations: resolveHeaderCitations(header, call.job?.evidence) };
58}
59function classifyResult(result: unknown, call: CallContext): HeaderOutcome {
60  const outcome = observedResult(result, call);
61  return { ...outcome, usageUnknown: outcome.usage === null && outcome.outcome !== "nothing-to-fork" };
62}
63function observedResult(result: unknown, call: CallContext): HeaderOutcome {
64  if (!object(result)) return failure(call, "unconfirmed");
65  if (nothingToFork(result, call.arm)) return failure(call, "nothing-to-fork");
66  const usage = readUsage(result.usage);
67  if (!usage) return failure(call, "unconfirmed");
68  if (result.isAnswered === true && typeof result.text === "string") return answered(call, result.text, usage);
69  return nonAnswer(result, call, usage);
70}
71function nothingToFork(result: Record<string, unknown>, arm: Arm): boolean {
72  return arm === "A" && result.isAnswered === false && result.reason === "nothing-to-fork";
73}
74function nonAnswer(result: Record<string, unknown>, call: CallContext, usage: HeaderUsage): HeaderOutcome {
75  if (result.isAnswered !== false) return failure(call, "invalid-output", usage);
76  if (result.reason === "aborted" || result.reason === "empty-reply") return failure(call, result.reason, usage);
77  if (result.reason !== "api-error") return failure(call, "unconfirmed", usage);
78  const status = typeof result.status === "number" && result.status >= 100 && result.status <= 599 && Number.isInteger(result.status) ? result.status : null;
79  const errorKind = typeof result.error === "string" && /^[a-z_]{1,80}$/.exec(result.error)?.[0] === result.error ? result.error : "unknown";
80  return { ...failure(call, "api-error", usage), status, errorKind };
81}
82function outputAccounting(results: readonly HeaderOutcome[]): { known: number; unknown: boolean } {
83  return { known: results.reduce((total, row) => total + (row.usage?.output_tokens ?? 0), 0),
84    unknown: results.some(row => row.usage === null && row.outcome !== "nothing-to-fork") };
85}
86export async function executeHeaderFork($: HeaderEngine, job: HeaderCall, { model, budget }: ExecutionContext): Promise<HeaderOutcome> {
87  const now = Date.now(), fixed = { ...job, evidence: freezeCitationEvidence(job.evidence) }, call: CallContext = { arm: "A", model: model.id, job: fixed, queuedAt: now, startedAt: now };
88  const lease = budget.reserveFork();
89  if (!lease) return { ...failure(call, budget.snapshot().available && !budget.snapshot().usageUnknown ? "unavailable" : "spend-cap"), ...refusal(budget), budget: budget.snapshot() };
90  let result: unknown;
91  try { result = await $.model.fork({ prompt: fixed.prompt }); }
92  catch { result = null; }
93  const outcome = classifyResult(result, call), accounted = outputAccounting([outcome]);
94  lease.settle(accounted.known, accounted.unknown);
95  return { ...outcome, budget: budget.snapshot() };
96}
97function refusal(budget: SessionOutputBudget): Pick<OutcomeMetadata, "refusalReason"> {
98  if (budget.snapshot().usageUnknown) return { refusalReason: "usageUnknown" };
99  return budget.snapshot().available ? {} : { refusalReason: "spendCap" };
100}
101async function completeMember($: HeaderEngine, call: CallContext): Promise<HeaderOutcome> {
102  let result: unknown;
103  try { result = await $.model.complete({ model: call.model, prompt: call.job!.prompt, maxTokens: call.maxTokens! }); }
104  catch { result = null; }
105  return classifyResult(result, call);
106}
107function unavailablePair(model: SessionModelAtCut, queuedAt: number, { job, outcome = "unavailable" }: { job?: HeaderCall; outcome?: HeaderFailure } = {}) {
108  const startedAt = Date.now();
109  return { B: failure({ arm: "B", model: model.id, job, queuedAt, startedAt }, outcome), C: failure({ arm: "C", model: "sonnet", job, queuedAt, startedAt }, outcome) };
110}
111export async function executeHeaderPair($: HeaderEngine, ready: Promise<HeaderCall>, { model, budget, canStart }: ExecutionContext): Promise<{ B: HeaderOutcome; C: HeaderOutcome }> {
112  const queuedAt = Date.now();
113  let job: HeaderCall;
114  try { const input = await ready; job = Object.freeze({ ...input, evidence: freezeCitationEvidence(input.evidence) }); }
115  catch { return unavailablePair(model, queuedAt, { outcome: canStart?.() === false ? "aborted" : "unavailable" }); }
116  const requested = job.maxTokens ?? DEFAULT_HEADER_MAX_TOKENS;
117  if (!Number.isSafeInteger(requested) || requested <= 0 || requested > MAX_COMPLETION_TOKENS) return unavailablePair(model, queuedAt, { job });
118  if (canStart?.() === false) return unavailablePair(model, queuedAt, { job, outcome: "aborted" });
119  const reservation = await reserveHeaderPair({ budget, canStart }, requested);
120  if (typeof reservation === "string") return refusedPair({ model, queuedAt, job, outcome: reservation }, budget);
121  const lease = reservation;
122  const common = { job, queuedAt, startedAt: Date.now(), maxTokens: lease.maxTokens! };
123  const [B, C] = await Promise.all([completeMember($, { ...common, arm: "B", model: model.id }), completeMember($, { ...common, arm: "C", model: "sonnet" })]);
124  const accounted = outputAccounting([B, C]); lease.settle(accounted.known, accounted.unknown);
125  return { B: { ...B, budget: budget.snapshot() }, C: { ...C, budget: budget.snapshot() } };
126}
127async function reserveHeaderPair({ budget, canStart }: Pick<ExecutionContext, "budget" | "canStart">, requested: number): Promise<BudgetLease | "aborted" | "spend-cap"> {
128  const lease = await budget.reserveComplete(requested, CLEAN_HEADER_ARMS);
129  if (canStart?.() === false) { lease?.settle(0); return "aborted"; }
130  return lease ?? "spend-cap";
131}
132function refusedPair(input: { model: SessionModelAtCut; queuedAt: number; job: HeaderCall; outcome: "aborted" | "spend-cap" }, budget: SessionOutputBudget) {
133  const pair = unavailablePair(input.model, input.queuedAt, { job: input.job, outcome: input.outcome });
134  if (input.outcome === "aborted") return pair;
135  return { B: { ...pair.B, ...refusal(budget), budget: budget.snapshot() }, C: { ...pair.C, ...refusal(budget), budget: budget.snapshot() } };
136}
137export const headerExecutor: HeaderExecutor = { captureModel: captureHeaderModel, fork: executeHeaderFork, pair: executeHeaderPair };
138
hooks/compaction-header-schema.ts 84 lines
1export const COMPACTION_HEADER_SECTIONS = ["intent", "instructionsInForce", "decisions", "taskState", "procedure", "nextSteps", "openThreads", "files", "errors"] as const;
2export type HeaderSource = string | { quote: string };
3export type HeaderProvenance = "authorized by the user" | "proposed by the assistant" | "observed" | "unresolved";
4export type WorkingItem = { text: string; sources: HeaderSource[]; state?: "reported" | "confirmed" | "unknown" };
5export type CompactionHeader = {
6  version: 2; intent: WorkingItem[]; instructionsInForce: { sources: string[] }[];
7  decisions: (WorkingItem & { supersedes?: HeaderSource[] })[];
8  taskState: (WorkingItem & { status: "done" | "in progress" | "blocked"; provenance: HeaderProvenance })[];
9  procedure: WorkingItem[]; nextSteps: (WorkingItem & { provenance: HeaderProvenance })[];
10  openThreads: WorkingItem[]; files: (WorkingItem & { status: string })[]; errors: (WorkingItem & { fix: string })[];
11};
12export type CompactionHeaderItem = CompactionHeader[typeof COMPACTION_HEADER_SECTIONS[number]][number];
13export function compactionHeaderItems(header: CompactionHeader): CompactionHeaderItem[] {
14  return COMPACTION_HEADER_SECTIONS.flatMap<CompactionHeaderItem>(key => header[key]);
15}
16const PROVENANCE = ["authorized by the user", "proposed by the assistant", "observed", "unresolved"];
17export function validModelName(value: unknown): value is string {
18  return typeof value === "string" && /^[A-Za-z0-9][A-Za-z0-9._:/-]{0,63}$/.exec(value)?.[0] === value;
19}
20const record = (value: unknown): value is Record<string, unknown> => value !== null && typeof value === "object" && !Array.isArray(value);
21const text = (value: unknown): value is string => typeof value === "string" && value.trim().length > 0;
22const only = (value: Record<string, unknown>, keys: readonly string[]) => Object.keys(value).every(key => keys.includes(key));
23export function validHeaderSource(value: unknown): value is HeaderSource {
24  if (record(value)) return only(value, ["quote"]) && text(value.quote);
25  if (typeof value !== "string") return false;
26  const match = /^\[(?:excerpt:[A-Za-z0-9_-]{1,140}|sum:sum_[A-Za-z0-9_-]{1,136}|raw:[A-Za-z0-9_-]{1,140}:[1-9]\d*)\]$/.exec(value);
27  if (match?.[0] !== value) return false;
28  const raw = /^\[raw:[A-Za-z0-9_-]+:([1-9]\d*)\]$/.exec(value);
29  return !raw || Number.isSafeInteger(Number(raw[1]));
30}
31function sourced(value: unknown): value is Record<string, unknown> & { sources: HeaderSource[] } {
32  return record(value) && Array.isArray(value.sources) && value.sources.length > 0 && value.sources.every(validHeaderSource);
33}
34function working(value: unknown): value is Record<string, unknown> & WorkingItem {
35  if (!sourced(value) || !text(value.text)) return false;
36  if (value.state !== undefined && !["reported", "confirmed", "unknown"].includes(value.state as string)) return false;
37  const summaryOnly = value.sources.every(source => typeof source === "string" && source.startsWith("[sum:"));
38  return !summaryOnly || value.state === "reported";
39}
40function instruction(value: unknown): boolean {
41  return sourced(value) && only(value, ["sources"]) && value.sources.every(source => typeof source === "string" && source.startsWith("[excerpt:"));
42}
43const ITEM_KEYS = ["text", "sources", "state"];
44const plain = (value: unknown): boolean => working(value) && only(value, ITEM_KEYS);
45function decision(value: unknown): boolean {
46  if (!working(value) || !only(value, [...ITEM_KEYS, "supersedes"])) return false;
47  return value.supersedes === undefined || Array.isArray(value.supersedes) && value.supersedes.every(validHeaderSource);
48}
49function task(value: unknown): boolean {
50  return working(value) && only(value, [...ITEM_KEYS, "status", "provenance"]) &&
51    ["done", "in progress", "blocked"].includes(value.status as string) && PROVENANCE.includes(value.provenance as string);
52}
53function nextStep(value: unknown): boolean {
54  return working(value) && only(value, [...ITEM_KEYS, "provenance"]) && PROVENANCE.includes(value.provenance as string);
55}
56function file(value: unknown): boolean { return working(value) && only(value, [...ITEM_KEYS, "status"]) && text(value.status); }
57function error(value: unknown): boolean { return working(value) && only(value, [...ITEM_KEYS, "fix"]) && text(value.fix); }
58const validators: Record<typeof COMPACTION_HEADER_SECTIONS[number], (value: unknown) => boolean> = {
59  intent: plain, instructionsInForce: instruction, decisions: decision, taskState: task, procedure: plain,
60  nextSteps: nextStep, openThreads: plain, files: file, errors: error,
61};
62export function validCompactionHeader(value: unknown): value is CompactionHeader {
63  if (!record(value) || value.version !== 2 || !only(value, ["version", ...COMPACTION_HEADER_SECTIONS])) return false;
64  return COMPACTION_HEADER_SECTIONS.every(key => Array.isArray(value[key]) && value[key].every(validators[key]));
65}
66/** String identifiers and typed state retain identity; only free text is mapped. */
67export function mapCompactionHeaderText(value: CompactionHeader, map: (text: string) => string): CompactionHeader {
68  const result = JSON.parse(JSON.stringify(value)) as CompactionHeader;
69  for (const key of COMPACTION_HEADER_SECTIONS) for (const item of result[key]) {
70    mapWorkingItem({ item, file: key === "files" }, map);
71  }
72  return result;
73}
74function mapWorkingItem({ item, file }: { item: CompactionHeaderItem; file: boolean }, map: (text: string) => string): void {
75  item.sources = item.sources.map(source => mapSourceQuote(source, map)) as typeof item.sources;
76  if ("text" in item) item.text = map(item.text);
77  if ("fix" in item) item.fix = map(item.fix);
78  if (file && "status" in item) item.status = map(item.status) as typeof item.status;
79  if ("supersedes" in item) item.supersedes = item.supersedes?.map(source => mapSourceQuote(source, map));
80}
81function mapSourceQuote(source: HeaderSource, map: (text: string) => string): HeaderSource {
82  return typeof source === "string" ? source : { quote: map(source.quote) };
83}
84
hooks/header-citations.ts 53 lines
1import { COMPACTION_HEADER_SECTIONS, type CompactionHeader, type CompactionHeaderItem, type HeaderSource } from "./compaction-header-schema.js";
2
3export type CitationEvidence = { cutId: string; originals: readonly { id: number; text: string }[]; excerpts: readonly { id: string; rawMessageId: number }[]; summaries: readonly string[] };
4export type CitationStatus = "resolved" | "missing" | "ambiguous";
5export type CitationCheck = { field: "sources" | "supersedes"; index: number; source: HeaderSource; status: CitationStatus; originalIds: number[] };
6export type ItemCitations = { section: typeof COMPACTION_HEADER_SECTIONS[number]; item: number; status: CitationStatus; citations: CitationCheck[] };
7export type HeaderCitations = { version: 1; items: ItemCitations[] };
8const object = (value: unknown): value is Record<string, unknown> => value !== null && typeof value === "object" && !Array.isArray(value);
9/** The live path uses this predicate; absent, missing or ambiguous evidence refuses. */
10export function allCitationsResolved(value: HeaderCitations | null | undefined): boolean {
11  return value?.version === 1 && value.items.every(item => item.status === "resolved" && item.citations.every(citation => citation.status === "resolved"));
12}
13export function resolveHeaderCitations(header: CompactionHeader, evidence?: CitationEvidence): HeaderCitations {
14  return { version: 1, items: COMPACTION_HEADER_SECTIONS.flatMap(section => header[section].map((item, index) => resolveItem({ section, item: index }, item, evidence))) };
15}
16function resolveItem(address: Pick<ItemCitations, "section" | "item">, item: CompactionHeaderItem, evidence?: CitationEvidence): ItemCitations {
17  const citations = item.sources.map((source, index) => resolveCitation({ source, field: "sources", index }, evidence));
18  if ("supersedes" in item) citations.push(...(item.supersedes ?? []).map((source, index) => resolveCitation({ source, field: "supersedes", index }, evidence)));
19  const status = itemStatus(citations);
20  return { ...address, status, citations };
21}
22function itemStatus(citations: readonly CitationCheck[]): CitationStatus {
23  if (citations.some(citation => citation.status === "ambiguous")) return "ambiguous";
24  return citations.every(citation => citation.status === "resolved") ? "resolved" : "missing";
25}
26function resolveCitation(citation: Pick<CitationCheck, "source" | "field" | "index">, evidence?: CitationEvidence): CitationCheck {
27  if (!evidence) return { ...citation, status: "missing", originalIds: [] };
28  const { source } = citation;
29  if (typeof source !== "string") return matchResult(citation, source.quote ? evidence.originals.filter(row => row.text.includes(source.quote)).map(row => row.id) : []);
30  const raw = /^\[raw:([A-Za-z0-9_-]+):([1-9]\d*)\]$/.exec(source);
31  if (raw) return matchResult(citation, raw[1] === evidence.cutId ? evidence.originals.filter(row => row.id === Number(raw[2])).map(row => row.id) : []);
32  return resolveNamed(citation, evidence);
33}
34function resolveNamed(citation: Pick<CitationCheck, "source" | "field" | "index">, evidence: CitationEvidence): CitationCheck {
35  const source = citation.source as string;
36  const excerpt = /^\[excerpt:([A-Za-z0-9_-]+)\]$/.exec(source);
37  if (excerpt) {
38    const rows = evidence.excerpts.filter(row => row.id === excerpt[1]);
39    return matchResult(citation, rows.flatMap(row => evidence.originals.filter(original => original.id === row.rawMessageId).map(original => original.id)));
40  }
41  const summary = /^\[sum:(sum_[A-Za-z0-9_-]+)\]$/.exec(source);
42  const count = summary ? evidence.summaries.filter(id => id === summary[1]).length : 0;
43  return { ...citation, status: matchStatus(count), originalIds: [] };
44}
45function matchResult(citation: Pick<CitationCheck, "source" | "field" | "index">, ids: number[]): CitationCheck {
46  return { ...citation, status: matchStatus(ids.length), originalIds: ids };
47}
48function matchStatus(count: number): CitationStatus { return count === 1 ? "resolved" : count === 0 ? "missing" : "ambiguous"; }
49/** Copy before calling a model; a later cut cannot mutate this resolution index. */
50export function freezeCitationEvidence(value?: CitationEvidence): CitationEvidence | undefined {
51  return value === undefined ? undefined : JSON.parse(JSON.stringify(value)) as CitationEvidence;
52}
53