diff --git a/bun.lock b/bun.lock index 64367ba33..ebb446c1a 100644 --- a/bun.lock +++ b/bun.lock @@ -361,6 +361,7 @@ "dependencies": { "@corbits/artifact-ui": "workspace:*", "@corbits/artifacts": "github:corbitsdev/corbits-artifacts#81049ed24a64e927498c7238bda6ffa66b63d2ab", + "@corbits/collections": "workspace:*", "@corbits/folded-runs": "workspace:*", "@intx/crypto": "0.3.0", "@intx/db": "workspace:*", @@ -3609,6 +3610,8 @@ "@babel/helper-compilation-targets/semver": ["semver@6.3.1", "", { "bin": { "semver": "bin/semver.js" } }, "sha512-BR7VvDCVHO+q2xBEWskxS6DJE1qRnb7DxzUrogb71CWoSficBxYsiAGd+Kl0mmq/MprG9yArRkyrQxTO6XjMzA=="], + "@corbits/memory-hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], + "@esbuild-kit/core-utils/esbuild": ["esbuild@0.18.20", "", { "optionalDependencies": { "@esbuild/android-arm": "0.18.20", "@esbuild/android-arm64": "0.18.20", "@esbuild/android-x64": "0.18.20", "@esbuild/darwin-arm64": "0.18.20", "@esbuild/darwin-x64": "0.18.20", "@esbuild/freebsd-arm64": "0.18.20", "@esbuild/freebsd-x64": "0.18.20", "@esbuild/linux-arm": "0.18.20", "@esbuild/linux-arm64": "0.18.20", "@esbuild/linux-ia32": "0.18.20", "@esbuild/linux-loong64": "0.18.20", "@esbuild/linux-mips64el": "0.18.20", "@esbuild/linux-ppc64": "0.18.20", "@esbuild/linux-riscv64": "0.18.20", "@esbuild/linux-s390x": "0.18.20", "@esbuild/linux-x64": "0.18.20", "@esbuild/netbsd-x64": "0.18.20", "@esbuild/openbsd-x64": "0.18.20", "@esbuild/sunos-x64": "0.18.20", "@esbuild/win32-arm64": "0.18.20", "@esbuild/win32-ia32": "0.18.20", "@esbuild/win32-x64": "0.18.20" }, "bin": { "esbuild": "bin/esbuild" } }, "sha512-ceqxoedUrcayh7Y7ZX6NdbbDzGROiyVBgC4PriJThBKSVPWnnFHZAkfI1lJT8QFkOwH4qOS2SJkS4wvpGl8BpA=="], "@eslint-community/eslint-utils/eslint-visitor-keys": ["eslint-visitor-keys@3.4.3", "", {}, "sha512-wpc+LXeiyiisxPlEkUzU6svyS1frIO3Mgxj1fdy7Pm8Ygzguax2N3Fa/D/ag1WqbOprdI+uY6wMUl8/a2G+iag=="], @@ -3631,6 +3634,8 @@ "@typescript-eslint/eslint-plugin/ignore": ["ignore@7.0.6", "", {}, "sha512-BAg6QkE8W+TuQLrrw0Ugr7HegXduRuuj8/ti2kSOc+jz1dmx8/WNcjr6XGnq5YpDWxFwwaavqD0+jIUOKelTsw=="], + "@workbench/hub/@corbits/memory": ["@corbits/memory@github:corbitsdev/corbits-memory#9e6f213", { "dependencies": { "@intx/agent": "0.2.2", "@intx/authz": "0.2.2", "@intx/hub-api": "0.2.2", "@intx/log": "0.2.2", "@intx/workflow": "0.2.2", "arktype": "^2.1.29", "drizzle-orm": "^0.45.1", "hono": "^4.9.0", "hono-openapi": "^1.3.1", "postgres": "^3.4.7" } }, "corbitsdev-corbits-memory-9e6f213", "sha512-utnM4ZT2zmslcPXYWAAqxlDNLcpGsXFiTOtj8h7+OXnhCP0Eaw8yl25+yCTyHpvt3jcdeG4h5uFsSj7ou0BZCA=="], + "ajv-formats/ajv": ["ajv@8.20.0", "", { "dependencies": { "fast-deep-equal": "^3.1.3", "fast-uri": "^3.0.1", "json-schema-traverse": "^1.0.0", "require-from-string": "^2.0.2" } }, "sha512-Thbli+OlOj+iMPYFBVBfJ3OmCAnaSyNn4M1vz9T6Gka5Jt9ba/HIR56joy65tY6kx/FCF5VXNB819Y7/GUrBGA=="], "better-call/@better-auth/utils": ["@better-auth/utils@0.5.0", "", { "dependencies": { "@noble/hashes": "^2.0.1" } }, "sha512-BL8W4EfIZFwlu0r54m3v1ztjDhu6dDe/amLTm0xybmbZaNgYUqhD3SjpAsnq0q8YD6/ki4iwIgxJNLP/N3TxiA=="], diff --git a/packages/artifacts-hub/package.json b/packages/artifacts-hub/package.json index d46424d77..8841166f6 100644 --- a/packages/artifacts-hub/package.json +++ b/packages/artifacts-hub/package.json @@ -15,6 +15,7 @@ "dependencies": { "@corbits/artifact-ui": "workspace:*", "@corbits/artifacts": "github:corbitsdev/corbits-artifacts#81049ed24a64e927498c7238bda6ffa66b63d2ab", + "@corbits/collections": "workspace:*", "@corbits/folded-runs": "workspace:*", "@intx/crypto": "0.3.0", "@intx/db": "workspace:*", diff --git a/packages/artifacts-hub/src/workflow-routes.rate-limiter.test.ts b/packages/artifacts-hub/src/workflow-routes.rate-limiter.test.ts new file mode 100644 index 000000000..fd1454a6b --- /dev/null +++ b/packages/artifacts-hub/src/workflow-routes.rate-limiter.test.ts @@ -0,0 +1,67 @@ +import { describe, expect, test } from "bun:test"; + +import { createRunCreateRateLimiter, RATE_WINDOW_MS } from "./workflow-routes"; + +describe("createRunCreateRateLimiter", () => { + test("evicts an idle run's entry instead of holding it for the process lifetime", () => { + let clock = 0; + const limiter = createRunCreateRateLimiter(3, () => clock); + + for (let i = 0; i < 500; i++) { + limiter.allow(`run-${i}`); + } + expect(limiter.trackedRunCount).toBe(500); + + // Every one of those runs has gone idle for a full window: their + // entries should be reclaimed, not carried forever. + clock += RATE_WINDOW_MS; + limiter.allow("run-fresh"); + expect(limiter.trackedRunCount).toBe(1); + }); + + test("a caller cannot exceed the rate by exploiting eviction timing", () => { + let clock = 0; + const maxPerWindow = 3; + const limiter = createRunCreateRateLimiter(maxPerWindow, () => clock); + const runId = "run-under-test"; + + for (let i = 0; i < maxPerWindow; i++) { + expect(limiter.allow(runId)).toBe(true); + } + expect(limiter.allow(runId)).toBe(false); + + // Advance right up to (but not past) the window boundary: the + // entry's TTL must not have lapsed yet, so the earlier timestamps + // are still counted and the limit still holds. + clock += RATE_WINDOW_MS - 1; + expect(limiter.allow(runId)).toBe(false); + + // Advance past the window: the original timestamps are now stale + // and a fresh budget opens up, which is the intended sliding-window + // behavior rather than a leak of the earlier eviction. + clock += 1; + for (let i = 0; i < maxPerWindow; i++) { + expect(limiter.allow(runId)).toBe(true); + } + expect(limiter.allow(runId)).toBe(false); + }); + + test("does not let concurrently active runs go unbounded by other idle ones", () => { + let clock = 0; + const limiter = createRunCreateRateLimiter(3, () => clock); + + // A burst of one-shot runs, each idle immediately after. + for (let i = 0; i < 200; i++) { + limiter.allow(`idle-run-${i}`); + clock += 1; + } + + // One run stays continuously active, well past when the idle runs' + // entries should have expired. + clock += RATE_WINDOW_MS; + const activeRunId = "active-run"; + limiter.allow(activeRunId); + + expect(limiter.trackedRunCount).toBe(1); + }); +}); diff --git a/packages/artifacts-hub/src/workflow-routes.ts b/packages/artifacts-hub/src/workflow-routes.ts index 735e9a7ea..2792c7559 100644 --- a/packages/artifacts-hub/src/workflow-routes.ts +++ b/packages/artifacts-hub/src/workflow-routes.ts @@ -22,6 +22,7 @@ * itself never holds a database handle. */ import { type } from "arktype"; +import { createExpiringMap } from "@corbits/collections"; import { ARTIFACT_UPLOAD_POLICY, anonymousIdentity, @@ -58,8 +59,8 @@ const MAX_ARTIFACT_CONTENT_CHARS = 64_000; // A finalized turn can legitimately persist a handful of artifacts in // one burst; 30/minute per run comfortably covers that while still // catching a runaway loop before it floods Library storage. -const MAX_CREATES_PER_RUN_PER_MINUTE = 30; -const RATE_WINDOW_MS = 60_000; +export const MAX_CREATES_PER_RUN_PER_MINUTE = 30; +export const RATE_WINDOW_MS = 60_000; // Same per-file ceiling the tenant Library's own `POST /upload` enforces // (`MAX_UPLOAD_BYTES`) — one number for "how big a file artifact may be" @@ -80,13 +81,30 @@ export const MAX_WORKFLOW_BINARY_BYTES = MAX_UPLOAD_BYTES; * only what it personally handled — a known fail-open gap, not a * fail-closed one, so it under-limits rather than wrongly rejecting a * caller a sibling replica hasn't seen yet. + * + * The per-run entry lives in a `createExpiringMap` (CL-7243) rather than + * a plain `Map`, so a run's entry is reclaimed once it goes idle instead + * of staying resident for the rest of the hub process's uptime. The TTL + * is exactly `RATE_WINDOW_MS`: every `allow()` call re-`set`s the entry, + * refreshing its expiry, so an entry can only lapse after a full window + * with no calls for that run — by which point every timestamp it held + * has already aged out of the sliding window's own `cutoff` filter below. + * A shorter TTL could evict an entry (and thus its still-in-window + * timestamps) before the window's filter would have dropped them, + * silently resetting a caller's quota early; this TTL can't. */ -function createRunCreateRateLimiter(maxPerWindow: number) { - const timestampsByRunId = new Map(); +export function createRunCreateRateLimiter( + maxPerWindow: number, + now: () => number = Date.now, +) { + const timestampsByRunId = createExpiringMap({ + ttlMs: RATE_WINDOW_MS, + now, + }); return { allow(runId: string): boolean { - const now = Date.now(); - const cutoff = now - RATE_WINDOW_MS; + const at = now(); + const cutoff = at - RATE_WINDOW_MS; const recent = (timestampsByRunId.get(runId) ?? []).filter( (timestamp) => timestamp > cutoff, ); @@ -94,10 +112,14 @@ function createRunCreateRateLimiter(maxPerWindow: number) { timestampsByRunId.set(runId, recent); return false; } - recent.push(now); + recent.push(at); timestampsByRunId.set(runId, recent); return true; }, + /** Live entry count, for tests asserting the map stays bounded. */ + get trackedRunCount(): number { + return timestampsByRunId.size; + }, }; }