fix(metrics): bound incremental snapshot recovery

Cap JSONL line buffering and per-hook catch-up work, persist discard cursors for oversized records, report retention failures, normalize malformed token totals, and strengthen bounded-read regression fixtures.
This commit is contained in:
wellkilo
2026-09-13 04:55:52 +08:00
parent c2405149f7
commit 987c1e103f
7 changed files with 424 additions and 107 deletions
+14 -2
View File
@@ -109,7 +109,17 @@ function isSonnet5(model) {
function toNumber(v) {
const n = Number(v);
return Number.isFinite(n) ? n : 0;
return Number.isFinite(n) && n >= 0 ? n : 0;
}
function normalizeUsageTotals(totals) {
return {
inputTokens: toNumber(totals.inputTokens),
outputTokens: toNumber(totals.outputTokens),
cacheWriteTokens: toNumber(totals.cacheWriteTokens),
cacheReadTokens: toNumber(totals.cacheReadTokens),
model: totals.model
};
}
/**
@@ -167,7 +177,9 @@ function sumUsageFromTranscript(transcriptPath) {
cacheReadTokens += toNumber(u.cache_read_input_tokens);
}
return { inputTokens, outputTokens, cacheWriteTokens, cacheReadTokens, model };
return normalizeUsageTotals({
inputTokens, outputTokens, cacheWriteTokens, cacheReadTokens, model
});
}
// 1MB, matching the other Stop hooks. The Stop payload carries
+2 -2
View File
@@ -155,7 +155,7 @@ function readSessionCost(sessionId) {
'malformed',
costsPath,
`${snapshotResult.malformed}:${snapshotResult.malformedSignature}`,
`[ecc-metrics-bridge] skipped ${snapshotResult.malformed} malformed line(s) in ${costsPath}\n`
`[ecc-metrics-bridge] skipped ${snapshotResult.malformed} malformed line(s) during the snapshot scan of ${costsPath}\n`
);
}
if (snapshotResult.invalid > 0) {
@@ -163,7 +163,7 @@ function readSessionCost(sessionId) {
'invalid-row',
costsPath,
`${snapshotResult.invalid}:${snapshotResult.invalidSignature}`,
`[ecc-metrics-bridge] skipped ${snapshotResult.invalid} invalid cumulative row(s) for ${sessionId} in ${costsPath}\n`
`[ecc-metrics-bridge] skipped ${snapshotResult.invalid} invalid cumulative row(s) for ${sessionId} during the snapshot scan of ${costsPath}\n`
);
}
if (snapshotResult.snapshotError) {
+189 -88
View File
@@ -11,6 +11,8 @@ const COST_SNAPSHOT_SCHEMA_VERSION = 'ecc.cost-snapshot.v1';
const COST_SNAPSHOT_DIRECTORY = 'cost-snapshots';
const COST_LOG_FILENAME = 'costs.jsonl';
const READ_CHUNK_BYTES = 64 * 1024;
const MAX_JSONL_LINE_BYTES = 1024 * 1024;
const MAX_SCAN_BYTES = 16 * 1024 * 1024;
const FINGERPRINT_WINDOW_BYTES = 256;
const PRUNE_INTERVAL_MS = 24 * 60 * 60 * 1000;
const SNAPSHOT_MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000;
@@ -109,42 +111,121 @@ function validSnapshotBase(snapshot, descriptor, stat, sessionId) {
source.fingerprint,
fingerprintProcessedPrefix(descriptor, source.offset_bytes)
)) return null;
return { row: snapshot.row, offset: source.offset_bytes };
return {
row: snapshot.row,
offset: source.offset_bytes,
discardingLine: source.discarding_line === true
};
}
function scanJsonlRange(descriptor, start, end, sessionId, initialRow) {
function createScanState(initialRow) {
return {
latestRow: initialRow,
committedRow: initialRow,
malformed: 0,
invalid: 0,
malformedHasher: crypto.createHash('sha256'),
invalidHasher: crypto.createHash('sha256')
};
}
function processCostLine(state, line, sessionId, committed = true) {
if (!line.trim()) return state;
try {
const row = JSON.parse(line);
if (row.session_id !== sessionId) return state;
if (!isValidCostRow(row, sessionId)) {
if (!committed) return state;
return {
...state,
invalid: state.invalid + 1,
invalidHasher: state.invalidHasher.copy().update(line).update('\0')
};
}
return {
...state,
latestRow: chooseNewerCumulativeRow(state.latestRow, row),
committedRow: committed
? chooseNewerCumulativeRow(state.committedRow, row)
: state.committedRow
};
} catch {
if (!committed) return state;
return {
...state,
malformed: state.malformed + 1,
malformedHasher: state.malformedHasher.copy().update(line).update('\0')
};
}
}
function markOversizedLine(state, pendingChunks, segment) {
const hasher = state.malformedHasher.copy();
for (const chunk of pendingChunks) hasher.update(chunk);
const remaining = Math.max(0, MAX_JSONL_LINE_BYTES - pendingChunks.reduce(
(total, chunk) => total + chunk.length,
0
));
hasher.update(segment.subarray(0, remaining)).update('\0<oversized>');
return { ...state, malformed: state.malformed + 1, malformedHasher: hasher };
}
function consumeLineSegment(scan, segment, terminated, sessionId) {
if (scan.discardingLine) {
return { ...scan, discardingLine: !terminated };
}
if (scan.pendingBytes + segment.length > MAX_JSONL_LINE_BYTES) {
return {
state: markOversizedLine(scan.state, scan.pendingChunks, segment),
pendingChunks: [],
pendingBytes: 0,
discardingLine: !terminated
};
}
const pendingChunks = segment.length > 0
? [...scan.pendingChunks, Buffer.from(segment)]
: scan.pendingChunks;
const pendingBytes = scan.pendingBytes + segment.length;
if (!terminated) return { ...scan, pendingChunks, pendingBytes };
const line = Buffer.concat(pendingChunks, pendingBytes).toString('utf8');
return {
state: processCostLine(scan.state, line, sessionId),
pendingChunks: [],
pendingBytes: 0,
discardingLine: false
};
}
function consumeJsonlChunk(scan, chunk, sessionId, absoluteStart, processedOffset) {
let nextScan = scan;
let nextOffset = processedOffset;
let segmentStart = 0;
for (;;) {
const newlineIndex = chunk.indexOf(0x0a, segmentStart);
if (newlineIndex < 0) break;
nextScan = consumeLineSegment(
nextScan, chunk.subarray(segmentStart, newlineIndex), true, sessionId
);
nextOffset = absoluteStart + newlineIndex + 1;
segmentStart = newlineIndex + 1;
}
nextScan = consumeLineSegment(
nextScan, chunk.subarray(segmentStart), false, sessionId
);
if (nextScan.discardingLine) nextOffset = absoluteStart + chunk.length;
return { scan: nextScan, processedOffset: nextOffset };
}
function scanJsonlRange(descriptor, start, end, sessionId, initialRow, initialDiscard = false) {
const buffer = Buffer.allocUnsafe(READ_CHUNK_BYTES);
let lineScan = {
state: createScanState(initialRow),
pendingChunks: [],
pendingBytes: 0,
discardingLine: initialDiscard
};
let position = start;
let processedOffset = start;
let pending = Buffer.alloc(0);
let latestRow = initialRow;
let committedRow = initialRow;
let malformed = 0;
let invalid = 0;
const malformedHasher = crypto.createHash('sha256');
const invalidHasher = crypto.createHash('sha256');
const processLine = (line, committed = true) => {
if (!line.trim()) return;
try {
const row = JSON.parse(line);
if (row.session_id !== sessionId) return;
if (!isValidCostRow(row, sessionId)) {
if (committed) {
invalid += 1;
invalidHasher.update(line).update('\0');
}
return;
}
latestRow = chooseNewerCumulativeRow(latestRow, row);
if (committed) committedRow = chooseNewerCumulativeRow(committedRow, row);
} catch {
if (committed) {
malformed += 1;
malformedHasher.update(line).update('\0');
}
}
};
while (position < end) {
const bytesRead = fs.readSync(
@@ -155,33 +236,40 @@ function scanJsonlRange(descriptor, start, end, sessionId, initialRow) {
position
);
if (bytesRead === 0) break;
const combined = pending.length > 0
? Buffer.concat([pending, buffer.subarray(0, bytesRead)])
: buffer.subarray(0, bytesRead);
let lineStart = 0;
for (;;) {
const newlineIndex = combined.indexOf(0x0a, lineStart);
if (newlineIndex < 0) break;
processLine(combined.subarray(lineStart, newlineIndex).toString('utf8'));
lineStart = newlineIndex + 1;
}
pending = Buffer.from(combined.subarray(lineStart));
const consumed = consumeJsonlChunk(
lineScan, buffer.subarray(0, bytesRead), sessionId, position, processedOffset
);
lineScan = consumed.scan;
processedOffset = consumed.processedOffset;
position += bytesRead;
processedOffset = position - pending.length;
}
if (pending.toString('utf8').trim()) processLine(pending.toString('utf8'), false);
if (lineScan.pendingBytes > 0) {
const line = Buffer.concat(lineScan.pendingChunks, lineScan.pendingBytes).toString('utf8');
lineScan = {
...lineScan,
state: processCostLine(lineScan.state, line, sessionId, false)
};
}
const { state } = lineScan;
return {
row: latestRow,
committedRow,
row: state.latestRow,
committedRow: state.committedRow,
processedOffset,
malformed,
invalid,
malformedSignature: malformed > 0 ? malformedHasher.digest('hex').slice(0, 16) : null,
invalidSignature: invalid > 0 ? invalidHasher.digest('hex').slice(0, 16) : null
malformed: state.malformed,
invalid: state.invalid,
malformedSignature: state.malformed > 0
? state.malformedHasher.digest('hex').slice(0, 16)
: null,
invalidSignature: state.invalid > 0
? state.invalidHasher.digest('hex').slice(0, 16)
: null,
discardingLine: lineScan.discardingLine
};
}
function writeSnapshotAtOffset(metricsDir, sessionId, row, descriptor, stat, offset) {
function writeSnapshotAtOffset(
metricsDir, sessionId, row, descriptor, stat, offset, discardingLine = false
) {
if (row !== null && !isValidCostRow(row, sessionId)) return false;
const snapshot = {
schema_version: COST_SNAPSHOT_SCHEMA_VERSION,
@@ -189,6 +277,7 @@ function writeSnapshotAtOffset(metricsDir, sessionId, row, descriptor, stat, off
identity: sourceIdentity(stat),
offset_bytes: offset,
mtime_ms: stat.mtimeMs,
discarding_line: discardingLine,
fingerprint: fingerprintProcessedPrefix(descriptor, offset)
},
row
@@ -199,9 +288,6 @@ function writeSnapshotAtOffset(metricsDir, sessionId, row, descriptor, stat, off
{
beforeRename() {
const current = fs.fstatSync(descriptor);
if (sourceIdentity(current) !== snapshot.source.identity) {
throw new Error('Cost log identity changed during snapshot publication');
}
if (current.size < offset) {
throw new Error('Cost log was truncated during snapshot publication');
}
@@ -220,6 +306,36 @@ function writeSnapshotAtOffset(metricsDir, sessionId, row, descriptor, stat, off
return true;
}
function emptySnapshotResult(row) {
return {
row,
scannedBytes: 0,
malformed: 0,
invalid: 0,
malformedSignature: null,
invalidSignature: null,
snapshotError: null
};
}
function publishScanSnapshot(metricsDir, sessionId, scan, descriptor, stat) {
if (!scan.committedRow && scan.processedOffset === 0) return null;
try {
writeSnapshotAtOffset(
metricsDir,
sessionId,
scan.committedRow,
descriptor,
stat,
scan.processedOffset,
scan.discardingLine
);
return null;
} catch (error) {
return error;
}
}
function refreshSessionCostSnapshot(metricsDir, sessionId) {
assertSafeSessionId(sessionId);
const costsPath = path.join(metricsDir, COST_LOG_FILENAME);
@@ -228,39 +344,19 @@ function refreshSessionCostSnapshot(metricsDir, sessionId) {
const stat = fs.fstatSync(descriptor);
const snapshot = readJsonFile(getCostSnapshotPath(metricsDir, sessionId));
const base = validSnapshotBase(snapshot, descriptor, stat, sessionId);
if (base?.offset === stat.size) {
return {
row: base.row,
scannedBytes: 0,
malformed: 0,
invalid: 0,
malformedSignature: null,
invalidSignature: null,
snapshotError: null
};
}
if (base?.offset === stat.size) return emptySnapshotResult(base.row);
const scanEnd = Math.min(stat.size, (base?.offset || 0) + MAX_SCAN_BYTES);
const scan = scanJsonlRange(
descriptor,
base?.offset || 0,
stat.size,
scanEnd,
sessionId,
base?.row || null
base?.row || null,
base?.discardingLine || false
);
const snapshotError = publishScanSnapshot(
metricsDir, sessionId, scan, descriptor, stat
);
let snapshotError = null;
if (scan.committedRow || scan.processedOffset > 0) {
try {
writeSnapshotAtOffset(
metricsDir,
sessionId,
scan.committedRow,
descriptor,
stat,
scan.processedOffset
);
} catch (error) {
snapshotError = error;
}
}
return {
row: scan.row,
scannedBytes: scan.processedOffset - (base?.offset || 0),
@@ -308,8 +404,10 @@ function appendSessionCostRow(metricsDir, sessionId, row) {
if (result.snapshotError) throw result.snapshotError;
try {
maybePruneSessionCostSnapshots(metricsDir);
} catch {
// Retention is opportunistic and retried by a later update.
} catch (error) {
// Retention is opportunistic and retried by a later update, but a
// persistent failure remains visible without rolling back the log append.
warnSessionCostSnapshotFailure('retention', metricsDir, sessionId, error);
}
return JSON.stringify(result.row) === JSON.stringify(row);
}
@@ -341,10 +439,10 @@ function maybePruneSessionCostSnapshots(metricsDir, options = {}) {
try {
const intervalIsFresh = now - fs.statSync(markerPath).mtimeMs < PRUNE_INTERVAL_MS;
if (intervalIsFresh && snapshotEntries.length <= maxSnapshots) return 0;
} catch { /* missing marker */ }
} catch (error) {
if (error.code !== 'ENOENT') throw error;
}
}
fs.writeFileSync(markerPath, String(now), { encoding: 'utf8', mode: 0o600 });
const snapshots = snapshotEntries
.map(entry => {
const filePath = path.join(snapshotDir, entry.name);
@@ -361,8 +459,11 @@ function maybePruneSessionCostSnapshots(metricsDir, options = {}) {
fs.rmSync(entry.filePath, { force: true });
removed += 1;
}
} catch { /* already replaced or removed */ }
} catch (error) {
if (error.code !== 'ENOENT') throw error;
}
}
fs.writeFileSync(markerPath, String(now), { encoding: 'utf8', mode: 0o600 });
return removed;
}