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

<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> • <a href="#runtime-model">Runtime Model</a> • <a href="#installation">Installation</a> • <a href="#mcp-tools">MCP Tools</a> • <a href="#development">Development</a>
lcm replaces sliding-window forgetfulness with a persistent memory runtime for both humans and agents.
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.
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"]
| Path | Restore | Prompt hints | Turn writeback | Automatic compaction | Notes |
|---|---|---|---|---|---|
| Claude Code | Yes | Yes | Yes, via transcript/hooks | Yes | Primary hook-based integration |
| GitHub Copilot (VS Code) | No | Yes, via skill/rules | No | No | Repo-local skill can teach Copilot to call lcm, but there is no automatic restore or turn capture yet |
| Codex | Yes | Yes | Yes, via native lifecycle hooks | LCM memory compacts on PreCompact; native compaction continues | lcm connectors install codex installs the hooks; see docs/vscode-codex.md. MCP config in .codex/config.toml is still manual |
| Oh My Pi | Yes | Yes | Yes, via native lifecycle hooks | Yes | lcm connectors install omp installs the hooks and --type mcp registers the MCP server; see docs/omp.md. |
| Phase | What happens |
|---|---|
| Persist | Raw messages are stored in SQLite per conversation |
| Summarize | Older messages are grouped into leaf summaries |
| Condense | Summaries roll up into higher-level DAG nodes |
| Promote | Durable insights are copied into cross-session memory |
| Restore | New sessions recover context from summaries and promoted memory |
| Recall | Agents 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"]
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.
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.
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.
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.
| Hook | Command | Purpose |
|---|---|---|
PreCompact | lcm compact --hook | Writes a DAG summary before compaction |
SessionStart | lcm restore | Restores project context, recent summaries, and promoted memory |
SessionEnd | lcm session-end | Ingests the completed Claude transcript |
UserPromptSubmit | lcm user-prompt | Searches memory and injects prompt-time hints |
Stop | lcm session-snapshot | Rolling transcript ingest, throttled |
PostToolUse | lcm post-tool | Passive learning: records decisions, plans, files, commands |
PostToolUseFailure | lcm post-tool | Passive 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"]
| Tool | Purpose |
|---|---|
lcm_search | Search across episodic memory (messages and summaries) and promoted memory |
lcm_grep | Regex or full-text search across raw messages and summaries |
lcm_expand | Decompress a summary node into its source content by traversing the DAG |
lcm_describe | Inspect session or summary metadata, lineage and explicit commit references |
lcm_store | Persist durable memory manually with optional tags |
lcm_stats | Show token savings, compression ratios, and usage statistics |
lcm_doctor | Diagnose setup and report stale stores, record-less stores and orphan summaries |
# 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.
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.
| Variable | Default | Description |
|---|---|---|
LCM_SUMMARY_PROVIDER | auto | auto, claude-process, codex-process, copilot-process, omp-process, anthropic, openai, disabled, or session |
ANTHROPIC_API_KEY | unset | Read only when llm.provider is anthropic (or session falling back to it) and llm.apiKey is unset |
LCM_HOME | ~/.lossless-claude | Where the daemon, databases, sidecars and logs live |
LCM_ENABLED | true | Set to false to make every Claude Code and Codex command hook a no-op while keeping the plugin registered |
LCM_CONTEXT_THRESHOLD | 0.75 | Context fill ratio that triggers compaction |
LCM_FRESH_TAIL_COUNT | 8 | Most recent raw messages protected from compaction |
LCM_LEAF_MIN_FANOUT | 3 | Minimum raw messages outside the fresh tail before a leaf pass runs |
LCM_CONDENSED_MIN_FANOUT | 2 | Minimum same-depth summaries before they are condensed |
LCM_CONDENSED_MIN_FANOUT_HARD | 1 | The same minimum during a hard-trigger sweep |
LCM_LEAF_CHUNK_TOKENS | 20000 | Maximum source tokens per leaf compaction pass |
LCM_CONDENSED_TARGET_TOKENS | 900 | Target size for condensed summaries |
auto resolves per caller:
claude-processcodex-processcopilot-processomp-processLCM_SUMMARY_PROVIDER override always takes precedenceSee docs/configuration.md for tuning notes and deeper operational guidance.
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.
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
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.
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.
MIT
hooks/lcm-hooks.ts 1030 lines1// 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}
1030hooks/model-budget.ts 120 lines1export 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}
120hooks/shadow-budget.ts 28 lines1import { 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()]; }
28hooks/shadow-deadline.ts 39 lines1import 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}
39hooks/shadow-boundaries.ts 15 lines1export 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}
15hooks/compaction-shadow.ts 253 lines1import 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}
253hooks/compaction-header.ts 138 lines1import 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 };
138hooks/compaction-header-schema.ts 84 lines1export 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}
84hooks/header-citations.ts 53 lines1import { 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