From a67cefbd8019b33e39f1c6e3eeba36263f00bb46 Mon Sep 17 00:00:00 2001 From: unnoq Date: Tue, 4 Nov 2025 10:00:47 +0700 Subject: [PATCH 01/23] init --- packages/ratelimit/.gitignore | 26 +++++++++ packages/ratelimit/README.md | 81 ++++++++++++++++++++++++++++ packages/ratelimit/package.json | 75 ++++++++++++++++++++++++++ packages/ratelimit/src/index.test.ts | 3 ++ packages/ratelimit/src/index.ts | 0 packages/ratelimit/tsconfig.json | 17 ++++++ pnpm-lock.yaml | 37 +++++++++++++ 7 files changed, 239 insertions(+) create mode 100644 packages/ratelimit/.gitignore create mode 100644 packages/ratelimit/README.md create mode 100644 packages/ratelimit/package.json create mode 100644 packages/ratelimit/src/index.test.ts create mode 100644 packages/ratelimit/src/index.ts create mode 100644 packages/ratelimit/tsconfig.json diff --git a/packages/ratelimit/.gitignore b/packages/ratelimit/.gitignore new file mode 100644 index 000000000..f3620b55e --- /dev/null +++ b/packages/ratelimit/.gitignore @@ -0,0 +1,26 @@ +# Hidden folders and files +.* +!.gitignore +!.*.example + +# Common generated folders +logs/ +node_modules/ +out/ +dist/ +dist-ssr/ +build/ +coverage/ +temp/ + +# Common generated files +*.log +*.log.* +*.tsbuildinfo +*.vitest-temp.json +vite.config.ts.timestamp-* +vitest.config.ts.timestamp-* + +# Common manual ignore files +*.local +*.pem \ No newline at end of file diff --git a/packages/ratelimit/README.md b/packages/ratelimit/README.md new file mode 100644 index 000000000..1d3532c85 --- /dev/null +++ b/packages/ratelimit/README.md @@ -0,0 +1,81 @@ +
+ oRPC logo +
+ +

+ +
+ + codecov + + + weekly downloads + + + MIT License + + + Discord + + + Ask DeepWiki + +
+ +

Typesafe APIs Made Simple 🪄

+ +**oRPC is a powerful combination of RPC and OpenAPI**, makes it easy to build APIs that are end-to-end type-safe and adhere to OpenAPI standards + +--- + +## Highlights + +- **🔗 End-to-End Type Safety**: Ensure type-safe inputs, outputs, and errors from client to server. +- **📘 First-Class OpenAPI**: Built-in support that fully adheres to the OpenAPI standard. +- **📝 Contract-First Development**: Optionally define your API contract before implementation. +- **🔍 First-Class OpenTelemetry**: Seamlessly integrate with OpenTelemetry for observability. +- **⚙️ Framework Integrations**: Seamlessly integrate with TanStack Query (React, Vue, Solid, Svelte, Angular), SWR, Pinia Colada, and more. +- **🚀 Server Actions**: Fully compatible with React Server Actions on Next.js, TanStack Start, and other platforms. +- **🔠 Standard Schema Support**: Works out of the box with Zod, Valibot, ArkType, and other schema validators. +- **🗃️ Native Types**: Supports native types like Date, File, Blob, BigInt, URL, and more. +- **⏱️ Lazy Router**: Enhance cold start times with our lazy routing feature. +- **📡 SSE & Streaming**: Enjoy full type-safe support for SSE and streaming. +- **🌍 Multi-Runtime Support**: Fast and lightweight on Cloudflare, Deno, Bun, Node.js, and beyond. +- **🔌 Extendability**: Easily extend functionality with plugins, middleware, and interceptors. + +## Documentation + +You can find the full documentation [here](https://orpc.unnoq.com). + +## Packages + +- [@orpc/contract](https://www.npmjs.com/package/@orpc/contract): Build your API contract. +- [@orpc/server](https://www.npmjs.com/package/@orpc/server): Build your API or implement API contract. +- [@orpc/client](https://www.npmjs.com/package/@orpc/client): Consume your API on the client with type-safety. +- [@orpc/openapi](https://www.npmjs.com/package/@orpc/openapi): Generate OpenAPI specs and handle OpenAPI requests. +- [@orpc/otel](https://www.npmjs.com/package/@orpc/otel): [OpenTelemetry](https://opentelemetry.io/) integration for observability. +- [@orpc/nest](https://www.npmjs.com/package/@orpc/nest): Deeply integrate oRPC with [NestJS](https://nestjs.com/). +- [@orpc/react](https://www.npmjs.com/package/@orpc/react): Utilities for integrating oRPC with React and React Server Actions. +- [@orpc/tanstack-query](https://www.npmjs.com/package/@orpc/tanstack-query): [TanStack Query](https://tanstack.com/query/latest) integration. +- [@orpc/experimental-react-swr](https://www.npmjs.com/package/@orpc/experimental-react-swr): [SWR](https://swr.vercel.app/) integration. +- [@orpc/vue-colada](https://www.npmjs.com/package/@orpc/vue-colada): Integration with [Pinia Colada](https://pinia-colada.esm.dev/). +- [@orpc/hey-api](https://www.npmjs.com/package/@orpc/hey-api): [Hey API](https://heyapi.dev/) integration. +- [@orpc/zod](https://www.npmjs.com/package/@orpc/zod): More schemas that [Zod](https://zod.dev/) doesn't support yet. +- [@orpc/valibot](https://www.npmjs.com/package/@orpc/valibot): OpenAPI spec generation from [Valibot](https://valibot.dev/). +- [@orpc/arktype](https://www.npmjs.com/package/@orpc/arktype): OpenAPI spec generation from [ArkType](https://arktype.io/). + +## `@orpc/experimental-ratelimit` + +Rate Limiting Feature for oRPC + +## Sponsors + +

+ + + +

