diff --git a/.claude-plugin/marketplace.json b/.claude-plugin/marketplace.json index 4cc57be..7fbfff4 100644 --- a/.claude-plugin/marketplace.json +++ b/.claude-plugin/marketplace.json @@ -9,7 +9,7 @@ "name": "agentic-control-plane", "source": "./", "description": "Control, audit, and cost-optimize every Claude Code tool call. Governance hook + bundled ACP MCP (cost X-ray, run traces, policy checks) + /cost-xray pre-ship report.", - "version": "0.23.0", + "version": "0.24.0", "author": { "name": "GatewayStack" }, diff --git a/bin/govern.mjs b/bin/govern.mjs index 347e2ad..295eb53 100644 --- a/bin/govern.mjs +++ b/bin/govern.mjs @@ -67,7 +67,7 @@ const ACP_GOVERN = process.env.ACP_API_BASE || "https://govern.agenticcontrolplane.com"; -const PLUGIN_VERSION = "0.23.0"; +const PLUGIN_VERSION = "0.24.0"; // Console base for user-facing deep links (session receipt, #606). const ACP_CONSOLE = @@ -1389,22 +1389,87 @@ async function handlePreToolUse() { // every branch: a missing, unreadable, truncated or malformed transcript // returns undefined (the field is simply absent — never an empty array, // never a throw) and can never touch the decision or delay the call. -// The offset advances past what is being sent before the request goes -// out, so a turn is normally reported once; the gateway's (session, id) -// dedupe absorbs retries. A partial trailing line is left for next time. +// A partial trailing line is left for next time. +// +// Send-after-ack (gatewaystack-connect#1279). The offset moves only after +// the gateway answered 2xx for the batch that covered it. Until 0.23.0 it +// moved BEFORE the request went out, and a 4xx/5xx/timeout on the +// tool-output call silently lost that window's turns for good — there was +// no retry, and the field is not in the offline ledger. Now a failed send +// leaves the offset where it was, so the next PostToolUse re-reads the same +// window and sends it again. Bounded: after TRANSCRIPT_SEND_ATTEMPTS +// failures on the same batch the offset moves past it anyway and one line +// goes to ~/.acp/lapse.log, so a gateway that rejects every call cannot pin +// a session on a 2 MB re-read forever. A retry after a 2xx whose response +// was lost in transit (the gateway priced the batch, we never heard) is a +// real re-send; the gateway's durable (session, id) dedupe is what absorbs +// that, and that is the only case it now has to absorb routinely. +// +// The offsets file (`{ [transcript]: { off, acked, attempts, dropped } }`; +// a bare number is the pre-0.24.0 shape and still reads) also keeps the +// last TRANSCRIPT_ACKED_KEEP turn ids the gateway acknowledged per +// transcript. They matter behind a held turn: `holdAt` pins the offset at +// the first incomplete turn (see below) so its remaining records are +// re-read, and every completed turn after it gets re-read too — 0.23.0 +// re-SENT those on every tool call for the rest of the session. Acked ids +// are filtered out before the batch is built, so the steady state behind a +// hold is a bounded re-read and an empty send, not a re-send. const TRANSCRIPT_OFFSETS = join(ACP_DIR, "transcript-offsets.json"); const TRANSCRIPT_READ_CAP = 2 * 1024 * 1024; +const TRANSCRIPT_HOLD_CAP = 256 * 1024; const TRANSCRIPT_MAX_TURNS = 50; const TRANSCRIPT_OFFSETS_KEEP = 40; +const TRANSCRIPT_ACKED_KEEP = 200; +const TRANSCRIPT_SEND_ATTEMPTS = 3; +function readTranscriptOffsets() { + let offsets = {}; + try { offsets = JSON.parse(readFileSync(TRANSCRIPT_OFFSETS, "utf8")) || {}; } catch { offsets = {}; } + if (!offsets || typeof offsets !== "object" || Array.isArray(offsets)) offsets = {}; + return offsets; +} +function transcriptEntry(raw) { + const empty = { off: 0, acked: [], attempts: 0, dropped: 0 }; + if (typeof raw === "number") return { ...empty, off: raw >= 0 ? raw : 0 }; + if (!raw || typeof raw !== "object") return empty; + const num = (v) => (typeof v === "number" && Number.isFinite(v) && v > 0 ? Math.floor(v) : 0); + return { + off: num(raw.off), + acked: Array.isArray(raw.acked) ? raw.acked.filter((x) => typeof x === "string").slice(-TRANSCRIPT_ACKED_KEEP) : [], + attempts: num(raw.attempts), + dropped: num(raw.dropped), + }; +} +// Keep the file small: transcripts that no longer exist are dropped, at +// most TRANSCRIPT_OFFSETS_KEEP others are kept. Best-effort — at worst a +// window is reported twice, and the gateway dedupes. +function writeTranscriptEntry(offsets, path, entry) { + try { + const keep = {}; + let n = 0; + for (const [k, v] of Object.entries(offsets)) { if (k !== path && existsSync(k) && n++ < TRANSCRIPT_OFFSETS_KEEP) keep[k] = v; } + keep[path] = entry; + mkdirSync(ACP_DIR, { recursive: true }); + writeFileSync(TRANSCRIPT_OFFSETS, JSON.stringify(keep)); + } catch { /* best-effort */ } +} +// Under-counted model usage is a lapse in the cost record; it goes where +// the other lapses go, one JSON line per event, never to the model. +function transcriptLapse(event, detail) { + try { + appendFileSync(join(ACP_DIR, "lapse.log"), JSON.stringify({ at: new Date().toISOString(), event, ...detail }) + "\n"); + } catch { /* best-effort */ } +} +// Returns undefined when there is nothing to send, else { turns, ack, fail }: +// the caller sends `turns` and reports the outcome with exactly one of +// ack() (gateway said 2xx) or fail() (anything else). Neither throws. function collectTranscriptUsage(path) { try { if (typeof path !== "string" || !path || !existsSync(path)) return undefined; - let offsets = {}; - try { offsets = JSON.parse(readFileSync(TRANSCRIPT_OFFSETS, "utf8")) || {}; } catch { offsets = {}; } - if (!offsets || typeof offsets !== "object" || Array.isArray(offsets)) offsets = {}; + const offsets = readTranscriptOffsets(); + const entry = transcriptEntry(offsets[path]); let size = 0; try { size = statSync(path).size; } catch { return undefined; } - let start = typeof offsets[path] === "number" && offsets[path] >= 0 && offsets[path] <= size ? offsets[path] : 0; + let start = entry.off <= size ? entry.off : 0; if (size - start > TRANSCRIPT_READ_CAP) start = size - TRANSCRIPT_READ_CAP; if (size <= start) return undefined; let chunk = ""; @@ -1419,7 +1484,8 @@ function collectTranscriptUsage(path) { // Only consume complete lines; a partial trailing line waits for next time. const lastNl = chunk.lastIndexOf("\n"); if (lastNl < 0) return undefined; - let consumed = Buffer.byteLength(chunk.slice(0, lastNl + 1), "utf8"); + const lineEnd = Buffer.byteLength(chunk.slice(0, lastNl + 1), "utf8"); + let consumed = lineEnd; // One API call is appended as MANY records — one per content block as it // streams. cache_read / cache_creation / input_tokens are fixed at request // time and identical on every record, but output_tokens is a placeholder @@ -1436,10 +1502,19 @@ function collectTranscriptUsage(path) { // first byte of the first incomplete turn so its remaining records are // re-read next time. A turn that never completes (an aborted or // interrupted request — ~0.7% of calls) is never sent; losing those beats - // billing a placeholder for them. The stall is bounded: TRANSCRIPT_READ_CAP - // eventually slides `start` past a turn that is never finished. + // billing a placeholder for them. The stall is bounded twice over: + // TRANSCRIPT_HOLD_CAP moves the offset past a held region that has grown + // beyond 256 KB (one lapse line says so), and TRANSCRIPT_READ_CAP slides + // `start` regardless. const turns = []; const complete = new Set(); + // First byte (relative to `start`) of each turn's first partial record. + // The hold is decided AFTER the loop, on turns still incomplete at the + // end of the window: 0.23.0 fixed it at the first partial record it met, + // so a turn whose final record arrived later in the SAME window kept the + // offset pinned anyway — and every completed turn was re-sent on every + // call for the rest of the session. + const pendingAt = new Map(); let holdAt = -1; // byte offset (relative to `start`) of the first incomplete turn let lineStart = 0; for (const line of chunk.slice(0, lastNl).split("\n")) { @@ -1454,7 +1529,7 @@ function collectTranscriptUsage(path) { if (!id || typeof m.model !== "string" || !m.model) continue; const done = typeof m.stop_reason === "string" && m.stop_reason !== ""; if (done) complete.add(id); - else if (!complete.has(id) && holdAt < 0) holdAt = thisLineStart; + else if (!complete.has(id) && !pendingAt.has(id)) pendingAt.set(id, thisLineStart); turns.push({ id, model: m.model, input_tokens: u.input_tokens || 0, @@ -1464,24 +1539,54 @@ function collectTranscriptUsage(path) { ts: typeof e.timestamp === "string" ? e.timestamp : undefined, }); } - if (holdAt >= 0) consumed = Math.min(consumed, holdAt); - // Advance the offset only past what we are sending; keep the file small - // (transcripts that no longer exist are dropped, at most - // TRANSCRIPT_OFFSETS_KEEP others are kept). - try { - const keep = {}; - let n = 0; - for (const [k, v] of Object.entries(offsets)) { if (k !== path && existsSync(k) && n++ < TRANSCRIPT_OFFSETS_KEEP) keep[k] = v; } - keep[path] = start + consumed; - mkdirSync(ACP_DIR, { recursive: true }); - writeFileSync(TRANSCRIPT_OFFSETS, JSON.stringify(keep)); - } catch { /* best-effort — at worst a window is reported twice, and the gateway dedupes */ } + for (const [id, at] of pendingAt) if (!complete.has(id) && (holdAt < 0 || at < holdAt)) holdAt = at; + if (holdAt >= 0) { + if (lineEnd - holdAt > TRANSCRIPT_HOLD_CAP) { + // The incomplete turn is abandoned (never sent) and the window + // moves on; if its final record does land later it is a complete + // turn in a later window and bills at its real total then. + transcriptLapse("transcript-hold-cap", { transcript: path, held_bytes: lineEnd - holdAt }); + } else { + consumed = holdAt; + } + } // Same-id turns can repeat across streamed chunks; keep the last (fullest), - // and only turns whose final record we actually saw (see above). + // only turns whose final record we actually saw, and none the gateway has + // already acknowledged (the re-read behind a hold). + const acked = new Set(entry.acked); const byId = new Map(); - for (const t of turns) if (complete.has(t.id)) byId.set(t.id, t); + for (const t of turns) if (complete.has(t.id) && !acked.has(t.id)) byId.set(t.id, t); const out = Array.from(byId.values()).slice(-TRANSCRIPT_MAX_TURNS); - return out.length ? out : undefined; + const next = start + consumed; + if (!out.length) { + // Nothing to send, so nothing to lose: move on now. No write when the + // offset would not change — the steady state behind a hold is a read. + if (next !== entry.off || entry.attempts) writeTranscriptEntry(offsets, path, { ...entry, off: next, attempts: 0 }); + return undefined; + } + const sentIds = out.map((t) => t.id); + return { + turns: out, + settled: false, + ack() { + this.settled = true; + writeTranscriptEntry(offsets, path, { ...entry, off: next, acked: entry.acked.concat(sentIds).slice(-TRANSCRIPT_ACKED_KEEP), attempts: 0 }); + }, + fail() { + this.settled = true; + const attempts = entry.attempts + 1; + if (attempts < TRANSCRIPT_SEND_ATTEMPTS) { + writeTranscriptEntry(offsets, path, { ...entry, attempts }); + return; + } + // Give up on this batch: the offset moves past it, the ids are NOT + // marked acked (they were not), and the loss is on record. Logged on + // the first drop per transcript; later drops are counted, not logged. + const dropped = entry.dropped + 1; + if (dropped === 1) transcriptLapse("transcript-usage-dropped", { transcript: path, turns: out.length, attempts }); + writeTranscriptEntry(offsets, path, { ...entry, off: next, attempts: 0, dropped }); + }, + }; } catch { return undefined; } } @@ -1510,7 +1615,8 @@ async function handlePostToolUse() { // until now — proxy and transcript rows ~100ms apart with identical cost on // every turn of a `claude-acp` session. The gateway carries the same guard // for clients that never upgrade (proxy/proxySessions.ts). - const modelUsage = modelCallsGoThroughProxy() ? undefined : collectTranscriptUsage(input.transcript_path); + const collected = modelCallsGoThroughProxy() ? undefined : collectTranscriptUsage(input.transcript_path); + const modelUsage = collected ? collected.turns : undefined; const body = JSON.stringify({ tool_name: input.tool_name, tool_input: input.tool_input, @@ -1529,7 +1635,10 @@ async function handlePostToolUse() { try { const res = await fetch(`${ACP_GOVERN}/govern/tool-output`, { method: "POST", headers, body, signal: controller.signal }); clearTimeout(timeout); - if (!res.ok) { process.exit(0); } + // The offset behind `model_usage` moves only on a 2xx (#1279); any + // other outcome leaves it for the next PostToolUse to retry, bounded. + if (!res.ok) { if (collected) collected.fail(); process.exit(0); } + if (collected) collected.ack(); if (preLapse) clearPendingLapse(input.session_id); const data = await res.json(); // Receipt bookkeeping (#606): one governed call, plus what ACP said @@ -1561,7 +1670,9 @@ async function handlePostToolUse() { process.stdout.write(JSON.stringify({ systemMessage: staleNotice })); } } catch { - // silent pass-through + // silent pass-through for the verdict; the usage batch is retried next + // time (a 2xx whose body failed to parse was already acked above). + if (collected && !collected.settled) collected.fail(); } finally { clearTimeout(timeout); } process.exit(0); } diff --git a/plugin.json b/plugin.json index 8b46401..029c0fa 100644 --- a/plugin.json +++ b/plugin.json @@ -1,6 +1,6 @@ { "name": "agentic-control-plane", - "version": "0.23.0", + "version": "0.24.0", "description": "Identity, governance, and audit for every Claude Code tool call. Logs all tool usage, enforces policies, and gives teams full visibility \u2014 without changing how you use Claude.", "author": { "name": "GatewayStack", diff --git a/test/transcript-model-usage.test.mjs b/test/transcript-model-usage.test.mjs index 68dbad2..23c0421 100644 --- a/test/transcript-model-usage.test.mjs +++ b/test/transcript-model-usage.test.mjs @@ -28,7 +28,7 @@ import { test, before, after, beforeEach } from "node:test"; import assert from "node:assert/strict"; import { spawn } from "node:child_process"; -import { mkdtempSync, mkdirSync, writeFileSync, appendFileSync, readFileSync, rmSync } from "node:fs"; +import { mkdtempSync, mkdirSync, writeFileSync, appendFileSync, readFileSync, rmSync, statSync } from "node:fs"; import { createServer } from "node:http"; import { tmpdir } from "node:os"; import { join, dirname } from "node:path"; @@ -55,6 +55,13 @@ before(async () => { try { body = JSON.parse(raw); } catch { /* keep null */ } seen.push({ url: req.url, body }); res.setHeader("content-type", "application/json"); + // A session whose id starts with "fail-" gets a 503 on tool-output: + // the gateway saw the batch and refused it (#1279). + if (req.url === "/govern/tool-output" && String(body?.session_id ?? "").startsWith("fail-")) { + res.statusCode = 503; + res.end(JSON.stringify({ error: "unavailable" })); + return; + } res.end(JSON.stringify(req.url === "/govern/tool-use" ? { decision: "allow" } : { action: "pass" })); }); }); @@ -205,7 +212,8 @@ test("offset advances: the second call sends only what was appended; nothing new await runHook(post("sess-usage-7", path)); assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_E"]); const offsets = JSON.parse(readFileSync(join(HOME, ".acp", "transcript-offsets.json"), "utf8")); - assert.equal(offsets[path], Buffer.byteLength(assistantLine("msg_E", usage(1, 0, 0, 1)))); + assert.equal(offsets[path].off, Buffer.byteLength(assistantLine("msg_E", usage(1, 0, 0, 1)))); + assert.deepEqual(offsets[path].acked, ["msg_E"], "the acked id rides in the same file"); seen = []; await runHook(post("sess-usage-7", path)); @@ -321,6 +329,14 @@ test("a turn still streaming is withheld, then billed at its real total once com { id: "msg_S", model: "claude-fable-5", input_tokens: 2, cache_read_input_tokens: 126148, cache_creation_input_tokens: 0, output_tokens: 10904, ts: "2026-09-16T21:00:00.000Z" }, ], "the completed turn must bill 10904 output tokens, not the 3-token placeholder"); + // The hold releases once the turn completes: 0.23.0 fixed holdAt at the + // first partial record it met, so the offset stayed pinned at 0 here and + // msg_S was re-sent on every later call. + const off = JSON.parse(readFileSync(join(HOME, ".acp", "transcript-offsets.json"), "utf8"))[path]; + assert.equal(off.off, statSync(path).size, "offset moves past a turn once its final record was seen"); + seen = []; + await runHook(post("sess-partial", path)); + assert.ok(!("model_usage" in lastToolOutput()), "no re-send after the hold released"); }); test("an incomplete turn does not block the completed turns before it", async () => { @@ -346,3 +362,191 @@ test("a plain session still reports, so the proxy check is not just disabling th await runHook(post("sess-plain-2", path), { ANTHROPIC_BASE_URL: "https://api.anthropic.com" }); assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_Q"]); }); + +// --------------------------------------------------------------------- +// Send-after-ack (gatewaystack-connect#1279). +// +// Until 0.23.0 the offset moved before the request went out, so any +// non-2xx, timeout or network error on /govern/tool-output lost that +// window's turns for good: no retry, and model_usage is not in the offline +// ledger. Now the offset moves only on a 2xx; a failed batch is re-sent by +// the next PostToolUse, at most TRANSCRIPT_SEND_ATTEMPTS times. +// --------------------------------------------------------------------- + +const offsetsFile = () => JSON.parse(readFileSync(join(HOME, ".acp", "transcript-offsets.json"), "utf8")); +const lapseLog = () => { try { return readFileSync(join(HOME, ".acp", "lapse.log"), "utf8"); } catch { return ""; } }; + +test("a non-2xx leaves the offset where it was; the next call re-sends the same turns", async () => { + const path = transcriptPath(); + writeFileSync(path, assistantLine("msg_R1", usage(1, 0, 0, 10)) + assistantLine("msg_R2", usage(1, 0, 0, 20))); + let r = await runHook(post("fail-1", path)); + assert.equal(r.code, 0, "a refused batch never changes the hook's exit"); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_R1", "msg_R2"]); + let off = offsetsFile()[path]; + assert.equal(off.off, 0, "offset must not move on a 503"); + assert.equal(off.attempts, 1); + assert.deepEqual(off.acked, []); + + // Gateway back: the identical window goes out again and is acked. + seen = []; + await runHook(post("sess-recovered", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_R1", "msg_R2"]); + off = offsetsFile()[path]; + assert.equal(off.off, Buffer.byteLength(assistantLine("msg_R1", usage(1, 0, 0, 10)) + assistantLine("msg_R2", usage(1, 0, 0, 20)))); + assert.equal(off.attempts, 0); + assert.deepEqual(off.acked, ["msg_R1", "msg_R2"]); + + // And nothing is sent twice after the ack. + seen = []; + await runHook(post("sess-recovered", path)); + assert.ok(!("model_usage" in lastToolOutput())); +}); + +test("a network error (connection refused) is a failed send too, not a silent advance", async () => { + const path = transcriptPath(); + writeFileSync(path, assistantLine("msg_N1", usage(1, 0, 0, 10))); + const r = await runHook(post("sess-net", path), { ACP_GOVERN_BASE: `http://${LOOPBACK}:1` }); + assert.equal(r.code, 0); + const off = offsetsFile()[path]; + assert.equal(off.off, 0); + assert.equal(off.attempts, 1); + seen = []; + await runHook(post("sess-net", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_N1"]); +}); + +test("after three failed attempts the batch is dropped, the offset moves on, and one lapse line records it", async () => { + const path = transcriptPath(); + const line = assistantLine("msg_X1", usage(1, 0, 0, 10)); + writeFileSync(path, line); + await runHook(post("fail-2", path)); + await runHook(post("fail-2", path)); + assert.equal(offsetsFile()[path].attempts, 2); + assert.equal(offsetsFile()[path].off, 0); + seen = []; + await runHook(post("fail-2", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_X1"], "the third attempt still tries"); + const off = offsetsFile()[path]; + assert.equal(off.off, Buffer.byteLength(line), "given up: the window is consumed"); + assert.equal(off.attempts, 0); + assert.equal(off.dropped, 1); + assert.deepEqual(off.acked, [], "dropped ids are not pretended acked"); + const dropLines = lapseLog().split("\n").filter((l) => l.includes("transcript-usage-dropped")); + assert.equal(dropLines.length, 1); + assert.ok(dropLines[0].includes(path)); + + // The gateway is back: the dropped turn is NOT re-sent (the offset is past it). + seen = []; + await runHook(post("sess-after-drop", path)); + assert.ok(!("model_usage" in lastToolOutput())); + + // A second dropped batch on the same transcript is counted, not logged again. + appendFileSync(path, assistantLine("msg_X2", usage(1, 0, 0, 10))); + for (let i = 0; i < 3; i++) await runHook(post("fail-2", path)); + assert.equal(offsetsFile()[path].dropped, 2); + assert.equal(lapseLog().split("\n").filter((l) => l.includes("transcript-usage-dropped")).length, 1); +}); + +test("a pre-0.24.0 offsets file (bare number) still reads as the offset", async () => { + const path = transcriptPath(); + const first = assistantLine("msg_L1", usage(1, 0, 0, 1)); + writeFileSync(path, first + assistantLine("msg_L2", usage(1, 0, 0, 2))); + writeFileSync(join(HOME, ".acp", "transcript-offsets.json"), JSON.stringify({ [path]: Buffer.byteLength(first) })); + await runHook(post("sess-legacy", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_L2"]); + assert.deepEqual(offsetsFile()[path].acked, ["msg_L2"]); +}); + +// --------------------------------------------------------------------- +// Acked ids behind a held turn (#1279 item 1). +// +// holdAt pins the offset at the first incomplete turn so its remaining +// records get re-read. 0.23.0 also re-SENT every completed turn after it +// on every later PostToolUse, for the rest of the session. The acked set in +// the offsets file makes those an empty send. +// --------------------------------------------------------------------- + +test("second PostToolUse after a held turn sends zero already-acked ids", async () => { + const path = transcriptPath(); + // Held turn first, then a completed one behind it: the offset must stay at + // msg_live, and msg_done1 must be sent exactly once. + writeFileSync(path, partialLine("msg_live", 99, 2) + assistantLine("msg_done1", usage(1, 10, 0, 500))); + await runHook(post("sess-held", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_done1"]); + let off = offsetsFile()[path]; + assert.equal(off.off, 0, "offset is pinned at the held turn"); + assert.deepEqual(off.acked, ["msg_done1"]); + + // Same window re-read, nothing new: nothing sent, nothing re-sent. + seen = []; + await runHook(post("sess-held", path)); + assert.ok(!("model_usage" in lastToolOutput()), "the re-read behind the hold must not re-send msg_done1"); + + // Another completed turn lands behind the hold: only IT goes out. + seen = []; + appendFileSync(path, assistantLine("msg_done2", usage(1, 10, 0, 600))); + await runHook(post("sess-held", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_done2"]); + off = offsetsFile()[path]; + assert.equal(off.off, 0); + assert.deepEqual(off.acked, ["msg_done1", "msg_done2"]); + + // The held turn finally completes: it is sent at its real total, alone, + // and the offset moves past everything. + seen = []; + appendFileSync(path, assistantLine("msg_live", usage(2, 99, 0, 4321))); + await runHook(post("sess-held", path)); + const sent = lastToolOutput().model_usage; + assert.deepEqual(sent.map((t) => t.id), ["msg_live"]); + assert.equal(sent[0].output_tokens, 4321); + off = offsetsFile()[path]; + assert.equal(off.off, statSync(path).size); + assert.deepEqual(off.acked, ["msg_done1", "msg_done2", "msg_live"]); +}); + +test("the acked set is capped at 200 ids, newest kept", async () => { + const path = transcriptPath(); + let s = partialLine("msg_hold", 1, 1); + for (let i = 0; i < 50; i++) s += assistantLine(`msg_a${i}`, usage(1, 0, 0, 1)); + writeFileSync(path, s); + for (let round = 0; round < 5; round++) { + await runHook(post("sess-cap", path)); + let more = ""; + for (let i = 0; i < 50; i++) more += assistantLine(`msg_r${round}_${i}`, usage(1, 0, 0, 1)); + appendFileSync(path, more); + } + const off = offsetsFile()[path]; + assert.equal(off.acked.length, 200); + assert.equal(off.acked.at(-1), "msg_r3_49"); + assert.ok(!off.acked.includes("msg_a0"), "the oldest ids are evicted first"); +}); + +test("a held region past 256 KB is advanced over, once, with a lapse line", async () => { + const path = transcriptPath(); + // 3 completed turns of ~100 KB each behind a turn that never finishes. + const big = "x".repeat(100 * 1024); + const bigLine = (id) => assistantLine(id, usage(1, 0, 0, 1), { padding: big }); + writeFileSync(path, partialLine("msg_stuck", 5, 3) + bigLine("msg_big1") + bigLine("msg_big2") + bigLine("msg_big3")); + await runHook(post("sess-huge", path)); + assert.deepEqual(lastToolOutput().model_usage.map((t) => t.id), ["msg_big1", "msg_big2", "msg_big3"]); + const off = offsetsFile()[path]; + assert.equal(off.off, statSync(path).size, "the window moved past the held turn"); + const capLines = lapseLog().split("\n").filter((l) => l.includes("transcript-hold-cap")); + assert.equal(capLines.length, 1); + assert.ok(capLines[0].includes(path)); + + // Nothing new → nothing read again, nothing logged again. + seen = []; + await runHook(post("sess-huge", path)); + assert.ok(!("model_usage" in lastToolOutput())); + assert.equal(lapseLog().split("\n").filter((l) => l.includes("transcript-hold-cap")).length, 1); + + // If the abandoned turn's final record does land, it is a complete turn in + // a fresh window and bills at its real total — nothing was lost. + seen = []; + appendFileSync(path, assistantLine("msg_stuck", usage(1, 5, 0, 777))); + await runHook(post("sess-huge", path)); + const sent = lastToolOutput().model_usage; + assert.deepEqual(sent.map((t) => t.id), ["msg_stuck"]); + assert.equal(sent[0].output_tokens, 777); +});