// aidlc-usage.ts - the single token-usage + cost extraction seam. // // One place owns: the model rate table (framework-default public list prices, // overridable), the Claude transcript readers (main + sibling sub-agent files), // the pure cost math, and the durable per-file-cursor ledger. Everything else // (audit rollup fields, metrics magnitude lines, the statusline segment) // CONSUMES this module and never re-parses a transcript itself. // // HARNESS SCOPE. The transcript reader is Claude-Code-format-specific: only the // Claude harness wires a producer (the PreToolUse + PostToolUse fold hook and // the Stop-hook flush). On Kiro / Codex / opencode no producer is wired, so the // ledger is never written and every consumer here degrades silently to no-data: // the statusline renders no cost segment, and the audit rollup adds no fields. // See docs/reference/06-hooks-and-tools.md. // // ROBUSTNESS CONTRACT. Nothing here throws on malformed or missing input. A // half-written last JSONL line is normal for a live transcript, so a bad line // is skipped silently; an absent or corrupt file yields [] / a fresh empty // ledger. Unknown model => tokens recorded, cost `null` - never a fabricated // number. // // The verified transcript facts this module encodes (separate sub-agent files, // .meta.json sidecars, per-token cache costs on the converse API, converse/ + // region model prefixes, and the per-file cursor requirement - uuids collide // across concurrent sub-agent files, so `(sourceFile, uuid)` is the only unique // key) are documented in docs/reference/06-hooks-and-tools.md. import { dlopen, ptr } from "bun:ffi"; import { createHash } from "node:crypto"; import { closeSync, existsSync, fstatSync, mkdirSync, openSync, readSync, readdirSync, readFileSync, } from "node:fs"; import { basename, dirname, join } from "node:path"; import { auditLockIdentity, intentUuidForSelection, modelRatesPath, readSessionIntentUuid, resolveProjectFlag, resolveWorkflowSelection, sessionsDir, withAuditLock, writeFileAtomic, } from "./aidlc-lib.ts"; // =========================================================================== // Task 1 - Rate table + pure cost math (no transcript I/O) // =========================================================================== // A per-million-token price row. `input`/`output` price ordinary prompt and // completion tokens; `cacheWrite5m`/`cacheWrite1h` price cache CREATION at the // two ephemeral TTLs; `cacheRead` prices a cache HIT. export type PriceRow = { input: number; output: number; cacheWrite5m: number; cacheWrite1h: number; cacheRead: number; }; // FRAMEWORK-DEFAULT rates, USD / 1e6 tokens. These are PUBLIC Anthropic list // prices, shipped as DEFAULTS only, and are the pure dev-checkout fallback used // when no shipped `tools/data/model-rates.json` is present (an authored core/ // tree carries the JSON, so an installed harness always reads that). The cache // fields follow the standard multipliers: cacheWrite5m = 1.25x input, // cacheWrite1h = 2x input, cacheRead = 0.1x input. // // GENERATION-DISCRETE keys (a row per model GENERATION, NOT one per family): // verified real sessions mix generations, and a family-collapse silently // misprices them onto whatever the "current" row happens to be. Key off the // generation token, not the account/region variant. // // To OVERRIDE: set AIDLC_MODEL_RATES to a rates file (same shape as // tools/data/model-rates.json), or edit the shipped model-rates.json in your // install. The override layers ON TOP of these defaults - a partial file only // changes the models it names; an unknown model stays tokens-with-null-cost. export const DEFAULT_RATES: Record = { "opus-5": { input: 5.0, output: 25.0, cacheWrite5m: 6.25, cacheWrite1h: 10.0, cacheRead: 0.5 }, "opus-4-8": { input: 5.0, output: 25.0, cacheWrite5m: 6.25, cacheWrite1h: 10.0, cacheRead: 0.5 }, "opus-4-7": { input: 5.0, output: 25.0, cacheWrite5m: 6.25, cacheWrite1h: 10.0, cacheRead: 0.5 }, "opus-4-6": { input: 5.0, output: 25.0, cacheWrite5m: 6.25, cacheWrite1h: 10.0, cacheRead: 0.5 }, "sonnet-5": { input: 3.0, output: 15.0, cacheWrite5m: 3.75, cacheWrite1h: 6.0, cacheRead: 0.3 }, "sonnet-4-6": { input: 3.0, output: 15.0, cacheWrite5m: 3.75, cacheWrite1h: 6.0, cacheRead: 0.3 }, "haiku-4-5": { input: 1.0, output: 5.0, cacheWrite5m: 1.25, cacheWrite1h: 2.0, cacheRead: 0.1 }, "fable-5": { input: 10.0, output: 50.0, cacheWrite5m: 12.5, cacheWrite1h: 20.0, cacheRead: 1.0 }, }; // Validate one parsed JSON object into a PriceRow, or null when it is not a // well-formed row (any missing/non-finite field). Keeps a malformed override // file from poisoning a real rate with NaN - the offending row is dropped and // the default (or nothing) stands. function coercePriceRow(raw: unknown): PriceRow | null { if (!raw || typeof raw !== "object") return null; const o = raw as Record; const num = (v: unknown): number | null => typeof v === "number" && Number.isFinite(v) ? v : null; const input = num(o.input); const output = num(o.output); const cacheWrite5m = num(o.cacheWrite5m); const cacheWrite1h = num(o.cacheWrite1h); const cacheRead = num(o.cacheRead); if ( input === null || output === null || cacheWrite5m === null || cacheWrite1h === null || cacheRead === null ) { return null; } return { input, output, cacheWrite5m, cacheWrite1h, cacheRead }; } // Read a rates file's `rates` map into overlay rows. Missing/corrupt/ill-shaped // file => {} (no overlay). Never throws. function readRatesFile(path: string): Record { let raw: string; try { raw = readFileSync(path, "utf-8"); } catch { return {}; } try { const parsed = JSON.parse(raw) as { rates?: unknown }; const rates = parsed?.rates; if (!rates || typeof rates !== "object") return {}; const out: Record = {}; for (const [key, value] of Object.entries(rates as Record)) { const row = coercePriceRow(value); if (row) out[key] = row; } return out; } catch { return {}; } } // The usage-tracking kill switch. Set AIDLC_DISABLE_USAGE_TRACKING=1 to stop // the ledger producer from writing and every consumer from reading: the fold // hooks and Stop-hook flush become no-ops, the statusline renders no cost // segment, and STAGE_COMPLETED / WORKFLOW_COMPLETED add no rollup fields. // Follows the AIDLC_DISABLE_* hook-flag convention (exact string "1"). Read at // call time, never cached - each hook/tool is a fresh process, and tests // toggle it mid-process. export function usageTrackingDisabled(): boolean { return resolveProjectFlag("AIDLC_DISABLE_USAGE_TRACKING") === "1"; } let _rates: Record | null = null; // The effective rate table, built in three layers (each overlays the previous, // per-model, so a partial file only changes the models it names): // 1. DEFAULT_RATES - hardcoded public list prices (the dev-checkout floor) // 2. tools/data/model-rates.json - the shipped framework default (an install edits this) // 3. $AIDLC_MODEL_RATES - a user/project-supplied override file // Cached at first use - a per-line computeCost never re-reads the files. Robust: // a missing or malformed source contributes nothing and the layers below stand. export function loadRates(): Record { if (_rates !== null) return _rates; const merged: Record = { ...DEFAULT_RATES }; // Layer 2: the shipped default file (present in an installed harness; absent // in a dev checkout's core/, where DEFAULT_RATES is the only source). try { for (const [k, v] of Object.entries(readRatesFile(modelRatesPath()))) merged[k] = v; } catch { /* no shipped file / unresolvable data dir - DEFAULT_RATES stands */ } // Layer 3: the AIDLC_MODEL_RATES override file, on top of everything. const override = process.env.AIDLC_MODEL_RATES; if (override) { for (const [k, v] of Object.entries(readRatesFile(override))) merged[k] = v; } _rates = merged; return _rates; } // TEST SEAM: drop the cached rate table so a test that mutates AIDLC_MODEL_RATES // mid-process re-reads it. Not used at runtime (each tool/hook is a fresh // process). Pure. export function _resetRatesCacheForTest(): void { _rates = null; } // Bare family aliases with NO generation token. Matched only on an EXACT // residual equality (never as a substring), so `opus` resolves but `opus-6` // (a would-be new generation) does NOT match and stays null. These are a // defensive convenience for hand-typed / test inputs; the wire form always // carries a generation token. `opus`/`sonnet`/`haiku` map to the generation // this harness ships as its default model set. const BARE_ALIASES: Record = { opus: "opus-4-8", sonnet: "sonnet-4-6", haiku: "haiku-4-5", fable: "fable-5", }; // Map a transcript `message.model` to a rate-table key, or null if unknown. // // Handles the real Bedrock forms verified in transcripts: // converse/us.anthropic.claude-opus-4-8 (main thread) // us.anthropic.claude-opus-4-8 (some sub-agents) // claude-haiku-4-5-20251001 (sub-agent, dated, no region) // converse/au.anthropic.claude-haiku-4-5-20251001-v1:0 // global.anthropic.claude-opus-4-8[1m] (this harness's settings alias) // // UNKNOWN-GENERATION POLICY: a Claude model whose GENERATION is not in the rate // table (e.g. a future `opus-6`, or an `opus-4-9` we haven't priced) returns // `null` => the caller records the tokens but withholds cost. An honest // "unknown" (made visible by the audit's `Cost USD: null`) beats a // confidently-wrong number from an old generation's rate. ``, empty, // malformed provider/model shapes, and non-Claude models => null too. Generation // keys come from the EFFECTIVE rate table, so an AIDLC_MODEL_RATES override can // add a new generation without a source-code matcher change. export function normalizeModel(modelId: string): string | null { if (!modelId || typeof modelId !== "string") return null; let s = modelId.trim().toLowerCase(); if (BARE_ALIASES[s]) return BARE_ALIASES[s]; // 1. Drop a leading `converse/` provider tag. MUST run before step 2 - the // region wildcard is anchored at the string start. const hadConversePrefix = s.startsWith("converse/"); if (hadConversePrefix) s = s.slice("converse/".length); // 2. Drop an inference-profile region prefix + `anthropic.`. WILDCARDED: // `` is any `[a-z0-9-]+` token, so a new region (eu/apac/global/ // ... or one not yet seen) never silently breaks normalization. The // `\.anthropic\.` anchor keeps the wildcard from eating a region-less model // (bare `claude-opus-4-8` has no `.anthropic.`, so it is untouched here). const provider = s.match(/^(?:[a-z0-9-]+\.)?anthropic\./); if (provider) { s = s.slice(provider[0].length); } else if (hadConversePrefix) { return null; } // 3. Accept only a real Anthropic/Claude shape or one of the documented bare // family aliases handled above. Exact rate keys are not wire model IDs: // accepting them after `converse/` or `anthropic.` would price malformed // provider shapes such as `converse/opus-4-8`. const rates = loadRates(); if (!s.startsWith("claude-")) return null; s = s.slice("claude-".length); // 4. Match a generation key only at a token boundary. Keys are longest-first // so an override containing related keys resolves the most specific one. for (const key of Object.keys(rates).sort((a, b) => b.length - a.length)) { if (s === key || s.startsWith(`${key}-`) || s.startsWith(`${key}[`)) { return key; } } // Unknown generation / non-Claude => null (tokens recorded, cost withheld). return null; } // The token magnitudes we price. Cache creation is split by ephemeral TTL; // cache reads are a single bucket. export type TokenCounts = { input: number; output: number; cacheCreate5m: number; cacheCreate1h: number; cacheRead: number; }; // Price a token bundle for a model. Unknown model => `{ usd: null, model: null }` // (tokens are preserved by the caller; only the cost is withheld). Cache reads // price at `cacheRead`; cache creation splits across the 5m / 1h rows. export function computeCost( counts: TokenCounts, modelId: string, ): { usd: number | null; model: string | null } { const key = normalizeModel(modelId); if (!key) return { usd: null, model: null }; const r = loadRates()[key]; if (!r) return { usd: null, model: null }; const usd = (counts.input / 1e6) * r.input + (counts.output / 1e6) * r.output + (counts.cacheCreate5m / 1e6) * r.cacheWrite5m + (counts.cacheCreate1h / 1e6) * r.cacheWrite1h + (counts.cacheRead / 1e6) * r.cacheRead; return { usd, model: key }; } // One priced assistant turn. `sourceKey` identifies the file the row came from // (`"main"` or `"agent-"`) and is the PER-FILE cursor key in Task 3. // `agentId` is filled during extraction (Task 2) for sub-agent rows; `agentType` // is filled during attribution (Task 4). export type UsageRow = { uuid: string; // The Claude `message.id`. Claude Code splits ONE llm call (one message.id) // across MULTIPLE JSONL lines - one per content block (thinking/text/tool_use) // - stamping the SAME `message.usage` on every line. `msgId` lets // dedupeByMessageId collapse a contiguous run of same-id lines into ONE row so // the shared usage is counted ONCE, not 2-3x. `""` when a line carries no // `message.id` (hand-authored fixtures, older transcripts): an id-less row is // NEVER merged - each such line stays its own row. msgId: string; timestamp: string; model: string; isSidechain: boolean; sourceKey: string; agentId: string | null; agentType: string | null; counts: TokenCounts; usd: number | null; }; // =========================================================================== // Task 2 - Transcript extraction (Claude main + sub-agent files) // =========================================================================== // Parse one JSONL object's `message.usage` into a TokenCounts. The 5m/1h split // lives under `cache_creation`; when that object is absent we treat the flat // `cache_creation_input_tokens` total as 5m (its pre-split form). function countsFromUsage(usage: Record): TokenCounts { const num = (v: unknown): number => (typeof v === "number" && Number.isFinite(v) ? v : 0); const cc = (usage.cache_creation ?? {}) as Record; const e5 = num(cc.ephemeral_5m_input_tokens); const e1 = num(cc.ephemeral_1h_input_tokens); const flat = num(usage.cache_creation_input_tokens); // Trust the nested 5m/1h split only when it actually accounts for cache // creation. Bedrock's converse API ships the split present-but-zeroed with // the real total in the flat cache_creation_input_tokens field; native // Anthropic always has flat == e5 + e1, so this is a no-op there. When the // split is absent or zeroed we fall back to the flat total, attributing it // to 5m: the harness only ever issues 5m-TTL cache writes, so a Bedrock // creation is a 5m write in practice. (A hypothetical 1h Bedrock write // would be under-priced at the 5m rate - but that beats the old behavior, // which dropped the flat value entirely and priced the write at zero.) const splitSum = e5 + e1; return { input: num(usage.input_tokens), output: num(usage.output_tokens), cacheCreate5m: splitSum > 0 ? e5 : flat, cacheCreate1h: splitSum > 0 ? e1 : 0, cacheRead: num(usage.cache_read_input_tokens), }; } // Parse ONE JSONL line into a UsageRow, or null when it is not a priced // assistant turn (non-assistant, no usage, malformed). Pure - the shared per- // line parse used by both the whole-file reader (readClaudeTranscript) and the // offset-aware fold reader. `sourceKey`/`fallbackAgentId` supply the file // identity when the line does not carry its own `agentId`. function parseTranscriptLine( line: string, sourceKey: string, fallbackAgentId: string | null, ): UsageRow | null { const trimmed = line.trim(); if (!trimmed) return null; let obj: Record; try { obj = JSON.parse(trimmed) as Record; } catch { return null; // half-written / malformed line - skip, don't throw } const message = obj.message as Record | undefined; if (!message) return null; if (message.role !== "assistant") return null; const usage = message.usage as Record | undefined; if (!usage || typeof usage !== "object") return null; const counts = countsFromUsage(usage); const model = typeof message.model === "string" ? message.model : ""; const agentId = typeof obj.agentId === "string" ? obj.agentId : fallbackAgentId; return { uuid: typeof obj.uuid === "string" ? obj.uuid : "", msgId: typeof message.id === "string" ? message.id : "", timestamp: typeof obj.timestamp === "string" ? obj.timestamp : "", model, isSidechain: obj.isSidechain === true, sourceKey, agentId, agentType: null, counts, usd: computeCost(counts, model).usd, }; } // covers: function:dedupeByMessageId // Collapse the split-line overcount (the CRITICAL bug): Claude Code writes one // llm call (one `message.id`) as MULTIPLE contiguous JSONL lines - one per // content block (thinking/text/tool_use). Left alone that inflates input ~2x and // output ~2.6x. Verified invariants this relies on: (a) split lines of one // message.id are always CONTIGUOUS; (b) the group's usage must be taken ONCE. // // TWO transcript styles exist (both handled by taking the group's usage as the // MAXIMUM-usage line of the run): // OLD-style: every split line repeats the SAME usage => max == any line. // NEW-style: leading lines (thinking/text) carry usage=0; the REAL usage // appears on the LAST line(s) of the run (identical when several trail) => // max == that real-usage line, so the zeros never mask the real count. // The representative is the max-usage row (ties => the LAST such row, so its // uuid/timestamp is the real end-of-turn line and the per-file cursor advances // correctly). On OLD-style all lines tie so this is the last line, matching the // prior behaviour exactly. // // FALLBACK: a row with an empty `msgId` (no `message.id` - hand-authored // fixtures, older transcripts) is NEVER merged; each such row passes through // as-is, preserving one-row-per-line behaviour. A `seen` Set guards against a // (non-contiguous) reappearance of an already-emitted non-empty msgId so a // second run of the same id can never double-count. Pure. export function dedupeByMessageId(rows: UsageRow[]): UsageRow[] { const out: UsageRow[] = []; const seen = new Set(); let i = 0; while (i < rows.length) { const row = rows[i]; if (!row.msgId) { // Id-less line - never merged, always its own row. out.push(row); i++; continue; } if (seen.has(row.msgId)) { // A non-contiguous reappearance of an already-collapsed id - drop it (its // usage was already counted once). Defensive; verified runs are contiguous. i++; continue; } // Consume the contiguous run sharing this msgId; take the group's usage ONCE // from the MAX-usage line (ties => last, the real end-of-turn line). let j = i; while (j + 1 < rows.length && rows[j + 1].msgId === row.msgId) j++; seen.add(row.msgId); out.push(representativeOfRun(rows, i, j)); i = j + 1; } return out; } // The representative of a contiguous same-msgId run `rows[start..end]`: the row // with the greatest total usage magnitude, breaking ties toward the LAST such // row. This is what makes both transcript styles collapse losslessly - NEW-style // leading usage=0 lines lose to the trailing real-usage line, OLD-style ties // resolve to the last line (the real end-of-turn uuid/timestamp). Pure. function representativeOfRun(rows: UsageRow[], start: number, end: number): UsageRow { let best = rows[start]; let bestMag = usageMagnitude(best.counts); for (let k = start + 1; k <= end; k++) { const mag = usageMagnitude(rows[k].counts); if (mag >= bestMag) { // `>=` breaks ties toward the later row (real end-of-turn line). best = rows[k]; bestMag = mag; } } return best; } // Total token magnitude of a bundle - the tie-break/selection key for a split // run's representative. Sums every priced magnitude so a line carrying real // usage always outranks a leading usage=0 line. Pure. function usageMagnitude(c: TokenCounts): number { return c.input + c.output + c.cacheCreate5m + c.cacheCreate1h + c.cacheRead; } // Read a single Claude JSONL transcript into priced UsageRows. `opts.sourceKey` // and `opts.agentId` default to a main-thread read; readClaudeSession supplies // the sub-agent identity. Emits one row per assistant LLM CALL (message.id) - // contiguous split lines sharing a message.id are collapsed via // dedupeByMessageId so the shared usage is counted once. Malformed lines (incl. // a half-written trailing line) are skipped silently - never throws on a bad // line or a missing file. export function readClaudeTranscript( path: string, opts?: { sourceKey?: string; agentId?: string | null }, ): UsageRow[] { const sourceKey = opts?.sourceKey ?? "main"; const fallbackAgentId = opts?.agentId ?? null; let raw: string; try { raw = readFileSync(path, "utf-8"); } catch { return []; } const rows: UsageRow[] = []; for (const line of raw.split("\n")) { const row = parseTranscriptLine(line, sourceKey, fallbackAgentId); if (row) rows.push(row); } return dedupeByMessageId(rows); } // Resolve the sibling `subagents/` directory for a main transcript. Layout: // main is `//.jsonl`; sub-agents are // `///subagents/agent-.jsonl`. export function subagentDir(mainTranscriptPath: string): string { const dir = dirname(mainTranscriptPath); const session = basename(mainTranscriptPath).replace(/\.jsonl$/, ""); return join(dir, session, "subagents"); } // Read one `agent-.meta.json` sidecar; return its parsed object or null. function readMetaSidecar(jsonlPath: string): { agentType?: string } | null { const metaPath = jsonlPath.replace(/\.jsonl$/, ".meta.json"); if (!existsSync(metaPath)) return null; try { return JSON.parse(readFileSync(metaPath, "utf-8")) as { agentType?: string }; } catch { return null; } } // Read a full Claude session: the main transcript plus every sibling // `subagents/agent-*.jsonl`. Main rows carry `agentId: null`, // `sourceKey: "main"`; sub-agent rows carry their file's `agentId` and // `sourceKey: "agent-"`. Each sub-agent's `.meta.json` sidecar is read // for its `agentType`, and agent attribution (Task 4) is applied before // returning. A missing `subagents/` dir => just the main rows. Never throws. export function readClaudeSession(mainPath: string): UsageRow[] { const rows: UsageRow[] = readClaudeTranscript(mainPath, { sourceKey: "main", agentId: null, }); const metaByAgentId: Record = {}; const subDir = subagentDir(mainPath); let entries: string[] = []; try { if (existsSync(subDir)) { entries = readdirSync(subDir).filter( (f) => f.startsWith("agent-") && f.endsWith(".jsonl"), ); } } catch { entries = []; } for (const file of entries) { const agentId = file.replace(/^agent-/, "").replace(/\.jsonl$/, ""); const full = join(subDir, file); const meta = readMetaSidecar(full); if (meta && typeof meta.agentType === "string" && meta.agentType) { metaByAgentId[agentId] = meta.agentType; } rows.push( ...readClaudeTranscript(full, { sourceKey: `agent-${agentId}`, agentId, }), ); } return attributeAgents(rows, metaByAgentId); } // Sniff the transcript format from its first non-empty line and dispatch. // Claude lines carry `message`/`uuid` => readClaudeSession. Any other shape // (empty, non-Claude) => []. Only the Claude Code transcript format is // supported - Kiro/Codex/opencode wire no producer, so their transcripts (if // any) are never read here. export function readTranscript(path: string): UsageRow[] { let raw: string; try { raw = readFileSync(path, "utf-8"); } catch { return []; } for (const line of raw.split("\n")) { const trimmed = line.trim(); if (!trimmed) continue; let obj: Record; try { obj = JSON.parse(trimmed) as Record; } catch { continue; } if ("message" in obj || "uuid" in obj) { return readClaudeSession(path); } // Sniffed a line we don't recognise - keep looking at subsequent lines. } return []; } // =========================================================================== // Task 4 - Agent-type attribution (declared before Task 3's fold so // readClaudeSession can call it) // =========================================================================== // Attribute each row to an agentType. Main-thread rows => "main". Sub-agent rows // => the sidecar `agentType` for their agentId, falling back to "subagent" when // no sidecar/field exists (never null, never dropped). Pure. export function attributeAgents( rows: UsageRow[], metaByAgentId: Record, ): UsageRow[] { for (const row of rows) { if (row.sourceKey === "main") { row.agentType = "main"; } else { const id = row.agentId ?? ""; row.agentType = metaByAgentId[id] ?? "subagent"; } } return rows; } // Fallback join: build `agentId -> subagent_type` from a parent transcript's // lines. A `Task`/`Agent` tool_use carries an `id` (tool_use_id) and // `input.subagent_type`; the matching `toolUseResult` line carries `agentId` // for that same tool_use_id. Used ONLY for agentIds absent from the sidecar map // (the sidecar `agentType` is present in ~100% of sampled sidecars; `toolUseId` // only ~86%). Never throws. export function buildAgentTypeMapFromParent( parentLines: string[], ): Record { const typeByToolUseId: Record = {}; const map: Record = {}; for (const line of parentLines) { const trimmed = line.trim(); if (!trimmed) continue; let obj: Record; try { obj = JSON.parse(trimmed) as Record; } catch { continue; } // (a) Collect subagent_type per tool_use id from assistant tool_use blocks. const message = obj.message as Record | undefined; const content = message?.content; if (Array.isArray(content)) { for (const block of content) { const b = block as Record; if ( b.type === "tool_use" && (b.name === "Task" || b.name === "Agent") && typeof b.id === "string" ) { const input = (b.input ?? {}) as Record; if (typeof input.subagent_type === "string") { typeByToolUseId[b.id] = input.subagent_type; } } } } // (b) Join a toolUseResult's agentId back to the tool_use's subagent_type. const tur = obj.toolUseResult as Record | undefined; const toolUseId = typeof obj.tool_use_id === "string" ? obj.tool_use_id : typeof tur?.tool_use_id === "string" ? (tur.tool_use_id as string) : undefined; const agentId = typeof tur?.agentId === "string" ? (tur.agentId as string) : undefined; if (toolUseId && agentId && typeByToolUseId[toolUseId]) { map[agentId] = typeByToolUseId[toolUseId]; } } return map; } // =========================================================================== // Task 3 - Per-file-cursor ledger // =========================================================================== // Summed token magnitudes + total USD for a bucket. export type Totals = { tokens: TokenCounts; usd: number; }; // A per-stage bucket. `totals` is the flat stage sum; `byModel`/`byAgent` are // the stage's OWN per-model / per-agent sub-splits - stage-scoped, NOT the // global session buckets. This is what lets STAGE_COMPLETED show a By Model line // that agrees with its stage-scoped Cost USD (a global By Model would sum every // stage and contradict a single-stage cost). export type StageBucket = { totals: Totals; byModel: Record; byAgent: Record; }; // One aggregate at a particular ownership boundary. The top-level ledger keeps // a workspace aggregate for diagnostics, while `workflows[workflow].sessions` // provides the authoritative workflow/session intersection. export type UsageAggregate = { totals: Totals; byStage: Record; byModel: Record; byAgent: Record; }; export type WorkflowUsage = UsageAggregate & { sessions: Record; }; // The current on-disk ledger schema version. BUMP THIS whenever the token // COUNTING semantics change so old, differently-counted totals are discarded // rather than re-added onto. v2 is the first schema whose totals were produced // by the HOLDBACK fold (a group is counted once, when provably complete); every // pre-v2 ledger was produced by the buggy per-split-line counter (input ~2x, // output ~2.6x inflated), so its totals are unreliable and loadLedger resets // them (see the migration in loadLedger). v3 adds session/intent ownership and // pending-group attribution; v2's workspace-only totals cannot be partitioned // retrospectively, so they are rebuilt too. const CURRENT_SCHEMA_VERSION = 3 as const; // The durable rollup. `cursors` is keyed by SOURCE FILE identity - the OFFSET // FOLD (foldFileIntoLedger) keys each cursor by the transcript FILE PATH (the // main transcript's path; each sub-agent file's full path), which is unique // ACROSS sessions so a cumulative workspace ledger never reuses one session's // byteOffset against a different session's file. The row-based updateLedger - // which has no file path - keys by the row's `sourceKey` (`"main"` / // `"agent-"`) instead; the two schemes coexist without collision // because a file path never equals `"main"`. Verified: uuids collide across // concurrent sub-agent files, so a global-uuid cursor would drop real turns or // count the 0-token broadcast copies - per-file cursors are the whole mechanism. // `byStage` carries stage-scoped sub-splits (StageBucket); top-level totals and // splits span the workspace and are diagnostic only. export type LedgerCursor = { lastUuid: string; lastTimestamp: string; // Byte position in the source file already folded. The offset-aware fold // reads only bytes `[byteOffset, size)` each call. byteOffset?: number; // The `message.id` of the last folded group (diagnostics only). lastMessageId?: string; // A non-flush fold retains the last message group. Its ownership is captured // NOW, before a lifecycle tool can advance state, and reused when the group // is eventually folded. pending?: { byteOffset: number; messageId: string; stageSlug: string | null; sessionKey: string; workflowKey: string; }; }; export type Ledger = UsageAggregate & { // Schema version of the totals below. Absent/`< CURRENT_SCHEMA_VERSION` on // disk => produced by an older, differently-counted fold => loadLedger // discards it and rebuilds from the transcript. schemaVersion: number; cursors: Record; workflows: Record; }; function emptyTokenCounts(): TokenCounts { return { input: 0, output: 0, cacheCreate5m: 0, cacheCreate1h: 0, cacheRead: 0 }; } function emptyTotals(): Totals { return { tokens: emptyTokenCounts(), usd: 0 }; } function emptyStageBucket(): StageBucket { return { totals: emptyTotals(), byModel: {}, byAgent: {} }; } function emptyUsageAggregate(): UsageAggregate { return { totals: emptyTotals(), byStage: {}, byModel: {}, byAgent: {} }; } function emptyLedger(): Ledger { return { schemaVersion: CURRENT_SCHEMA_VERSION, cursors: {}, ...emptyUsageAggregate(), workflows: {}, }; } // The gitignored runtime ledger path: `aidlc/.aidlc-sessions/usage-ledger.json`. export function ledgerPath(projectDir: string): string { return join(sessionsDir(projectDir), "usage-ledger.json"); } export type UsageContext = { sessionId?: string; sessionKey?: string; workflowKey?: string; }; // The transcript path is the stable session identity available to both the fold // hook and statusline. A caller without a transcript can supply an explicit // session key (the row-based test/utility seam) or fall back to session_id. export function sessionUsageKey( transcriptPath?: string, sessionId?: string, ): string { const transcript = transcriptPath?.trim(); if (transcript) return `transcript:${transcript}`; const session = safeSessionSegment(sessionId ?? ""); return session ? `session:${session}` : "session:unknown"; } // Resolve the active intent to a stable UUID when possible. Legacy/orphan // records fall back to their space + record-dir identity. export function intentUsageKey( projectDir: string, sessionId?: string, ): string { try { if (sessionId) { const stamped = readSessionIntentUuid(projectDir, sessionId); if (stamped) return `intent:${stamped}`; } const selection = resolveWorkflowSelection(projectDir, { sessionId }); const uuid = intentUuidForSelection(projectDir, selection); if (uuid) return `intent:${uuid}`; return `record:${selection.space}/${selection.intent ?? "legacy"}`; } catch { return "record:default/legacy"; } } // Coerce a possibly-old-shape byStage map into the StageBucket shape. An older // ledger stored a flat Totals per stage; treat that as `{ totals, byModel:{}, // byAgent:{} }`. A missing sub-map on a new-shape entry becomes {}. Never throws. function normalizeByStage( raw: Record | undefined, ): Record { const out: Record = {}; if (!raw || typeof raw !== "object") return out; for (const [slug, val] of Object.entries(raw)) { const v = (val ?? {}) as Record; if (v.totals && typeof v.totals === "object") { // New shape (StageBucket) - tolerate missing sub-maps. out[slug] = { totals: v.totals as Totals, byModel: (v.byModel as Record) ?? {}, byAgent: (v.byAgent as Record) ?? {}, }; } else if (v.tokens && typeof v.tokens === "object") { // Old shape (flat Totals) - wrap it, no per-model/agent history available. out[slug] = { totals: v as unknown as Totals, byModel: {}, byAgent: {} }; } else { out[slug] = emptyStageBucket(); } } return out; } function normalizeUsageAggregate(raw: unknown): UsageAggregate { if (!raw || typeof raw !== "object") return emptyUsageAggregate(); const value = raw as Partial & { byStage?: Record; }; return { totals: value.totals ?? emptyTotals(), byStage: normalizeByStage(value.byStage), byModel: value.byModel ?? {}, byAgent: value.byAgent ?? {}, }; } function normalizeAggregateMap( raw: Record | undefined, ): Record { const out: Record = {}; if (!raw || typeof raw !== "object") return out; for (const [key, value] of Object.entries(raw)) { out[key] = normalizeUsageAggregate(value); } return out; } function normalizeWorkflowMap( raw: Record | undefined, ): Record { const out: Record = {}; if (!raw || typeof raw !== "object") return out; for (const [key, value] of Object.entries(raw)) { const aggregate = normalizeUsageAggregate(value); const sessions = value && typeof value === "object" ? normalizeAggregateMap( (value as { sessions?: Record }).sessions, ) : {}; out[key] = { ...aggregate, sessions }; } return out; } // Whether an on-disk ledger's cursors are all pre-v2 shaped (belt-and-suspenders // migration signal). A current fold cursor ALWAYS carries a numeric // `byteOffset`; a // pre-v2 ledger's cursors carry only `lastUuid`/`lastMessageId`. If ANY cursor // lacks a byteOffset the totals below it were produced (at least partly) by the // legacy per-line counter, so we treat the whole ledger as pre-v2. An empty // cursors map is NOT a downgrade signal (a fresh ledger has none yet); the // schemaVersion check is authoritative there. function cursorsLackByteOffset( cursors: Record | undefined, ): boolean { if (!cursors || typeof cursors !== "object") return false; for (const c of Object.values(cursors)) { if ( c === null || typeof c !== "object" || typeof (c as { byteOffset?: unknown }).byteOffset !== "number" ) { return true; } } return false; } // Load the ledger, or a fresh empty one on missing/corrupt/invalid-shape input // (rebuild-on-corruption). Never throws, logs nothing. // // MIGRATION. A ledger whose on-disk `schemaVersion` is missing or // `< CURRENT_SCHEMA_VERSION` - or (belt-and-suspenders) whose cursors lack a // numeric byteOffset - had its totals produced by the buggy per-split-line // counter. Those totals are UNRELIABLE and must not be re-added onto, so we // DISCARD the whole file and return a fresh current-schema ledger; the next fold // rebuilds totals/byStage/byModel/byAgent/cursors cleanly from the // transcript(s). This resets, rather than migrates, because the old numbers // cannot be corrected in place. export function loadLedger(projectDir: string): Ledger { let raw: string; try { raw = readFileSync(ledgerPath(projectDir), "utf-8"); } catch { return emptyLedger(); } try { const parsed = JSON.parse(raw) as Partial & { byStage?: Record; schemaVersion?: unknown; cursors?: Record; workflows?: Record; }; if (!parsed || typeof parsed !== "object") return emptyLedger(); // Migration reset: an older ledger cannot supply the current ownership and // counting guarantees, so discard and rebuild rather than folding onto it. const onDiskVersion = typeof parsed.schemaVersion === "number" ? parsed.schemaVersion : 0; if ( onDiskVersion < CURRENT_SCHEMA_VERSION || cursorsLackByteOffset(parsed.cursors) ) { return emptyLedger(); } // Merge onto an empty ledger so a partial/old file can't crash a consumer. return { schemaVersion: CURRENT_SCHEMA_VERSION, cursors: (parsed.cursors as Record) ?? {}, totals: parsed.totals ?? emptyTotals(), byStage: normalizeByStage(parsed.byStage), byModel: parsed.byModel ?? {}, byAgent: parsed.byAgent ?? {}, workflows: normalizeWorkflowMap(parsed.workflows), }; } catch { return emptyLedger(); } } // Fold one row's counts + usd into a Totals accumulator (mutates + returns it). function addInto(t: Totals, row: UsageRow): Totals { t.tokens.input += row.counts.input; t.tokens.output += row.counts.output; t.tokens.cacheCreate5m += row.counts.cacheCreate5m; t.tokens.cacheCreate1h += row.counts.cacheCreate1h; t.tokens.cacheRead += row.counts.cacheRead; if (typeof row.usd === "number" && Number.isFinite(row.usd)) t.usd += row.usd; return t; } // The bucket key for byModel: the normalized rate-table key when known, else the // raw model string, else "unknown". Keeps unknown-model tokens visible. function modelBucketKey(row: UsageRow): string { return normalizeModel(row.model) ?? row.model ?? "unknown"; } function foldRowIntoAggregate( aggregate: UsageAggregate, row: UsageRow, stageSlug: string | null, ): void { addInto(aggregate.totals, row); const mk = modelBucketKey(row); const ak = row.agentType ?? "subagent"; aggregate.byModel[mk] = addInto(aggregate.byModel[mk] ?? emptyTotals(), row); aggregate.byAgent[ak] = addInto(aggregate.byAgent[ak] ?? emptyTotals(), row); if (stageSlug) { const stage = aggregate.byStage[stageSlug] ?? emptyStageBucket(); addInto(stage.totals, row); // Stage-scoped sub-splits mirror the aggregate folds so By Model / By Agent // on STAGE_COMPLETED agree with the stage's own Cost USD. stage.byModel[mk] = addInto(stage.byModel[mk] ?? emptyTotals(), row); stage.byAgent[ak] = addInto(stage.byAgent[ak] ?? emptyTotals(), row); aggregate.byStage[stageSlug] = stage; } } // Fold one row into the workspace diagnostic aggregate and the authoritative // session + intent aggregates. function foldRowIntoLedger( ledger: Ledger, row: UsageRow, stageSlug: string | null, sessionKey: string, workflowKey: string, ): void { foldRowIntoAggregate(ledger, row, stageSlug); const workflow = ledger.workflows[workflowKey] ?? { ...emptyUsageAggregate(), sessions: {}, }; foldRowIntoAggregate(workflow, row, stageSlug); const session = workflow.sessions[sessionKey] ?? emptyUsageAggregate(); foldRowIntoAggregate(session, row, stageSlug); workflow.sessions[sessionKey] = session; ledger.workflows[workflowKey] = workflow; } const USAGE_LOCK_INTENT = "__usage-ledger__"; const USAGE_LOCK_SPACE = "__runtime__"; // The directory lock's stale-reaper restore gap can admit two holders on Win32. // A kernel mutex closes that gap and transfers ownership after an abandoned owner. const WIN32_USAGE_MUTEX = process.platform === "win32" ? dlopen("kernel32.dll", { CreateMutexW: { args: ["ptr", "i32", "ptr"], returns: "ptr" }, WaitForSingleObject: { args: ["ptr", "u32"], returns: "u32" }, ReleaseMutex: { args: ["ptr"], returns: "i32" }, CloseHandle: { args: ["ptr"], returns: "i32" }, }) : null; const WAIT_OBJECT_0 = 0; const WAIT_ABANDONED = 0x80; const USAGE_MUTEX_WAIT_MS = 5000; function withUsageLedgerLock(projectDir: string, fn: () => Ledger): Ledger { if (WIN32_USAGE_MUTEX !== null) { const identity = auditLockIdentity( projectDir, USAGE_LOCK_INTENT, USAGE_LOCK_SPACE, ); const hash = createHash("sha256").update(identity).digest("hex").slice(0, 32); const name = Buffer.from(`Global\\aidlc-usage-${hash}\0`, "utf16le"); const handle = WIN32_USAGE_MUTEX.symbols.CreateMutexW(null, 0, ptr(name)); if (handle === null) { throw new Error("Failed to create the Windows usage-ledger mutex"); } const waitResult = WIN32_USAGE_MUTEX.symbols.WaitForSingleObject( handle, USAGE_MUTEX_WAIT_MS, ); if (waitResult !== WAIT_OBJECT_0 && waitResult !== WAIT_ABANDONED) { WIN32_USAGE_MUTEX.symbols.CloseHandle(handle); throw new Error( `Failed to acquire the Windows usage-ledger mutex: ${waitResult}`, ); } try { return fn(); } finally { WIN32_USAGE_MUTEX.symbols.ReleaseMutex(handle); WIN32_USAGE_MUTEX.symbols.CloseHandle(handle); } } return withAuditLock( projectDir, fn, USAGE_LOCK_INTENT, USAGE_LOCK_SPACE, 200, 25, ); } // Incrementally fold new rows into the ledger and advance the per-file cursors. // // 1. Group rows by sourceKey (file identity). // 2. Rows are already in file order from extraction - do NOT re-sort across // files by timestamp (timestamps collide across concurrent sub-agents). // 3. If a group has a cursor, drop rows up to AND INCLUDING lastUuid (append-only // files => everything after the checkpoint is new). No cursor => all new. // 4. Fold survivors into totals / byModel / byAgent, and byStage[stageSlug] when // a stage is supplied. // 5. Advance cursors[sourceKey] to each group's last folded row. // 6. Write atomically. // // Duplicate-turn guard: a broadcast turn appears in several files with // output_tokens:0 in the copies; folding stays correct because each file is // counted ONLY against its own cursor. Do NOT dedupe by uuid across files - that // would drop the one real-token copy. export function updateLedger( projectDir: string, rows: UsageRow[], stageSlug: string | null, context: UsageContext = {}, ): Ledger { return withUsageLedgerLock(projectDir, () => { const ledger = loadLedger(projectDir); const sessionKey = context.sessionKey ?? sessionUsageKey(undefined, context.sessionId); const workflowKey = context.workflowKey ?? intentUsageKey(projectDir, context.sessionId); // 1. Group by sourceKey, preserving file order within each group. const groups = new Map(); for (const row of rows) { const arr = groups.get(row.sourceKey) ?? []; arr.push(row); groups.set(row.sourceKey, arr); } for (const [sourceKey, rawGroupRows] of groups) { // Collapse split-line overcount WITHIN this file's group (a caller may pass // raw split rows). Per-group so the dedup never crosses files - a broadcast // turn sharing an id across concurrent sub-agent files must still count once // per file (the per-file-cursor invariant). Reuses the same pure helper the // whole-file reader uses. const groupRows = dedupeByMessageId(rawGroupRows); const cursor = ledger.cursors[sourceKey]; let fresh = groupRows; if (cursor?.lastUuid) { const idx = groupRows.findIndex((r) => r.uuid === cursor.lastUuid); // Everything AFTER the checkpoint row is new. If the checkpoint uuid isn't // in this batch (e.g. the batch is entirely older, or the file was // truncated), keep all rows - a re-fold is harmless per idempotency. fresh = idx >= 0 ? groupRows.slice(idx + 1) : groupRows; } if (fresh.length === 0) continue; for (const row of fresh) { foldRowIntoLedger(ledger, row, stageSlug, sessionKey, workflowKey); } const last = fresh[fresh.length - 1]; // Preserve any byteOffset the file-fold layer set for this source - the // row-based API advances uuid/timestamp/msgId but does not track bytes. // ALWAYS write a numeric byteOffset (default 0) so this API produces a // cursor-complete current-schema ledger. This is safe alongside the fold // layer because the two keyspaces are disjoint: updateLedger keys by // row.sourceKey (`"main"`/`"agent-"`) while foldFileIntoLedger keys by // FILE PATH, so a 0 here is never read by a fold. const prevOffset = ledger.cursors[sourceKey]?.byteOffset ?? 0; ledger.cursors[sourceKey] = { lastUuid: last.uuid, lastTimestamp: last.timestamp, lastMessageId: last.msgId || ledger.cursors[sourceKey]?.lastMessageId || "", byteOffset: prevOffset, }; } // 6. Persist atomically (mkdir the sessions dir first). try { mkdirSync(sessionsDir(projectDir), { recursive: true }); } catch { /* dir may already exist */ } writeFileAtomic(ledgerPath(projectDir), JSON.stringify(ledger, null, 2)); return ledger; }); } // =========================================================================== // Task 5 - Audit per-stage rollup fields (ledger read only, NO transcript I/O) // =========================================================================== // Compact token count for audit field values: 1234 -> "1.2k", 3_400_000 -> // "3.4M", sub-1000 -> the integer, 0/negative/non-finite -> "0". Mirrors // aidlc-statusline.ts fmtTokens exactly (kept a LOCAL copy - usage.ts must not // import a hook, and the metrics reader un-abbreviates this same form). One // decimal at k/M scale, trailing ".0" trimmed. function fmtTokensCompact(n: number): string { if (!Number.isFinite(n) || n <= 0) return "0"; const trim = (s: string): string => s.replace(/\.0$/, ""); if (n >= 1e6) return `${trim((n / 1e6).toFixed(1))}M`; if (n >= 1e3) return `${trim((n / 1e3).toFixed(1))}k`; return String(Math.round(n)); } // Format a per-key COST breakdown for the `By Agent` line - `k=1.23; k2=0.45`, // USD to 2dp, only buckets with a positive cost. Empty -> "". (Agent buckets // can't classify unknown-vs-zero at agent granularity, so zero-cost agents are // simply omitted here; their TOKENS still surface via the token-breakdown // fields below, so no agent is hidden.) function formatByCost(buckets: Record): string { const parts: string[] = []; for (const [key, t] of Object.entries(buckets)) { if (t.usd > 0) parts.push(`${key}=${t.usd.toFixed(2)}`); } return parts.join("; "); } // Format the per-MODEL cost breakdown with null-visibility: a known // (rate-keyed) model prints its cost `opus-4-8=1.23` (even `=0.00`); an // unknown-model bucket prints `=null` so an unpriceable slice is // VISIBLE, never silently dropped. Empty map -> "". function formatModelCost(buckets: Record): string { const rates = loadRates(); const parts: string[] = []; for (const [key, t] of Object.entries(buckets)) { parts.push(key in rates ? `${key}=${t.usd.toFixed(2)}` : `${key}=null`); } return parts.join("; "); } // Format a per-key TOKEN breakdown - `k=in/out/cacheRead/cacheWrite; k2=...` // with each count in compact form (`12.3k/4.1k/2k/79.6k`). The quad order is // input / output / cacheRead / cacheWrite, documented in // docs/reference/06-hooks-and-tools.md. cacheWrite = cacheCreate5m + // cacheCreate1h - both TTLs summed into one bucket, mirroring the single Cache // Read bucket (the harness issues 5m-TTL writes in practice, but summing is // correct for either TTL and matches computeCost, which prices both). A bucket // whose four counts are all zero is skipped. Empty -> "". This is what makes the // four token counts visible per model AND per agent, not just cost. The metrics // reader parses this exact shape (un-abbreviating the compact values) into // per-model / per-agent token magnitude lines. function formatByTokens(buckets: Record): string { const parts: string[] = []; for (const [key, t] of Object.entries(buckets)) { const { input, output, cacheRead } = t.tokens; const cacheWrite = t.tokens.cacheCreate5m + t.tokens.cacheCreate1h; if (input <= 0 && output <= 0 && cacheRead <= 0 && cacheWrite <= 0) continue; parts.push( `${key}=${fmtTokensCompact(input)}/${fmtTokensCompact(output)}/${fmtTokensCompact(cacheRead)}/${fmtTokensCompact(cacheWrite)}`, ); } return parts.join("; "); } // Whether a stage bucket carries any priceable (known-model) sub-bucket. function hasKnownModel(byModel: Record): boolean { const rates = loadRates(); return Object.keys(byModel).some((k) => k in rates); } // Whether a stage recorded any tokens at all. function hasAnyTokens(t: Totals): boolean { const c = t.tokens; return ( c.input > 0 || c.output > 0 || c.cacheCreate5m > 0 || c.cacheCreate1h > 0 || c.cacheRead > 0 ); } // Compact usage summaries for STAGE_COMPLETED / WORKFLOW_COMPLETED, read from // the authoritative workflow aggregate - no transcript re-parse or time-slicing. // // THREE cost states are distinguished: // (a) no usage data for the stage => {} (no fields; caller stays clean) // (b) usage recorded but UNPRICEABLE => `Cost USD: null` (all unknown-model) // (c) priced => `Cost USD: 1.23` // A mixed known+unknown stage prices the KNOWN portion (state c) and the // `By Model` line shows the unknown slice as `=null` - no fabricated cost. // // `Tokens By Model` / `Tokens By Agent` carry the four token counts // (in/out/cacheRead/cacheWrite) per model and per agent, so the breakdowns are // token-aware, not cost-only. function aggregateUsageAuditFields( stage: Pick, ): Record { const t = stage.totals; const fields: Record = { "Tokens In": String(t.tokens.input), "Tokens Out": String(t.tokens.output), "Cache Read": String(t.tokens.cacheRead), "Cache Write": String(t.tokens.cacheCreate5m + t.tokens.cacheCreate1h), }; // Cost USD - three-state. Priced => 2dp; unpriceable-but-used => "null"; // genuinely no usage => omitted. if (t.usd > 0 || hasKnownModel(stage.byModel)) { fields["Cost USD"] = t.usd.toFixed(2); } else if (hasAnyTokens(t)) { fields["Cost USD"] = "null"; // recorded usage, all unknown-model => unpriceable } // Stage-scoped breakdowns - read the stage's OWN sub-maps, never the global // byModel/byAgent (those sum every stage and would contradict Cost USD). const byModel = formatModelCost(stage.byModel); if (byModel) fields["By Model"] = byModel; const byAgent = formatByCost(stage.byAgent); if (byAgent) fields["By Agent"] = byAgent; // Token counts grouped by model and by agent (in/out/cacheRead/cacheWrite). const tokModel = formatByTokens(stage.byModel); if (tokModel) fields["Tokens By Model"] = tokModel; const tokAgent = formatByTokens(stage.byAgent); if (tokAgent) fields["Tokens By Agent"] = tokAgent; return fields; } export function stageUsageAuditFields( projectDir: string, stageSlug: string, workflowKey: string = intentUsageKey(projectDir), ): Record { if (usageTrackingDisabled()) return {}; const stage = loadLedger(projectDir).workflows[workflowKey]?.byStage[stageSlug]; return stage ? aggregateUsageAuditFields(stage) : {}; } export function workflowUsageAuditFields( projectDir: string, workflowKey: string = intentUsageKey(projectDir), ): Record { if (usageTrackingDisabled()) return {}; const workflow = loadLedger(projectDir).workflows[workflowKey]; return workflow && hasAnyTokens(workflow.totals) ? aggregateUsageAuditFields(workflow) : {}; } export function sessionUsageAggregate( projectDir: string, transcriptPath?: string, workflowKey?: string, sessionId?: string, ): UsageAggregate | null { if (usageTrackingDisabled()) return null; const resolvedTranscript = transcriptPath ?? (sessionId ? readCurrentTranscriptPath(projectDir, sessionId) : null); const sessionKey = sessionUsageKey(resolvedTranscript ?? undefined, sessionId); const resolvedWorkflowKey = workflowKey ?? intentUsageKey(projectDir, sessionId); return ( loadLedger(projectDir).workflows[resolvedWorkflowKey]?.sessions[sessionKey] ?? null ); } // =========================================================================== // Task 6 - Persisted transcript path + runtime ledger producer (usage.ts half) // =========================================================================== // Normalise a session id to a safe filename segment (mirrors aidlc-lib's // sessionRecordPath sanitisation) so a host-supplied id can't escape the dir. function safeSessionSegment(sessionId: string): string { return sessionId.replace(/[^A-Za-z0-9._-]+/g, "-").replace(/^-+|-+$/g, ""); } // Atomically persist the live transcript path so transcript-blind CLI tools // (state, statusline) can find it. Writes both a session-scoped // `.transcript` (safe for concurrent sessions) AND a // `current.transcript` convenience pointer for readers that lack a session id. // Best-effort: swallows write errors (never breaks a hook). export function writeCurrentTranscriptPath( projectDir: string, sessionId: string, transcriptPath: string, ): void { // Only usage consumers read these pointers, so the kill switch covers them. if (usageTrackingDisabled()) return; if (!transcriptPath) return; try { mkdirSync(sessionsDir(projectDir), { recursive: true }); } catch { /* dir may already exist */ } const seg = safeSessionSegment(sessionId); try { if (seg) { writeFileAtomic( join(sessionsDir(projectDir), `${seg}.transcript`), transcriptPath, ); } writeFileAtomic( join(sessionsDir(projectDir), "current.transcript"), transcriptPath, ); } catch { /* best-effort */ } } // Read back the persisted transcript path. With a sessionId, reads the // session-scoped file; without one, falls back to `current.transcript`. Returns // null if absent/empty. Never throws. export function readCurrentTranscriptPath( projectDir: string, sessionId?: string, ): string | null { const candidates: string[] = []; if (sessionId) { const seg = safeSessionSegment(sessionId); if (seg) candidates.push(join(sessionsDir(projectDir), `${seg}.transcript`)); } else { candidates.push(join(sessionsDir(projectDir), "current.transcript")); } for (const path of candidates) { try { const raw = readFileSync(path, "utf-8").trim(); if (raw) return raw; } catch { /* try next candidate */ } } return null; } // The result of an offset-aware chunk read: the COMPLETE (newline-terminated) // lines newly available since `byteOffset`, the ABSOLUTE byte start position of // each of those lines (parallel to `lines`, so the caller can rewind the offset // to where a held-back group begins), and the byte position up to which those // complete lines run. A trailing partial line (no final `\n` - a half-written // append) is DROPPED and left for the next fold. `reset` signals a // truncation/rotation was detected (size < byteOffset) so the caller resets the // per-file checkpoint too. `null` means nothing to do (file gone, or // size === byteOffset with no new bytes). type ChunkRead = { lines: string[]; lineByteStarts: number[]; newByteOffset: number; reset: boolean; trailingPartial: boolean; }; // Read only bytes `[byteOffset, size)` of a file, BYTE-accurately (UTF-8 multi- // byte safe - offsets are byte positions, never char indices). Returns the // complete lines in that window, each line's absolute byte start, and the // advanced byte offset. Never throws. function readChunkFromOffset( path: string, byteOffset: number, flush: boolean, ): ChunkRead | null { let fd: number; try { fd = openSync(path, "r"); } catch { return null; // file absent / unreadable } try { const size = fstatSync(fd).size; let reset = false; let start = byteOffset; if (size === start) return null; // no new bytes - zero parse if (size < start) { // Truncation / rotation - start over from the top. reset = true; start = 0; } const len = size - start; const buf = Buffer.allocUnsafe(len); let read = 0; while (read < len) { const n = readSync(fd, buf, read, len - read, start + read); if (n <= 0) break; read += n; } const chunk = buf.subarray(0, read); // Find the LAST newline byte in the chunk. Everything up to and including it // is complete lines; any bytes after it are a partial trailing write we drop // until the next fold completes the line. const lastNl = chunk.lastIndexOf(0x0a); // '\n' // Split the complete region on newline bytes, recording each line's ABSOLUTE // byte start (start + relative offset). Byte-accurate: we scan the raw buffer // rather than string-splitting so a multi-byte char never skews an offset. const lines: string[] = []; const lineByteStarts: number[] = []; let lineStart = 0; for (let idx = 0; idx <= lastNl; idx++) { if (chunk[idx] === 0x0a) { lines.push(chunk.subarray(lineStart, idx).toString("utf-8")); lineByteStarts.push(start + lineStart); lineStart = idx + 1; } } let newByteOffset = lastNl >= 0 ? start + lastNl + 1 : start; let trailingPartial = lineStart < chunk.length; // A final flush may observe a fully-written JSON object whose writer never // appended `\n`. Parse only syntactically complete JSON at EOF; malformed // trailing bytes keep the old offset so a later append can complete them. if (flush && trailingPartial) { const trailing = chunk.subarray(lineStart).toString("utf-8"); try { JSON.parse(trailing); lines.push(trailing); lineByteStarts.push(start + lineStart); newByteOffset = start + chunk.length; trailingPartial = false; } catch { // partial final JSON - retain it for a later fold } } return { lines, lineByteStarts, newByteOffset, reset, trailingPartial }; } catch { return null; } finally { try { closeSync(fd); } catch { /* best-effort */ } } } // Offset-aware fold of ONE source file into the ledger. Reads only the new bytes // since the file's cursor, parses+groups them, HOLDS BACK the last (not-yet- // provably-complete) group, folds the rest via the shared accumulation path, and // advances the per-file cursor (byteOffset + lastUuid/lastTimestamp). Mutates // `ledger`. // // CURSOR KEY. The cursor is keyed by the transcript FILE PATH (`path`), NOT the // constant `sourceKey`. The workspace ledger is cumulative across sessions and // each session is a DIFFERENT file, so a constant `"main"` key made one // session's byteOffset get read against another session's longer file, skipping // its early turns. A file path is unique per session, so cursors never collide. // `sourceKey` ("main"/"agent-") is still stamped on every ROW for // byModel/byAgent attribution - only the CURSOR key changed. // // HOLDBACK. Claude Code splits one llm call (one message.id) across contiguous // JSONL lines; a mid-turn fold may land on a chunk boundary INSIDE such a group, // before its real-usage line is written (new-style transcripts carry usage=0 on // the leading lines and the real usage on the LAST line of the run). The fix: // never count a group until it is provably COMPLETE. A group is complete only if // a DIFFERENT-msgId group follows it in this chunk (so a new group has started, // closing the prior one), OR the caller passed `flush=true` (the turn truly // ended - the Stop hook). The LAST group in a non-flush chunk is held back: it // is NOT folded, and byteOffset is rewound to the BYTE START of that group so // its lines are re-read in full next fold. Because a held-back group is never // counted until complete and is always re-read whole, incremental folding equals // a single whole-file flush fold for BOTH transcript styles, regardless of where // a chunk boundary falls (mid-group or mid-line). function foldFileIntoLedger( ledger: Ledger, sourceKey: string, path: string, fallbackAgentId: string | null, metaByAgentId: Record, stageSlug: string | null, sessionKey: string, workflowKey: string, flush: boolean, ): void { // The cursor is keyed by FILE PATH, unique across sessions. const cursorKey = path; const cursor = ledger.cursors[cursorKey]; const prevOffset = cursor?.byteOffset ?? 0; const chunk = readChunkFromOffset(path, prevOffset, flush); if (chunk === null) return; // no new bytes / unreadable // Parse each complete line, keeping the parallel byte-start so a held-back // group can rewind the offset to where its FIRST line begins. const parsed: UsageRow[] = []; const parsedByteStarts: number[] = []; for (let i = 0; i < chunk.lines.length; i++) { const row = parseTranscriptLine(chunk.lines[i], sourceKey, fallbackAgentId); if (row) { parsed.push(row); parsedByteStarts.push(chunk.lineByteStarts[i]); } } // Group contiguous same-msgId runs (the split-line collapse), recording each // group's representative row (max-usage line of the run) and its first line's // byte start. Id-less rows are singleton groups (never merged). Within one // chunk, same-msgId lines are always contiguous (verified), so a contiguous // grouping matches dedupeByMessageId exactly here. const groups: { rep: UsageRow; byteStart: number }[] = []; { let i = 0; while (i < parsed.length) { const row = parsed[i]; const byteStart = parsedByteStarts[i]; if (!row.msgId) { groups.push({ rep: row, byteStart }); i++; continue; } let j = i; while (j + 1 < parsed.length && parsed[j + 1].msgId === row.msgId) j++; groups.push({ rep: representativeOfRun(parsed, i, j), byteStart }); i = j + 1; } } // HOLDBACK. Fold every complete group; on a non-flush fold the LAST group is // not provably complete => hold it back and rewind the offset to its start so // it is re-read whole next time. A flush closes it only when EOF is clean: a // malformed trailing fragment may be another split line for the same message. const closeLastGroup = flush && !chunk.trailingPartial; let foldCount = groups.length; let newByteOffset = chunk.newByteOffset; if (!closeLastGroup && groups.length > 0) { foldCount = groups.length - 1; newByteOffset = groups[groups.length - 1].byteStart; } const toFold = groups.slice(0, foldCount).map((g) => g.rep); // Attribute agent types BEFORE folding (main -> "main", sub -> sidecar type, // else "subagent"), reusing the shared attributor. attributeAgents(toFold, metaByAgentId); const newCursor = { lastUuid: cursor?.lastUuid ?? "", lastTimestamp: cursor?.lastTimestamp ?? "", lastMessageId: cursor?.lastMessageId ?? "", byteOffset: newByteOffset, pending: chunk.reset ? undefined : cursor?.pending, }; const priorPending = chunk.reset ? undefined : cursor?.pending; for (let i = 0; i < toFold.length; i++) { const group = groups[i]; const captured = i === 0 && priorPending !== undefined && priorPending.byteOffset === group.byteStart ? priorPending : null; foldRowIntoLedger( ledger, toFold[i], captured?.stageSlug ?? stageSlug, captured?.sessionKey ?? sessionKey, captured?.workflowKey ?? workflowKey, ); } if (toFold.length > 0) { const last = toFold[toFold.length - 1]; newCursor.lastUuid = last.uuid; newCursor.lastTimestamp = last.timestamp; if (last.msgId) newCursor.lastMessageId = last.msgId; // diagnostics only } if (!closeLastGroup && groups.length > 0) { const held = groups[groups.length - 1]; newCursor.pending = priorPending?.byteOffset === held.byteStart ? priorPending : { byteOffset: held.byteStart, messageId: held.rep.msgId, stageSlug, sessionKey, workflowKey, }; } else if (closeLastGroup) { newCursor.pending = undefined; } ledger.cursors[cursorKey] = newCursor; } // The runtime ledger PRODUCER: read the transcript (main + sub-agents) and fold // its NEW bytes into the ledger under the current stage. This is the runtime // caller wired from the PreToolUse, PostToolUse, and Stop hooks (Claude only). // Offset-aware: each fold reads only bytes appended since the last fold per file // (main + each sub-agent sidecar), so a per-tool-call fold never re-parses whole // files. Fully guarded - never throws into a hook; on any failure returns the // existing ledger unchanged and persists nothing. // // Fold modes. `holdback` retains every file's last group for PostToolUse; // `seal-main` closes only the main group; `flush-all` closes every complete // group at an engine boundary or Stop. export type FoldMode = "holdback" | "seal-main" | "flush-all"; export function foldTranscriptIntoLedger( projectDir: string, transcriptPath: string, stageSlug: string | null, modeOrFlush: FoldMode | boolean = "holdback", context: UsageContext = {}, ): Ledger { // Kill switch: the producer writes nothing (existing ledger left untouched). if (usageTrackingDisabled()) return emptyLedger(); try { return withUsageLedgerLock(projectDir, () => { const mode: FoldMode = typeof modeOrFlush === "boolean" ? modeOrFlush ? "flush-all" : "holdback" : modeOrFlush; const sessionKey = context.sessionKey ?? sessionUsageKey(transcriptPath, context.sessionId); const workflowKey = context.workflowKey ?? intentUsageKey(projectDir, context.sessionId); const ledger = loadLedger(projectDir); // Main transcript. foldFileIntoLedger( ledger, "main", transcriptPath, null, {}, stageSlug, sessionKey, workflowKey, mode !== "holdback", ); // Each sibling sub-agent file, with its sidecar agentType map (rebuilt each // fold - sidecars are tiny and the set can grow between folds). const subDir = subagentDir(transcriptPath); let entries: string[] = []; try { if (existsSync(subDir)) { entries = readdirSync(subDir).filter( (f) => f.startsWith("agent-") && f.endsWith(".jsonl"), ); } } catch { entries = []; } for (const file of entries) { const agentId = file.replace(/^agent-/, "").replace(/\.jsonl$/, ""); const full = join(subDir, file); const metaByAgentId: Record = {}; const meta = readMetaSidecar(full); if (meta && typeof meta.agentType === "string" && meta.agentType) { metaByAgentId[agentId] = meta.agentType; } foldFileIntoLedger( ledger, `agent-${agentId}`, full, agentId, metaByAgentId, stageSlug, sessionKey, workflowKey, mode === "flush-all", ); } // Persist atomically (mkdir the sessions dir first). try { mkdirSync(sessionsDir(projectDir), { recursive: true }); } catch { /* dir may already exist */ } writeFileAtomic(ledgerPath(projectDir), JSON.stringify(ledger, null, 2)); return ledger; }); } catch { return loadLedger(projectDir); } }