+ +## License + +Distributed under the MIT License. See [LICENSE](https://github.com/unnoq/orpc/blob/main/LICENSE) for more information. diff --git a/packages/ratelimit/package.json b/packages/ratelimit/package.json new file mode 100644 index 000000000..7cd6c5779 --- /dev/null +++ b/packages/ratelimit/package.json @@ -0,0 +1,75 @@ +{ + "name": "@orpc/experimental-ratelimit", + "type": "module", + "version": "0.0.0", + "license": "MIT", + "homepage": "https://orpc.unnoq.com", + "repository": { + "type": "git", + "url": "git+https://github.com/unnoq/orpc.git", + "directory": "packages/ratelimit" + }, + "keywords": [ + "unnoq", + "orpc" + ], + "publishConfig": { + "exports": { + ".": { + "types": "./dist/index.d.mts", + "import": "./dist/index.mjs", + "default": "./dist/index.mjs" + }, + "./memory": { + "types": "./dist/adapters/memory.d.mts", + "import": "./dist/adapters/memory.mjs", + "default": "./dist/adapters/memory.mjs" + }, + "./ioredis": { + "types": "./dist/adapters/ioredis.d.mts", + "import": "./dist/adapters/ioredis.mjs", + "default": "./dist/adapters/ioredis.mjs" + }, + "./upstash-ratelimit": { + "types": "./dist/adapters/upstash-ratelimit.d.mts", + "import": "./dist/adapters/upstash-ratelimit.mjs", + "default": "./dist/adapters/upstash-ratelimit.mjs" + } + } + }, + "exports": { + ".": "./src/index.ts", + "./memory": "./src/adapters/memory.ts", + "./ioredis": "./src/adapters/ioredis.ts", + "./upstash-ratelimit": "./src/adapters/upstash-ratelimit.ts" + }, + "files": [ + "dist" + ], + "scripts": { + "build": "unbuild", + "build:watch": "pnpm run build --watch", + "type:check": "tsc -b" + }, + "peerDependencies": { + "@upstash/ratelimit": ">=2.0.7", + "ioredis": ">=5.8.1" + }, + "peerDependenciesMeta": { + "@upstash/ratelimit": { + "optional": true + }, + "ioredis": { + "optional": true + } + }, + "dependencies": { + "@orpc/client": "workspace:*", + "@orpc/shared": "workspace:*", + "@orpc/standard-server": "workspace:*" + }, + "devDependencies": { + "@upstash/ratelimit": "^2.0.7", + "ioredis": "^5.8.2" + } +} diff --git a/packages/ratelimit/src/index.test.ts b/packages/ratelimit/src/index.test.ts new file mode 100644 index 000000000..8297b9d7e --- /dev/null +++ b/packages/ratelimit/src/index.test.ts @@ -0,0 +1,3 @@ +it('exports Publisher', async () => { + expect(Object.keys(await import('./index'))).toContain('Publisher') +}) diff --git a/packages/ratelimit/src/index.ts b/packages/ratelimit/src/index.ts new file mode 100644 index 000000000..e69de29bb diff --git a/packages/ratelimit/tsconfig.json b/packages/ratelimit/tsconfig.json new file mode 100644 index 000000000..954616514 --- /dev/null +++ b/packages/ratelimit/tsconfig.json @@ -0,0 +1,17 @@ +{ + "extends": "../../tsconfig.lib.json", + "references": [ + { "path": "../shared" }, + { "path": "../client" }, + { "path": "../standard-server" } + ], + "include": ["src"], + "exclude": [ + "**/*.bench.*", + "**/*.test.*", + "**/*.test-d.ts", + "**/__tests__/**", + "**/__mocks__/**", + "**/__snapshots__/**" + ] +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index de755b092..9b9e494c1 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -553,6 +553,25 @@ importers: specifier: ^5.8.2 version: 5.8.2 + packages/ratelimit: + dependencies: + '@orpc/client': + specifier: workspace:* + version: link:../client + '@orpc/shared': + specifier: workspace:* + version: link:../shared + '@orpc/standard-server': + specifier: workspace:* + version: link:../standard-server + devDependencies: + '@upstash/ratelimit': + specifier: ^2.0.7 + version: 2.0.7(@upstash/redis@1.35.6) + ioredis: + specifier: ^5.8.2 + version: 5.8.2 + packages/react: dependencies: '@orpc/client': @@ -6560,6 +6579,15 @@ packages: peerDependencies: vue: '>=3.5.18' + '@upstash/core-analytics@0.0.10': + resolution: {integrity: sha512-7qJHGxpQgQr9/vmeS1PktEwvNAF7TI4iJDi8Pu2CFZ9YUGHZH4fOP5TfYlZ4aVxfopnELiE4BS4FBjyK7V1/xQ==} + engines: {node: '>=16.0.0'} + + '@upstash/ratelimit@2.0.7': + resolution: {integrity: sha512-qNQW4uBPKVk8c4wFGj2S/vfKKQxXx1taSJoSGBN36FeiVBBKHQgsjPbKUijZ9Xu5FyVK+pfiXWKIsQGyoje8Fw==} + peerDependencies: + '@upstash/redis': ^1.34.3 + '@upstash/redis@1.35.6': resolution: {integrity: sha512-aSEIGJgJ7XUfTYvhQcQbq835re7e/BXjs8Janq6Pvr6LlmTZnyqwT97RziZLO/8AVUL037RLXqqiQC6kCt+5pA==} @@ -20265,6 +20293,15 @@ snapshots: unhead: 2.0.19 vue: 3.5.22(typescript@5.8.3) + '@upstash/core-analytics@0.0.10': + dependencies: + '@upstash/redis': 1.35.6 + + '@upstash/ratelimit@2.0.7(@upstash/redis@1.35.6)': + dependencies: + '@upstash/core-analytics': 0.0.10 + '@upstash/redis': 1.35.6 + '@upstash/redis@1.35.6': dependencies: uncrypto: 0.1.3 From c518f421027c776e22da03b82f07fc39d7c8f289 Mon Sep 17 00:00:00 2001 From: unnoq Date: Tue, 4 Nov 2025 15:44:58 +0700 Subject: [PATCH 02/23] upstash ratelimit adapter --- packages/ratelimit/package.json | 1 + .../src/adapters/upstash-ratelimit.test.ts | 80 +++++++++++++++++++ .../src/adapters/upstash-ratelimit.ts | 48 +++++++++++ packages/ratelimit/src/index.ts | 1 + packages/ratelimit/src/types.ts | 28 +++++++ pnpm-lock.yaml | 3 + 6 files changed, 161 insertions(+) create mode 100644 packages/ratelimit/src/adapters/upstash-ratelimit.test.ts create mode 100644 packages/ratelimit/src/adapters/upstash-ratelimit.ts create mode 100644 packages/ratelimit/src/types.ts diff --git a/packages/ratelimit/package.json b/packages/ratelimit/package.json index 7cd6c5779..ed4d86fc6 100644 --- a/packages/ratelimit/package.json +++ b/packages/ratelimit/package.json @@ -70,6 +70,7 @@ }, "devDependencies": { "@upstash/ratelimit": "^2.0.7", + "@upstash/redis": "^1.35.6", "ioredis": "^5.8.2" } } diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts new file mode 100644 index 000000000..f667f5f01 --- /dev/null +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts @@ -0,0 +1,80 @@ +import { Ratelimit } from '@upstash/ratelimit' +import { Redis } from '@upstash/redis' +import { UpstashRatelimiter } from './upstash-ratelimit' + +const UPSTASH_REDIS_REST_URL = process.env.UPSTASH_REDIS_REST_URL +const UPSTASH_REDIS_REST_TOKEN = process.env.UPSTASH_REDIS_REST_TOKEN + +/** + * These tests depend on a real Upstash redis server — make sure to set the `UPSTASH_REDIS_REST_URL`, `UPSTASH_REDIS_REST_TOKEN` envs. + * When writing new tests, always use unique keys to avoid conflicts with other test cases. + */ +describe.concurrent( + 'upstash ratelimit adapter', + { skip: !UPSTASH_REDIS_REST_URL || !UPSTASH_REDIS_REST_TOKEN, timeout: 20_000 }, + () => { + function createTestingRatelimiter(options: ConstructorParameters[1] = {}) { + const redis = new Redis({ + url: UPSTASH_REDIS_REST_URL, + token: UPSTASH_REDIS_REST_TOKEN, + }) + + const ratelimit = new Ratelimit({ + redis, + limiter: Ratelimit.slidingWindow(3, '10 s'), + prefix: `rate-limit${crypto.randomUUID()}`, + }) + + return new UpstashRatelimiter(ratelimit, options) + } + + it('should limit key successfully and return valid result', async () => { + const ratelimiter = createTestingRatelimiter() + const key = `test-key-${crypto.randomUUID()}` + + const result1 = await ratelimiter.limit(key) + + expect(result1).toMatchObject({ + success: true, + limit: 3, + remaining: 2, + }) + }) + + // retry: 5 to deal with upstash's rare condition limitation + it('should block when blockingUntilReady is enabled', { retry: 5 }, async () => { + const timeoutMs = 2000 + const ratelimiter = createTestingRatelimiter({ + blockingUntilReady: { + enabled: true, + timeoutInMs: timeoutMs, + }, + }) + const key = `test-blocking-${crypto.randomUUID()}` + + // Fill up the rate limit first + while (true) { + const result = await ratelimiter.limit(key) + if (!result.success) { + break + } + } + + const startTime = Date.now() + await expect(ratelimiter.limit(key)).resolves.toMatchObject({ success: false }) + expect(Date.now() - startTime).toBeGreaterThanOrEqual(timeoutMs) // should wait until reach timeout + }) + + it('should use waitUntil callback when provided', async () => { + const waitUntilSpy = vi.fn() + const ratelimiter = createTestingRatelimiter({ + waitUtil: waitUntilSpy, + }) + const key = `test-waituntil-${crypto.randomUUID()}` + + const result1 = await ratelimiter.limit(key) + expect(waitUntilSpy).toHaveBeenCalledTimes(1) + expect(waitUntilSpy).toHaveBeenCalledWith(result1.pending) + }) + }, +) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts new file mode 100644 index 000000000..a5f70238b --- /dev/null +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -0,0 +1,48 @@ +import type { Ratelimit } from '@upstash/ratelimit' +import type { Ratelimiter, RatelimiterLimitResult } from '../types' + +export interface UpstashRatelimiterOptions { + /** + * Block until the request may pass or timeout is reached. + */ + blockingUntilReady?: { + enabled: boolean + timeoutInMs: number + } + + /** + * For the MultiRegion setup we do some synchronizing in the background, after returning the current limit. + * Or when analytics is enabled, we send the analytics asynchronously after returning the limit. + * In most case you can simply ignore this. + * + * On Vercel Edge or Cloudflare workers, you might need `.bind` before assign: + * ```ts + * const ratelimiter = new UpstashRatelimiter(ratelimit, { + * waitUtil: context.waitUntil.bind(context), + * }) + * ``` + */ + waitUtil?: (promise: Promise) => any +} + +export class UpstashRatelimiter implements Ratelimiter { + protected blockingUntilReady: UpstashRatelimiterOptions['blockingUntilReady'] + protected waitUtil: UpstashRatelimiterOptions['waitUtil'] + + constructor( + protected readonly ratelimit: Ratelimit, + options: UpstashRatelimiterOptions = {}, + ) { + this.blockingUntilReady = options.blockingUntilReady + this.waitUtil = options.waitUtil + } + + async limit(key: string): Promise { + const result = this.blockingUntilReady?.enabled + ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeoutInMs) + : await this.ratelimit.limit(key) + + this.waitUtil?.(result.pending) + return result + } +} diff --git a/packages/ratelimit/src/index.ts b/packages/ratelimit/src/index.ts index e69de29bb..c9f6f047d 100644 --- a/packages/ratelimit/src/index.ts +++ b/packages/ratelimit/src/index.ts @@ -0,0 +1 @@ +export * from './types' diff --git a/packages/ratelimit/src/types.ts b/packages/ratelimit/src/types.ts new file mode 100644 index 000000000..c3634d6e0 --- /dev/null +++ b/packages/ratelimit/src/types.ts @@ -0,0 +1,28 @@ +export interface RatelimiterLimitResult { + /** + * Whether the request may pass(true) or exceeded the limit(false) + */ + success: boolean + /** + * Maximum number of requests allowed within a window. + */ + limit?: number + /** + * How many requests the user has left within the current window. + */ + remaining?: number + /** + * Unix timestamp in milliseconds when the limits are reset. + */ + reset?: number + /** + * For the MultiRegion setup we do some synchronizing in the background, after returning the current limit. + * Or when analytics is enabled, we send the analytics asynchronously after returning the limit. + * In most case you can simply ignore this. + */ + pending?: Promise +} + +export interface Ratelimiter { + limit(key: string): Promise +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 9b9e494c1..6ede7c5df 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -568,6 +568,9 @@ importers: '@upstash/ratelimit': specifier: ^2.0.7 version: 2.0.7(@upstash/redis@1.35.6) + '@upstash/redis': + specifier: ^1.35.6 + version: 1.35.6 ioredis: specifier: ^5.8.2 version: 5.8.2 From 2f4f0820c06c7c7b38cddfec4e99f7b0d419434b Mon Sep 17 00:00:00 2001 From: unnoq Date: Tue, 4 Nov 2025 16:18:38 +0700 Subject: [PATCH 03/23] improve --- packages/ratelimit/src/adapters/upstash-ratelimit.test.ts | 2 +- packages/ratelimit/src/adapters/upstash-ratelimit.ts | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts index f667f5f01..159c74c31 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts @@ -47,7 +47,7 @@ describe.concurrent( const ratelimiter = createTestingRatelimiter({ blockingUntilReady: { enabled: true, - timeoutInMs: timeoutMs, + timeoutMs, }, }) const key = `test-blocking-${crypto.randomUUID()}` diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index a5f70238b..9047858f8 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -7,7 +7,7 @@ export interface UpstashRatelimiterOptions { */ blockingUntilReady?: { enabled: boolean - timeoutInMs: number + timeoutMs: number } /** @@ -39,7 +39,7 @@ export class UpstashRatelimiter implements Ratelimiter { async limit(key: string): Promise { const result = this.blockingUntilReady?.enabled - ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeoutInMs) + ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeoutMs) : await this.ratelimit.limit(key) this.waitUtil?.(result.pending) From 2d8b9c6e8cb8998dd58d5c022cf2169bb9c58e5e Mon Sep 17 00:00:00 2001 From: unnoq Date: Tue, 4 Nov 2025 16:19:05 +0700 Subject: [PATCH 04/23] improve --- packages/ratelimit/src/adapters/upstash-ratelimit.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 9047858f8..67e24dfb0 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -18,7 +18,7 @@ export interface UpstashRatelimiterOptions { * On Vercel Edge or Cloudflare workers, you might need `.bind` before assign: * ```ts * const ratelimiter = new UpstashRatelimiter(ratelimit, { - * waitUtil: context.waitUntil.bind(context), + * waitUtil: ctx.waitUntil.bind(ctx), * }) * ``` */ From b24783466d7f96a8501fa564cfc15564bea9c9e7 Mon Sep 17 00:00:00 2001 From: unnoq Date: Tue, 4 Nov 2025 20:09:47 +0700 Subject: [PATCH 05/23] improve --- packages/ratelimit/src/adapters/upstash-ratelimit.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 67e24dfb0..9a30bae60 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -26,11 +26,11 @@ export interface UpstashRatelimiterOptions { } export class UpstashRatelimiter implements Ratelimiter { - protected blockingUntilReady: UpstashRatelimiterOptions['blockingUntilReady'] - protected waitUtil: UpstashRatelimiterOptions['waitUtil'] + private blockingUntilReady: UpstashRatelimiterOptions['blockingUntilReady'] + private waitUtil: UpstashRatelimiterOptions['waitUtil'] constructor( - protected readonly ratelimit: Ratelimit, + private readonly ratelimit: Ratelimit, options: UpstashRatelimiterOptions = {}, ) { this.blockingUntilReady = options.blockingUntilReady From c5026676c602edaa6283c52a170d1ccf903bf374 Mon Sep 17 00:00:00 2001 From: unnoq Date: Wed, 5 Nov 2025 10:38:49 +0700 Subject: [PATCH 06/23] ratelimit: add ioredis adapter + tests; unify result shape; adapt upstash adapter & tests - Add IORedisRatelimiter implementation and full test suite for ioredis adapter. - Normalize RatelimiterLimitResult: rename `reset` -> `resetAtMs` and remove `pending`. - Update Upstash adapter to map Upstash response to the unified result shape and keep invoking waitUtil. - Adjust Upstash tests to expect waitUtil called with a Promise. --- .../ratelimit/src/adapters/ioredis.test.ts | 181 ++++++++++++++++++ packages/ratelimit/src/adapters/ioredis.ts | 155 +++++++++++++++ .../src/adapters/upstash-ratelimit.test.ts | 2 +- .../src/adapters/upstash-ratelimit.ts | 9 +- packages/ratelimit/src/types.ts | 8 +- 5 files changed, 345 insertions(+), 10 deletions(-) create mode 100644 packages/ratelimit/src/adapters/ioredis.test.ts create mode 100644 packages/ratelimit/src/adapters/ioredis.ts diff --git a/packages/ratelimit/src/adapters/ioredis.test.ts b/packages/ratelimit/src/adapters/ioredis.test.ts new file mode 100644 index 000000000..6e4f9e2b1 --- /dev/null +++ b/packages/ratelimit/src/adapters/ioredis.test.ts @@ -0,0 +1,181 @@ +import type { IORedisRatelimiterOptions } from './ioredis' +import { Redis } from 'ioredis' +import { IORedisRatelimiter } from './ioredis' + +const REDIS_URL = process.env.REDIS_URL + +/** + * These tests depend on a real Redis server — make sure to set the `REDIS_URL` env. + * When writing new tests, always use unique keys to avoid conflicts with other test cases. + */ +describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, () => { + let redis: Redis + + function createTestingRatelimiter(options: Partial = {}) { + const ratelimiter = new IORedisRatelimiter(redis, { + prefix: `test:${crypto.randomUUID()}:`, // isolated from other tests + maxRequests: 10, + windowMs: 60000, + ...options, + }) + return ratelimiter + } + + beforeAll(() => { + redis = new Redis(REDIS_URL!) + }) + + afterAll(async () => { + await redis.quit() + expect(redis.listenerCount('error')).toEqual(0) + }) + + describe('without blocking', () => { + it('allows requests within the limit', async () => { + const ratelimiter = createTestingRatelimiter({ maxRequests: 2, windowMs: 5000 }) + const key = 'user1' + + const result1 = await ratelimiter.limit(key) + expect(result1.success).toBe(true) + expect(result1.limit).toBe(2) + expect(result1.remaining).toBe(1) + expect(result1.resetAtMs).toBeGreaterThan(Date.now()) + + const result2 = await ratelimiter.limit(key) + expect(result2.success).toBe(true) + expect(result2.limit).toBe(2) + expect(result2.remaining).toBe(0) + expect(result2.resetAtMs).toEqual(result1.resetAtMs) + }) + + it('denies requests exceeding the limit', async () => { + const ratelimiter = createTestingRatelimiter({ maxRequests: 1, windowMs: 5000 }) + const key = 'user2' + + // reach the limit + const result1 = await ratelimiter.limit(key) + expect(result1.remaining).toBe(0) + + const result2 = await ratelimiter.limit(key) + expect(result2.success).toBe(false) + }) + + it('resets the limit after the window expires', async () => { + const ratelimiter = createTestingRatelimiter({ maxRequests: 1, windowMs: 2000 }) + const key = 'user3' + + // reach the limit + const result1 = await ratelimiter.limit(key) + expect(result1.remaining).toBe(0) + const result2 = await ratelimiter.limit(key) + expect(result2.success).toBe(false) + + // wait for the window to expire + await new Promise(resolve => setTimeout(resolve, 2000)) + + const result3 = await ratelimiter.limit(key) + expect(result3.success).toBe(true) + }) + }) + + describe('with blocking', () => { + it('blocks until the rate limit resets', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 1, + windowMs: 2000, + blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + }) + const key = 'user-blocking-1' + + // reach the limit + const result1 = await ratelimiter.limit(key) + expect(result1.remaining).toBe(0) + + const startTime = Date.now() + const result2 = await ratelimiter.limit(key) + const endTime = Date.now() + + expect(result2.success).toBe(true) + expect(endTime - startTime).toBeLessThanOrEqual(2000) + }) + + it('times out if the reset time is beyond the timeout', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 1, + windowMs: 60_000, + blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + }) + const key = 'user-blocking-2' + + await ratelimiter.limit(key) // Consume the first request + + const startTime = Date.now() + const result = await ratelimiter.limit(key) + const endTime = Date.now() + + expect(result.success).toBe(false) // Should fail as it times out + expect(endTime - startTime).toBeLessThan(2000) + }) + + it('handles concurrent blocking requests correctly', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 2, + windowMs: 1000, + blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + }) + const key = 'user-concurrent-blocking' + + const promises = [ + ratelimiter.limit(key), + ratelimiter.limit(key), + ratelimiter.limit(key), // This one should block + ] + + const results = await Promise.all(promises) + expect(results.filter(r => r.success).length).toBe(3) + }) + }) + + describe('edge cases', () => { + it('uses the correct prefix for keys & auto expiration', async () => { + const prefix = `custom-prefix:${crypto.randomUUID()}:` + const ratelimiter = createTestingRatelimiter({ + windowMs: 1000, + prefix, + maxRequests: 1, + }) + const key = 'user4' + + await ratelimiter.limit(key) + + const keys = await redis.keys(`${prefix}${key}`) + expect(keys).toHaveLength(1) + + // Wait until the key is auto-expired + await vi.waitFor(async () => { + const keysAfterExpiry = await redis.keys(`${prefix}${key}`) + expect(keysAfterExpiry).toHaveLength(0) + }, { timeout: 10_000 }) + }) + + it('handles Redis errors gracefully', async () => { + const mockRedis = { + ...redis, + eval: async () => { throw new Error('Redis error') }, + } as any + + const ratelimiter = new IORedisRatelimiter(mockRedis, { maxRequests: 10, windowMs: 60000 }) + await expect(ratelimiter.limit('some-key')).rejects.toThrow('Redis error') + }) + + it('handles invalid script response', async () => { + const mockRedis = { + ...redis, + eval: async () => [1, 2, 3], // Invalid response, should have 4 elements + } as any + + const ratelimiter = new IORedisRatelimiter(mockRedis, { maxRequests: 10, windowMs: 60000 }) + await expect(ratelimiter.limit('some-key')).rejects.toThrow('Invalid response from rate limit script') + }) + }) +}) diff --git a/packages/ratelimit/src/adapters/ioredis.ts b/packages/ratelimit/src/adapters/ioredis.ts new file mode 100644 index 000000000..536dee8ff --- /dev/null +++ b/packages/ratelimit/src/adapters/ioredis.ts @@ -0,0 +1,155 @@ +import type Redis from 'ioredis' +import type { Ratelimiter, RatelimiterLimitResult } from '../types' +import { fallback } from '@orpc/shared' + +/** + * Sliding window Lua script for Redis. + * + * This script implements atomic sliding window rate limiting using Redis sorted sets. + * It removes expired entries, checks the current count, and adds new requests atomically. + * + * @returns A tuple with [success, limit, remaining, resetAtMs] where: + * - success: 1 if request is allowed, 0 if rate limited + * - limit: The maximum number of requests allowed + * - remaining: Number of requests remaining in the window + * - resetAtMs: Unix timestamp (in milliseconds) when the window resets + */ +const SLIDING_WINDOW_LUA_SCRIPT = ` + local key = KEYS[1] + local now = tonumber(ARGV[1]) + local window = tonumber(ARGV[2]) + local limit = tonumber(ARGV[3]) + + local windowStart = now - window + -- Set TTL to window + small buffer (converted to seconds for EXPIRE) + -- The buffer ensures the key doesn't expire while still in use + local ttl = math.ceil(window / 1000) + 1 + + -- Remove expired entries + redis.call('ZREMRANGEBYSCORE', key, 0, windowStart) + + -- Get all valid entries with scores in one call + -- This replaces separate ZCARD and ZRANGE operations + local entries = redis.call('ZRANGE', key, 0, -1, 'WITHSCORES') + local current = #entries / 2 -- Each entry has value + score + + -- Calculate reset time (when oldest entry expires) + local resetAtMs + if current > 0 then + resetAtMs = tonumber(entries[2]) + window -- entries[2] is the oldest score + else + resetAtMs = now + window + end + + -- Check if limit is exceeded + if current >= limit then + return {0, limit, 0, resetAtMs} + end + + -- Add current request + redis.call('ZADD', key, now, now) + redis.call('EXPIRE', key, ttl) + + -- Calculate remaining requests + local remaining = limit - current - 1 + + return {1, limit, remaining, resetAtMs} +` + +export class IORedisRatelimiterError extends Error {} + +export interface IORedisRatelimiterOptions { + /** + * Block until the request may pass or timeout is reached. + */ + blockingUntilReady?: { + enabled: boolean + timeoutMs: number + } + + /** + * The prefix to use for Redis keys. + * + * @default orpc:ratelimit: + */ + prefix?: string + + /** + * Maximum number of requests allowed within the window. + */ + maxRequests: number + + /** + * The duration of the sliding window in milliseconds. + */ + windowMs: number + +} + +export class IORedisRatelimiter implements Ratelimiter { + private readonly prefix: string + private readonly maxRequests: number + private readonly windowMs: number + private readonly blockingUntilReady: IORedisRatelimiterOptions['blockingUntilReady'] + + constructor( + private readonly redis: Redis, + options: IORedisRatelimiterOptions, + ) { + this.prefix = fallback(options.prefix, 'orpc:ratelimit:') + this.maxRequests = options.maxRequests + this.windowMs = options.windowMs + this.blockingUntilReady = options.blockingUntilReady + } + + async limit(key: string): Promise> { + const prefixedKey = `${this.prefix}${key}` + + if (this.blockingUntilReady?.enabled) { + return await this.blockUntilReady(prefixedKey, this.blockingUntilReady.timeoutMs) + } + + return await this.checkLimit(prefixedKey) + } + + private async checkLimit(key: string) { + const result = await this.redis.eval( + SLIDING_WINDOW_LUA_SCRIPT, + 1, + key, + Date.now().toString(), + this.windowMs.toString(), + this.maxRequests.toString(), + ) as unknown + + if (!Array.isArray(result) || result.length !== 4) { + throw new IORedisRatelimiterError('Invalid response from rate limit script') + } + + const [success, limit, remaining, resetAtMs] = result as [number, number, number, number] + + return { + success: success === 1, + limit, + remaining, + resetAtMs, + } + } + + private async blockUntilReady(key: string, timeoutMs: number) { + const deadlineAtMs = Date.now() + timeoutMs + let result: Awaited> + + while (true) { + result = await this.checkLimit(key) + + if (result.success || result.resetAtMs > deadlineAtMs) { + break + } + + await new Promise(resolve => setTimeout(resolve, result.resetAtMs - Date.now())) + } + + return result + } +} diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts index 159c74c31..ccb44780a 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts @@ -74,7 +74,7 @@ describe.concurrent( const result1 = await ratelimiter.limit(key) expect(waitUntilSpy).toHaveBeenCalledTimes(1) - expect(waitUntilSpy).toHaveBeenCalledWith(result1.pending) + expect(waitUntilSpy).toHaveBeenCalledWith(expect.any(Promise)) }) }, ) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 9a30bae60..7a8b0cb3a 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -37,12 +37,17 @@ export class UpstashRatelimiter implements Ratelimiter { this.waitUtil = options.waitUtil } - async limit(key: string): Promise { + async limit(key: string): Promise> { const result = this.blockingUntilReady?.enabled ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeoutMs) : await this.ratelimit.limit(key) this.waitUtil?.(result.pending) - return result + return { + success: result.success, + limit: result.limit, + remaining: result.remaining, + resetAtMs: result.reset, + } } } diff --git a/packages/ratelimit/src/types.ts b/packages/ratelimit/src/types.ts index c3634d6e0..93d446bb1 100644 --- a/packages/ratelimit/src/types.ts +++ b/packages/ratelimit/src/types.ts @@ -14,13 +14,7 @@ export interface RatelimiterLimitResult { /** * Unix timestamp in milliseconds when the limits are reset. */ - reset?: number - /** - * For the MultiRegion setup we do some synchronizing in the background, after returning the current limit. - * Or when analytics is enabled, we send the analytics asynchronously after returning the limit. - * In most case you can simply ignore this. - */ - pending?: Promise + resetAtMs?: number } export interface Ratelimiter { From 63e5751c7925e9cdd23de542fdc497506d3d190e Mon Sep 17 00:00:00 2001 From: unnoq Date: Wed, 5 Nov 2025 10:52:40 +0700 Subject: [PATCH 07/23] improve --- packages/ratelimit/src/adapters/ioredis.ts | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/packages/ratelimit/src/adapters/ioredis.ts b/packages/ratelimit/src/adapters/ioredis.ts index 536dee8ff..127ab8a15 100644 --- a/packages/ratelimit/src/adapters/ioredis.ts +++ b/packages/ratelimit/src/adapters/ioredis.ts @@ -138,18 +138,15 @@ export class IORedisRatelimiter implements Ratelimiter { private async blockUntilReady(key: string, timeoutMs: number) { const deadlineAtMs = Date.now() + timeoutMs - let result: Awaited> while (true) { - result = await this.checkLimit(key) + const result = await this.checkLimit(key) if (result.success || result.resetAtMs > deadlineAtMs) { - break + return result } await new Promise(resolve => setTimeout(resolve, result.resetAtMs - Date.now())) } - - return result } } From 728d884f3504870a2b3574834ee1254650fe5cc5 Mon Sep 17 00:00:00 2001 From: unnoq Date: Wed, 5 Nov 2025 15:10:23 +0700 Subject: [PATCH 08/23] memory adapter --- .../ratelimit/src/adapters/memory.test.ts | 158 ++++++++++++++++++ packages/ratelimit/src/adapters/memory.ts | 119 +++++++++++++ 2 files changed, 277 insertions(+) create mode 100644 packages/ratelimit/src/adapters/memory.test.ts create mode 100644 packages/ratelimit/src/adapters/memory.ts diff --git a/packages/ratelimit/src/adapters/memory.test.ts b/packages/ratelimit/src/adapters/memory.test.ts new file mode 100644 index 000000000..00ec6b789 --- /dev/null +++ b/packages/ratelimit/src/adapters/memory.test.ts @@ -0,0 +1,158 @@ +import { sleep } from '@orpc/shared' +import { describe, expect, it } from 'vitest' +import { MemoryRatelimiter } from './memory' + +describe('memoryRatelimiter', () => { + function createTestingRatelimiter(options: Partial[0]> = {}) { + return new MemoryRatelimiter({ + maxRequests: 2, + windowMs: 1000, + ...options, + }) + } + + describe('basic rate limiting', () => { + it('should allow requests within limit', async () => { + const ratelimiter = createTestingRatelimiter() + + const result1 = await ratelimiter.limit('test') + expect(result1.success).toBe(true) + expect(result1.remaining).toBe(1) + expect(result1.limit).toBe(2) + expect(result1.resetAtMs).toBeGreaterThan(Date.now()) + + const result2 = await ratelimiter.limit('test') + expect(result2.success).toBe(true) + expect(result2.remaining).toBe(0) + expect(result2.limit).toBe(2) + expect(result2.resetAtMs).toEqual(result1.resetAtMs) + + const result3 = await ratelimiter.limit('test') + expect(result3.success).toBe(false) + expect(result3.remaining).toBe(0) + expect(result3.limit).toBe(2) + expect(result3.resetAtMs).toEqual(result1.resetAtMs) + }) + + it('should reset after window expires', async () => { + const ratelimiter = createTestingRatelimiter({ + windowMs: 200, + }) + + const result1 = await ratelimiter.limit('test') + expect(result1.remaining).toBe(1) + + await sleep(210) + + const result2 = await ratelimiter.limit('test') + expect(result2.success).toBe(true) + expect(result2.remaining).toBe(1) + }) + + it('should handle multiple keys independently', async () => { + const ratelimiter = createTestingRatelimiter() + + const result1 = await ratelimiter.limit('test1') + expect(result1.success).toBe(true) + expect(result1.remaining).toBe(1) + + const result2 = await ratelimiter.limit('test2') + expect(result2.success).toBe(true) + expect(result2.remaining).toBe(1) + }) + }) + + describe('blocking mode', () => { + it('should block until ready', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 1, + windowMs: 1000, + blockingUntilReady: { + enabled: true, + timeoutMs: 2000, + }, + }) + + const result1 = await ratelimiter.limit('test') + expect(result1.remaining).toEqual(0) + + const start = Date.now() + const result2 = await ratelimiter.limit('test') + const end = Date.now() + expect(result2.success).toEqual(true) + expect(end - start).toBeGreaterThanOrEqual(100) // actually await + expect(end - start).toBeLessThanOrEqual(1100) + }) + + it('should respect timeout', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 1, + windowMs: 2000, + blockingUntilReady: { + enabled: true, + timeoutMs: 1000, + }, + }) + + const result1 = await ratelimiter.limit('test') + expect(result1.remaining).toBe(0) + + const result2 = await ratelimiter.limit('test') + expect(result2.success).toBe(false) + }) + }) + + it('should handle concurrent requests correctly', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 3, + windowMs: 1000, + }) + + const test = async (key: string, request: number) => { + const results = await Promise.all( + Array.from({ length: 5 }, () => ratelimiter.limit(key)), + ) + + // Count successful and failed requests + const successful = results.filter(r => r.success).length + const failed = results.filter(r => !r.success).length + + // Should have exactly maxRequests successful requests + expect(successful).toBe(3) + expect(failed).toBe(2) + + // Verify remaining counts are consistent + const successfulResults = results.filter(r => r.success) + expect(successfulResults[0]!.remaining).toBe(2) + expect(successfulResults[1]!.remaining).toBe(1) + expect(successfulResults[2]!.remaining).toBe(0) + } + + await Promise.all( + Array.from({ length: 5 }, (_, i) => test(`test${i}`, i + 1)), + ) + }) + + describe('cleanup', () => { + it('should cleanup expired entries on next limit call', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 2, + windowMs: 200, + }) + + await ratelimiter.limit('test') + + // @ts-expect-error accessing private property for testing + expect(ratelimiter.store.size).toBe(1) + + // move time forward to pass window + await sleep(210) + + // Make another limit call to different key to trigger cleanup + await ratelimiter.limit('test2') + + // @ts-expect-error accessing private property for testing + expect(ratelimiter.store.size).toBe(1) // Only test2 should remain + }) + }) +}) diff --git a/packages/ratelimit/src/adapters/memory.ts b/packages/ratelimit/src/adapters/memory.ts new file mode 100644 index 000000000..699ff0754 --- /dev/null +++ b/packages/ratelimit/src/adapters/memory.ts @@ -0,0 +1,119 @@ +import type { Ratelimiter, RatelimiterLimitResult } from '../types' + +export interface MemoryRatelimiterOptions { + /** + * Block until the request may pass or timeout is reached. + */ + blockingUntilReady?: { + enabled: boolean + timeoutMs: number + } + + /** + * Maximum number of requests allowed within the window. + */ + maxRequests: number + + /** + * The duration of the sliding window in milliseconds. + */ + windowMs: number + +} + +export class MemoryRatelimiter implements Ratelimiter { + private readonly maxRequests: number + private readonly windowMs: number + private readonly blockingUntilReady: MemoryRatelimiterOptions['blockingUntilReady'] + private readonly store: Map + private lastCleanupTime: number | null = null + + constructor(options: MemoryRatelimiterOptions) { + this.maxRequests = options.maxRequests + this.windowMs = options.windowMs + this.blockingUntilReady = options.blockingUntilReady + this.store = new Map() + } + + limit(key: string): Promise> { + this.cleanup() + + if (this.blockingUntilReady?.enabled) { + return this.blockUntilReady(key, this.blockingUntilReady.timeoutMs) + } + + return this.checkLimit(key) + } + + private cleanup(): void { + const now = Date.now() + + // Only clean up once per window to avoid excessive processing + if (this.lastCleanupTime !== null && this.lastCleanupTime + this.windowMs > now) { + return + } + + this.lastCleanupTime = now + const windowStart = now - this.windowMs + + for (const [key, timestamps] of this.store) { + // remove expired timestamps + timestamps.splice(0, timestamps.findIndex(timestamp => timestamp < windowStart) + 1) + + if (timestamps.length === 0) { + this.store.delete(key) + } + } + } + + private async checkLimit(key: string): Promise> { + const now = Date.now() + const windowStart = now - this.windowMs + + let timestamps = this.store.get(key) + if (timestamps) { + // Remove expired timestamps + timestamps.splice(0, timestamps.findIndex(timestamp => timestamp < windowStart) + 1) + } + else { + this.store.set(key, timestamps = []) + } + + // Calculate reset time based on oldest timestamp or current time if no timestamps + const resetAtMs = timestamps[0] !== undefined + ? timestamps[0] + this.windowMs + : now + this.windowMs + + if (timestamps.length >= this.maxRequests) { + return { + success: false, + limit: this.maxRequests, + remaining: 0, + resetAtMs, + } + } + + timestamps.push(now) + + return { + success: true, + limit: this.maxRequests, + remaining: this.maxRequests - timestamps.length, + resetAtMs, + } + } + + private async blockUntilReady(key: string, timeoutMs: number): Promise> { + const deadlineAtMs = Date.now() + timeoutMs + + while (true) { + const result = await this.checkLimit(key) + + if (result.success || result.resetAtMs > deadlineAtMs) { + return result + } + + await new Promise(resolve => setTimeout(resolve, result.resetAtMs - Date.now())) + } + } +} From abdb1a689c1da3b544e47320fd5a02e1e00b0d40 Mon Sep 17 00:00:00 2001 From: unnoq Date: Wed, 5 Nov 2025 16:10:47 +0700 Subject: [PATCH 09/23] improve --- .../ratelimit/src/adapters/ioredis.test.ts | 34 ++++++++- packages/ratelimit/src/adapters/ioredis.ts | 76 +++++++++---------- 2 files changed, 67 insertions(+), 43 deletions(-) diff --git a/packages/ratelimit/src/adapters/ioredis.test.ts b/packages/ratelimit/src/adapters/ioredis.test.ts index 6e4f9e2b1..c83226318 100644 --- a/packages/ratelimit/src/adapters/ioredis.test.ts +++ b/packages/ratelimit/src/adapters/ioredis.test.ts @@ -78,6 +78,37 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) }) + it('handles concurrent requests correctly', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 3, + windowMs: 5000, + }) + + const test = async (key: string) => { + const results = await Promise.all( + Array.from({ length: 5 }, () => ratelimiter.limit(key)), + ) + + // Count successful and failed requests + const successful = results.filter(r => r.success).length + const failed = results.filter(r => !r.success).length + + // Should have exactly maxRequests successful requests + expect(successful).toBe(3) + expect(failed).toBe(2) + + // Verify remaining counts are consistent + const successfulResults = results.filter(r => r.success) + expect(successfulResults[0]!.remaining).toBe(2) + expect(successfulResults[1]!.remaining).toBe(1) + expect(successfulResults[2]!.remaining).toBe(0) + } + + await Promise.all( + Array.from({ length: 5 }, (_, i) => test(`test${i}`)), + ) + }) + describe('with blocking', () => { it('blocks until the rate limit resets', async () => { const ratelimiter = createTestingRatelimiter({ @@ -96,6 +127,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, const endTime = Date.now() expect(result2.success).toBe(true) + expect(endTime - startTime).toBeGreaterThanOrEqual(500) // actually waited expect(endTime - startTime).toBeLessThanOrEqual(2000) }) @@ -155,7 +187,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, await vi.waitFor(async () => { const keysAfterExpiry = await redis.keys(`${prefix}${key}`) expect(keysAfterExpiry).toHaveLength(0) - }, { timeout: 10_000 }) + }, { timeout: 20_000, interval: 1000 }) }) it('handles Redis errors gracefully', async () => { diff --git a/packages/ratelimit/src/adapters/ioredis.ts b/packages/ratelimit/src/adapters/ioredis.ts index 127ab8a15..604cbf570 100644 --- a/packages/ratelimit/src/adapters/ioredis.ts +++ b/packages/ratelimit/src/adapters/ioredis.ts @@ -3,9 +3,9 @@ import type { Ratelimiter, RatelimiterLimitResult } from '../types' import { fallback } from '@orpc/shared' /** - * Sliding window Lua script for Redis. + * Sliding window log Lua script for Redis. * - * This script implements atomic sliding window rate limiting using Redis sorted sets. + * This script implements atomic sliding window log rate limiting using Redis sorted sets. * It removes expired entries, checks the current count, and adds new requests atomically. * * @returns A tuple with [success, limit, remaining, resetAtMs] where: @@ -14,46 +14,38 @@ import { fallback } from '@orpc/shared' * - remaining: Number of requests remaining in the window * - resetAtMs: Unix timestamp (in milliseconds) when the window resets */ -const SLIDING_WINDOW_LUA_SCRIPT = ` - local key = KEYS[1] - local now = tonumber(ARGV[1]) - local window = tonumber(ARGV[2]) - local limit = tonumber(ARGV[3]) - - local windowStart = now - window - -- Set TTL to window + small buffer (converted to seconds for EXPIRE) - -- The buffer ensures the key doesn't expire while still in use - local ttl = math.ceil(window / 1000) + 1 - - -- Remove expired entries - redis.call('ZREMRANGEBYSCORE', key, 0, windowStart) - - -- Get all valid entries with scores in one call - -- This replaces separate ZCARD and ZRANGE operations - local entries = redis.call('ZRANGE', key, 0, -1, 'WITHSCORES') - local current = #entries / 2 -- Each entry has value + score - - -- Calculate reset time (when oldest entry expires) - local resetAtMs - if current > 0 then - resetAtMs = tonumber(entries[2]) + window -- entries[2] is the oldest score - else - resetAtMs = now + window - end - - -- Check if limit is exceeded - if current >= limit then - return {0, limit, 0, resetAtMs} - end - +const SLIDING_WINDOW_LOG_LUA_SCRIPT = ` +local key = KEYS[1] +local now_ms = tonumber(ARGV[1]) +local window_ms = tonumber(ARGV[2]) +local limit = tonumber(ARGV[3]) + +local window_start_ms = now_ms - window_ms + +-- Remove old entries outside the current window +redis.call('ZREMRANGEBYSCORE', key, '-inf', window_start_ms) + +-- Count requests in the current window +local current_count = redis.call('ZCARD', key) + +-- Calculate reset time (end of current window) +local oldest_entry = redis.call('ZRANGE', key, 0, 0, 'WITHSCORES') +local reset_at +if #oldest_entry > 0 then + reset_at = tonumber(oldest_entry[2]) + window_ms +else + reset_at = now_ms + window_ms +end + +if current_count < limit then -- Add current request - redis.call('ZADD', key, now, now) - redis.call('EXPIRE', key, ttl) - - -- Calculate remaining requests - local remaining = limit - current - 1 - - return {1, limit, remaining, resetAtMs} + redis.call('ZADD', key, now_ms, now_ms .. ':' .. math.random()) + redis.call('PEXPIRE', key, window_ms) + + return {1, limit, limit - current_count - 1, reset_at} +else + return {0, limit, 0, reset_at} +end ` export class IORedisRatelimiterError extends Error {} @@ -114,7 +106,7 @@ export class IORedisRatelimiter implements Ratelimiter { private async checkLimit(key: string) { const result = await this.redis.eval( - SLIDING_WINDOW_LUA_SCRIPT, + SLIDING_WINDOW_LOG_LUA_SCRIPT, 1, key, Date.now().toString(), From f40724780cb4347f4eeac2f3ce9ea4a4d95cfe31 Mon Sep 17 00:00:00 2001 From: unnoq Date: Wed, 5 Nov 2025 16:16:02 +0700 Subject: [PATCH 10/23] improve --- packages/ratelimit/package.json | 20 ++++++--------- .../{ioredis.test.ts => redis.test.ts} | 25 +++++++++---------- .../src/adapters/{ioredis.ts => redis.ts} | 9 ++++--- 3 files changed, 25 insertions(+), 29 deletions(-) rename packages/ratelimit/src/adapters/{ioredis.test.ts => redis.test.ts} (93%) rename packages/ratelimit/src/adapters/{ioredis.ts => redis.ts} (94%) diff --git a/packages/ratelimit/package.json b/packages/ratelimit/package.json index ed4d86fc6..2e58ae1f4 100644 --- a/packages/ratelimit/package.json +++ b/packages/ratelimit/package.json @@ -10,8 +10,8 @@ "directory": "packages/ratelimit" }, "keywords": [ - "unnoq", - "orpc" + "orpc", + "ratelimit" ], "publishConfig": { "exports": { @@ -25,10 +25,10 @@ "import": "./dist/adapters/memory.mjs", "default": "./dist/adapters/memory.mjs" }, - "./ioredis": { - "types": "./dist/adapters/ioredis.d.mts", - "import": "./dist/adapters/ioredis.mjs", - "default": "./dist/adapters/ioredis.mjs" + "./redis": { + "types": "./dist/adapters/redis.d.mts", + "import": "./dist/adapters/redis.mjs", + "default": "./dist/adapters/redis.mjs" }, "./upstash-ratelimit": { "types": "./dist/adapters/upstash-ratelimit.d.mts", @@ -40,7 +40,7 @@ "exports": { ".": "./src/index.ts", "./memory": "./src/adapters/memory.ts", - "./ioredis": "./src/adapters/ioredis.ts", + "./redis": "./src/adapters/redis.ts", "./upstash-ratelimit": "./src/adapters/upstash-ratelimit.ts" }, "files": [ @@ -52,15 +52,11 @@ "type:check": "tsc -b" }, "peerDependencies": { - "@upstash/ratelimit": ">=2.0.7", - "ioredis": ">=5.8.1" + "@upstash/ratelimit": ">=2.0.7" }, "peerDependenciesMeta": { "@upstash/ratelimit": { "optional": true - }, - "ioredis": { - "optional": true } }, "dependencies": { diff --git a/packages/ratelimit/src/adapters/ioredis.test.ts b/packages/ratelimit/src/adapters/redis.test.ts similarity index 93% rename from packages/ratelimit/src/adapters/ioredis.test.ts rename to packages/ratelimit/src/adapters/redis.test.ts index c83226318..b0c347eda 100644 --- a/packages/ratelimit/src/adapters/ioredis.test.ts +++ b/packages/ratelimit/src/adapters/redis.test.ts @@ -1,6 +1,6 @@ -import type { IORedisRatelimiterOptions } from './ioredis' +import type { IORedisRatelimiterOptions } from './redis' import { Redis } from 'ioredis' -import { IORedisRatelimiter } from './ioredis' +import { IORedisRatelimiter } from './redis' const REDIS_URL = process.env.REDIS_URL @@ -12,7 +12,8 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, let redis: Redis function createTestingRatelimiter(options: Partial = {}) { - const ratelimiter = new IORedisRatelimiter(redis, { + const ratelimiter = new IORedisRatelimiter({ + eval: redis.eval.bind(redis), prefix: `test:${crypto.randomUUID()}:`, // isolated from other tests maxRequests: 10, windowMs: 60000, @@ -191,22 +192,20 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) it('handles Redis errors gracefully', async () => { - const mockRedis = { - ...redis, + const ratelimiter = new IORedisRatelimiter({ eval: async () => { throw new Error('Redis error') }, - } as any - - const ratelimiter = new IORedisRatelimiter(mockRedis, { maxRequests: 10, windowMs: 60000 }) + maxRequests: 10, + windowMs: 60000, + }) await expect(ratelimiter.limit('some-key')).rejects.toThrow('Redis error') }) it('handles invalid script response', async () => { - const mockRedis = { - ...redis, + const ratelimiter = new IORedisRatelimiter({ eval: async () => [1, 2, 3], // Invalid response, should have 4 elements - } as any - - const ratelimiter = new IORedisRatelimiter(mockRedis, { maxRequests: 10, windowMs: 60000 }) + maxRequests: 10, + windowMs: 60000, + }) await expect(ratelimiter.limit('some-key')).rejects.toThrow('Invalid response from rate limit script') }) }) diff --git a/packages/ratelimit/src/adapters/ioredis.ts b/packages/ratelimit/src/adapters/redis.ts similarity index 94% rename from packages/ratelimit/src/adapters/ioredis.ts rename to packages/ratelimit/src/adapters/redis.ts index 604cbf570..69772f877 100644 --- a/packages/ratelimit/src/adapters/ioredis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -1,4 +1,3 @@ -import type Redis from 'ioredis' import type { Ratelimiter, RatelimiterLimitResult } from '../types' import { fallback } from '@orpc/shared' @@ -51,6 +50,8 @@ end export class IORedisRatelimiterError extends Error {} export interface IORedisRatelimiterOptions { + eval: (script: string, numKeys: number, ...args: string[]) => Promise + /** * Block until the request may pass or timeout is reached. */ @@ -75,19 +76,19 @@ export interface IORedisRatelimiterOptions { * The duration of the sliding window in milliseconds. */ windowMs: number - } export class IORedisRatelimiter implements Ratelimiter { + private readonly eval: IORedisRatelimiterOptions['eval'] private readonly prefix: string private readonly maxRequests: number private readonly windowMs: number private readonly blockingUntilReady: IORedisRatelimiterOptions['blockingUntilReady'] constructor( - private readonly redis: Redis, options: IORedisRatelimiterOptions, ) { + this.eval = options.eval this.prefix = fallback(options.prefix, 'orpc:ratelimit:') this.maxRequests = options.maxRequests this.windowMs = options.windowMs @@ -105,7 +106,7 @@ export class IORedisRatelimiter implements Ratelimiter { } private async checkLimit(key: string) { - const result = await this.redis.eval( + const result = await this.eval( SLIDING_WINDOW_LOG_LUA_SCRIPT, 1, key, From d500f2b6b1998dc1424bedacaaa2cdaecf4a4677 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 15:24:31 +0700 Subject: [PATCH 11/23] improve --- .../ratelimit/src/adapters/memory.test.ts | 22 ++++++------- packages/ratelimit/src/adapters/memory.ts | 30 ++++++++--------- packages/ratelimit/src/adapters/redis.test.ts | 32 +++++++++---------- packages/ratelimit/src/adapters/redis.ts | 24 +++++++------- .../src/adapters/upstash-ratelimit.test.ts | 2 +- .../src/adapters/upstash-ratelimit.ts | 11 ++----- packages/ratelimit/src/types.ts | 2 +- 7 files changed, 58 insertions(+), 65 deletions(-) diff --git a/packages/ratelimit/src/adapters/memory.test.ts b/packages/ratelimit/src/adapters/memory.test.ts index 00ec6b789..3cec26315 100644 --- a/packages/ratelimit/src/adapters/memory.test.ts +++ b/packages/ratelimit/src/adapters/memory.test.ts @@ -6,7 +6,7 @@ describe('memoryRatelimiter', () => { function createTestingRatelimiter(options: Partial[0]> = {}) { return new MemoryRatelimiter({ maxRequests: 2, - windowMs: 1000, + window: 1000, ...options, }) } @@ -19,24 +19,24 @@ describe('memoryRatelimiter', () => { expect(result1.success).toBe(true) expect(result1.remaining).toBe(1) expect(result1.limit).toBe(2) - expect(result1.resetAtMs).toBeGreaterThan(Date.now()) + expect(result1.reset).toBeGreaterThan(Date.now()) const result2 = await ratelimiter.limit('test') expect(result2.success).toBe(true) expect(result2.remaining).toBe(0) expect(result2.limit).toBe(2) - expect(result2.resetAtMs).toEqual(result1.resetAtMs) + expect(result2.reset).toEqual(result1.reset) const result3 = await ratelimiter.limit('test') expect(result3.success).toBe(false) expect(result3.remaining).toBe(0) expect(result3.limit).toBe(2) - expect(result3.resetAtMs).toEqual(result1.resetAtMs) + expect(result3.reset).toEqual(result1.reset) }) it('should reset after window expires', async () => { const ratelimiter = createTestingRatelimiter({ - windowMs: 200, + window: 200, }) const result1 = await ratelimiter.limit('test') @@ -66,10 +66,10 @@ describe('memoryRatelimiter', () => { it('should block until ready', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 1, - windowMs: 1000, + window: 1000, blockingUntilReady: { enabled: true, - timeoutMs: 2000, + timeout: 2000, }, }) @@ -87,10 +87,10 @@ describe('memoryRatelimiter', () => { it('should respect timeout', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 1, - windowMs: 2000, + window: 2000, blockingUntilReady: { enabled: true, - timeoutMs: 1000, + timeout: 1000, }, }) @@ -105,7 +105,7 @@ describe('memoryRatelimiter', () => { it('should handle concurrent requests correctly', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 3, - windowMs: 1000, + window: 1000, }) const test = async (key: string, request: number) => { @@ -137,7 +137,7 @@ describe('memoryRatelimiter', () => { it('should cleanup expired entries on next limit call', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 2, - windowMs: 200, + window: 200, }) await ratelimiter.limit('test') diff --git a/packages/ratelimit/src/adapters/memory.ts b/packages/ratelimit/src/adapters/memory.ts index 699ff0754..863012d25 100644 --- a/packages/ratelimit/src/adapters/memory.ts +++ b/packages/ratelimit/src/adapters/memory.ts @@ -6,7 +6,7 @@ export interface MemoryRatelimiterOptions { */ blockingUntilReady?: { enabled: boolean - timeoutMs: number + timeout: number } /** @@ -17,20 +17,20 @@ export interface MemoryRatelimiterOptions { /** * The duration of the sliding window in milliseconds. */ - windowMs: number + window: number } export class MemoryRatelimiter implements Ratelimiter { private readonly maxRequests: number - private readonly windowMs: number + private readonly window: number private readonly blockingUntilReady: MemoryRatelimiterOptions['blockingUntilReady'] private readonly store: Map private lastCleanupTime: number | null = null constructor(options: MemoryRatelimiterOptions) { this.maxRequests = options.maxRequests - this.windowMs = options.windowMs + this.window = options.window this.blockingUntilReady = options.blockingUntilReady this.store = new Map() } @@ -39,7 +39,7 @@ export class MemoryRatelimiter implements Ratelimiter { this.cleanup() if (this.blockingUntilReady?.enabled) { - return this.blockUntilReady(key, this.blockingUntilReady.timeoutMs) + return this.blockUntilReady(key, this.blockingUntilReady.timeout) } return this.checkLimit(key) @@ -49,12 +49,12 @@ export class MemoryRatelimiter implements Ratelimiter { const now = Date.now() // Only clean up once per window to avoid excessive processing - if (this.lastCleanupTime !== null && this.lastCleanupTime + this.windowMs > now) { + if (this.lastCleanupTime !== null && this.lastCleanupTime + this.window > now) { return } this.lastCleanupTime = now - const windowStart = now - this.windowMs + const windowStart = now - this.window for (const [key, timestamps] of this.store) { // remove expired timestamps @@ -68,7 +68,7 @@ export class MemoryRatelimiter implements Ratelimiter { private async checkLimit(key: string): Promise> { const now = Date.now() - const windowStart = now - this.windowMs + const windowStart = now - this.window let timestamps = this.store.get(key) if (timestamps) { @@ -80,16 +80,16 @@ export class MemoryRatelimiter implements Ratelimiter { } // Calculate reset time based on oldest timestamp or current time if no timestamps - const resetAtMs = timestamps[0] !== undefined - ? timestamps[0] + this.windowMs - : now + this.windowMs + const reset = timestamps[0] !== undefined + ? timestamps[0] + this.window + : now + this.window if (timestamps.length >= this.maxRequests) { return { success: false, limit: this.maxRequests, remaining: 0, - resetAtMs, + reset, } } @@ -99,7 +99,7 @@ export class MemoryRatelimiter implements Ratelimiter { success: true, limit: this.maxRequests, remaining: this.maxRequests - timestamps.length, - resetAtMs, + reset, } } @@ -109,11 +109,11 @@ export class MemoryRatelimiter implements Ratelimiter { while (true) { const result = await this.checkLimit(key) - if (result.success || result.resetAtMs > deadlineAtMs) { + if (result.success || result.reset > deadlineAtMs) { return result } - await new Promise(resolve => setTimeout(resolve, result.resetAtMs - Date.now())) + await new Promise(resolve => setTimeout(resolve, result.reset - Date.now())) } } } diff --git a/packages/ratelimit/src/adapters/redis.test.ts b/packages/ratelimit/src/adapters/redis.test.ts index b0c347eda..4d163ae36 100644 --- a/packages/ratelimit/src/adapters/redis.test.ts +++ b/packages/ratelimit/src/adapters/redis.test.ts @@ -16,7 +16,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, eval: redis.eval.bind(redis), prefix: `test:${crypto.randomUUID()}:`, // isolated from other tests maxRequests: 10, - windowMs: 60000, + window: 60000, ...options, }) return ratelimiter @@ -33,24 +33,24 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, describe('without blocking', () => { it('allows requests within the limit', async () => { - const ratelimiter = createTestingRatelimiter({ maxRequests: 2, windowMs: 5000 }) + const ratelimiter = createTestingRatelimiter({ maxRequests: 2, window: 5000 }) const key = 'user1' const result1 = await ratelimiter.limit(key) expect(result1.success).toBe(true) expect(result1.limit).toBe(2) expect(result1.remaining).toBe(1) - expect(result1.resetAtMs).toBeGreaterThan(Date.now()) + expect(result1.reset).toBeGreaterThan(Date.now()) const result2 = await ratelimiter.limit(key) expect(result2.success).toBe(true) expect(result2.limit).toBe(2) expect(result2.remaining).toBe(0) - expect(result2.resetAtMs).toEqual(result1.resetAtMs) + expect(result2.reset).toEqual(result1.reset) }) it('denies requests exceeding the limit', async () => { - const ratelimiter = createTestingRatelimiter({ maxRequests: 1, windowMs: 5000 }) + const ratelimiter = createTestingRatelimiter({ maxRequests: 1, window: 5000 }) const key = 'user2' // reach the limit @@ -62,7 +62,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) it('resets the limit after the window expires', async () => { - const ratelimiter = createTestingRatelimiter({ maxRequests: 1, windowMs: 2000 }) + const ratelimiter = createTestingRatelimiter({ maxRequests: 1, window: 2000 }) const key = 'user3' // reach the limit @@ -82,7 +82,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, it('handles concurrent requests correctly', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 3, - windowMs: 5000, + window: 5000, }) const test = async (key: string) => { @@ -114,8 +114,8 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, it('blocks until the rate limit resets', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 1, - windowMs: 2000, - blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + window: 2000, + blockingUntilReady: { enabled: true, timeout: 2000 }, }) const key = 'user-blocking-1' @@ -135,8 +135,8 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, it('times out if the reset time is beyond the timeout', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 1, - windowMs: 60_000, - blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + window: 60_000, + blockingUntilReady: { enabled: true, timeout: 2000 }, }) const key = 'user-blocking-2' @@ -153,8 +153,8 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, it('handles concurrent blocking requests correctly', async () => { const ratelimiter = createTestingRatelimiter({ maxRequests: 2, - windowMs: 1000, - blockingUntilReady: { enabled: true, timeoutMs: 2000 }, + window: 1000, + blockingUntilReady: { enabled: true, timeout: 2000 }, }) const key = 'user-concurrent-blocking' @@ -173,7 +173,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, it('uses the correct prefix for keys & auto expiration', async () => { const prefix = `custom-prefix:${crypto.randomUUID()}:` const ratelimiter = createTestingRatelimiter({ - windowMs: 1000, + window: 1000, prefix, maxRequests: 1, }) @@ -195,7 +195,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, const ratelimiter = new IORedisRatelimiter({ eval: async () => { throw new Error('Redis error') }, maxRequests: 10, - windowMs: 60000, + window: 60000, }) await expect(ratelimiter.limit('some-key')).rejects.toThrow('Redis error') }) @@ -204,7 +204,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, const ratelimiter = new IORedisRatelimiter({ eval: async () => [1, 2, 3], // Invalid response, should have 4 elements maxRequests: 10, - windowMs: 60000, + window: 60000, }) await expect(ratelimiter.limit('some-key')).rejects.toThrow('Invalid response from rate limit script') }) diff --git a/packages/ratelimit/src/adapters/redis.ts b/packages/ratelimit/src/adapters/redis.ts index 69772f877..bc68974d7 100644 --- a/packages/ratelimit/src/adapters/redis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -47,8 +47,6 @@ else end ` -export class IORedisRatelimiterError extends Error {} - export interface IORedisRatelimiterOptions { eval: (script: string, numKeys: number, ...args: string[]) => Promise @@ -57,7 +55,7 @@ export interface IORedisRatelimiterOptions { */ blockingUntilReady?: { enabled: boolean - timeoutMs: number + timeout: number } /** @@ -75,14 +73,14 @@ export interface IORedisRatelimiterOptions { /** * The duration of the sliding window in milliseconds. */ - windowMs: number + window: number } export class IORedisRatelimiter implements Ratelimiter { private readonly eval: IORedisRatelimiterOptions['eval'] private readonly prefix: string private readonly maxRequests: number - private readonly windowMs: number + private readonly window: number private readonly blockingUntilReady: IORedisRatelimiterOptions['blockingUntilReady'] constructor( @@ -91,7 +89,7 @@ export class IORedisRatelimiter implements Ratelimiter { this.eval = options.eval this.prefix = fallback(options.prefix, 'orpc:ratelimit:') this.maxRequests = options.maxRequests - this.windowMs = options.windowMs + this.window = options.window this.blockingUntilReady = options.blockingUntilReady } @@ -99,7 +97,7 @@ export class IORedisRatelimiter implements Ratelimiter { const prefixedKey = `${this.prefix}${key}` if (this.blockingUntilReady?.enabled) { - return await this.blockUntilReady(prefixedKey, this.blockingUntilReady.timeoutMs) + return await this.blockUntilReady(prefixedKey, this.blockingUntilReady.timeout) } return await this.checkLimit(prefixedKey) @@ -111,21 +109,21 @@ export class IORedisRatelimiter implements Ratelimiter { 1, key, Date.now().toString(), - this.windowMs.toString(), + this.window.toString(), this.maxRequests.toString(), ) as unknown if (!Array.isArray(result) || result.length !== 4) { - throw new IORedisRatelimiterError('Invalid response from rate limit script') + throw new TypeError('Invalid response from rate limit script') } - const [success, limit, remaining, resetAtMs] = result as [number, number, number, number] + const [success, limit, remaining, reset] = result as [number, number, number, number] return { success: success === 1, limit, remaining, - resetAtMs, + reset, } } @@ -135,11 +133,11 @@ export class IORedisRatelimiter implements Ratelimiter { while (true) { const result = await this.checkLimit(key) - if (result.success || result.resetAtMs > deadlineAtMs) { + if (result.success || result.reset > deadlineAtMs) { return result } - await new Promise(resolve => setTimeout(resolve, result.resetAtMs - Date.now())) + await new Promise(resolve => setTimeout(resolve, result.reset - Date.now())) } } } diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts index ccb44780a..602d37218 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts @@ -47,7 +47,7 @@ describe.concurrent( const ratelimiter = createTestingRatelimiter({ blockingUntilReady: { enabled: true, - timeoutMs, + timeout: timeoutMs, }, }) const key = `test-blocking-${crypto.randomUUID()}` diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 7a8b0cb3a..5251ffcea 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -7,7 +7,7 @@ export interface UpstashRatelimiterOptions { */ blockingUntilReady?: { enabled: boolean - timeoutMs: number + timeout: number } /** @@ -39,15 +39,10 @@ export class UpstashRatelimiter implements Ratelimiter { async limit(key: string): Promise> { const result = this.blockingUntilReady?.enabled - ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeoutMs) + ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeout) : await this.ratelimit.limit(key) this.waitUtil?.(result.pending) - return { - success: result.success, - limit: result.limit, - remaining: result.remaining, - resetAtMs: result.reset, - } + return result } } diff --git a/packages/ratelimit/src/types.ts b/packages/ratelimit/src/types.ts index 93d446bb1..7401f4467 100644 --- a/packages/ratelimit/src/types.ts +++ b/packages/ratelimit/src/types.ts @@ -14,7 +14,7 @@ export interface RatelimiterLimitResult { /** * Unix timestamp in milliseconds when the limits are reset. */ - resetAtMs?: number + reset?: number } export interface Ratelimiter { From 382abb3e3ac7fdbd6b56c948825ca10de4897784 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 17:03:48 +0700 Subject: [PATCH 12/23] wip --- packages/ratelimit/package.json | 1 + packages/ratelimit/src/handler-plugin.ts | 56 +++++++++++++++ packages/ratelimit/src/index.test.ts | 4 +- packages/ratelimit/src/index.ts | 2 + packages/ratelimit/src/middleware.ts | 90 ++++++++++++++++++++++++ pnpm-lock.yaml | 3 + 6 files changed, 154 insertions(+), 2 deletions(-) create mode 100644 packages/ratelimit/src/handler-plugin.ts create mode 100644 packages/ratelimit/src/middleware.ts diff --git a/packages/ratelimit/package.json b/packages/ratelimit/package.json index 2e58ae1f4..fe08dba59 100644 --- a/packages/ratelimit/package.json +++ b/packages/ratelimit/package.json @@ -61,6 +61,7 @@ }, "dependencies": { "@orpc/client": "workspace:*", + "@orpc/server": "workspace:*", "@orpc/shared": "workspace:*", "@orpc/standard-server": "workspace:*" }, diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts new file mode 100644 index 000000000..efedb1c2e --- /dev/null +++ b/packages/ratelimit/src/handler-plugin.ts @@ -0,0 +1,56 @@ +import type { Context } from '@orpc/server' +import type { StandardHandlerOptions, StandardHandlerPlugin } from '@orpc/server/standard' +import type { RatelimiterLimitResult } from './types' + +export const RATELIMIT_HANDLER_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_HANDLER_CONTEXT') + +export interface RatelimitHandlerPluginContext { + /** + * The result of the ratelimiter after applying limits + */ + ratelimitResult?: RatelimiterLimitResult +} + +export class RatelimitHandlerPlugin implements StandardHandlerPlugin { + /** + * this plugin should lower priority than response headers plugin, + * if user want override rate limit headers + */ + order = 100_000 + + init(options: StandardHandlerOptions): void { + options.rootInterceptors ??= [] + + options.rootInterceptors.push(async (interceptorOptions) => { + const handlerContext: RatelimitHandlerPluginContext = {} + + const result = await interceptorOptions.next({ + ...interceptorOptions, + context: { + ...interceptorOptions.context, + [RATELIMIT_HANDLER_CONTEXT_SYMBOL]: handlerContext, + }, + }) + + if (result.matched && handlerContext.ratelimitResult) { + return { + ...result, + response: { + ...result.response, + headers: { + ...result.response.headers, + 'rateLimit-limit': handlerContext.ratelimitResult.limit?.toString(), + 'rateLimit-remaining': handlerContext.ratelimitResult.remaining?.toString(), + 'rateLimit-reset': handlerContext.ratelimitResult.reset?.toString(), + 'retry-after': !handlerContext.ratelimitResult.success && result.response.status === 429 && handlerContext.ratelimitResult.reset !== undefined + ? Math.ceil((handlerContext.ratelimitResult.reset - Date.now()) / 1000).toString() + : undefined, + }, + }, + } + } + + return result + }) + } +} diff --git a/packages/ratelimit/src/index.test.ts b/packages/ratelimit/src/index.test.ts index 8297b9d7e..8a4926dd3 100644 --- a/packages/ratelimit/src/index.test.ts +++ b/packages/ratelimit/src/index.test.ts @@ -1,3 +1,3 @@ -it('exports Publisher', async () => { - expect(Object.keys(await import('./index'))).toContain('Publisher') +it('exports createRatelimitMiddleware', async () => { + expect(Object.keys(await import('./index'))).toContain('createRatelimitMiddleware') }) diff --git a/packages/ratelimit/src/index.ts b/packages/ratelimit/src/index.ts index c9f6f047d..3c5daeee8 100644 --- a/packages/ratelimit/src/index.ts +++ b/packages/ratelimit/src/index.ts @@ -1 +1,3 @@ +export * from './handler-plugin' +export * from './middleware' export * from './types' diff --git a/packages/ratelimit/src/middleware.ts b/packages/ratelimit/src/middleware.ts new file mode 100644 index 000000000..aceb37259 --- /dev/null +++ b/packages/ratelimit/src/middleware.ts @@ -0,0 +1,90 @@ +import type { Context, Meta, Middleware, MiddlewareOptions } from '@orpc/server' +import type { Promisable, Value } from '@orpc/shared' +import type { RatelimitHandlerPluginContext } from './handler-plugin' +import type { Ratelimiter } from './types' +import { ORPCError } from '@orpc/server' +import { toArray, value } from '@orpc/shared' +import { RATELIMIT_HANDLER_CONTEXT_SYMBOL } from './handler-plugin' + +const RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_MIDDLEWARE_CONTEXT') + +export interface RatelimiterMiddlewareContext { + /** + * The applied limits in this request, mainly for deduplication purposes + */ + limits: { limiter: Ratelimiter, key: string }[] +} + +export interface CreateRatelimitMiddlewareOptions< + TInContext extends Context, + TInput, + TMeta extends Meta, +> { + /** + * The rule set to use for rate limiting + */ + limiter: Value, [middlewareOptions: MiddlewareOptions, TMeta>, input: TInput]> + + /** + * The key to identify the user/requester + */ + key: Value, [middlewareOptions: MiddlewareOptions, TMeta>, input: TInput]> + /** + * If you ratelimit middleware is used multiple times + * or you invoke a procedure inside another procedure (shared the same context) that also has + * ratelimit middleware **with the same limiter and key**, this option + * will ensure that the limit is only applied once per request. + * + * @default true + */ + dedupe?: boolean +} + +export function createRatelimitMiddleware< + TInContext extends Context, + TInput, + TMeta extends Meta, +>( + { dedupe = true, ...options }: CreateRatelimitMiddlewareOptions, +): Middleware, TInput, any, any, TMeta> { + return async function ratelimit(middlewareOptions, input) { + const [limiter, key] = await Promise.all([ + value(options.limiter, middlewareOptions, input), + value(options.key, middlewareOptions, input), + ]) + + const middlewareContext: RatelimiterMiddlewareContext | undefined = middlewareOptions.context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] + if (dedupe && middlewareContext?.limits.some(l => l.key === key && l.limiter === limiter)) { + return middlewareOptions.next() + } + + const result = await limiter.limit(key) + + const pluginContext: RatelimitHandlerPluginContext | undefined = middlewareOptions.context[RATELIMIT_HANDLER_CONTEXT_SYMBOL] + if (pluginContext) { + pluginContext.ratelimitResult = result + } + + if (!result.success) { + throw new ORPCError('TOO_MANY_REQUESTS', { + data: { + limit: result.limit, + remaining: result.remaining, + reset: result.reset, + }, + }) + } + + return middlewareOptions.next({ + context: { + [RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL]: { + ...middlewareOptions, + limits: [ + ...toArray(middlewareContext?.limits), + { limiter, key }, + ], + }, + }, + }) + } +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 6ede7c5df..f057cddd6 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -558,6 +558,9 @@ importers: '@orpc/client': specifier: workspace:* version: link:../client + '@orpc/server': + specifier: workspace:* + version: link:../server '@orpc/shared': specifier: workspace:* version: link:../shared From 069ff2bd030fc3281007eb30ed075223a7c8ee7e Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 17:32:20 +0700 Subject: [PATCH 13/23] handler plugin tests --- packages/ratelimit/src/handler-plugin.test.ts | 184 ++++++++++++++++++ packages/ratelimit/src/handler-plugin.ts | 6 +- 2 files changed, 187 insertions(+), 3 deletions(-) create mode 100644 packages/ratelimit/src/handler-plugin.test.ts diff --git a/packages/ratelimit/src/handler-plugin.test.ts b/packages/ratelimit/src/handler-plugin.test.ts new file mode 100644 index 000000000..b218b631a --- /dev/null +++ b/packages/ratelimit/src/handler-plugin.test.ts @@ -0,0 +1,184 @@ +import { StandardHandler, StandardRPCMatcher } from '@orpc/server/standard' +import { describe, expect, it } from 'vitest' +import { RATELIMIT_HANDLER_CONTEXT_SYMBOL, RatelimitHandlerPlugin } from './handler-plugin' + +describe('ratelimitHandlerPlugin', () => { + const createMockRequest = (url: string) => ({ + method: 'GET', + url: new URL(url), + headers: {}, + signal: new AbortController().signal, + body: () => Promise.resolve(''), + }) + + it('adds rate limit headers', async () => { + const resetTime = Date.now() + 60000 + + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { + limit: 100, + remaining: 50, + reset: resetTime, + } + + return { + matched: true as const, + response: { status: 200, headers: {}, body: 'ok' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + expect(result.response.headers['ratelimit-limit']).toBe('100') + expect(result.response.headers['ratelimit-remaining']).toBe('50') + expect(result.response.headers['ratelimit-reset']).toBe(resetTime.toString()) + }) + + it('adds retry-after header on 429 status', async () => { + const resetTime = Date.now() + 30000 + + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { + success: false, + reset: resetTime, + } + + return { + matched: true as const, + response: { status: 429, headers: {}, body: 'too many requests' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + const retryAfter = Number.parseInt(result.response.headers['retry-after'] as string) + expect(retryAfter).toBeGreaterThan(0) + expect(retryAfter).toBeLessThanOrEqual(31) + }) + + it('does not add retry-after when success is true', async () => { + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { + success: true, + reset: Date.now() + 60000, + } + + return { + matched: true as const, + response: { status: 429, headers: {}, body: 'ok' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + expect(result.response.headers['retry-after']).toBeUndefined() + }) + + it('does not add retry-after when status is not 429', async () => { + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { + success: false, + reset: Date.now() + 60000, + } + + return { + matched: true as const, + response: { status: 200, headers: {}, body: 'ok' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + expect(result.response.headers['retry-after']).toBeUndefined() + }) + + it('does not add headers when ratelimitResult is undefined', async () => { + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async () => ({ + matched: true as const, + response: { status: 200, headers: {}, body: 'ok' }, + })) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + expect(result.response.headers['ratelimit-limit']).toBeUndefined() + expect(result.response.headers['ratelimit-remaining']).toBeUndefined() + }) + + it('handles partial ratelimitResult', async () => { + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { limit: 100 } + + return { + matched: true as const, + response: { status: 200, headers: {}, body: 'ok' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) + throw new Error('request should match') + + expect(result.response.headers['ratelimit-limit']).toBe('100') + expect(result.response.headers['ratelimit-remaining']).toBeUndefined() + }) + + it('injects context with symbol', async () => { + let capturedContext: any + + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + capturedContext = context + return { + matched: true as const, + response: { status: 200, headers: {}, body: 'ok' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + expect(capturedContext[RATELIMIT_HANDLER_CONTEXT_SYMBOL]).toBeDefined() + }) +}) diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts index efedb1c2e..1977c9f06 100644 --- a/packages/ratelimit/src/handler-plugin.ts +++ b/packages/ratelimit/src/handler-plugin.ts @@ -39,9 +39,9 @@ export class RatelimitHandlerPlugin implements StandardHandle ...result.response, headers: { ...result.response.headers, - 'rateLimit-limit': handlerContext.ratelimitResult.limit?.toString(), - 'rateLimit-remaining': handlerContext.ratelimitResult.remaining?.toString(), - 'rateLimit-reset': handlerContext.ratelimitResult.reset?.toString(), + 'ratelimit-limit': handlerContext.ratelimitResult.limit?.toString(), + 'ratelimit-remaining': handlerContext.ratelimitResult.remaining?.toString(), + 'ratelimit-reset': handlerContext.ratelimitResult.reset?.toString(), 'retry-after': !handlerContext.ratelimitResult.success && result.response.status === 429 && handlerContext.ratelimitResult.reset !== undefined ? Math.ceil((handlerContext.ratelimitResult.reset - Date.now()) / 1000).toString() : undefined, From 211d1932c5c3d4fc442d9018502cb8c34fd0bc62 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 19:33:36 +0700 Subject: [PATCH 14/23] tests middleware --- packages/ratelimit/src/middleware.test-d.ts | 45 +++ packages/ratelimit/src/middleware.test.ts | 295 ++++++++++++++++++++ packages/ratelimit/src/middleware.ts | 6 +- 3 files changed, 343 insertions(+), 3 deletions(-) create mode 100644 packages/ratelimit/src/middleware.test-d.ts create mode 100644 packages/ratelimit/src/middleware.test.ts diff --git a/packages/ratelimit/src/middleware.test-d.ts b/packages/ratelimit/src/middleware.test-d.ts new file mode 100644 index 000000000..b41ba4bfe --- /dev/null +++ b/packages/ratelimit/src/middleware.test-d.ts @@ -0,0 +1,45 @@ +import type { Ratelimiter } from './types' +import { os, type } from '@orpc/server' +import { createRatelimitMiddleware } from './middleware' + +describe('createRatelimitMiddleware', () => { + it('can infer context & input & meta types', () => { + const procedure = os + .$context<{ userId: string, ratelimiter: Ratelimiter }>() + .$meta<{ meta?: string }>({}) + .input(type<{ amount: number }>()) + .use(({ next }) => { + return next({ + context: { + db: 'postgres', + }, + }) + }) + .use( + createRatelimitMiddleware({ + limiter: async ({ context, procedure }, input) => { + expectTypeOf(input.amount).toBeNumber() + expectTypeOf(context.userId).toBeString() + expectTypeOf(context.db).toBeString() + expectTypeOf(procedure['~orpc'].meta.meta).toEqualTypeOf() + return context.ratelimiter + }, + key: ({ context, procedure }, input) => { + expectTypeOf(input.amount).toBeNumber() + expectTypeOf(context.userId).toBeString() + expectTypeOf(context.db).toBeString() + expectTypeOf(procedure['~orpc'].meta.meta).toEqualTypeOf() + return context.userId + }, + }), + ) + .handler(({ context, input, procedure }) => { + expectTypeOf(context.ratelimiter).toEqualTypeOf() + expectTypeOf(context.userId).toBeString() + expectTypeOf(context.db).toBeString() + expectTypeOf(input.amount).toBeNumber() + expectTypeOf(procedure['~orpc'].meta.meta).toEqualTypeOf() + return 'ok' + }) + }) +}) diff --git a/packages/ratelimit/src/middleware.test.ts b/packages/ratelimit/src/middleware.test.ts new file mode 100644 index 000000000..257da157c --- /dev/null +++ b/packages/ratelimit/src/middleware.test.ts @@ -0,0 +1,295 @@ +import type { RatelimiterMiddlewareContext } from './middleware' +import type { Ratelimiter, RatelimiterLimitResult } from './types' +import { call, os, type } from '@orpc/server' +import { describe, expect, it, vi } from 'vitest' +import { RATELIMIT_HANDLER_CONTEXT_SYMBOL } from './handler-plugin' +import { createRatelimitMiddleware, RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL } from './middleware' + +describe('createRatelimitMiddleware', () => { + const createLimiter = (result: RatelimiterLimitResult): Ratelimiter => ({ + limit: vi.fn().mockResolvedValue(result), + }) + const success: RatelimiterLimitResult = { success: true, limit: 10, remaining: 5, reset: Date.now() + 60000 } + + describe('basic', () => { + it('applies rate limit successfully', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'key' }) + await expect(call(os.use(mw).handler(() => 'ok'), undefined, { context: {} })).resolves.toBe('ok') + expect(limiter.limit).toHaveBeenCalledWith('key') + }) + + it('throws TOO_MANY_REQUESTS when limit exceeded', async () => { + const reset = Date.now() + 60000 + const limiter = createLimiter({ success: false, limit: 10, remaining: 0, reset }) + await expect(call(os.use(createRatelimitMiddleware({ limiter, key: 'k' })).handler(() => 'ok'), undefined, { context: {} })) + .rejects + .toMatchObject({ code: 'TOO_MANY_REQUESTS', data: { limit: 10, remaining: 0, reset } }) + }) + + it('passes result to handler plugin context', async () => { + const limiter = createLimiter(success) + const ctx = {} + await call(os.use(createRatelimitMiddleware({ limiter, key: 'k' })).handler(() => 'ok'), undefined, { + context: { [RATELIMIT_HANDLER_CONTEXT_SYMBOL]: ctx }, + }) + expect(ctx).toHaveProperty('ratelimitResult', success) + }) + + it('handles missing plugin context', async () => { + await expect(call(os.use(createRatelimitMiddleware({ limiter: createLimiter(success), key: 'k' })).handler(() => 'ok'), undefined, { context: {} })) + .resolves + .toBe('ok') + }) + + it('handles partial result', async () => { + await expect(call(os.use(createRatelimitMiddleware({ limiter: createLimiter({ success: true }), key: 'k' })).handler(() => 'ok'), undefined, { context: {} })) + .resolves + .toBe('ok') + }) + }) + + describe('dynamic', () => { + it('supports function limiter', async () => { + const l1 = createLimiter(success) + const l2 = createLimiter(success) + const fn = vi.fn((_, input) => input === 'u1' ? l1 : l2) + + const proc = os + .use(createRatelimitMiddleware({ limiter: fn, key: 'k' })) + .handler(({ input }) => input) + + await call(proc, 'u1', { context: {} }) + expect(l1.limit).toHaveBeenCalledWith('k') + expect(l2.limit).not.toHaveBeenCalled() + + await call(proc, 'u2', { context: {} }) + expect(l2.limit).toHaveBeenCalledWith('k') + }) + + it('supports async functions', async () => { + const limiter = createLimiter(success) + const lFn = vi.fn(async () => limiter) + const kFn = vi.fn(async (_, i: string) => `user:${i}`) + + const procedure = os + .input(type()) + .use(createRatelimitMiddleware({ limiter: lFn, key: kFn })) + .handler(({ input }) => input) + + await call(procedure, 'a', { context: {} }) + expect(lFn).toHaveBeenCalledTimes(1) + expect(kFn).toHaveBeenCalledTimes(1) + expect(limiter.limit).toHaveBeenCalledWith('user:a') + }) + + it('passes middleware options', async () => { + const limiter = createLimiter(success) + const lFn = vi.fn(() => limiter) + const kFn = vi.fn(() => 'k') + const procedure = os.use(createRatelimitMiddleware({ limiter: lFn, key: kFn })).handler(({ input }) => input) + await call(procedure, 'in', { context: {} }) + expect(lFn).toHaveBeenCalledWith(expect.objectContaining({ context: {}, next: expect.any(Function) }), 'in') + expect(kFn).toHaveBeenCalledWith(expect.objectContaining({ context: {}, next: expect.any(Function) }), 'in') + }) + + it('resolves in parallel', async () => { + const limiter = createLimiter(success) + let lt = 0 + let kt = 0 + const lFn = async () => { + await new Promise(r => setTimeout(r, 10)) + lt = Date.now() + return limiter + } + const kFn = async () => { + await new Promise(r => setTimeout(r, 10)) + kt = Date.now() + return 'k' + } + await call(os.use(createRatelimitMiddleware({ limiter: lFn, key: kFn })).handler(() => 'ok'), undefined, { context: {} }) + expect(Math.abs(lt - kt)).toBeLessThan(5) + }) + }) + + describe('context', () => { + it('defines ratelimit context if does not exist', async () => { + const limiter = createLimiter(success) + let ctx: any + const proc = os + .use(createRatelimitMiddleware({ limiter, key: 'k' })) + .handler(({ context }) => { + ctx = context + return 'ok' + }) + await call(proc, undefined, { context: {} }) + expect(ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL].limits).toEqual([{ limiter, key: 'k' }]) + }) + + it('extends ratelimit context if exists', async () => { + let ctx: any + await call( + os + .use(createRatelimitMiddleware({ + limiter: createLimiter(success), + key: 'k', + })) + .handler(({ context }) => { + ctx = context + return 'ok' + }), + undefined, + { + context: { + [RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL]: { + limits: [{ limiter: createLimiter(success), key: 'k' }], + }, + userId: '1', + db: 'pg', + }, + }, + ) + expect(ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL].limits.length).toBe(2) + }) + + it('isolate ratelimit contexts', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'k', dedupe: false }) + const innerHandlerFn = vi.fn().mockReturnValue('in') + const inner = os + .use(mw) + .handler(innerHandlerFn) + const outer = os + .use(mw) + .handler(async ({ context }) => { + return `out:${await call(inner, undefined, { context })}:${await call(inner, undefined, { context })}` + }) + + await call(outer, undefined, { context: {} }) + expect(limiter.limit).toHaveBeenCalledTimes(3) + expect(innerHandlerFn).toHaveBeenCalledTimes(2) + expect( + innerHandlerFn.mock.calls[0]![0].context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL], + ).not.toBe( + innerHandlerFn.mock.calls[1]![0].context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL], + ) + }) + }) + + describe('dedupe', () => { + it('deduplicates by default', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'k' }) + await call(os.use(mw).use(mw).handler(() => 'ok'), undefined, { context: {} }) + expect(limiter.limit).toHaveBeenCalledTimes(1) + }) + + it('skips dedupe when disabled', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'k', dedupe: false }) + await call( + os.use(mw).use(mw).handler(() => 'ok'), + undefined, + { context: {} }, + ) + expect(limiter.limit).toHaveBeenCalledTimes(2) + }) + + it('dedupes only same limiter+key', async () => { + const l1 = createLimiter(success) + const l2 = createLimiter(success) + await call( + os.use(createRatelimitMiddleware({ limiter: l1, key: 'k' })) + .use(createRatelimitMiddleware({ limiter: l2, key: 'k' })) + .handler(() => 'ok'), + undefined, + { context: {} }, + ) + expect(l1.limit).toHaveBeenCalledTimes(1) + expect(l2.limit).toHaveBeenCalledTimes(1) + + const l3 = createLimiter(success) + await call( + os.use(createRatelimitMiddleware({ limiter: l3, key: 'k1' })) + .use(createRatelimitMiddleware({ limiter: l3, key: 'k2' })).handler(() => 'ok'), + undefined, + { context: {} }, + ) + expect(l3.limit).toHaveBeenCalledTimes(2) + }) + + it('accumulates limits', async () => { + const l1 = createLimiter(success) + const l2 = createLimiter(success) + let ctx: any + await call( + os.use(createRatelimitMiddleware({ limiter: l1, key: 'k1' })) + .use(createRatelimitMiddleware({ limiter: l2, key: 'k2' })) + .handler(({ context }) => { + ctx = context + return 'ok' + }), + undefined, + { context: {} }, + ) + expect((ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] as RatelimiterMiddlewareContext).limits) + .toEqual([{ limiter: l1, key: 'k1' }, { limiter: l2, key: 'k2' }]) + }) + + it('dedupes in nested calls', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'k' }) + const inner = os + .use(mw) + .handler(() => 'in') + const outer = os + .use(mw) + .handler(async ({ context }) => `out:${await call(inner, undefined, { context })}`) + + await call(outer, undefined, { context: {} }) + expect(limiter.limit).toHaveBeenCalledTimes(1) + }) + + it('handles partial context', async () => { + const limiter = createLimiter(success) + await call( + os.use(createRatelimitMiddleware({ limiter, key: 'k' })) + .handler(() => 'ok'), + undefined, + { + context: { [RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL]: { limits: [] } as RatelimiterMiddlewareContext }, + }, + ) + expect(limiter.limit).toHaveBeenCalledTimes(1) + }) + + it('dedupes with existing limit', async () => { + const limiter = createLimiter(success) + const mw = createRatelimitMiddleware({ limiter, key: 'k' }) + let ctx: any + const inner = os + .use(mw) + .handler(({ context }) => { + ctx = context + return 'in' + }) + const outer = os + .use(mw) + .handler(async ({ context }) => await call(inner, undefined, { context })) + await call(outer, undefined, { context: {} }) + expect(limiter.limit).toHaveBeenCalledTimes(1) + expect(ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL].limits).toHaveLength(1) + }) + + it('respects per-instance dedupe', async () => { + const limiter = createLimiter(success) + await call( + os.use(createRatelimitMiddleware({ limiter, key: 'k', dedupe: true })) + .use(createRatelimitMiddleware({ limiter, key: 'k', dedupe: false })).handler(() => 'ok'), + undefined, + { context: {} }, + ) + expect(limiter.limit).toHaveBeenCalledTimes(2) + }) + }) +}) diff --git a/packages/ratelimit/src/middleware.ts b/packages/ratelimit/src/middleware.ts index aceb37259..60b6449dd 100644 --- a/packages/ratelimit/src/middleware.ts +++ b/packages/ratelimit/src/middleware.ts @@ -6,7 +6,7 @@ import { ORPCError } from '@orpc/server' import { toArray, value } from '@orpc/shared' import { RATELIMIT_HANDLER_CONTEXT_SYMBOL } from './handler-plugin' -const RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_MIDDLEWARE_CONTEXT') +export const RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_MIDDLEWARE_CONTEXT') export interface RatelimiterMiddlewareContext { /** @@ -42,8 +42,8 @@ export interface CreateRatelimitMiddlewareOptions< export function createRatelimitMiddleware< TInContext extends Context, - TInput, - TMeta extends Meta, + TInput = unknown, + TMeta extends Meta = Record, >( { dedupe = true, ...options }: CreateRatelimitMiddlewareOptions, ): Middleware, TInput, any, any, TMeta> { From 6ed445720a135e2c028422ada07938fd9b8bf6a9 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 20:18:31 +0700 Subject: [PATCH 15/23] ete test --- packages/ratelimit/tests/e2e.test.ts | 67 ++++++++++++++++++++++++++++ 1 file changed, 67 insertions(+) create mode 100644 packages/ratelimit/tests/e2e.test.ts diff --git a/packages/ratelimit/tests/e2e.test.ts b/packages/ratelimit/tests/e2e.test.ts new file mode 100644 index 000000000..87224fbd7 --- /dev/null +++ b/packages/ratelimit/tests/e2e.test.ts @@ -0,0 +1,67 @@ +import type { Ratelimiter } from '../src' +import { os } from '@orpc/server' +import { RPCHandler } from '@orpc/server/fetch' +import { z } from 'zod' +import { createRatelimitMiddleware, RatelimitHandlerPlugin } from '../src' +import { MemoryRatelimiter } from '../src/adapters/memory' + +it('works', async () => { + const router = { + login: os + .$context<{ limiter: Ratelimiter }>() + .input(z.object({ email: z.email() })) + .use( + createRatelimitMiddleware({ + limiter: ({ context }) => context.limiter, + key: (_, input) => `ping:${input.email}`, + }), + ) + .handler(({ input }) => { + return { success: true } + }), + } + + const handler = new RPCHandler(router, { + plugins: [ + new RatelimitHandlerPlugin(), + ], + }) + + const request = new Request('https://example.com/login', { + method: 'POST', + body: JSON.stringify({ json: { email: 'test@example.com' } }), + headers: { + 'Content-Type': 'application/json', + }, + }) + + const limiter = new MemoryRatelimiter({ + maxRequests: 5, + window: 1000, + }) + + for (let i = 0; i < 5; i++) { + const { response } = await handler.handle(request.clone(), { + context: { + limiter, + }, + }) + + expect(response?.status).toBe(200) + expect(response?.headers.get('RateLimit-Limit')).toBe('5') + expect(response?.headers.get('RateLimit-Remaining')).toBe(String(4 - i)) + expect(response?.headers.get('RateLimit-Reset')).toBeTypeOf('string') + } + + const { response } = await handler.handle(request.clone(), { + context: { + limiter, + }, + }) + + expect(response?.status).toBe(429) + expect(response?.headers.get('RateLimit-Limit')).toBe('5') + expect(response?.headers.get('RateLimit-Remaining')).toBe('0') + expect(response?.headers.get('RateLimit-Reset')).toBeTypeOf('string') + expect(response?.headers.get('Retry-After')).toBeTypeOf('string') +}) From 7a2c5214099b3e643fb03a1b153d71eeaf359043 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 21:02:40 +0700 Subject: [PATCH 16/23] docs --- apps/content/.vitepress/config.ts | 1 + apps/content/docs/helpers/ratelimit.md | 229 ++++++++++++++++++ packages/ratelimit/src/adapters/redis.test.ts | 12 +- packages/ratelimit/src/adapters/redis.ts | 10 +- .../src/adapters/upstash-ratelimit.test.ts | 8 +- .../src/adapters/upstash-ratelimit.ts | 8 +- 6 files changed, 249 insertions(+), 19 deletions(-) create mode 100644 apps/content/docs/helpers/ratelimit.md diff --git a/apps/content/.vitepress/config.ts b/apps/content/.vitepress/config.ts index 3379e62c3..5220bcd81 100644 --- a/apps/content/.vitepress/config.ts +++ b/apps/content/.vitepress/config.ts @@ -164,6 +164,7 @@ export default withMermaid(defineConfig({ { text: 'Encryption', link: '/docs/helpers/encryption' }, { text: 'Form Data', link: '/docs/helpers/form-data' }, { text: 'Publisher', link: '/docs/helpers/publisher' }, + { text: 'Ratelimit', link: '/docs/helpers/ratelimit' }, { text: 'Signing', link: '/docs/helpers/signing' }, ], }, diff --git a/apps/content/docs/helpers/ratelimit.md b/apps/content/docs/helpers/ratelimit.md new file mode 100644 index 000000000..066612db2 --- /dev/null +++ b/apps/content/docs/helpers/ratelimit.md @@ -0,0 +1,229 @@ +--- +title: Rate Limit +description: Rate limiting features for oRPC with multiple adapters support. +--- + +# Rate Limit + +The Rate Limit package provides flexible rate limiting for oRPC with multiple storage backend support. It includes adapters for in-memory, Redis, and Upstash, along with middleware and plugin helpers for seamless integration. + +## Installation + +::: code-group + +```sh [npm] +npm install @orpc/experimental-ratelimit@latest +``` + +```sh [yarn] +yarn add @orpc/experimental-ratelimit@latest +``` + +```sh [pnpm] +pnpm add @orpc/experimental-ratelimit@latest +``` + +```sh [bun] +bun add @orpc/experimental-ratelimit@latest +``` + +```sh [deno] +deno add npm:@orpc/experimental-ratelimit@latest +``` + +::: + +## Available Adapters + +### Memory Adapter + +A simple in-memory rate limiter using a sliding window log algorithm. Ideal for single-instance applications or development. + +```ts +import { MemoryRatelimiter } from '@orpc/experimental-ratelimit/memory' + +const limiter = new MemoryRatelimiter({ + maxRequests: 10, // Maximum requests allowed + window: 60000, // Time window in milliseconds (60 seconds) +}) +``` + +### Redis Adapter + +Redis-based rate limiter using atomic Lua scripts for distributed rate limiting. + +```ts +import { RedisRatelimiter } from '@orpc/experimental-ratelimit/redis' +import { Redis } from 'ioredis' + +const redis = new Redis('redis://localhost:6379') + +const limiter = new RedisRatelimiter({ + eval: async (script: string, numKeys: number, ...args: string[]) => { + return redis.eval(script, numKeys, ...args) + }, + maxRequests: 100, + window: 60000, + prefix: 'orpc:ratelimit:', // Optional key prefix +}) +``` + +::: info +You can use any Redis client that supports Lua script evaluation by providing an `eval` function. +::: + +### Upstash Adapter + +Adapter for [@upstash/ratelimit](https://www.npmjs.com/package/@upstash/ratelimit), optimized for serverless environments like Vercel Edge and Cloudflare Workers. + +```ts +import { Ratelimit } from '@upstash/ratelimit' +import { Redis } from '@upstash/redis' +import { UpstashRatelimiter } from '@orpc/experimental-ratelimit/upstash-ratelimit' + +const redis = Redis.fromEnv() + +const ratelimit = new Ratelimit({ + redis, + limiter: Ratelimit.slidingWindow(10, '60 s'), + prefix: 'my-app:', +}) + +const limiter = new UpstashRatelimiter(ratelimit) +``` + +::: tip Edge Runtime Support +For Edge runtime like Vercel Edge or Cloudflare Workers, pass the `waitUntil` function to enable background analytics: + +```ts +const limiter = new UpstashRatelimiter(ratelimit, { + waitUntil: ctx.waitUntil.bind(ctx), +}) +``` + +::: + +## Blocking Mode + +Some adapters support blocking mode, which waits for the rate limit to reset instead of immediately rejecting requests. + +```ts +const limiter = new MemoryRatelimiter({ + maxRequests: 10, + window: 60000, + blockingUntilReady: { + enabled: true, + timeout: 5000, // Wait up to 5 seconds + }, +}) +``` + +## Manual Usage + +You can use adapters directly without middleware for custom rate limiting logic: + +```ts twoslash +import { MemoryRatelimiter } from '@orpc/experimental-ratelimit/memory' +import { ORPCError } from '@orpc/server' + +const limiter = new MemoryRatelimiter({ + maxRequests: 5, + window: 60000, +}) + +const result = await limiter.limit('user:123') + +if (!result.success) { + throw new ORPCError('TOO_MANY_REQUESTS', { + data: { + limit: result.limit, + remaining: result.remaining, + reset: result.reset, + }, + }) +} +``` + +## `createRatelimitMiddleware` + +The `createRatelimitMiddleware` helper simplifies rate limiting in oRPC procedures. + +```ts twoslash +import { call, os } from '@orpc/server' +import { MemoryRatelimiter } from '@orpc/experimental-ratelimit/memory' +import { createRatelimitMiddleware, Ratelimiter } from '@orpc/experimental-ratelimit' +import { z } from 'zod' + +const loginProcedure = os + .$context<{ ratelimiter: Ratelimiter }>() + .input(z.object({ email: z.email() })) + .use( + createRatelimitMiddleware({ + limiter: ({ context }) => context.ratelimiter, + key: ({ context }, input) => `login:${input.email}`, + }), + ) + .handler(({ input }) => { + return { success: true } + }) + +const ratelimiter = new MemoryRatelimiter({ + maxRequests: 10, + window: 60000, +}) + +const result = await call( + loginProcedure, + { email: 'user@example.com' }, + { context: { ratelimiter } } +) +``` + +::: info Automatic Deduplication +The `createRatelimitMiddleware` automatically deduplicates rate limit checks when the same `limiter` and `key` combination is used multiple times in a request chain. This behavior follows the [Dedupe Middleware Best Practice](/docs/best-practices/dedupe-middleware). To disable deduplication, set the `dedupe: false` option. +::: + +::: tip Conditional Limiter +You can dynamically choose different limiters based on context: + +```ts +const premiumLimiter = new MemoryRatelimiter({ + maxRequests: 100, + window: 60000, +}) + +const standardLimiter = new MemoryRatelimiter({ + maxRequests: 10, + window: 60000, +}) + +const result = await call( + loginProcedure, + { email: 'user@example.com' }, + { + context: { + ratelimiter: isPremiumUser ? premiumLimiter : standardLimiter, + }, + }, +) +``` + +::: + +## Handler Plugin + +The `RatelimitHandlerPlugin` automatically adds rate limit headers (`RateLimit-*` and `Retry-After`) to HTTP responses when using middleware created with `createRatelimitMiddleware`. + +```ts +import { RatelimitHandlerPlugin } from '@orpc/experimental-ratelimit' + +const handler = new RPCHandler(router, { + plugins: [ + new RatelimitHandlerPlugin(), + ], +}) +``` + +::: info +The `handler` can be any supported oRPC handler, such as [RPCHandler](/docs/rpc-handler), [OpenAPIHandler](/docs/openapi/openapi-handler), or other custom handlers. +::: diff --git a/packages/ratelimit/src/adapters/redis.test.ts b/packages/ratelimit/src/adapters/redis.test.ts index 4d163ae36..1819bd859 100644 --- a/packages/ratelimit/src/adapters/redis.test.ts +++ b/packages/ratelimit/src/adapters/redis.test.ts @@ -1,6 +1,6 @@ -import type { IORedisRatelimiterOptions } from './redis' +import type { RedisRatelimiterOptions } from './redis' import { Redis } from 'ioredis' -import { IORedisRatelimiter } from './redis' +import { RedisRatelimiter } from './redis' const REDIS_URL = process.env.REDIS_URL @@ -11,8 +11,8 @@ const REDIS_URL = process.env.REDIS_URL describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, () => { let redis: Redis - function createTestingRatelimiter(options: Partial = {}) { - const ratelimiter = new IORedisRatelimiter({ + function createTestingRatelimiter(options: Partial = {}) { + const ratelimiter = new RedisRatelimiter({ eval: redis.eval.bind(redis), prefix: `test:${crypto.randomUUID()}:`, // isolated from other tests maxRequests: 10, @@ -192,7 +192,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) it('handles Redis errors gracefully', async () => { - const ratelimiter = new IORedisRatelimiter({ + const ratelimiter = new RedisRatelimiter({ eval: async () => { throw new Error('Redis error') }, maxRequests: 10, window: 60000, @@ -201,7 +201,7 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) it('handles invalid script response', async () => { - const ratelimiter = new IORedisRatelimiter({ + const ratelimiter = new RedisRatelimiter({ eval: async () => [1, 2, 3], // Invalid response, should have 4 elements maxRequests: 10, window: 60000, diff --git a/packages/ratelimit/src/adapters/redis.ts b/packages/ratelimit/src/adapters/redis.ts index bc68974d7..20b2978bf 100644 --- a/packages/ratelimit/src/adapters/redis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -47,7 +47,7 @@ else end ` -export interface IORedisRatelimiterOptions { +export interface RedisRatelimiterOptions { eval: (script: string, numKeys: number, ...args: string[]) => Promise /** @@ -76,15 +76,15 @@ export interface IORedisRatelimiterOptions { window: number } -export class IORedisRatelimiter implements Ratelimiter { - private readonly eval: IORedisRatelimiterOptions['eval'] +export class RedisRatelimiter implements Ratelimiter { + private readonly eval: RedisRatelimiterOptions['eval'] private readonly prefix: string private readonly maxRequests: number private readonly window: number - private readonly blockingUntilReady: IORedisRatelimiterOptions['blockingUntilReady'] + private readonly blockingUntilReady: RedisRatelimiterOptions['blockingUntilReady'] constructor( - options: IORedisRatelimiterOptions, + options: RedisRatelimiterOptions, ) { this.eval = options.eval this.prefix = fallback(options.prefix, 'orpc:ratelimit:') diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts index 602d37218..bc9b29e05 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.test.ts @@ -66,15 +66,15 @@ describe.concurrent( }) it('should use waitUntil callback when provided', async () => { - const waitUntilSpy = vi.fn() + const waitUntil = vi.fn() const ratelimiter = createTestingRatelimiter({ - waitUtil: waitUntilSpy, + waitUntil, }) const key = `test-waituntil-${crypto.randomUUID()}` const result1 = await ratelimiter.limit(key) - expect(waitUntilSpy).toHaveBeenCalledTimes(1) - expect(waitUntilSpy).toHaveBeenCalledWith(expect.any(Promise)) + expect(waitUntil).toHaveBeenCalledTimes(1) + expect(waitUntil).toHaveBeenCalledWith(expect.any(Promise)) }) }, ) diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 5251ffcea..661f67801 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -22,19 +22,19 @@ export interface UpstashRatelimiterOptions { * }) * ``` */ - waitUtil?: (promise: Promise) => any + waitUntil?: (promise: Promise) => any } export class UpstashRatelimiter implements Ratelimiter { private blockingUntilReady: UpstashRatelimiterOptions['blockingUntilReady'] - private waitUtil: UpstashRatelimiterOptions['waitUtil'] + private waitUntil: UpstashRatelimiterOptions['waitUntil'] constructor( private readonly ratelimit: Ratelimit, options: UpstashRatelimiterOptions = {}, ) { this.blockingUntilReady = options.blockingUntilReady - this.waitUtil = options.waitUtil + this.waitUntil = options.waitUntil } async limit(key: string): Promise> { @@ -42,7 +42,7 @@ export class UpstashRatelimiter implements Ratelimiter { ? await this.ratelimit.blockUntilReady(key, this.blockingUntilReady.timeout) : await this.ratelimit.limit(key) - this.waitUtil?.(result.pending) + this.waitUntil?.(result.pending) return result } } From 1d7953e8a4b3b7c724b1caa7692b39f364011e89 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 21:08:39 +0700 Subject: [PATCH 17/23] improve --- packages/ratelimit/src/adapters/redis.ts | 4 ++-- packages/ratelimit/src/handler-plugin.ts | 20 +++++++++++--------- packages/ratelimit/src/middleware.test.ts | 2 +- packages/ratelimit/src/middleware.ts | 16 +++++++++------- 4 files changed, 23 insertions(+), 19 deletions(-) diff --git a/packages/ratelimit/src/adapters/redis.ts b/packages/ratelimit/src/adapters/redis.ts index 20b2978bf..57df4ad54 100644 --- a/packages/ratelimit/src/adapters/redis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -7,11 +7,11 @@ import { fallback } from '@orpc/shared' * This script implements atomic sliding window log rate limiting using Redis sorted sets. * It removes expired entries, checks the current count, and adds new requests atomically. * - * @returns A tuple with [success, limit, remaining, resetAtMs] where: + * @returns A tuple with [success, limit, remaining, reset] where: * - success: 1 if request is allowed, 0 if rate limited * - limit: The maximum number of requests allowed * - remaining: Number of requests remaining in the window - * - resetAtMs: Unix timestamp (in milliseconds) when the window resets + * - reset: Unix timestamp (in milliseconds) when the window resets */ const SLIDING_WINDOW_LOG_LUA_SCRIPT = ` local key = KEYS[1] diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts index 1977c9f06..4326970f7 100644 --- a/packages/ratelimit/src/handler-plugin.ts +++ b/packages/ratelimit/src/handler-plugin.ts @@ -2,16 +2,18 @@ import type { Context } from '@orpc/server' import type { StandardHandlerOptions, StandardHandlerPlugin } from '@orpc/server/standard' import type { RatelimiterLimitResult } from './types' -export const RATELIMIT_HANDLER_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_HANDLER_CONTEXT') - -export interface RatelimitHandlerPluginContext { - /** - * The result of the ratelimiter after applying limits - */ - ratelimitResult?: RatelimiterLimitResult +export const RATELIMIT_HANDLER_CONTEXT_SYMBOL: unique symbol = Symbol('ORPC_RATE_LIMIT_HANDLER_CONTEXT') + +export interface RatelimitHandlerPluginContext extends Context { + [RATELIMIT_HANDLER_CONTEXT_SYMBOL]?: { + /** + * The result of the ratelimiter after applying limits + */ + ratelimitResult?: RatelimiterLimitResult + } } -export class RatelimitHandlerPlugin implements StandardHandlerPlugin { +export class RatelimitHandlerPlugin implements StandardHandlerPlugin { /** * this plugin should lower priority than response headers plugin, * if user want override rate limit headers @@ -22,7 +24,7 @@ export class RatelimitHandlerPlugin implements StandardHandle options.rootInterceptors ??= [] options.rootInterceptors.push(async (interceptorOptions) => { - const handlerContext: RatelimitHandlerPluginContext = {} + const handlerContext: Exclude = {} const result = await interceptorOptions.next({ ...interceptorOptions, diff --git a/packages/ratelimit/src/middleware.test.ts b/packages/ratelimit/src/middleware.test.ts index 257da157c..6b34ca6be 100644 --- a/packages/ratelimit/src/middleware.test.ts +++ b/packages/ratelimit/src/middleware.test.ts @@ -232,7 +232,7 @@ describe('createRatelimitMiddleware', () => { undefined, { context: {} }, ) - expect((ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] as RatelimiterMiddlewareContext).limits) + expect(ctx[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL].limits) .toEqual([{ limiter: l1, key: 'k1' }, { limiter: l2, key: 'k2' }]) }) diff --git a/packages/ratelimit/src/middleware.ts b/packages/ratelimit/src/middleware.ts index 60b6449dd..47ea6a210 100644 --- a/packages/ratelimit/src/middleware.ts +++ b/packages/ratelimit/src/middleware.ts @@ -6,13 +6,15 @@ import { ORPCError } from '@orpc/server' import { toArray, value } from '@orpc/shared' import { RATELIMIT_HANDLER_CONTEXT_SYMBOL } from './handler-plugin' -export const RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL = Symbol('ORPC_RATE_LIMIT_MIDDLEWARE_CONTEXT') +export const RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL: unique symbol = Symbol('ORPC_RATE_LIMIT_MIDDLEWARE_CONTEXT') export interface RatelimiterMiddlewareContext { - /** - * The applied limits in this request, mainly for deduplication purposes - */ - limits: { limiter: Ratelimiter, key: string }[] + [RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL]?: { + /** + * The applied limits in this request, mainly for deduplication purposes + */ + limits: { limiter: Ratelimiter, key: string }[] + } } export interface CreateRatelimitMiddlewareOptions< @@ -53,14 +55,14 @@ export function createRatelimitMiddleware< value(options.key, middlewareOptions, input), ]) - const middlewareContext: RatelimiterMiddlewareContext | undefined = middlewareOptions.context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] + const middlewareContext: RatelimiterMiddlewareContext[typeof RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] = middlewareOptions.context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] if (dedupe && middlewareContext?.limits.some(l => l.key === key && l.limiter === limiter)) { return middlewareOptions.next() } const result = await limiter.limit(key) - const pluginContext: RatelimitHandlerPluginContext | undefined = middlewareOptions.context[RATELIMIT_HANDLER_CONTEXT_SYMBOL] + const pluginContext: Exclude = middlewareOptions.context[RATELIMIT_HANDLER_CONTEXT_SYMBOL] if (pluginContext) { pluginContext.ratelimitResult = result } From cdd1bd2e23a60993292be1ac6584621d9caeb065 Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 21:12:13 +0700 Subject: [PATCH 18/23] fix --- apps/content/package.json | 1 + pnpm-lock.yaml | 3 +++ 2 files changed, 4 insertions(+) diff --git a/apps/content/package.json b/apps/content/package.json index d60616e3d..9dac5f3c1 100644 --- a/apps/content/package.json +++ b/apps/content/package.json @@ -18,6 +18,7 @@ "@orpc/client": "workspace:*", "@orpc/contract": "workspace:*", "@orpc/experimental-publisher": "workspace:*", + "@orpc/experimental-ratelimit": "workspace:*", "@orpc/experimental-react-swr": "workspace:*", "@orpc/openapi": "workspace:*", "@orpc/openapi-client": "workspace:*", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index f057cddd6..1fa97619b 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -132,6 +132,9 @@ importers: '@orpc/experimental-publisher': specifier: workspace:* version: link:../../packages/publisher + '@orpc/experimental-ratelimit': + specifier: workspace:* + version: link:../../packages/ratelimit '@orpc/experimental-react-swr': specifier: workspace:* version: link:../../packages/react-swr From 5b915451bd91dbf55a75167931e69c5b9457489f Mon Sep 17 00:00:00 2001 From: unnoq Date: Thu, 6 Nov 2025 22:02:25 +0700 Subject: [PATCH 19/23] fix --- .../ratelimit/src/adapters/memory.test.ts | 62 ++++++++++++++++--- packages/ratelimit/src/adapters/memory.ts | 6 +- 2 files changed, 57 insertions(+), 11 deletions(-) diff --git a/packages/ratelimit/src/adapters/memory.test.ts b/packages/ratelimit/src/adapters/memory.test.ts index 3cec26315..ed36f0025 100644 --- a/packages/ratelimit/src/adapters/memory.test.ts +++ b/packages/ratelimit/src/adapters/memory.test.ts @@ -37,16 +37,21 @@ describe('memoryRatelimiter', () => { it('should reset after window expires', async () => { const ratelimiter = createTestingRatelimiter({ window: 200, + maxRequests: 3.0, }) const result1 = await ratelimiter.limit('test') - expect(result1.remaining).toBe(1) + expect(result1.remaining).toBe(2) + const result2 = await ratelimiter.limit('test') + expect(result2.remaining).toBe(1) + const result3 = await ratelimiter.limit('test') + expect(result3.remaining).toBe(0) await sleep(210) - const result2 = await ratelimiter.limit('test') - expect(result2.success).toBe(true) - expect(result2.remaining).toBe(1) + const result4 = await ratelimiter.limit('test') + expect(result4.success).toBe(true) + expect(result4.remaining).toBe(2) }) it('should handle multiple keys independently', async () => { @@ -140,19 +145,58 @@ describe('memoryRatelimiter', () => { window: 200, }) - await ratelimiter.limit('test') + await ratelimiter.limit('test1') + await ratelimiter.limit('test2') // @ts-expect-error accessing private property for testing - expect(ratelimiter.store.size).toBe(1) + expect(ratelimiter.store.size).toBe(2) - // move time forward to pass window await sleep(210) - // Make another limit call to different key to trigger cleanup + // Trigger cleanup - test1 should be removed + await ratelimiter.limit('test2') + + // @ts-expect-error accessing private property for testing + expect(ratelimiter.store.size).toBe(1) + }) + + it('should handle cleanup with all expired timestamps', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 2, + window: 150, + }) + + await ratelimiter.limit('test1') + await sleep(160) + + // @ts-expect-error accessing private property for testing + ratelimiter.lastCleanupTime = Date.now() - 200 + + // Trigger cleanup - test1 has all timestamps expired (idx === -1 branch) + await ratelimiter.limit('test2') + + // @ts-expect-error accessing private property for testing + expect(ratelimiter.store.has('test1')).toBe(false) + }) + + it('should handle cleanup with partially expired timestamps', async () => { + const ratelimiter = createTestingRatelimiter({ + maxRequests: 3, + window: 150, + }) + + await ratelimiter.limit('test1') + await sleep(160) + await ratelimiter.limit('test1') + + // @ts-expect-error accessing private property for testing + ratelimiter.lastCleanupTime = Date.now() - 200 + + // Trigger cleanup - test1 has partial timestamps expired (idx !== -1 branch) await ratelimiter.limit('test2') // @ts-expect-error accessing private property for testing - expect(ratelimiter.store.size).toBe(1) // Only test2 should remain + expect(ratelimiter.store.get('test1')?.length).toBe(1) }) }) }) diff --git a/packages/ratelimit/src/adapters/memory.ts b/packages/ratelimit/src/adapters/memory.ts index 863012d25..1eec014d4 100644 --- a/packages/ratelimit/src/adapters/memory.ts +++ b/packages/ratelimit/src/adapters/memory.ts @@ -58,7 +58,8 @@ export class MemoryRatelimiter implements Ratelimiter { for (const [key, timestamps] of this.store) { // remove expired timestamps - timestamps.splice(0, timestamps.findIndex(timestamp => timestamp < windowStart) + 1) + const idx = timestamps.findIndex(timestamp => timestamp >= windowStart) + timestamps.splice(0, idx === -1 ? timestamps.length : idx) if (timestamps.length === 0) { this.store.delete(key) @@ -73,7 +74,8 @@ export class MemoryRatelimiter implements Ratelimiter { let timestamps = this.store.get(key) if (timestamps) { // Remove expired timestamps - timestamps.splice(0, timestamps.findIndex(timestamp => timestamp < windowStart) + 1) + const idx = timestamps.findIndex(timestamp => timestamp >= windowStart) + timestamps.splice(0, idx === -1 ? timestamps.length : idx) } else { this.store.set(key, timestamps = []) From 200e6510b875c5cba838556f66a210859cbac231 Mon Sep 17 00:00:00 2001 From: unnoq Date: Fri, 7 Nov 2025 09:34:21 +0700 Subject: [PATCH 20/23] fix --- apps/content/docs/helpers/ratelimit.md | 6 +++--- packages/ratelimit/src/adapters/redis.ts | 4 ++-- .../ratelimit/src/adapters/upstash-ratelimit.ts | 2 +- packages/ratelimit/src/handler-plugin.ts | 16 +++++++--------- packages/ratelimit/src/middleware.ts | 8 ++++---- packages/ratelimit/tsconfig.json | 1 + 6 files changed, 18 insertions(+), 19 deletions(-) diff --git a/apps/content/docs/helpers/ratelimit.md b/apps/content/docs/helpers/ratelimit.md index 066612db2..15e88e51b 100644 --- a/apps/content/docs/helpers/ratelimit.md +++ b/apps/content/docs/helpers/ratelimit.md @@ -59,8 +59,8 @@ import { Redis } from 'ioredis' const redis = new Redis('redis://localhost:6379') const limiter = new RedisRatelimiter({ - eval: async (script: string, numKeys: number, ...args: string[]) => { - return redis.eval(script, numKeys, ...args) + eval: async (script, numKeys, ...rest) => { + return redis.eval(script, numKeys, ...rest) }, maxRequests: 100, window: 60000, @@ -93,7 +93,7 @@ const limiter = new UpstashRatelimiter(ratelimit) ``` ::: tip Edge Runtime Support -For Edge runtime like Vercel Edge or Cloudflare Workers, pass the `waitUntil` function to enable background analytics: +For Edge runtime like Vercel Edge or Cloudflare Workers, pass the `waitUntil` function to better handle background tasks: ```ts const limiter = new UpstashRatelimiter(ratelimit, { diff --git a/packages/ratelimit/src/adapters/redis.ts b/packages/ratelimit/src/adapters/redis.ts index 57df4ad54..88da70331 100644 --- a/packages/ratelimit/src/adapters/redis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -48,7 +48,7 @@ end ` export interface RedisRatelimiterOptions { - eval: (script: string, numKeys: number, ...args: string[]) => Promise + eval: (script: string, numKeys: number, ...rest: string[]) => Promise /** * Block until the request may pass or timeout is reached. @@ -111,7 +111,7 @@ export class RedisRatelimiter implements Ratelimiter { Date.now().toString(), this.window.toString(), this.maxRequests.toString(), - ) as unknown + ) if (!Array.isArray(result) || result.length !== 4) { throw new TypeError('Invalid response from rate limit script') diff --git a/packages/ratelimit/src/adapters/upstash-ratelimit.ts b/packages/ratelimit/src/adapters/upstash-ratelimit.ts index 661f67801..09549ea8a 100644 --- a/packages/ratelimit/src/adapters/upstash-ratelimit.ts +++ b/packages/ratelimit/src/adapters/upstash-ratelimit.ts @@ -18,7 +18,7 @@ export interface UpstashRatelimiterOptions { * On Vercel Edge or Cloudflare workers, you might need `.bind` before assign: * ```ts * const ratelimiter = new UpstashRatelimiter(ratelimit, { - * waitUtil: ctx.waitUntil.bind(ctx), + * waitUntil: ctx.waitUntil.bind(ctx), * }) * ``` */ diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts index 4326970f7..18daa23f7 100644 --- a/packages/ratelimit/src/handler-plugin.ts +++ b/packages/ratelimit/src/handler-plugin.ts @@ -4,7 +4,7 @@ import type { RatelimiterLimitResult } from './types' export const RATELIMIT_HANDLER_CONTEXT_SYMBOL: unique symbol = Symbol('ORPC_RATE_LIMIT_HANDLER_CONTEXT') -export interface RatelimitHandlerPluginContext extends Context { +export interface RatelimitHandlerPluginContext { [RATELIMIT_HANDLER_CONTEXT_SYMBOL]?: { /** * The result of the ratelimiter after applying limits @@ -13,17 +13,15 @@ export interface RatelimitHandlerPluginContext extends Context { } } -export class RatelimitHandlerPlugin implements StandardHandlerPlugin { - /** - * this plugin should lower priority than response headers plugin, - * if user want override rate limit headers - */ - order = 100_000 - +export class RatelimitHandlerPlugin implements StandardHandlerPlugin { init(options: StandardHandlerOptions): void { options.rootInterceptors ??= [] - options.rootInterceptors.push(async (interceptorOptions) => { + /** + * This plugin should set headers before "response headers" plugin or user defined interceptors + * In case user wants to override ratelimit headers + */ + options.rootInterceptors.unshift(async (interceptorOptions) => { const handlerContext: Exclude = {} const result = await interceptorOptions.next({ diff --git a/packages/ratelimit/src/middleware.ts b/packages/ratelimit/src/middleware.ts index 47ea6a210..80344621c 100644 --- a/packages/ratelimit/src/middleware.ts +++ b/packages/ratelimit/src/middleware.ts @@ -32,7 +32,7 @@ export interface CreateRatelimitMiddlewareOptions< */ key: Value, [middlewareOptions: MiddlewareOptions, TMeta>, input: TInput]> /** - * If you ratelimit middleware is used multiple times + * If your ratelimit middleware is used multiple times * or you invoke a procedure inside another procedure (shared the same context) that also has * ratelimit middleware **with the same limiter and key**, this option * will ensure that the limit is only applied once per request. @@ -55,14 +55,14 @@ export function createRatelimitMiddleware< value(options.key, middlewareOptions, input), ]) - const middlewareContext: RatelimiterMiddlewareContext[typeof RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] = middlewareOptions.context[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] + const middlewareContext = (middlewareOptions.context as RatelimiterMiddlewareContext)[RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL] if (dedupe && middlewareContext?.limits.some(l => l.key === key && l.limiter === limiter)) { return middlewareOptions.next() } const result = await limiter.limit(key) - const pluginContext: Exclude = middlewareOptions.context[RATELIMIT_HANDLER_CONTEXT_SYMBOL] + const pluginContext = (middlewareOptions.context as RatelimitHandlerPluginContext)[RATELIMIT_HANDLER_CONTEXT_SYMBOL] if (pluginContext) { pluginContext.ratelimitResult = result } @@ -80,7 +80,7 @@ export function createRatelimitMiddleware< return middlewareOptions.next({ context: { [RATELIMIT_MIDDLEWARE_CONTEXT_SYMBOL]: { - ...middlewareOptions, + ...middlewareContext, limits: [ ...toArray(middlewareContext?.limits), { limiter, key }, diff --git a/packages/ratelimit/tsconfig.json b/packages/ratelimit/tsconfig.json index 954616514..51be5001d 100644 --- a/packages/ratelimit/tsconfig.json +++ b/packages/ratelimit/tsconfig.json @@ -3,6 +3,7 @@ "references": [ { "path": "../shared" }, { "path": "../client" }, + { "path": "../server" }, { "path": "../standard-server" } ], "include": ["src"], From 73eeb5e08b925c6ffeae34c17d54047e8c2662c9 Mon Sep 17 00:00:00 2001 From: unnoq Date: Fri, 7 Nov 2025 10:02:57 +0700 Subject: [PATCH 21/23] improve --- packages/ratelimit/src/adapters/redis.test.ts | 9 +++++++++ packages/ratelimit/src/adapters/redis.ts | 10 +++++++++- 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/packages/ratelimit/src/adapters/redis.test.ts b/packages/ratelimit/src/adapters/redis.test.ts index 1819bd859..7274010c7 100644 --- a/packages/ratelimit/src/adapters/redis.test.ts +++ b/packages/ratelimit/src/adapters/redis.test.ts @@ -208,5 +208,14 @@ describe.concurrent('ioredis ratelimiter', { skip: !REDIS_URL, timeout: 20000 }, }) await expect(ratelimiter.limit('some-key')).rejects.toThrow('Invalid response from rate limit script') }) + + it('handles invalid script response 2', async () => { + const ratelimiter = new RedisRatelimiter({ + eval: async () => ['a', 'b', 'c', 'd'], // should be integers + maxRequests: 10, + window: 60000, + }) + await expect(ratelimiter.limit('some-key')).rejects.toThrow('Invalid response from rate limit script') + }) }) }) diff --git a/packages/ratelimit/src/adapters/redis.ts b/packages/ratelimit/src/adapters/redis.ts index 88da70331..61ceaab37 100644 --- a/packages/ratelimit/src/adapters/redis.ts +++ b/packages/ratelimit/src/adapters/redis.ts @@ -117,7 +117,15 @@ export class RedisRatelimiter implements Ratelimiter { throw new TypeError('Invalid response from rate limit script') } - const [success, limit, remaining, reset] = result as [number, number, number, number] + const numbers = result.map((item) => { + const num = Number(item) + if (!Number.isInteger(num)) { + throw new TypeError('Invalid response from rate limit script') + } + return num + }) + + const [success, limit, remaining, reset] = numbers as [number, number, number, number] return { success: success === 1, From a296ec0282eacc50cf2276187d7f6aea5fd311cb Mon Sep 17 00:00:00 2001 From: unnoq Date: Fri, 7 Nov 2025 12:07:59 +0700 Subject: [PATCH 22/23] improve docs --- apps/content/docs/helpers/ratelimit.md | 4 ++-- packages/ratelimit/src/handler-plugin.ts | 6 ++++++ packages/ratelimit/src/middleware.ts | 6 ++++++ 3 files changed, 14 insertions(+), 2 deletions(-) diff --git a/apps/content/docs/helpers/ratelimit.md b/apps/content/docs/helpers/ratelimit.md index 15e88e51b..0f2201eb8 100644 --- a/apps/content/docs/helpers/ratelimit.md +++ b/apps/content/docs/helpers/ratelimit.md @@ -146,7 +146,7 @@ if (!result.success) { ## `createRatelimitMiddleware` -The `createRatelimitMiddleware` helper simplifies rate limiting in oRPC procedures. +The `createRatelimitMiddleware` helper creates middleware for oRPC procedures to enforce rate limits. ```ts twoslash import { call, os } from '@orpc/server' @@ -212,7 +212,7 @@ const result = await call( ## Handler Plugin -The `RatelimitHandlerPlugin` automatically adds rate limit headers (`RateLimit-*` and `Retry-After`) to HTTP responses when using middleware created with `createRatelimitMiddleware`. +The `RatelimitHandlerPlugin` automatically adds HTTP rate-limiting headers (`RateLimit-*` and `Retry-After`) to responses when used with middleware created by [`createRatelimitMiddleware`](#createratelimitmiddleware). ```ts import { RatelimitHandlerPlugin } from '@orpc/experimental-ratelimit' diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts index 18daa23f7..023e40d22 100644 --- a/packages/ratelimit/src/handler-plugin.ts +++ b/packages/ratelimit/src/handler-plugin.ts @@ -13,6 +13,12 @@ export interface RatelimitHandlerPluginContext { } } +/** + * Automatically adds HTTP rate-limiting headers (RateLimit-* and Retry-After) to responses + * when used with middleware created by createRatelimitMiddleware. + * + * @see {@link https://orpc.unnoq.com/docs/helpers/ratelimit#handler-plugin Ratelimit handler plugin} + */ export class RatelimitHandlerPlugin implements StandardHandlerPlugin { init(options: StandardHandlerOptions): void { options.rootInterceptors ??= [] diff --git a/packages/ratelimit/src/middleware.ts b/packages/ratelimit/src/middleware.ts index 80344621c..4af501a7b 100644 --- a/packages/ratelimit/src/middleware.ts +++ b/packages/ratelimit/src/middleware.ts @@ -42,6 +42,12 @@ export interface CreateRatelimitMiddlewareOptions< dedupe?: boolean } +/** + * Creates a middleware that enforces rate limits in oRPC procedures. + * Supports per-request deduplication and integrates with the ratelimit handler plugin. + * + * @see {@link https://orpc.unnoq.com/docs/helpers/ratelimit#createratelimitmiddleware Ratelimit middleware} + */ export function createRatelimitMiddleware< TInContext extends Context, TInput = unknown, From 7a99cb78d9e3fcb6c1589e35498268ba0ce4e08b Mon Sep 17 00:00:00 2001 From: unnoq Date: Fri, 7 Nov 2025 14:48:32 +0700 Subject: [PATCH 23/23] fix: clamp retry-after header to 0 when reset time is in the past --- packages/ratelimit/src/handler-plugin.test.ts | 28 +++++++++++++++++++ packages/ratelimit/src/handler-plugin.ts | 2 +- 2 files changed, 29 insertions(+), 1 deletion(-) diff --git a/packages/ratelimit/src/handler-plugin.test.ts b/packages/ratelimit/src/handler-plugin.test.ts index b218b631a..594a083c0 100644 --- a/packages/ratelimit/src/handler-plugin.test.ts +++ b/packages/ratelimit/src/handler-plugin.test.ts @@ -120,6 +120,34 @@ describe('ratelimitHandlerPlugin', () => { expect(result.response.headers['retry-after']).toBeUndefined() }) + it('clamps retry-after to 0 when reset is in the past', async () => { + const resetTime = Date.now() - 5000 // 5 seconds ago + + const options: any = { rootInterceptors: [] } + new RatelimitHandlerPlugin().init(options) + + options.rootInterceptors.push(async ({ context }: any) => { + context[RATELIMIT_HANDLER_CONTEXT_SYMBOL].ratelimitResult = { + success: false, + reset: resetTime, + } + + return { + matched: true as const, + response: { status: 429, headers: {}, body: 'too many requests' }, + } + }) + + const handler = new StandardHandler({}, new StandardRPCMatcher(), {} as any, options) + const result = await handler.handle(createMockRequest('https://example.com/ping'), { context: {} }) + + if (!result.matched) { + throw new Error('request should match') + } + + expect(result.response.headers['retry-after']).toBe('0') + }) + it('does not add headers when ratelimitResult is undefined', async () => { const options: any = { rootInterceptors: [] } new RatelimitHandlerPlugin().init(options) diff --git a/packages/ratelimit/src/handler-plugin.ts b/packages/ratelimit/src/handler-plugin.ts index 023e40d22..747faa758 100644 --- a/packages/ratelimit/src/handler-plugin.ts +++ b/packages/ratelimit/src/handler-plugin.ts @@ -49,7 +49,7 @@ export class RatelimitHandlerPlugin implements StandardHandle 'ratelimit-remaining': handlerContext.ratelimitResult.remaining?.toString(), 'ratelimit-reset': handlerContext.ratelimitResult.reset?.toString(), 'retry-after': !handlerContext.ratelimitResult.success && result.response.status === 429 && handlerContext.ratelimitResult.reset !== undefined - ? Math.ceil((handlerContext.ratelimitResult.reset - Date.now()) / 1000).toString() + ? Math.max(0, Math.ceil((handlerContext.ratelimitResult.reset - Date.now()) / 1000)).toString() : undefined, }, },