From cd7548b3f9f4cc70d6a4ef9e3bb485fe509128d1 Mon Sep 17 00:00:00 2001 From: David Crowe Date: Mon, 21 Sep 2026 10:50:18 -0700 Subject: [PATCH] Send transcript usage after the gateway acks it; remember acked ids behind a held turn; cap the held region (0.24.0) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit gatewaystack-connect#1279, items 1 and 2, client half. Two accounting defects in the hook's transcript collector, both introduced or exposed by #40. 1. Silent under-count on any failure. collectTranscriptUsage wrote the new offset BEFORE the request went out; handlePostToolUse then exited on !res.ok and swallowed abort/network errors. There were no retries, and model_usage is not in the offline ledger, so a 4xx/5xx/timeout on /govern/tool-output lost that window's turns for good. Now the collector returns { turns, ack, fail } and the offset moves only in ack(), which runs after a 2xx. Any other outcome — non-2xx, abort, connection refused — calls fail(): the offset stays, `attempts` is recorded in transcript-offsets.json, and the next PostToolUse re-reads the same window and sends it again. Bounded: after 3 attempts on the same batch the offset moves past it, the loss is one JSON line in ~/.acp/lapse.log (transcript-usage-dropped, logged on the first drop per transcript and counted after), and the ids are NOT pretended acked. 2. Re-sends were the steady state. #40's holdAt pinned the offset at the first incomplete turn so its remaining records get re-read — and every completed turn behind it was re-SENT on every later tool call for the rest of the session, relying on the gateway's per-process Set to absorb it. Worse than the issue describes: holdAt was fixed at the first partial record the loop met, so a turn whose final record landed later in the SAME window still pinned the offset. The existing "withheld, then billed" test now also asserts the offset releases. The offsets file entry is now { off, acked, attempts, dropped } (a bare number, the pre-0.24.0 shape, still reads). `acked` keeps the last 200 turn ids the gateway acknowledged for that transcript; they are filtered out before a batch is built, so the steady state behind a hold is a bounded re-read and an EMPTY send. A re-send now happens in exactly one case: a 2xx whose response was lost in transit. The gateway's durable (session, id) dedupe (gatewaystack-connect PR, same issue) is what absorbs that. 3. The held region is capped at 256 KB. Past that the offset moves on, one lapse line says so (transcript-hold-cap), and the incomplete turn is abandoned — unless its final record does land later, in which case it is a complete turn in a fresh window and bills at its real total. Previously only TRANSCRIPT_READ_CAP (2 MB) bounded the stall, which a normal session never reaches, so every tool call paid the re-read. Tests (test/transcript-model-usage.test.mjs, 23 cases): a 503 leaves the offset and the next call re-sends the identical window; connection refused is a failed send too; three failures drop the batch with one lapse line and the gateway coming back does not re-send it; the legacy bare-number offsets file still reads; the second PostToolUse after a held turn sends zero already-acked ids and the hold releases when the turn completes; the acked set is capped at 200 newest; a 300 KB held region is advanced over once. Note for release: the hook's content hash changes with this, so the registry needs the 0.24.0 hash AFTER merge, never before: fddeccaa3b1db29fb48d1ebd95fc6a81333576105fc19df0ae1c11c6ea639ca5 Co-Authored-By: Claude Fable 5.1 --- .claude-plugin/marketplace.json | 2 +- bin/govern.mjs | 171 ++++++++++++++++++---- plugin.json | 2 +- test/transcript-model-usage.test.mjs | 208 ++++++++++++++++++++++++++- 4 files changed, 349 insertions(+), 34 deletions(-) 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); +});