From b2bc8051f16b0221f502bd627544789548b133c6 Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Mon, 21 Sep 2026 04:56:48 -0700 Subject: [PATCH 1/2] Cut Laya fast latency below 150ms --- docs/postgres-setup.md | 12 ++++---- inference/laya/README.md | 16 +++++------ inference/laya/deploy.py | 7 +++-- inference/laya/smoke.py | 2 +- src/docs.ts | 2 +- src/index.ts | 54 ++++++++++++++++++++++++++++-------- src/laya.ts | 7 +++-- test/laya.test.ts | 46 +++++++++++++++++++++++++++--- tests/runtime-config.test.ts | 7 ++--- wrangler.example.toml | 19 ++++++++----- 10 files changed, 124 insertions(+), 48 deletions(-) diff --git a/docs/postgres-setup.md b/docs/postgres-setup.md index 78b6ddd..66e057a 100644 --- a/docs/postgres-setup.md +++ b/docs/postgres-setup.md @@ -101,12 +101,12 @@ connection strings into commands that will remain in shell history. Migrations are transactional, ordered, and checksum-verified; editing an already-applied migration is rejected. -The production template targets `aws:us-west-2`, supported by Wrangler 4.122's -placement schema. Cloudflare runs fetch handlers in a nearby Cloudflare data -center, not inside AWS itself. Placement does not move an existing Neon project -or Durable Object, and is not a residency guarantee. Confirm the Neon primary -and upstream latency before activating the production migration. See -[Cloudflare placement](https://developers.cloudflare.com/workers/configuration/placement/). +The production Worker executes at Cloudflare's ingress edge. Neon remains in +`aws:us-west-2`; removing the former Worker placement avoids forwarding every +latency-sensitive classification request through Oregon. This is not a data +residency guarantee and does not move the Neon primary or a Durable Object. +Confirm database-dependent account latency as well as inference latency before +changing placement again. See [Cloudflare placement](https://developers.cloudflare.com/workers/configuration/placement/). Do not enable the new account offering until retail token rates, Autumn events, legacy Pro handling, reconciliation, and end-to-end production checks are ready. diff --git a/inference/laya/README.md b/inference/laya/README.md index 8532449..7e91d0d 100644 --- a/inference/laya/README.md +++ b/inference/laya/README.md @@ -21,15 +21,15 @@ One single-label decision is one question. Multi-label asks one question per lab Text ≤2,000 characters, 2–16 labels ≤100 characters each, instructions ≤400 characters. The combined model input must also fit the selected checkpoint's context (512 English, 1,024 multilingual); long labels/instructions may hit the smaller question budget. Reject instead of truncate. -Modal app `classifier-laya-router-trial` has global compute placement (no west-coast pin), one ingress in `us-east`, and private proxy authentication. This is not a GPU replica in every region, and does not promise 40–70ms worldwide. Pools are isolated. Large flashes are shed, not absorbed into an unbounded backlog. Accepted payloads are only held in memory, never stored as durable jobs. +Modal app `classifier-laya-router-trial-west` colocates its compute and ingress in `us-west`, near west-coast Worker execution, and uses private proxy authentication. This avoids Modal's former east-to-global inter-region hop, but it is not a GPU replica in every region and does not promise the same client round-trip worldwide. Pools are isolated. Large flashes are shed, not absorbed into an unbounded backlog. Accepted payloads are only held in memory, never stored as durable jobs. ## Cost -At [Modal list prices](https://modal.com/pricing), checked 2026-09-20, approximate allocation cost is $0.80/hour L4 + 2 × $0.0473/hour CPU + 4 × $0.008/hour memory = **$0.9266/hour per lane**. +At [Modal list prices](https://modal.com/pricing), checked 2026-09-20, base allocation is $0.80/hour L4 + 2 × $0.0473/hour CPU + 4 × $0.008/hour memory. The narrow-region 1.75× multiplier makes the colocated deployment approximately **$1.6216/hour per lane**. -- Warm fast: about **$22.24/day or $667 per 30-day month**. -- Bulk: about **$0.93 per allocated hour**, including startup/idle tails; another ~$667/month if continuously active. -- Both continuously active: about **$44.48/day or $1,334/month**. Fleet caps are not a hard dollar budget; CPU overage, other infrastructure, reviews, taxes and credits are separate. +- Warm fast: about **$38.92/day or $1,168 per 30-day month**. +- Bulk: about **$1.62 per allocated hour**, including startup/idle tails; another ~$1,168/month if continuously active. +- Both continuously active: about **$77.83/day or $2,335/month**. Fleet caps are not a hard dollar budget; CPU overage, other infrastructure, reviews, taxes and credits are separate. Laya's explicit retail token rate is zero during the trial. Smart reviews retain existing prices. Provider-token spend is not Modal GPU hosting spend: inspect Modal usage for the actual bill. Account analytics marks combined provider cost unknown when Modal is involved. @@ -54,7 +54,7 @@ Worker configuration lives in `wrangler.example.toml`; production deploys from m To pause new requests, set `LAYA_ENABLED = "false"` in the canonical config and deploy. **Disabling the Worker route does not stop the warm GPU bill.** To stop both trial pools immediately: ```sh -uvx --from modal modal app stop classifier-laya-router-trial +uvx --from modal modal app stop classifier-laya-router-trial-west ``` Existing Jev remains available. Redeploy `deploy.py` and re-enable the Worker flag to resume. Do not scale up to hide overload without revisiting costs. @@ -63,13 +63,13 @@ Watch model labels `laya-0.3.4--` in existing classifier analy ## Quota admission rollout -`QuotaCoordinator` keeps each caller's existing fast/smart tier counters and fast/bulk Laya counters in one object. A warm Laya request checks and writes both quotas in one durable operation. Jev still shares its tier allowance with Laya, and a Laya lane still shares its allowance across tiers. Decision and question costs remain separate. A tier refusal spends nothing; a lane refusal retains the tier debit. Combined Laya admission fails closed if either counter is unavailable; Jev retains its existing fail-open behavior. +`QuotaCoordinator` keeps each caller's existing fast/smart tier counters and fast/bulk Laya counters in one object. A warm Laya request checks and writes both quotas in one durable operation. Anonymous fast calls first pass a local Cloudflare burst shield, then exact admission overlaps inference; an exact refusal aborts the in-flight Modal request and is still authoritative. Jev still shares its tier allowance with Laya, and a Laya lane still shares its allowance across tiers. Decision and question costs remain separate. A tier refusal spends nothing; a lane refusal retains the tier debit. Combined Laya admission fails closed if either counter is unavailable; Jev retains its existing fail-open behavior. Deploy the transfer-aware `RateLimiter`, coordinator binding and migration with `QUOTA_COORDINATOR_ENABLED = "false"` first. After that deployment completes, enable the flag in a second deployment. On first use, each old object freezes its counters and forwards later requests. The coordinator imports the frozen snapshot without resetting the minute or day allowance. Interrupted imports retry the same snapshot. Only opaque Durable Object IDs, scope names and counters are stored. Migration adds latency on the first call, not every call. Rollback by disabling `QUOTA_COORDINATOR_ENABLED` while retaining the new classes, bindings and forwarding code. **Do not roll back to code predating the transfer protocol:** its old counters are frozen and no longer authoritative. The disabled path continues to follow transferred counters. Neither successful admission nor forwarding uses unconfirmed storage writes. -Run `node inference/laya/latency.mjs ` before and after deployment from the same client. It records 20 sequential synthetic requests, separates the first call, and reports end-to-end, Worker, quota, Modal and backend timings. Do not equate backend compute time or the Worker-to-Modal span with client latency. The main Worker is in Oregon for account database access; the Modal ingress is in Virginia. Moving every request east would also move database-dependent account traffic, so placement must be evaluated per path. +Run `node inference/laya/latency.mjs ` before and after deployment from the same client. It records 20 sequential synthetic requests, separates the first call, and reports end-to-end, Worker, quota, Modal and backend timings. Do not equate backend compute time or the Worker-to-Modal span with client latency. The Worker executes at Cloudflare's ingress edge; Modal ingress and compute are colocated in `us-west`. Measure other client regions independently. ## Evidence diff --git a/inference/laya/deploy.py b/inference/laya/deploy.py index 3825d04..2e75d90 100644 --- a/inference/laya/deploy.py +++ b/inference/laya/deploy.py @@ -18,9 +18,12 @@ def download(): .env({"USE_TF": "0", "LAYA_REVISION": REVISION, "HF_HUB_OFFLINE": "1", "TRANSFORMERS_OFFLINE": "1"}) .add_local_file(root / "adapter.py", "/root/adapter.py") .add_local_file(root / "runtime.py", "/root/runtime.py")) -app = modal.App("classifier-laya-router-trial") +app = modal.App("classifier-laya-router-trial-west") resources = dict(image=image, gpu="L4", cpu=2, memory=4096, - compute_region=None, routing_region="us-east", unauthenticated=False, + # Keep Modal's ingress and container together near west-coast + # traffic: an east-coast ingress adds a continent-scale hop + # that dwarfs ~30 ms inference. + compute_region="us-west", routing_region="us-west", unauthenticated=False, max_containers=1, target_concurrency=1, startup_timeout=240) diff --git a/inference/laya/smoke.py b/inference/laya/smoke.py index 4364513..018d2b4 100644 --- a/inference/laya/smoke.py +++ b/inference/laya/smoke.py @@ -17,7 +17,7 @@ async def main(): report = {} async with httpx.AsyncClient(headers=headers, timeout=30) as client: for lane in ("fast", "bulk"): - url = f"https://miryaboy--classifier-laya-router-trial-{lane}.us-east.modal.direct" + url = f"https://miryaboy--classifier-laya-router-trial-west-{lane}.us-west.modal.direct" started = time.monotonic() while time.monotonic() - started < 240: try: diff --git a/src/docs.ts b/src/docs.ts index f4c350b..4cec3df 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -95,7 +95,7 @@ LAYA TRIAL Fast stays warm. Bulk starts on demand and can return 503 while starting. On 429 or 503, respect Retry-After and use bounded retries with backoff. The repository CLI retries Laya for up to three minutes per batch. - Global GPU placement is not a replica in every region or a latency promise. + One regional GPU deployment is not a replica in every region or a latency promise. Accepted work is held only in memory; there is no durable batch-job service. Overload never silently switches the model or processing lane. diff --git a/src/index.ts b/src/index.ts index 017dcbd..2031ec0 100644 --- a/src/index.ts +++ b/src/index.ts @@ -930,6 +930,14 @@ async function escalate(env: Env, inputs: string[], labels: string[], instructio * fast-tier answers under a smart-tier label. */ type LayaRequestTiming = LayaTiming & { runMs?: number }; +type LayaRun = Promise>>; + +function startLaya(env: Env, plan: LayaPlan, meter: Meter | undefined, timing: LayaRequestTiming | undefined, signal?: AbortSignal): LayaRun { + const started = performance.now(); + return runLaya(env, plan, meter, timing, signal).finally(() => { + if (timing) timing.runMs = performance.now() - started; + }); +} async function classifyMany( env: Env, @@ -941,6 +949,7 @@ async function classifyMany( meter?: Meter, layaPlan?: LayaPlan, layaTiming?: LayaRequestTiming, + layaRun?: LayaRun, ): Promise<{ results: Result[]; escalationFailed: number }> { const keys = jevKeys(env); if (keys || layaPlan) { @@ -948,9 +957,7 @@ async function classifyMany( let jev: Awaited> | null = null; try { if (layaPlan) { - const runStarted = performance.now(); - jev = await runLaya(env, layaPlan, meter, layaTiming); - if (layaTiming) layaTiming.runMs = performance.now() - runStarted; + jev = await (layaRun ?? startLaya(env, layaPlan, meter, layaTiming)); } else jev = await jevClassify(keys!, inputs, labels, instructions, !!multi, meter); } catch (e) { if (layaPlan) throw e; @@ -995,14 +1002,12 @@ async function classifyMany( } /** Keep each field independent, including smart escalation and the bounded LLM fallback. */ -async function classifyMatrix(env: Env, inputs: string[], dimensions: Dimension[], batches: DimensionBatch[], tier: Tier, instructions: string | undefined, meter: Meter, layaPlan?: LayaPlan, layaTiming?: LayaRequestTiming) { +async function classifyMatrix(env: Env, inputs: string[], dimensions: Dimension[], batches: DimensionBatch[], tier: Tier, instructions: string | undefined, meter: Meter, layaPlan?: LayaPlan, layaTiming?: LayaRequestTiming, layaRun?: LayaRun) { const started = Date.now(); let jev: Awaited> | undefined; const keys = jevKeys(env); if (layaPlan) { - const runStarted = performance.now(); - const flat = await runLaya(env, layaPlan, meter, layaTiming); - if (layaTiming) layaTiming.runMs = performance.now() - runStarted; + const flat = await (layaRun ?? startLaya(env, layaPlan, meter, layaTiming)); jev = inputs.map((_, i) => flat.slice(i * dimensions.length, (i + 1) * dimensions.length)); } else if (keys) { try { jev = await classifyDimensions(keys, batches, meter); } @@ -1987,21 +1992,45 @@ const worker = { 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"); const regularQuotaStarted = performance.now(); const combinedQuota = !!layaPlan && env.QUOTA_COORDINATOR_ENABLED === "true" && !!env.QUOTAS; + let layaRun: LayaRun | undefined; + let layaAbort: AbortController | undefined; + if (combinedQuota && layaTiming && processing === "fast" && env.LAYA_FAST_ADMISSION) { + try { + const key = env.LIMITER.idFromName(`laya:fast:${quotaOwner}`).toString(); + const checks = await Promise.all(Array.from({length:layaPlan!.cost}, () => env.LAYA_FAST_ADMISSION!.limit({key}))); + if (checks.some(check => !check.success)) + return fail("Laya trial limit reached; retry after the window resets",429,"laya_rate_limit",{"retry-after":"1"}); + } catch { + return fail("Laya admission control is temporarily unavailable",503,"laya_unavailable",{"retry-after":"1"}); + } + layaAbort = new AbortController(); + layaRun = startLaya(env, layaPlan!, meter, layaTiming, layaAbort.signal); + // Exact admission can still reject before this promise is awaited. + void layaRun.catch(() => {}); + } let gate: AdmissionResult; if (combinedQuota) { const quotas: Omit[] = enterprise ? [] : [{scope:tier,cost:decisions,limit:rpm,daily:TIERS[tier].daily * multiplier}]; quotas.push({scope:`laya:${processing}`,cost:layaPlan!.cost,limit:LAYA_LIMITS[processing].rpm,daily:LAYA_LIMITS[processing].daily}); try { - const response = await admit({LIMITER:env.LIMITER,QUOTAS:env.QUOTAS!},quotaOwner,quotas,!!layaTiming); + // The coordinator always returns its stage timings. `timing=1` also + // forces storage.sync(), turning observability into ~40 ms of latency + // on every anonymous Laya call even though the awaited transactional + // put already preserves the counters. + const response = await admit({LIMITER:env.LIMITER,QUOTAS:env.QUOTAS!},quotaOwner,quotas,false); readQuotaTiming(response, layaTiming ? regularQuotaTiming : undefined); gate = await response.json() as AdmissionResult; if (typeof gate?.limited !== "boolean" || !Number.isFinite(gate.remaining) || (gate.limited ? !["tier","lane"].includes(gate.limitedBy!) : !Number.isFinite(gate.laneRemaining))) throw new Error("Invalid quota admission response"); - if (gate.limited && gate.limitedBy === "lane") return fail("Laya trial limit reached; retry after the window resets",429, - gate.scope === "day" ? "rate_limit_day" : "laya_rate_limit", {"retry-after":String(gate.resetIn ?? 60)}); + if (gate.limited && gate.limitedBy === "lane") { + layaAbort?.abort(); + return fail("Laya trial limit reached; retry after the window resets",429, + gate.scope === "day" ? "rate_limit_day" : "laya_rate_limit", {"retry-after":String(gate.resetIn ?? 60)}); + } layaRemaining = gate.laneRemaining ?? -1; } catch { + layaAbort?.abort(); return fail("Laya admission control is temporarily unavailable",503,"laya_unavailable", {"retry-after":"1"}); } } else gate = enterprise @@ -2009,6 +2038,7 @@ const worker = { : await limited(env, tier, quotaOwner, decisions, multiplier, layaTiming ? regularQuotaTiming : undefined); regularQuotaMs = performance.now() - regularQuotaStarted; if (gate.limited) { + layaAbort?.abort(); const perDay = gate.scope === "day"; // The moment someone runs out of room is the moment to say where more is. // A free caller is told the plan that lifts this exact limit and gets its @@ -2053,12 +2083,12 @@ const worker = { let escalationFailed = 0; try { if (dimensions) { - const r = await classifyMatrix(env, inputs, dimensions, dimensionBatches, tier, instructions, meter, layaPlan, layaTiming); + const r = await classifyMatrix(env, inputs, dimensions, dimensionBatches, tier, instructions, meter, layaPlan, layaTiming, layaRun); matrix = r.results; results = matrix.flat(); escalationFailed = r.escalationFailed; fallbackDecisions = r.fallbackDecisions; - } else ({ results, escalationFailed } = await classifyMany(env, inputs, labels, tier, instructions, multi, meter, layaPlan, layaTiming)); + } else ({ results, escalationFailed } = await classifyMany(env, inputs, labels, tier, instructions, multi, meter, layaPlan, layaTiming, layaRun)); } catch (e) { if (e instanceof LayaError) return fail(e.message, e.status, e.status === 400 ? "laya_input" : e.status === 429 ? "laya_rate_limit" : "laya_unavailable", diff --git a/src/laya.ts b/src/laya.ts index 8cf52d5..2fd5da3 100644 --- a/src/laya.ts +++ b/src/laya.ts @@ -8,6 +8,7 @@ export type LayaEnv = { LAYA_MODAL_KEY?: string; LAYA_MODAL_SECRET?: string; LAYA_ENABLED?: string; + LAYA_FAST_ADMISSION?: RateLimit; LIMITER: DurableObjectNamespace; }; export const LAYA_LIMITS = { @@ -84,7 +85,7 @@ const probability = (v: unknown): v is number => typeof v === "number" && Number export type LayaTiming = { fetchMs?: number; headersMs?: number; backendMs?: number }; -export async function runLaya(env: LayaEnv, plan: LayaPlan, meter?: Meter, timing?: LayaTiming): Promise { +export async function runLaya(env: LayaEnv, plan: LayaPlan, meter?: Meter, timing?: LayaTiming, signal?: AbortSignal): Promise { const url = plan.processing === "fast" ? env.LAYA_FAST_URL : env.LAYA_BULK_URL; if (env.LAYA_ENABLED !== "true" || !url || !env.LAYA_MODAL_KEY || !env.LAYA_MODAL_SECRET) throw new LayaError("Laya trial is currently unavailable", 503); @@ -100,7 +101,9 @@ export async function runLaya(env: LayaEnv, plan: LayaPlan, meter?: Meter, timin try { response = await fetch(url + "/predict", { method: "POST", headers: { "content-type": "application/json", "Modal-Key": env.LAYA_MODAL_KEY, "Modal-Secret": env.LAYA_MODAL_SECRET }, - body: JSON.stringify({ batch }), signal: AbortSignal.timeout(Math.min(15_000, deadline - Date.now())) }); + body: JSON.stringify({ batch }), signal: signal + ? AbortSignal.any([signal, AbortSignal.timeout(Math.min(15_000, deadline - Date.now()))]) + : AbortSignal.timeout(Math.min(15_000, deadline - Date.now())) }); } catch { throw new LayaError("Laya timed out or could not be reached", 503, 5); } if (timing) timing.headersMs = (timing.headersMs ?? 0) + performance.now() - started; if (response.status === 429 || response.status === 503) diff --git a/test/laya.test.ts b/test/laya.test.ts index edd1a59..d8ea36a 100644 --- a/test/laya.test.ts +++ b/test/laya.test.ts @@ -32,17 +32,19 @@ function mockLaya() { return calls; } -function combinedEnv(result: unknown = {limited:false,remaining:2900,laneRemaining:56}) { - const admissions: any[] = []; +function combinedEnv(result: unknown | ((url: string, init: RequestInit) => Promise) = {limited:false,remaining:2900,laneRemaining:56}) { + const admissions: any[] = [], admissionUrls: string[] = []; let legacyCalls = 0; const bindings = {...env, QUOTA_COORDINATOR_ENABLED:"true", LIMITER:{ idFromName:(s: string) => s, get:() => ({fetch:async () => {legacyCalls++; throw new Error("unexpected legacy call");}}), - }, QUOTAS:{idFromName:(s: string) => s,get:() => ({fetch:async (_url: string, init: RequestInit) => { + }, QUOTAS:{idFromName:(s: string) => s,get:() => ({fetch:async (url: string, init: RequestInit) => { + admissionUrls.push(url); admissions.push(JSON.parse(String(init.body))); if (result instanceof Error) throw result; + if (typeof result === "function") return result(url,init); return Response.json(result); }})}} as unknown as Env; - return {bindings,admissions,legacyCalls:() => legacyCalls}; + return {bindings,admissions,admissionUrls,legacyCalls:() => legacyCalls}; } test("combined admission makes one RPC with separate decision and multi-label question costs", async () => { @@ -52,6 +54,7 @@ test("combined admission makes one RPC with separate decision and multi-label qu expect(response.status).toBe(200); expect(response.headers.get("ratelimit-remaining")).toBe("56"); expect(f.admissions).toHaveLength(1); + expect(f.admissionUrls).toEqual(["https://quota/admit"]); expect(f.admissions[0].quotas).toMatchObject([ {scope:"fast",cost:1,limit:3000,daily:20000}, {scope:"laya:fast",cost:2,limit:60,daily:2000}, @@ -73,6 +76,41 @@ test.each([ expect(calls).toHaveLength(0); }); +test("fast local admission lets exact admission and inference overlap", async () => { + mockLaya(); + const predict = globalThis.fetch; + let started!: () => void; + const modalStarted = new Promise(resolve => { started = resolve; }); + globalThis.fetch = (async (...args: Parameters) => { started(); return predict(...args); }) as typeof fetch; + const f = combinedEnv(async () => { + await modalStarted; + return Response.json({limited:false,remaining:2999,laneRemaining:59}); + }); + const bindings = {...f.bindings, LAYA_FAST_ADMISSION:{limit:async () => ({success:true})}} as unknown as Env; + const response = await request(task,bindings); + expect(response.status).toBe(200); + expect(f.admissionUrls).toEqual(["https://quota/admit"]); +}); + +test("exact admission rejection aborts speculative fast inference", async () => { + let started!: () => void, aborted!: () => void; + const modalStarted = new Promise(resolve => { started = resolve; }); + const modalAborted = new Promise(resolve => { aborted = resolve; }); + globalThis.fetch = ((_url: string | URL | Request, init?: RequestInit) => new Promise((_resolve,reject) => { + started(); + init?.signal?.addEventListener("abort",() => { aborted(); reject(new DOMException("aborted","AbortError")); },{once:true}); + })) as typeof fetch; + const f = combinedEnv(async () => { + await modalStarted; + return Response.json({limited:true,remaining:2999,limitedBy:"lane",scope:"minute",resetIn:12}); + }); + const bindings = {...f.bindings, LAYA_FAST_ADMISSION:{limit:async () => ({success:true})}} as unknown as Env; + const response = await request(task,bindings); + expect(response.status).toBe(429); + expect(await response.json()).toMatchObject({code:"laya_rate_limit"}); + await modalAborted; +}); + test("bulk chunks by question count, retains order and meters the selected lane", async () => { const calls = mockLaya(), meter = newMeter(); const tasks = Array.from({ length: 1000 }, (_, i) => ({ ...task, input: String(i) })); diff --git a/tests/runtime-config.test.ts b/tests/runtime-config.test.ts index 1bda5a7..3b43228 100644 --- a/tests/runtime-config.test.ts +++ b/tests/runtime-config.test.ts @@ -3,7 +3,6 @@ import { copyFileSync, existsSync, mkdtempSync, readFileSync, rmSync } from "nod import { tmpdir } from "node:os"; import { join } from "node:path"; import { spawnSync } from "node:child_process"; -import Ajv from "ajv"; const root = new URL("../", import.meta.url); const renderer = new URL(".github/render-wrangler.mjs", root); @@ -36,10 +35,8 @@ test("production rendering uses PostgreSQL secrets and activates verified accoun expect(config.vars.DATABASE_URL).toBeUndefined(); expect(config.secrets.required).toContain("DATABASE_URL"); expect(config.secrets.required).not.toContain("DATABASE_URL_UNPOOLED"); - expect(config.placement).toEqual({ region: "aws:us-west-2" }); - const schema = JSON.parse(readFileSync(new URL("node_modules/wrangler/config-schema.json", root), "utf8")); - const validate = new Ajv({ strict: false }).compile(schema.definitions.RawConfig.properties.placement); - expect(validate(config.placement)).toBe(true); + expect(config.placement).toBeUndefined(); + expect(config.ratelimits).toContainEqual({name:"LAYA_FAST_ADMISSION",namespace_id:"910001",simple:{limit:60,period:60}}); }); test("renderer needs no D1 identifier but still refuses missing Cloudflare bindings", () => { diff --git a/wrangler.example.toml b/wrangler.example.toml index 1363d3e..76f8b2f 100644 --- a/wrangler.example.toml +++ b/wrangler.example.toml @@ -10,11 +10,6 @@ compatibility_flags = ["nodejs_compat"] compatibility_date = "2026-08-01" account_id = "YOUR_CLOUDFLARE_ACCOUNT_ID" -# Run fetch handlers near the Oregon Neon primary. This is a placement target, -# not a data-residency guarantee; assets remain globally served. -[placement] -region = "aws:us-west-2" - # Set the pooled runtime URL as a Worker secret, never a plaintext var. CI uses # a separate direct DATABASE_URL_UNPOOLED secret for migrations. [secrets] @@ -42,12 +37,22 @@ dataset = "classifier_jev_attempts" binding = "ACCOUNT_AE" dataset = "classifier_account_events" +# A local, permissive burst shield lets exact Durable Object admission overlap +# the fast inference call without exposing the single warm GPU to obvious floods. +[[ratelimits]] +name = "LAYA_FAST_ADMISSION" +namespace_id = "910001" + + [ratelimits.simple] + limit = 60 + period = 60 + [vars] # Transfer-aware legacy objects were deployed before enabling admission here. QUOTA_COORDINATOR_ENABLED = "true" LAYA_ENABLED = "true" -LAYA_FAST_URL = "https://miryaboy--classifier-laya-router-trial-fast.us-east.modal.direct" -LAYA_BULK_URL = "https://miryaboy--classifier-laya-router-trial-bulk.us-east.modal.direct" +LAYA_FAST_URL = "https://miryaboy--classifier-laya-router-trial-west-fast.us-west.modal.direct" +LAYA_BULK_URL = "https://miryaboy--classifier-laya-router-trial-west-bulk.us-west.modal.direct" # LAYA_MODAL_KEY and LAYA_MODAL_SECRET are Worker secrets, never plaintext vars. # Enable only after account migration, pricing and Autumn reconciliation checks. APP_ACCOUNTS_ENABLED = "true" From 9cb023139297ae480c4dc493afc357f226b79583 Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Mon, 21 Sep 2026 05:13:03 -0700 Subject: [PATCH 2/2] Retain Oregon placement and remove redundant account plan lookup --- docs/postgres-setup.md | 11 +++++------ inference/laya/README.md | 2 +- src/http/classification.ts | 3 +-- src/index.ts | 6 ++---- src/server/usage.ts | 8 ++++++-- test/laya.test.ts | 19 ++++++++++++++++--- tests/accounts.test.ts | 8 ++++++++ tests/runtime-config.test.ts | 6 +++++- wrangler.example.toml | 5 +++++ 9 files changed, 49 insertions(+), 19 deletions(-) diff --git a/docs/postgres-setup.md b/docs/postgres-setup.md index 66e057a..2097aff 100644 --- a/docs/postgres-setup.md +++ b/docs/postgres-setup.md @@ -101,12 +101,11 @@ connection strings into commands that will remain in shell history. Migrations are transactional, ordered, and checksum-verified; editing an already-applied migration is rejected. -The production Worker executes at Cloudflare's ingress edge. Neon remains in -`aws:us-west-2`; removing the former Worker placement avoids forwarding every -latency-sensitive classification request through Oregon. This is not a data -residency guarantee and does not move the Neon primary or a Durable Object. -Confirm database-dependent account latency as well as inference latency before -changing placement again. See [Cloudflare placement](https://developers.cloudflare.com/workers/configuration/placement/). +The production Worker targets Oregon (`aws:us-west-2`) for fetch execution. +This is not a data-residency guarantee and does not move an existing Neon +database or Durable Object. Verify the active database's region separately; +Worker placement alone does not establish full colocation. +See [Cloudflare placement](https://developers.cloudflare.com/workers/configuration/placement/). Do not enable the new account offering until retail token rates, Autumn events, legacy Pro handling, reconciliation, and end-to-end production checks are ready. diff --git a/inference/laya/README.md b/inference/laya/README.md index 7e91d0d..2fb1d90 100644 --- a/inference/laya/README.md +++ b/inference/laya/README.md @@ -69,7 +69,7 @@ Deploy the transfer-aware `RateLimiter`, coordinator binding and migration with Rollback by disabling `QUOTA_COORDINATOR_ENABLED` while retaining the new classes, bindings and forwarding code. **Do not roll back to code predating the transfer protocol:** its old counters are frozen and no longer authoritative. The disabled path continues to follow transferred counters. Neither successful admission nor forwarding uses unconfirmed storage writes. -Run `node inference/laya/latency.mjs ` before and after deployment from the same client. It records 20 sequential synthetic requests, separates the first call, and reports end-to-end, Worker, quota, Modal and backend timings. Do not equate backend compute time or the Worker-to-Modal span with client latency. The Worker executes at Cloudflare's ingress edge; Modal ingress and compute are colocated in `us-west`. Measure other client regions independently. +Run `node inference/laya/latency.mjs ` before and after deployment from the same client. It records 20 sequential synthetic requests, separates the first call, and reports end-to-end, Worker, quota, Modal and backend timings. Do not equate backend compute time or the Worker-to-Modal span with client latency. The Worker targets Oregon (`aws:us-west-2`); Modal ingress and compute use `us-west`. Existing Durable Objects and databases are not relocated by these settings. Measure other client regions independently. ## Evidence diff --git a/src/http/classification.ts b/src/http/classification.ts index 35edc58..bffde56 100644 --- a/src/http/classification.ts +++ b/src/http/classification.ts @@ -34,7 +34,6 @@ export async function accountClassification(request: Request, env: AppEnv & Part type: `${source} · ${body.dimensions ? "Dimensions" : body.multi ? "Multi-label" : "Single-label"}`, meteringMode: "tokens", }); if (!reservation) throw new AppError(401, "Missing account credential."); - const plan = await env.APP_DB.prepare("SELECT billing_plan FROM app_accounts WHERE id=?").bind(accountId).first<{ billing_plan: string }>(); const meter = newMeter(); let admissionError: AppError | undefined; let reservationQueue = Promise.resolve(); @@ -67,7 +66,7 @@ export async function accountClassification(request: Request, env: AppEnv & Part let response: Response; try { response = await worker.fetch(new Request(request.url, { method: "POST", headers: request.headers, body: text }), env as Env, ctx, - { account: { id: accountId, multiplier: plan?.billing_plan === "free" ? 1 : 10 }, meter }); + { account: { id: accountId, multiplier: reservation.billingPlan === "free" ? 1 : 10 }, meter }); await reservationQueue; if (admissionError) throw admissionError; } catch (error) { diff --git a/src/index.ts b/src/index.ts index 2031ec0..f3d4221 100644 --- a/src/index.ts +++ b/src/index.ts @@ -2013,10 +2013,8 @@ const worker = { const quotas: Omit[] = enterprise ? [] : [{scope:tier,cost:decisions,limit:rpm,daily:TIERS[tier].daily * multiplier}]; quotas.push({scope:`laya:${processing}`,cost:layaPlan!.cost,limit:LAYA_LIMITS[processing].rpm,daily:LAYA_LIMITS[processing].daily}); try { - // The coordinator always returns its stage timings. `timing=1` also - // forces storage.sync(), turning observability into ~40 ms of latency - // on every anonymous Laya call even though the awaited transactional - // put already preserves the counters. + // The coordinator returns stage timings without the diagnostic + // storage.sync(); its output gate still preserves counter durability. const response = await admit({LIMITER:env.LIMITER,QUOTAS:env.QUOTAS!},quotaOwner,quotas,false); readQuotaTiming(response, layaTiming ? regularQuotaTiming : undefined); gate = await response.json() as AdmissionResult; diff --git a/src/server/usage.ts b/src/server/usage.ts index 371229b..cd95539 100644 --- a/src/server/usage.ts +++ b/src/server/usage.ts @@ -4,6 +4,7 @@ export interface Reservation { accountId: string; agentId: string; cost: number; + billingPlan: string; } /** Serializable transaction: conditional ledger insertion, account and agent debit. */ export async function authorizeAndReserve( @@ -62,7 +63,7 @@ export async function authorizeAndReserve( agent.account_id, ), env.APP_DB.prepare( - "UPDATE app_accounts SET balance=balance-?,paid_balance=paid_balance-(SELECT paid_credits FROM app_usage WHERE id=?) WHERE id=? AND EXISTS(SELECT 1 FROM app_usage WHERE id=?)", + "UPDATE app_accounts SET balance=balance-?,paid_balance=paid_balance-(SELECT paid_credits FROM app_usage WHERE id=?) WHERE id=? AND EXISTS(SELECT 1 FROM app_usage WHERE id=?) RETURNING billing_plan", ).bind(cost, id, agent.account_id, id), env.APP_DB.prepare( "UPDATE app_agents SET used=used+? WHERE id=? AND EXISTS(SELECT 1 FROM app_usage WHERE id=?)", @@ -73,7 +74,10 @@ export async function authorizeAndReserve( 403, "Credential is paused or revoked, or the workspace balance is too low.", ); - return { id, accountId: agent.account_id, agentId: agent.id, cost }; + // Reuse the account row already locked by reservation instead of making + // classification pay another database round trip to look up its plan. + return { id, accountId: agent.account_id, agentId: agent.id, cost, + billingPlan: results[1].results[0].billing_plan as string }; } /** Idempotent settlement; failures restore both account and agent reservations. */ export async function completeReservation( diff --git a/test/laya.test.ts b/test/laya.test.ts index d8ea36a..f6b88f4 100644 --- a/test/laya.test.ts +++ b/test/laya.test.ts @@ -92,7 +92,7 @@ test("fast local admission lets exact admission and inference overlap", async () expect(f.admissionUrls).toEqual(["https://quota/admit"]); }); -test("exact admission rejection aborts speculative fast inference", async () => { +test.each(["lane", "tier"])("exact %s admission rejection aborts speculative fast inference", async limitedBy => { let started!: () => void, aborted!: () => void; const modalStarted = new Promise(resolve => { started = resolve; }); const modalAborted = new Promise(resolve => { aborted = resolve; }); @@ -102,15 +102,28 @@ test("exact admission rejection aborts speculative fast inference", async () => })) as typeof fetch; const f = combinedEnv(async () => { await modalStarted; - return Response.json({limited:true,remaining:2999,limitedBy:"lane",scope:"minute",resetIn:12}); + return Response.json({limited:true,remaining:2999,limitedBy,scope:"minute",resetIn:12}); }); const bindings = {...f.bindings, LAYA_FAST_ADMISSION:{limit:async () => ({success:true})}} as unknown as Env; const response = await request(task,bindings); expect(response.status).toBe(429); - expect(await response.json()).toMatchObject({code:"laya_rate_limit"}); + expect(await response.json()).toMatchObject({code:limitedBy === "lane" ? "laya_rate_limit" : "rate_limit_minute"}); await modalAborted; }); +test.each([false, true])("local burst shield blocks inference when unavailable=%s", async unavailable => { + const calls = mockLaya(); + const f = combinedEnv(); + const bindings = {...f.bindings, LAYA_FAST_ADMISSION:{limit:async () => { + if (unavailable) throw new Error("unavailable"); + return {success:false}; + }}} as unknown as Env; + const response = await request(task,bindings); + expect(response.status).toBe(unavailable ? 503 : 429); + expect(calls).toHaveLength(0); + expect(f.admissions).toHaveLength(0); +}); + test("bulk chunks by question count, retains order and meters the selected lane", async () => { const calls = mockLaya(), meter = newMeter(); const tasks = Array.from({ length: 1000 }, (_, i) => ({ ...task, input: String(i) })); diff --git a/tests/accounts.test.ts b/tests/accounts.test.ts index 6aeb140..836a37c 100644 --- a/tests/accounts.test.ts +++ b/tests/accounts.test.ts @@ -33,6 +33,14 @@ beforeEach(async () => { ); }); describe("account lifecycle", () => { + test.each(["free", "pro"])("reservation returns the current %s plan with its account debit", async (plan) => { + const enrolled = await enroll(); + await env.APP_DB.prepare("UPDATE app_accounts SET billing_plan=? WHERE id='local-demo'").bind(plan).run(); + const reservation = await authorizeAndReserve(request(enrolled.secret), env, 1); + expect(reservation?.billingPlan).toBe(plan); + expect(reservation?.accountId).toBe("local-demo"); + expect((await env.APP_DB.prepare("SELECT balance FROM app_accounts WHERE id='local-demo'").first())?.balance).toBe(499999); + }); test("authentication transitions clear the previous workspace selection", () => { const callback = clearWorkspaceSelection( new Response(null, { diff --git a/tests/runtime-config.test.ts b/tests/runtime-config.test.ts index 3b43228..fe9c230 100644 --- a/tests/runtime-config.test.ts +++ b/tests/runtime-config.test.ts @@ -3,6 +3,7 @@ import { copyFileSync, existsSync, mkdtempSync, readFileSync, rmSync } from "nod import { tmpdir } from "node:os"; import { join } from "node:path"; import { spawnSync } from "node:child_process"; +import Ajv from "ajv"; const root = new URL("../", import.meta.url); const renderer = new URL(".github/render-wrangler.mjs", root); @@ -35,7 +36,10 @@ test("production rendering uses PostgreSQL secrets and activates verified accoun expect(config.vars.DATABASE_URL).toBeUndefined(); expect(config.secrets.required).toContain("DATABASE_URL"); expect(config.secrets.required).not.toContain("DATABASE_URL_UNPOOLED"); - expect(config.placement).toBeUndefined(); + expect(config.placement).toEqual({ region: "aws:us-west-2" }); + const schema = JSON.parse(readFileSync(new URL("node_modules/wrangler/config-schema.json", root), "utf8")); + const validate = new Ajv({ strict: false }).compile(schema.definitions.RawConfig.properties.placement); + expect(validate(config.placement)).toBe(true); expect(config.ratelimits).toContainEqual({name:"LAYA_FAST_ADMISSION",namespace_id:"910001",simple:{limit:60,period:60}}); }); diff --git a/wrangler.example.toml b/wrangler.example.toml index 76f8b2f..f3f2080 100644 --- a/wrangler.example.toml +++ b/wrangler.example.toml @@ -10,6 +10,11 @@ compatibility_flags = ["nodejs_compat"] compatibility_date = "2026-08-01" account_id = "YOUR_CLOUDFLARE_ACCOUNT_ID" +# Run fetch handlers near Oregon. This is a placement target, not a +# data-residency guarantee; assets remain globally served. +[placement] +region = "aws:us-west-2" + # Set the pooled runtime URL as a Worker secret, never a plaintext var. CI uses # a separate direct DATABASE_URL_UNPOOLED secret for migrations. [secrets]