diff --git a/CHANGELOG.md b/CHANGELOG.md index a0a8923..89361f8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Fixed + +- **Requests no longer wait on an unavailable WebDecoy.** After a 429, a 5xx or no answer, the client pauses calls to WebDecoy (honouring `Retry-After`, otherwise 1s doubling to 60s) and `protect()` fails open at once with an `ERROR` decision instead of waiting out the timeout on every request. The failure that starts a pause is logged as an error; requests during the pause log at debug, so an outage no longer floods your logs. +- **Bounded memory during an outage.** Violation events are held (up to 1,000, oldest dropped) while WebDecoy is unavailable and sent when it returns; the IP enrichment cache is capped at 10,000 IPs; AI referral counting keeps existing pairs counting but stops adding new ones while an unsent batch waits. +- **AI referral counts refused with 429 are kept and retried** under the same batch id instead of being treated as delivered. + ## [0.18.3] - 2026-09-30 ### Fixed diff --git a/packages/webdecoy/src/client.ts b/packages/webdecoy/src/client.ts index 1ddd219..e27a666 100644 --- a/packages/webdecoy/src/client.ts +++ b/packages/webdecoy/src/client.ts @@ -34,10 +34,36 @@ interface ApiResponse { data: T | undefined; } +/** + * Thrown without touching the network while the client is paused after + * WebDecoy refused work or did not answer. Callers fail open exactly as they + * do for any other error; the type lets them stay quiet about it, because the + * failure that started the pause was already reported. + */ +export class WebDecoyUnavailableError extends Error { + constructor(readonly retryAt: number) { + super('WebDecoy is temporarily unavailable; requests are paused'); + this.name = 'WebDecoyUnavailableError'; + // Keeps `instanceof` working when compiled to ES5, where extending Error + // loses the subclass prototype. + Object.setPrototypeOf(this, new.target.prototype); + } +} + +/** First pause after a failure; doubles per consecutive failure. */ +const PAUSE_BASE_MS = 1_000; +/** Longest pause without a Retry-After. */ +const PAUSE_MAX_MS = 60_000; +/** Longest pause a Retry-After can ask for. */ +const RETRY_AFTER_MAX_MS = 300_000; + export class WebDecoyClient { private config: ClientConfig; private baseUrl: string; private headers: Record; + /** No request is made before this time (epoch ms). */ + private pausedUntil = 0; + private consecutiveFailures = 0; constructor(config: ClientConfig) { this.config = config; @@ -56,7 +82,35 @@ export class WebDecoyClient { } } + /** + * Whether a request would be attempted right now. False while paused after + * WebDecoy refused work (429, 5xx) or did not answer, so a caller with + * something to keep (a buffer) can hold it rather than lose it. + */ + isAvailable(now = Date.now()): boolean { + return now >= this.pausedUntil; + } + + /** + * Pause every call to WebDecoy after a refusal or no answer: the Retry-After + * when one was given, otherwise 1s doubling per consecutive failure up to + * 60s. Without this, each request during an outage waited out the full + * timeout in the caller's request path. + */ + private pause(retryAfter: string | null): void { + const seconds = retryAfter ? Number(retryAfter) : NaN; + const ms = + Number.isFinite(seconds) && seconds > 0 + ? Math.min(Math.max(seconds * 1000, PAUSE_BASE_MS), RETRY_AFTER_MAX_MS) + : Math.min(PAUSE_BASE_MS * 2 ** this.consecutiveFailures, PAUSE_MAX_MS); + this.consecutiveFailures++; + this.pausedUntil = Date.now() + ms; + } + private async request(method: 'GET' | 'POST', path: string, body?: unknown): Promise> { + if (!this.isAvailable()) { + throw new WebDecoyUnavailableError(this.pausedUntil); + } const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), this.config.timeout); // Node returns a Timeout object that would otherwise hold the process open; @@ -91,8 +145,16 @@ export class WebDecoyClient { console.log('[WebDecoy] Response:', { status: response.status, data }); } + if (response.status === 429 || response.status >= 500) { + this.pause(response.headers.get('retry-after')); + } else { + this.consecutiveFailures = 0; + } + return { status: response.status, data }; } catch (error) { + // No answer at all: a timeout or a connection failure. + this.pause(null); if (this.config.debug) { console.error('[WebDecoy] Error:', { message: error instanceof Error ? error.message : String(error), @@ -167,7 +229,8 @@ export class WebDecoyClient { async sendAIReferrals(batch: AIReferralBatch): Promise { try { const response = await this.request('POST', '/api/v1/sdk/ai-referrals', batch); - return response.status < 500; + // 429 is "not now", not "never": keep the batch for the next flush. + return response.status < 500 && response.status !== 429; } catch (error) { if (this.config.debug) { console.error('[WebDecoy] Failed to send AI referrals:', error); diff --git a/packages/webdecoy/src/ip-enrichment.ts b/packages/webdecoy/src/ip-enrichment.ts index b329552..484ead5 100644 --- a/packages/webdecoy/src/ip-enrichment.ts +++ b/packages/webdecoy/src/ip-enrichment.ts @@ -11,6 +11,9 @@ interface CachedEntry { expiresAt: number; } +/** Most IPs cached at once; the oldest entry is evicted past this. */ +export const MAX_ENRICHMENT_ENTRIES = 10_000; + export class IPEnrichmentClient { private client: WebDecoyClient; private cache = new Map(); @@ -54,10 +57,16 @@ export class IPEnrichmentClient { const data = await this.client.getIPEnrichment(ip); if (data) { + this.cache.delete(ip); this.cache.set(ip, { data, expiresAt: Date.now() + this.ttlMs, }); + if (this.cache.size > MAX_ENRICHMENT_ENTRIES) { + // Map iterates in insertion order: this is the oldest entry. + const oldest = this.cache.keys().next(); + if (!oldest.done) this.cache.delete(oldest.value); + } } return data; diff --git a/packages/webdecoy/src/outage.test.ts b/packages/webdecoy/src/outage.test.ts new file mode 100644 index 0000000..4657632 --- /dev/null +++ b/packages/webdecoy/src/outage.test.ts @@ -0,0 +1,176 @@ +/** + * Behavior while WebDecoy is refusing work or not answering. + * + * Every call goes through one client, which pauses after a 429, a 5xx or no + * answer (honouring Retry-After). While paused, calls fail at once without + * touching the network, so a request never waits out the timeout on a dead + * service; buffers stay bounded; and the operator's log gets the failure once, + * not once per request. + */ + +import { WebDecoy } from './sdk'; +import { WebDecoyClient, WebDecoyUnavailableError } from './client'; +import { ViolationReporter, MAX_BUFFERED_VIOLATIONS } from './violation-reporter'; +import { AIReferralCounter } from './referrals/referral-counter'; +import { IPEnrichmentClient, MAX_ENRICHMENT_ENTRIES } from './ip-enrichment'; +import type { RequestMetadata } from './types'; + +const realFetch = global.fetch; +let calls = 0; +function serve(answer: () => Response | Promise) { + calls = 0; + global.fetch = jest.fn(async () => { + calls++; + return answer(); + }) as any; +} + +function client() { + return new WebDecoyClient({ apiKey: 'k', apiUrl: 'https://ingest.example', timeout: 5000, debug: false, tlsRejectUnauthorized: true }); +} + +const okDetection = () => + new Response( + JSON.stringify({ decision: 'allow', confidence: 0, threat_level: 'MINIMAL', bot_detected: false, detection_id: 'd', rule_enforced: false }), + { status: 200 }, + ); + +beforeEach(() => { + jest.useFakeTimers({ doNotFake: ['nextTick', 'setImmediate', 'queueMicrotask'] }); + jest.setSystemTime(new Date('2026-10-02T12:00:00Z')); +}); +afterEach(() => { + global.fetch = realFetch; + jest.useRealTimers(); +}); + +describe('client pause', () => { + it('pauses after a 503: the next call fails at once without a request', async () => { + const c = client(); + serve(() => new Response('', { status: 503 })); + await expect(c.getIPEnrichment('198.51.100.1')).resolves.toBeNull(); + expect(c.isAvailable()).toBe(false); + await expect(c.detect({ request_metadata: {} } as any)).rejects.toBeInstanceOf(WebDecoyUnavailableError); + expect(calls).toBe(1); + }); + + it('honours Retry-After on a 429', async () => { + const c = client(); + serve(() => new Response('', { status: 429, headers: { 'retry-after': '30' } })); + await c.getIPEnrichment('198.51.100.1'); + jest.setSystemTime(Date.now() + 29_000); + expect(c.isAvailable()).toBe(false); + jest.setSystemTime(Date.now() + 2_000); + expect(c.isAvailable()).toBe(true); + }); + + it('backs off exponentially when nothing answers, and a success resets it', async () => { + const c = client(); + serve(() => { + throw new TypeError('fetch failed'); + }); + await c.getIPEnrichment('a'); // pause 1s + jest.setSystemTime(Date.now() + 1_001); + await c.getIPEnrichment('b'); // pause 2s + jest.setSystemTime(Date.now() + 1_500); + expect(c.isAvailable()).toBe(false); + jest.setSystemTime(Date.now() + 600); + serve(okDetection); + await c.getIPEnrichment('c'); + expect(c.isAvailable()).toBe(true); + }); + + it('a 4xx answer is not an outage', async () => { + const c = client(); + serve(() => new Response('{"error":"bad"}', { status: 400 })); + await c.getIPEnrichment('a'); + expect(c.isAvailable()).toBe(true); + }); +}); + +describe('protect() during an outage', () => { + const bot: RequestMetadata = { + method: 'GET', + path: '/', + ip: '203.0.113.9', + user_agent: 'python-requests/2.31.0', + headers: { 'user-agent': 'python-requests/2.31.0' }, + timestamp: Date.now(), + }; + + it('fails open without waiting, and logs the outage once, not per request', async () => { + const errors: string[] = []; + const logger = { debug() {}, info() {}, warn() {}, error: (m: string) => errors.push(m) }; + const wd = new WebDecoy({ apiKey: 'sk_test_outage', apiUrl: 'https://ingest.example', logger } as any); + serve(() => new Response('', { status: 503 })); + const decisions = []; + for (let i = 0; i < 10; i++) decisions.push(await wd.protect({ ...bot, ip: `203.0.113.${i + 10}` })); + expect(decisions.every((d) => d.conclusion === 'ERROR')).toBe(true); + expect(calls).toBe(1); + expect(errors).toHaveLength(1); + }); +}); + +describe('bounded buffers', () => { + it('holds violations while paused, capped, and sends them once WebDecoy is back', async () => { + const c = client(); + serve(() => new Response('', { status: 503 })); + await c.getIPEnrichment('a'); // paused + const reporter = new ViolationReporter(c, { flushInterval: 3_600_000, maxBufferSize: 1_000_000 }); + for (let i = 0; i < MAX_BUFFERED_VIOLATIONS + 500; i++) reporter.report([{ rule: 'r', ip: String(i) } as any]); + await reporter.flush(); + expect(calls).toBe(1); // nothing sent while paused + jest.setSystemTime(Date.now() + 2_000); + const sent: number[] = []; + global.fetch = jest.fn(async (_u: any, init: any) => { + sent.push(JSON.parse(init.body).events.length); + return new Response('{}', { status: 202 }); + }) as any; + await reporter.destroy(); + expect(sent.reduce((a, b) => a + b, 0)).toBe(MAX_BUFFERED_VIOLATIONS); + }); + + it('keeps an AI referral batch refused with 429, and stops adding pairs past the cap', async () => { + const batches: any[] = []; + let ok = false; + const counter = new AIReferralCounter({ + async sendAIReferrals(b) { + batches.push(b); + return ok; + }, + }, { flushInterval: 3_600_000 }); + const visit = (path: string) => + counter.observe({ + method: 'GET', + path, + ip: '1.1.1.1', + headers: { 'sec-fetch-mode': 'navigate', 'sec-fetch-dest': 'document', referer: 'https://chatgpt.com/' }, + timestamp: Date.now(), + }); + visit('/a'); + await counter.flush(); // refused: kept as pending + for (let i = 0; i < 2_000; i++) visit(`/p${i}`); + await new Promise((r) => setImmediate(r)); + expect(batches).toHaveLength(1); // no send per new pair while one is pending + ok = true; + await counter.flush(); // the pending batch, same id + await counter.flush(); // what was counted meanwhile + expect(batches[1].report_id).toBe(batches[0].report_id); + expect(batches[2].referrals.length).toBeLessThanOrEqual(500); + await counter.destroy(); + }); + + it('a 429 on AI referrals is not treated as delivered', async () => { + const c = client(); + serve(() => new Response('', { status: 429 })); + await expect(c.sendAIReferrals({ report_id: 'x', source: 'sdk', referrals: [] })).resolves.toBe(false); + }); + + it('caps the IP enrichment cache', async () => { + const c = client(); + serve(() => new Response(JSON.stringify({ security: {} }), { status: 200 })); + const e = new IPEnrichmentClient(c); + for (let i = 0; i < MAX_ENRICHMENT_ENTRIES + 50; i++) await e.enrich(`10.0.${i >> 8}.${i & 255}`); + expect((e as any).cache.size).toBe(MAX_ENRICHMENT_ENTRIES); + }); +}); diff --git a/packages/webdecoy/src/referrals/referral-counter.ts b/packages/webdecoy/src/referrals/referral-counter.ts index 355b760..dcda36b 100644 --- a/packages/webdecoy/src/referrals/referral-counter.ts +++ b/packages/webdecoy/src/referrals/referral-counter.ts @@ -79,8 +79,12 @@ export class AIReferralCounter { if (entry) { entry.count++; } else { + // While an unsent batch is waiting, new pairs past the cap are not + // kept: memory stays bounded and each new pair does not trigger + // another send attempt. Existing pairs keep counting. + if (this.pending && this.counts.size >= MAX_ENTRIES) return; this.counts.set(key, { platform, path: landing, count: 1 }); - if (this.counts.size >= MAX_ENTRIES) void this.flush(); + if (this.counts.size >= MAX_ENTRIES && !this.pending) void this.flush(); } } catch { // Counting must never affect the request. diff --git a/packages/webdecoy/src/sdk.ts b/packages/webdecoy/src/sdk.ts index 09a2b00..05e829a 100644 --- a/packages/webdecoy/src/sdk.ts +++ b/packages/webdecoy/src/sdk.ts @@ -3,7 +3,7 @@ * Main SDK class for bot detection and protection */ -import { WebDecoyClient } from './client'; +import { WebDecoyClient, WebDecoyUnavailableError } from './client'; import { analyzeRequest } from './local-analysis'; import { RuleEngine } from './rules/rule-engine'; import { tripwire } from './rules'; @@ -579,7 +579,11 @@ export class WebDecoy { } catch (error) { // An error here means no verdict was reached, which the operator wants to // know about whether or not they turned debug on. - this.log.error('Protection error', { + // While the client is paused after an outage, every request lands here. + // The failure that started the pause was already logged as an error; + // logging each paused request too would flood the operator's logs. + const log = error instanceof WebDecoyUnavailableError ? this.log.debug.bind(this.log) : this.log.error.bind(this.log); + log('Protection error', { error: error instanceof Error ? error.message : String(error), }); diff --git a/packages/webdecoy/src/violation-reporter.ts b/packages/webdecoy/src/violation-reporter.ts index 20771c4..68a232e 100644 --- a/packages/webdecoy/src/violation-reporter.ts +++ b/packages/webdecoy/src/violation-reporter.ts @@ -15,6 +15,12 @@ export interface ViolationReporterConfig { debug?: boolean; } +/** + * Most events held while WebDecoy is unavailable. The oldest are dropped past + * this, so an outage can never grow the process's memory without bound. + */ +export const MAX_BUFFERED_VIOLATIONS = 1_000; + export class ViolationReporter { private client: WebDecoyClient; private buffer: ViolationEvent[] = []; @@ -43,6 +49,9 @@ export class ViolationReporter { */ report(violations: ViolationEvent[]): void { this.buffer.push(...violations); + if (this.buffer.length > MAX_BUFFERED_VIOLATIONS) { + this.buffer.splice(0, this.buffer.length - MAX_BUFFERED_VIOLATIONS); + } if (this.buffer.length >= this.maxBufferSize) { // Fire-and-forget flush @@ -54,7 +63,9 @@ export class ViolationReporter { * Flush buffered violations to the backend */ async flush(): Promise { - if (this.flushing || this.buffer.length === 0) return; + // While the client is paused, hold the buffer (capped above) rather than + // drain it into requests that cannot be made. + if (this.flushing || this.buffer.length === 0 || !this.client.isAvailable()) return; this.flushing = true; // Drain the buffer