Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
65 changes: 64 additions & 1 deletion packages/webdecoy/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,10 +34,36 @@ interface ApiResponse<T> {
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<string, string>;
/** No request is made before this time (epoch ms). */
private pausedUntil = 0;
private consecutiveFailures = 0;

constructor(config: ClientConfig) {
this.config = config;
Expand All @@ -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<T>(method: 'GET' | 'POST', path: string, body?: unknown): Promise<ApiResponse<T>> {
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;
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -167,7 +229,8 @@ export class WebDecoyClient {
async sendAIReferrals(batch: AIReferralBatch): Promise<boolean> {
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);
Expand Down
9 changes: 9 additions & 0 deletions packages/webdecoy/src/ip-enrichment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, CachedEntry>();
Expand Down Expand Up @@ -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;
Expand Down
176 changes: 176 additions & 0 deletions packages/webdecoy/src/outage.test.ts
Original file line number Diff line number Diff line change
@@ -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<Response>) {
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);
});
});
6 changes: 5 additions & 1 deletion packages/webdecoy/src/referrals/referral-counter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
8 changes: 6 additions & 2 deletions packages/webdecoy/src/sdk.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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),
});

Expand Down
13 changes: 12 additions & 1 deletion packages/webdecoy/src/violation-reporter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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[] = [];
Expand Down Expand Up @@ -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
Expand All @@ -54,7 +63,9 @@ export class ViolationReporter {
* Flush buffered violations to the backend
*/
async flush(): Promise<void> {
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
Expand Down
Loading