From 3ccc6ffe506b13a5d300a0bce0e8048f6eef581d Mon Sep 17 00:00:00 2001 From: Dave Jong Date: Mon, 28 Sep 2026 16:06:51 +0200 Subject: [PATCH 1/2] Account for reporting outcomes and compare refresh secrets in full Report each failing host callback once, rather than once per hook name for the whole process, with a process-wide cap on warnings. Count block-log records as delivered, failed or dropped, exposed as protection.blockLogHealth(). Compare the refresh secret through fixed-length digests. Co-Authored-By: Claude Opus 5.5 --- AGENT-INSTALL.md | 5 +- src/protect/firewall-log.js | 48 +++++- src/protect/notify.js | 45 +++-- src/protect/protect.d.ts | 20 ++- src/protect/rules/refresh.js | 43 ++++- src/protect/runtime.js | 5 + tests/protect/reporting-accounting.test.ts | 183 +++++++++++++++++++++ tests/protect/rules-refresh.test.ts | 4 +- 8 files changed, 329 insertions(+), 24 deletions(-) create mode 100644 tests/protect/reporting-accounting.test.ts diff --git a/AGENT-INSTALL.md b/AGENT-INSTALL.md index dc40cdbc..c1ba8d97 100644 --- a/AGENT-INSTALL.md +++ b/AGENT-INSTALL.md @@ -457,8 +457,9 @@ requests are aborted, it starts nothing further, and it discards what it was hol request to stop, not a guarantee — a transport that ignores it is detached rather than completed, so "resolved" means the reporter is finished with it, and a runtime that kills the process still wins regardless. Every detection event ends up delivered, refused or dropped and is reported in the health -counts; block-log records have no counters, so one lost to a failed send or an expired shutdown is -reported nowhere. +counts. The block log keeps the same kind of local counts: `protection.blockLogHealth()` returns how many +block records were accepted, delivered, failed or dropped, and how many are still queued. Those counts +carry no request data and stay in your process. The client address is reported with its **provenance**, because an address is only as trustworthy as whatever supplied it. `client_ip_source` is one of `runtime` (the address the transport observed), diff --git a/src/protect/firewall-log.js b/src/protect/firewall-log.js index 48ade62c..64cd5821 100644 --- a/src/protect/firewall-log.js +++ b/src/protect/firewall-log.js @@ -79,7 +79,12 @@ export function resolveApiBase(pulseOrManifestUrl) { export function createFirewallLogReporter(opts) { const creds = parseApiKey(opts.apiKey); if (!creds) { - return { record() {}, flush: () => Promise.resolve(), stop: () => Promise.resolve() }; + return { + record() {}, + flush: () => Promise.resolve(), + stop: () => Promise.resolve(), + health: () => ({ recorded: 0, delivered: 0, failed: 0, dropped: 0, queued: 0 }), + }; } const apiBase = safeBaseUrl(opts.apiBase, DEFAULT_API_BASE, 'block-log API').replace(/\/$/, ''); @@ -102,6 +107,15 @@ export function createFirewallLogReporter(opts) { /** Set when a shutdown gives up waiting: nothing may start, continue, or be retained after it. */ let ended = false; + // Where every record went. `recorded` is what the queue accepted, and each accepted record ends up + // delivered (the endpoint acknowledged its batch), failed (its batch was refused or could not be sent), + // or dropped (a shutdown ran out of time before it was sent). `dropped` also counts records turned + // away because the queue was full, which were never accepted. + let recorded = 0; + let delivered = 0; + let failed = 0; + let dropped = 0; + /** @type {{ token: string, expiresAt: number } | null} */ let cachedToken = null; /** @type {Promise | null} */ @@ -168,11 +182,18 @@ export function createFirewallLogReporter(opts) { // whether it succeeded. /** @type {Promise} */ let entry; + let settled = false; entry = (async () => { try { const token = await fetchAccessToken(controller?.signal); // Not after a shutdown gave up: it has already reported itself finished. - if (!token || ended) return; + if (ended) return; + if (!token) { + failed += batch.length; + settled = true; + + return; + } const body = new URLSearchParams(); body.set('type', 'firewall'); @@ -191,10 +212,19 @@ export function createFirewallLogReporter(opts) { // Both phases carry the attempt controller, so slow transports cannot accumulate work. ...(controller ? { signal: controller.signal } : {}), }); - if (p && typeof p.then === 'function') await p.catch(() => {}); + const res = p && typeof p.then === 'function' ? await p.catch(() => null) : p; + if (ended) return; + if (res && res.ok) delivered += batch.length; + else failed += batch.length; + settled = true; } catch { /* A delivery problem is never worth disturbing the app over. */ } finally { + // Anything not accounted for above was abandoned by a shutdown, or failed before a verdict. + if (!settled) { + if (ended) dropped += batch.length; + else failed += batch.length; + } clearTimeout(attemptTimer); if (activeController === controller) activeController = null; if (inFlight === entry) inFlight = null; @@ -225,7 +255,12 @@ export function createFirewallLogReporter(opts) { const fid = event?.rule?.id; if (fid === undefined || fid === null || fid === '') return; - if (queue.length >= MAX_QUEUE) return; + if (queue.length >= MAX_QUEUE) { + dropped++; + + return; + } + recorded++; queue.push({ fid, method: event.method ?? null, @@ -242,6 +277,10 @@ export function createFirewallLogReporter(opts) { if (!timer) timer = setTimeout(flush, flushMs); }, flush, + /** Where the records went so far; see the counters above. `queued` is what is waiting now. */ + health() { + return { recorded, delivered, failed, dropped, queued: queue.length }; + }, /** * Stop, and hand back a wait for what was outstanding. * @@ -275,6 +314,7 @@ export function createFirewallLogReporter(opts) { ended = true; activeController?.abort(); activeController = null; + dropped += queue.length; queue = []; inFlight = null; }; diff --git a/src/protect/notify.js b/src/protect/notify.js index ec3b7177..a299c96c 100644 --- a/src/protect/notify.js +++ b/src/protect/notify.js @@ -12,9 +12,11 @@ * applied to the hooks that were missed rather than restated for one of them. * * Not silent, though. A hook that throws is a bug in the host's code and swallowing it entirely would - * hide it forever, so the first failure per hook is reported — once, because these run per request and a - * persistently broken hook would otherwise print on every one. Same reasoning as the engine's - * report-once for a persistently broken rule. + * hide it forever, so the first failure of each callback is reported — once, because these run per request + * and a persistently broken hook would otherwise print on every one. Same reasoning as the engine's + * report-once for a persistently broken rule. "Each callback" is each function under each hook name, so a + * second guard's broken hook is reported even after the first guard's was; the total is capped, so a host + * that creates a fresh callback per request cannot turn this into per-request output. * * @param {unknown} fn the callback, or anything that is not a function (then this is a no-op) * @param {unknown} arg the single argument to hand it @@ -25,8 +27,12 @@ * caller has already decided not to fall back. Synchronous handlers get the stronger answer. */ -/** Hooks already reported as broken. Module-scoped: one warning per hook per process, not per guard. */ -const reported = new Set(); +/** Callbacks already reported as broken, by function, then by hook name. */ +let reported = new WeakMap(); + +/** Warnings written so far, and the most this process writes before saying the rest are suppressed. */ +let warnings = 0; +const MAX_WARNINGS = 20; export function notify(fn, arg, label) { if (typeof fn !== 'function') return false; @@ -40,7 +46,7 @@ export function notify(fn, arg, label) { // able to kill the app, which is a worse outcome than the throw we set out to contain. if (result !== null && typeof result === 'object' && typeof result.then === 'function') { try { - result.then(undefined, (err) => warnOnce(label, err)); + result.then(undefined, (err) => warnOnce(fn, label, err)); } catch { // A `then` that throws on access. Nothing more to attach to; the value is not a usable promise. } @@ -48,27 +54,39 @@ export function notify(fn, arg, label) { return true; } catch (err) { - warnOnce(label, err); + warnOnce(fn, label, err); return false; } } /** - * Report a broken callback once per process. + * Report a broken callback once. * * Must not throw: it runs inside the containment, so its own failure would be the thing that breaks the * guarantee it exists to report on. */ -function warnOnce(label, err) { - if (reported.has(label)) return; - reported.add(label); +function warnOnce(fn, label, err) { + let labels = reported.get(fn); + if (labels?.has(label)) return; + if (!labels) { + labels = new Set(); + reported.set(fn, labels); + } + labels.add(label); + if (warnings > MAX_WARNINGS) return; + warnings++; try { + if (warnings > MAX_WARNINGS) { + console.warn('Patchstack: further failing callbacks passed to createProtection are not reported in this process.'); + + return; + } // Named as the host's callback, not as a Patchstack failure: pointing at ourselves for someone // else's throw sends them reading the wrong code. console.warn( `Patchstack: the ${label} callback passed to createProtection failed and was ignored. ` + - `Protection is unaffected; this is reported once per process. ` + + `Protection is unaffected; this callback's failures are reported once. ` + `Cause: ${err && err.message ? err.message : String(err)}`, ); } catch { @@ -78,5 +96,6 @@ function warnOnce(label, err) { /** Test seam: forget which hooks have been reported, so warn-once is assertable more than once. */ export function resetNotifyWarnings() { - reported.clear(); + reported = new WeakMap(); + warnings = 0; } diff --git a/src/protect/protect.d.ts b/src/protect/protect.d.ts index 1047bb30..0b210818 100644 --- a/src/protect/protect.d.ts +++ b/src/protect/protect.d.ts @@ -62,9 +62,9 @@ export interface Protection { * reporter is finished with it, not that the underlying request has ended. * - A runtime that terminates the process regardless still wins, whatever this resolves. * - * What is accounted for also differs by reporter. Every detection event ends up delivered, refused or - * dropped, and `detectionHealth()` reports each. Block-log records have no counters at all, so one lost - * to a failed token exchange, a failed post, or a shutdown that ran out of time is reported nowhere. + * Both reporters account for what they held. Every detection event ends up delivered, refused or + * dropped, and `detectionHealth()` reports each; every block-log record ends up delivered, failed or + * dropped, and `blockLogHealth()` reports those. */ stop: () => Promise; /** Alias of `stop`, under the name callers already have. */ @@ -128,6 +128,20 @@ export interface Protection { lastAcknowledgedAt: string | null; }; }; + /** Present when block-log reporting is on — where the block records went, in records. Carries no + * request data. */ + blockLogHealth?: () => { + /** Records the queue accepted. Each ends up delivered, failed, or dropped by a shutdown. */ + recorded: number; + /** Records in a batch the endpoint acknowledged. */ + delivered: number; + /** Records in a batch that was refused, or that could not be sent (including a failed token exchange). */ + failed: number; + /** Records a shutdown discarded, plus records turned away because the queue was full. */ + dropped: number; + /** Records waiting to be sent now. */ + queued: number; + }; } /** diff --git a/src/protect/rules/refresh.js b/src/protect/rules/refresh.js index 7d3f941a..d8f1bfc4 100644 --- a/src/protect/rules/refresh.js +++ b/src/protect/rules/refresh.js @@ -159,7 +159,7 @@ export function makeRefreshHandler(tick, secret) { // No secret configured → the endpoint doesn't exist (never an open refresh-DoS surface). if (!secret) return new Response('not found', { status: 404 }); const provided = request?.headers?.get?.('x-patchstack-refresh') ?? null; - if (provided !== secret) return new Response('forbidden', { status: 403 }); + if (!(await sameSecret(provided, secret))) return new Response('forbidden', { status: 403 }); let refreshed = true; try { // `{ ok: false }` means the tick ran but the rules did not come from the source, which is not a @@ -172,3 +172,44 @@ export function makeRefreshHandler(tick, secret) { return new Response(JSON.stringify({ refreshed }), { status: 200, headers: { 'content-type': 'application/json' } }); }; } + +/** + * Whether a presented refresh secret equals the configured one, in time that does not depend on where + * they first differ. + * + * Both are digested and the fixed-length digests compared in full, so neither the position of the first + * differing character nor the secret's length shapes the time taken. Without Web Crypto the strings are + * compared in full over the longer length instead. + */ +async function sameSecret(provided, secret) { + if (typeof provided !== 'string' || typeof secret !== 'string') return false; + const subtle = globalThis.crypto?.subtle; + if (subtle && typeof TextEncoder === 'function') { + try { + const encoder = new TextEncoder(); + const [a, b] = await Promise.all([ + subtle.digest('SHA-256', encoder.encode(provided)), + subtle.digest('SHA-256', encoder.encode(secret)), + ]); + + return equalBytes(new Uint8Array(a), new Uint8Array(b)); + } catch { + // Fall through to the full-length comparison. + } + } + let diff = provided.length ^ secret.length; + const length = Math.max(provided.length, secret.length); + for (let i = 0; i < length; i++) { + diff |= (provided.charCodeAt(i) || 0) ^ (secret.charCodeAt(i) || 0); + } + + return diff === 0; +} + +function equalBytes(a, b) { + if (a.length !== b.length) return false; + let diff = 0; + for (let i = 0; i < a.length; i++) diff |= a[i] ^ b[i]; + + return diff === 0; +} diff --git a/src/protect/runtime.js b/src/protect/runtime.js index 51575f02..3588e779 100644 --- a/src/protect/runtime.js +++ b/src/protect/runtime.js @@ -1725,6 +1725,11 @@ export async function createProtection(options = {}) { get: () => (detections ? () => detections.health() : undefined), enumerable: true, }); + // The same for the block log, when there is one: accepted, delivered, failed, dropped, and still queued. + Object.defineProperty(protection, 'blockLogHealth', { + get: () => (firewallLog ? () => firewallLog.health() : undefined), + enumerable: true, + }); return protection; } diff --git a/tests/protect/reporting-accounting.test.ts b/tests/protect/reporting-accounting.test.ts new file mode 100644 index 00000000..88aecb33 --- /dev/null +++ b/tests/protect/reporting-accounting.test.ts @@ -0,0 +1,183 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { notify, resetNotifyWarnings } from '../../src/protect/notify.js'; +import { makeRefreshHandler } from '../../src/protect/rules/refresh.js'; +import { createFirewallLogReporter } from '../../src/protect/firewall-log.js'; +import { createProtection } from '../../src/protect/runtime.js'; + +const API_KEY = 'samplesamplesamplesamplesamplesamplesamp-7'; + +describe('callback failure reporting', () => { + beforeEach(() => resetNotifyWarnings()); + afterEach(() => vi.restoreAllMocks()); + + const broken = () => () => { + throw new Error('sample failure'); + }; + + it('reports a second callback under the same hook name', () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const first = broken(); + const second = broken(); + notify(first, {}, 'onDetect'); + notify(first, {}, 'onDetect'); + notify(second, {}, 'onDetect'); + notify(second, {}, 'onDetect'); + expect(warn).toHaveBeenCalledTimes(2); + }); + + it('reports one callback once per hook name it is passed as', () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const shared = broken(); + notify(shared, {}, 'onDetect'); + notify(shared, {}, 'onError'); + notify(shared, {}, 'onError'); + expect(warn).toHaveBeenCalledTimes(2); + }); + + it('reports a failing async callback once per callback', async () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const first = async () => { throw new Error('sample failure'); }; + const second = async () => { throw new Error('sample failure'); }; + notify(first, {}, 'onSkip'); + notify(first, {}, 'onSkip'); + notify(second, {}, 'onSkip'); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(warn).toHaveBeenCalledTimes(2); + }); + + it('caps the number of warnings a process writes', () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + for (let i = 0; i < 40; i++) notify(broken(), {}, 'onDetect'); + expect(warn).toHaveBeenCalledTimes(21); + expect(String(warn.mock.calls[20][0])).toContain('not reported'); + }); + + it('reports each guard whose hook fails', async () => { + const warn = vi.spyOn(console, 'warn').mockImplementation(() => {}); + const rules = { + firewall: [{ id: 1, title: 'sample', rule_v2: [{ parameter: 'get.q', match: { type: 'contains', value: 'SAMPLE' } }] }], + whitelists: [], + whitelist_keys: {}, + }; + for (let i = 0; i < 2; i++) { + const protection: any = await createProtection({ rules, mode: 'dry-run', onDetect: broken() }); + await protection.fetchGuard()(new Request('https://app.example.test/?q=SAMPLE')); + await protection.stop(); + } + expect(warn.mock.calls.filter(([m]) => String(m).includes('onDetect'))).toHaveLength(2); + }); +}); + +describe('refresh secret comparison', () => { + afterEach(() => { + vi.restoreAllMocks(); + vi.unstubAllGlobals(); + }); + + const call = (secret: string, provided?: string) => { + const tick = vi.fn(async () => ({ ok: true })); + const headers = provided === undefined ? {} : { 'x-patchstack-refresh': provided }; + return makeRefreshHandler(tick, secret)(new Request('https://app.example.test/refresh', { method: 'POST', headers })) + .then((response) => ({ status: response.status, ticks: tick.mock.calls.length })); + }; + + it.each([ + ['the configured secret', 'sample-secret-value', 200, 1], + ['a different secret of the same length', 'sample-secret-valuf', 403, 0], + ['a prefix', 'sample-secret', 403, 0], + ['a longer value', 'sample-secret-value-2', 403, 0], + ['an empty value', '', 403, 0], + ['no header', undefined, 403, 0], + ])('answers %s', async (_label, provided, status, ticks) => { + expect(await call('sample-secret-value', provided)).toEqual({ status, ticks }); + }); + + it('compares fixed-length digests of both values', async () => { + const digest = vi.spyOn(globalThis.crypto.subtle, 'digest'); + expect((await call('sample-secret-value', 'sample-other')).status).toBe(403); + const inputs = digest.mock.calls.map(([, data]) => new TextDecoder().decode(data as Uint8Array)); + expect(inputs.sort()).toEqual(['sample-other', 'sample-secret-value']); + }); + + it('still decides correctly without Web Crypto', async () => { + vi.stubGlobal('crypto', undefined); + expect((await call('sample-secret-value', 'sample-secret-value')).status).toBe(200); + expect((await call('sample-secret-value', 'sample-secret-valuf')).status).toBe(403); + expect((await call('sample-secret-value', 'sample-secret')).status).toBe(403); + // Characters past the shorter value's end are not read as NUL. + expect((await call('sample\u0000', 'sample')).status).toBe(403); + }); +}); + +describe('block-log accounting', () => { + beforeEach(() => vi.useFakeTimers()); + afterEach(() => vi.useRealTimers()); + + const tokenResponse = () => new Response(JSON.stringify({ access_token: 'sample-token', expires_in: 3600 }), { status: 200 }); + + const reporterWith = (logStatus: number | 'hang', tokenOk = true) => { + const fetchImpl = vi.fn(async (url: string, init?: RequestInit) => { + if (String(url).includes('/oauth/token')) return tokenOk ? tokenResponse() : new Response('', { status: 401 }); + if (logStatus === 'hang') { + return new Promise((_resolve, reject) => { + init?.signal?.addEventListener('abort', () => reject(new Error('aborted'))); + }); + } + return new Response('{}', { status: logStatus }); + }); + return createFirewallLogReporter({ apiKey: API_KEY, apiBase: 'https://api.example.test', fetchImpl, flushMs: 10 }); + }; + + const record = (reporter: any, n: number) => { + for (let i = 0; i < n; i++) reporter.record({ rule: { id: 1 }, method: 'GET', path: '/' }); + }; + + it('counts delivered records', async () => { + const reporter: any = reporterWith(200); + record(reporter, 3); + await vi.advanceTimersByTimeAsync(50); + expect(reporter.health()).toEqual({ recorded: 3, delivered: 3, failed: 0, dropped: 0, queued: 0 }); + }); + + it('counts records in a refused batch as failed', async () => { + const reporter: any = reporterWith(500); + record(reporter, 2); + await vi.advanceTimersByTimeAsync(50); + expect(reporter.health()).toMatchObject({ recorded: 2, delivered: 0, failed: 2, dropped: 0 }); + }); + + it('counts records as failed when no token can be obtained', async () => { + const reporter: any = reporterWith(200, false); + record(reporter, 2); + await vi.advanceTimersByTimeAsync(50); + expect(reporter.health()).toMatchObject({ recorded: 2, delivered: 0, failed: 2, dropped: 0 }); + }); + + it('counts records turned away by a full queue', async () => { + const reporter: any = reporterWith('hang'); + record(reporter, 600); + const health = reporter.health(); + expect(health.recorded + health.dropped).toBe(600); + expect(health.dropped).toBeGreaterThan(0); + expect(health.recorded).toBe(health.queued + 50); + }); + + it('counts what a shutdown discards as dropped', async () => { + const reporter: any = reporterWith('hang'); + record(reporter, 120); + const stopped = reporter.stop(); + await vi.advanceTimersByTimeAsync(10_000); + await stopped; + expect(reporter.health()).toEqual({ recorded: 120, delivered: 0, failed: 0, dropped: 120, queued: 0 }); + }); + + it('is exposed on the protection object only when the block log is on', async () => { + const rules = { firewall: [], whitelists: [], whitelist_keys: {} }; + const withLog: any = await createProtection({ rules, apiKey: API_KEY, fetchImpl: async () => new Response('{}') }); + expect(withLog.blockLogHealth()).toEqual({ recorded: 0, delivered: 0, failed: 0, dropped: 0, queued: 0 }); + await withLog.stop(); + const without: any = await createProtection({ rules, reportFirewallLog: false, apiKey: API_KEY }); + expect(without.blockLogHealth).toBeUndefined(); + await without.stop(); + }); +}); diff --git a/tests/protect/rules-refresh.test.ts b/tests/protect/rules-refresh.test.ts index ee9261f7..a2bd0f3c 100644 --- a/tests/protect/rules-refresh.test.ts +++ b/tests/protect/rules-refresh.test.ts @@ -86,7 +86,9 @@ describe('push refresh endpoint (refreshHandler)', () => { const first = handler(request()); const second = handler(request()); - await Promise.resolve(); + // Both requests pass the secret check before the first refresh settles; they share one refresh. + await vi.waitFor(() => expect(tick).toHaveBeenCalled()); + await new Promise((resolve) => setTimeout(resolve, 20)); expect(tick).toHaveBeenCalledOnce(); release({ ok: true }); expect((await first).status).toBe(200); From 29413e0feb5157b8e9e93f85cf4a64ae4a546385 Mon Sep 17 00:00:00 2001 From: Dave Jong Date: Mon, 28 Sep 2026 16:40:07 +0200 Subject: [PATCH 2/2] Count the block-log batch in flight when a shutdown gives up A transport can ignore its abort signal and never settle. The batch it held is now counted as dropped when the shutdown budget runs out, and an answer that arrives afterwards does not count it again, so accepted records always equal delivered + failed + dropped + queued. Co-Authored-By: Claude Opus 5.5 --- src/protect/firewall-log.js | 33 +++++++---- tests/protect/reporting-accounting.test.ts | 65 ++++++++++++++++++++++ 2 files changed, 86 insertions(+), 12 deletions(-) diff --git a/src/protect/firewall-log.js b/src/protect/firewall-log.js index 64cd5821..fc28fd55 100644 --- a/src/protect/firewall-log.js +++ b/src/protect/firewall-log.js @@ -115,6 +115,18 @@ export function createFirewallLogReporter(opts) { let delivered = 0; let failed = 0; let dropped = 0; + /** The batch being sent, until its outcome is counted. A shutdown that gives up counts it as dropped. */ + /** @type {{ size: number, counted: boolean } | null} */ + let activeBatch = null; + /** Count a batch's outcome once: whichever of the send and a shutdown decides first. */ + const settle = (sent, outcome) => { + if (sent.counted) return; + sent.counted = true; + if (outcome === 'delivered') delivered += sent.size; + else if (outcome === 'failed') failed += sent.size; + else dropped += sent.size; + if (activeBatch === sent) activeBatch = null; + }; /** @type {{ token: string, expiresAt: number } | null} */ let cachedToken = null; @@ -182,15 +194,15 @@ export function createFirewallLogReporter(opts) { // whether it succeeded. /** @type {Promise} */ let entry; - let settled = false; + const sent = { size: batch.length, counted: false }; + activeBatch = sent; entry = (async () => { try { const token = await fetchAccessToken(controller?.signal); // Not after a shutdown gave up: it has already reported itself finished. if (ended) return; if (!token) { - failed += batch.length; - settled = true; + settle(sent, 'failed'); return; } @@ -213,18 +225,13 @@ export function createFirewallLogReporter(opts) { ...(controller ? { signal: controller.signal } : {}), }); const res = p && typeof p.then === 'function' ? await p.catch(() => null) : p; - if (ended) return; - if (res && res.ok) delivered += batch.length; - else failed += batch.length; - settled = true; + settle(sent, res && res.ok ? 'delivered' : 'failed'); } catch { /* A delivery problem is never worth disturbing the app over. */ } finally { - // Anything not accounted for above was abandoned by a shutdown, or failed before a verdict. - if (!settled) { - if (ended) dropped += batch.length; - else failed += batch.length; - } + // Anything not counted above failed before a verdict. A shutdown that gave up has already counted + // this batch as dropped, so an answer arriving later changes nothing. + settle(sent, 'failed'); clearTimeout(attemptTimer); if (activeController === controller) activeController = null; if (inFlight === entry) inFlight = null; @@ -314,6 +321,8 @@ export function createFirewallLogReporter(opts) { ended = true; activeController?.abort(); activeController = null; + // The batch in flight may never settle: a transport can ignore its abort signal. + if (activeBatch) settle(activeBatch, 'dropped'); dropped += queue.length; queue = []; inFlight = null; diff --git a/tests/protect/reporting-accounting.test.ts b/tests/protect/reporting-accounting.test.ts index 88aecb33..4e9918c8 100644 --- a/tests/protect/reporting-accounting.test.ts +++ b/tests/protect/reporting-accounting.test.ts @@ -153,6 +153,17 @@ describe('block-log accounting', () => { expect(reporter.health()).toMatchObject({ recorded: 2, delivered: 0, failed: 2, dropped: 0 }); }); + it('counts records as failed when the transport throws', async () => { + const fetchImpl = vi.fn((url: string) => { + if (String(url).includes('/oauth/token')) return Promise.resolve(tokenResponse()); + throw new Error('sample transport failure'); + }); + const reporter: any = createFirewallLogReporter({ apiKey: API_KEY, apiBase: 'https://api.example.test', fetchImpl, flushMs: 10 }); + record(reporter, 2); + await vi.advanceTimersByTimeAsync(50); + expect(reporter.health()).toMatchObject({ recorded: 2, delivered: 0, failed: 2, dropped: 0, queued: 0 }); + }); + it('counts records turned away by a full queue', async () => { const reporter: any = reporterWith('hang'); record(reporter, 600); @@ -171,6 +182,60 @@ describe('block-log accounting', () => { expect(reporter.health()).toEqual({ recorded: 120, delivered: 0, failed: 0, dropped: 120, queued: 0 }); }); + describe('with a transport that ignores cancellation', () => { + const balanced = (h: any) => h.recorded === h.delivered + h.failed + h.dropped + h.queued; + + // The post never settles on its own and does not listen to the abort signal; `release` settles it later. + const ignoringReporter = () => { + let release: (response: Response) => void = () => {}; + const fetchImpl = vi.fn(async (url: string) => { + if (String(url).includes('/oauth/token')) return tokenResponse(); + return new Promise((resolve) => { release = resolve; }); + }); + const reporter: any = createFirewallLogReporter({ apiKey: API_KEY, apiBase: 'https://api.example.test', fetchImpl, flushMs: 10 }); + return { reporter, release: (response: Response) => release(response) }; + }; + + it('accounts for the batch in flight when the shutdown budget runs out', async () => { + const { reporter } = ignoringReporter(); + record(reporter, 80); + await vi.advanceTimersByTimeAsync(0); + const stopped = reporter.stop(); + await vi.advanceTimersByTimeAsync(10_000); + await stopped; + const health = reporter.health(); + expect(health).toEqual({ recorded: 80, delivered: 0, failed: 0, dropped: 80, queued: 0 }); + expect(balanced(health)).toBe(true); + }); + + it('does not count the batch again when the transport answers after shutdown', async () => { + const { reporter, release } = ignoringReporter(); + record(reporter, 80); + await vi.advanceTimersByTimeAsync(0); + const stopped = reporter.stop(); + await vi.advanceTimersByTimeAsync(10_000); + await stopped; + release(new Response('{}', { status: 200 })); + await vi.advanceTimersByTimeAsync(20_000); + const health = reporter.health(); + expect(health).toEqual({ recorded: 80, delivered: 0, failed: 0, dropped: 80, queued: 0 }); + expect(balanced(health)).toBe(true); + }); + + it('keeps the counts balanced when the answer arrives before the budget runs out', async () => { + const { reporter, release } = ignoringReporter(); + record(reporter, 80); + await vi.advanceTimersByTimeAsync(0); + const stopped = reporter.stop(); + release(new Response('{}', { status: 200 })); + await vi.advanceTimersByTimeAsync(10_000); + await stopped; + const health = reporter.health(); + expect(health.delivered).toBeGreaterThanOrEqual(50); + expect(balanced(health)).toBe(true); + }); + }); + it('is exposed on the protection object only when the block log is on', async () => { const rules = { firewall: [], whitelists: [], whitelist_keys: {} }; const withLog: any = await createProtection({ rules, apiKey: API_KEY, fetchImpl: async () => new Response('{}') });