From dc5d954a7165fb0cd5865e759b960b32fc044a7c Mon Sep 17 00:00:00 2001 From: dietrich ayala Date: Thu, 12 Feb 2026 23:49:27 +0100 Subject: [PATCH] Add rate limiting and abuse prevention across HTTP, gossipsub, and libp2p Defense-in-depth rate limiting with zero-config defaults: - Sliding window rate limiter (src/rate-limiter.ts) with per-pool isolation - HTTP middleware: per-route rate limits (meta/sync/session/read/write/challenge/admin) - Body size limits: 1MB JSON, 64KB challenge, 60MB blob, 100MB CAR - Gossipsub: 8KB message size cap, per-topic rate limiting (60/min commits, 10/min identity) - libp2p: stream size caps (64KB inbound challenges, 1MB responses) - libp2p: connection manager limits (100 max, 10 pending, 5/s inbound threshold) - WebSocket firehose: per-IP connection limits (default 3) - Challenge validation: targetDid check, path/CID count caps, expiration rejection - All configurable via env vars, disabled by default in tests --- src/config.ts | 18 +++ src/index.ts | 134 +++++++++++++++++ src/ipfs.test.ts | 8 + src/ipfs.ts | 40 +++++ src/middleware/body-limit.ts | 41 +++++ src/middleware/rate-limit.test.ts | 125 ++++++++++++++++ src/middleware/rate-limit.ts | 62 ++++++++ src/rate-limiter.test.ts | 136 +++++++++++++++++ src/rate-limiter.ts | 141 ++++++++++++++++++ .../challenge-response.test.ts | 8 + .../challenge-response/e2e-challenge.test.ts | 8 + .../challenge-response/libp2p-transport.ts | 27 ++-- src/replication/challenge-response/types.ts | 11 ++ src/replication/e2e-multi-node.test.ts | 8 + src/replication/firehose-incremental.test.ts | 8 + .../gossipsub-notifications.test.ts | 24 +++ src/replication/mst-proof.test.ts | 8 + src/replication/offer-manager.test.ts | 8 + src/replication/peer-freshness.test.ts | 8 + src/replication/policy-integration.test.ts | 8 + src/replication/replication.test.ts | 8 + src/server.ts | 41 ++++- src/xrpc/admin-e2e.test.ts | 8 + src/xrpc/admin.test.ts | 8 + 24 files changed, 886 insertions(+), 10 deletions(-) create mode 100644 src/middleware/body-limit.ts create mode 100644 src/middleware/rate-limit.test.ts create mode 100644 src/middleware/rate-limit.ts create mode 100644 src/rate-limiter.test.ts create mode 100644 src/rate-limiter.ts diff --git a/src/config.ts b/src/config.ts index d55655c..7e19c93 100644 --- a/src/config.ts +++ b/src/config.ts @@ -30,6 +30,16 @@ export interface Config { NODE_MANAGERS: string[]; /** Optional human-readable node name. */ NODE_NAME?: string; + /** Whether rate limiting is enabled (default true). */ + RATE_LIMIT_ENABLED: boolean; + /** Per-pool rate limit overrides (requests per minute). */ + RATE_LIMIT_READ_PER_MIN: number; + RATE_LIMIT_SYNC_PER_MIN: number; + RATE_LIMIT_SESSION_PER_MIN: number; + RATE_LIMIT_WRITE_PER_MIN: number; + RATE_LIMIT_CHALLENGE_PER_MIN: number; + RATE_LIMIT_MAX_CONNECTIONS: number; + RATE_LIMIT_FIREHOSE_PER_IP: number; } const REQUIRED_KEYS = [ @@ -119,6 +129,14 @@ export function loadConfig(envPath?: string): Config { NODE_DID: `did:web:${pdsHostname.replace(/:/g, "%3A")}`, NODE_MANAGERS: (process.env.NODE_MANAGERS ?? "").split(",").map(s => s.trim()).filter(Boolean), NODE_NAME: process.env.NODE_NAME || undefined, + RATE_LIMIT_ENABLED: process.env.RATE_LIMIT_ENABLED !== "false", + RATE_LIMIT_READ_PER_MIN: parseInt(process.env.RATE_LIMIT_READ_PER_MIN ?? "300", 10), + RATE_LIMIT_SYNC_PER_MIN: parseInt(process.env.RATE_LIMIT_SYNC_PER_MIN ?? "30", 10), + RATE_LIMIT_SESSION_PER_MIN: parseInt(process.env.RATE_LIMIT_SESSION_PER_MIN ?? "10", 10), + RATE_LIMIT_WRITE_PER_MIN: parseInt(process.env.RATE_LIMIT_WRITE_PER_MIN ?? "200", 10), + RATE_LIMIT_CHALLENGE_PER_MIN: parseInt(process.env.RATE_LIMIT_CHALLENGE_PER_MIN ?? "20", 10), + RATE_LIMIT_MAX_CONNECTIONS: parseInt(process.env.RATE_LIMIT_MAX_CONNECTIONS ?? "100", 10), + RATE_LIMIT_FIREHOSE_PER_IP: parseInt(process.env.RATE_LIMIT_FIREHOSE_PER_IP ?? "3", 10), }; } diff --git a/src/index.ts b/src/index.ts index f9ba488..09cd44f 100644 --- a/src/index.ts +++ b/src/index.ts @@ -1,6 +1,14 @@ import { Hono } from "hono"; import { cors } from "hono/cors"; import { requireAuth } from "./middleware/auth.js"; +import { rateLimitMiddleware } from "./middleware/rate-limit.js"; +import { + jsonBodyLimit, + challengeBodyLimit, + blobBodyLimit, + carBodyLimit, +} from "./middleware/body-limit.js"; +import type { RateLimiter } from "./rate-limiter.js"; import type { RepoManager } from "./repo-manager.js"; import type { Firehose } from "./firehose.js"; import type { Config } from "./config.js"; @@ -15,6 +23,7 @@ import * as admin from "./xrpc/admin.js"; import { respondToChallenge } from "./replication/challenge-response/challenge-responder.js"; import { serializeResponse } from "./replication/challenge-response/http-transport.js"; import type { StorageChallenge } from "./replication/challenge-response/types.js"; +import { MAX_RECORD_PATHS, MAX_BLOCK_CIDS } from "./replication/challenge-response/types.js"; import { generateMstProof } from "./replication/mst-proof.js"; import { generateNodeDidDocument } from "./node-identity.js"; @@ -42,6 +51,7 @@ export function createApp( replicationManager?: ReplicationManager, replicatedRepoReader?: ReplicatedRepoReader, nodeOpts?: NodeIdentityOpts, + rateLimiter?: RateLimiter, ) { const nodeDid = nodeOpts?.nodeDid ?? config.NODE_DID; const nodePublicKeyMultibase = nodeOpts?.nodePublicKeyMultibase ?? ""; @@ -72,6 +82,101 @@ export function createApp( }), ); + // ============================================ + // Rate limit + body size middleware (per route group) + // ============================================ + if (rateLimiter && config.RATE_LIMIT_ENABLED) { + const w = 60_000; // 1 minute window + + // Meta: .well-known, health, describeServer, resolveHandle + const metaRL = rateLimitMiddleware(rateLimiter, { + pool: "meta", + rule: { maxRequests: 600, windowMs: w }, + }); + app.use("/.well-known/*", metaRL); + app.use("/xrpc/_health", metaRL); + app.use("/xrpc/com.atproto.server.describeServer", metaRL); + app.use("/xrpc/com.atproto.identity.resolveHandle", metaRL); + + // RASL + app.use( + "/.well-known/rasl/*", + rateLimitMiddleware(rateLimiter, { + pool: "rasl", + rule: { maxRequests: 600, windowMs: w }, + }), + ); + + // Sync endpoints + app.use( + "/xrpc/com.atproto.sync.*", + rateLimitMiddleware(rateLimiter, { + pool: "sync", + rule: { maxRequests: config.RATE_LIMIT_SYNC_PER_MIN, windowMs: w }, + }), + ); + + // Session endpoints (login/refresh) + const sessionRL = rateLimitMiddleware(rateLimiter, { + pool: "session", + rule: { maxRequests: config.RATE_LIMIT_SESSION_PER_MIN, windowMs: w }, + }); + app.use("/xrpc/com.atproto.server.createSession", sessionRL); + app.use("/xrpc/com.atproto.server.refreshSession", sessionRL); + + // Read endpoints (repo reads) + const readRL = rateLimitMiddleware(rateLimiter, { + pool: "read", + rule: { maxRequests: config.RATE_LIMIT_READ_PER_MIN, windowMs: w }, + authRule: { maxRequests: 1000, windowMs: w }, + }); + app.use("/xrpc/com.atproto.repo.getRecord", readRL); + app.use("/xrpc/com.atproto.repo.listRecords", readRL); + app.use("/xrpc/com.atproto.repo.describeRepo", readRL); + + // Write endpoints (authed only) + const writeRL = rateLimitMiddleware(rateLimiter, { + pool: "write", + rule: { maxRequests: config.RATE_LIMIT_WRITE_PER_MIN, windowMs: w }, + }); + app.use("/xrpc/com.atproto.repo.createRecord", writeRL); + app.use("/xrpc/com.atproto.repo.deleteRecord", writeRL); + app.use("/xrpc/com.atproto.repo.putRecord", writeRL); + app.use("/xrpc/com.atproto.repo.applyWrites", writeRL); + + // Challenge endpoint + app.use( + "/xrpc/org.p2pds.verification.challenge", + rateLimitMiddleware(rateLimiter, { + pool: "challenge", + rule: { maxRequests: config.RATE_LIMIT_CHALLENGE_PER_MIN, windowMs: w }, + }), + ); + + // MST proof + app.use( + "/xrpc/org.p2pds.verification.getMstProof", + rateLimitMiddleware(rateLimiter, { + pool: "mstProof", + rule: { maxRequests: config.RATE_LIMIT_SYNC_PER_MIN, windowMs: w }, + }), + ); + + // Admin endpoints + const adminRL = rateLimitMiddleware(rateLimiter, { + pool: "admin", + rule: { maxRequests: 300, windowMs: w }, + }); + app.use("/xrpc/org.p2pds.admin.*", adminRL); + } + + // Body size limits (always active, independent of rate limiting) + app.use("/xrpc/org.p2pds.verification.challenge", challengeBodyLimit()); + app.use("/xrpc/com.atproto.repo.uploadBlob", blobBodyLimit()); + app.use("/xrpc/com.atproto.repo.importRepo", carBodyLimit()); + // Default 1MB JSON body limit for all other POST endpoints + app.post("/xrpc/*", jsonBodyLimit()); + // DID document for did:web resolution — serves the node's DID document app.get("/.well-known/did.json", (c) => { const didDocument = generateNodeDidDocument( @@ -463,6 +568,35 @@ ${handleHtml} } const challenge = (await c.req.json()) as StorageChallenge; + + // Validate challenge is targeted at this node + if (challenge.targetDid !== nodeDid) { + return c.json( + { error: "InvalidChallenge", message: "Challenge is not targeted at this node" }, + 400, + ); + } + // Limit work: cap record paths and block CIDs + if (challenge.recordPaths && challenge.recordPaths.length > MAX_RECORD_PATHS) { + return c.json( + { error: "InvalidChallenge", message: `Too many recordPaths (max ${MAX_RECORD_PATHS})` }, + 400, + ); + } + if (challenge.blockCids && challenge.blockCids.length > MAX_BLOCK_CIDS) { + return c.json( + { error: "InvalidChallenge", message: `Too many blockCids (max ${MAX_BLOCK_CIDS})` }, + 400, + ); + } + // Reject expired challenges + if (challenge.expiresAt && new Date(challenge.expiresAt).getTime() < Date.now()) { + return c.json( + { error: "InvalidChallenge", message: "Challenge has expired" }, + 400, + ); + } + const response = await respondToChallenge( challenge, blockStore, diff --git a/src/ipfs.test.ts b/src/ipfs.test.ts index 2bb6722..745b26f 100644 --- a/src/ipfs.test.ts +++ b/src/ipfs.test.ts @@ -45,6 +45,14 @@ function testConfig(dataDir: string): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/ipfs.ts b/src/ipfs.ts index de827c9..5947435 100644 --- a/src/ipfs.ts +++ b/src/ipfs.ts @@ -4,6 +4,11 @@ import { FsDatastore } from "datastore-fs"; import type { Helia } from "@helia/interface"; import type { BlockMap } from "@atproto/repo"; import { encode as cborEncode, decode as cborDecode } from "./cbor-compat.js"; +import type { RateLimiter } from "./rate-limiter.js"; +import { DEFAULT_RATE_LIMIT_CONFIG } from "./rate-limiter.js"; + +/** Maximum gossipsub message size (8 KB). Messages larger than this are dropped before decoding. */ +const MAX_GOSSIPSUB_MESSAGE_SIZE = 8192; /** * Pure storage: put, get, has blocks by CID string. @@ -104,11 +109,20 @@ export class IpfsService implements BlockStore, NetworkService { private commitHandlers: CommitNotificationHandler[] = []; private identityHandlers: IdentityNotificationHandler[] = []; private subscribedTopics: Set = new Set(); + private rateLimiter: RateLimiter | null = null; constructor(config: IpfsConfig) { this.config = config; } + /** + * Set a rate limiter for gossipsub message rate limiting. + * Called from server.ts after construction. + */ + setRateLimiter(limiter: RateLimiter): void { + this.rateLimiter = limiter; + } + async start(): Promise { this.blockstore = new FsBlockstore(this.config.blocksPath); @@ -121,7 +135,15 @@ export class IpfsService implements BlockStore, NetworkService { libp2pConfig.services.pubsub = gossipsub({ emitSelf: false, allowPublishToZeroTopicPeers: true, + maxInboundDataLength: MAX_GOSSIPSUB_MESSAGE_SIZE, }); + // Connection manager limits + libp2pConfig.connectionManager = { + ...libp2pConfig.connectionManager, + maxConnections: 100, + maxIncomingPendingConnections: 10, + inboundConnectionThreshold: 5, + }; this.helia = await createHelia({ libp2p: libp2pConfig, @@ -394,6 +416,24 @@ export class IpfsService implements BlockStore, NetworkService { try { const detail = (evt as { detail: { topic: string; data: Uint8Array } }).detail; + // Drop oversized messages before decoding + if (detail.data.length > MAX_GOSSIPSUB_MESSAGE_SIZE) { + return; + } + + // Per-topic rate limiting via the shared RateLimiter + if (this.rateLimiter) { + const isCommit = detail.topic.startsWith(COMMIT_TOPIC_PREFIX); + const isIdentity = detail.topic.startsWith(IDENTITY_TOPIC_PREFIX); + if (isCommit || isIdentity) { + const rule = isCommit + ? DEFAULT_RATE_LIMIT_CONFIG.gossipsubCommit + : DEFAULT_RATE_LIMIT_CONFIG.gossipsubIdentity; + const result = this.rateLimiter.check("gossipsub", detail.topic, rule); + if (!result.allowed) return; // silently drop + } + } + if (detail.topic.startsWith(COMMIT_TOPIC_PREFIX)) { const notification = cborDecode(detail.data) as CommitNotification; if ( diff --git a/src/middleware/body-limit.ts b/src/middleware/body-limit.ts new file mode 100644 index 0000000..326de9c --- /dev/null +++ b/src/middleware/body-limit.ts @@ -0,0 +1,41 @@ +/** + * Body size limit middleware using Hono's built-in bodyLimit. + * + * Returns 413 Payload Too Large with a JSON error body when exceeded. + */ + +import { bodyLimit } from "hono/body-limit"; + +function limitHandler(maxSize: string) { + return () => + new Response( + JSON.stringify({ + error: "PayloadTooLarge", + message: `Request body exceeds ${maxSize} limit`, + }), + { + status: 413, + headers: { "Content-Type": "application/json" }, + }, + ); +} + +/** 1 MB — default for JSON endpoints. */ +export function jsonBodyLimit() { + return bodyLimit({ maxSize: 1024 * 1024, onError: limitHandler("1 MB") }); +} + +/** 64 KB — challenge POST payload. */ +export function challengeBodyLimit() { + return bodyLimit({ maxSize: 64 * 1024, onError: limitHandler("64 KB") }); +} + +/** 60 MB — blob uploads (matches existing uploadBlob check). */ +export function blobBodyLimit() { + return bodyLimit({ maxSize: 60 * 1024 * 1024, onError: limitHandler("60 MB") }); +} + +/** 100 MB — CAR file imports (matches existing importRepo check). */ +export function carBodyLimit() { + return bodyLimit({ maxSize: 100 * 1024 * 1024, onError: limitHandler("100 MB") }); +} diff --git a/src/middleware/rate-limit.test.ts b/src/middleware/rate-limit.test.ts new file mode 100644 index 0000000..54305ad --- /dev/null +++ b/src/middleware/rate-limit.test.ts @@ -0,0 +1,125 @@ +import { describe, it, expect, beforeEach, afterEach } from "vitest"; +import { Hono } from "hono"; +import { RateLimiter } from "../rate-limiter.js"; +import { rateLimitMiddleware, getClientIp } from "./rate-limit.js"; + +describe("rateLimitMiddleware", () => { + let limiter: RateLimiter; + + beforeEach(() => { + limiter = new RateLimiter(); + }); + + afterEach(() => { + limiter.stop(); + }); + + function createTestApp() { + const app = new Hono(); + app.use( + "/api/*", + rateLimitMiddleware(limiter, { + pool: "test", + rule: { maxRequests: 3, windowMs: 60_000 }, + authRule: { maxRequests: 10, windowMs: 60_000 }, + }), + ); + app.get("/api/hello", (c) => c.json({ ok: true })); + return app; + } + + it("allows requests under the limit", async () => { + const app = createTestApp(); + const res = await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + expect(res.status).toBe(200); + expect(res.headers.get("X-RateLimit-Limit")).toBe("3"); + expect(res.headers.get("X-RateLimit-Remaining")).toBe("2"); + }); + + it("returns 429 when limit exceeded", async () => { + const app = createTestApp(); + for (let i = 0; i < 3; i++) { + await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + } + const res = await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + expect(res.status).toBe(429); + const body = (await res.json()) as { error: string }; + expect(body.error).toBe("RateLimitExceeded"); + expect(res.headers.get("Retry-After")).toBeTruthy(); + }); + + it("uses auth rule for authenticated requests", async () => { + const app = createTestApp(); + // Exhaust unauth limit + for (let i = 0; i < 3; i++) { + await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + } + // Unauth should be blocked + const unauthRes = await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + expect(unauthRes.status).toBe(429); + + // Auth requests use a separate pool with higher limit + const authRes = await app.request("/api/hello", { + headers: { + "X-Forwarded-For": "1.2.3.4", + Authorization: "Bearer test-token", + }, + }); + expect(authRes.status).toBe(200); + expect(authRes.headers.get("X-RateLimit-Limit")).toBe("10"); + }); + + it("isolates different IPs", async () => { + const app = createTestApp(); + for (let i = 0; i < 3; i++) { + await app.request("/api/hello", { + headers: { "X-Forwarded-For": "1.2.3.4" }, + }); + } + const res = await app.request("/api/hello", { + headers: { "X-Forwarded-For": "5.6.7.8" }, + }); + expect(res.status).toBe(200); + }); +}); + +describe("getClientIp", () => { + function mockContext(headers: Record) { + return { + req: { + header: (name: string) => headers[name], + }, + } as unknown as Parameters[0]; + } + + it("prefers X-Forwarded-For", () => { + expect( + getClientIp( + mockContext({ + "X-Forwarded-For": "1.2.3.4, 5.6.7.8", + "X-Real-IP": "9.0.0.1", + }), + ), + ).toBe("1.2.3.4"); + }); + + it("falls back to X-Real-IP", () => { + expect(getClientIp(mockContext({ "X-Real-IP": "9.0.0.1" }))).toBe( + "9.0.0.1", + ); + }); + + it("returns unknown when no headers", () => { + expect(getClientIp(mockContext({}))).toBe("unknown"); + }); +}); diff --git a/src/middleware/rate-limit.ts b/src/middleware/rate-limit.ts new file mode 100644 index 0000000..e4e0ef2 --- /dev/null +++ b/src/middleware/rate-limit.ts @@ -0,0 +1,62 @@ +/** + * HTTP rate limit middleware for Hono. + * + * Keys by client IP (X-Forwarded-For → X-Real-IP → "unknown"). + * Supports separate rules for authenticated vs unauthenticated requests. + * Sets standard rate limit response headers. + */ + +import type { Context, Next, MiddlewareHandler } from "hono"; +import type { RateLimiter, RateLimitRule } from "../rate-limiter.js"; + +export interface RateLimitOptions { + pool: string; + rule: RateLimitRule; + /** Higher limits for authenticated requests (optional). */ + authRule?: RateLimitRule; +} + +/** + * Extract client IP from request headers. + * Trusts X-Forwarded-For (first hop) → X-Real-IP → "unknown". + */ +export function getClientIp(c: Context): string { + const xff = c.req.header("X-Forwarded-For"); + if (xff) { + const first = xff.split(",")[0]?.trim(); + if (first) return first; + } + const xri = c.req.header("X-Real-IP"); + if (xri) return xri.trim(); + return "unknown"; +} + +export function rateLimitMiddleware( + rateLimiter: RateLimiter, + options: RateLimitOptions, +): MiddlewareHandler { + return async (c: Context, next: Next) => { + const ip = getClientIp(c); + const isAuthed = c.req.header("Authorization")?.startsWith("Bearer ") ?? false; + const rule = isAuthed && options.authRule ? options.authRule : options.rule; + const poolName = isAuthed && options.authRule ? `${options.pool}:auth` : options.pool; + + const result = rateLimiter.check(poolName, ip, rule); + + c.header("X-RateLimit-Limit", String(rule.maxRequests)); + c.header("X-RateLimit-Remaining", String(result.remaining)); + + if (!result.allowed) { + c.header( + "Retry-After", + String(Math.ceil((result.retryAfterMs ?? 1000) / 1000)), + ); + return c.json( + { error: "RateLimitExceeded", message: "Too many requests" }, + 429, + ); + } + + await next(); + }; +} diff --git a/src/rate-limiter.test.ts b/src/rate-limiter.test.ts new file mode 100644 index 0000000..4302ffe --- /dev/null +++ b/src/rate-limiter.test.ts @@ -0,0 +1,136 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; +import { RateLimiter, DEFAULT_RATE_LIMIT_CONFIG } from "./rate-limiter.js"; +import type { RateLimitRule } from "./rate-limiter.js"; + +describe("RateLimiter", () => { + let limiter: RateLimiter; + + beforeEach(() => { + limiter = new RateLimiter(); + }); + + afterEach(() => { + limiter.stop(); + }); + + describe("check()", () => { + const rule: RateLimitRule = { maxRequests: 5, windowMs: 1000 }; + + it("allows requests under the limit", () => { + for (let i = 0; i < 5; i++) { + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(true); + expect(result.remaining).toBe(4 - i); + } + }); + + it("rejects requests at the limit", () => { + for (let i = 0; i < 5; i++) { + limiter.check("test", "ip1", rule); + } + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(false); + expect(result.remaining).toBe(0); + expect(result.retryAfterMs).toBeGreaterThan(0); + }); + + it("isolates different keys", () => { + for (let i = 0; i < 5; i++) { + limiter.check("test", "ip1", rule); + } + const result = limiter.check("test", "ip2", rule); + expect(result.allowed).toBe(true); + }); + + it("isolates different pools", () => { + for (let i = 0; i < 5; i++) { + limiter.check("pool1", "ip1", rule); + } + const result = limiter.check("pool2", "ip1", rule); + expect(result.allowed).toBe(true); + }); + + it("resets after window expires", () => { + vi.useFakeTimers(); + try { + for (let i = 0; i < 5; i++) { + limiter.check("test", "ip1", rule); + } + expect(limiter.check("test", "ip1", rule).allowed).toBe(false); + + // Advance past 2 full windows so previous count is also cleared + vi.advanceTimersByTime(2001); + + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(true); + expect(result.remaining).toBe(4); + } finally { + vi.useRealTimers(); + } + }); + + it("uses sliding window weight for gradual recovery", () => { + vi.useFakeTimers(); + try { + // Fill up the window + for (let i = 0; i < 5; i++) { + limiter.check("test", "ip1", rule); + } + expect(limiter.check("test", "ip1", rule).allowed).toBe(false); + + // Advance to 80% through the next window + // Previous window count (5) gets weighted by 0.2 = 1.0 effective + // So we should have ~4 remaining + vi.advanceTimersByTime(1800); + + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(true); + } finally { + vi.useRealTimers(); + } + }); + }); + + describe("cleanup", () => { + it("removes stale entries", () => { + vi.useFakeTimers(); + try { + const rule: RateLimitRule = { maxRequests: 10, windowMs: 1000 }; + limiter.check("test", "ip1", rule); + limiter.startCleanup(100); + + // Advance past 2 minutes (cleanup threshold) + vi.advanceTimersByTime(130_000); + + // Entry should be cleaned up, next check starts fresh + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(true); + expect(result.remaining).toBe(9); + } finally { + vi.useRealTimers(); + } + }); + }); + + describe("stop()", () => { + it("clears all state", () => { + const rule: RateLimitRule = { maxRequests: 1, windowMs: 60_000 }; + limiter.check("test", "ip1", rule); + limiter.stop(); + + // After stop, state is cleared — new check should succeed + const result = limiter.check("test", "ip1", rule); + expect(result.allowed).toBe(true); + }); + }); + + describe("DEFAULT_RATE_LIMIT_CONFIG", () => { + it("has expected pools", () => { + expect(DEFAULT_RATE_LIMIT_CONFIG.httpUnauthenticated.read.maxRequests).toBe(300); + expect(DEFAULT_RATE_LIMIT_CONFIG.httpUnauthenticated.sync.maxRequests).toBe(30); + expect(DEFAULT_RATE_LIMIT_CONFIG.httpUnauthenticated.session.maxRequests).toBe(10); + expect(DEFAULT_RATE_LIMIT_CONFIG.httpAuthenticated.write.maxRequests).toBe(200); + expect(DEFAULT_RATE_LIMIT_CONFIG.firehosePerIp.maxConnections).toBe(3); + }); + }); +}); diff --git a/src/rate-limiter.ts b/src/rate-limiter.ts new file mode 100644 index 0000000..f5e3bd0 --- /dev/null +++ b/src/rate-limiter.ts @@ -0,0 +1,141 @@ +/** + * Sliding window rate limiter. + * + * Uses a weighted sliding window counter: tracks the previous and current + * window counts, then computes an effective count by weighting the previous + * window's count by how far into the current window we are. + * + * Each (pool, key) pair has its own counter. Pools separate different + * endpoint groups (e.g. "sync", "read", "session") so their limits + * are independent. Keys are typically client IPs. + */ + +export interface RateLimitRule { + maxRequests: number; + windowMs: number; +} + +export interface RateLimitResult { + allowed: boolean; + remaining: number; + retryAfterMs?: number; +} + +interface WindowEntry { + prevCount: number; + currCount: number; + windowStart: number; +} + +export interface RateLimitConfig { + httpUnauthenticated: Record; + httpAuthenticated: Record; + challengeIncoming: RateLimitRule; + gossipsubCommit: RateLimitRule; + gossipsubIdentity: RateLimitRule; + firehosePerIp: { maxConnections: number }; + cleanupIntervalMs: number; +} + +export const DEFAULT_RATE_LIMIT_CONFIG: RateLimitConfig = { + httpUnauthenticated: { + read: { maxRequests: 300, windowMs: 60_000 }, + sync: { maxRequests: 30, windowMs: 60_000 }, + session: { maxRequests: 10, windowMs: 60_000 }, + challenge: { maxRequests: 20, windowMs: 60_000 }, + rasl: { maxRequests: 600, windowMs: 60_000 }, + mstProof: { maxRequests: 30, windowMs: 60_000 }, + meta: { maxRequests: 600, windowMs: 60_000 }, + }, + httpAuthenticated: { + read: { maxRequests: 1000, windowMs: 60_000 }, + write: { maxRequests: 200, windowMs: 60_000 }, + admin: { maxRequests: 300, windowMs: 60_000 }, + }, + challengeIncoming: { maxRequests: 10, windowMs: 300_000 }, + gossipsubCommit: { maxRequests: 60, windowMs: 60_000 }, + gossipsubIdentity: { maxRequests: 10, windowMs: 60_000 }, + firehosePerIp: { maxConnections: 3 }, + cleanupIntervalMs: 60_000, +}; + +export class RateLimiter { + private pools: Map> = new Map(); + private cleanupTimer: ReturnType | null = null; + + check(pool: string, key: string, rule: RateLimitRule): RateLimitResult { + const now = Date.now(); + let poolMap = this.pools.get(pool); + if (!poolMap) { + poolMap = new Map(); + this.pools.set(pool, poolMap); + } + + let entry = poolMap.get(key); + if (!entry) { + entry = { prevCount: 0, currCount: 0, windowStart: now }; + poolMap.set(key, entry); + } + + const elapsed = now - entry.windowStart; + + // If we've moved past the current window, slide forward + if (elapsed >= rule.windowMs) { + // How many full windows have passed? + if (elapsed >= rule.windowMs * 2) { + // Two or more windows: previous window is also gone + entry.prevCount = 0; + } else { + // One window: current becomes previous + entry.prevCount = entry.currCount; + } + entry.currCount = 0; + entry.windowStart = now - (elapsed % rule.windowMs); + } + + // Weighted sliding window: weight previous window by how far into current window + const currentElapsed = now - entry.windowStart; + const weight = 1 - currentElapsed / rule.windowMs; + const effectiveCount = entry.prevCount * weight + entry.currCount; + + if (effectiveCount >= rule.maxRequests) { + const retryAfterMs = rule.windowMs - currentElapsed; + return { + allowed: false, + remaining: 0, + retryAfterMs: Math.max(retryAfterMs, 1), + }; + } + + entry.currCount++; + const remaining = Math.max(0, Math.floor(rule.maxRequests - effectiveCount - 1)); + return { allowed: true, remaining }; + } + + startCleanup(intervalMs: number = 60_000): void { + if (this.cleanupTimer) return; + this.cleanupTimer = setInterval(() => { + const now = Date.now(); + for (const [poolName, poolMap] of this.pools) { + for (const [key, entry] of poolMap) { + // Remove entries that are older than 2 windows (conservative) + // Use the largest common window (60s) as baseline + if (now - entry.windowStart > 120_000) { + poolMap.delete(key); + } + } + if (poolMap.size === 0) { + this.pools.delete(poolName); + } + } + }, intervalMs); + } + + stop(): void { + if (this.cleanupTimer) { + clearInterval(this.cleanupTimer); + this.cleanupTimer = null; + } + this.pools.clear(); + } +} diff --git a/src/replication/challenge-response/challenge-response.test.ts b/src/replication/challenge-response/challenge-response.test.ts index c5a3048..7b5700a 100644 --- a/src/replication/challenge-response/challenge-response.test.ts +++ b/src/replication/challenge-response/challenge-response.test.ts @@ -37,6 +37,14 @@ function testConfig(dataDir: string): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/challenge-response/e2e-challenge.test.ts b/src/replication/challenge-response/e2e-challenge.test.ts index 3f24312..afb9035 100644 --- a/src/replication/challenge-response/e2e-challenge.test.ts +++ b/src/replication/challenge-response/e2e-challenge.test.ts @@ -45,6 +45,14 @@ function testConfig(dataDir: string): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/challenge-response/libp2p-transport.ts b/src/replication/challenge-response/libp2p-transport.ts index 1fbfd48..fb1abf5 100644 --- a/src/replication/challenge-response/libp2p-transport.ts +++ b/src/replication/challenge-response/libp2p-transport.ts @@ -14,29 +14,37 @@ import type { Libp2p, Stream } from "@libp2p/interface"; import { multiaddr } from "@multiformats/multiaddr"; import type { StorageChallenge, StorageChallengeResponse } from "./types.js"; +import { MAX_CHALLENGE_SIZE } from "./types.js"; import type { ChallengeTransport } from "./transport.js"; import { serializeResponse, deserializeResponse } from "./http-transport.js"; +/** Maximum response size (1 MB) — responses can include proof data. */ +const MAX_RESPONSE_SIZE = 1024 * 1024; + export const CHALLENGE_PROTOCOL = "/p2pds/challenge/1.0.0"; /** * Collect all chunks from a libp2p stream into a single Uint8Array. * Stream chunks may be Uint8Array or Uint8ArrayList; normalize via subarray(). + * Throws if accumulated bytes exceed maxSize (abuse prevention). */ -async function collectStream(stream: AsyncIterable): Promise { +async function collectStream( + stream: AsyncIterable, + maxSize: number = MAX_CHALLENGE_SIZE, +): Promise { const chunks: Uint8Array[] = []; + let totalSize = 0; for await (const chunk of stream) { - if (chunk instanceof Uint8Array) { - chunks.push(chunk); - } else { - // Uint8ArrayList — convert to Uint8Array - chunks.push(chunk.subarray()); + const bytes = chunk instanceof Uint8Array ? chunk : chunk.subarray(); + totalSize += bytes.length; + if (totalSize > maxSize) { + throw new Error(`Stream exceeded maximum size of ${maxSize} bytes`); } + chunks.push(bytes); } if (chunks.length === 0) return new Uint8Array(0); if (chunks.length === 1) return chunks[0]!; - const total = chunks.reduce((acc, c) => acc + c.length, 0); - const result = new Uint8Array(total); + const result = new Uint8Array(totalSize); let offset = 0; for (const c of chunks) { result.set(c, offset); @@ -67,9 +75,10 @@ export class Libp2pChallengeTransport implements ChallengeTransport { stream.send(challengeBytes); await stream.close(); // flush + close write; stream remains readable - // Read response + // Read response (allow up to 1 MB for proof data) const responseBytes = await collectStream( stream as unknown as AsyncIterable, + MAX_RESPONSE_SIZE, ); const raw = JSON.parse(new TextDecoder().decode(responseBytes)); return deserializeResponse(raw); diff --git a/src/replication/challenge-response/types.ts b/src/replication/challenge-response/types.ts index 7c25d8e..0245604 100644 --- a/src/replication/challenge-response/types.ts +++ b/src/replication/challenge-response/types.ts @@ -17,6 +17,17 @@ import type { MstProof } from "../mst-proof.js"; export const CHALLENGE_PROTOCOL_VERSION = 1; +// ============================================ +// Size / count limits for abuse prevention +// ============================================ + +/** Maximum number of record paths in a single challenge. */ +export const MAX_RECORD_PATHS = 10; +/** Maximum number of block CIDs in a single challenge. */ +export const MAX_BLOCK_CIDS = 20; +/** Maximum serialized challenge size in bytes (64 KB). */ +export const MAX_CHALLENGE_SIZE = 65536; + // ============================================ // Challenge types // ============================================ diff --git a/src/replication/e2e-multi-node.test.ts b/src/replication/e2e-multi-node.test.ts index 930ea51..f0c056e 100644 --- a/src/replication/e2e-multi-node.test.ts +++ b/src/replication/e2e-multi-node.test.ts @@ -50,6 +50,14 @@ function testConfig(dataDir: string): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/firehose-incremental.test.ts b/src/replication/firehose-incremental.test.ts index cd91dcb..0951b49 100644 --- a/src/replication/firehose-incremental.test.ts +++ b/src/replication/firehose-incremental.test.ts @@ -64,6 +64,14 @@ function testConfig(dataDir: string, replicateDids: string[] = []): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:local.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/gossipsub-notifications.test.ts b/src/replication/gossipsub-notifications.test.ts index c425d34..6d0c0db 100644 --- a/src/replication/gossipsub-notifications.test.ts +++ b/src/replication/gossipsub-notifications.test.ts @@ -299,6 +299,14 @@ describe("ReplicationManager gossipsub integration", () => { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; const { RepoManager } = await import("../repo-manager.js"); @@ -359,6 +367,14 @@ describe("ReplicationManager gossipsub integration", () => { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; const { RepoManager } = await import("../repo-manager.js"); @@ -441,6 +457,14 @@ describe("ReplicationManager gossipsub integration", () => { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; const { RepoManager } = await import("../repo-manager.js"); diff --git a/src/replication/mst-proof.test.ts b/src/replication/mst-proof.test.ts index 4d83b42..10a2d44 100644 --- a/src/replication/mst-proof.test.ts +++ b/src/replication/mst-proof.test.ts @@ -30,6 +30,14 @@ function testConfig(dataDir: string): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/offer-manager.test.ts b/src/replication/offer-manager.test.ts index 02b0591..3fc7ea6 100644 --- a/src/replication/offer-manager.test.ts +++ b/src/replication/offer-manager.test.ts @@ -32,6 +32,14 @@ function testConfig(dataDir: string, did = "did:plc:local"): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/peer-freshness.test.ts b/src/replication/peer-freshness.test.ts index f4036fa..4e6d824 100644 --- a/src/replication/peer-freshness.test.ts +++ b/src/replication/peer-freshness.test.ts @@ -46,6 +46,14 @@ function testConfig(dataDir: string, replicateDids: string[] = []): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/policy-integration.test.ts b/src/replication/policy-integration.test.ts index 6634505..bcaf1db 100644 --- a/src/replication/policy-integration.test.ts +++ b/src/replication/policy-integration.test.ts @@ -54,6 +54,14 @@ function testConfig(dataDir: string, replicateDids: string[] = []): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/replication/replication.test.ts b/src/replication/replication.test.ts index a5834a5..c76aa0a 100644 --- a/src/replication/replication.test.ts +++ b/src/replication/replication.test.ts @@ -57,6 +57,14 @@ function testConfig(dataDir: string, replicateDids: string[] = []): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/server.ts b/src/server.ts index 33200ec..bdfa997 100644 --- a/src/server.ts +++ b/src/server.ts @@ -23,10 +23,18 @@ import { FailoverChallengeTransport } from "./replication/challenge-response/fai import type { ChallengeTransport } from "./replication/challenge-response/transport.js"; import type { Libp2p } from "@libp2p/interface"; import { loadOrCreateNodeIdentity, getPublicKeyMultibase } from "./node-identity.js"; +import { RateLimiter } from "./rate-limiter.js"; // Load configuration const config = loadConfig(); +// Initialize rate limiter +let rateLimiter: RateLimiter | undefined; +if (config.RATE_LIMIT_ENABLED) { + rateLimiter = new RateLimiter(); + rateLimiter.startCleanup(60_000); +} + // Ensure data directory exists const dataDir = resolve(config.DATA_DIR); mkdirSync(dataDir, { recursive: true }); @@ -128,6 +136,11 @@ if (ipfsService && hasReplicateDids) { replicationManager.setReplicatedRepoReader(replicatedRepoReader); } +// Pass rate limiter to IPFS service for gossipsub rate limiting +if (rateLimiter && ipfsService) { + ipfsService.setRateLimiter(rateLimiter); +} + // Create Hono app const app = createApp( config, @@ -143,20 +156,43 @@ const app = createApp( nodeRepoManager, repoManager, }, + rateLimiter, ); // Create HTTP server using @hono/node-server's request listener const requestListener = getRequestListener(app.fetch); const httpServer = createServer(requestListener); -// Set up WebSocket server for firehose +// Set up WebSocket server for firehose with per-IP connection limits const wss = new WebSocketServer({ noServer: true }); +const firehoseConnections = new Map(); +const maxFirehosePerIp = config.RATE_LIMIT_FIREHOSE_PER_IP; httpServer.on("upgrade", (request, socket, head) => { const url = new URL(request.url ?? "/", `http://localhost:${config.PORT}`); if (url.pathname === "/xrpc/com.atproto.sync.subscribeRepos") { + const ip = + (request.headers["x-forwarded-for"] as string)?.split(",")[0]?.trim() ?? + request.headers["x-real-ip"] as string ?? + "unknown"; + + const current = firehoseConnections.get(ip) ?? 0; + if (config.RATE_LIMIT_ENABLED && current >= maxFirehosePerIp) { + socket.destroy(); + return; + } + wss.handleUpgrade(request, socket, head, (ws) => { + firehoseConnections.set(ip, (firehoseConnections.get(ip) ?? 0) + 1); + ws.on("close", () => { + const count = (firehoseConnections.get(ip) ?? 1) - 1; + if (count <= 0) { + firehoseConnections.delete(ip); + } else { + firehoseConnections.set(ip, count); + } + }); firehose.handleConnection(ws, request); }); } else { @@ -267,6 +303,9 @@ httpServer.listen(config.PORT, async () => { function shutdown() { console.log(pc.dim("\nShutting down...")); const cleanup = async () => { + if (rateLimiter) { + rateLimiter.stop(); + } if (replicationManager) { replicationManager.stop(); } diff --git a/src/xrpc/admin-e2e.test.ts b/src/xrpc/admin-e2e.test.ts index b27fce0..da7ab44 100644 --- a/src/xrpc/admin-e2e.test.ts +++ b/src/xrpc/admin-e2e.test.ts @@ -47,6 +47,14 @@ function makeConfig( FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } diff --git a/src/xrpc/admin.test.ts b/src/xrpc/admin.test.ts index dae4947..be723d2 100644 --- a/src/xrpc/admin.test.ts +++ b/src/xrpc/admin.test.ts @@ -33,6 +33,14 @@ function testConfig(dataDir: string, replicateDids: string[] = []): Config { FIREHOSE_ENABLED: false, NODE_DID: "did:web:test.example.com", NODE_MANAGERS: [], + RATE_LIMIT_ENABLED: false, + RATE_LIMIT_READ_PER_MIN: 300, + RATE_LIMIT_SYNC_PER_MIN: 30, + RATE_LIMIT_SESSION_PER_MIN: 10, + RATE_LIMIT_WRITE_PER_MIN: 200, + RATE_LIMIT_CHALLENGE_PER_MIN: 20, + RATE_LIMIT_MAX_CONNECTIONS: 100, + RATE_LIMIT_FIREHOSE_PER_IP: 3, }; } -- 2.51.2