From 4f6c3711105fdcd27e77b68f55ccaf442d38f079 Mon Sep 17 00:00:00 2001 From: Alexander Eichhorn Date: Mon, 9 Feb 2026 19:18:22 +0100 Subject: [PATCH 1/5] added agent skills --- .agents/skills/durable-objects/SKILL.md | 186 ++++ .../durable-objects/references/rules.md | 286 ++++++ .../durable-objects/references/testing.md | 264 ++++++ .../durable-objects/references/workers.md | 346 +++++++ .agents/skills/wrangler/SKILL.md | 887 ++++++++++++++++++ 5 files changed, 1969 insertions(+) create mode 100644 .agents/skills/durable-objects/SKILL.md create mode 100644 .agents/skills/durable-objects/references/rules.md create mode 100644 .agents/skills/durable-objects/references/testing.md create mode 100644 .agents/skills/durable-objects/references/workers.md create mode 100644 .agents/skills/wrangler/SKILL.md diff --git a/.agents/skills/durable-objects/SKILL.md b/.agents/skills/durable-objects/SKILL.md new file mode 100644 index 0000000..75cbe06 --- /dev/null +++ b/.agents/skills/durable-objects/SKILL.md @@ -0,0 +1,186 @@ +--- +name: durable-objects +description: Create and review Cloudflare Durable Objects. Use when building stateful coordination (chat rooms, multiplayer games, booking systems), implementing RPC methods, SQLite storage, alarms, WebSockets, or reviewing DO code for best practices. Covers Workers integration, wrangler config, and testing with Vitest. +--- + +# Durable Objects + +Build stateful, coordinated applications on Cloudflare's edge using Durable Objects. + +## Retrieval-First Development + +**Prefer retrieval from official docs over pre-training for Durable Objects tasks.** + +| Resource | URL | +|----------|-----| +| Docs | https://developers.cloudflare.com/durable-objects/ | +| API Reference | https://developers.cloudflare.com/durable-objects/api/ | +| Best Practices | https://developers.cloudflare.com/durable-objects/best-practices/ | +| Examples | https://developers.cloudflare.com/durable-objects/examples/ | + +Fetch the relevant doc page when implementing features. + +## When to Use + +- Creating new Durable Object classes for stateful coordination +- Implementing RPC methods, alarms, or WebSocket handlers +- Reviewing existing DO code for best practices +- Configuring wrangler.jsonc/toml for DO bindings and migrations +- Writing tests with `@cloudflare/vitest-pool-workers` +- Designing sharding strategies and parent-child relationships + +## Reference Documentation + +- `./references/rules.md` - Core rules, storage, concurrency, RPC, alarms +- `./references/testing.md` - Vitest setup, unit/integration tests, alarm testing +- `./references/workers.md` - Workers handlers, types, wrangler config, observability + +Search: `blockConcurrencyWhile`, `idFromName`, `getByName`, `setAlarm`, `sql.exec` + +## Core Principles + +### Use Durable Objects For + +| Need | Example | +|------|---------| +| Coordination | Chat rooms, multiplayer games, collaborative docs | +| Strong consistency | Inventory, booking systems, turn-based games | +| Per-entity storage | Multi-tenant SaaS, per-user data | +| Persistent connections | WebSockets, real-time notifications | +| Scheduled work per entity | Subscription renewals, game timeouts | + +### Do NOT Use For + +- Stateless request handling (use plain Workers) +- Maximum global distribution needs +- High fan-out independent requests + +## Quick Reference + +### Wrangler Configuration + +```jsonc +// wrangler.jsonc +{ + "durable_objects": { + "bindings": [{ "name": "MY_DO", "class_name": "MyDurableObject" }] + }, + "migrations": [{ "tag": "v1", "new_sqlite_classes": ["MyDurableObject"] }] +} +``` + +### Basic Durable Object Pattern + +```typescript +import { DurableObject } from "cloudflare:workers"; + +export interface Env { + MY_DO: DurableObjectNamespace; +} + +export class MyDurableObject extends DurableObject { + constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + ctx.blockConcurrencyWhile(async () => { + this.ctx.storage.sql.exec(` + CREATE TABLE IF NOT EXISTS items ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + data TEXT NOT NULL + ) + `); + }); + } + + async addItem(data: string): Promise { + const result = this.ctx.storage.sql.exec<{ id: number }>( + "INSERT INTO items (data) VALUES (?) RETURNING id", + data + ); + return result.one().id; + } +} + +export default { + async fetch(request: Request, env: Env): Promise { + const stub = env.MY_DO.getByName("my-instance"); + const id = await stub.addItem("hello"); + return Response.json({ id }); + }, +}; +``` + +## Critical Rules + +1. **Model around coordination atoms** - One DO per chat room/game/user, not one global DO +2. **Use `getByName()` for deterministic routing** - Same input = same DO instance +3. **Use SQLite storage** - Configure `new_sqlite_classes` in migrations +4. **Initialize in constructor** - Use `blockConcurrencyWhile()` for schema setup only +5. **Use RPC methods** - Not fetch() handler (compatibility date >= 2024-04-03) +6. **Persist first, cache second** - Always write to storage before updating in-memory state +7. **One alarm per DO** - `setAlarm()` replaces any existing alarm + +## Anti-Patterns (NEVER) + +- Single global DO handling all requests (bottleneck) +- Using `blockConcurrencyWhile()` on every request (kills throughput) +- Storing critical state only in memory (lost on eviction/crash) +- Using `await` between related storage writes (breaks atomicity) +- Holding `blockConcurrencyWhile()` across `fetch()` or external I/O + +## Stub Creation + +```typescript +// Deterministic - preferred for most cases +const stub = env.MY_DO.getByName("room-123"); + +// From existing ID string +const id = env.MY_DO.idFromString(storedIdString); +const stub = env.MY_DO.get(id); + +// New unique ID - store mapping externally +const id = env.MY_DO.newUniqueId(); +const stub = env.MY_DO.get(id); +``` + +## Storage Operations + +```typescript +// SQL (synchronous, recommended) +this.ctx.storage.sql.exec("INSERT INTO t (c) VALUES (?)", value); +const rows = this.ctx.storage.sql.exec("SELECT * FROM t").toArray(); + +// KV (async) +await this.ctx.storage.put("key", value); +const val = await this.ctx.storage.get("key"); +``` + +## Alarms + +```typescript +// Schedule (replaces existing) +await this.ctx.storage.setAlarm(Date.now() + 60_000); + +// Handler +async alarm(): Promise { + // Process scheduled work + // Optionally reschedule: await this.ctx.storage.setAlarm(...) +} + +// Cancel +await this.ctx.storage.deleteAlarm(); +``` + +## Testing Quick Start + +```typescript +import { env } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; + +describe("MyDO", () => { + it("should work", async () => { + const stub = env.MY_DO.getByName("test"); + const result = await stub.addItem("test"); + expect(result).toBe(1); + }); +}); +``` diff --git a/.agents/skills/durable-objects/references/rules.md b/.agents/skills/durable-objects/references/rules.md new file mode 100644 index 0000000..e2ce48f --- /dev/null +++ b/.agents/skills/durable-objects/references/rules.md @@ -0,0 +1,286 @@ +# Durable Objects Rules & Best Practices + +## Design & Sharding + +### Model Around Coordination Atoms + +Create one DO per logical unit needing coordination: chat room, game session, document, user, tenant. + +```typescript +// ✅ Good: One DO per chat room +const stub = env.CHAT_ROOM.getByName(roomId); + +// ❌ Bad: Single global DO +const stub = env.CHAT_ROOM.getByName("global"); // Bottleneck! +``` + +### Parent-Child Relationships + +For hierarchical data, create separate child DOs. Parent tracks references, children handle own state. + +```typescript +// Parent: GameServer tracks match references +// Children: GameMatch handles individual match state +async createMatch(name: string): Promise { + const matchId = crypto.randomUUID(); + this.ctx.storage.sql.exec( + "INSERT INTO matches (id, name) VALUES (?, ?)", + matchId, name + ); + const child = this.env.GAME_MATCH.getByName(matchId); + await child.init(matchId, name); + return matchId; +} +``` + +### Location Hints + +Influence DO creation location for latency-sensitive apps: + +```typescript +const id = env.GAME.idFromName(gameId, { locationHint: "wnam" }); +``` + +Available hints: `wnam`, `enam`, `sam`, `weur`, `eeur`, `apac`, `oc`, `afr`, `me`. + +## Storage + +### SQLite (Recommended) + +Configure in wrangler: +```jsonc +{ "migrations": [{ "tag": "v1", "new_sqlite_classes": ["MyDO"] }] } +``` + +SQL API is synchronous: +```typescript +// Write +this.ctx.storage.sql.exec( + "INSERT INTO items (name, value) VALUES (?, ?)", + name, value +); + +// Read +const rows = this.ctx.storage.sql.exec<{ id: number; name: string }>( + "SELECT * FROM items WHERE name = ?", name +).toArray(); + +// Single row +const row = this.ctx.storage.sql.exec<{ count: number }>( + "SELECT COUNT(*) as count FROM items" +).one(); +``` + +### Migrations + +Use `PRAGMA user_version` for schema versioning: + +```typescript +constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + ctx.blockConcurrencyWhile(async () => this.migrate()); +} + +private async migrate() { + const version = this.ctx.storage.sql + .exec<{ user_version: number }>("PRAGMA user_version") + .one().user_version; + + if (version < 1) { + this.ctx.storage.sql.exec(` + CREATE TABLE IF NOT EXISTS items (id INTEGER PRIMARY KEY, data TEXT); + CREATE INDEX IF NOT EXISTS idx_items_data ON items(data); + PRAGMA user_version = 1; + `); + } + if (version < 2) { + this.ctx.storage.sql.exec(` + ALTER TABLE items ADD COLUMN created_at INTEGER; + PRAGMA user_version = 2; + `); + } +} +``` + +### State Types + +| Type | Speed | Persistence | Use Case | +|------|-------|-------------|----------| +| Class properties | Fastest | Lost on eviction | Caching, active connections | +| SQLite storage | Fast | Durable | Primary data | +| External (R2, D1) | Variable | Durable, cross-DO | Large files, shared data | + +**Rule**: Always persist critical state to SQLite first, then update in-memory cache. + +## Concurrency + +### Input/Output Gates + +Storage operations automatically block other requests (input gates). Responses wait for writes (output gates). + +```typescript +async increment(): Promise { + // Safe: input gates block interleaving during storage ops + const val = (await this.ctx.storage.get("count")) ?? 0; + await this.ctx.storage.put("count", val + 1); + return val + 1; +} +``` + +### Write Coalescing + +Multiple writes without `await` between them are batched atomically: + +```typescript +// ✅ Good: All three writes commit atomically +this.ctx.storage.sql.exec("UPDATE accounts SET balance = balance - ? WHERE id = ?", amount, fromId); +this.ctx.storage.sql.exec("UPDATE accounts SET balance = balance + ? WHERE id = ?", amount, toId); +this.ctx.storage.sql.exec("INSERT INTO transfers (from_id, to_id, amount) VALUES (?, ?, ?)", fromId, toId, amount); + +// ❌ Bad: await breaks coalescing +await this.ctx.storage.put("key1", val1); +await this.ctx.storage.put("key2", val2); // Separate transaction! +``` + +### Race Conditions with External I/O + +`fetch()` and other non-storage I/O allows interleaving: + +```typescript +// ⚠️ Race condition possible +async processItem(id: string) { + const item = await this.ctx.storage.get(`item:${id}`); + if (item?.status === "pending") { + await fetch("https://api.example.com/process"); // Other requests can run here! + await this.ctx.storage.put(`item:${id}`, { status: "completed" }); + } +} +``` + +**Solution**: Use optimistic locking (version numbers) or `transaction()`. + +### blockConcurrencyWhile() + +Blocks ALL concurrency. Use sparingly - only for initialization: + +```typescript +// ✅ Good: One-time init +constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + ctx.blockConcurrencyWhile(async () => this.migrate()); +} + +// ❌ Bad: On every request (kills throughput) +async handleRequest() { + await this.ctx.blockConcurrencyWhile(async () => { + // ~5ms = max 200 req/sec + }); +} +``` + +**Never** hold across external I/O (fetch, R2, KV). + +## RPC Methods + +Use RPC (compatibility date >= 2024-04-03) instead of fetch() handler: + +```typescript +export class ChatRoom extends DurableObject { + async sendMessage(userId: string, content: string): Promise { + // Public methods are RPC endpoints + const result = this.ctx.storage.sql.exec<{ id: number }>( + "INSERT INTO messages (user_id, content) VALUES (?, ?) RETURNING id", + userId, content + ); + return { id: result.one().id, userId, content }; + } +} + +// Caller +const stub = env.CHAT_ROOM.getByName(roomId); +const msg = await stub.sendMessage("user-123", "Hello!"); // Typed! +``` + +### Explicit init() Method + +DOs don't know their own ID. Pass identity explicitly: + +```typescript +async init(entityId: string, metadata: Metadata): Promise { + await this.ctx.storage.put("entityId", entityId); + await this.ctx.storage.put("metadata", metadata); +} +``` + +## Alarms + +One alarm per DO. `setAlarm()` replaces existing. + +```typescript +// Schedule +await this.ctx.storage.setAlarm(Date.now() + 60_000); + +// Handler +async alarm(): Promise { + const tasks = this.ctx.storage.sql.exec( + "SELECT * FROM tasks WHERE due_at <= ?", Date.now() + ).toArray(); + + for (const task of tasks) { + await this.processTask(task); + } + + // Reschedule if more work + const next = this.ctx.storage.sql.exec<{ due_at: number }>( + "SELECT MIN(due_at) as due_at FROM tasks WHERE due_at > ?", Date.now() + ).one(); + if (next?.due_at) { + await this.ctx.storage.setAlarm(next.due_at); + } +} + +// Get/Delete +const alarm = await this.ctx.storage.getAlarm(); +await this.ctx.storage.deleteAlarm(); +``` + +**Retry**: Alarms auto-retry on failure. Use idempotent handlers. + +## WebSockets (Hibernation API) + +```typescript +async fetch(request: Request): Promise { + const pair = new WebSocketPair(); + this.ctx.acceptWebSocket(pair[1]); + return new Response(null, { status: 101, webSocket: pair[0] }); +} + +async webSocketMessage(ws: WebSocket, message: string | ArrayBuffer) { + const data = JSON.parse(message as string); + // Handle message + ws.send(JSON.stringify({ type: "ack" })); +} + +async webSocketClose(ws: WebSocket, code: number, reason: string) { + // Cleanup +} + +// Broadcast +getWebSockets().forEach(ws => ws.send(JSON.stringify(payload))); +``` + +## Error Handling + +```typescript +async safeOperation(): Promise { + try { + return await this.riskyOperation(); + } catch (error) { + console.error("Operation failed:", error); + // Log to external service if needed + throw error; // Re-throw to signal failure to caller + } +} +``` + +**Note**: Uncaught exceptions may terminate the DO instance. In-memory state is lost, but SQLite storage persists. diff --git a/.agents/skills/durable-objects/references/testing.md b/.agents/skills/durable-objects/references/testing.md new file mode 100644 index 0000000..c4fb7e9 --- /dev/null +++ b/.agents/skills/durable-objects/references/testing.md @@ -0,0 +1,264 @@ +# Testing Durable Objects + +Use `@cloudflare/vitest-pool-workers` to test DOs inside the Workers runtime. + +## Setup + +### Install Dependencies + +```bash +npm i -D vitest@~3.2.0 @cloudflare/vitest-pool-workers +``` + +### vitest.config.ts + +```typescript +import { defineWorkersConfig } from "@cloudflare/vitest-pool-workers/config"; + +export default defineWorkersConfig({ + test: { + poolOptions: { + workers: { + wrangler: { configPath: "./wrangler.toml" }, + }, + }, + }, +}); +``` + +### TypeScript Config (test/tsconfig.json) + +```jsonc +{ + "extends": "../tsconfig.json", + "compilerOptions": { + "moduleResolution": "bundler", + "types": ["@cloudflare/vitest-pool-workers"] + }, + "include": ["./**/*.ts", "../src/worker-configuration.d.ts"] +} +``` + +### Environment Types (env.d.ts) + +```typescript +declare module "cloudflare:test" { + interface ProvidedEnv extends Env {} +} +``` + +## Unit Tests (Direct DO Access) + +```typescript +import { env } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; + +describe("Counter DO", () => { + it("should increment", async () => { + const stub = env.COUNTER.getByName("test-counter"); + + expect(await stub.increment()).toBe(1); + expect(await stub.increment()).toBe(2); + expect(await stub.getCount()).toBe(2); + }); + + it("isolates different instances", async () => { + const stub1 = env.COUNTER.getByName("counter-1"); + const stub2 = env.COUNTER.getByName("counter-2"); + + await stub1.increment(); + await stub1.increment(); + await stub2.increment(); + + expect(await stub1.getCount()).toBe(2); + expect(await stub2.getCount()).toBe(1); + }); +}); +``` + +## Integration Tests (HTTP via SELF) + +```typescript +import { SELF } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; + +describe("Worker HTTP", () => { + it("should increment via POST", async () => { + const res = await SELF.fetch("http://example.com?id=test", { + method: "POST", + }); + + expect(res.status).toBe(200); + const data = await res.json<{ count: number }>(); + expect(data.count).toBe(1); + }); + + it("should get count via GET", async () => { + await SELF.fetch("http://example.com?id=get-test", { method: "POST" }); + await SELF.fetch("http://example.com?id=get-test", { method: "POST" }); + + const res = await SELF.fetch("http://example.com?id=get-test"); + const data = await res.json<{ count: number }>(); + expect(data.count).toBe(2); + }); +}); +``` + +## Direct Internal Access + +Use `runInDurableObject()` to access instance internals and storage: + +```typescript +import { env, runInDurableObject } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; +import { Counter } from "../src"; + +describe("DO internals", () => { + it("can verify storage directly", async () => { + const stub = env.COUNTER.getByName("direct-test"); + await stub.increment(); + await stub.increment(); + + await runInDurableObject(stub, async (instance: Counter, state) => { + expect(instance).toBeInstanceOf(Counter); + + const result = state.storage.sql + .exec<{ value: number }>( + "SELECT value FROM counters WHERE name = ?", + "default" + ) + .one(); + expect(result.value).toBe(2); + }); + }); +}); +``` + +## List DO IDs + +```typescript +import { env, listDurableObjectIds } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; + +describe("DO listing", () => { + it("can list all IDs in namespace", async () => { + const id1 = env.COUNTER.idFromName("list-1"); + const id2 = env.COUNTER.idFromName("list-2"); + + await env.COUNTER.get(id1).increment(); + await env.COUNTER.get(id2).increment(); + + const ids = await listDurableObjectIds(env.COUNTER); + expect(ids.length).toBe(2); + expect(ids.some(id => id.equals(id1))).toBe(true); + expect(ids.some(id => id.equals(id2))).toBe(true); + }); +}); +``` + +## Testing Alarms + +Use `runDurableObjectAlarm()` to trigger alarms immediately: + +```typescript +import { env, runInDurableObject, runDurableObjectAlarm } from "cloudflare:test"; +import { describe, it, expect } from "vitest"; + +describe("DO alarms", () => { + it("can trigger alarms immediately", async () => { + const stub = env.COUNTER.getByName("alarm-test"); + await stub.increment(); + await stub.increment(); + expect(await stub.getCount()).toBe(2); + + // Schedule alarm + await runInDurableObject(stub, async (instance, state) => { + await state.storage.setAlarm(Date.now() + 60_000); + }); + + // Execute immediately without waiting + const ran = await runDurableObjectAlarm(stub); + expect(ran).toBe(true); + + // Verify alarm handler ran (if it resets counter) + expect(await stub.getCount()).toBe(0); + + // No alarm scheduled now + const ranAgain = await runDurableObjectAlarm(stub); + expect(ranAgain).toBe(false); + }); +}); +``` + +Example alarm handler: +```typescript +async alarm(): Promise { + this.ctx.storage.sql.exec("DELETE FROM counters"); +} +``` + +## Test Isolation + +Each test gets isolated storage automatically. DOs from one test don't affect others: + +```typescript +describe("Isolation", () => { + it("first test creates DO", async () => { + const stub = env.COUNTER.getByName("isolated"); + await stub.increment(); + expect(await stub.getCount()).toBe(1); + }); + + it("second test has fresh state", async () => { + const ids = await listDurableObjectIds(env.COUNTER); + expect(ids.length).toBe(0); // Previous test's DO is gone + + const stub = env.COUNTER.getByName("isolated"); + expect(await stub.getCount()).toBe(0); // Fresh instance + }); +}); +``` + +## SQLite Storage Testing + +```typescript +describe("SQLite", () => { + it("can verify SQL storage", async () => { + const stub = env.COUNTER.getByName("sqlite-test"); + await stub.increment("page-views"); + await stub.increment("page-views"); + await stub.increment("api-calls"); + + await runInDurableObject(stub, async (instance, state) => { + const rows = state.storage.sql + .exec<{ name: string; value: number }>( + "SELECT name, value FROM counters ORDER BY name" + ) + .toArray(); + + expect(rows).toEqual([ + { name: "api-calls", value: 1 }, + { name: "page-views", value: 2 }, + ]); + + expect(state.storage.sql.databaseSize).toBeGreaterThan(0); + }); + }); +}); +``` + +## Running Tests + +```bash +npx vitest # Watch mode +npx vitest run # Single run +``` + +package.json: +```json +{ + "scripts": { + "test": "vitest" + } +} +``` diff --git a/.agents/skills/durable-objects/references/workers.md b/.agents/skills/durable-objects/references/workers.md new file mode 100644 index 0000000..bd115a5 --- /dev/null +++ b/.agents/skills/durable-objects/references/workers.md @@ -0,0 +1,346 @@ +# Cloudflare Workers Best Practices + +High-level guidance for Workers that invoke Durable Objects. + +## Wrangler Configuration + +### wrangler.jsonc (Recommended) + +```jsonc +{ + "$schema": "node_modules/wrangler/config-schema.json", + "name": "my-worker", + "main": "src/index.ts", + "compatibility_date": "2024-12-01", + "compatibility_flags": ["nodejs_compat"], + + "durable_objects": { + "bindings": [ + { "name": "CHAT_ROOM", "class_name": "ChatRoom" }, + { "name": "USER_SESSION", "class_name": "UserSession" } + ] + }, + + "migrations": [ + { "tag": "v1", "new_sqlite_classes": ["ChatRoom", "UserSession"] } + ], + + // Environment variables + "vars": { + "ENVIRONMENT": "production" + }, + + // KV namespaces + "kv_namespaces": [ + { "binding": "CONFIG", "id": "abc123" } + ], + + // R2 buckets + "r2_buckets": [ + { "binding": "UPLOADS", "bucket_name": "my-uploads" } + ], + + // D1 databases + "d1_databases": [ + { "binding": "DB", "database_id": "xyz789" } + ] +} +``` + +### wrangler.toml (Alternative) + +```toml +name = "my-worker" +main = "src/index.ts" +compatibility_date = "2024-12-01" +compatibility_flags = ["nodejs_compat"] + +[[durable_objects.bindings]] +name = "CHAT_ROOM" +class_name = "ChatRoom" + +[[migrations]] +tag = "v1" +new_sqlite_classes = ["ChatRoom"] + +[vars] +ENVIRONMENT = "production" +``` + +## TypeScript Types + +### Environment Interface + +```typescript +// src/types.ts +import { ChatRoom } from "./durable-objects/chat-room"; +import { UserSession } from "./durable-objects/user-session"; + +export interface Env { + // Durable Objects + CHAT_ROOM: DurableObjectNamespace; + USER_SESSION: DurableObjectNamespace; + + // KV + CONFIG: KVNamespace; + + // R2 + UPLOADS: R2Bucket; + + // D1 + DB: D1Database; + + // Environment variables + ENVIRONMENT: string; + API_KEY: string; // From secrets +} +``` + +### Export Durable Object Classes + +```typescript +// src/index.ts +export { ChatRoom } from "./durable-objects/chat-room"; +export { UserSession } from "./durable-objects/user-session"; + +export default { + async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise { + // Worker handler + }, +}; +``` + +## Worker Handler Pattern + +```typescript +export default { + async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise { + const url = new URL(request.url); + + try { + // Route to appropriate handler + if (url.pathname.startsWith("/api/rooms")) { + return handleRooms(request, env); + } + if (url.pathname.startsWith("/api/users")) { + return handleUsers(request, env); + } + + return new Response("Not Found", { status: 404 }); + } catch (error) { + console.error("Request failed:", error); + return new Response("Internal Server Error", { status: 500 }); + } + }, +}; + +async function handleRooms(request: Request, env: Env): Promise { + const url = new URL(request.url); + const roomId = url.searchParams.get("room"); + + if (!roomId) { + return Response.json({ error: "Missing room parameter" }, { status: 400 }); + } + + const stub = env.CHAT_ROOM.getByName(roomId); + + if (request.method === "POST") { + const body = await request.json<{ userId: string; message: string }>(); + const result = await stub.sendMessage(body.userId, body.message); + return Response.json(result); + } + + const messages = await stub.getMessages(); + return Response.json(messages); +} +``` + +## Request Validation + +```typescript +import { z } from "zod"; + +const SendMessageSchema = z.object({ + userId: z.string().min(1), + message: z.string().min(1).max(1000), +}); + +async function handleSendMessage(request: Request, env: Env): Promise { + const body = await request.json(); + const result = SendMessageSchema.safeParse(body); + + if (!result.success) { + return Response.json( + { error: "Validation failed", details: result.error.issues }, + { status: 400 } + ); + } + + const stub = env.CHAT_ROOM.getByName(result.data.userId); + const message = await stub.sendMessage(result.data.userId, result.data.message); + return Response.json(message); +} +``` + +## Observability & Logging + +### Structured Logging + +```typescript +function log(level: "info" | "warn" | "error", message: string, data?: Record) { + console.log(JSON.stringify({ + level, + message, + timestamp: new Date().toISOString(), + ...data, + })); +} + +// Usage +log("info", "Request received", { path: url.pathname, method: request.method }); +log("error", "DO call failed", { roomId, error: String(error) }); +``` + +### Request Tracing + +```typescript +async function handleRequest(request: Request, env: Env): Promise { + const requestId = crypto.randomUUID(); + const startTime = Date.now(); + + try { + const response = await processRequest(request, env); + + log("info", "Request completed", { + requestId, + duration: Date.now() - startTime, + status: response.status, + }); + + return response; + } catch (error) { + log("error", "Request failed", { + requestId, + duration: Date.now() - startTime, + error: String(error), + }); + throw error; + } +} +``` + +### Tail Workers (Production) + +For production logging, use Tail Workers to forward logs: + +```jsonc +// wrangler.jsonc +{ + "tail_consumers": [ + { "service": "log-collector" } + ] +} +``` + +## Error Handling + +### Graceful DO Errors + +```typescript +async function callDO(stub: DurableObjectStub, method: string): Promise { + try { + const result = await stub.getMessages(); + return Response.json(result); + } catch (error) { + if (error instanceof Error) { + // DO threw an error + log("error", "DO operation failed", { error: error.message }); + return Response.json( + { error: "Service temporarily unavailable" }, + { status: 503 } + ); + } + throw error; + } +} +``` + +### Timeout Handling + +```typescript +async function withTimeout(promise: Promise, ms: number): Promise { + const timeout = new Promise((_, reject) => + setTimeout(() => reject(new Error("Timeout")), ms) + ); + return Promise.race([promise, timeout]); +} + +// Usage +const result = await withTimeout(stub.processData(data), 5000); +``` + +## CORS Handling + +```typescript +function corsHeaders(): HeadersInit { + return { + "Access-Control-Allow-Origin": "*", + "Access-Control-Allow-Methods": "GET, POST, PUT, DELETE, OPTIONS", + "Access-Control-Allow-Headers": "Content-Type, Authorization", + }; +} + +export default { + async fetch(request: Request, env: Env): Promise { + if (request.method === "OPTIONS") { + return new Response(null, { headers: corsHeaders() }); + } + + const response = await handleRequest(request, env); + + // Add CORS headers to response + const newHeaders = new Headers(response.headers); + Object.entries(corsHeaders()).forEach(([k, v]) => newHeaders.set(k, v)); + + return new Response(response.body, { + status: response.status, + headers: newHeaders, + }); + }, +}; +``` + +## Secrets Management + +Set secrets via wrangler CLI (not in config files): + +```bash +wrangler secret put API_KEY +wrangler secret put DATABASE_URL +``` + +Access in code: +```typescript +export default { + async fetch(request: Request, env: Env): Promise { + const apiKey = env.API_KEY; // From secret + // ... + }, +}; +``` + +## Development Commands + +```bash +# Local development +wrangler dev + +# Deploy +wrangler deploy + +# Tail logs +wrangler tail + +# List DOs +wrangler d1 execute DB --command "SELECT * FROM _cf_DO" +``` diff --git a/.agents/skills/wrangler/SKILL.md b/.agents/skills/wrangler/SKILL.md new file mode 100644 index 0000000..76d2030 --- /dev/null +++ b/.agents/skills/wrangler/SKILL.md @@ -0,0 +1,887 @@ +--- +name: wrangler +description: Cloudflare Workers CLI for deploying, developing, and managing Workers, KV, R2, D1, Vectorize, Hyperdrive, Workers AI, Containers, Queues, Workflows, Pipelines, and Secrets Store. Load before running wrangler commands to ensure correct syntax and best practices. +--- + +# Wrangler CLI + +Deploy, develop, and manage Cloudflare Workers and associated resources. + +## FIRST: Verify Wrangler Installation + +```bash +wrangler --version # Requires v4.x+ +``` + +If not installed: +```bash +npm install -D wrangler@latest +``` + +## Key Guidelines + +- **Use `wrangler.jsonc`**: Prefer JSON config over TOML. Newer features are JSON-only. +- **Set `compatibility_date`**: Use a recent date (within 30 days). Check https://developers.cloudflare.com/workers/configuration/compatibility-dates/ +- **Generate types after config changes**: Run `wrangler types` to update TypeScript bindings. +- **Local dev defaults to local storage**: Bindings use local simulation unless `remote: true`. +- **Validate config before deploy**: Run `wrangler check` to catch errors early. +- **Use environments for staging/prod**: Define `env.staging` and `env.production` in config. + +## Quick Start: New Worker + +```bash +# Initialize new project +npx wrangler init my-worker + +# Or with a framework +npx create-cloudflare@latest my-app +``` + +## Quick Reference: Core Commands + +| Task | Command | +|------|---------| +| Start local dev server | `wrangler dev` | +| Deploy to Cloudflare | `wrangler deploy` | +| Deploy dry run | `wrangler deploy --dry-run` | +| Generate TypeScript types | `wrangler types` | +| Validate configuration | `wrangler check` | +| View live logs | `wrangler tail` | +| Delete Worker | `wrangler delete` | +| Auth status | `wrangler whoami` | + +--- + +## Configuration (wrangler.jsonc) + +### Minimal Config + +```jsonc +{ + "$schema": "./node_modules/wrangler/config-schema.json", + "name": "my-worker", + "main": "src/index.ts", + "compatibility_date": "2026-01-01" +} +``` + +### Full Config with Bindings + +```jsonc +{ + "$schema": "./node_modules/wrangler/config-schema.json", + "name": "my-worker", + "main": "src/index.ts", + "compatibility_date": "2026-01-01", + "compatibility_flags": ["nodejs_compat_v2"], + + // Environment variables + "vars": { + "ENVIRONMENT": "production" + }, + + // KV Namespace + "kv_namespaces": [ + { "binding": "KV", "id": "" } + ], + + // R2 Bucket + "r2_buckets": [ + { "binding": "BUCKET", "bucket_name": "my-bucket" } + ], + + // D1 Database + "d1_databases": [ + { "binding": "DB", "database_name": "my-db", "database_id": "" } + ], + + // Workers AI (always remote) + "ai": { "binding": "AI" }, + + // Vectorize + "vectorize": [ + { "binding": "VECTOR_INDEX", "index_name": "my-index" } + ], + + // Hyperdrive + "hyperdrive": [ + { "binding": "HYPERDRIVE", "id": "" } + ], + + // Durable Objects + "durable_objects": { + "bindings": [ + { "name": "COUNTER", "class_name": "Counter" } + ] + }, + + // Cron triggers + "triggers": { + "crons": ["0 * * * *"] + }, + + // Environments + "env": { + "staging": { + "name": "my-worker-staging", + "vars": { "ENVIRONMENT": "staging" } + } + } +} +``` + +### Generate Types from Config + +```bash +# Generate worker-configuration.d.ts +wrangler types + +# Custom output path +wrangler types ./src/env.d.ts + +# Check types are up to date (CI) +wrangler types --check +``` + +--- + +## Local Development + +### Start Dev Server + +```bash +# Local mode (default) - uses local storage simulation +wrangler dev + +# With specific environment +wrangler dev --env staging + +# Force local-only (disable remote bindings) +wrangler dev --local + +# Remote mode - runs on Cloudflare edge (legacy) +wrangler dev --remote + +# Custom port +wrangler dev --port 8787 + +# Live reload for HTML changes +wrangler dev --live-reload + +# Test scheduled/cron handlers +wrangler dev --test-scheduled +# Then visit: http://localhost:8787/__scheduled +``` + +### Remote Bindings for Local Dev + +Use `remote: true` in binding config to connect to real resources while running locally: + +```jsonc +{ + "r2_buckets": [ + { "binding": "BUCKET", "bucket_name": "my-bucket", "remote": true } + ], + "ai": { "binding": "AI", "remote": true }, + "vectorize": [ + { "binding": "INDEX", "index_name": "my-index", "remote": true } + ] +} +``` + +**Recommended remote bindings**: AI (required), Vectorize, Browser Rendering, mTLS, Images. + +### Local Secrets + +Create `.dev.vars` for local development secrets: + +``` +API_KEY=local-dev-key +DATABASE_URL=postgres://localhost:5432/dev +``` + +--- + +## Deployment + +### Deploy Worker + +```bash +# Deploy to production +wrangler deploy + +# Deploy specific environment +wrangler deploy --env staging + +# Dry run (validate without deploying) +wrangler deploy --dry-run + +# Keep dashboard-set variables +wrangler deploy --keep-vars + +# Minify code +wrangler deploy --minify +``` + +### Manage Secrets + +```bash +# Set secret interactively +wrangler secret put API_KEY + +# Set from stdin +echo "secret-value" | wrangler secret put API_KEY + +# List secrets +wrangler secret list + +# Delete secret +wrangler secret delete API_KEY + +# Bulk secrets from JSON file +wrangler secret bulk secrets.json +``` + +### Versions and Rollback + +```bash +# List recent versions +wrangler versions list + +# View specific version +wrangler versions view + +# Rollback to previous version +wrangler rollback + +# Rollback to specific version +wrangler rollback +``` + +--- + +## KV (Key-Value Store) + +### Manage Namespaces + +```bash +# Create namespace +wrangler kv namespace create MY_KV + +# List namespaces +wrangler kv namespace list + +# Delete namespace +wrangler kv namespace delete --namespace-id +``` + +### Manage Keys + +```bash +# Put value +wrangler kv key put --namespace-id "key" "value" + +# Put with expiration (seconds) +wrangler kv key put --namespace-id "key" "value" --expiration-ttl 3600 + +# Get value +wrangler kv key get --namespace-id "key" + +# List keys +wrangler kv key list --namespace-id + +# Delete key +wrangler kv key delete --namespace-id "key" + +# Bulk put from JSON +wrangler kv bulk put --namespace-id data.json +``` + +### Config Binding + +```jsonc +{ + "kv_namespaces": [ + { "binding": "CACHE", "id": "" } + ] +} +``` + +--- + +## R2 (Object Storage) + +### Manage Buckets + +```bash +# Create bucket +wrangler r2 bucket create my-bucket + +# Create with location hint +wrangler r2 bucket create my-bucket --location wnam + +# List buckets +wrangler r2 bucket list + +# Get bucket info +wrangler r2 bucket info my-bucket + +# Delete bucket +wrangler r2 bucket delete my-bucket +``` + +### Manage Objects + +```bash +# Upload object +wrangler r2 object put my-bucket/path/file.txt --file ./local-file.txt + +# Download object +wrangler r2 object get my-bucket/path/file.txt + +# Delete object +wrangler r2 object delete my-bucket/path/file.txt +``` + +### Config Binding + +```jsonc +{ + "r2_buckets": [ + { "binding": "ASSETS", "bucket_name": "my-bucket" } + ] +} +``` + +--- + +## D1 (SQL Database) + +### Manage Databases + +```bash +# Create database +wrangler d1 create my-database + +# Create with location +wrangler d1 create my-database --location wnam + +# List databases +wrangler d1 list + +# Get database info +wrangler d1 info my-database + +# Delete database +wrangler d1 delete my-database +``` + +### Execute SQL + +```bash +# Execute SQL command (remote) +wrangler d1 execute my-database --remote --command "SELECT * FROM users" + +# Execute SQL file (remote) +wrangler d1 execute my-database --remote --file ./schema.sql + +# Execute locally +wrangler d1 execute my-database --local --command "SELECT * FROM users" +``` + +### Migrations + +```bash +# Create migration +wrangler d1 migrations create my-database create_users_table + +# List pending migrations +wrangler d1 migrations list my-database --local + +# Apply migrations locally +wrangler d1 migrations apply my-database --local + +# Apply migrations to remote +wrangler d1 migrations apply my-database --remote +``` + +### Export/Backup + +```bash +# Export schema and data +wrangler d1 export my-database --remote --output backup.sql + +# Export schema only +wrangler d1 export my-database --remote --output schema.sql --no-data +``` + +### Config Binding + +```jsonc +{ + "d1_databases": [ + { + "binding": "DB", + "database_name": "my-database", + "database_id": "", + "migrations_dir": "./migrations" + } + ] +} +``` + +--- + +## Vectorize (Vector Database) + +### Manage Indexes + +```bash +# Create index with dimensions +wrangler vectorize create my-index --dimensions 768 --metric cosine + +# Create with preset (auto-configures dimensions/metric) +wrangler vectorize create my-index --preset @cf/baai/bge-base-en-v1.5 + +# List indexes +wrangler vectorize list + +# Get index info +wrangler vectorize get my-index + +# Delete index +wrangler vectorize delete my-index +``` + +### Manage Vectors + +```bash +# Insert vectors from NDJSON file +wrangler vectorize insert my-index --file vectors.ndjson + +# Query vectors +wrangler vectorize query my-index --vector "[0.1, 0.2, ...]" --top-k 10 +``` + +### Config Binding + +```jsonc +{ + "vectorize": [ + { "binding": "SEARCH_INDEX", "index_name": "my-index" } + ] +} +``` + +--- + +## Hyperdrive (Database Accelerator) + +### Manage Configs + +```bash +# Create config +wrangler hyperdrive create my-hyperdrive \ + --connection-string "postgres://user:pass@host:5432/database" + +# List configs +wrangler hyperdrive list + +# Get config details +wrangler hyperdrive get + +# Update config +wrangler hyperdrive update --origin-password "new-password" + +# Delete config +wrangler hyperdrive delete +``` + +### Config Binding + +```jsonc +{ + "compatibility_flags": ["nodejs_compat_v2"], + "hyperdrive": [ + { "binding": "HYPERDRIVE", "id": "" } + ] +} +``` + +--- + +## Workers AI + +### List Models + +```bash +# List available models +wrangler ai models + +# List finetunes +wrangler ai finetune list +``` + +### Config Binding + +```jsonc +{ + "ai": { "binding": "AI" } +} +``` + +**Note**: Workers AI always runs remotely and incurs usage charges even in local dev. + +--- + +## Queues + +### Manage Queues + +```bash +# Create queue +wrangler queues create my-queue + +# List queues +wrangler queues list + +# Delete queue +wrangler queues delete my-queue + +# Add consumer to queue +wrangler queues consumer add my-queue my-worker + +# Remove consumer +wrangler queues consumer remove my-queue my-worker +``` + +### Config Binding + +```jsonc +{ + "queues": { + "producers": [ + { "binding": "MY_QUEUE", "queue": "my-queue" } + ], + "consumers": [ + { + "queue": "my-queue", + "max_batch_size": 10, + "max_batch_timeout": 30 + } + ] + } +} +``` + +--- + +## Containers + +### Build and Push Images + +```bash +# Build container image +wrangler containers build -t my-app:latest . + +# Build and push in one command +wrangler containers build -t my-app:latest . --push + +# Push existing image to Cloudflare registry +wrangler containers push my-app:latest +``` + +### Manage Containers + +```bash +# List containers +wrangler containers list + +# Get container info +wrangler containers info + +# Delete container +wrangler containers delete +``` + +### Manage Images + +```bash +# List images in registry +wrangler containers images list + +# Delete image +wrangler containers images delete my-app:latest +``` + +### Manage External Registries + +```bash +# List configured registries +wrangler containers registries list + +# Configure external registry (e.g., ECR) +wrangler containers registries configure \ + --public-credential + +# Delete registry configuration +wrangler containers registries delete +``` + +--- + +## Workflows + +### Manage Workflows + +```bash +# List workflows +wrangler workflows list + +# Describe workflow +wrangler workflows describe my-workflow + +# Trigger workflow instance +wrangler workflows trigger my-workflow + +# Trigger with parameters +wrangler workflows trigger my-workflow --params '{"key": "value"}' + +# Delete workflow +wrangler workflows delete my-workflow +``` + +### Manage Workflow Instances + +```bash +# List instances +wrangler workflows instances list my-workflow + +# Describe instance +wrangler workflows instances describe my-workflow + +# Terminate instance +wrangler workflows instances terminate my-workflow +``` + +### Config Binding + +```jsonc +{ + "workflows": [ + { + "binding": "MY_WORKFLOW", + "name": "my-workflow", + "class_name": "MyWorkflow" + } + ] +} +``` + +--- + +## Pipelines + +### Manage Pipelines + +```bash +# Create pipeline +wrangler pipelines create my-pipeline --r2 my-bucket + +# List pipelines +wrangler pipelines list + +# Show pipeline details +wrangler pipelines show my-pipeline + +# Update pipeline +wrangler pipelines update my-pipeline --batch-max-mb 100 + +# Delete pipeline +wrangler pipelines delete my-pipeline +``` + +### Config Binding + +```jsonc +{ + "pipelines": [ + { "binding": "MY_PIPELINE", "pipeline": "my-pipeline" } + ] +} +``` + +--- + +## Secrets Store + +### Manage Stores + +```bash +# Create store +wrangler secrets-store store create my-store + +# List stores +wrangler secrets-store store list + +# Delete store +wrangler secrets-store store delete +``` + +### Manage Secrets in Store + +```bash +# Add secret to store +wrangler secrets-store secret put my-secret + +# List secrets in store +wrangler secrets-store secret list + +# Get secret +wrangler secrets-store secret get my-secret + +# Delete secret from store +wrangler secrets-store secret delete my-secret +``` + +### Config Binding + +```jsonc +{ + "secrets_store_secrets": [ + { + "binding": "MY_SECRET", + "store_id": "", + "secret_name": "my-secret" + } + ] +} +``` + +--- + +## Pages (Frontend Deployment) + +```bash +# Create Pages project +wrangler pages project create my-site + +# Deploy directory to Pages +wrangler pages deploy ./dist + +# Deploy with specific branch +wrangler pages deploy ./dist --branch main + +# List deployments +wrangler pages deployment list --project-name my-site +``` + +--- + +## Observability + +### Tail Logs + +```bash +# Stream live logs +wrangler tail + +# Tail specific Worker +wrangler tail my-worker + +# Filter by status +wrangler tail --status error + +# Filter by search term +wrangler tail --search "error" + +# JSON output +wrangler tail --format json +``` + +### Config Logging + +```jsonc +{ + "observability": { + "enabled": true, + "head_sampling_rate": 1 + } +} +``` + +--- + +## Testing + +### Local Testing with Vitest + +```bash +npm install -D @cloudflare/vitest-pool-workers vitest +``` + +`vitest.config.ts`: +```typescript +import { defineWorkersConfig } from "@cloudflare/vitest-pool-workers/config"; + +export default defineWorkersConfig({ + test: { + poolOptions: { + workers: { + wrangler: { configPath: "./wrangler.jsonc" }, + }, + }, + }, +}); +``` + +### Test Scheduled Events + +```bash +# Enable in dev +wrangler dev --test-scheduled + +# Trigger via HTTP +curl http://localhost:8787/__scheduled +``` + +--- + +## Troubleshooting + +### Common Issues + +| Issue | Solution | +|-------|----------| +| `command not found: wrangler` | Install: `npm install -D wrangler` | +| Auth errors | Run `wrangler login` | +| Config validation errors | Run `wrangler check` | +| Type errors after config change | Run `wrangler types` | +| Local storage not persisting | Check `.wrangler/state` directory | +| Binding undefined in Worker | Verify binding name matches config exactly | + +### Debug Commands + +```bash +# Check auth status +wrangler whoami + +# Validate config +wrangler check + +# View config schema +wrangler docs configuration +``` + +--- + +## Best Practices + +1. **Version control `wrangler.jsonc`**: Treat as source of truth for Worker config. +2. **Use automatic provisioning**: Omit resource IDs for auto-creation on deploy. +3. **Run `wrangler types` in CI**: Add to build step to catch binding mismatches. +4. **Use environments**: Separate staging/production with `env.staging`, `env.production`. +5. **Set `compatibility_date`**: Update quarterly to get new runtime features. +6. **Use `.dev.vars` for local secrets**: Never commit secrets to config. +7. **Test locally first**: `wrangler dev` with local bindings before deploying. +8. **Use `--dry-run` before major deploys**: Validate changes without deployment. From c4f54e6994e31d404edc3d63d59af2f90f6ba06c Mon Sep 17 00:00:00 2001 From: Alexander Eichhorn Date: Mon, 9 Feb 2026 19:57:22 +0100 Subject: [PATCH 2/5] Add configurable rolling rate limits --- README.md | 14 +- package.json | 3 +- src/durable-objects/app-rate-limiter.ts | 198 ++++++++++++++++++++++++ src/index.ts | 56 ++++++- worker-configuration.d.ts | 83 ++++++---- wrangler.toml | 24 ++- 6 files changed, 347 insertions(+), 31 deletions(-) create mode 100644 src/durable-objects/app-rate-limiter.ts diff --git a/README.md b/README.md index d449460..500206e 100644 --- a/README.md +++ b/README.md @@ -21,8 +21,8 @@ This architecture ensures streams work with the client's IP address and location ### Prerequisites -- [Cloudflare Workers](https://workers.cloudflare.com/) account -- [Wrangler CLI](https://developers.cloudflare.com/workers/wrangler/) installed +- [Cloudflare Workers](https://workers.cloudflare.com/) account +- [Wrangler CLI](https://developers.cloudflare.com/workers/wrangler/) installed ### Deploy to Cloudflare Workers @@ -36,3 +36,13 @@ For local development: ```bash npm run dev ``` + +## Rate Limiting + +Requests are limited per app ID (`X-AppID-v1`) using a Durable Object with rolling windows. + +Configure limits using environment variables in `wrangler.toml`: + +- `RATE_LIMIT_DAILY_REQUESTS`: max requests per app ID in rolling 24h window. +- `RATE_LIMIT_WEEKLY_REQUESTS`: max requests per app ID in rolling 7d window. +- `RATE_LIMIT_BUCKET_SECONDS`: bucket size used for rolling-window aggregation (default: `300`). diff --git a/package.json b/package.json index 4b237ea..a8e11ad 100644 --- a/package.json +++ b/package.json @@ -7,7 +7,8 @@ "dev": "wrangler dev", "start": "wrangler dev", "test": "vitest", - "cf-typegen": "wrangler types" + "cf-typegen": "wrangler types", + "typecheck": "tsc --noEmit" }, "devDependencies": { "@cloudflare/vitest-pool-workers": "^0.8.19", diff --git a/src/durable-objects/app-rate-limiter.ts b/src/durable-objects/app-rate-limiter.ts new file mode 100644 index 0000000..bb8a01c --- /dev/null +++ b/src/durable-objects/app-rate-limiter.ts @@ -0,0 +1,198 @@ +import { DurableObject } from 'cloudflare:workers'; + +interface AdmitRequest { + cost?: number; + nowMs?: number; +} + +export interface RateLimitDecision { + allowed: boolean; + limitDaily: number; + limitWeekly: number; + remainingDaily: number; + remainingWeekly: number; + retryAfterSeconds: number; +} + +interface RateLimitPolicy { + dailyLimit: number; + weeklyLimit: number; + bucketSizeMs: number; +} + +const DAY_MS = 24 * 60 * 60 * 1000; +const WEEK_MS = 7 * DAY_MS; +const DEFAULT_DAILY_LIMIT = 5000; +const DEFAULT_WEEKLY_LIMIT = 20000; +const DEFAULT_BUCKET_SECONDS = 300; + +export class AppRateLimiter extends DurableObject { + private readonly sql = this.ctx.storage.sql; + + constructor(ctx: DurableObjectState, env: Env) { + super(ctx, env); + this.ctx.blockConcurrencyWhile(async () => { + this.initializeSchema(); + }); + } + + async fetch(request: Request): Promise { + const url = new URL(request.url); + if (request.method !== 'POST' || url.pathname !== '/admit') { + return new Response('Not found', { status: 404 }); + } + + const payload = (await request.json().catch(() => null)) as AdmitRequest | null; + if (!payload) { + return new Response('Invalid JSON payload', { status: 400 }); + } + + const nowMs = Number.isFinite(payload.nowMs) ? Number(payload.nowMs) : Date.now(); + const requestedCost = Number.isFinite(payload.cost) ? Number(payload.cost) : 1; + const cost = Math.max(1, Math.floor(requestedCost)); + + const decision = this.admit(cost, nowMs); + return Response.json(decision); + } + + private initializeSchema() { + this.sql.exec(` + CREATE TABLE IF NOT EXISTS request_buckets ( + bucket_start INTEGER PRIMARY KEY, + count INTEGER NOT NULL + ) + `); + } + + private admit(cost: number, nowMs: number): RateLimitDecision { + const policy = this.getPolicy(); + const currentBucketStart = Math.floor(nowMs / policy.bucketSizeMs) * policy.bucketSizeMs; + const dayWindowStart = nowMs - DAY_MS; + const weekWindowStart = nowMs - WEEK_MS; + + this.cleanupOldBuckets(weekWindowStart, policy.bucketSizeMs); + + const usedDaily = this.getUsageSince(dayWindowStart); + const usedWeekly = this.getUsageSince(weekWindowStart); + const nextDaily = usedDaily + cost; + const nextWeekly = usedWeekly + cost; + + if (nextDaily > policy.dailyLimit || nextWeekly > policy.weeklyLimit) { + const retryAfterSeconds = this.calculateRetryAfterSeconds({ + nowMs, + dayWindowStart, + weekWindowStart, + dayExceeded: nextDaily > policy.dailyLimit, + weekExceeded: nextWeekly > policy.weeklyLimit, + bucketSizeMs: policy.bucketSizeMs, + }); + + return { + allowed: false, + limitDaily: policy.dailyLimit, + limitWeekly: policy.weeklyLimit, + remainingDaily: Math.max(0, policy.dailyLimit - usedDaily), + remainingWeekly: Math.max(0, policy.weeklyLimit - usedWeekly), + retryAfterSeconds, + }; + } + + this.sql.exec( + ` + INSERT INTO request_buckets (bucket_start, count) + VALUES (?1, ?2) + ON CONFLICT(bucket_start) DO UPDATE SET count = count + excluded.count + `, + currentBucketStart, + cost + ); + + return { + allowed: true, + limitDaily: policy.dailyLimit, + limitWeekly: policy.weeklyLimit, + remainingDaily: Math.max(0, policy.dailyLimit - nextDaily), + remainingWeekly: Math.max(0, policy.weeklyLimit - nextWeekly), + retryAfterSeconds: 0, + }; + } + + private cleanupOldBuckets(weekWindowStart: number, bucketSizeMs: number) { + // Keep one extra bucket outside the 7-day range for stable boundary behavior. + this.sql.exec('DELETE FROM request_buckets WHERE bucket_start <= ?1', weekWindowStart - bucketSizeMs); + } + + private getUsageSince(windowStart: number): number { + const row = this.sql + .exec<{ total: number | null }>('SELECT COALESCE(SUM(count), 0) AS total FROM request_buckets WHERE bucket_start > ?1', windowStart) + .one(); + + return Number(row.total ?? 0); + } + + private calculateRetryAfterSeconds(params: { + nowMs: number; + dayWindowStart: number; + weekWindowStart: number; + dayExceeded: boolean; + weekExceeded: boolean; + bucketSizeMs: number; + }): number { + const retries: number[] = []; + + if (params.dayExceeded) { + retries.push(this.getWindowRetryAfterSeconds(params.dayWindowStart, DAY_MS, params.nowMs, params.bucketSizeMs)); + } + + if (params.weekExceeded) { + retries.push(this.getWindowRetryAfterSeconds(params.weekWindowStart, WEEK_MS, params.nowMs, params.bucketSizeMs)); + } + + if (retries.length === 0) { + return Math.max(1, Math.ceil(params.bucketSizeMs / 1000)); + } + + return Math.max(...retries); + } + + private getWindowRetryAfterSeconds(windowStart: number, windowMs: number, nowMs: number, bucketSizeMs: number): number { + const row = this.sql + .exec<{ bucket_start: number }>( + 'SELECT bucket_start FROM request_buckets WHERE bucket_start > ?1 ORDER BY bucket_start ASC LIMIT 1', + windowStart + ) + .toArray()[0]; + + if (!row) { + return Math.max(1, Math.ceil(bucketSizeMs / 1000)); + } + + const retryAt = Number(row.bucket_start) + windowMs; + return Math.max(1, Math.ceil((retryAt - nowMs) / 1000)); + } + + private getPolicy(): RateLimitPolicy { + const dailyLimit = this.parsePositiveInt(this.env.RATE_LIMIT_DAILY_REQUESTS, DEFAULT_DAILY_LIMIT); + const weeklyLimit = this.parsePositiveInt(this.env.RATE_LIMIT_WEEKLY_REQUESTS, DEFAULT_WEEKLY_LIMIT); + const bucketSeconds = this.parsePositiveInt(this.env.RATE_LIMIT_BUCKET_SECONDS, DEFAULT_BUCKET_SECONDS); + + return { + dailyLimit, + weeklyLimit, + bucketSizeMs: bucketSeconds * 1000, + }; + } + + private parsePositiveInt(value: string | undefined, fallback: number): number { + if (!value) { + return fallback; + } + + const parsed = Number.parseInt(value, 10); + if (!Number.isFinite(parsed) || parsed <= 0) { + return fallback; + } + + return parsed; + } +} diff --git a/src/index.ts b/src/index.ts index 78bd727..7f7b539 100644 --- a/src/index.ts +++ b/src/index.ts @@ -12,6 +12,11 @@ */ import { YouTubeService } from './youtube/service'; +export { AppRateLimiter } from './durable-objects/app-rate-limiter'; +import { RateLimitDecision } from './durable-objects/app-rate-limiter'; + +const APP_ID_HEADER = 'X-AppID-v1'; +const APP_ID_MAX_LENGTH = 128; export default { async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise { @@ -22,11 +27,21 @@ export default { console.log(`User Agent: ${userAgent}`); // Log the App ID header for debugging purposes - const appID = request.headers.get('X-AppID-v1') ?? 'unknown'; + const appID = normalizeAppID(request.headers.get(APP_ID_HEADER)); console.log(`App ID: ${appID}`); // Only handle GET /v1?videoID=... as WebSocket upgrades if (url.pathname === '/v1' && request.headers.get('Upgrade') === 'websocket') { + try { + const decision = await checkRateLimit(appID, env); + if (!decision.allowed) { + return buildRateLimitResponse(decision); + } + } catch (error) { + console.error('Rate limiter unavailable:', error); + return new Response('Rate limiter unavailable', { status: 503 }); + } + const videoID = url.searchParams.get('videoID'); if (!videoID) { return new Response('Missing videoID', { status: 400 }); @@ -46,3 +61,42 @@ export default { return new Response('Not found', { status: 404 }); }, } satisfies ExportedHandler; + +function normalizeAppID(rawAppID: string | null): string { + const trimmed = rawAppID?.trim(); + if (!trimmed) { + return 'unknown'; + } + + return trimmed.slice(0, APP_ID_MAX_LENGTH); +} + +async function checkRateLimit(appID: string, env: Env): Promise { + const objectID = env.APP_RATE_LIMITER.idFromName(appID); + const stub = env.APP_RATE_LIMITER.get(objectID); + + const response = await stub.fetch('https://limiter/admit', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ cost: 1, nowMs: Date.now() }), + }); + + if (!response.ok) { + throw new Error(`Rate limit DO request failed with status ${response.status}`); + } + + return (await response.json()) as RateLimitDecision; +} + +function buildRateLimitResponse(decision: RateLimitDecision): Response { + return new Response('Too many requests', { + status: 429, + headers: { + 'Retry-After': decision.retryAfterSeconds.toString(), + 'X-RateLimit-Limit-Day': decision.limitDaily.toString(), + 'X-RateLimit-Limit-Week': decision.limitWeekly.toString(), + 'X-RateLimit-Remaining-Day': decision.remainingDaily.toString(), + 'X-RateLimit-Remaining-Week': decision.remainingWeekly.toString(), + }, + }); +} diff --git a/worker-configuration.d.ts b/worker-configuration.d.ts index 8d0ec36..287a9e7 100644 --- a/worker-configuration.d.ts +++ b/worker-configuration.d.ts @@ -1,8 +1,12 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types` (hash: 869ac3b4ce0f52ba3b2e0bc70c49089e) -// Runtime types generated with workerd@1.20250428.0 2025-05-03 +// Generated by Wrangler by running `wrangler types` (hash: 92e805438630fb89e1fb7b6e0ae6d95e) +// Runtime types generated with workerd@1.20250604.0 2025-05-03 declare namespace Cloudflare { interface Env { + RATE_LIMIT_DAILY_REQUESTS: "5000"; + RATE_LIMIT_WEEKLY_REQUESTS: "20000"; + RATE_LIMIT_BUCKET_SECONDS: "300"; + APP_RATE_LIMITER: DurableObjectNamespace; } } interface Env extends Cloudflare.Env {} @@ -87,7 +91,7 @@ interface Console { clear(): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/count_static) */ count(label?: string): void; - /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/countreset_static) */ + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/countReset_static) */ countReset(label?: string): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/debug_static) */ debug(...data: any[]): void; @@ -99,9 +103,9 @@ interface Console { error(...data: any[]): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/group_static) */ group(...data: any[]): void; - /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/groupcollapsed_static) */ + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/groupCollapsed_static) */ groupCollapsed(...data: any[]): void; - /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/groupend_static) */ + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/groupEnd_static) */ groupEnd(): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/info_static) */ info(...data: any[]): void; @@ -111,9 +115,9 @@ interface Console { table(tabularData?: any, properties?: string[]): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/time_static) */ time(label?: string): void; - /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/timeend_static) */ + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/timeEnd_static) */ timeEnd(label?: string): void; - /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/timelog_static) */ + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/timeLog_static) */ timeLog(label?: string, ...data: any[]): void; timeStamp(label?: string): void; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/console/trace_static) */ @@ -288,25 +292,25 @@ declare function dispatchEvent(event: WorkerGlobalScopeEventMap[keyof WorkerGlob declare function btoa(data: string): string; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/atob) */ declare function atob(data: string): string; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/setTimeout) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/setTimeout) */ declare function setTimeout(callback: (...args: any[]) => void, msDelay?: number): number; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/setTimeout) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/setTimeout) */ declare function setTimeout(callback: (...args: Args) => void, msDelay?: number, ...args: Args): number; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/clearTimeout) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/clearTimeout) */ declare function clearTimeout(timeoutId: number | null): void; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/setInterval) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/setInterval) */ declare function setInterval(callback: (...args: any[]) => void, msDelay?: number): number; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/setInterval) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/setInterval) */ declare function setInterval(callback: (...args: Args) => void, msDelay?: number, ...args: Args): number; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/clearInterval) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/clearInterval) */ declare function clearInterval(timeoutId: number | null): void; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/queueMicrotask) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/queueMicrotask) */ declare function queueMicrotask(task: Function): void; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/structuredClone) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/structuredClone) */ declare function structuredClone(value: T, options?: StructuredSerializeOptions): T; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/reportError) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/reportError) */ declare function reportError(error: any): void; -/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/fetch) */ +/* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Window/fetch) */ declare function fetch(input: RequestInfo | URL, init?: RequestInit): Promise; declare const self: ServiceWorkerGlobalScope; /** @@ -785,6 +789,7 @@ declare class Blob { slice(start?: number, end?: number, type?: string): Blob; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Blob/arrayBuffer) */ arrayBuffer(): Promise; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Blob/bytes) */ bytes(): Promise; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Blob/text) */ text(): Promise; @@ -1079,10 +1084,15 @@ interface TextEncoderEncodeIntoResult { */ declare class ErrorEvent extends Event { constructor(type: string, init?: ErrorEventErrorEventInit); + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/ErrorEvent/filename) */ get filename(): string; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/ErrorEvent/message) */ get message(): string; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/ErrorEvent/lineno) */ get lineno(): number; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/ErrorEvent/colno) */ get colno(): number; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/ErrorEvent/error) */ get error(): any; } interface ErrorEventErrorEventInit { @@ -1256,6 +1266,7 @@ declare abstract class Body { get bodyUsed(): boolean; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Request/arrayBuffer) */ arrayBuffer(): Promise; + /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Request/bytes) */ bytes(): Promise; /* [MDN Reference](https://developer.mozilla.org/docs/Web/API/Request/text) */ text(): Promise; @@ -1367,7 +1378,11 @@ interface Request> e * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Request/integrity) */ integrity: string; - /* Returns a boolean indicating whether or not request can outlive the global in which it was created. */ + /** + * Returns a boolean indicating whether or not request can outlive the global in which it was created. + * + * [MDN Reference](https://developer.mozilla.org/docs/Web/API/Request/keepalive) + */ keepalive: boolean; /** * Returns the cache mode associated with request, which is a string indicating how the request will interact with the browser's cache when fetching. @@ -1567,6 +1582,7 @@ interface R2ObjectBody extends R2Object { get body(): ReadableStream; get bodyUsed(): boolean; arrayBuffer(): Promise; + bytes(): Promise; text(): Promise; json(): Promise; blob(): Promise; @@ -5174,6 +5190,7 @@ declare module 'cloudflare:workers' { export type WorkflowSleepDuration = `${number} ${WorkflowDurationLabel}${'s' | ''}` | number; export type WorkflowDelayDuration = WorkflowSleepDuration; export type WorkflowTimeoutDuration = WorkflowSleepDuration; + export type WorkflowRetentionDuration = WorkflowSleepDuration; export type WorkflowBackoff = 'constant' | 'linear' | 'exponential'; export type WorkflowStepConfig = { retries?: { @@ -5304,6 +5321,7 @@ declare namespace TailStream { readonly type: "onset"; readonly dispatchNamespace?: string; readonly entrypoint?: string; + readonly executionModel: string; readonly scriptName?: string; readonly scriptTags?: string[]; readonly scriptVersion?: ScriptVersion; @@ -5321,8 +5339,8 @@ declare namespace TailStream { } interface SpanOpen { readonly type: "spanOpen"; - readonly op?: string; - readonly info?: FetchEventInfo | JsRpcEventInfo | Attribute[]; + readonly name: string; + readonly info?: FetchEventInfo | JsRpcEventInfo | Attributes; } interface SpanClose { readonly type: "spanClose"; @@ -5346,7 +5364,7 @@ declare namespace TailStream { } interface Return { readonly type: "return"; - readonly info?: FetchResponseInfo | Attribute[]; + readonly info?: FetchResponseInfo; } interface Link { readonly type: "link"; @@ -5356,21 +5374,23 @@ declare namespace TailStream { readonly spanId: string; } interface Attribute { - readonly type: "attribute"; readonly name: string; - readonly value: string | string[] | boolean | boolean[] | number | number[]; + readonly value: string | string[] | boolean | boolean[] | number | number[] | bigint | bigint[]; + } + interface Attributes { + readonly type: "attributes"; + readonly info: Attribute[]; } - type Mark = DiagnosticChannelEvent | Exception | Log | Return | Link | Attribute[]; interface TailEvent { readonly traceId: string; readonly invocationId: string; readonly spanId: string; readonly timestamp: Date; readonly sequence: number; - readonly event: Onset | Outcome | Hibernate | SpanOpen | SpanClose | Mark; + readonly event: Onset | Outcome | Hibernate | SpanOpen | SpanClose | DiagnosticChannelEvent | Exception | Log | Return | Link | Attributes; } type TailEventHandler = (event: TailEvent) => void | Promise; - type TailEventHandlerName = "onset" | "outcome" | "hibernate" | "spanOpen" | "spanClose" | "diagnosticChannel" | "exception" | "log" | "return" | "link" | "attribute"; + type TailEventHandlerName = "outcome" | "hibernate" | "spanOpen" | "spanClose" | "diagnosticChannel" | "exception" | "log" | "return" | "link" | "attributes"; type TailEventHandlerObject = Record; type TailEventHandlerType = TailEventHandler | TailEventHandlerObject; } @@ -5684,6 +5704,9 @@ declare abstract class Workflow { */ public createBatch(batch: WorkflowInstanceCreateOptions[]): Promise; } +type WorkflowDurationLabel = 'second' | 'minute' | 'hour' | 'day' | 'week' | 'month' | 'year'; +type WorkflowSleepDuration = `${number} ${WorkflowDurationLabel}${'s' | ''}` | number; +type WorkflowRetentionDuration = WorkflowSleepDuration; interface WorkflowInstanceCreateOptions { /** * An id for your Workflow instance. Must be unique within the Workflow. @@ -5693,6 +5716,14 @@ interface WorkflowInstanceCreateOptions { * The event payload the Workflow instance is triggered with */ params?: PARAMS; + /** + * The retention policy for Workflow instance. + * Defaults to the maximum retention period available for the owner's account. + */ + retention?: { + successRetention?: WorkflowRetentionDuration; + errorRetention?: WorkflowRetentionDuration; + }; } type InstanceStatus = { status: 'queued' // means that instance is waiting to be started (see concurrency limits) diff --git a/wrangler.toml b/wrangler.toml index ab06058..1dd1563 100644 --- a/wrangler.toml +++ b/wrangler.toml @@ -14,11 +14,33 @@ enabled = true # Variables will be loaded from .dev.vars for local development [vars] +RATE_LIMIT_DAILY_REQUESTS = "5000" +RATE_LIMIT_WEEKLY_REQUESTS = "20000" +RATE_LIMIT_BUCKET_SECONDS = "300" +[[durable_objects.bindings]] +name = "APP_RATE_LIMITER" +class_name = "AppRateLimiter" + +[[migrations]] +tag = "v1" +new_sqlite_classes = [ "AppRateLimiter" ] # Production environment -------------------------------------------------------- [env.production] -name = "youtubekit-server-production" \ No newline at end of file +name = "youtubekit-server-production" +[env.production.vars] +RATE_LIMIT_DAILY_REQUESTS = "5000" +RATE_LIMIT_WEEKLY_REQUESTS = "20000" +RATE_LIMIT_BUCKET_SECONDS = "300" + +[[env.production.durable_objects.bindings]] +name = "APP_RATE_LIMITER" +class_name = "AppRateLimiter" + +[[env.production.migrations]] +tag = "v1" +new_sqlite_classes = [ "AppRateLimiter" ] From 8b8cde3a7cf7c79e3bb77f617dcb9da37261d0dd Mon Sep 17 00:00:00 2001 From: Alexander Eichhorn Date: Mon, 9 Feb 2026 20:13:21 +0100 Subject: [PATCH 3/5] Clarify anchored rate limits --- README.md | 11 +- src/durable-objects/app-rate-limiter.ts | 147 ++++++++++++++---------- worker-configuration.d.ts | 3 +- wrangler.toml | 2 - 4 files changed, 93 insertions(+), 70 deletions(-) diff --git a/README.md b/README.md index 500206e..e65934b 100644 --- a/README.md +++ b/README.md @@ -39,10 +39,13 @@ npm run dev ## Rate Limiting -Requests are limited per app ID (`X-AppID-v1`) using a Durable Object with rolling windows. +Requests are limited per app ID (`X-AppID-v1`) using a Durable Object with anchored windows: + +- daily window starts at first request and runs for 24h +- weekly window starts at first request and runs for 7d +- when a window expires, the next request starts a new window Configure limits using environment variables in `wrangler.toml`: -- `RATE_LIMIT_DAILY_REQUESTS`: max requests per app ID in rolling 24h window. -- `RATE_LIMIT_WEEKLY_REQUESTS`: max requests per app ID in rolling 7d window. -- `RATE_LIMIT_BUCKET_SECONDS`: bucket size used for rolling-window aggregation (default: `300`). +- `RATE_LIMIT_DAILY_REQUESTS`: max requests per app ID in each 24h window. +- `RATE_LIMIT_WEEKLY_REQUESTS`: max requests per app ID in each 7d window. diff --git a/src/durable-objects/app-rate-limiter.ts b/src/durable-objects/app-rate-limiter.ts index bb8a01c..b7e14fe 100644 --- a/src/durable-objects/app-rate-limiter.ts +++ b/src/durable-objects/app-rate-limiter.ts @@ -17,14 +17,19 @@ export interface RateLimitDecision { interface RateLimitPolicy { dailyLimit: number; weeklyLimit: number; - bucketSizeMs: number; } const DAY_MS = 24 * 60 * 60 * 1000; const WEEK_MS = 7 * DAY_MS; const DEFAULT_DAILY_LIMIT = 5000; const DEFAULT_WEEKLY_LIMIT = 20000; -const DEFAULT_BUCKET_SECONDS = 300; + +interface CounterState { + dayWindowStartMs: number; + dayCount: number; + weekWindowStartMs: number; + weekCount: number; +} export class AppRateLimiter extends DurableObject { private readonly sql = this.ctx.storage.sql; @@ -57,129 +62,147 @@ export class AppRateLimiter extends DurableObject { private initializeSchema() { this.sql.exec(` - CREATE TABLE IF NOT EXISTS request_buckets ( - bucket_start INTEGER PRIMARY KEY, - count INTEGER NOT NULL + CREATE TABLE IF NOT EXISTS limiter_state ( + id INTEGER PRIMARY KEY CHECK (id = 1), + day_window_start_ms INTEGER NOT NULL, + day_count INTEGER NOT NULL, + week_window_start_ms INTEGER NOT NULL, + week_count INTEGER NOT NULL ) `); } private admit(cost: number, nowMs: number): RateLimitDecision { const policy = this.getPolicy(); - const currentBucketStart = Math.floor(nowMs / policy.bucketSizeMs) * policy.bucketSizeMs; - const dayWindowStart = nowMs - DAY_MS; - const weekWindowStart = nowMs - WEEK_MS; + const state = this.getOrCreateState(nowMs); + const nextState = this.rollExpiredWindows(state, nowMs); - this.cleanupOldBuckets(weekWindowStart, policy.bucketSizeMs); - - const usedDaily = this.getUsageSince(dayWindowStart); - const usedWeekly = this.getUsageSince(weekWindowStart); - const nextDaily = usedDaily + cost; - const nextWeekly = usedWeekly + cost; + const nextDaily = nextState.dayCount + cost; + const nextWeekly = nextState.weekCount + cost; if (nextDaily > policy.dailyLimit || nextWeekly > policy.weeklyLimit) { const retryAfterSeconds = this.calculateRetryAfterSeconds({ nowMs, - dayWindowStart, - weekWindowStart, + dayWindowStartMs: nextState.dayWindowStartMs, + weekWindowStartMs: nextState.weekWindowStartMs, dayExceeded: nextDaily > policy.dailyLimit, weekExceeded: nextWeekly > policy.weeklyLimit, - bucketSizeMs: policy.bucketSizeMs, }); return { allowed: false, limitDaily: policy.dailyLimit, limitWeekly: policy.weeklyLimit, - remainingDaily: Math.max(0, policy.dailyLimit - usedDaily), - remainingWeekly: Math.max(0, policy.weeklyLimit - usedWeekly), + remainingDaily: Math.max(0, policy.dailyLimit - nextState.dayCount), + remainingWeekly: Math.max(0, policy.weeklyLimit - nextState.weekCount), retryAfterSeconds, }; } this.sql.exec( ` - INSERT INTO request_buckets (bucket_start, count) - VALUES (?1, ?2) - ON CONFLICT(bucket_start) DO UPDATE SET count = count + excluded.count + INSERT INTO limiter_state ( + id, + day_window_start_ms, + day_count, + week_window_start_ms, + week_count + ) + VALUES (1, ?1, ?2, ?3, ?4) + ON CONFLICT(id) DO UPDATE SET + day_window_start_ms = excluded.day_window_start_ms, + day_count = excluded.day_count, + week_window_start_ms = excluded.week_window_start_ms, + week_count = excluded.week_count `, - currentBucketStart, - cost + nextState.dayWindowStartMs, + nextDaily, + nextState.weekWindowStartMs, + nextWeekly ); return { allowed: true, limitDaily: policy.dailyLimit, limitWeekly: policy.weeklyLimit, - remainingDaily: Math.max(0, policy.dailyLimit - nextDaily), - remainingWeekly: Math.max(0, policy.weeklyLimit - nextWeekly), + remainingDaily: Math.max(0, policy.dailyLimit - (nextState.dayCount + cost)), + remainingWeekly: Math.max(0, policy.weeklyLimit - (nextState.weekCount + cost)), retryAfterSeconds: 0, }; } - private cleanupOldBuckets(weekWindowStart: number, bucketSizeMs: number) { - // Keep one extra bucket outside the 7-day range for stable boundary behavior. - this.sql.exec('DELETE FROM request_buckets WHERE bucket_start <= ?1', weekWindowStart - bucketSizeMs); + private getOrCreateState(nowMs: number): CounterState { + const row = this.sql + .exec<{ + day_window_start_ms: number; + day_count: number; + week_window_start_ms: number; + week_count: number; + }>('SELECT day_window_start_ms, day_count, week_window_start_ms, week_count FROM limiter_state WHERE id = 1') + .toArray()[0]; + + if (!row) { + return { + dayWindowStartMs: nowMs, + dayCount: 0, + weekWindowStartMs: nowMs, + weekCount: 0, + }; + } + + return { + dayWindowStartMs: Number(row.day_window_start_ms), + dayCount: Number(row.day_count), + weekWindowStartMs: Number(row.week_window_start_ms), + weekCount: Number(row.week_count), + }; } - private getUsageSince(windowStart: number): number { - const row = this.sql - .exec<{ total: number | null }>('SELECT COALESCE(SUM(count), 0) AS total FROM request_buckets WHERE bucket_start > ?1', windowStart) - .one(); + private rollExpiredWindows(state: CounterState, nowMs: number): CounterState { + const nextState = { ...state }; - return Number(row.total ?? 0); + if (nowMs >= nextState.dayWindowStartMs + DAY_MS) { + nextState.dayWindowStartMs = nowMs; + nextState.dayCount = 0; + } + + if (nowMs >= nextState.weekWindowStartMs + WEEK_MS) { + nextState.weekWindowStartMs = nowMs; + nextState.weekCount = 0; + } + + return nextState; } private calculateRetryAfterSeconds(params: { nowMs: number; - dayWindowStart: number; - weekWindowStart: number; + dayWindowStartMs: number; + weekWindowStartMs: number; dayExceeded: boolean; weekExceeded: boolean; - bucketSizeMs: number; }): number { const retries: number[] = []; if (params.dayExceeded) { - retries.push(this.getWindowRetryAfterSeconds(params.dayWindowStart, DAY_MS, params.nowMs, params.bucketSizeMs)); + const dayResetAt = params.dayWindowStartMs + DAY_MS; + retries.push(Math.max(1, Math.ceil((dayResetAt - params.nowMs) / 1000))); } if (params.weekExceeded) { - retries.push(this.getWindowRetryAfterSeconds(params.weekWindowStart, WEEK_MS, params.nowMs, params.bucketSizeMs)); - } - - if (retries.length === 0) { - return Math.max(1, Math.ceil(params.bucketSizeMs / 1000)); - } - - return Math.max(...retries); - } - - private getWindowRetryAfterSeconds(windowStart: number, windowMs: number, nowMs: number, bucketSizeMs: number): number { - const row = this.sql - .exec<{ bucket_start: number }>( - 'SELECT bucket_start FROM request_buckets WHERE bucket_start > ?1 ORDER BY bucket_start ASC LIMIT 1', - windowStart - ) - .toArray()[0]; - - if (!row) { - return Math.max(1, Math.ceil(bucketSizeMs / 1000)); + const weekResetAt = params.weekWindowStartMs + WEEK_MS; + retries.push(Math.max(1, Math.ceil((weekResetAt - params.nowMs) / 1000))); } - const retryAt = Number(row.bucket_start) + windowMs; - return Math.max(1, Math.ceil((retryAt - nowMs) / 1000)); + return retries.length > 0 ? Math.max(...retries) : 1; } private getPolicy(): RateLimitPolicy { const dailyLimit = this.parsePositiveInt(this.env.RATE_LIMIT_DAILY_REQUESTS, DEFAULT_DAILY_LIMIT); const weeklyLimit = this.parsePositiveInt(this.env.RATE_LIMIT_WEEKLY_REQUESTS, DEFAULT_WEEKLY_LIMIT); - const bucketSeconds = this.parsePositiveInt(this.env.RATE_LIMIT_BUCKET_SECONDS, DEFAULT_BUCKET_SECONDS); return { dailyLimit, weeklyLimit, - bucketSizeMs: bucketSeconds * 1000, }; } diff --git a/worker-configuration.d.ts b/worker-configuration.d.ts index 287a9e7..bdfcbc9 100644 --- a/worker-configuration.d.ts +++ b/worker-configuration.d.ts @@ -1,11 +1,10 @@ /* eslint-disable */ -// Generated by Wrangler by running `wrangler types` (hash: 92e805438630fb89e1fb7b6e0ae6d95e) +// Generated by Wrangler by running `wrangler types` (hash: 1fcade365cd6b1ad09020ecd1b8a7043) // Runtime types generated with workerd@1.20250604.0 2025-05-03 declare namespace Cloudflare { interface Env { RATE_LIMIT_DAILY_REQUESTS: "5000"; RATE_LIMIT_WEEKLY_REQUESTS: "20000"; - RATE_LIMIT_BUCKET_SECONDS: "300"; APP_RATE_LIMITER: DurableObjectNamespace; } } diff --git a/wrangler.toml b/wrangler.toml index 1dd1563..9f2b837 100644 --- a/wrangler.toml +++ b/wrangler.toml @@ -16,7 +16,6 @@ enabled = true [vars] RATE_LIMIT_DAILY_REQUESTS = "5000" RATE_LIMIT_WEEKLY_REQUESTS = "20000" -RATE_LIMIT_BUCKET_SECONDS = "300" [[durable_objects.bindings]] name = "APP_RATE_LIMITER" @@ -35,7 +34,6 @@ name = "youtubekit-server-production" [env.production.vars] RATE_LIMIT_DAILY_REQUESTS = "5000" RATE_LIMIT_WEEKLY_REQUESTS = "20000" -RATE_LIMIT_BUCKET_SECONDS = "300" [[env.production.durable_objects.bindings]] name = "APP_RATE_LIMITER" From ffe6353ae6a32f8f4d7f2852ea6ab784e85050f3 Mon Sep 17 00:00:00 2001 From: Alexander Eichhorn Date: Mon, 9 Feb 2026 20:18:29 +0100 Subject: [PATCH 4/5] Explain rate limiter window change --- src/durable-objects/app-rate-limiter.ts | 21 +++++---------------- src/index.ts | 15 ++------------- 2 files changed, 7 insertions(+), 29 deletions(-) diff --git a/src/durable-objects/app-rate-limiter.ts b/src/durable-objects/app-rate-limiter.ts index b7e14fe..635cd41 100644 --- a/src/durable-objects/app-rate-limiter.ts +++ b/src/durable-objects/app-rate-limiter.ts @@ -41,23 +41,12 @@ export class AppRateLimiter extends DurableObject { }); } - async fetch(request: Request): Promise { - const url = new URL(request.url); - if (request.method !== 'POST' || url.pathname !== '/admit') { - return new Response('Not found', { status: 404 }); - } - - const payload = (await request.json().catch(() => null)) as AdmitRequest | null; - if (!payload) { - return new Response('Invalid JSON payload', { status: 400 }); - } - - const nowMs = Number.isFinite(payload.nowMs) ? Number(payload.nowMs) : Date.now(); - const requestedCost = Number.isFinite(payload.cost) ? Number(payload.cost) : 1; + admit(payload?: AdmitRequest): RateLimitDecision { + const nowMs = Number.isFinite(payload?.nowMs) ? Number(payload?.nowMs) : Date.now(); + const requestedCost = Number.isFinite(payload?.cost) ? Number(payload?.cost) : 1; const cost = Math.max(1, Math.floor(requestedCost)); - const decision = this.admit(cost, nowMs); - return Response.json(decision); + return this.evaluateAdmission(cost, nowMs); } private initializeSchema() { @@ -72,7 +61,7 @@ export class AppRateLimiter extends DurableObject { `); } - private admit(cost: number, nowMs: number): RateLimitDecision { + private evaluateAdmission(cost: number, nowMs: number): RateLimitDecision { const policy = this.getPolicy(); const state = this.getOrCreateState(nowMs); const nextState = this.rollExpiredWindows(state, nowMs); diff --git a/src/index.ts b/src/index.ts index 7f7b539..3152521 100644 --- a/src/index.ts +++ b/src/index.ts @@ -73,19 +73,8 @@ function normalizeAppID(rawAppID: string | null): string { async function checkRateLimit(appID: string, env: Env): Promise { const objectID = env.APP_RATE_LIMITER.idFromName(appID); - const stub = env.APP_RATE_LIMITER.get(objectID); - - const response = await stub.fetch('https://limiter/admit', { - method: 'POST', - headers: { 'content-type': 'application/json' }, - body: JSON.stringify({ cost: 1, nowMs: Date.now() }), - }); - - if (!response.ok) { - throw new Error(`Rate limit DO request failed with status ${response.status}`); - } - - return (await response.json()) as RateLimitDecision; + const rate_limiter = env.APP_RATE_LIMITER.get(objectID); + return await rate_limiter.admit({ cost: 1, nowMs: Date.now() }); } function buildRateLimitResponse(decision: RateLimitDecision): Response { From 6a7ed360d155c64615b6d95646b23cbbd3173b20 Mon Sep 17 00:00:00 2001 From: Alexander Eichhorn Date: Mon, 9 Feb 2026 20:30:19 +0100 Subject: [PATCH 5/5] Log rate limit rejects --- src/index.ts | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/src/index.ts b/src/index.ts index 3152521..d552aa0 100644 --- a/src/index.ts +++ b/src/index.ts @@ -35,6 +35,18 @@ export default { try { const decision = await checkRateLimit(appID, env); if (!decision.allowed) { + console.warn( + 'Rate limit rejected request', + JSON.stringify({ + appID, + path: url.pathname, + limitDay: decision.limitDaily, + remainingDay: decision.remainingDaily, + limitWeek: decision.limitWeekly, + remainingWeek: decision.remainingWeekly, + retryAfterSeconds: decision.retryAfterSeconds, + }) + ); return buildRateLimitResponse(decision); } } catch (error) {