From 747bf8c503103a5b713dd89f85e3f1b99d710a65 Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Tue, 22 Sep 2026 20:21:02 -0700 Subject: [PATCH 1/3] Limit anonymous traffic by label set --- README.md | 5 +- src/admin.ts | 30 +++++++--- src/docs.ts | 22 ++++--- src/features/admin/client.tsx | 6 +- src/features/admin/reasons.ts | 1 + src/index.ts | 110 +++++++++++++++++++++++++++++++--- src/openapi.ts | 6 +- src/pages.ts | 24 +++++--- src/privacy.ts | 36 ++++++++--- src/report.ts | 16 ++++- src/wellknown.ts | 2 +- test/admin-fixture.ts | 4 +- test/record.test.ts | 15 ++--- tests/spending.e2e.test.ts | 66 +++++++++++++++++++- wrangler.example.toml | 3 + 15 files changed, 285 insertions(+), 61 deletions(-) diff --git a/README.md b/README.md index cc7f879..5af1bee 100644 --- a/README.md +++ b/README.md @@ -266,7 +266,10 @@ The report includes status, reason, count and mean latency for each provider; `/alerts` shows current incidents even when their notification is suppressed. The existing 15-minute alert check warns when at least three attempts fail and failures exceed 5% for either provider. Counts account for Analytics Engine -sampling. No input text, labels, caller identifiers or upstream messages are stored. +sampling. No input text, caller identifiers or upstream messages are stored. +Successful simple and multi-label classifier names are retained for 90 days in +a separate aggregate KV record that is not joined to a caller or source text; +per-request analytics contain only their keyed fingerprint. `AI_GATEWAY_DISABLED = "true"` in `wrangler.example.toml` keeps production on TypeSafe directly after the gateway repeatedly returned 429 on September 19. diff --git a/src/admin.ts b/src/admin.ts index fd97e99..6a4ecf0 100644 --- a/src/admin.ts +++ b/src/admin.ts @@ -12,10 +12,10 @@ * does not answer at all, because a page that says what it is is a page worth * attacking. * - * None of these figures is anybody's data. The caller column is a day-scoped - * hash and the label column a keyed fingerprint, both made in src/privacy.ts - * before anything is written, so there is nothing here to leak and nothing to - * hand over if somebody asks. + * The caller column is a day-scoped hash and the label column a keyed + * fingerprint, both made in src/privacy.ts before per-request analytics are + * written. A separate 90-day registry resolves successful classifier + * fingerprints to aggregate label names; it contains no caller or source text. */ import type { Env } from "./index"; @@ -23,6 +23,7 @@ import { sql } from "./report"; import { reasonText } from "./features/admin/reasons"; import { esc, BASE_CSS } from "./ui"; import { secretEquals, deriveSigningKey, hmacHex } from "./secrets"; +import { recordedClassifierLabels } from "./privacy"; const DATASET = "classifier_events"; // The __Secure- prefix is enforced by the browser, not by us: it refuses to @@ -392,6 +393,21 @@ async function load(env: Env, range: RangeKey) { ), ]); + const classifierNames = new Map(); + const fingerprints = [...new Set( + [...topLabels, ...failLabels] + .map((row) => String(row.labels ?? "")) + .filter(Boolean), + )]; + await Promise.all(fingerprints.map(async (fingerprint) => { + const labels = await recordedClassifierLabels(env.STATS, fingerprint); + if (labels.length) classifierNames.set(fingerprint, labels.join(" · ")); + })); + const nameClassifiers = (rows: Row[]) => rows.map((row) => ({ + ...row, + label_names: classifierNames.get(String(row.labels ?? "")) ?? "—", + })); + return { totals, series, @@ -400,12 +416,12 @@ async function load(env: Env, range: RangeKey) { byCountry, byStatus, byClient, - topLabels, + topLabels: nameClassifiers(topLabels), visitors, labelSets, byReason, byAgent, - failLabels, + failLabels: nameClassifiers(failLabels), dimensionTraffic, dimensionCallers, dimensionSeries, @@ -541,7 +557,7 @@ ${table("Failure causes", d.byReason, ["reason", "status", "agent", "requests"], ${table("Dimension outcomes", d.dimensionTraffic, ["status", "reason", "requests", "uncertain", "fallback"], "dimensionTraffic")} ${table("Clients", d.byAgent, ["agent", "requests", "classifications"], "byAgent")} ${table("Countries", d.byCountry, ["country", "requests"], "byCountry")} -${table("Busiest classifiers", d.topLabels, ["labels", "requests", "classifications", "usd"], "topLabels")} +${table("Busiest classifiers", d.topLabels, ["label_names", "labels", "requests", "classifications", "usd"], "topLabels")} `; return shell( "admin · classifier.dev", diff --git a/src/docs.ts b/src/docs.ts index f3286a0..69ca5cd 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -475,8 +475,11 @@ LIMITS 3,000 per minute and 20,000 per day; the smart tier 200 per minute and 2,000 per day. A batch must fit the remaining quota in full. Public smart requests accept at most 200 inputs; larger batches return 400 so callers - can split them. Pro workspaces allow 30,000/minute and 200,000/day on - fast, 2,000/minute and 20,000/day on smart, shared across keys and agents. + can split them. Anonymous traffic also shares a 5,000/minute and 50,000/day + allowance across every caller using the same label set. Rotating IPs does + not reset it; workspace, operator and partner keys bypass it. Pro workspaces + allow 30,000/minute and 200,000/day on fast, 2,000/minute and 20,000/day on + smart, shared across keys and agents. Pro, operator and partner keys have a 1,000-input ceiling. Workspace keys use the workspace credit balance and share workspace quotas. Free workspaces have the same ceilings as public access. Current plans are at @@ -509,9 +512,9 @@ ERRORS 402 insufficient workspace balance for inference 403 the key is inactive or the workspace cannot authorize usage 404 not_found - 429 rate_limit_minute, rate_limit_day, chunklaya_busy, with Retry-After; on the free - tier the body also carries upgrade, the URL of the plan that lifts - the limit (https://classifier.dev/pricing) + 429 rate_limit_minute, rate_limit_day, label_set_limit, chunklaya_busy, + with Retry-After; on the free tier the body also carries upgrade, the + URL of the plan that lifts the limit (https://classifier.dev/pricing) 502 typesafe or typesafe_ when the decision model failed; openrouter_, chain_exhausted or timeout when the fallback chain did; upstream_other. Retry with backoff. @@ -541,10 +544,11 @@ ${roadmapDoc()} PRIVACY The text you send is never stored or logged. It goes to the model provider - for the classification and nowhere else. What is recorded: a keyed - fingerprint of the label set, never the labels, plus the tier, the model, - the latency, the status and a coarse country. The usage counts are built - from those. + for the classification and nowhere else. Per-request analytics record a + keyed fingerprint of the label set, plus the tier, model, latency, status and + coarse country. Successful simple and multi-label classifier names are also + kept for 90 days in a separate aggregate registry with no caller identity or + source text. The usage counts are built from the fingerprinted records. Built by @michael_chomsky — https://x.com/michael_chomsky diff --git a/src/features/admin/client.tsx b/src/features/admin/client.tsx index 9203d4c..6d8a1ee 100644 --- a/src/features/admin/client.tsx +++ b/src/features/admin/client.tsx @@ -699,13 +699,14 @@ function Dashboard({ data: d, range }: { data: AdminData; range: RangeKey }) {
= { input_too_long: "an input was over 32,000 characters", rate_limit_minute: "per-minute rate limit", rate_limit_day: "daily rate limit", + label_set_limit: "anonymous label-set limit", bad_dimensions: "invalid dimension definitions or conflicting options", too_many_decisions: "too many item × dimension decisions", dimension_context_too_large: "input and dimension exceed the model context", diff --git a/src/index.ts b/src/index.ts index ec456df..3ef752a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -75,6 +75,9 @@ export interface Env extends LayaEnv, SpendingEnv { LIMITER: DurableObjectNamespace; AE: AnalyticsEngineDataset; REPORT_KEY: string; + /** Shared anonymous decision ceilings for one stable label-set fingerprint. */ + FREE_LABEL_RPM?: string; + FREE_LABEL_DAILY?: string; /** Gates /admin. Set with `npx wrangler secret put ADMIN_PASSWORD`. */ ADMIN_PASSWORD?: string; /** Signs the /admin session cookie. Random, and unrelated to the password. */ @@ -1098,14 +1101,82 @@ function upstreamReason(msg: string): ErrorCode { } /** - * Name a label set so we can count distinct classifiers without keeping one. - * This used to be the labels themselves, lowercased and joined, which meant - * every caller's wording sat in analytics and on the dashboard for 90 days. + * Name a label set without putting its words in per-request analytics. The + * separate registry can resolve successful classifiers for aggregate review. */ function classifierId(env: Env, labels: string[]) { return labelFingerprint(env, labels); } +const LABEL_LIMITS = { rpm: 5_000, daily: 50_000 } as const; +const CLASSIFIER_TTL = 90 * 24 * 60 * 60; +const MAX_RECORDED_LABEL_CHARS = 4_000; + +function configuredLabelLimit(value: string | undefined, fallback: number) { + const parsed = value === undefined ? fallback : Number(value); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : fallback; +} + +function labelLimitEnabled(env: Env) { + return env.FREE_LABEL_RPM !== undefined || env.FREE_LABEL_DAILY !== undefined; +} + +/** + * One public classifier gets one quota, regardless of how many IPs carry it. + * This is deliberately separate from per-IP admission: the existing limiter + * object gives it the same atomic minute/day semantics without linking callers. + */ +async function limitClassifier(env: Env, fingerprint: string, cost: number) { + const rpm = configuredLabelLimit(env.FREE_LABEL_RPM, LABEL_LIMITS.rpm); + const daily = configuredLabelLimit(env.FREE_LABEL_DAILY, LABEL_LIMITS.daily); + try { + const id = env.LIMITER.idFromName(`classifier:${fingerprint}`); + const response = await env.LIMITER.get(id).fetch( + `https://limiter/?limit=${rpm}&daily=${daily}&cost=${cost}`, + ); + const result = await response.json() as { limited?: unknown; scope?: unknown; resetIn?: unknown; remaining?: unknown }; + if (typeof result.limited !== "boolean") throw new Error("Invalid classifier quota response"); + return { + limited: result.limited, + scope: result.scope === "day" ? "day" as const : "minute" as const, + resetIn: typeof result.resetIn === "number" && Number.isFinite(result.resetIn) ? result.resetIn : 60, + remaining: typeof result.remaining === "number" && Number.isFinite(result.remaining) ? result.remaining : -1, + limit: result.scope === "day" ? daily : rpm, + rpm, + daily, + }; + } catch { + // The ordinary per-IP limiter already fails open on infrastructure errors; + // this additional abuse shield must not become a new availability dependency. + return { limited: false, scope: "minute" as const, resetIn: 60, remaining: -1, limit: rpm, rpm, daily }; + } +} + +function retainedLabels(labels: string[]) { + const normalized = [...labels] + .map((label) => label.trim().toLowerCase().replace(/\s+/g, " ")) + .sort(); + if (!normalized.length || normalized.some((label) => !label) || JSON.stringify(normalized).length > MAX_RECORDED_LABEL_CHARS) + return null; + return normalized; +} + +async function rememberClassifier(env: Env, fingerprint: string, labels: string[]) { + if (!fingerprint) return; + const kept = retainedLabels(labels); + if (!kept) return; + const key = `cls:${fingerprint}`; + const existing = await env.STATS.get(key); + if (existing) { + try { + if (Array.isArray((JSON.parse(existing) as { labels?: unknown }).labels)) return; + } catch { + // Legacy entries were timestamps; replace them on the next success. + } + } + await env.STATS.put(key, JSON.stringify({ labels: kept, firstSeen: new Date().toISOString() }), { expirationTtl: CLASSIFIER_TTL }); +} + /** Enterprise and operator-agent callers use dedicated unmetered bearer credentials. */ async function hasEnterpriseAccess(req: Request, env: Env) { const keys = [env.ENTERPRISE_API_KEY, env.AGENT_API_KEY].filter((key): key is string => !!key); @@ -1162,13 +1233,10 @@ export function record(env: Env, ctx: ExecutionContext, d: { } catch { /* analytics must never break a request */ } - // Distinct classifier registry, for the daily digest. Cheap: one write per new label set. + // Distinct classifier registry, for abuse review and the daily digest. + // Labels are aggregate configuration, never joined to caller or input data. try { - if (!labels) return; - const key = `cls:${labels}`; - if (!(await env.STATS.get(key))) { - await env.STATS.put(key, new Date().toISOString(), { expirationTtl: 60 * 60 * 24 * 90 }); - } + if (d.status === 200 && d.mode !== "dimensions") await rememberClassifier(env, labels, d.labels); } catch { /* ignore */ } @@ -2142,6 +2210,30 @@ const worker = { ); } + // Per-IP admission cannot stop a fleet that rotates addresses. The stable, + // keyed classifier fingerprint supplies the missing cross-IP boundary while + // funded and operator traffic retain the capacity it paid for. + if (!account && !enterprise && !execution?.internal && labelLimitEnabled(env)) { + const fingerprint = await classifierId(env, labels); + const classifierGate = await limitClassifier(env, fingerprint, decisions); + if (classifierGate.limited) { + layaAbort?.abort(); + return fail( + `The shared free allowance for this label set has reached ${classifierGate.limit.toLocaleString("en-US")} classifications ${classifierGate.scope === "day" ? "today" : "this minute"}. Use a funded workspace key for a separate allowance.`, + 429, + "label_set_limit", + { + "retry-after": String(Math.max(1, Math.ceil(classifierGate.resetIn))), + "ratelimit-limit": String(classifierGate.limit), + "ratelimit-policy": `${classifierGate.rpm};w=60, ${classifierGate.daily};w=86400`, + }, + 0, + classifierGate.remaining, + { upgrade: "https://classifier.dev/pricing" }, + ); + } + } + // Validate both request shapes and pass the normal tier gate before spending // the separate, deliberately small Laya allowance. if (layaPlan && !combinedQuota) { diff --git a/src/openapi.ts b/src/openapi.ts index a120b7f..ae37039 100644 --- a/src/openapi.ts +++ b/src/openapi.ts @@ -24,7 +24,7 @@ export const ERROR_CODES = [ // 409: the same skill text is already listed "duplicate_skill", // 429 - "rate_limit_minute", "rate_limit_day", "rate_limit_hour", + "rate_limit_minute", "rate_limit_day", "rate_limit_hour", "label_set_limit", // 502: the model provider failed after retries "typesafe", "chain_exhausted", "batch_unavailable", "timeout", "upstream_other", // 500 @@ -121,7 +121,7 @@ const errors = (plain: boolean) => ({ "403": err("The key is inactive or the workspace cannot authorize usage."), "503": err("Workspace billing or Laya inference is temporarily unavailable. A cold bulk worker can return laya_unavailable; respect Retry-After and retry with backoff."), "404": err("No such path. The body points at the docs, llms.txt, the spec and the sitemap.", undefined, plain), - "429": err("Quota or shared Laya capacity reached. Wait Retry-After seconds. Laya trial caps also apply to paid keys and cannot be lifted by upgrading; code laya_rate_limit identifies that lane's admission limit. Daily limits use rate_limit_day.", { + "429": err("Quota or shared Laya capacity reached. Wait Retry-After seconds. Anonymous requests also share a global allowance with every request using the same label set; code label_set_limit identifies that limit. Laya trial caps also apply to paid keys and cannot be lifted by upgrading; code laya_rate_limit identifies that lane's admission limit. Daily per-caller limits use rate_limit_day.", { "Retry-After": { schema: { type: "integer" }, description: "Seconds until the window resets." }, ...RATE_LIMIT_HEADERS, }, plain), @@ -1207,7 +1207,7 @@ export const OPENAPI = { description: "Stable machine-readable code: one of the listed values, or typesafe_ / openrouter_ carrying the upstream HTTP status. " + "400: bad_dimensions, too_many_decisions, dimension_context_too_large, bad_json, no_input, too_many_inputs, too_few_labels, too_many_labels, empty_label, duplicate_labels, empty_input, input_too_long, bad_tier, bad_cursor, invalid_submission, skill_invalid, account_route_required (use POST /v1/classify with a workspace key). " + - "404: not_found. 409: duplicate_skill. 429: rate_limit_minute, rate_limit_day, rate_limit_hour. 502: typesafe, typesafe_, openrouter_, chain_exhausted, batch_unavailable, timeout, upstream_other. 500: internal. 503: review_unavailable, inference_unavailable (provider credentials are not configured).", + "404: not_found. 409: duplicate_skill. 429: rate_limit_minute, rate_limit_day, rate_limit_hour, label_set_limit. 502: typesafe, typesafe_, openrouter_, chain_exhausted, batch_unavailable, timeout, upstream_other. 500: internal. 503: review_unavailable, inference_unavailable (provider credentials are not configured).", anyOf: [{ enum: [...ERROR_CODES] }, { pattern: UPSTREAM_CODE_PATTERN }], }, retryable: { type: "boolean", description: "Whether retrying later can resolve a spending refusal. Use backoff and Retry-After; do not loop on false." }, diff --git a/src/pages.ts b/src/pages.ts index 6bf4fa5..6896087 100644 --- a/src/pages.ts +++ b/src/pages.ts @@ -522,7 +522,11 @@ RATE LIMITS Pro Smart 2,000/minute; 20,000/day Limits count classifications and are shared across workspace keys and - agents. Public access is limited per IP. Laya trial limits apply to every plan. + agents. Public access is limited per IP. Anonymous traffic also shares a + 5,000/minute and 50,000/day allowance with every caller using the same label + set, so rotating addresses does not create a new budget. Workspace and + enterprise keys do not use that anonymous label-set allowance. Laya trial + limits apply to every plan. `; export const ABOUT = `About classifier.dev @@ -622,9 +626,11 @@ SECURITY export const PRIVACY = `classifier.dev privacy The short version: public, keyless classification does not store your texts. -Signed-in workspaces store account details, API keys and usage records so you -can manage access and billing. Optional workspace request-content logging is -described below and is disabled by default. +Successful classifier label names are retained for 90 days as one aggregate +record per label set, without caller identity or source text. Signed-in +workspaces store account details, API keys and usage records so you can manage +access and billing. Optional workspace request-content logging is described +below and is disabled by default. WHAT IS SENT WHERE @@ -642,9 +648,13 @@ PUBLIC SERVICE LOGS Per request, for rate limiting and operations: which tier ran, which model answered, the latency, the response status, a coarse request-country and a client family derived from the User-Agent (curl, python, browser, MCP, ...). - Not the text and not the labels: a keyed fingerprint of the label set counts - the distinct classifiers in use without recording anyone's wording. The - caller is a keyed hash of the IP that changes daily, so a record cannot be + Not the text: a keyed fingerprint of the label set counts distinct + classifiers in per-request analytics. Separately, successful simple and + multi-label classifier names are retained in an aggregate registry for 90 + days so operators can understand use and enforce one anonymous allowance per + label set. That registry has no caller fingerprint, request ID or source text. + Dimension definitions remain fingerprint-only. The caller is a keyed hash of + the IP that changes daily, so a record cannot be read back to an address or followed across days. The address itself serves the per-IP limits while the request is in flight and is not written down. To prevent abuse of free inference, the caller IP is sent to Spur's diff --git a/src/privacy.ts b/src/privacy.ts index 0aaa288..5a9cb99 100644 --- a/src/privacy.ts +++ b/src/privacy.ts @@ -1,16 +1,16 @@ /** * Pseudonyms, so that nothing kept points back at a caller. * - * Two pieces of caller data used to be written verbatim into Analytics Engine - * and KV: the IP address, and the label set. Labels are the caller's own words, - * and a set like "biopsy benign, biopsy malignant" says more about whoever sent - * it than every count around it put together. Both go through a keyed hash now. - * The dashboard and the digest can still count them; nobody can read them back. + * Caller addresses and label sets go through keyed hashes before they reach + * per-request analytics. The address never lands in storage. Successful + * classifier label names are also kept in a separate, bounded KV registry so + * operators can understand aggregate use; that record is never joined to the + * caller fingerprint or classified text and expires after 90 days. * * The key is a secret, because an unkeyed hash of either one is not a * pseudonym. The whole IPv4 space hashes in seconds on a laptop, and common - * label sets are a short word list, so anyone holding the dataset and this file - * — which is public — could invert both. Keyed, they cannot. + * label sets are a short word list, so anyone holding the analytics dataset and + * this file — which is public — could invert both. Keyed, they cannot. * * A caller pseudonym also takes the UTC day, so it is a different value * tomorrow and nothing accumulates into a profile of one person over 90 days of @@ -86,11 +86,29 @@ export function normalizeLabels(labels: string[]): string { /** * A stable, opaque name for one label set. Same set, same fingerprint, for as - * long as the salt lives — which is what lets the registry count distinct - * classifiers without ever holding one. + * long as the salt lives. Analytics contains only this value; the separate + * short-lived registry may resolve it to aggregate label names. */ export async function labelFingerprint(env: PrivacyEnv, labels: string[]): Promise { const normalized = normalizeLabels(labels); if (!normalized) return ""; return `ls_${await digest(env, `labels:${normalized}`)}`; } + +/** Read a current registry record; legacy timestamp-only entries stay opaque. */ +export async function recordedClassifierLabels( + stats: { get(key: string): Promise }, + fingerprint: string, +): Promise { + if (!fingerprint) return []; + try { + const raw = await stats.get(`cls:${fingerprint}`); + if (!raw) return []; + const labels = (JSON.parse(raw) as { labels?: unknown }).labels; + return Array.isArray(labels) && labels.every((label) => typeof label === "string") + ? labels as string[] + : []; + } catch { + return []; + } +} diff --git a/src/report.ts b/src/report.ts index f09584b..189931a 100644 --- a/src/report.ts +++ b/src/report.ts @@ -1,5 +1,6 @@ import { jevAttemptsQuery } from "./jev-observability"; import type { Env } from "./index"; +import { recordedClassifierLabels } from "./privacy"; const DATASET = "classifier_events"; const WINDOW_HOURS = 8; // three reports a day @@ -134,6 +135,15 @@ export async function dailyReport( (r) => r, ); + const nameClassifiers = async (rows: Record[]) => Promise.all(rows.map(async (row) => { + const labels = await recordedClassifierLabels(env.STATS, String(row.labels ?? "")); + return { ...row, label_names: labels.join(" · ") }; + })); + [enterpriseClassifiers, topClassifiers] = await Promise.all([ + nameClassifiers(enterpriseClassifiers), + nameClassifiers(topClassifiers), + ]); + const t = totals[0] ?? {}; const requests = num(t.requests); const classifications = num(t.classifications); @@ -207,15 +217,15 @@ export async function dailyReport( if (enterpriseClassifiers.length) { lines.push("ENTERPRISE CLASSIFIERS"); for (const r of enterpriseClassifiers) { - lines.push(` ${pad(String(num(r.requests)), 6)} ${pad(String(r.tier ?? "?"), 8)} ${String(r.labels)}`); + lines.push(` ${pad(String(num(r.requests)), 6)} ${pad(String(r.tier ?? "?"), 8)} ${String(r.label_names || r.labels)}`); } lines.push(""); } if (topClassifiers.length) { - lines.push("TOP CLASSIFIERS (fingerprints — the labels themselves are never kept)"); + lines.push("TOP CLASSIFIERS (aggregate label names; source text and callers are not kept here)"); for (const r of topClassifiers) { - lines.push(` ${pad(String(num(r.requests)), 6)} ${String(r.labels)}`); + lines.push(` ${pad(String(num(r.requests)), 6)} ${String(r.label_names || r.labels)}`); } lines.push(""); } diff --git a/src/wellknown.ts b/src/wellknown.ts index 6091bc3..45651c3 100644 --- a/src/wellknown.ts +++ b/src/wellknown.ts @@ -18,7 +18,7 @@ export const MCP_REGISTRY_AUTH = "v=MCPv1; k=ed25519; p=aWvcKpRNSyAPr+bh7ba+Hiiy export const MCP_REGISTRY_ENTRY = "https://registry.modelcontextprotocol.io/v0/servers?search=dev.classifier"; /** Bumped when any public page changes materially; feeds sitemap lastmod. */ -export const SITE_UPDATED = "2026-09-21"; +export const SITE_UPDATED = "2026-09-22"; export const SITE = { name: "classifier.dev", diff --git a/test/admin-fixture.ts b/test/admin-fixture.ts index ed38cb4..6f2360d 100644 --- a/test/admin-fixture.ts +++ b/test/admin-fixture.ts @@ -122,12 +122,14 @@ export function adminFixture(): AdminData { ], topLabels: [ { + label_names: "billing · product question · technical support", labels: "ls_9f3c7a8d", requests: 19002, classifications: 152016, usd: 0.4, }, { + label_names: "ham · spam", labels: "ls_64b1e7af", requests: 9211, classifications: 73688, @@ -160,7 +162,7 @@ export function adminFixture(): AdminData { }, ], failLabels: [ - { labels: "ls_9f3c7a8d", reason: "rate_limit_minute", requests: 80 }, + { label_names: "billing · product question · technical support", labels: "ls_9f3c7a8d", reason: "label_set_limit", requests: 80 }, ], dimensionTraffic: [ { diff --git a/test/record.test.ts b/test/record.test.ts index 9a590e9..225702d 100644 --- a/test/record.test.ts +++ b/test/record.test.ts @@ -9,12 +9,12 @@ const LABELS = ["invoice", "receipt", "payslip"]; /** Collects what the request would have written, and waits for it to be written. */ function spy() { const points: { blobs: string[]; indexes: string[] }[] = []; - const kv: string[] = []; + const kv: { key: string; value: string }[] = []; const pending: Promise[] = []; const env = { PRIVACY_SALT: "salt-under-test-0000000000000000", AE: { writeDataPoint: (p: { blobs: string[]; indexes: string[] }) => points.push(p) }, - STATS: { get: async () => null, put: async (k: string) => void kv.push(k) }, + STATS: { get: async () => null, put: async (key: string, value: string) => void kv.push({ key, value }) }, } as unknown as Env; const ctx = { waitUntil: (p: Promise) => pending.push(p) } as unknown as ExecutionContext; return { env, ctx, points, kv, settled: () => Promise.all(pending) }; @@ -33,19 +33,20 @@ describe("what a request leaves behind", () => { test("is never the caller's address", async () => { const s = spy(); await write(s); - const written = JSON.stringify(s.points) + s.kv.join(" "); + const written = JSON.stringify(s.points) + JSON.stringify(s.kv); expect(written).not.toContain(IP); expect(written).not.toContain("203.0.113"); expect(s.points[0].indexes[0]).toMatch(/^c_[0-9a-f]{16}$/); }); - test("is never the caller's labels", async () => { + test("keeps labels only in the separate aggregate registry", async () => { const s = spy(); await write(s); - const written = JSON.stringify(s.points) + s.kv.join(" "); - for (const label of LABELS) expect(written).not.toContain(label); + const analytics = JSON.stringify(s.points); + for (const label of LABELS) expect(analytics).not.toContain(label); expect(s.points[0].blobs[1]).toMatch(/^ls_[0-9a-f]{16}$/); - expect(s.kv[0]).toMatch(/^cls:ls_[0-9a-f]{16}$/); + expect(s.kv[0].key).toMatch(/^cls:ls_[0-9a-f]{16}$/); + expect(JSON.parse(s.kv[0].value)).toMatchObject({ labels: ["invoice", "payslip", "receipt"] }); }); test("still carries everything the dashboard counts", async () => { diff --git a/tests/spending.e2e.test.ts b/tests/spending.e2e.test.ts index 921e9dc..bb03088 100644 --- a/tests/spending.e2e.test.ts +++ b/tests/spending.e2e.test.ts @@ -53,15 +53,32 @@ function setup(overrides: Record = {}) { }, }; const pending: Promise[] = []; + const kv = new Map(); const ctx = { waitUntil(p: Promise) { pending.push(p); } } as ExecutionContext; const env = { SPENDING_ENABLED: "true", PRIVACY_SALT: "test-private-identity", SPUR_API_KEY: "fixture", TYPESAFE_API_KEY: "fixture", OPENROUTER_API_KEY: "fixture", AI_GATEWAY_DISABLED: "true", - STATS: { get: async () => null, put: async () => {} }, ...overrides, + STATS: { get: async (key: string) => kv.get(key) ?? null, put: async (key: string, value: string) => { kv.set(key, value); } }, ...overrides, } as unknown as Env; const budget = new FreeBudget({ storage } as unknown as DurableObjectState, env); env.FREE_BUDGET = { idFromName: () => "one", get: () => ({ fetch: (input: string | Request, init?: RequestInit) => budget.fetch(new Request(input, init)) }) } as unknown as DurableObjectNamespace; - return { env, ctx, stored, async flush() { while (pending.length) await Promise.all(pending.splice(0)); } }; + return { env, ctx, stored, kv, async flush() { while (pending.length) await Promise.all(pending.splice(0)); } }; +} + +function quotaNamespace() { + const used = new Map(); + return { + idFromName: (name: string) => name, + get: (name: string) => ({ fetch: async (input: string | Request) => { + const url = new URL(String(input)); + const cost = Number(url.searchParams.get("cost") ?? 1); + const limit = Number(url.searchParams.get("daily") ?? 5000); + const next = (used.get(name) ?? 0) + cost; + if (next > limit) return Response.json({ limited: true, scope: "day", remaining: 0, resetIn: 60 }); + used.set(name, next); + return Response.json({ limited: false, remaining: limit - next, dailyRemaining: limit - next }); + } }), + } as unknown as DurableObjectNamespace; } function request(ip = "203.0.113.1", path = "/v1/classify", body: unknown = { inputs: ["hello"], labels: ["a", "b"] }, extra: Record = {}) { return new Request(`https://classifier.dev${path}`, { method: "POST", headers: { "cf-connecting-ip": ip, "content-type": "application/json", ...extra }, body: JSON.stringify(body) }); @@ -136,6 +153,51 @@ test("known anonymous proxies are refused, enterprise VPNs accepted, and Spur ha expect(calls.length).toBe(2); }); +test("one anonymous label set has a global quota across rotating IPs and records its labels without caller data", async () => { + const s = setup({ LIMITER: quotaNamespace(), FREE_LABEL_RPM: "4", FREE_LABEL_DAILY: "4" }); + const calls = providers(); + + const malformed = await worker.fetch(request("203.0.113.1", undefined, { inputs: ["one"], labels: ["only-one"] }), s.env, s.ctx); + expect(malformed.status).toBe(400); + await s.flush(); + expect(s.kv.size).toBe(0); + + const responses = await Promise.all([ + worker.fetch(request("203.0.113.1", undefined, { inputs: ["one", "two"], labels: ["Spam", "Not spam"] }), s.env, s.ctx), + worker.fetch(request("203.0.113.2", undefined, { inputs: ["three", "four"], labels: ["not spam", "spam"] }), s.env, s.ctx), + worker.fetch(request("203.0.113.3", undefined, { inputs: ["five", "six"], labels: ["spam", "not spam"] }), s.env, s.ctx), + ]); + expect(responses.filter(response => response.status === 200)).toHaveLength(2); + const denied = responses.find(response => response.status === 429)!; + expect((await denied.json() as { code: string }).code).toBe("label_set_limit"); + expect(denied.headers.get("ratelimit-limit")).toBe("4"); + expect(denied.headers.get("ratelimit-policy")).toBe("4;w=60, 4;w=86400"); + expect(calls).toHaveLength(2); + await s.flush(); + + const records = [...s.kv.entries()].filter(([key]) => key.startsWith("cls:")); + expect(records).toHaveLength(1); + expect(records[0][0]).toMatch(/^cls:ls_[0-9a-f]{16}$/); + expect(JSON.parse(records[0][1])).toEqual(expect.objectContaining({ labels: ["not spam", "spam"] })); + expect(records[0][1]).not.toContain("203.0.113"); + + expect((await worker.fetch(request("203.0.113.4", undefined, { inputs: ["seven"], labels: ["ham", "eggs"] }), s.env, s.ctx)).status).toBe(200); + await s.flush(); +}); + +test("funded and enterprise traffic bypasses the anonymous label-set quota, and label storage is best effort", async () => { + const brokenStats = { get: async () => { throw new Error("KV down"); }, put: async () => { throw new Error("KV down"); } }; + const s = setup({ LIMITER: quotaNamespace(), FREE_LABEL_RPM: "1", FREE_LABEL_DAILY: "1", STATS: brokenStats, ENTERPRISE_API_KEY: "partner" }); + const calls = providers(); + const body = { inputs: ["one"], labels: ["a", "b"] }; + + expect((await worker.fetch(request("203.0.113.1", undefined, body), s.env, s.ctx)).status).toBe(200); + expect((await worker.fetch(request("203.0.113.2", undefined, body), s.env, s.ctx)).status).toBe(429); + expect((await worker.fetch(request("203.0.113.3", undefined, body, { authorization: "Bearer partner" }), s.env, s.ctx)).status).toBe(200); + await s.flush(); + expect(calls).toHaveLength(2); +}); + test("duplicate idempotency keys never execute twice and body conflicts are rejected", async () => { const s = setup(); const calls = providers(); const headers = { "idempotency-key": "one-operation" }; diff --git a/wrangler.example.toml b/wrangler.example.toml index dddc134..d428439 100644 --- a/wrangler.example.toml +++ b/wrangler.example.toml @@ -61,6 +61,9 @@ FREE_IP_DAILY_USD = "0.50" FREE_REQUEST_USD = "0.01" FREE_IP_CONCURRENCY = "4" FREE_CONCURRENCY = "128" +# One shared budget follows an anonymous label set across rotating callers. +FREE_LABEL_RPM = "5000" +FREE_LABEL_DAILY = "50000" PAID_REQUEST_USD = "10" SPUR_MONTHLY_LOOKUPS = "45000" # Laya and Kev are hosted by Beam on its shared endpoints; there is no From 945dc458c7ec76b65ce2a903c22039d87d830cbc Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Tue, 22 Sep 2026 20:23:15 -0700 Subject: [PATCH 2/3] Preserve public quota documentation contract --- src/docs.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/docs.ts b/src/docs.ts index 69ca5cd..822e4fb 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -477,9 +477,9 @@ LIMITS requests accept at most 200 inputs; larger batches return 400 so callers can split them. Anonymous traffic also shares a 5,000/minute and 50,000/day allowance across every caller using the same label set. Rotating IPs does - not reset it; workspace, operator and partner keys bypass it. Pro workspaces - allow 30,000/minute and 200,000/day on fast, 2,000/minute and 20,000/day on - smart, shared across keys and agents. + not reset it; workspace, operator and partner keys bypass it. + Pro workspaces allow 30,000/minute and 200,000/day on fast, 2,000/minute and + 20,000/day on smart, shared across keys and agents. Pro, operator and partner keys have a 1,000-input ceiling. Workspace keys use the workspace credit balance and share workspace quotas. Free workspaces have the same ceilings as public access. Current plans are at From 28c3c999036f0703bf0d5bcec3d2eb380f076e62 Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Tue, 22 Sep 2026 20:49:06 -0700 Subject: [PATCH 3/3] Close cross-endpoint label allowance bypasses --- README.md | 5 +- e2e/spending.mjs | 34 +++++++++- src/docs.ts | 10 ++- src/features/admin/reasons.ts | 1 + src/index.ts | 121 +++++++++++++++++++++------------- src/openapi.ts | 6 +- src/pages.ts | 5 +- src/privacy.ts | 16 +++-- src/typesafe-compat.ts | 23 +++++++ tests/spending.e2e.test.ts | 96 +++++++++++++++++++++++---- wrangler.local.toml | 2 + 11 files changed, 245 insertions(+), 74 deletions(-) diff --git a/README.md b/README.md index 5af1bee..6a79f49 100644 --- a/README.md +++ b/README.md @@ -268,8 +268,9 @@ The existing 15-minute alert check warns when at least three attempts fail and failures exceed 5% for either provider. Counts account for Analytics Engine sampling. No input text, caller identifiers or upstream messages are stored. Successful simple and multi-label classifier names are retained for 90 days in -a separate aggregate KV record that is not joined to a caller or source text; -per-request analytics contain only their keyed fingerprint. +a separate aggregate KV record with no caller or source-text fields. Its keyed +fingerprint also appears in per-request analytics, allowing operators to resolve +label names in those pseudonymous records. `AI_GATEWAY_DISABLED = "true"` in `wrangler.example.toml` keeps production on TypeSafe directly after the gateway repeatedly returned 429 on September 19. diff --git a/e2e/spending.mjs b/e2e/spending.mjs index 813b50c..42cf1d4 100644 --- a/e2e/spending.mjs +++ b/e2e/spending.mjs @@ -10,7 +10,7 @@ let barrier = Promise.resolve(); const modules = (await readdir('dist/server', { recursive: true })).filter(p => p.endsWith('.js')).sort((a, b) => a === 'index.js' ? -1 : b === 'index.js' ? 1 : a.localeCompare(b)).map(p => ({ type: 'ESModule', path: `dist/server/${p}` })); const mf = new Miniflare(convertV4MiniflareOptions({ workers: [{ name: "spending", modules, modulesRoot: 'dist/server', compatibilityDate: '2026-08-01', compatibilityFlags: ['nodejs_compat'], durableObjects: { FREE_BUDGET: { className: 'FreeBudget', useSQLite: true }, LIMITER: { className: 'RateLimiter', useSQLite: true } }, kvNamespaces: ['STATS'], - bindings: { SPENDING_ENABLED: 'true', PRIVACY_SALT: 'private-e2e-fixture', SPUR_API_KEY: 'fixture', TYPESAFE_API_KEY: 'fixture', OPENROUTER_API_KEY: 'fixture', INTERNAL_API_KEY: 'private-fixture' }, + bindings: { SPENDING_ENABLED: 'true', PRIVACY_SALT: 'private-e2e-fixture', SPUR_API_KEY: 'fixture', TYPESAFE_API_KEY: 'fixture', OPENROUTER_API_KEY: 'fixture', INTERNAL_API_KEY: 'private-fixture', FREE_LABEL_RPM: '200', FREE_LABEL_DAILY: '200' }, outboundService: async request => { if (request.url.startsWith('https://api.spur.us/')) { lookups++; return WorkerResponse.json({}); } calls++; await barrier; @@ -46,6 +46,38 @@ try { const duplicate = await request('203.0.113.4', undefined, undefined, { 'idempotency-key': 'one' }); assert.equal(duplicate.status, 409); report.results.push({ name: 'durable idempotency', first: first.status, duplicate: duplicate.status }); + const beforeLabels = calls; + const labels = ['allow', 'deny']; + const labelResponses = await Promise.all(Array.from({ length: 3 }, (_, index) => request(`203.0.113.${20 + index}`, { + inputs: Array(100).fill('short text'), labels: index === 1 ? ['DENY', 'ALLOW'] : labels, + }))); + const labelStatuses = labelResponses.map(response => response.status).sort(); + assert.deepEqual(labelStatuses, [200, 200, 429]); + const labelDenied = labelResponses.find(response => response.status === 429); + assert.equal((await labelDenied.json()).code, 'label_set_limit'); + assert.equal(labelDenied.headers.get('ratelimit-remaining'), '0'); + const afterLabels = calls; + const dimensionDenied = await request('203.0.113.23', { items: ['text'], dimensions: { renamed: labels } }); + assert.equal(dimensionDenied.status, 429); + const sdkDenied = await request('203.0.113.24', { state: 'text', questions: { renamed: { type: 'choice', criteria: { allow: null, deny: null } } } }, '/v1/systemone'); + assert.equal(sdkDenied.status, 429); + assert.equal(calls, afterLabels); + assert.equal((await request('203.0.113.25', { input: 'text', labels: ['independent', 'other'] })).status, 200); + const registry = await mf.getKVNamespace('STATS'); + let storedLabels; + for (let attempt = 0; attempt < 50 && !storedLabels; attempt++) { + const entries = await registry.list({ prefix: 'cls:' }); + for (const key of entries.keys) { + const record = await registry.get(key.name, 'json'); + if (record?.labels?.join(',') === 'allow,deny') storedLabels = record; + } + if (!storedLabels) await new Promise(resolve => setTimeout(resolve, 20)); + } + assert.deepEqual(storedLabels?.labels, labels); + assert.equal(JSON.stringify(storedLabels).includes('short text'), false); + assert.equal(JSON.stringify(storedLabels).includes('203.0.113'), false); + report.results.push({ name: 'global label allowance across concurrent IPs, dimensions and SDK', statuses: labelStatuses, + decisionsAccepted: 200, providerCalls: afterLabels - beforeLabels, dimensionStatus: dimensionDenied.status, sdkStatus: sdkDenied.status, registry: storedLabels }); jevUnavailable = true; const recovered = await request('203.0.113.5', { inputs: Array(119).fill('Invoice'), labels: ['billing', 'support'] }); assert.equal(recovered.status, 200); diff --git a/src/docs.ts b/src/docs.ts index 822e4fb..ee86f1f 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -478,6 +478,10 @@ LIMITS can split them. Anonymous traffic also shares a 5,000/minute and 50,000/day allowance across every caller using the same label set. Rotating IPs does not reset it; workspace, operator and partner keys bypass it. + REST, MCP, dimensions and TypeSafe choice questions share these counters. + Each dimension debits its labels by the item count; each SDK choice question + debits its labels once. These count attempts admitted by this gate, including + attempts subsequently refused by another quota or a provider. Pro workspaces allow 30,000/minute and 200,000/day on fast, 2,000/minute and 20,000/day on smart, shared across keys and agents. Pro, operator and partner keys have a 1,000-input ceiling. @@ -518,6 +522,8 @@ ERRORS 502 typesafe or typesafe_ when the decision model failed; openrouter_, chain_exhausted or timeout when the fallback chain did; upstream_other. Retry with backoff. + 503 label_set_unavailable when label admission cannot be checked; + no inference starts. Respect Retry-After and retry with backoff. 402 request_spending_limit: send fewer or shorter inputs, or use a funded workspace key. Do not repeatedly retry an unchanged over-budget request. @@ -548,7 +554,9 @@ PRIVACY keyed fingerprint of the label set, plus the tier, model, latency, status and coarse country. Successful simple and multi-label classifier names are also kept for 90 days in a separate aggregate registry with no caller identity or - source text. The usage counts are built from the fingerprinted records. + source text. Its shared fingerprint lets operators associate label names + with pseudonymous usage records. The same collection applies to TypeSafe + choice requests with one distinct label set. Built by @michael_chomsky — https://x.com/michael_chomsky diff --git a/src/features/admin/reasons.ts b/src/features/admin/reasons.ts index a8d97d8..ec32109 100644 --- a/src/features/admin/reasons.ts +++ b/src/features/admin/reasons.ts @@ -12,6 +12,7 @@ const REASON_TEXT: Record = { rate_limit_minute: "per-minute rate limit", rate_limit_day: "daily rate limit", label_set_limit: "anonymous label-set limit", + label_set_unavailable: "label-set admission unavailable", bad_dimensions: "invalid dimension definitions or conflicting options", too_many_decisions: "too many item × dimension decisions", dimension_context_too_large: "input and dimension exceed the model context", diff --git a/src/index.ts b/src/index.ts index 3ef752a..6a350be 100644 --- a/src/index.ts +++ b/src/index.ts @@ -34,7 +34,7 @@ export { RateLimiter } from "./limiter"; export { QuotaCoordinator } from "./admission"; import { admit, type AdmissionResult, type Quota } from "./admission"; import { pricingHtml } from "./pricingui"; -import { typeSafeCompatibleResponse, typeSafeDecisionCount } from "./typesafe-compat"; +import { typeSafeCompatibleResponse, typeSafeDecisionCount, typeSafeLabelSets } from "./typesafe-compat"; export interface Env extends LayaEnv, SpendingEnv { QUOTAS?: DurableObjectNamespace; @@ -1129,32 +1129,57 @@ function labelLimitEnabled(env: Env) { async function limitClassifier(env: Env, fingerprint: string, cost: number) { const rpm = configuredLabelLimit(env.FREE_LABEL_RPM, LABEL_LIMITS.rpm); const daily = configuredLabelLimit(env.FREE_LABEL_DAILY, LABEL_LIMITS.daily); + const id = env.LIMITER.idFromName(`classifier:${fingerprint}`); + const response = await env.LIMITER.get(id).fetch( + `https://limiter/?limit=${rpm}&daily=${daily}&cost=${cost}`, + ); + const result = await response.json() as { limited?: unknown; scope?: unknown; resetIn?: unknown; remaining?: unknown }; + if (!response.ok || typeof result.limited !== "boolean" || + typeof result.remaining !== "number" || !Number.isFinite(result.remaining) || result.remaining < 0 || + (result.limited && (!["day", "minute"].includes(String(result.scope)) || + typeof result.resetIn !== "number" || !Number.isFinite(result.resetIn) || result.resetIn <= 0))) + throw new Error("Invalid classifier quota response"); + return { + limited: result.limited, + scope: result.scope === "day" ? "day" as const : "minute" as const, + resetIn: typeof result.resetIn === "number" ? result.resetIn : 60, + limit: result.scope === "day" ? daily : rpm, + rpm, + daily, + }; +} + +type ClassifierRefusal = { message: string; status: number; code: ErrorCode; headers: Record }; + +async function classifierRefusal(env: Env, sets: { labels: string[]; cost: number }[]): Promise { + if (!labelLimitEnabled(env)) return null; try { - const id = env.LIMITER.idFromName(`classifier:${fingerprint}`); - const response = await env.LIMITER.get(id).fetch( - `https://limiter/?limit=${rpm}&daily=${daily}&cost=${cost}`, - ); - const result = await response.json() as { limited?: unknown; scope?: unknown; resetIn?: unknown; remaining?: unknown }; - if (typeof result.limited !== "boolean") throw new Error("Invalid classifier quota response"); - return { - limited: result.limited, - scope: result.scope === "day" ? "day" as const : "minute" as const, - resetIn: typeof result.resetIn === "number" && Number.isFinite(result.resetIn) ? result.resetIn : 60, - remaining: typeof result.remaining === "number" && Number.isFinite(result.remaining) ? result.remaining : -1, - limit: result.scope === "day" ? daily : rpm, - rpm, - daily, - }; + const grouped = new Map(); + for (const set of sets) { + const fingerprint = await classifierId(env, set.labels); + if (fingerprint) grouped.set(fingerprint, (grouped.get(fingerprint) ?? 0) + set.cost); + } + for (const [fingerprint, cost] of grouped) { + const gate = await limitClassifier(env, fingerprint, cost); + if (gate.limited) return { + status: 429, code: "label_set_limit", + message: `The shared free allowance for this label set has reached ${gate.limit.toLocaleString("en-US")} classifications ${gate.scope === "day" ? "today" : "this minute"}. Use a funded workspace key for a separate allowance.`, + headers: { + "retry-after": String(Math.max(1, Math.ceil(gate.resetIn))), + "ratelimit-limit": String(gate.limit), "ratelimit-remaining": "0", + "ratelimit-policy": `${gate.rpm};w=60, ${gate.daily};w=86400`, + }, + }; + } + return null; } catch { - // The ordinary per-IP limiter already fails open on infrastructure errors; - // this additional abuse shield must not become a new availability dependency. - return { limited: false, scope: "minute" as const, resetIn: 60, remaining: -1, limit: rpm, rpm, daily }; + return { status: 503, code: "label_set_unavailable", message: "Label-set admission is temporarily unavailable; retry with backoff.", headers: { "retry-after": "5" } }; } } function retainedLabels(labels: string[]) { - const normalized = [...labels] - .map((label) => label.trim().toLowerCase().replace(/\s+/g, " ")) + const normalized = [...new Set(labels + .map((label) => label.trim().toLowerCase().replace(/\s+/g, " ")))] .sort(); if (!normalized.length || normalized.some((label) => !label) || JSON.stringify(normalized).length > MAX_RECORDED_LABEL_CHARS) return null; @@ -1217,7 +1242,7 @@ export function record(env: Env, ctx: ExecutionContext, d: { // Hashing is async, so the write moved off the response path entirely. The // point is unchanged apart from the two columns that used to carry the // caller: blob2 is a label-set fingerprint, index1 a day-scoped pseudonym. - // Neither can be read back into what the caller sent. See src/privacy.ts. + // The label registry can resolve blob2; it cannot recover raw IPs or inputs. ctx.waitUntil( (async () => { const [caller, labels] = await Promise.all([callerId(env, d.ip), classifierId(env, d.labels)]); @@ -1234,7 +1259,7 @@ export function record(env: Env, ctx: ExecutionContext, d: { /* analytics must never break a request */ } // Distinct classifier registry, for abuse review and the daily digest. - // Labels are aggregate configuration, never joined to caller or input data. + // The registry has no caller or input fields; analytics shares its key. try { if (d.status === 200 && d.mode !== "dimensions") await rememberClassifier(env, labels, d.labels); } catch { @@ -1472,6 +1497,18 @@ const worker = { const decisions = path === "v1/systemone" && req.method === "POST" ? typeSafeDecisionCount(body ?? "") : 0; + const labelSets = decisions ? typeSafeLabelSets(body ?? "") : []; + const sdkStarted = Date.now(); + const sdkRecord = (status: number, reason = "") => { + if (!decisions) return; + record(env, ctx, { + tier: "fast", n: status === 200 ? decisions : 0, ms: Date.now() - sdkStarted, + labels: labelSets.length === 1 ? labelSets[0].labels : labelSets.map(set => JSON.stringify(set.labels)), + mode: labelSets.length === 1 ? "single" : "dimensions", + ip, country, client, status, model: "typesafe", usd: meter.usd, reason, agent, + attempted: decisions, escalationFailed: 0, + }); + }; const multiplier = account?.multiplier ?? 1; const quotaOwner = account ? `account:${account.id}` : ip; const rpm = TIERS.fast.rpm * multiplier; @@ -1495,12 +1532,20 @@ const worker = { } } + if (decisions && !account && !enterprise && !execution?.internal) { + const refusal = await classifierRefusal(env, labelSets); + if (refusal) { + sdkRecord(refusal.status, refusal.code); + return json({ error: refusal.message, code: refusal.code }, refusal.status, refusal.headers); + } + } const response = await typeSafeCompatibleResponse( req, env.TYPESAFE_API_KEY, body, decisions ? meter : undefined, ); + sdkRecord(response.status, response.ok ? "" : `typesafe_${response.status}`); const headers = new Headers(response.headers); // Own the browser policy at this boundary. An upstream credentialed CORS // header combined with our wildcard origin would make an otherwise valid @@ -2140,6 +2185,12 @@ const worker = { const rpm = TIERS[tier].rpm * multiplier; // Waiting cannot make a batch larger than the entire window fit. if (!enterprise && decisions > rpm) return fail(`Maximum ${rpm} ${dimensions ? "decisions" : "inputs"} per ${account ? "account" : "public"} ${tier} request; split the batch to fit the per-minute quota`, 400, dimensions ? "too_many_decisions" : "too_many_inputs"); + if (!account && !enterprise && !execution?.internal) { + const refusal = await classifierRefusal(env, dimensions + ? dimensions.map(dimension => ({ labels: dimension.labels, cost: inputs.length })) + : [{ labels, cost: decisions }]); + if (refusal) return fail(refusal.message, refusal.status, refusal.code, refusal.headers); + } const regularQuotaStarted = performance.now(); const combinedQuota = !!layaPlan && env.QUOTA_COORDINATOR_ENABLED === "true" && !!env.QUOTAS; let layaRun: LayaRun | undefined; @@ -2210,30 +2261,6 @@ const worker = { ); } - // Per-IP admission cannot stop a fleet that rotates addresses. The stable, - // keyed classifier fingerprint supplies the missing cross-IP boundary while - // funded and operator traffic retain the capacity it paid for. - if (!account && !enterprise && !execution?.internal && labelLimitEnabled(env)) { - const fingerprint = await classifierId(env, labels); - const classifierGate = await limitClassifier(env, fingerprint, decisions); - if (classifierGate.limited) { - layaAbort?.abort(); - return fail( - `The shared free allowance for this label set has reached ${classifierGate.limit.toLocaleString("en-US")} classifications ${classifierGate.scope === "day" ? "today" : "this minute"}. Use a funded workspace key for a separate allowance.`, - 429, - "label_set_limit", - { - "retry-after": String(Math.max(1, Math.ceil(classifierGate.resetIn))), - "ratelimit-limit": String(classifierGate.limit), - "ratelimit-policy": `${classifierGate.rpm};w=60, ${classifierGate.daily};w=86400`, - }, - 0, - classifierGate.remaining, - { upgrade: "https://classifier.dev/pricing" }, - ); - } - } - // Validate both request shapes and pass the normal tier gate before spending // the separate, deliberately small Laya allowance. if (layaPlan && !combinedQuota) { diff --git a/src/openapi.ts b/src/openapi.ts index ae37039..31d7c4e 100644 --- a/src/openapi.ts +++ b/src/openapi.ts @@ -30,7 +30,7 @@ export const ERROR_CODES = [ // 500 "internal", // 503: required review or inference providers are unavailable - "review_unavailable", "inference_unavailable", + "review_unavailable", "inference_unavailable", "label_set_unavailable", ] as const; export type ErrorCode = (typeof ERROR_CODES)[number] | `typesafe_${number}` | `openrouter_${number}`; export const UPSTREAM_CODE_PATTERN = "^(typesafe|openrouter)_[0-9]{3}$"; @@ -119,7 +119,7 @@ const errors = (plain: boolean) => ({ "401": err("Invalid API key. Create a workspace key at /app/keys."), "402": err("Insufficient workspace balance to reserve inference usage."), "403": err("The key is inactive or the workspace cannot authorize usage."), - "503": err("Workspace billing or Laya inference is temporarily unavailable. A cold bulk worker can return laya_unavailable; respect Retry-After and retry with backoff."), + "503": err("Workspace billing, label-set admission or Laya inference is temporarily unavailable. label_set_unavailable means the label allowance could not be checked and inference did not start. A cold bulk worker can return laya_unavailable; respect Retry-After and retry with backoff."), "404": err("No such path. The body points at the docs, llms.txt, the spec and the sitemap.", undefined, plain), "429": err("Quota or shared Laya capacity reached. Wait Retry-After seconds. Anonymous requests also share a global allowance with every request using the same label set; code label_set_limit identifies that limit. Laya trial caps also apply to paid keys and cannot be lifted by upgrading; code laya_rate_limit identifies that lane's admission limit. Daily per-caller limits use rate_limit_day.", { "Retry-After": { schema: { type: "integer" }, description: "Seconds until the window resets." }, @@ -1207,7 +1207,7 @@ export const OPENAPI = { description: "Stable machine-readable code: one of the listed values, or typesafe_ / openrouter_ carrying the upstream HTTP status. " + "400: bad_dimensions, too_many_decisions, dimension_context_too_large, bad_json, no_input, too_many_inputs, too_few_labels, too_many_labels, empty_label, duplicate_labels, empty_input, input_too_long, bad_tier, bad_cursor, invalid_submission, skill_invalid, account_route_required (use POST /v1/classify with a workspace key). " + - "404: not_found. 409: duplicate_skill. 429: rate_limit_minute, rate_limit_day, rate_limit_hour, label_set_limit. 502: typesafe, typesafe_, openrouter_, chain_exhausted, batch_unavailable, timeout, upstream_other. 500: internal. 503: review_unavailable, inference_unavailable (provider credentials are not configured).", + "404: not_found. 409: duplicate_skill. 429: rate_limit_minute, rate_limit_day, rate_limit_hour, label_set_limit. 502: typesafe, typesafe_, openrouter_, chain_exhausted, batch_unavailable, timeout, upstream_other. 500: internal. 503: review_unavailable, inference_unavailable (provider credentials are not configured), label_set_unavailable (label allowance could not be checked; inference did not start).", anyOf: [{ enum: [...ERROR_CODES] }, { pattern: UPSTREAM_CODE_PATTERN }], }, retryable: { type: "boolean", description: "Whether retrying later can resolve a spending refusal. Use backoff and Retry-After; do not loop on false." }, diff --git a/src/pages.ts b/src/pages.ts index 6896087..48840bf 100644 --- a/src/pages.ts +++ b/src/pages.ts @@ -653,7 +653,10 @@ PUBLIC SERVICE LOGS multi-label classifier names are retained in an aggregate registry for 90 days so operators can understand use and enforce one anonymous allowance per label set. That registry has no caller fingerprint, request ID or source text. - Dimension definitions remain fingerprint-only. The caller is a keyed hash of + Its shared label-set fingerprint can associate label names with pseudonymous + usage records. Label names may themselves contain information you supply. + Dimension definitions remain fingerprint-only. TypeSafe choice requests with + one distinct label set use the same registry. The caller is a keyed hash of the IP that changes daily, so a record cannot be read back to an address or followed across days. The address itself serves the per-IP limits while the request is in flight and is not written down. diff --git a/src/privacy.ts b/src/privacy.ts index 5a9cb99..07820de 100644 --- a/src/privacy.ts +++ b/src/privacy.ts @@ -1,16 +1,18 @@ /** - * Pseudonyms, so that nothing kept points back at a caller. + * Day-scoped caller pseudonyms and stable classifier fingerprints. * * Caller addresses and label sets go through keyed hashes before they reach * per-request analytics. The address never lands in storage. Successful * classifier label names are also kept in a separate, bounded KV registry so - * operators can understand aggregate use; that record is never joined to the - * caller fingerprint or classified text and expires after 90 days. + * operators can understand aggregate use. It expires after 90 days and has no + * caller or input fields, but its shared fingerprint resolves label names in + * pseudonymous analytics. * * The key is a secret, because an unkeyed hash of either one is not a * pseudonym. The whole IPv4 space hashes in seconds on a laptop, and common - * label sets are a short word list, so anyone holding the analytics dataset and - * this file — which is public — could invert both. Keyed, they cannot. + * label sets are a short word list, so anyone holding only the analytics + * dataset and this public source could invert an unkeyed hash. The separate + * label registry intentionally resolves fingerprints for operators. * * A caller pseudonym also takes the UTC day, so it is a different value * tomorrow and nothing accumulates into a profile of one person over 90 days of @@ -73,10 +75,10 @@ export async function callerId(env: PrivacyEnv, ip: string): Promise { /** Order and case never distinguished two classifiers, so neither does this. */ export function normalizeLabels(labels: string[]): string { - return [...labels] + return [...new Set(labels .filter((l) => typeof l === "string") .map((l) => l.toLowerCase().trim()) - .filter(Boolean) + .filter(Boolean))] .sort() // Keep existing fingerprints for ordinary labels; escape delimiter-bearing // labels so ["a|b", "c"] and ["a", "b|c"] are different classifiers. diff --git a/src/typesafe-compat.ts b/src/typesafe-compat.ts index dba45af..1ee5958 100644 --- a/src/typesafe-compat.ts +++ b/src/typesafe-compat.ts @@ -1,5 +1,6 @@ import { providerFetch } from "./spending/permit"; import { SpendingError } from "./spending/policy"; +import { normalizeLabels } from "./privacy"; /** * Wire-compatible TypeSafe API surface. * @@ -33,6 +34,28 @@ export function typeSafeDecisionCount(body: string): number { } } +/** Choice keys are labels; question names, instructions and state are not. */ +export function typeSafeLabelSets(body: string): { labels: string[]; cost: number }[] { + try { + const questions = object(object(JSON.parse(body))?.questions); + if (!questions) return []; + const sets = new Map(); + for (const value of Object.values(questions)) { + const question = object(value); + const criteria = question?.type === "choice" ? object(question.criteria) : null; + const labels = criteria ? Object.keys(criteria) : []; + if (labels.length < 2 || labels.some(label => !label.trim())) continue; + const key = normalizeLabels(labels); + const set = sets.get(key) ?? { labels, cost: 0 }; + set.cost++; + sets.set(key, set); + } + return [...sets.values()]; + } catch { + return []; + } +} + /** Forward one official TypeSafe SDK request using the service credential. */ export async function typeSafeCompatibleResponse( request: Request, diff --git a/tests/spending.e2e.test.ts b/tests/spending.e2e.test.ts index bb03088..b6c854f 100644 --- a/tests/spending.e2e.test.ts +++ b/tests/spending.e2e.test.ts @@ -3,6 +3,8 @@ import { SQL } from "bun"; import { readFileSync, readdirSync } from "node:fs"; import worker, { type Env } from "../src/index"; import { FreeBudget } from "../src/spending/free-budget"; +import { RateLimiter } from "../src/limiter"; +import { labelFingerprint } from "../src/privacy"; import { database as portableDatabase } from "./support/postgres"; import { provisionTestAccount } from "./support/account"; import { performAction } from "../src/server/agents"; @@ -66,18 +68,29 @@ function setup(overrides: Record = {}) { } function quotaNamespace() { - const used = new Map(); + const objects = new Map(); return { idFromName: (name: string) => name, - get: (name: string) => ({ fetch: async (input: string | Request) => { - const url = new URL(String(input)); - const cost = Number(url.searchParams.get("cost") ?? 1); - const limit = Number(url.searchParams.get("daily") ?? 5000); - const next = (used.get(name) ?? 0) + cost; - if (next > limit) return Response.json({ limited: true, scope: "day", remaining: 0, resetIn: 60 }); - used.set(name, next); - return Response.json({ limited: false, remaining: limit - next, dailyRemaining: limit - next }); - } }), + get: (name: string) => { + if (!objects.has(name)) { + const data = new Map(); + let pending = Promise.resolve(); + objects.set(name, new RateLimiter({ + storage: { + get: async (key: string) => structuredClone(data.get(key)), + put: async (entries: Record) => { + for (const [key, value] of Object.entries(entries)) data.set(key, structuredClone(value)); + }, + }, + blockConcurrencyWhile(fn: () => Promise) { + const next = pending.then(fn); + pending = next.then(() => {}, () => {}); + return next; + }, + } as unknown as DurableObjectState)); + } + return { fetch: (input: string | Request) => objects.get(name)!.fetch(new Request(input)) }; + }, } as unknown as DurableObjectNamespace; } function request(ip = "203.0.113.1", path = "/v1/classify", body: unknown = { inputs: ["hello"], labels: ["a", "b"] }, extra: Record = {}) { @@ -181,11 +194,16 @@ test("one anonymous label set has a global quota across rotating IPs and records expect(JSON.parse(records[0][1])).toEqual(expect.objectContaining({ labels: ["not spam", "spam"] })); expect(records[0][1]).not.toContain("203.0.113"); + const duplicateCase = await worker.fetch(request("203.0.113.5", undefined, { + inputs: ["extra"], labels: ["Spam", "Not spam", "SPAM"], + }), s.env, s.ctx); + expect(duplicateCase.status).toBe(429); + expect((await worker.fetch(request("203.0.113.4", undefined, { inputs: ["seven"], labels: ["ham", "eggs"] }), s.env, s.ctx)).status).toBe(200); await s.flush(); }); -test("funded and enterprise traffic bypasses the anonymous label-set quota, and label storage is best effort", async () => { +test("enterprise traffic bypasses the anonymous label-set quota, and label storage is best effort", async () => { const brokenStats = { get: async () => { throw new Error("KV down"); }, put: async () => { throw new Error("KV down"); } }; const s = setup({ LIMITER: quotaNamespace(), FREE_LABEL_RPM: "1", FREE_LABEL_DAILY: "1", STATS: brokenStats, ENTERPRISE_API_KEY: "partner" }); const calls = providers(); @@ -198,6 +216,58 @@ test("funded and enterprise traffic bypasses the anonymous label-set quota, and expect(calls).toHaveLength(2); }); +test("REST, dimensions and TypeSafe choice questions share a label allowance despite renamed fields", async () => { + const s = setup({ LIMITER: quotaNamespace(), FREE_LABEL_RPM: "10", FREE_LABEL_DAILY: "2" }); + const calls = providers(); + const labels = ["billing", "support"]; + expect((await worker.fetch(request("203.0.113.10", undefined, { inputs: ["first"], labels }), s.env, s.ctx)).status).toBe(200); + const dimension = await worker.fetch(request("203.0.113.11", undefined, { items: ["second"], dimensions: { renamed: labels } }), s.env, s.ctx); + expect(dimension.status).toBe(200); + const sdk = await worker.fetch(request("203.0.113.12", "/v1/systemone", { + state: "third", questions: { arbitrary: { type: "choice", instructions: "changed", criteria: { support: null, billing: null } } }, + }), s.env, s.ctx); + expect(sdk.status).toBe(429); + expect(await sdk.json()).toMatchObject({ code: "label_set_limit" }); + const renamed = await worker.fetch(request("203.0.113.13", undefined, { items: ["fourth"], dimensions: { another: labels } }), s.env, s.ctx); + expect(renamed.status).toBe(429); + expect(calls).toHaveLength(2); + await s.flush(); +}); + +test("label quotas reject before speculative Laya inference and fail closed on unavailable storage", async () => { + const s = setup({ SPENDING_ENABLED: "false", LIMITER: quotaNamespace(), FREE_LABEL_RPM: "1", FREE_LABEL_DAILY: "1", + LAYA_ENABLED: "true", BEAM_API_KEY: "fixture", QUOTA_COORDINATOR_ENABLED: "true", + QUOTAS: { idFromName: (name: string) => name, get: () => ({ fetch: async () => Response.json({ limited: false, remaining: 50, laneRemaining: 50 }) }) }, + LAYA_FAST_ADMISSION: { limit: async () => ({ success: true }) }, + }); + const calls = providers(); + const body = { inputs: ["first"], labels: ["a", "b"] }; + expect((await worker.fetch(request("203.0.113.1", undefined, body), s.env, s.ctx)).status).toBe(200); + const laya = await worker.fetch(request("203.0.113.2", undefined, { ...body, model: "laya", processing: "fast" }), s.env, s.ctx); + expect(laya.status).toBe(429); + expect(calls).toHaveLength(1); + const broken = { ...s.env, LIMITER: { idFromName: (name: string) => name, get: () => ({ fetch: async () => { throw new Error("offline"); } }) } } as unknown as Env; + const unavailable = await worker.fetch(request("203.0.113.3", undefined, body), broken, s.ctx); + expect(unavailable.status).toBe(503); + expect(await unavailable.json()).toMatchObject({ code: "label_set_unavailable" }); + expect(calls).toHaveLength(1); + await s.flush(); +}); + +test("successful SDK choice labels upgrade legacy records without storing state or question instructions", async () => { + const s = setup({ LIMITER: quotaNamespace(), FREE_LABEL_RPM: "10", FREE_LABEL_DAILY: "10" }); + providers(); + const fingerprint = await labelFingerprint(s.env, ["support", "billing"]); + s.kv.set(`cls:${fingerprint}`, "2026-09-01T00:00:00.000Z"); + const response = await worker.fetch(request("203.0.113.14", "/v1/systemone", { + state: "PRIVATE_STATE", questions: { category: { type: "choice", instructions: "PRIVATE_INSTRUCTION", criteria: { support: null, billing: null } } }, + }), s.env, s.ctx); + expect(response.status).toBe(200); + await s.flush(); + expect(JSON.parse(s.kv.get(`cls:${fingerprint}`)!)).toMatchObject({ labels: ["billing", "support"] }); + expect(JSON.stringify([...s.kv])).not.toMatch(/PRIVATE_|203\.0\.113/); +}); + test("duplicate idempotency keys never execute twice and body conflicts are rejected", async () => { const s = setup(); const calls = providers(); const headers = { "idempotency-key": "one-operation" }; @@ -222,7 +292,9 @@ test("uncertain attempts consume allowance, retries cannot exceed a penny, and f }); test("funded HTTP classification skips Spur/free budget and bills the published input-token price", async () => { - const s = setup(); + const s = setup({ FREE_LABEL_RPM: "1", FREE_LABEL_DAILY: "1", LIMITER: quotaNamespace() }); + const fingerprint = await labelFingerprint(s.env, ["a", "b"]); + await s.env.LIMITER.get(s.env.LIMITER.idFromName(`classifier:${fingerprint}`)).fetch("https://limiter/?limit=1&daily=1&cost=1"); const env = { ...s.env, APP_DB: database(), APP_ACCOUNTS_ENABLED: "true", API_KEY_ENCRYPTION_KEY: "test-only-key-encryption-secret-32-characters" } as Env & AppEnv; await provisionTestAccount(new Request("http://localhost/auth/demo", { headers: { origin: "http://localhost" } }), env); await env.APP_DB.prepare("UPDATE app_accounts SET paid_balance=balance WHERE id='local-demo'").run(); diff --git a/wrangler.local.toml b/wrangler.local.toml index 3b45b13..d0a4a49 100644 --- a/wrangler.local.toml +++ b/wrangler.local.toml @@ -32,6 +32,8 @@ dataset = "classifier_account_events" [vars] SPENDING_ENABLED = "true" +FREE_LABEL_RPM = "5000" +FREE_LABEL_DAILY = "50000" APP_ACCOUNTS_ENABLED = "true" # Gateway free-tier 429s observed in production. Re-enable only after capacity is verified. AI_GATEWAY_DISABLED = "true"