From 04ff0b631fd90fe31b7679483b0e9f27af0c7624 Mon Sep 17 00:00:00 2001 From: Liam Steiner Date: Tue, 14 Jul 2026 12:13:15 +0300 Subject: [PATCH 1/2] refactor(rate-limit): consolidate 3x-duplicated sliding-window limiter MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit credential-proxy.ts, odysseus-server.ts, and ingress/gateway.ts each independently implemented the same Map sliding-window rate limiter — the first two even commented "mirrors credential-proxy.ts" but were never consolidated. Extract the shared logic into src/rate-limiter.ts's createRateLimiter(), with an opt-in cleanupInterval (default false) so gateway.ts's no-interval behavior and credential-proxy.ts/odysseus-server.ts's periodic-prune behavior are both preserved exactly. All 3 call sites keep their existing thresholds, window semantics, and test-reset export names/signatures (_resetRateLimiterForTest, _resetServerStateForTest). The LIA-363 close-handler ordering (register only after a successful bind, so an EADDRINUSE retry's server.close() can't kill the timer early) is preserved verbatim in credential-proxy.ts. Slice 1 of an opportunistic dedup cleanup; not part of the tracked V2 migration roadmap. Design finalized after 8 rounds of adversarial plan review (native + GPT-5.6-Sol) before implementation. Co-Authored-By: Claude Sonnet 5 --- src/credential-proxy.ts | 45 +++++------------ src/ingress/gateway.ts | 22 ++------- src/odysseus-server.ts | 35 +++++--------- src/rate-limiter.test.ts | 102 +++++++++++++++++++++++++++++++++++++++ src/rate-limiter.ts | 60 +++++++++++++++++++++++ 5 files changed, 189 insertions(+), 75 deletions(-) create mode 100644 src/rate-limiter.test.ts create mode 100644 src/rate-limiter.ts diff --git a/src/credential-proxy.ts b/src/credential-proxy.ts index 2b543334..cf70ee6f 100644 --- a/src/credential-proxy.ts +++ b/src/credential-proxy.ts @@ -25,6 +25,7 @@ import { envPositiveInt } from './env-utils.js'; import { validateGroupToken } from './group-tokens.js'; import { readEnvFile } from './env.js'; import { logger } from './logger.js'; +import { createRateLimiter } from './rate-limiter.js'; import { AuthProviderRegistry, AnthropicAuthProvider, @@ -105,45 +106,21 @@ function bridgeRecallBoundArgs(): string[] { const RATE_LIMIT_MAX = 20; const RATE_LIMIT_WINDOW_MS = 60_000; -interface RateBucket { - timestamps: number[]; -} - -const rateBuckets = new Map(); - -/** Prune expired entries periodically to prevent unbounded growth. */ -const rateLimitCleanupInterval = setInterval(() => { - const now = Date.now(); - for (const [key, bucket] of rateBuckets) { - bucket.timestamps = bucket.timestamps.filter( - (t) => now - t < RATE_LIMIT_WINDOW_MS, - ); - if (bucket.timestamps.length === 0) rateBuckets.delete(key); - } -}, RATE_LIMIT_WINDOW_MS); - -// Prevent the cleanup timer from keeping Node alive after tests/shutdown -rateLimitCleanupInterval.unref(); +// Cleanup interval prunes expired entries periodically to prevent unbounded +// growth; registered at module scope (not inside startCredentialProxy/ +// tryListen) so it is created exactly once at import, never recreated on an +// EADDRINUSE retry (LIA-363). +const limiter = createRateLimiter(RATE_LIMIT_MAX, RATE_LIMIT_WINDOW_MS, { + cleanupInterval: true, +}); /** @internal exposed for testing only */ export function _resetRateLimiterForTest(): void { - rateBuckets.clear(); + limiter.resetForTest(); } function isRateLimited(sourceKey: string): boolean { - const now = Date.now(); - let bucket = rateBuckets.get(sourceKey); - if (!bucket) { - bucket = { timestamps: [] }; - rateBuckets.set(sourceKey, bucket); - } - // Prune expired timestamps for this source - bucket.timestamps = bucket.timestamps.filter( - (t) => now - t < RATE_LIMIT_WINDOW_MS, - ); - if (bucket.timestamps.length >= RATE_LIMIT_MAX) return true; - bucket.timestamps.push(now); - return false; + return limiter.isRateLimited(sourceKey); } /** @@ -497,7 +474,7 @@ export function startCredentialProxy( // binds on a later attempt with proactive OAuth refresh permanently // dead — the exact "next-morning 401" class the timer prevents (LIA-363). server.on('close', () => { - clearInterval(rateLimitCleanupInterval); + limiter.dispose(); if (proactiveRefreshTimer) clearInterval(proactiveRefreshTimer); }); logger.info( diff --git a/src/ingress/gateway.ts b/src/ingress/gateway.ts index 7eeada9f..45259397 100644 --- a/src/ingress/gateway.ts +++ b/src/ingress/gateway.ts @@ -14,6 +14,7 @@ import { type ServerResponse, } from 'http'; import { logger } from '../logger.js'; +import { createRateLimiter } from '../rate-limiter.js'; import type { VerifyResult } from './hmac.js'; /** @@ -78,26 +79,13 @@ function forwardedFor(req: IncomingMessage): string | undefined { return undefined; } -/** Sliding-window per-key rate limiter (same shape as odysseus-server.ts:83-95). */ -function makeRateLimiter(max: number, windowMs: number) { - const buckets = new Map(); - return function isRateLimited(key: string, now: number): boolean { - const ts = (buckets.get(key) ?? []).filter((t) => now - t < windowMs); - if (ts.length >= max) { - buckets.set(key, ts); - return true; - } - ts.push(now); - buckets.set(key, ts); - return false; - }; -} - /** Build (but do not start) the gateway server. Exposed for tests. */ export function createIngressGateway(deps: GatewayDeps): Server { const { handlers, config } = deps; const allow = new Set(config.ipAllowlist); - const isRateLimited = makeRateLimiter( + // No cleanupInterval — gateway.ts prunes inline during isRateLimited only, + // matching its behavior before the shared-limiter consolidation. + const limiter = createRateLimiter( config.rateLimitMax, config.rateLimitWindowMs, ); @@ -131,7 +119,7 @@ export function createIngressGateway(deps: GatewayDeps): Server { } // 2) Rate limit per peer IP. - if (isRateLimited(ip, started)) { + if (limiter.isRateLimited(ip, started)) { writeStatus(res, 429, 'rate limited'); return; } diff --git a/src/odysseus-server.ts b/src/odysseus-server.ts index f149c515..087d7338 100644 --- a/src/odysseus-server.ts +++ b/src/odysseus-server.ts @@ -43,6 +43,7 @@ import { GroupQueue } from './group-queue.js'; import { scanForInjection } from './guardrails/injection-scanner.js'; import { logger } from './logger.js'; import { messageText } from './openai-messages.js'; +import { createRateLimiter } from './rate-limiter.js'; import { getAvailableGroups } from './router-state.js'; import { RegisteredGroup } from './types.js'; import { consolidateWebConversation } from './webui-consolidation.js'; @@ -62,41 +63,27 @@ const MAX_HISTORY_CHARS = (() => { return Number.isFinite(n) && n > 0 ? n : 24000; })(); -/* ── Rate limiter (mirrors credential-proxy.ts:69-113) ─────────────────── */ +/* ── Rate limiter (shared implementation — src/rate-limiter.ts) ────────── */ const RATE_LIMIT_MAX = 5; const RATE_LIMIT_WINDOW_MS = 60_000; -const rateBuckets = new Map(); -const rateLimitCleanupInterval = setInterval(() => { - const now = Date.now(); - for (const [key, ts] of rateBuckets) { - const kept = ts.filter((t) => now - t < RATE_LIMIT_WINDOW_MS); - if (kept.length === 0) rateBuckets.delete(key); - else rateBuckets.set(key, kept); - } -}, RATE_LIMIT_WINDOW_MS); -rateLimitCleanupInterval.unref(); +// Registered at module scope (not inside startOdysseusServer) so the cleanup +// interval is created exactly once at import (LIA-363 lesson: avoid a timer +// created/leaked on every retry of a listen-success path). +const limiter = createRateLimiter(RATE_LIMIT_MAX, RATE_LIMIT_WINDOW_MS, { + cleanupInterval: true, +}); let activeSse = 0; /** main jid → true while a turn is queued/running (one in-flight per jid). */ const inFlight = new Set(); function isRateLimited(key: string): boolean { - const now = Date.now(); - const ts = (rateBuckets.get(key) ?? []).filter( - (t) => now - t < RATE_LIMIT_WINDOW_MS, - ); - if (ts.length >= RATE_LIMIT_MAX) { - rateBuckets.set(key, ts); - return true; - } - ts.push(now); - rateBuckets.set(key, ts); - return false; + return limiter.isRateLimited(key); } /** @internal exposed for testing only */ export function _resetServerStateForTest(): void { - rateBuckets.clear(); + limiter.resetForTest(); inFlight.clear(); activeSse = 0; } @@ -396,7 +383,7 @@ export function startOdysseusServer( return new Promise((resolve, reject) => { const server = createOdysseusServer(deps, token); - server.on('close', () => clearInterval(rateLimitCleanupInterval)); + server.on('close', () => limiter.dispose()); server.on('error', (err: NodeJS.ErrnoException) => reject(err)); server.listen(ODYSSEUS_HTTP_PORT, ODYSSEUS_BIND_HOST, () => { logger.info( diff --git a/src/rate-limiter.test.ts b/src/rate-limiter.test.ts new file mode 100644 index 00000000..64268064 --- /dev/null +++ b/src/rate-limiter.test.ts @@ -0,0 +1,102 @@ +import { describe, it, expect, vi, afterEach } from 'vitest'; +import { createRateLimiter } from './rate-limiter.js'; + +describe('createRateLimiter', () => { + afterEach(() => { + vi.useRealTimers(); + }); + + it('allows requests within the window up to max, then blocks', () => { + const limiter = createRateLimiter(3, 60_000); + const now = 1_000_000; + expect(limiter.isRateLimited('a', now)).toBe(false); + expect(limiter.isRateLimited('a', now)).toBe(false); + expect(limiter.isRateLimited('a', now)).toBe(false); + // 4th request within the same window is blocked. + expect(limiter.isRateLimited('a', now)).toBe(true); + }); + + it('allows a request once earlier timestamps fall outside the window', () => { + const limiter = createRateLimiter(2, 60_000); + const t0 = 1_000_000; + expect(limiter.isRateLimited('a', t0)).toBe(false); + expect(limiter.isRateLimited('a', t0)).toBe(false); + expect(limiter.isRateLimited('a', t0)).toBe(true); + // Advance past the window — the two earlier timestamps are now expired. + const t1 = t0 + 60_001; + expect(limiter.isRateLimited('a', t1)).toBe(false); + }); + + it('tracks separate keys independently', () => { + const limiter = createRateLimiter(1, 60_000); + const now = 1_000_000; + expect(limiter.isRateLimited('a', now)).toBe(false); + expect(limiter.isRateLimited('a', now)).toBe(true); + // A different key has its own bucket. + expect(limiter.isRateLimited('b', now)).toBe(false); + }); + + it('defaults `now` to Date.now() when omitted', () => { + vi.useFakeTimers(); + vi.setSystemTime(1_000_000); + const limiter = createRateLimiter(1, 60_000); + expect(limiter.isRateLimited('a')).toBe(false); + expect(limiter.isRateLimited('a')).toBe(true); + }); + + it('cleanupInterval: true creates an interval that dispose() clears', () => { + vi.useFakeTimers(); + const setIntervalSpy = vi.spyOn(global, 'setInterval'); + const clearIntervalSpy = vi.spyOn(global, 'clearInterval'); + + const limiter = createRateLimiter(5, 60_000, { cleanupInterval: true }); + expect(setIntervalSpy).toHaveBeenCalledTimes(1); + expect(setIntervalSpy).toHaveBeenCalledWith(expect.any(Function), 60_000); + + limiter.dispose(); + expect(clearIntervalSpy).toHaveBeenCalledTimes(1); + }); + + it('cleanupInterval: false (default) creates no interval', () => { + const setIntervalSpy = vi.spyOn(global, 'setInterval'); + const limiter = createRateLimiter(5, 60_000); + expect(setIntervalSpy).not.toHaveBeenCalled(); + // dispose() is still safe to call (no-op) even though no interval exists. + expect(() => limiter.dispose()).not.toThrow(); + }); + + it('the cleanup interval prunes expired entries across buckets', () => { + vi.useFakeTimers(); + vi.setSystemTime(1_000_000); + const limiter = createRateLimiter(1, 60_000, { cleanupInterval: true }); + expect(limiter.isRateLimited('a', Date.now())).toBe(false); + expect(limiter.isRateLimited('a', Date.now())).toBe(true); + + // Advance past the window and let the prune interval fire. + vi.setSystemTime(1_000_000 + 60_001); + vi.advanceTimersByTime(60_000); + + // The bucket was pruned, so a fresh request at the current time succeeds. + expect(limiter.isRateLimited('a', Date.now())).toBe(false); + limiter.dispose(); + }); + + it('resetForTest() clears all bucket state', () => { + const limiter = createRateLimiter(1, 60_000); + const now = 1_000_000; + expect(limiter.isRateLimited('a', now)).toBe(false); + expect(limiter.isRateLimited('a', now)).toBe(true); + limiter.resetForTest(); + expect(limiter.isRateLimited('a', now)).toBe(false); + }); + + it('dispose() is idempotent — safe to call multiple times', () => { + vi.useFakeTimers(); + const limiter = createRateLimiter(5, 60_000, { cleanupInterval: true }); + expect(() => { + limiter.dispose(); + limiter.dispose(); + limiter.dispose(); + }).not.toThrow(); + }); +}); diff --git a/src/rate-limiter.ts b/src/rate-limiter.ts new file mode 100644 index 00000000..10f35ba8 --- /dev/null +++ b/src/rate-limiter.ts @@ -0,0 +1,60 @@ +/** + * Shared sliding-window rate limiter (dedup of 3 independent implementations: + * credential-proxy.ts, odysseus-server.ts, ingress/gateway.ts — the first two + * commented "mirrors credential-proxy.ts" but were never consolidated). + * + * `Map` timestamp-array shape, matching all 3 prior + * implementations. `opts.cleanupInterval` is opt-in (defaults false) — callers + * that already run a periodic prune (credential-proxy.ts, odysseus-server.ts) + * ask for one; gateway.ts relies on inline pruning during `isRateLimited` only, + * matching its current no-interval behavior exactly. + */ +export function createRateLimiter( + max: number, + windowMs: number, + opts?: { cleanupInterval?: boolean }, +): { + isRateLimited(key: string, now?: number): boolean; + dispose(): void; + resetForTest(): void; +} { + const buckets = new Map(); + + let cleanupTimer: ReturnType | undefined; + if (opts?.cleanupInterval) { + cleanupTimer = setInterval(() => { + const now = Date.now(); + for (const [key, ts] of buckets) { + const kept = ts.filter((t) => now - t < windowMs); + if (kept.length === 0) buckets.delete(key); + else buckets.set(key, kept); + } + }, windowMs); + // Prevent the cleanup timer from keeping Node alive after tests/shutdown. + cleanupTimer.unref(); + } + + function isRateLimited(key: string, now: number = Date.now()): boolean { + const ts = (buckets.get(key) ?? []).filter((t) => now - t < windowMs); + if (ts.length >= max) { + buckets.set(key, ts); + return true; + } + ts.push(now); + buckets.set(key, ts); + return false; + } + + function dispose(): void { + if (cleanupTimer) { + clearInterval(cleanupTimer); + cleanupTimer = undefined; + } + } + + function resetForTest(): void { + buckets.clear(); + } + + return { isRateLimited, dispose, resetForTest }; +} From 09bca40f7e9546f084d6cb1adb98e80509b553f0 Mon Sep 17 00:00:00 2001 From: Liam Steiner Date: Fri, 17 Jul 2026 19:14:55 +0300 Subject: [PATCH 2/2] chore(patterns): bump drift check --- patterns/general-code.md | 2 +- patterns/security-review.md | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/patterns/general-code.md b/patterns/general-code.md index 55016f1a..96156945 100644 --- a/patterns/general-code.md +++ b/patterns/general-code.md @@ -1,5 +1,5 @@ --- -last_verified: "2026-07-17" # auto-bump @1784303093 +last_verified: "2026-07-17" # auto-bump @1784304895 governs: - src/ - evolution/ diff --git a/patterns/security-review.md b/patterns/security-review.md index 390e85a5..691a4422 100644 --- a/patterns/security-review.md +++ b/patterns/security-review.md @@ -5,7 +5,7 @@ governs: - src/ipc.ts - src/sender-allowlist.ts - src/mount-security.ts -last_verified: "2026-07-13" # auto-bump @1783974306 +last_verified: "2026-07-17" # auto-bump @1784304895 test_tasks: - "Add a new mount in src/container-mounter.ts for per-group config files" - "Update src/mount-security.ts to permit reading a new credential path"