diff --git a/cli/classify.js b/cli/classify.js index 47e7e04..efc12fa 100755 --- a/cli/classify.js +++ b/cli/classify.js @@ -121,8 +121,8 @@ function parseArgs(argv) { if (!(Number.isInteger(o.max) && o.max > 0)) fail("--max takes a whole number above 0, e.g. --max 3"); o.multi = true; // the API reads max_labels as multi-label; a single label cannot be capped } - if (!["jev", "laya"].includes(o.model)) fail("--model must be jev or laya"); - if (!["fast", "bulk"].includes(o.processing) || o.model !== "laya" && o.processing !== "fast") fail("--processing bulk requires --model laya"); + if (!["jev", "laya", "kev"].includes(o.model)) fail("--model must be jev, laya or kev"); + if (!["fast", "bulk"].includes(o.processing) || o.model === "jev" && o.processing !== "fast") fail("--processing bulk requires --model laya or --model kev"); return o; } @@ -175,7 +175,7 @@ async function readStdin() { async function post(o, inputs) { const body = { inputs, labels: o.labels }; - if (o.model === "laya") { body.model = "laya"; body.processing = o.processing; } + if (o.model !== "jev") { body.model = o.model; body.processing = o.processing; } if (o.multi) body.multi = true; if (o.max) body.max_labels = o.max; if (o.smart) body.tier = "smart"; @@ -184,8 +184,8 @@ async function post(o, inputs) { if (o.apiKey) headers.authorization = `Bearer ${o.apiKey}`; let last = ""; - const attempts = o.model === "laya" ? 20 : ATTEMPTS; - const deadline = o.model === "laya" ? Date.now() + 180_000 : Infinity; + const attempts = o.model !== "jev" ? 20 : ATTEMPTS; + const deadline = o.model !== "jev" ? Date.now() + 180_000 : Infinity; for (let attempt = 0; attempt < attempts && Date.now() < deadline; attempt++) { let res; try { @@ -213,7 +213,7 @@ async function post(o, inputs) { // The API says exactly how long; a batch that trips the minute window // resumes on its own rather than dying at item 7,400. The first wait // also says, once, that the wait is optional. - if (o.model !== "laya" && !o.apiKey && !hinted) { hinted = true; process.stderr.write(`classify: free limits are per IP; Pro lifts them 10x for $20/month: ${payload.upgrade || "https://classifier.dev/pro"}\n`); } + if (o.model === "jev" && !o.apiKey && !hinted) { hinted = true; process.stderr.write(`classify: free limits are per IP; Pro lifts them 10x for $20/month: ${payload.upgrade || "https://classifier.dev/pro"}\n`); } await retry("rate limited", Number(res.headers.get("retry-after")) || 5, attempt, attempts, deadline); continue; } @@ -244,7 +244,7 @@ function configuredBatch() { } function batchSize(o) { - if (o.model === "laya") return Math.min(configuredBatch(), o.processing === "fast" ? 1 : Math.floor(1000 / (o.multi ? o.labels.length : 1)), o.smart ? SMART_BATCH : MAX_BATCH); + if (o.model !== "jev") return Math.min(configuredBatch(), o.processing === "fast" ? 1 : Math.floor(1000 / (o.multi ? o.labels.length : 1)), o.smart ? SMART_BATCH : MAX_BATCH); return Math.min(configuredBatch(), o.smart && !o.apiKey ? SMART_BATCH : MAX_BATCH); } @@ -294,7 +294,7 @@ export async function classifyAll(o, items, onReady, onProgress) { // Public smart traffic is limited to 200 classifications/minute. One // worker keeps concurrent batches from spending that whole window at once; // partner keys can use the normal four workers. - const concurrency = o.model === "laya" || o.smart && !o.apiKey ? 1 : CONCURRENCY; + const concurrency = o.model !== "jev" || o.smart && !o.apiKey ? 1 : CONCURRENCY; await Promise.all( Array.from({ length: Math.min(concurrency, batches.length) }, async () => { while (next < batches.length) { diff --git a/inference/laya/LICENSE.upstream b/inference/laya/LICENSE.upstream deleted file mode 100644 index d9a10c0..0000000 --- a/inference/laya/LICENSE.upstream +++ /dev/null @@ -1,176 +0,0 @@ - Apache License - Version 2.0, January 2004 - http://www.apache.org/licenses/ - - TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION - - 1. Definitions. - - "License" shall mean the terms and conditions for use, reproduction, - and distribution as defined by Sections 1 through 9 of this document. - - "Licensor" shall mean the copyright owner or entity authorized by - the copyright owner that is granting the License. - - "Legal Entity" shall mean the union of the acting entity and all - other entities that control, are controlled by, or are under common - control with that entity. For the purposes of this definition, - "control" means (i) the power, direct or indirect, to cause the - direction or management of such entity, whether by contract or - otherwise, or (ii) ownership of fifty percent (50%) or more of the - outstanding shares, or (iii) beneficial ownership of such entity. - - "You" (or "Your") shall mean an individual or Legal Entity - exercising permissions granted by this License. - - "Source" form shall mean the preferred form for making modifications, - including but not limited to software source code, documentation - source, and configuration files. - - "Object" form shall mean any form resulting from mechanical - transformation or translation of a Source form, including but - not limited to compiled object code, generated documentation, - and conversions to other media types. - - "Work" shall mean the work of authorship, whether in Source or - Object form, made available under the License, as indicated by a - copyright notice that is included in or attached to the work - (an example is provided in the Appendix below). - - "Derivative Works" shall mean any work, whether in Source or Object - form, that is based on (or derived from) the Work and for which the - editorial revisions, annotations, elaborations, or other modifications - represent, as a whole, an original work of authorship. For the purposes - of this License, Derivative Works shall not include works that remain - separable from, or merely link (or bind by name) to the interfaces of, - the Work and Derivative Works thereof. - - "Contribution" shall mean any work of authorship, including - the original version of the Work and any modifications or additions - to that Work or Derivative Works thereof, that is intentionally - submitted to Licensor for inclusion in the Work by the copyright owner - or by an individual or Legal Entity authorized to submit on behalf of - the copyright owner. For the purposes of this definition, "submitted" - means any form of electronic, verbal, or written communication sent - to the Licensor or its representatives, including but not limited to - communication on electronic mailing lists, source code control systems, - and issue tracking systems that are managed by, or on behalf of, the - Licensor for the purpose of discussing and improving the Work, but - excluding communication that is conspicuously marked or otherwise - designated in writing by the copyright owner as "Not a Contribution." - - "Contributor" shall mean Licensor and any individual or Legal Entity - on behalf of whom a Contribution has been received by Licensor and - subsequently incorporated within the Work. - - 2. Grant of Copyright License. Subject to the terms and conditions of - this License, each Contributor hereby grants to You a perpetual, - worldwide, non-exclusive, no-charge, royalty-free, irrevocable - copyright license to reproduce, prepare Derivative Works of, - publicly display, publicly perform, sublicense, and distribute the - Work and such Derivative Works in Source or Object form. - - 3. Grant of Patent License. Subject to the terms and conditions of - this License, each Contributor hereby grants to You a perpetual, - worldwide, non-exclusive, no-charge, royalty-free, irrevocable - (except as stated in this section) patent license to make, have made, - use, offer to sell, sell, import, and otherwise transfer the Work, - where such license applies only to those patent claims licensable - by such Contributor that are necessarily infringed by their - Contribution(s) alone or by combination of their Contribution(s) - with the Work to which such Contribution(s) was submitted. If You - institute patent litigation against any entity (including a - cross-claim or counterclaim in a lawsuit) alleging that the Work - or a Contribution incorporated within the Work constitutes direct - or contributory patent infringement, then any patent licenses - granted to You under this License for that Work shall terminate - as of the date such litigation is filed. - - 4. Redistribution. You may reproduce and distribute copies of the - Work or Derivative Works thereof in any medium, with or without - modifications, and in Source or Object form, provided that You - meet the following conditions: - - (a) You must give any other recipients of the Work or - Derivative Works a copy of this License; and - - (b) You must cause any modified files to carry prominent notices - stating that You changed the files; and - - (c) You must retain, in the Source form of any Derivative Works - that You distribute, all copyright, patent, trademark, and - attribution notices from the Source form of the Work, - excluding those notices that do not pertain to any part of - the Derivative Works; and - - (d) If the Work includes a "NOTICE" text file as part of its - distribution, then any Derivative Works that You distribute must - include a readable copy of the attribution notices contained - within such NOTICE file, excluding those notices that do not - pertain to any part of the Derivative Works, in at least one - of the following places: within a NOTICE text file distributed - as part of the Derivative Works; within the Source form or - documentation, if provided along with the Derivative Works; or, - within a display generated by the Derivative Works, if and - wherever such third-party notices normally appear. The contents - of the NOTICE file are for informational purposes only and - do not modify the License. You may add Your own attribution - notices within Derivative Works that You distribute, alongside - or as an addendum to the NOTICE text from the Work, provided - that such additional attribution notices cannot be construed - as modifying the License. - - You may add Your own copyright statement to Your modifications and - may provide additional or different license terms and conditions - for use, reproduction, or distribution of Your modifications, or - for any such Derivative Works as a whole, provided Your use, - reproduction, and distribution of the Work otherwise complies with - the conditions stated in this License. - - 5. Submission of Contributions. Unless You explicitly state otherwise, - any Contribution intentionally submitted for inclusion in the Work - by You to the Licensor shall be under the terms and conditions of - this License, without any additional terms or conditions. - Notwithstanding the above, nothing herein shall supersede or modify - the terms of any separate license agreement you may have executed - with Licensor regarding such Contributions. - - 6. Trademarks. This License does not grant permission to use the trade - names, trademarks, service marks, or product names of the Licensor, - except as required for reasonable and customary use in describing the - origin of the Work and reproducing the content of the NOTICE file. - - 7. Disclaimer of Warranty. Unless required by applicable law or - agreed to in writing, Licensor provides the Work (and each - Contributor provides its Contributions) on an "AS IS" BASIS, - WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or - implied, including, without limitation, any warranties or conditions - of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A - PARTICULAR PURPOSE. You are solely responsible for determining the - appropriateness of using or redistributing the Work and assume any - risks associated with Your exercise of permissions under this License. - - 8. Limitation of Liability. In no event and under no legal theory, - whether in tort (including negligence), contract, or otherwise, - unless required by applicable law (such as deliberate and grossly - negligent acts) or agreed to in writing, shall any Contributor be - liable to You for damages, including any direct, indirect, special, - incidental, or consequential damages of any character arising as a - result of this License or out of the use or inability to use the - Work (including but not limited to damages for loss of goodwill, - work stoppage, computer failure or malfunction, or any and all - other commercial damages or losses), even if such Contributor - has been advised of the possibility of such damages. - - 9. Accepting Warranty or Additional Liability. While redistributing - the Work or Derivative Works thereof, You may choose to offer, - and charge a fee for, acceptance of support, warranty, indemnity, - or other liability obligations and/or rights consistent with this - License. However, in accepting such obligations, You may act only - on Your own behalf and on Your sole responsibility, not on behalf - of any other Contributor, and only if You agree to indemnify, - defend, and hold each Contributor harmless for any liability - incurred by, or claims asserted against, such Contributor by reason - of your accepting any such warranty or additional liability. - - END OF TERMS AND CONDITIONS diff --git a/inference/laya/README.md b/inference/laya/README.md index 2fb1d90..142e0ba 100644 --- a/inference/laya/README.md +++ b/inference/laya/README.md @@ -1,80 +1,76 @@ -# Laya fast/bulk trial +# Laya and Kev on Beam -Opt-in POST `/v1/classify` with `model: "laya"` and `processing: "fast"` (default) or `"bulk"`. Existing clients still use Jev. Both lanes use the same preloaded Laya 0.3.4 Router; `tier: "smart"` independently reviews uncertain answers. Laya accuracy/calibration has not been established by the Jev evaluation. Each result's existing `model` field identifies its actual checkpoint and lane; mixed-language batches can return `model: "mixed"` at the top level. +Both models are hosted by Beam on its shared inference endpoints. There is no +deployment of our own any more: the Modal app `classifier-laya-router-trial-west`, +its adapter, its image build and its smoke tests were removed when production +moved to Beam. ```sh curl https://classifier.dev/v1/classify -H 'Content-Type: application/json' \ -d '{"model":"laya","processing":"fast","input":"Please refund this charge","labels":["billing","technical"]}' -node cli/classify.js billing,technical --model laya --processing bulk < tickets.txt +curl https://classifier.dev/v1/classify -H 'Content-Type: application/json' \ + -d '{"model":"kev","input":"Please refund this charge","labels":["billing","technical"]}' ``` -## Capacity and limits +## What Beam serves -| | Fast | Bulk | +| | `jev/laya` | `jev/kev` | |---|---|---| -| GPU pool | One warm L4; max one | Zero to one L4; idle shutdown after 60s | -| Request | One decision, up to four questions | Up to 1,000 questions; sequential chunks of at most 64 | -| Per caller | 60 questions/min, 2,000/day | 1,000 questions/min, 20,000/day | -| GPU admission | One request at a time; 20 questions/s, burst four | At most four requests in memory, one GPU forward at a time; queue wait at most 5s | - -One single-label decision is one question. Multi-label asks one question per label. Lane quotas apply to paid/operator keys too, in addition to normal tier quotas, and count attempts (including failed inference and retries). Rejected work returns 429 with Retry-After; cold/unavailable bulk returns 503. No silent lane/model fallback. A cold 1,000-question attempt can spend that minute's quota; retries must respect the next window. Source CLI retries are bounded at three minutes per batch. SDK/CLI additions are in this repository, not a new package release. - -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-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, 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 **$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. - -## Model-card compliance - -Checked against the [upstream model card](https://huggingface.co/convaiinnovations/laya) on 2026-09-20. This trial uses its recommended **preloaded Router**. English and multilingual inputs follow the SDK's unchanged routing heuristic. All three checkpoints are resident; typed-decisions is loaded but the default upstream Router does not automatically select it, and this API adds no checkpoint override. Language detection is heuristic, not an accuracy guarantee. Mixed-language rows are grouped by selected checkpoint, questions share one forward pass per checkpoint, and results are restored to input order. Fast and bulk use identical routing. - -The image sets `USE_TF=0` to avoid the documented Transformers/TensorFlow initialization issue. We pin checkpoint revision `1c5edc17a7acd8701df6fc341c0d179f1c62c982`, bake all checkpoints into the image, and attach the locally loaded agents to `Router(max_loaded=3)` before preloading. HF/Transformers offline mode prevents downloads at startup or inference; language switches do not rebuild or evict models. CUDA failure aborts startup instead of silently serving on CPU. The adapter uses Laya's own sequence/marker construction, collation, option-count temperature buckets and probability decoding. Its intentional change is rejecting oversized inputs instead of silently truncating them. - -The published fine-tuned benchmark wins are **not general zero-shot accuracy claims for this API**. Choice confidence is entropy-based certainty, not the winning label's probability. Shipped confidence can be overconfident; no classifier.dev-specific temperature fitting has been performed. Smart's existing confidence threshold is experimental with Laya and cannot catch confidently wrong answers. Evaluate and calibrate on representative, consented data before expanding this trial or trusting score thresholds. - -## Operate and stop - -Deploy the backend from the repository root: - -```sh -uvx --from modal modal deploy inference/laya/deploy.py -``` - -Worker configuration lives in `wrangler.example.toml`; production deploys from main. Worker secrets `LAYA_MODAL_KEY` and `LAYA_MODAL_SECRET` contain a dedicated Modal proxy credential. Do not place them in source or logs. - -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: +| Model | ModernBERT-large, trained decision head | Qwen2.5-0.5B, LoRA + pointer head | +| Context | 512 tokens (state plus one question) | 8,192 tokens (one packed sequence) | +| Questions per request | 32 upstream; we pack at most 16 items | 32 | +| Price | $0.021 per 1M input tokens, output free | same | + +Beam caps every request at 32 named questions whatever the model, so a large +batch becomes several requests, run concurrently and reassembled in input +order by `runJevBatches`. Laya's item ceiling is 16, chosen from what Beam +actually accepts: 20 items of ordinary support text passed and 25 were +refused. A context refusal is translated to `max_tokens_exceeded`, which makes +the batch halve and retry, so an underestimate costs a round trip rather than +the request. + +## Credential and configuration + +`BEAM_API_KEY` is a Worker secret holding a Beam workspace token. It is the +only credential either model needs; requests are billed per token to that +workspace. `LAYA_ENABLED` must be `"true"`. There are no URLs to configure — +`src/jev.ts` owns the endpoint — and no `LAYA_MODAL_KEY`/`LAYA_MODAL_SECRET` +any more. ```sh -uvx --from modal modal app stop classifier-laya-router-trial-west +npx wrangler secret put BEAM_API_KEY ``` -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. - -Watch model labels `laya-0.3.4--` in existing classifier analytics for successes and latency; rejected calls use `laya-0.3.4-routed-`. Account token usage aggregates under the routed lane model at the explicit zero trial rate. Compare with Modal allocation/invocation metrics. Measure client round-trip separately from the response's server processing time. Never add caller text or raw IP to diagnostic logs. - -## 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. 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. +## Capacity and limits -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. +| | Fast | Bulk | +|---|---|---| +| Request | One decision, up to four questions | Up to 1,000 questions | +| Per caller | 60 questions/min, 2,000/day | 1,000 questions/min, 20,000/day | -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. +Lane quotas are a product decision about shared capacity, not a property of +Beam, and they are unchanged by the move. They apply to paid and operator keys +too, and count attempts. Neither lane has a pool to start, so cold-start 503s +are gone; a Beam refusal is returned as it happened rather than retried, +because the quota has already counted the attempt and a model at capacity will +not clear inside a backoff. -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. +## Cost -## Evidence +Usage-priced, with no idle cost. At $0.021 per 1M input tokens and roughly 50 +input tokens for a short single-label decision, a million such decisions cost +about **$1.05**. The Modal deployment this replaced was billed by allocation: +$1.6216/hour per lane after the narrow-region multiplier, about **$1,168 per +30-day month** for the warm fast lane alone, whether or not it served traffic. -`live-checks.json` records synthetic direct-backend checks: both lanes accepted valid input, rejected invalid/oversized context, rejected unauthenticated requests, shed a 24-request concurrent flash with 429 rather than 5xx, and recovered. This is a bounded smoke test, not a sustained throughput or global latency guarantee. +Retail is zero during the trial, so the input-token spend is ours. Because +Beam reports token counts, account analytics now records a real provider cost +for these models instead of marking it unknown as it did for Modal. -`sdk-check.json` compares the batched adapter with official `Router.predict` on an L4 using English, Hindi and Spanish. Routing, answers, token counts and ordering match; checkpoint identities remain resident. Small BF16 batch-shape differences are checked within a 0.01 score tolerance. `public-smoke.mjs` checks the actual classifier.dev path, including 128-row bulk and mixed-language result models; it writes a local report to ignored `eval/data/laya-public-checks.json`. +## Measuring -Local regression checks include worker routing/order/failures/quota/billing, CLI forwarding, Python/Go forwarding, and asynchronous backend admission. For authenticated backend rechecks, pass a mode-600 JSON file containing the two Worker secret names to `smoke.py`; never commit that file. +`node inference/laya/latency.mjs ` records 20 sequential +synthetic requests from one client, separates the first call, and reports +end-to-end, Worker and quota timings. Do not equate the Worker-to-Beam span +with client latency, and measure other client regions independently: one +shared endpoint is not a replica in every region. diff --git a/inference/laya/adapter.py b/inference/laya/adapter.py deleted file mode 100644 index 74b8bf0..0000000 --- a/inference/laya/adapter.py +++ /dev/null @@ -1,121 +0,0 @@ -# SPDX-License-Identifier: Apache-2.0 -# Adapted from Laya 0.3.4 (https://github.com/NandhaKishorM/laya). -# Modified for multi-state batching, checkpoint grouping and strict context rejection. -# See LICENSE.upstream for the upstream license. - -def predict_routed_batch(router, requests: list[dict], predict=None) -> list[dict]: - """Use the SDK route decision, batch by resident checkpoint, restore caller order.""" - predict = predict or predict_batch - groups = {} - for index, request in enumerate(requests): - decision = dict(router.route(request["state"], request["questions"])) - name = decision["model"] - if name not in router.loaded: - raise RuntimeError("Required Laya checkpoint was not preloaded") - groups.setdefault(name, []).append((index, request, decision)) - results = [None] * len(requests) - for name, rows in groups.items(): - answers = predict(router.load(name), [row[1] for row in rows]) - if len(answers) != len(rows): - raise RuntimeError("Incomplete checkpoint batch") - for (index, _, decision), answer in zip(rows, answers): - answer["routing"] = decision - answer["model"] = "laya-0.3.4/" + name - results[index] = answer - return results - - -def predict_batch(agent, requests: list[dict]) -> list[dict]: - """Flatten independent requests into one encoder batch, then split their answers.""" - import numpy as np - import torch - from laya.common import ( - QTYPES, - build_sequence, - collate_items, - confidence_from_probs, - render_options, - temp_bucket, - ) - - items = [] - metadata = [] - max_len = agent.cfg.get("max_len", 512) - head_max_len = agent.cfg.get("head_max_len", 192) - results = [{"model": "laya-0.3.4/english", "answers": {}} for _ in requests] - for request_index, request in enumerate(requests): - state, questions = request.get("state"), request.get("questions") - if state is None or not isinstance(questions, dict) or not questions: - raise ValueError(f"batch item {request_index} must contain state and non-empty questions") - for question_id, question in questions.items(): - internal = agent._to_internal(question) - # Reject instead of silently truncating caller text or label definitions. - opts = render_options(internal) - # Match the SDK's special-mask sanitization before checking budgets. - mask = agent.tok.mask_token - option_ids = [agent.tok(" " + o.replace(mask, " "), add_special_tokens=False)["input_ids"] for o in opts] - head_ids = agent.tok("%s question: %s" % (internal["t"], internal["ins"].replace(mask, " ")), add_special_tokens=False)["input_ids"] - if any(len(o) > 48 for o in option_ids) or sum(len(o) + 1 for o in option_ids) + max(16, len(head_ids)) > head_max_len: - raise ValueError("question or labels exceed Laya's context budget") - sequence, markers = build_sequence(agent.tok, state, internal, 100_000, head_max_len) - if len(sequence) > max_len: - raise ValueError(f"text and question exceed Laya's {max_len}-token context") - if len(markers) != len(render_options(internal)): - raise ValueError(f"question {question_id!r} exceeds head_max_len={head_max_len}") - items.append({"ids": sequence, "markers": markers, "qtype": QTYPES[internal["t"]]}) - metadata.append((request_index, question_id, internal, len(markers))) - - batch = collate_items([items], agent.tok.pad_token_id) - with torch.no_grad(), torch.autocast(device_type=agent.device.type, dtype=agent.dtype, enabled=agent.device.type == "cuda"): - logits, actions = agent.model( - batch["input_ids"].to(agent.device), - batch["attention_mask"].to(agent.device), - batch["marker_pos"].to(agent.device), - batch["marker_mask"].to(agent.device), - batch["qtype"].to(agent.device), - ) - logits = logits.float().cpu().numpy() - actions = torch.softmax(actions.float(), -1).cpu().numpy() - - token_counts = [0] * len(requests) - for row, (request_index, question_id, question, option_count) in enumerate(metadata): - token_counts[request_index] += int(batch["attention_mask"][row].sum()) - question_type = QTYPES[question["t"]] - scale = agent.temperature_by_options.get( - temp_bucket(question_type, option_count), agent.temperature[question_type] - ) - values = logits[row, :option_count] / max(1e-3, float(scale)) - probabilities = np.exp(values - values.max()) - probabilities /= probabilities.sum() - confidence = round(confidence_from_probs(probabilities, option_count), 4) - action = {"act_probability": round(float(actions[row, 0]), 4)} - if question["t"] == "choice": - keys = list(question["crit"].keys()) - answer = { - "type": "choice", - "choice": keys[int(probabilities.argmax())], - "probabilities": {key: round(float(value), 4) for key, value in zip(keys, probabilities)}, - "confidence": confidence, - "action": action, - } - elif question["t"] == "score": - answer = { - "type": "score", - "score": round(float((np.arange(option_count) * probabilities).sum()), 4), - "legend": {str(index): criterion for index, criterion in enumerate(question["crit"])}, - "probabilities": {str(index): round(float(value), 4) for index, value in enumerate(probabilities)}, - "confidence": confidence, - "action": action, - } - else: - answer = { - "type": "noul", - "noul": round(float(probabilities[1]), 4), - "confidence": round(max(float(probabilities[1]), 1.0 - float(probabilities[1])), 4), - "action": action, - } - results[request_index]["answers"][question_id] = answer - - for result, tokens in zip(results, token_counts): - result["usage"] = {"input_tokens": tokens, "output_tokens": 0} - return results diff --git a/inference/laya/deploy.py b/inference/laya/deploy.py deleted file mode 100644 index 2e75d90..0000000 --- a/inference/laya/deploy.py +++ /dev/null @@ -1,43 +0,0 @@ -"""Global trial: one warm fast GPU, bulk scales to zero; one GPU maximum per lane.""" -from pathlib import Path -import modal - -MODEL = "convaiinnovations/laya" -REVISION = "1c5edc17a7acd8701df6fc341c0d179f1c62c982" - - -def download(): - from huggingface_hub import snapshot_download - snapshot_download(MODEL, revision=REVISION) - - -root = Path(__file__).parent -image = (modal.Image.debian_slim(python_version="3.11") - .pip_install("fastapi[standard]==0.116.1", "laya==0.3.4", "torch==2.7.1", "transformers==4.56.2") - .run_function(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-west") -resources = dict(image=image, gpu="L4", cpu=2, memory=4096, - # 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) - - -@app.server(**resources, min_containers=1, scaledown_window=300) -class Fast: - @modal.enter() - def start(self): - from runtime import start_server - self.server = start_server("fast") - - -@app.server(**resources, min_containers=0, scaledown_window=60) -class Bulk: - @modal.enter() - def start(self): - from runtime import start_server - self.server = start_server("bulk") diff --git a/inference/laya/live-checks.json b/inference/laya/live-checks.json deleted file mode 100644 index 0917489..0000000 --- a/inference/laya/live-checks.json +++ /dev/null @@ -1,28 +0,0 @@ -{ - "fast": { - "valid_rows": 1, - "invalid": 400, - "oversized_context": 400, - "burst": { - "200": 3, - "429": 21 - }, - "burst_elapsed_seconds": 0.48, - "recovery": 200, - "unauthenticated": 401, - "english_and_hindi_routes": "passed" - }, - "bulk": { - "valid_rows": 64, - "invalid": 400, - "oversized_context": 400, - "burst": { - "200": 7, - "429": 17 - }, - "burst_elapsed_seconds": 0.88, - "recovery": 200, - "unauthenticated": 401, - "english_and_hindi_routes": "passed" - } -} diff --git a/inference/laya/public-smoke.mjs b/inference/laya/public-smoke.mjs deleted file mode 100644 index ed6a836..0000000 --- a/inference/laya/public-smoke.mjs +++ /dev/null @@ -1,73 +0,0 @@ -// Synthetic public-path checks; no key, customer payload, or account identifier. -import assert from "node:assert/strict"; -import { mkdir, writeFile } from "node:fs/promises"; -const endpoint = "https://classifier.dev/v1/classify"; -const base = { input: "Please refund the duplicate invoice payment.", labels: ["billing", "technical"] }; -const report = { checkedAt: new Date().toISOString(), endpoint, checks: {} }; - -async function call(body, retry = false) { - const deadline = Date.now() + 180_000; - let attempts = 0; - while (true) { - const started = performance.now(); - const response = await fetch(endpoint, { method: "POST", headers: { "content-type": "application/json" }, - body: JSON.stringify(body), signal: AbortSignal.timeout(30_000) }); - const data = await response.json(); - const ms = Math.round(performance.now() - started); - attempts++; - if (retry && [429, 503].includes(response.status) && Date.now() < deadline) { - const wait = Number(response.headers.get("retry-after")) || 5; - if (Date.now() + wait * 1000 >= deadline) throw new Error("retry deadline exceeded"); - await new Promise(resolve => setTimeout(resolve, wait * 1000)); - continue; - } - return { status: response.status, data, ms, attempts, lane: response.headers.get("x-classifier-processing"), limit: response.headers.get("ratelimit-limit") }; - } -} - -const jev = await call(base); -assert.equal(jev.status, 200); -assert.match(jev.data.model, /^jev/); -report.checks.jevDefault = { status: jev.status, model: jev.data.model }; -const timings = []; -const serverTimings = []; -for (let i = 0; i < 8; i++) { - const result = await call({ ...base, model: "laya", processing: "fast" }, true); - assert.equal(result.status, 200); - assert.equal(result.data.model, "laya-0.3.4-english-fast"); - assert.equal(result.lane, "fast"); - assert.equal(result.limit, "60"); - assert.equal(result.data.results.length, 1); - timings.push(result.ms); - serverTimings.push(result.data.usage.ms); -} -report.checks.fast = { status: 200, samples: 8, roundTripMs: timings, serverProcessingMs: serverTimings, scope: "Single client, sequential requests; not a global latency benchmark" }; -console.log("Fast public route passed"); -const invalid = await call({ ...base, model: "laya", processing: ["bulk"] }); -assert.equal(invalid.status, 400); -report.checks.invalidLane = invalid.status; -const tooMany = await call({ labels: base.labels, inputs: [base.input, base.input], model: "laya", processing: "fast" }); -assert.equal(tooMany.status, 400); -report.checks.fastBatchRejected = tooMany.status; -const ready = await call({ ...base, model: "laya", processing: "bulk" }, true); -assert.equal(ready.status, 200); -assert.equal(ready.data.model, "laya-0.3.4-english-bulk"); -report.checks.bulkStartupAttempts = ready.attempts; -const bulk = await call({ labels: base.labels, inputs: Array(128).fill(base.input), model: "laya", processing: "bulk" }, true); -assert.equal(bulk.status, 200); -assert.equal(bulk.data.results.length, 128); -assert.ok(bulk.data.results.every(row => row.model === "laya-0.3.4-english-bulk")); -report.checks.bulk = { status: bulk.status, rows: 128, roundTripMs: bulk.ms, serverProcessingMs: bulk.data.usage.ms, attempts: bulk.attempts }; -const languages = [base.input, "मुझसे दो बार शुल्क लिया गया है। कृपया मेरा पैसा वापस कर दें।", - "Me han cobrado dos veces en mi cuenta. Por favor, quiero que me devuelvan el dinero.", "The application crashes every time I log in."]; -const mixed = await call({ inputs: languages, labels: base.labels, model: "laya", processing: "bulk" }, true); -assert.equal(mixed.status, 200); -assert.equal(mixed.data.model, "mixed"); -assert.deepEqual(mixed.data.results.map(row => row.model), ["english", "multilingual", "multilingual", "english"].map(checkpoint => `laya-0.3.4-${checkpoint}-bulk`)); -const hindi = await call({ input: languages[1], labels: base.labels, model: "laya", processing: "fast" }, true); -assert.equal(hindi.status, 200); -assert.equal(hindi.data.model, "laya-0.3.4-multilingual-fast"); -report.checks.routing = { fastHindi: hindi.data.model, mixedBulkModels: mixed.data.results.map(row => row.model), orderPreserved: true }; -await mkdir(new URL("../../eval/data/", import.meta.url), { recursive: true }); -await writeFile(new URL("../../eval/data/laya-public-checks.json", import.meta.url), JSON.stringify(report, null, 2) + "\n"); -console.log(JSON.stringify(report, null, 2)); diff --git a/inference/laya/runtime.py b/inference/laya/runtime.py deleted file mode 100644 index a1d68b0..0000000 --- a/inference/laya/runtime.py +++ /dev/null @@ -1,136 +0,0 @@ -"""Bounded, in-memory serving. Overload is explicit, never an unbounded GPU queue.""" -import asyncio -import time - - -class Admission: - def __init__(self, lane): - self.capacity = 1 if lane == "fast" else 4 - self.pending = 0 - self.lock = asyncio.Lock() - self.tokens = 4.0 - self.updated = time.monotonic() - - def enter(self, lane, questions): - now = time.monotonic() - self.tokens = min(4.0, self.tokens + (now - self.updated) * 20) - self.updated = now - if self.pending >= self.capacity or (lane == "fast" and questions > self.tokens): - return False - if lane == "fast": - self.tokens -= questions - self.pending += 1 - return True - - def leave(self): - self.pending -= 1 - - -def make_api(lane, predict): - from fastapi import FastAPI, HTTPException, Request - api = FastAPI() - admission = Admission(lane) - limit = 4 if lane == "fast" else 64 - - @api.get("/health") - async def health(): - return {"ok": True, "lane": lane, "max_questions": limit, "pending": admission.pending} - - @api.post("/predict") - async def classify(request: Request): - raw = bytearray() - async for chunk in request.stream(): - raw.extend(chunk) - if len(raw) > 256_000: - raise HTTPException(413, "request too large") - import json - try: - body = json.loads(raw) - except (ValueError, UnicodeDecodeError): - raise HTTPException(400, "invalid JSON") - rows = body.get("batch") if isinstance(body, dict) else None - if not isinstance(rows, list) or not 1 <= len(rows) <= (1 if lane == "fast" else 64): - raise HTTPException(400, "invalid batch size") - count = 0 - for row in rows: - if not isinstance(row, dict) or not isinstance(row.get("state"), str) or not row["state"].strip(): - raise HTTPException(400, "state must be non-empty text") - questions = row.get("questions") - if not isinstance(questions, dict) or not questions: - raise HTTPException(400, "questions are required") - for q in questions.values(): - if not isinstance(q, dict) or q.get("type") not in ("choice", "noul") or not isinstance(q.get("instructions"), str): - raise HTTPException(400, "invalid question") - if q["type"] == "choice" and (not isinstance(q.get("criteria"), dict) or not 2 <= len(q["criteria"]) <= 16): - raise HTTPException(400, "choice requires 2–16 labels") - count += len(questions) - if count > limit: - raise HTTPException(400, "too many questions") - if not admission.enter(lane, count): - raise HTTPException(429, "lane busy; retry later", headers={"Retry-After": "1"}) - started = time.monotonic() - try: - try: - await asyncio.wait_for(admission.lock.acquire(), timeout=5) - except TimeoutError: - raise HTTPException(429, "lane busy; retry later", headers={"Retry-After": "1"}) - try: - task = asyncio.create_task(asyncio.to_thread(predict, rows)) - try: - results = await asyncio.shield(task) - except asyncio.CancelledError: - # Do not release the GPU lock while its thread still runs. - await task - raise - finally: - admission.lock.release() - return {"results": results, "lane": lane, "inference_ms": round((time.monotonic() - started) * 1000, 2)} - except (ValueError, TypeError, KeyError): - raise HTTPException(400, "input exceeds Laya context or has an invalid question") - finally: - admission.leave() - - return api - - -def load_router(): - """Preload all pinned checkpoints offline, retaining the SDK's public route IDs.""" - import laya - import os - from huggingface_hub import snapshot_download - checkpoint = snapshot_download("convaiinnovations/laya", revision=os.environ["LAYA_REVISION"], local_files_only=True) - router = laya.Router(device="cuda", max_loaded=3) - for name, subfolder in (("english", None), ("multilingual", "multilingual"), ("typed-decisions", "typed-decisions")): - agent = laya.load(checkpoint, subfolder=subfolder, device="cuda") - if agent.device.type != "cuda": - raise RuntimeError("Laya GPU initialization failed; refusing CPU fallback") - router.attach(name, agent) - router.preload() # Already attached: no downloads, rebuilding, or eviction. - return router - - -def start_server(lane): - import threading - import torch - import uvicorn - from adapter import predict_routed_batch - router = load_router() - - def predict(rows): - with torch.inference_mode(): - return predict_routed_batch(router, rows) - - # Force tokenizer, CUDA kernels, and allocator setup before Modal marks the - # container ready. Otherwise the first real fast-lane request pays seconds - # of one-time GPU initialization despite the container being "warm". - warmup_question = {"q": {"type": "choice", "instructions": "Choose one.", - "criteria": {"yes": None, "no": None}}} - predict([ - {"state": "warmup", "questions": warmup_question}, - {"state": "तैयार", "questions": warmup_question}, - ]) - - server = uvicorn.Server(uvicorn.Config(make_api(lane, predict), host="0.0.0.0", port=8000, - log_level="critical", access_log=False)) - threading.Thread(target=server.run, daemon=True).start() - return server diff --git a/inference/laya/sdk-check.json b/inference/laya/sdk-check.json deleted file mode 100644 index 54a00bc..0000000 --- a/inference/laya/sdk-check.json +++ /dev/null @@ -1,23 +0,0 @@ -{ - "rows": 4, - "questions": 8, - "device": "cuda", - "routes": [ - "english", - "multilingual", - "multilingual", - "english" - ], - "resident_checkpoints": [ - "english", - "multilingual", - "typed-decisions" - ], - "same_routing_and_order": true, - "same_token_counts": true, - "no_model_reloads": true, - "same_choices": true, - "typed_checkpoint_checked": true, - "max_score_delta": 0.0012, - "tolerance": 0.01 -} diff --git a/inference/laya/smoke.py b/inference/laya/smoke.py deleted file mode 100644 index 018d2b4..0000000 --- a/inference/laya/smoke.py +++ /dev/null @@ -1,69 +0,0 @@ -"""Live isolated lane checks. Credential JSON path is supplied, never printed.""" -import asyncio -from collections import Counter -import json -from pathlib import Path -import sys -import time -import httpx - -ROW = {"state": "Please refund the duplicate invoice payment.", "questions": { - "q": {"type": "choice", "instructions": "Which department should handle this?", "criteria": {"billing": None, "technical": None}}}} - - -async def main(): - keys = json.loads(Path(sys.argv[1]).read_text()) - headers = {"Modal-Key": keys["LAYA_MODAL_KEY"], "Modal-Secret": keys["LAYA_MODAL_SECRET"]} - report = {} - async with httpx.AsyncClient(headers=headers, timeout=30) as client: - for lane in ("fast", "bulk"): - url = f"https://miryaboy--classifier-laya-router-trial-west-{lane}.us-west.modal.direct" - started = time.monotonic() - while time.monotonic() - started < 240: - try: - ready = await client.get(url + "/health") - if ready.status_code == 200: - break - assert ready.status_code == 503, ready.status_code - except httpx.TransportError: - pass - await asyncio.sleep(2) - else: - raise TimeoutError(lane + " not ready") - count = 1 if lane == "fast" else 64 - response = await client.post(url + "/predict", json={"batch": [ROW] * count}) - assert response.status_code == 200, (lane, response.status_code, response.text[:80]) - rows = response.json()["results"] - assert len(rows) == count and all(r["answers"]["q"]["choice"] == "billing" for r in rows) - assert all(r["routing"]["model"] == "english" for r in rows) - await asyncio.sleep(.3) - hindi = await client.post(url + "/predict", json={"batch": [{**ROW, "state": "मुझसे दो बार शुल्क लिया गया है। कृपया मेरा पैसा वापस कर दें।"}]}) - assert hindi.status_code == 200 - assert hindi.json()["results"][0]["routing"]["model"] == "multilingual" - invalid = await client.post(url + "/predict", json={"batch": []}) - assert invalid.status_code == 400 - await asyncio.sleep(.3) - long = await client.post(url + "/predict", json={"batch": [{**ROW, "state": "text " * 3000}]}) - assert long.status_code == 400, (lane, long.status_code) - await asyncio.sleep(.3) - started = time.monotonic() - burst = await asyncio.gather(*(client.post(url + "/predict", json={"batch": [ROW] * count}) for _ in range(24))) - statuses = dict(Counter(str(r.status_code) for r in burst)) - assert all(r.status_code in (200, 429) for r in burst), statuses - assert statuses.get("200", 0) > 0 and statuses.get("429", 0) > 0, statuses - elapsed = time.monotonic() - started - await asyncio.sleep(.3) - recovery = await client.post(url + "/predict", json={"batch": [ROW]}) - assert recovery.status_code == 200 - async with httpx.AsyncClient() as anonymous: - assert (await anonymous.get(url + "/health")).status_code == 401 - report[lane] = {"valid_rows": count, "invalid": invalid.status_code, - "oversized_context": long.status_code, "burst": statuses, - "burst_elapsed_seconds": round(elapsed, 2), "recovery": recovery.status_code, - "unauthenticated": 401, "english_and_hindi_routes": "passed"} - print(json.dumps({lane: report[lane]}), flush=True) - Path("inference/laya/live-checks.json").write_text(json.dumps(report, indent=2) + "\n") - - -if __name__ == "__main__": - asyncio.run(main()) diff --git a/inference/laya/test_runtime.py b/inference/laya/test_runtime.py deleted file mode 100644 index 393ea15..0000000 --- a/inference/laya/test_runtime.py +++ /dev/null @@ -1,100 +0,0 @@ -import asyncio -import unittest -import threading -import contextlib -import types -from unittest.mock import patch -import httpx -from runtime import make_api, Admission, start_server -from adapter import predict_routed_batch - -ROW = {"state": "refund please", "questions": {"q": {"type": "choice", "instructions": "Choose", "criteria": {"billing": None, "tech": None}}}} - - -class Tests(unittest.IsolatedAsyncioTestCase): - def test_checkpoints_are_warmed_before_server_readiness(self): - events = [] - def predict(router, rows): - events.append([row["state"] for row in rows]) - return [] - server = types.SimpleNamespace(run=lambda: None) - thread = types.SimpleNamespace(start=lambda: events.append("listening")) - with patch("runtime.load_router", return_value=object()), \ - patch("adapter.predict_routed_batch", side_effect=predict), \ - patch("threading.Thread", return_value=thread), \ - patch.dict("sys.modules", { - "torch": types.SimpleNamespace(inference_mode=contextlib.nullcontext), - "uvicorn": types.SimpleNamespace(Config=lambda *a, **k: None, Server=lambda config: server), - }): - self.assertIs(start_server("fast"), server) - self.assertEqual(events, [["warmup", "तैयार"], "listening"]) - - async def test_validation(self): - calls = [] - def predict(rows): - calls.append(rows) - return [{"answers": {}} for _ in rows] - async with httpx.AsyncClient(transport=httpx.ASGITransport(app=make_api("fast", predict)), base_url="http://test") as client: - for body in ({"batch": []}, {"batch": [ROW, ROW]}, {"batch": [{"state": "text"}]}): - self.assertEqual((await client.post("/predict", json=body)).status_code, 400) - self.assertEqual((await client.post("/predict", content="{" )).status_code, 400) - self.assertEqual((await client.post("/predict", content="x" * 256001)).status_code, 413) - self.assertEqual(len(calls), 0) - - async def test_fast_overload_does_not_block_health_and_recovers(self): - started, release = threading.Event(), threading.Event() - def predict(rows): - started.set() - release.wait(3) - return [{"answers": {}}] - async with httpx.AsyncClient(transport=httpx.ASGITransport(app=make_api("fast", predict)), base_url="http://test") as client: - first = asyncio.create_task(client.post("/predict", json={"batch": [ROW]})) - await asyncio.to_thread(started.wait, 2) - try: - self.assertEqual((await client.get("/health")).status_code, 200) - response = await client.post("/predict", json={"batch": [ROW]}) - self.assertEqual(response.status_code, 429) - self.assertEqual(response.headers["retry-after"], "1") - finally: - release.set() - self.assertEqual((await first).status_code, 200) - self.assertEqual((await client.post("/predict", json={"batch": [ROW]})).status_code, 200) - - def test_admission_is_bounded(self): - gate = Admission("bulk") - self.assertTrue(all(gate.enter("bulk", 64) for _ in range(4))) - self.assertFalse(gate.enter("bulk", 1)) - gate.leave() - self.assertTrue(gate.enter("bulk", 64)) - - def test_routed_batch_groups_checkpoints_and_restores_order(self): - class Router: - loaded = ["english", "multilingual", "typed-decisions"] - def route(self, state, questions): - return {"model": "multilingual" if state.startswith("ml") else "english", "reason": "fixture"} - def load(self, name): - return name - calls = [] - def predict(agent, rows): - calls.append((agent, [row["state"] for row in rows])) - return [{"answers": {"state": row["state"]}, "usage": {"input_tokens": len(row["state"])}} for row in rows] - states = ["en-first", "ml-second", "en-third", "ml-fourth"] - actual = predict_routed_batch(Router(), [{**ROW, "state": state} for state in states], predict) - self.assertEqual(calls, [("english", ["en-first", "en-third"]), ("multilingual", ["ml-second", "ml-fourth"])]) - self.assertEqual([row["answers"]["state"] for row in actual], states) - self.assertEqual([row["routing"]["model"] for row in actual], ["english", "multilingual", "english", "multilingual"]) - self.assertEqual([row["usage"]["input_tokens"] for row in actual], list(map(len, states))) - - def test_missing_resident_checkpoint_never_loads_on_request(self): - class Router: - loaded = ["english"] - def route(self, state, questions): - return {"model": "multilingual", "reason": "fixture"} - def load(self, name): - self.fail("must not load") - with self.assertRaisesRegex(RuntimeError, "not preloaded"): - predict_routed_batch(Router(), [ROW]) - - -if __name__ == "__main__": - unittest.main() diff --git a/inference/laya/verify_sdk.py b/inference/laya/verify_sdk.py deleted file mode 100644 index 06f1fed..0000000 --- a/inference/laya/verify_sdk.py +++ /dev/null @@ -1,58 +0,0 @@ -"""One-off GPU equivalence check: uvx --from modal modal run inference/laya/verify_sdk.py.""" -import modal -from deploy import image - -app = modal.App("classifier-laya-sdk-check") - - -@app.function(image=image, gpu="L4", cpu=2, memory=4096, max_containers=1, timeout=480) -def verify(): - from adapter import predict_routed_batch, predict_batch - from runtime import load_router - router = load_router() - questions = { - "department": {"type": "choice", "instructions": "Which department should handle this?", - "criteria": {"billing": "invoices, payments, refunds", "technical": "bugs, outages"}}, - "refund": {"type": "noul", "instructions": "Does the user explicitly request a refund?"}, - } - states = ["Please refund the duplicate invoice payment.", - "मुझसे दो बार शुल्क लिया गया है। कृपया मेरा पैसा वापस कर दें।", - "Me han cobrado dos veces en mi cuenta. Por favor, quiero que me devuelvan el dinero.", - "The application crashes every time I log in."] - resident = {name: id(router.load(name)) for name in router.loaded} - expected = [router.predict(state, questions) for state in states] - actual = predict_routed_batch(router, [{"state": state, "questions": questions} for state in states]) - assert [row["routing"]["model"] for row in actual] == ["english", "multilingual", "multilingual", "english"] - assert all(id(router.load(name)) == identity for name, identity in resident.items()) - # Typed-decisions is resident but never automatically substituted by the SDK default. - typed_expected = router.predict(states[0], questions, model="typed-decisions") - typed_actual = predict_batch(router.load("typed-decisions"), [{"state": states[0], "questions": questions}])[0] - max_delta = 0.0 - for official, adapted in zip(expected, actual): - assert official["routing"] == adapted["routing"] - assert official["usage"] == adapted["usage"] - for official, adapted in list(zip(expected, actual)) + [(typed_expected, typed_actual)]: - for key in questions: - a, b = official["answers"][key], adapted["answers"][key] - assert a.get("choice") == b.get("choice") - for field in ("confidence", "noul"): - if field in a: - max_delta = max(max_delta, abs(a[field] - b[field])) - for label, value in a.get("probabilities", {}).items(): - max_delta = max(max_delta, abs(value - b["probabilities"][label])) - max_delta = max(max_delta, abs(a["action"]["act_probability"] - b["action"]["act_probability"])) - assert max_delta <= .01, max_delta # BF16 padding/batch shape can change rounding. - return {"rows": len(states), "questions": len(states) * len(questions), "device": "cuda", - "routes": [row["routing"]["model"] for row in actual], "resident_checkpoints": sorted(resident), - "same_routing_and_order": True, "same_token_counts": True, "no_model_reloads": True, - "same_choices": True, "typed_checkpoint_checked": True, - "max_score_delta": round(max_delta, 6), "tolerance": .01} - - -@app.local_entrypoint() -def main(): - import json - from pathlib import Path - result = verify.remote() - Path(__file__).with_name("sdk-check.json").write_text(json.dumps(result, indent=2) + "\n") - print(result) diff --git a/src/cost.ts b/src/cost.ts index 640eba6..53099e9 100644 --- a/src/cost.ts +++ b/src/cost.ts @@ -16,6 +16,8 @@ /** TypeSafe bills Jev per input token: $0.042 per million. */ export const JEV_USD_PER_MTOK = 0.042; +/** Beam bills every hosted decision model per input token: $0.021 per million, output free. */ +export const BEAM_USD_PER_MTOK = 0.021; /** Versioned Jev model whose context window and retail rate are pinned for account billing. */ export const JEV_ACCOUNT_MODEL = "jev-1.13.0"; @@ -28,7 +30,7 @@ export type TokenCounts = { }; export type ModelTokenUsage = TokenCounts & { - provider: "typesafe" | "vercel" | "openrouter" | "modal"; + provider: "typesafe" | "vercel" | "openrouter" | "beam"; model: string; calls: number; }; @@ -85,6 +87,12 @@ export function addJevCost(meter: Meter | undefined, inputTokens: unknown) { if (meter && Number.isFinite(n) && n > 0) meter.usd += (n * JEV_USD_PER_MTOK) / 1e6; } +/** Beam prices its decision models on input tokens alone; output is always zero. */ +export function addBeamCost(meter: Meter | undefined, inputTokens: unknown) { + const n = Number(inputTokens); + if (meter && Number.isFinite(n) && n > 0) meter.usd += (n * BEAM_USD_PER_MTOK) / 1e6; +} + /** OpenRouter reports the charge for the call directly, already in USD. */ export function addUsd(meter: Meter | undefined, usd: unknown) { const n = Number(usd); diff --git a/src/docs.ts b/src/docs.ts index 21014e5..63bad5e 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -68,17 +68,23 @@ TYPESAFE SDK COMPATIBILITY The corresponding HTTP resources are POST /v1/systemone and GET /v1/models. -LAYA TRIAL - - Calls with neither model nor processing use Jev. To try Laya, send a POST - with model: "laya". Omit processing to automatically choose fast for one - decision (up to four multi-label questions), or bulk for larger work. - Supplying processing without model implies Laya. Explicit lanes are honored. - With explicit model: "jev", processing is accepted but has no effect. - Both lanes use the same preloaded Router: - English text uses the English checkpoint; other languages use multilingual. - Processing changes batching and capacity, not checkpoint selection. - Each result's model identifies the checkpoint and lane that answered. +LAYA AND KEV + + Calls with neither model nor processing use Jev. Two other models are + available, both hosted by Beam on shared inference endpoints: + + model: "laya" ModernBERT-large with a trained decision head. Answers + every question in one batched forward pass. State plus one + question must fit 512 tokens. + model: "kev" Qwen2.5-0.5B with released LoRA weights and a pointer + head. Encodes state once and scores question branches + together in one 8,192-token packed sequence. + + Omit processing to automatically choose fast for one decision (up to four + multi-label questions), or bulk for larger work. Supplying processing + without model implies Laya. Explicit lanes are honored. With explicit + model: "jev", processing is accepted but has no effect. Each result is + labelled jev/laya or jev/kev; neither model reports a checkpoint. {"model":"laya","processing":"fast","input":"Please refund this charge", "labels":["billing","technical"]} @@ -89,34 +95,34 @@ LAYA TRIAL decision is one question; multi-label uses one question per label. Fast accepts at most four questions. Multiple dimensions normally need bulk. Trial limits also apply to paid and operator keys; existing Smart quotas - still apply. Shared GPU capacity can return 429 even with quota remaining. - Quotas count attempted questions, including failed inference and retries. + still apply. Shared capacity can return 429 even with quota remaining. + Quotas count attempted questions, including failed inference. - Laya accepts short text: at most 2,000 characters, 2–16 labels of + Both models accept short text: at most 2,000 characters, 2–16 labels of at most 100 characters each, and instructions up to 400 characters. The - text, question and labels must also fit the selected checkpoint's context - (512 tokens for English, 1,024 for multilingual); - oversized content is rejected, not silently shortened. Jev's published - accuracy and calibration measurements do not describe Laya. - The upstream model card warns that shipped confidence can be overconfident. - We have not fitted temperatures on classifier.dev traffic: treat scores and - Smart's confidence-triggered reviews as experimental, not a quality guarantee. - Language routing is a heuristic, not a guarantee of language or task accuracy. - - 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. - One regional GPU deployment is not a replica in every region or a latency promise. + text, question and labels must also fit the model's context — 512 tokens + for Laya, 8,192 for Kev — and oversized content is rejected, not silently + shortened. Upstream caps a request at 32 questions, so a large batch is + split into several requests and reassembled in input order. + + Jev's published accuracy and calibration measurements do not describe + either model. We have not fitted temperatures on classifier.dev traffic: + treat scores and Smart's confidence-triggered reviews as experimental, + not a quality guarantee. + + Neither model has a warm pool to start, so there is no cold-start 503. + On 429, respect Retry-After and use bounded retries with backoff. A refused + request is not retried upstream: the refusal is returned as it happened. + Shared endpoints are 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. - Laya inference has no retail charge during this trial. Optional tier: "smart" - reviews remain separately priced as before and can be slower. GPU hosting - costs are separate from the API's provider-reported spend totals. + Laya and Kev inference has no retail charge during this trial. Optional + tier: "smart" reviews remain separately priced as before and can be slower. From this repository's CLI: node cli/classify.js billing,technical --model laya "Please refund this" - node cli/classify.js billing,technical --model laya --processing bulk < tickets.txt + node cli/classify.js billing,technical --model kev --processing bulk < tickets.txt AGAINST THE MODEL IT RUNS ON diff --git a/src/http/classification.ts b/src/http/classification.ts index 1958b93..b932590 100644 --- a/src/http/classification.ts +++ b/src/http/classification.ts @@ -69,7 +69,7 @@ export async function accountClassification(request: Request, env: AppEnv & Part const analytics = (success: boolean, retailCostUsd: number | null) => writeAccountAnalytics(env, { accountId, keyId: reservation.agentId, requestId: reservation.id, source, tier, status: success ? "success" : "error", items: itemCount, ...tokens(), - model: meter.tokens.map((row) => row.model).join(","), providerCostUsd: meter.tokens.length && !meter.tokens.some(row => row.provider === "modal") ? meter.usd : null, + model: meter.tokens.map((row) => row.model).join(","), providerCostUsd: meter.tokens.length ? meter.usd : null, retailCostUsd, latencyMs: Date.now() - started, escalations: meter.tokens.filter((row) => row.provider === "openrouter").reduce((total, row) => total + row.calls, 0), }); diff --git a/src/index.ts b/src/index.ts index 8376615..f6dd63a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -19,7 +19,7 @@ import * as feedback from "./feedback"; import * as newsletter from "./newsletter"; import * as skills from "./skills"; import { jevClassify, jevKeys, MULTI_THRESHOLD } from "./jev"; -import { LAYA_LIMITS, LayaError, planLaya, runLaya, limitLaya, layaModel, readQuotaTiming, type QuotaTiming, type LayaEnv, type LayaPlan, type LayaTiming, type Processing } from "./laya"; +import { LAYA_LIMITS, LayaError, planLaya, runLaya, limitLaya, layaModel, readQuotaTiming, type QuotaTiming, type LayaEnv, type LayaPlan, type LayaTiming, type Processing, type LayaModel } from "./laya"; import { readDimensions, packDimensions, classifyDimensions, dimensionInstructions, MAX_DECISIONS, type Dimension, type DimensionBatch } from "./dimensions"; import { newMeter, addUsd, addTokens, type Meter } from "./cost"; import { secretEquals } from "./secrets"; @@ -37,6 +37,8 @@ export interface Env extends LayaEnv { QUOTA_COORDINATOR_ENABLED?: string; OPENROUTER_API_KEY: string; TYPESAFE_API_KEY?: string; + /** Beam workspace token; the credential for the Beam-hosted Laya and Kev. */ + BEAM_API_KEY?: string; /** * Vercel AI Gateway, which serves Jev on a free monthly credit. With it set, * Jev is asked there first and TYPESAFE_API_KEY catches what the gateway @@ -934,7 +936,9 @@ 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(() => { + const keys = jevKeys(env); + if (!keys?.beam) return Promise.reject(new LayaError("Laya trial is currently unavailable", 503)); + return runLaya(env, keys, plan, meter, timing, signal).finally(() => { if (timing) timing.runMs = performance.now() - started; }); } @@ -1826,7 +1830,7 @@ const worker = { let inputs: string[] = []; let labels: string[] = []; let tier: Tier = "fast"; - let selectedModel: "jev" | "laya" = "jev"; + let selectedModel: "jev" | LayaModel = "jev"; let processing: Processing = "fast"; let automaticProcessing = true; let layaPlan: LayaPlan | undefined; @@ -1849,7 +1853,7 @@ const worker = { // key if the caller sent one — classification has no side effects, so the // echo is all a retrying client needs. const apiHeaders = (remaining = -1): Record => { - const isLaya = selectedModel === "laya"; + const isLaya = selectedModel !== "jev"; const limit = isLaya ? Math.min(LAYA_LIMITS[processing].rpm, TIERS[tier].rpm * multiplier) : TIERS[tier].rpm * multiplier; const daily = isLaya ? Math.min(LAYA_LIMITS[processing].daily, TIERS[tier].daily * multiplier) : TIERS[tier].daily * multiplier; const h: Record = { @@ -1867,7 +1871,7 @@ const worker = { return h; }; const fail = (msg: string, status: number, reason: ErrorCode, extra: Record = {}, ms = 0, remaining = -1, more: Record = {}) => { - record(env, ctx, { tier, n: 0, ms, labels, ip, country, status, client, model: selectedModel === "laya" ? layaModel(processing) : "", + record(env, ctx, { tier, n: 0, ms, labels, ip, country, status, client, model: selectedModel === "jev" ? "" : layaModel(selectedModel), usd: meter.usd, reason, agent, attempted: inputs.length, escalationFailed: 0, mode, dimensions: dimensions?.length ?? 0 }); const headers = { ...apiHeaders(remaining), ...extra }; // A GET that was malformed gets back a URL that would have worked, in the spelling it used. @@ -1895,8 +1899,13 @@ const worker = { return fail('Body must be a JSON object such as {"input":"...","labels":["a","b"]}. See https://classifier.dev', 400, "bad_json"); } const b = body as Record; - if (b.model !== undefined && b.model !== "jev" && b.model !== "laya") return fail('model must be "jev" or "laya"', 400, "bad_model"); - selectedModel = b.model === "laya" || (b.model === undefined && b.processing !== undefined) ? "laya" : "jev"; + if (b.model !== undefined && b.model !== "jev" && b.model !== "laya" && b.model !== "kev") + return fail('model must be "jev", "laya" or "kev"', 400, "bad_model"); + // `processing` without a model still means Laya: it is the lane selector + // for the Beam-hosted trial and Jev has no lanes. + selectedModel = b.model === "laya" || b.model === "kev" + ? (b.model as LayaModel) + : b.model === undefined && b.processing !== undefined ? "laya" : "jev"; if (b.processing !== undefined && b.processing !== "fast" && b.processing !== "bulk") return fail('processing must be "fast" or "bulk"', 400, "bad_processing"); processing = b.processing === "bulk" ? "bulk" : "fast"; @@ -1975,20 +1984,20 @@ const worker = { if (dimensions) { if (inputs.length * dimensions.length > MAX_DECISIONS) return fail(`Maximum ${MAX_DECISIONS} decisions (items × dimensions) per request`, 400, "too_many_decisions"); - try { if (selectedModel !== "laya") dimensionBatches = packDimensions(inputs, dimensions, instructions); } + try { if (selectedModel === "jev") dimensionBatches = packDimensions(inputs, dimensions, instructions); } catch (e) { return fail((e as Error).message, 400, "dimension_context_too_large"); } } else mode = multi ? "multi" : "single"; const decisions = inputs.length * (dimensions?.length ?? 1); - if (selectedModel === "laya" && automaticProcessing) { + if (selectedModel !== "jev" && automaticProcessing) { processing = decisions > 1 || (!!multi && labels.length > LAYA_LIMITS.fast.questions) ? "bulk" : "fast"; } - if (selectedModel === "laya") { + if (selectedModel !== "jev") { if (!account && !req.headers.has("authorization")) layaTiming = {}; if (env.LAYA_ENABLED !== "true") return fail("Laya trial is currently unavailable", 503, "laya_unavailable"); try { layaPlan = planLaya(dimensions ? inputs.flatMap(input => dimensions!.map(d => ({ input, labels: d.labels, instructions: dimensionInstructions(d, instructions) }))) - : inputs.map(input => ({ input, labels, instructions, multi: !!multi })), processing); + : inputs.map(input => ({ input, labels, instructions, multi: !!multi })), processing, selectedModel as LayaModel); } catch (error) { if (error instanceof LayaError) return fail(error.message, error.status, error.status === 400 ? "laya_input" : error.scope === "day" ? "rate_limit_day" : error.status === 429 ? "laya_rate_limit" : "laya_unavailable", @@ -2128,8 +2137,7 @@ const worker = { quota_combined: combinedQuota ? regularQuotaMs : undefined, regular_handler: regularQuotaTiming.handler, regular_read: regularQuotaTiming.read, regular_write: regularQuotaTiming.write, lane_handler: laneQuotaTiming.handler, lane_read: laneQuotaTiming.read, lane_write: laneQuotaTiming.write, - laya_run: layaTiming.runMs, modal_fetch: layaTiming.fetchMs, - modal_headers: layaTiming.headersMs, backend: layaTiming.backendMs, + laya_run: layaTiming.runMs, beam_fetch: layaTiming.fetchMs, }; headers["server-timing"] = Object.entries(durations) .filter((entry): entry is [string, number] => typeof entry[1] === "number" && Number.isFinite(entry[1]) && entry[1] >= 0) diff --git a/src/jev-observability.ts b/src/jev-observability.ts index 9246b41..8f99b80 100644 --- a/src/jev-observability.ts +++ b/src/jev-observability.ts @@ -2,7 +2,7 @@ * must remain visible without inflating request counts. Never record payloads, * provider messages, keys, caller identifiers or label sets here. */ export type JevAttempt = { - provider: "gateway" | "typesafe"; + provider: "gateway" | "typesafe" | "beam"; outcome: "success" | "failure" | "skipped"; reason: string; status: number; diff --git a/src/jev.ts b/src/jev.ts index 18af37c..03c3876 100644 --- a/src/jev.ts +++ b/src/jev.ts @@ -23,7 +23,7 @@ */ import { recordJevAttempt } from "./jev-observability"; -import { addJevCost, addUsd, addTokens, JEV_ACCOUNT_MODEL, type Meter } from "./cost"; +import { addJevCost, addBeamCost, addUsd, addTokens, JEV_ACCOUNT_MODEL, type Meter } from "./cost"; const API = "https://api.typesafe.ai/v1/systemone"; const MODEL = "jev-latest"; @@ -73,12 +73,12 @@ export function resetGatewayPause() { } /** Where Jev can be asked. Either key alone works; with both, the gateway goes first and TypeSafe catches what it drops. */ -export type JevKeys = { typesafe?: string; gateway?: string; analytics?: AnalyticsEngineDataset }; +export type JevKeys = { typesafe?: string; gateway?: string; beam?: string; analytics?: AnalyticsEngineDataset }; -export const jevKeys = (env: { TYPESAFE_API_KEY?: string; AI_GATEWAY_API_KEY?: string; AI_GATEWAY_DISABLED?: string; JEV_AE?: AnalyticsEngineDataset }): JevKeys | null => { +export const jevKeys = (env: { TYPESAFE_API_KEY?: string; AI_GATEWAY_API_KEY?: string; AI_GATEWAY_DISABLED?: string; BEAM_API_KEY?: string; JEV_AE?: AnalyticsEngineDataset }): JevKeys | null => { const gateway = env.AI_GATEWAY_DISABLED === "true" ? undefined : env.AI_GATEWAY_API_KEY; - return env.TYPESAFE_API_KEY || gateway - ? { typesafe: env.TYPESAFE_API_KEY, gateway, ...(env.JEV_AE ? { analytics: env.JEV_AE } : {}) } + return env.TYPESAFE_API_KEY || gateway || env.BEAM_API_KEY + ? { typesafe: env.TYPESAFE_API_KEY, gateway, beam: env.BEAM_API_KEY, ...(env.JEV_AE ? { analytics: env.JEV_AE } : {}) } : null; }; @@ -98,6 +98,70 @@ const MAX_ITEMS = 1000; const CONCURRENCY = 8; const UPSTREAM_TIMEOUT_MS = 10_000; +/** + * Beam serves the same System One protocol as TypeSafe, so Laya and Kev are a + * third transport here rather than a second client: one packer, one retry + * policy, one validator, one meter. Only the arithmetic differs. + * + * Two Beam limits have no analogue at TypeSafe. Every model refuses more than + * 32 named questions in a request whatever its context ("questions must + * contain 1..32 named questions"), and the contexts are small — 512 tokens for + * Laya, 8,192 for Kev, against Jev's 64k. Beam rejects an overflow rather than + * truncating it, which is the behaviour we want, so the budgets below sit well + * under the documented ceilings: `estimateTokens` is an estimate, and a Beam + * context refusal is translated to `max_tokens_exceeded` so `runJevBatches` + * halves the batch and recovers instead of failing the request. + * + * Laya's item cap is empirical, not documented: 20 items of ordinary support + * text were accepted and 25 were refused, so 16 is the conservative floor. + */ +const BEAM_API = "https://app.beam.cloud/v1/systemone"; +/** Beam's hard per-request cap, identical across its models. */ +export const BEAM_MAX_QUESTIONS = 32; + +export type Limits = { tokenBudget: number; stateQuestionBudget: number; maxItems: number; maxQuestions: number }; +export type BackendId = "jev" | "laya" | "kev"; +export type Backend = { + id: BackendId; + via: "typesafe" | "beam"; + url: string; + /** Sent as `model`, and the label an answer is attributed to. */ + model: string; + /** Pinned, versioned id used when an account is being billed. */ + accountModel: string; + /** + * Attempts per request. Jev retries because a retried 429 is not billed and + * a second door exists. The Beam lanes do not: their quota counts attempts, + * including failed inference, and a model that is at capacity will not clear + * inside a 300ms backoff. A caller gets the refusal and its Retry-After. + */ + attempts: number; + limits: Limits; +}; + +export const JEV_BACKEND: Backend = { + id: "jev", via: "typesafe", url: API, model: MODEL, accountModel: JEV_ACCOUNT_MODEL, attempts: 3, + limits: { tokenBudget: TOKEN_BUDGET, stateQuestionBudget: STATE_QUESTION_BUDGET, maxItems: MAX_ITEMS, maxQuestions: Number.MAX_SAFE_INTEGER }, +}; +/** + * Laya's documented budget is "state plus one question must fit 512 tokens", + * which is `stateQuestionBudget`, not a total across questions: it answers + * every question in one batched forward pass over shared state. Capping the + * total as well would fragment a thousand short items into hundreds of + * requests for no reason, so only the documented constraint binds, under a + * 16-item ceiling taken from what Beam actually accepts (20 passed, 25 did not). + */ +export const LAYA_BACKEND: Backend = { + id: "laya", via: "beam", url: BEAM_API, model: "jev/laya", accountModel: "jev/laya", attempts: 1, + limits: { tokenBudget: Number.MAX_SAFE_INTEGER, stateQuestionBudget: 460, maxItems: 16, maxQuestions: 16 }, +}; +export const KEV_BACKEND: Backend = { + id: "kev", via: "beam", url: BEAM_API, model: "jev/kev", accountModel: "jev/kev", attempts: 1, + limits: { tokenBudget: 7_200, stateQuestionBudget: 7_200, maxItems: BEAM_MAX_QUESTIONS, maxQuestions: BEAM_MAX_QUESTIONS }, +}; +export const BACKENDS: Record = { jev: JEV_BACKEND, laya: LAYA_BACKEND, kev: KEV_BACKEND }; +export const isBackendId = (v: unknown): v is BackendId => v === "jev" || v === "laya" || v === "kev"; + /** * Tokens in a string, estimated. ASCII runs at about 3.5 characters a token; * everything else is counted at two tokens a character, which is what CJK @@ -225,33 +289,43 @@ function batchFrom(groups: JevQuestionGroup[]): JevBatch { return { groups, state: [...states.values()], questions }; } -/** Own both documented context budgets and greedily preserve group order. */ -export function prepareJevBatches(groups: JevQuestionGroup[], options: { rejectOversized?: boolean } = {}): JevBatch[] { +/** + * Own both documented context budgets and greedily preserve group order. + * + * `limits` selects the backend's arithmetic; it defaults to Jev's. Beam adds a + * hard ceiling on the number of named questions per request, which Jev has no + * equivalent of, so a batch flushes on question count as well as on tokens. + */ +export function prepareJevBatches(groups: JevQuestionGroup[], options: { rejectOversized?: boolean; limits?: Limits } = {}): JevBatch[] { + const { tokenBudget, stateQuestionBudget, maxItems, maxQuestions } = options.limits ?? JEV_BACKEND.limits; const batches: JevBatch[] = []; let current: JevQuestionGroup[] = []; let included = new Set(); let stateTokens = 0; let questionTokens = 0; + let questionCount = 0; let longestQuestion = 0; const flush = () => { if (current.length) batches.push(batchFrom(current)); current = []; included = new Set(); - stateTokens = questionTokens = longestQuestion = 0; + stateTokens = questionTokens = questionCount = longestQuestion = 0; }; for (const group of groups) { const addedState = included.has(group.state.id) ? 0 : stateCost(group.state); const costs = Object.values(group.questions).map(questionCost); const addedQuestions = costs.reduce((sum, cost) => sum + cost, 0); + const addedCount = costs.length; const groupLongest = Math.max(0, ...costs); - if (options.rejectOversized && stateCost(group.state) + groupLongest > STATE_QUESTION_BUDGET) { - throw new JevContextError("A state item and question exceed Jev's context budget"); + if (options.rejectOversized && stateCost(group.state) + groupLongest > stateQuestionBudget) { + throw new JevContextError("A state item and question exceed the model's context budget"); } - const overBudget = stateTokens + addedState + questionTokens + addedQuestions > TOKEN_BUDGET || - stateTokens + addedState + Math.max(longestQuestion, groupLongest) > STATE_QUESTION_BUDGET || - (!included.has(group.state.id) && included.size >= MAX_ITEMS); + const overBudget = stateTokens + addedState + questionTokens + addedQuestions > tokenBudget || + stateTokens + addedState + Math.max(longestQuestion, groupLongest) > stateQuestionBudget || + questionCount + addedCount > maxQuestions || + (!included.has(group.state.id) && included.size >= maxItems); if (current.length && overBudget) flush(); current.push(group); @@ -260,6 +334,7 @@ export function prepareJevBatches(groups: JevQuestionGroup[], options: { r stateTokens += stateCost(group.state); } questionTokens += addedQuestions; + questionCount += addedCount; longestQuestion = Math.max(longestQuestion, groupLongest); } flush(); @@ -384,12 +459,19 @@ async function postGateway(key: string, body: JevBody, meter?: Meter, analytics? } /** The gateway when it is configured and not paused, TypeSafe otherwise; and TypeSafe again when the gateway drops a request. */ -async function post(keys: JevKeys, body: JevBody, meter?: Meter): Promise { +async function post(keys: JevKeys, body: JevBody, meter?: Meter, backend: Backend = JEV_BACKEND, signal?: AbortSignal): Promise { + // Beam hosts its own models; there is no gateway door and nothing to fall + // back to, so an unconfigured key is an error rather than a silent reroute + // to Jev, which would answer with a different model than the caller asked for. + if (backend.via === "beam") { + if (!keys.beam) throw new JevError("beam: no key configured", 0, "unconfigured"); + return postBeam(keys.beam, { ...body, model: backend.model }, backend, meter, keys.analytics, signal); + } // Paid work has an explicit versioned price. Do not route it through the // unversioned gateway or silently follow a new model behind jev-latest. if (meter?.beforeCall) { if (!keys.typesafe) throw new JevError("typesafe: no key configured", 0, "unconfigured"); - return postTypesafe(keys.typesafe, { ...body, model: JEV_ACCOUNT_MODEL }, meter, keys.analytics); + return postTypesafe(keys.typesafe, { ...body, model: backend.accountModel }, meter, keys.analytics); } // Without a TypeSafe key there is nothing to pause towards, so the gateway is always tried. const tryGateway = keys.gateway && (!keys.typesafe || Date.now() >= gatewayPausedUntil); @@ -410,6 +492,82 @@ async function post(keys: JevKeys, body: JevBody, meter?: Meter): Promise { + let last: Error = new Error("beam: no attempt made"); + for (let attempt = 0; attempt < backend.attempts; attempt++) { + await meter?.beforeCall?.("beam", body.model, 0); + const started = Date.now(); + const observe = (status: number, reason = "") => recordJevAttempt(analytics, { provider: "beam", outcome: reason ? "failure" : "success", reason, status, ms: Date.now() - started, items: body.state.length, attempt: attempt + 1 }); + let res: Response; + try { + res = await fetch(backend.url, { + method: "POST", + headers: { authorization: `Bearer ${key}`, "content-type": "application/json" }, + body: JSON.stringify(body), + signal: signal ? AbortSignal.any([signal, AbortSignal.timeout(UPSTREAM_TIMEOUT_MS)]) : AbortSignal.timeout(UPSTREAM_TIMEOUT_MS), + }); + } catch (e) { + const timeout = e instanceof Error && (e.name === "AbortError" || e.name === "TimeoutError"); + observe(timeout ? 504 : 0, timeout ? "timeout" : "network"); + last = new JevError(`beam ${timeout ? "timeout" : "network failure"}`, timeout ? 504 : 0, timeout ? "timeout" : "network"); + // A caller that has withdrawn the request gets no further attempts. + if (signal?.aborted || attempt >= backend.attempts - 1) break; + await new Promise((r) => setTimeout(r, 300 * 2 ** attempt + Math.random() * 200)); + continue; + } + const rawPayload = await res.json().catch(() => null); + const payload = (isRecord(rawPayload) ? rawPayload : {}) as Partial & { error?: unknown; detail?: unknown }; + if (res.ok && !payload.error && validPayload(payload, body)) { + addBeamCost(meter, payload.usage?.input_tokens); + addTokens(meter, "beam", payload.model, { inputTokens: payload.usage?.input_tokens, outputTokens: payload.usage?.output_tokens, cachedInputTokens: 0 }); + observe(res.status); + return payload; + } + observe(res.status, beamFailureReason(res.status, rawPayload)); + last = new JevError(`beam ${res.status}: ${res.ok ? "malformed response" : beamErrorType(res.status, rawPayload)}`, res.status, beamErrorType(res.status, rawPayload)); + const retryable = res.status === 408 || res.status === 425 || res.status === 429 || res.status === 529 || res.status >= 500 || res.ok; + if (!retryable || attempt >= backend.attempts - 1) break; + await new Promise((r) => setTimeout(r, 300 * 2 ** attempt + Math.random() * 200)); + } + throw last; +} + +/** + * Beam's 422 covers both "too many questions" and "this did not fit the + * context". Only the second is recoverable by halving, and the two are told + * apart by the message rather than the status, so the match is deliberately + * narrow: anything unrecognised stays a plain invalid request. + */ +function beamErrorType(status: number, raw: unknown): string { + if (status === 0) return "network"; + const detail = isRecord(raw) && typeof raw.detail === "string" ? raw.detail : ""; + const error = isRecord(raw) && isRecord(raw.error) ? raw.error : {}; + const code = safeErrorType(error.code); + if (status === 422) { + return /\b(exceeds?|max_length|too long|context)\b/i.test(detail) && /\btokens?\b/i.test(detail) + ? "max_tokens_exceeded" + : "invalid_request"; + } + if (status === 429) return code || "rate_limit_exceeded"; + return code || (status >= 500 ? "upstream_error" : ""); +} + +function beamFailureReason(status: number, raw: unknown) { + const type = beamErrorType(status, raw); + if (["timeout", "network", "malformed_response", "max_tokens_exceeded"].includes(type)) return type; + if (status === 429) return "rate_limit"; + if ([401, 402, 403].includes(status)) return "credentials_or_credit"; + if (status === 422) return "invalid_request"; + return "upstream_error"; +} + async function postTypesafe(key: string, body: JevBody, meter?: Meter, analytics?: AnalyticsEngineDataset): Promise { let last: Error = new Error("typesafe: no attempt made"); for (let attempt = 0; attempt < 3; attempt++) { @@ -471,7 +629,7 @@ export async function jevAsk(keys: JevKeys, state: { id: string; text: string }[ export type JevBatchResult = { value: T; model: string; answers: Record }; /** Execute prepared requests with bounded concurrency and recover from an underestimated context by halving groups. */ -export async function runJevBatches(keys: JevKeys, prepared: JevBatch[], meter?: Meter): Promise[]> { +export async function runJevBatches(keys: JevKeys, prepared: JevBatch[], meter?: Meter, backend: Backend = JEV_BACKEND, signal?: AbortSignal): Promise[]> { const queue = [...prepared]; const positions = new Map(prepared.flatMap((batch) => batch.groups).map((group, index) => [group, index])); const out: JevBatchResult[] = new Array(positions.size); @@ -482,7 +640,7 @@ export async function runJevBatches(keys: JevKeys, prepared: JevBatch[], m while (!failure && next < queue.length) { const batch = queue[next++]; try { - const res = await post(keys, { state: batch.state, model: MODEL, questions: batch.questions }, meter); + const res = await post(keys, { state: batch.state, model: backend.model, questions: batch.questions }, meter, backend, signal); for (const group of batch.groups) out[positions.get(group)!] = { value: group.value, model: res.model, answers: res.answers }; } catch (error) { if (error instanceof JevError && error.errorType === "max_tokens_exceeded" && batch.groups.length > 1) { @@ -505,6 +663,7 @@ export async function jevClassify( instructions: string | undefined, multi: boolean, meter?: Meter, + backend: Backend = JEV_BACKEND, ): Promise { const groups = inputs.map((text, index): JevQuestionGroup => { const id = `i${index}`; @@ -513,7 +672,7 @@ export async function jevClassify( : undefined; return { state: { id, text, ...(rubric ? { rubric } : {}) }, questions: questionsFor(id, labels, instructions, multi), value: index }; }); - const answered = await runJevBatches(keys, prepareJevBatches(groups), meter); + const answered = await runJevBatches(keys, prepareJevBatches(groups, { limits: backend.limits }), meter, backend); return answered.map(({ value: index, model, answers }) => { const id = `i${index}`; if (multi) { diff --git a/src/laya.ts b/src/laya.ts index 2fd5da3..c8042df 100644 --- a/src/laya.ts +++ b/src/laya.ts @@ -1,57 +1,72 @@ -import { addTokens, type Meter } from "./cost"; -import type { JevResult } from "./jev"; +import type { Meter } from "./cost"; +import { + type JevKeys, type JevResult, type JevQuestionGroup, type Question, type Backend, + LAYA_BACKEND, KEV_BACKEND, JevError, prepareJevBatches, runJevBatches, +} from "./jev"; + +/** + * Laya and Kev, both hosted by Beam. + * + * Beam speaks TypeSafe's System One protocol, so there is no separate client + * here any more: this module owns the product contract — which lanes exist, + * what a caller may send, and what a lane costs against quota — and hands the + * actual request to `src/jev.ts`, which owns every System One transport. + * + * What the move off Modal changed, visibly: + * - There is no GPU pool to keep warm, so `fast` and `bulk` are quota lanes + * rather than separate deployments. Neither can return "still starting". + * - Beam does not report which checkpoint answered, so a result is labelled + * `jev/laya` or `jev/kev` rather than `laya-0.3.4--`. + * - Spend is per input token and reported, so account analytics no longer + * marks provider cost unknown. + */ export type Processing = "fast" | "bulk"; +export type LayaModel = "laya" | "kev"; + export type LayaEnv = { - LAYA_FAST_URL?: string; - LAYA_BULK_URL?: string; - LAYA_MODAL_KEY?: string; - LAYA_MODAL_SECRET?: string; LAYA_ENABLED?: string; LAYA_FAST_ADMISSION?: RateLimit; LIMITER: DurableObjectNamespace; }; + +/** + * Lane quotas. These are a product decision about shared capacity, not a + * property of Beam, so they survived the move unchanged: existing callers keep + * the limits they were issued against. + */ export const LAYA_LIMITS = { fast: { rpm: 60, daily: 2000, questions: 4 }, bulk: { rpm: 1000, daily: 20_000, questions: 64 }, } as const; -export const layaModel = (lane: Processing, checkpoint = "routed") => `laya-0.3.4-${checkpoint}-${lane}`; + +export const BACKEND_FOR: Record = { laya: LAYA_BACKEND, kev: KEV_BACKEND }; +/** The label a result carries; Beam returns the same string. */ +export const layaModel = (model: LayaModel = "laya") => BACKEND_FOR[model].model; + export class LayaError extends Error { constructor(message: string, readonly status: 400 | 429 | 502 | 503, readonly retryAfter = 1, readonly scope?: "day") { super(message); } } + export type LayaTask = { input: string; labels: string[]; instructions?: string; multi?: boolean }; -type Row = { state: string; questions: Record }; -export type LayaPlan = { batches: Row[][]; tasks: LayaTask[]; cost: number; processing: Processing }; +export type LayaPlan = { tasks: LayaTask[]; cost: number; processing: Processing; model: LayaModel }; -/** Validate and chunk before any quota, billing, or inference side effects. */ -export function planLaya(tasks: LayaTask[], processing: Processing): LayaPlan { +/** Validate and price before any quota, billing, or inference side effects. */ +export function planLaya(tasks: LayaTask[], processing: Processing, model: LayaModel = "laya"): LayaPlan { if (!tasks.length || tasks.length > 1000) throw new LayaError("Laya accepts 1–1,000 decisions per request", 400); if (processing === "fast" && tasks.length > 1) throw new LayaError("Laya fast accepts one decision per call; use processing: bulk for batches or dimensions", 400); - const batches: Row[][] = []; - let batch: Row[] = [], size = 0, bytes = 0, cost = 0; + let cost = 0; for (const task of tasks) { if (task.labels.length < 2 || task.labels.length > 16 || task.labels.some(l => l.length > 100) || task.input.length > 2000 || (task.instructions?.length ?? 0) > 400) { throw new LayaError("Laya trial accepts short text (≤2,000 characters), 2–16 short labels, and instructions ≤400 characters; the full question must also fit the selected checkpoint's token budget", 400); } - const questions: Record = {}; - if (task.multi) task.labels.forEach((label, i) => { - questions[`q${i}`] = { type: "noul", instructions: `Does this text belong to the category ${JSON.stringify(label)}? ${task.instructions ?? ""}` }; - }); - else questions.q = { type: "choice", instructions: `Which category fits this text? ${task.instructions ?? ""}`, - criteria: Object.fromEntries(task.labels.map(label => [label, null])) }; - const n = Object.keys(questions).length; - if (n > LAYA_LIMITS[processing].questions) throw new LayaError("Too many questions for Laya fast; use processing: bulk", 400); - const row = { state: task.input, questions }; - const rowBytes = new TextEncoder().encode(JSON.stringify(row)).length; - if (batch.length && (size + n > LAYA_LIMITS[processing].questions || bytes + rowBytes > 200_000)) { - batches.push(batch); batch = []; size = 0; bytes = 0; - } - batch.push(row); size += n; bytes += rowBytes; cost += n; + const questions = task.multi ? task.labels.length : 1; + if (questions > LAYA_LIMITS[processing].questions) throw new LayaError("Too many questions for Laya fast; use processing: bulk", 400); + cost += questions; } if (cost > LAYA_LIMITS[processing].rpm) throw new LayaError("Batch exceeds the Laya per-minute question quota; split it into smaller calls", 400); - if (batch.length) batches.push(batch); - return { batches, tasks, cost, processing }; + return { tasks, cost, processing, model }; } export type QuotaTiming = Partial>; @@ -76,80 +91,98 @@ export async function limitLaya(env: LayaEnv, lane: Processing, owner: string, c return result.remaining; } catch (error) { if (error instanceof LayaError) throw error; - // Unlike the legacy quota, this protects a deliberately small GPU pool. + // Unlike the legacy quota, this protects a deliberately small shared pool. throw new LayaError("Laya admission control is temporarily unavailable", 503); } } -const probability = (v: unknown): v is number => typeof v === "number" && Number.isFinite(v) && v >= 0 && v <= 1; - -export type LayaTiming = { fetchMs?: number; headersMs?: number; backendMs?: number }; +/** + * Beam reports no server-side duration, so there is nothing to record beyond + * the client-side span. The Modal deployment's `inference_ms` and the + * headers/body split went with it rather than being reported as zero. + */ +export type LayaTiming = { fetchMs?: number }; -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); - const model = layaModel(plan.processing); - const output: JevResult[] = []; - let completeBackendTiming = true; - const deadline = Date.now() + 90_000; - for (const batch of plan.batches) { - if (Date.now() >= deadline) throw new LayaError("Laya request deadline exceeded; use smaller batches", 503); - await meter?.beforeCall?.("modal", model, 0); - let response: Response; - const started = performance.now(); - 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: 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) - throw new LayaError("Laya is busy or starting; retry with backoff", response.status, response.status === 503 ? 10 : 1); - if (response.status === 400 || response.status === 413) - throw new LayaError("Laya rejected the input: shorten text, instructions, or labels to fit the selected checkpoint's context", 400); - if (!response.ok) throw new LayaError("Laya inference failed", 502); - let body: { inference_ms?: unknown; results?: Array<{ routing?: { model?: string }; answers?: Record; noul?: number }>; usage?: { input_tokens?: number } }> }; - try { body = await response.json(); } catch { throw new LayaError("Invalid Laya response", 502); } - if (timing) { - timing.fetchMs = (timing.fetchMs ?? 0) + performance.now() - started; - if (typeof body?.inference_ms === "number" && Number.isFinite(body.inference_ms) && body.inference_ms >= 0) - timing.backendMs = (timing.backendMs ?? 0) + body.inference_ms; - else completeBackendTiming = false; +/** One question per label for multi-label, one choice otherwise — Jev's own shape. */ +function groupsFor(plan: LayaPlan): JevQuestionGroup[] { + return plan.tasks.map((task, index) => { + const id = `i${index}`; + const questions: Record = {}; + if (task.multi) { + task.labels.forEach((label, i) => { + questions[`${id}_${i}`] = { type: "noul", instructions: `Does this text belong to the category ${JSON.stringify(label)}? ${task.instructions ?? ""}`.trim() }; + }); + } else { + questions[id] = { + type: "choice", + instructions: `Which category fits this text? ${task.instructions ?? ""}`.trim(), + criteria: Object.fromEntries(task.labels.map(label => [label, null])), + }; } - if (!Array.isArray(body?.results) || body.results.length !== batch.length) throw new LayaError("Incomplete Laya response", 502); - for (const row of body.results) { - if (!row || typeof row !== "object") throw new LayaError("Invalid Laya response", 502); - const checkpoint = row.routing?.model; - if (!checkpoint || !["english", "multilingual", "typed-decisions"].includes(checkpoint)) - throw new LayaError("Invalid Laya routing metadata", 502); - const task = plan.tasks[output.length]; - let scores: Record, label: string, confidence: number; - if (task.multi) { - scores = {}; - task.labels.forEach((name, i) => { - const value = row.answers?.[`q${i}`]?.noul; - if (!probability(value)) throw new LayaError("Invalid Laya probabilities", 502); - Object.defineProperty(scores, name, { value, enumerable: true }); - }); - label = task.labels.reduce((a, b) => scores[a] >= scores[b] ? a : b); - confidence = scores[label]; - } else { - const answer = row.answers?.q; - if (!answer || !task.labels.includes(answer.choice ?? "") || !probability(answer.confidence) || - !answer.probabilities || !task.labels.every(l => probability(answer.probabilities?.[l]))) + return { state: { id, text: task.input }, questions, value: index }; + }); +} + +/** + * A Beam refusal in the vocabulary callers already handle. The mapping is + * deliberately lossy — upstream strings can echo caller text, so only the + * status and our own message travel outward. + */ +function asLayaError(error: unknown): LayaError { + if (!(error instanceof JevError)) return new LayaError("Laya inference failed", 502); + // Input the model will not accept, however Beam phrased the refusal. + if (error.errorType === "max_tokens_exceeded" || error.errorType === "invalid_request" || + error.status === 400 || error.status === 413 || error.status === 422) + return new LayaError("Laya rejected the input: shorten text, instructions, or labels to fit the selected checkpoint's context", 400); + if (error.status === 429) return new LayaError("Laya is busy; retry with backoff", 429, 1); + // A key or balance problem is ours, not the caller's: it reads as unavailable. + if (error.status === 401 || error.status === 402 || error.status === 403 || error.status === 503) + return new LayaError("Laya trial is currently unavailable", 503, 10); + if (error.errorType === "timeout" || error.errorType === "network") + return new LayaError("Laya timed out or could not be reached", 503, 5); + return new LayaError("Laya inference failed", 502); +} + +export async function runLaya(env: LayaEnv, keys: JevKeys, plan: LayaPlan, meter?: Meter, timing?: LayaTiming, signal?: AbortSignal): Promise { + if (env.LAYA_ENABLED !== "true" || !keys.beam) throw new LayaError("Laya trial is currently unavailable", 503); + const backend = BACKEND_FOR[plan.model]; + const started = performance.now(); + let answered; + try { + answered = await runJevBatches(keys, prepareJevBatches(groupsFor(plan), { limits: backend.limits }), meter, backend, signal); + } catch (error) { + throw asLayaError(error); + } + if (timing) timing.fetchMs = performance.now() - started; + + const out: JevResult[] = new Array(plan.tasks.length); + for (const { value: index, model, answers } of answered) { + const task = plan.tasks[index]; + const id = `i${index}`; + let scores: Record, label: string, confidence: number; + if (task.multi) { + scores = {}; + task.labels.forEach((name, i) => { + const value = answers[`${id}_${i}`]?.noul; + if (typeof value !== "number" || !Number.isFinite(value) || value < 0 || value > 1) throw new LayaError("Invalid Laya probabilities", 502); - label = answer.choice!; confidence = answer.confidence; - scores = Object.fromEntries(task.labels.map(l => [l, answer.probabilities![l]])); - } - output.push({ label, confidence, scores, model: layaModel(plan.processing, checkpoint) }); + // Define own properties so a label such as __proto__ stays data. + Object.defineProperty(scores, name, { value: Number(value.toFixed(4)), enumerable: true, writable: true, configurable: true }); + }); + label = task.labels.reduce((a, b) => (scores[a] >= scores[b] ? a : b)); + confidence = scores[label]; + } else { + const answer = answers[id]; + const ok = answer && task.labels.includes(answer.choice ?? "") && + typeof answer.confidence === "number" && answer.confidence >= 0 && answer.confidence <= 1 && + answer.probabilities && task.labels.every(l => typeof answer.probabilities![l] === "number"); + if (!ok) throw new LayaError("Invalid Laya probabilities", 502); + label = answer!.choice!; + confidence = Number(answer!.confidence!.toFixed(4)); + scores = Object.fromEntries(task.labels.map(l => [l, Number(answer!.probabilities![l].toFixed(4))])); } - const counts = body.results.map(r => r.usage?.input_tokens); - addTokens(meter, "modal", model, { inputTokens: counts.every(n => Number.isSafeInteger(n) && n! >= 0) ? counts.reduce((sum, n) => sum + n!, 0) : undefined, - outputTokens: 0, cachedInputTokens: 0 }); + out[index] = { label, confidence, scores, model }; } - if (timing && !completeBackendTiming) delete timing.backendMs; - return output; + if (out.some(row => row === undefined)) throw new LayaError("Incomplete Laya response", 502); + return out; } diff --git a/src/mcp.ts b/src/mcp.ts index d4e89fa..d583b02 100644 --- a/src/mcp.ts +++ b/src/mcp.ts @@ -136,7 +136,7 @@ export function productServer(classify: ClassifyFn): McpServer { inputSchema: { type: "object", properties: { inputs: INPUTS_SCHEMA, labels: LABELS_SCHEMA, instructions: INSTRUCTIONS_SCHEMA, tier: TIER_SCHEMA, - model: { type: "string", enum: ["jev", "laya"], description: "Optional Laya trial with automatic English/multilingual checkpoint routing; Jev remains the default." }, + model: { type: "string", enum: ["jev", "laya", "kev"], description: "Optional Beam-hosted models: 'laya' (ModernBERT-large, 512-token context) or 'kev' (Qwen2.5-0.5B with a pointer head, 8K context). Jev remains the default." }, processing: { type: "string", enum: ["fast", "bulk"], description: "Optional. Implies Laya if model is omitted; has no effect with explicit Jev. With Laya, omit to select fast for one decision or bulk for batches automatically. Explicit fast accepts one decision. Shared capacity limits can return 429." } }, required: ["inputs", "labels"], additionalProperties: false, @@ -189,7 +189,7 @@ export function productServer(classify: ClassifyFn): McpServer { inputSchema: { type: "object", required: ["items", "dimensions"], additionalProperties: false, properties: { items: INPUTS_SCHEMA, dimensions: DIMENSIONS_SCHEMA, instructions: { ...INSTRUCTIONS_SCHEMA, maxLength: 4000 }, tier: TIER_SCHEMA, - model: { type: "string", enum: ["jev", "laya"] }, processing: { type: "string", enum: ["fast", "bulk"], description: "Optional. Implies Laya if model is omitted; has no effect with explicit Jev. Omit for automatic fast/bulk selection based on item × dimension decisions." } }, + model: { type: "string", enum: ["jev", "laya", "kev"] }, processing: { type: "string", enum: ["fast", "bulk"], description: "Optional. Implies Laya if model is omitted; has no effect with explicit Jev. Omit for automatic fast/bulk selection based on item × dimension decisions." } }, }, annotations: { readOnlyHint: true, destructiveHint: false, idempotentHint: true, openWorldHint: false }, async run(a, ctx) { @@ -213,7 +213,7 @@ export function productServer(classify: ClassifyFn): McpServer { labels: LABELS_SCHEMA, instructions: INSTRUCTIONS_SCHEMA, max_labels: { type: "integer", minimum: 1, maximum: 100, description: "At most this many labels per text, most likely first." }, - model: { type: "string", enum: ["jev", "laya"] }, + model: { type: "string", enum: ["jev", "laya", "kev"] }, processing: { type: "string", enum: ["fast", "bulk"], description: "Optional. Implies Laya if model is omitted; has no effect with explicit Jev. Omit to select fast for up to four labels on one text, or bulk for larger work automatically." }, }, required: ["inputs", "labels"], diff --git a/src/openapi.ts b/src/openapi.ts index b328193..47f8766 100644 --- a/src/openapi.ts +++ b/src/openapi.ts @@ -994,7 +994,7 @@ export const OPENAPI = { { inputs: ["postgres index tuning for ML feature stores"], labels: ["databases", "ml", "frontend"], multi: true, max_labels: 2 }, ], properties: { - model: { type: "string", enum: ["jev", "laya"], description: "When omitted, uses Jev unless processing is supplied, which implies Laya. Opt into experimental Laya with automatic English/multilingual routing. Context is 512 tokens for English or 1,024 for multilingual, with 2–16 short labels, text ≤2,000 characters and instructions ≤400 characters. Results identify the actual checkpoint and lane. Jev calibration claims do not apply to Laya." }, + model: { type: "string", enum: ["jev", "laya", "kev"], description: "When omitted, uses Jev unless processing is supplied, which implies Laya. Laya and Kev are experimental models hosted by Beam: \"laya\" is ModernBERT-large with a 512-token context, \"kev\" is Qwen2.5-0.5B with a pointer head and an 8,192-token context. Both take 2–16 short labels, text ≤2,000 characters and instructions ≤400 characters. A result is labelled jev/laya or jev/kev; Beam does not report which checkpoint answered. Jev calibration claims do not apply to either." }, processing: { type: "string", enum: ["fast", "bulk"], description: "Implies Laya when model is omitted. Accepted but has no effect with explicit model jev, which handles batching automatically. When omitted for Laya, automatically selects fast for one decision with up to 4 yes/no questions, otherwise bulk. Explicit Laya lanes are honored. Fast allows 60 questions/min and 2,000/day per caller. Bulk chunks batches up to 1,000 questions per call, 1,000/min and 20,000/day. These caps also apply to paid/operator keys. Same model weights in both lanes. Overload returns 429; a cold bulk worker returns 503 with Retry-After. Smart review is independent." }, dimensions: DIMENSIONS_SCHEMA, items: { type: "array", minItems: 1, maxItems: 1000, items: { type: "string", minLength: 1, maxLength: 32000 }, description: "Alias for inputs in dimensions mode. Do not combine with input or inputs." }, diff --git a/src/retail-rates.json b/src/retail-rates.json index ebf6849..3d0cdf4 100644 --- a/src/retail-rates.json +++ b/src/retail-rates.json @@ -1,5 +1,5 @@ { - "version": "2026-09-20-laya-router-trial-v2", + "version": "2026-09-21-beam-laya-kev-v1", "models": [ { "provider": "typesafe", @@ -15,7 +15,7 @@ "outputUsdPerMillion": "4.500", "cachedInputUsdPerMillion": "0.090" }, - { "provider": "modal", "model": "laya-0.3.4-routed-fast", "inputUsdPerMillion": "0", "outputUsdPerMillion": "0", "cachedInputUsdPerMillion": "0" }, - { "provider": "modal", "model": "laya-0.3.4-routed-bulk", "inputUsdPerMillion": "0", "outputUsdPerMillion": "0", "cachedInputUsdPerMillion": "0" } + { "provider": "beam", "model": "jev/laya", "inputUsdPerMillion": "0", "outputUsdPerMillion": "0", "cachedInputUsdPerMillion": "0" }, + { "provider": "beam", "model": "jev/kev", "inputUsdPerMillion": "0", "outputUsdPerMillion": "0", "cachedInputUsdPerMillion": "0" } ] } diff --git a/src/server/token-pricing.ts b/src/server/token-pricing.ts index 2c2a1ce..8df3399 100644 --- a/src/server/token-pricing.ts +++ b/src/server/token-pricing.ts @@ -44,7 +44,7 @@ export function parseTokenRateCard(json: string | undefined): TokenRateCard | nu const seen = new Set(); const models = data.models.map(value => { const row = object(value); - if (!row || !["typesafe", "vercel", "openrouter", "modal"].includes(String(row.provider)) || typeof row.model !== "string" || !row.model.trim() || row.model.length > 200) { + if (!row || !["typesafe", "vercel", "openrouter", "beam"].includes(String(row.provider)) || typeof row.model !== "string" || !row.model.trim() || row.model.length > 200) { throw new Error("Invalid provider or model in token rate card."); } const key = `${row.provider}:${row.model}`; diff --git a/src/server/token-reservation.ts b/src/server/token-reservation.ts index d35fa03..309dee3 100644 --- a/src/server/token-reservation.ts +++ b/src/server/token-reservation.ts @@ -4,8 +4,11 @@ import type { TokenRateCard } from "./token-pricing"; /** Provider context limits, not token estimates. Unknown models cannot spend. */ export function providerCallBound(card: TokenRateCard, provider: ModelTokenUsage["provider"], model: string, maxOutput: number): number { + // Each bound is the model's own context window, so a single call can never + // reserve more than that model could physically consume. const inputLimit = provider === "typesafe" && model === JEV_ACCOUNT_MODEL ? 65_536 - : provider === "modal" && ["laya-0.3.4-routed-fast", "laya-0.3.4-routed-bulk"].includes(model) ? 64 * 1024 + : provider === "beam" && model === "jev/laya" ? 512 + : provider === "beam" && model === "jev/kev" ? 8_192 : provider === "openrouter" && model === "google/gemini-3.8-flash" ? 1_048_576 : null; const rate = card.models.find((row) => row.provider === provider && row.model === model); if (inputLimit === null || !rate || !Number.isSafeInteger(maxOutput) || maxOutput < 0 || maxOutput > 65_536) diff --git a/test/laya-timing.test.ts b/test/laya-timing.test.ts index a27e6be..ea1d69d 100644 --- a/test/laya-timing.test.ts +++ b/test/laya-timing.test.ts @@ -16,65 +16,56 @@ afterEach(() => { clock?.mockRestore(); clock = undefined; }); -const env = { LAYA_ENABLED: "true", LAYA_FAST_URL: "https://fast.example", LAYA_BULK_URL: "https://bulk.example", - LAYA_MODAL_KEY: "fixture", LAYA_MODAL_SECRET: "fixture" } as LayaEnv; +const env = { LAYA_ENABLED: "true" } as LayaEnv; +const keys = { beam: "beam-key" }; const task = { input: "refund", labels: ["billing", "support"] }; -function mockBackend(durations: unknown[]) { +/** + * Beam answers in System One's shape and reports no server-side duration, so + * the only span there is to record is the client-side one. + */ +function mockBeam() { let calls = 0; globalThis.fetch = (async (_url, init) => { - const { batch } = JSON.parse(String(init?.body)); - const body = { inference_ms: durations[calls++], results: batch.map(() => ({ - routing: { model: "english" }, usage: { input_tokens: 20 }, answers: { - q: { choice: "billing", confidence: .9, probabilities: { billing: .9, support: .1 } }, - }, - })) }; - // Preserve non-finite values here to exercise the parser's runtime guard, - // rather than JSON.stringify silently converting them to null. - const response = Response.json({}); - response.json = async () => body; - return response; + calls++; + const body = JSON.parse(String(init?.body)) as { state: { id: string }[] }; + const answers = Object.fromEntries(body.state.map(item => [item.id, { + choice: "billing", confidence: .9, probabilities: { billing: .9, support: .1 }, + }])); + return Response.json({ model: "jev/laya", answers, usage: { input_tokens: 20 * body.state.length, output_tokens: 0 } }); }) as typeof fetch; return () => calls; } -test.each([0, 12.5])("records finite nonnegative backend duration %s", async backendMs => { - mockBackend([backendMs]); +test("records a client-side span for a fast call", async () => { + mockBeam(); const timing: LayaTiming = {}; - const results = await runLaya(env, planLaya([task], "fast"), undefined, timing); + const results = await runLaya(env, keys, planLaya([task], "fast"), undefined, timing); expect(results[0].label).toBe("billing"); - expect(timing.backendMs).toBe(backendMs); expect(Number.isFinite(timing.fetchMs)).toBe(true); - expect(timing.fetchMs!).toBeGreaterThanOrEqual(timing.headersMs!); + expect(timing.fetchMs!).toBeGreaterThanOrEqual(0); }); -test.each([undefined, null, -1, "12.5", NaN, Infinity, {}])("ignores absent or invalid backend duration %j", async duration => { - mockBackend([duration]); +test("reports no backend duration, because Beam does not return one", async () => { + mockBeam(); const timing: LayaTiming = {}; - expect(await runLaya(env, planLaya([task], "fast"), undefined, timing)).toHaveLength(1); + await runLaya(env, keys, planLaya([task], "fast"), undefined, timing); expect(timing).not.toHaveProperty("backendMs"); - expect(Number.isFinite(timing.fetchMs)).toBe(true); + expect(timing).not.toHaveProperty("headersMs"); }); -test("aggregates bulk chunk durations instead of overwriting with the last chunk", async () => { - const calls = mockBackend([2.25, 4.5]); - const timestamps = [0, 3, 5, 10, 14, 17]; - clock = spyOn(performance, "now").mockImplementation(() => timestamps.shift()!); +test("a bulk span covers every chunk, not just the last", async () => { + const calls = mockBeam(); + const timestamps = [0, 17]; + clock = spyOn(performance, "now").mockImplementation(() => timestamps.shift() ?? 17); const timing: LayaTiming = {}; - expect(await runLaya(env, planLaya(Array.from({ length: 65 }, () => task), "bulk"), undefined, timing)).toHaveLength(65); - expect(calls()).toBe(2); - expect(timing).toEqual({ headersMs: 7, fetchMs: 12, backendMs: 6.75 }); + // 65 items exceed Laya's 16-item ceiling, so this is several Beam requests. + expect(await runLaya(env, keys, planLaya(Array.from({ length: 65 }, () => task), "bulk"), undefined, timing)).toHaveLength(65); + expect(calls()).toBeGreaterThan(1); + expect(timing).toEqual({ fetchMs: 17 }); }); test("timing remains optional for existing callers", async () => { - mockBackend([12]); - expect(await runLaya(env, planLaya([task], "fast"))).toHaveLength(1); -}); - -test.each([{ durations: [2.25, undefined] }, { durations: [undefined, 2.25] }])("omits an incomplete bulk backend total %j", async ({ durations }) => { - mockBackend(durations); - const timing: LayaTiming = {}; - expect(await runLaya(env, planLaya(Array.from({ length: 65 }, () => task), "bulk"), undefined, timing)).toHaveLength(65); - expect(timing).not.toHaveProperty("backendMs"); - expect(Number.isFinite(timing.fetchMs)).toBe(true); + mockBeam(); + expect(await runLaya(env, keys, planLaya([task], "fast"))).toHaveLength(1); }); diff --git a/test/laya.test.ts b/test/laya.test.ts index f6b88f4..825cda2 100644 --- a/test/laya.test.ts +++ b/test/laya.test.ts @@ -1,6 +1,8 @@ import { test, expect, afterEach } from "bun:test"; import worker, { type Env } from "../src/index"; import { planLaya, runLaya, limitLaya } from "../src/laya"; + +const keys = { beam: "beam-key" }; import { newMeter } from "../src/cost"; import { parseTokenRateCard, priceTokens } from "../src/server/token-pricing"; import { providerCallBound } from "../src/server/token-reservation"; @@ -8,30 +10,40 @@ import rates from "../src/retail-rates.json"; const original = globalThis.fetch; afterEach(() => { globalThis.fetch = original; }); -const env = { LAYA_ENABLED: "true", LAYA_FAST_URL: "https://fast.example", LAYA_BULK_URL: "https://bulk.example", - LAYA_MODAL_KEY: "key", LAYA_MODAL_SECRET: "secret", STATS: { get: async () => null, put: async () => {} }, +const env = { LAYA_ENABLED: "true", BEAM_API_KEY: "beam-key", STATS: { get: async () => null, put: async () => {} }, LIMITER: { idFromName: (s: string) => s, get: () => ({ fetch: async () => Response.json({ limited: false, remaining: 59 }) }) } } as unknown as Env; const ctx = { waitUntil: () => {} } as unknown as ExecutionContext; const task = { input: "refund", labels: ["billing", "tech"] }; const request = (body: object, bindings = env) => worker.fetch(new Request("https://classifier.dev/v1/classify", { method: "POST", body: JSON.stringify({ model: "laya", ...body }) }), bindings, ctx); -function mockLaya() { - const calls: { url: string; rows: { state: string; questions: Record }[] }[] = []; +function mockLaya(model = "jev/laya") { + // Beam speaks System One: state items and named questions, not Modal rows. + const calls: { url: string; model: string; rows: { state: string; questions: Record }[] }[] = []; globalThis.fetch = (async (url, init) => { - const { batch } = JSON.parse(String(init?.body)); - calls.push({ url: String(url), rows: batch }); - return Response.json({ results: batch.map((row: typeof calls[number]["rows"][number]) => ({ - answers: Object.fromEntries(Object.entries(row.questions).map(([id, q]) => { - const labels = Object.keys(q.criteria ?? {}); - return [id, q.type === "noul" ? { noul: .9 } : { choice: labels[0], confidence: .9, - probabilities: Object.fromEntries(labels.map((label, i) => [label, i ? .1 / (labels.length - 1) : .9])) }]; - })), usage: { input_tokens: 20 }, routing: { model: "english" }, - })) }); + const body = JSON.parse(String(init?.body)) as { + model: string; + state: { id: string; text: string }[]; + questions: Record }>; + }; + // Regroup by state id so chunking assertions still read as rows. + const rows = body.state.map(item => ({ + state: item.text, + questions: Object.fromEntries(Object.entries(body.questions).filter(([id]) => id === item.id || id.startsWith(`${item.id}_`))), + })); + calls.push({ url: String(url), model: body.model, rows }); + const answers = Object.fromEntries(Object.entries(body.questions).map(([id, q]) => { + const labels = Object.keys(q.criteria ?? {}); + return [id, q.type === "noul" ? { noul: .9 } : { choice: labels[0], confidence: .9, + probabilities: Object.fromEntries(labels.map((label, i) => [label, i ? .1 / (labels.length - 1) : .9])) }]; + })); + return Response.json({ model, answers, usage: { input_tokens: 20 * body.state.length, output_tokens: 0 } }); }) as typeof fetch; return calls; } +const BEAM_URL = "https://app.beam.cloud/v1/systemone"; + function combinedEnv(result: unknown | ((url: string, init: RequestInit) => Promise) = {limited:false,remaining:2900,laneRemaining:56}) { const admissions: any[] = [], admissionUrls: string[] = []; let legacyCalls = 0; @@ -127,39 +139,44 @@ test.each([false, true])("local burst shield blocks inference when unavailable=% 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) })); - const result = await runLaya(env, planLaya(tasks, "bulk"), meter); - expect(calls).toHaveLength(16); - expect(calls.every(c => c.url === "https://bulk.example/predict" && c.rows.length <= 64)).toBe(true); + const result = await runLaya(env, keys, planLaya(tasks, "bulk"), meter); + expect(calls).toHaveLength(63); + expect(calls.every(c => c.url === BEAM_URL && c.rows.length <= 32)).toBe(true); expect(calls.flatMap(c => c.rows.map(r => r.state))).toEqual(tasks.map(t => t.input)); expect(result).toHaveLength(1000); - expect(meter.tokens[0]).toMatchObject({ provider: "modal", model: "laya-0.3.4-routed-bulk", calls: 16, inputTokens: 20_000 }); + expect(meter.tokens[0]).toMatchObject({ provider: "beam", model: "jev/laya", calls: 63, inputTokens: 20_000 }); }); -test("mixed Router checkpoints preserve row order and meter one routed wrapper", async () => { - const checkpoints = ["english", "multilingual", "typed-decisions", "multilingual"]; +test("bulk preserves row order across items and meters one Beam model", async () => { + const meter = newMeter(), bounds: string[] = []; const expected = ["billing", "tech", "billing", "tech"]; - const meter = newMeter(); - const bounds: string[] = []; meter.beforeCall = async (_provider, model) => { bounds.push(model); }; - globalThis.fetch = (async () => Response.json({ results: checkpoints.map((model, i) => ({ - routing: { model }, usage: { input_tokens: 20 + i }, answers: { q: { + globalThis.fetch = (async (_url, init) => { + const body = JSON.parse(String(init?.body)) as { state: { id: string }[] }; + const answers = Object.fromEntries(body.state.map((item, i) => [item.id, { choice: expected[i], confidence: .9, probabilities: { billing: expected[i] === "billing" ? .9 : .1, tech: expected[i] === "tech" ? .9 : .1 }, - } }, - })) })) as typeof fetch; - const results = await runLaya(env, planLaya(["refund", "无法登录", "invoice 123", "connexion impossible"].map(input => ({ ...task, input })), "bulk"), meter); + }])); + return Response.json({ model: "jev/laya", answers, usage: { input_tokens: 86, output_tokens: 0 } }); + }) as typeof fetch; + const results = await runLaya(env, keys, planLaya(["refund", "无法登录", "invoice 123", "connexion impossible"].map(input => ({ ...task, input })), "bulk"), meter); expect(results.map(result => result.label)).toEqual(expected); - expect(results.map(result => result.model)).toEqual(checkpoints.map(model => `laya-0.3.4-${model}-bulk`)); - expect(bounds).toEqual(["laya-0.3.4-routed-bulk"]); - expect(meter.tokens).toEqual([{ provider: "modal", model: "laya-0.3.4-routed-bulk", calls: 1, + // Beam never says which checkpoint answered, so every row carries the model id. + expect(results.map(result => result.model)).toEqual(Array(4).fill("jev/laya")); + expect(bounds).toEqual(["jev/laya"]); + expect(meter.tokens).toEqual([{ provider: "beam", model: "jev/laya", calls: 1, inputTokens: 86, outputTokens: 0, cachedInputTokens: 0 }]); }); -test.each([undefined, {}, { model: "unknown" }, { model: ["english"] }])("missing or unknown Router checkpoint fails closed: %j", async routing => { - globalThis.fetch = (async () => Response.json({ results: [{ routing, usage: { input_tokens: 20 }, answers: { - q: { choice: "billing", confidence: .9, probabilities: { billing: .9, tech: .1 } }, - } }] })) as typeof fetch; - expect((await request({ input: task.input, labels: task.labels })).status).toBe(502); +test.each([ + ["no answers at all", {}], + ["a choice outside the label set", { i0: { choice: "elsewhere", confidence: .9, probabilities: { billing: .9, tech: .1 } } }], + ["a missing probability", { i0: { choice: "billing", confidence: .9, probabilities: { billing: .9 } } }], + ["a non-numeric confidence", { i0: { choice: "billing", confidence: "high", probabilities: { billing: .9, tech: .1 } } }], +])("a malformed Beam answer fails closed rather than answering: %s", async (_name, answers) => { + globalThis.fetch = (async () => Response.json({ model: "jev/laya", answers, usage: { input_tokens: 20 } })) as typeof fetch; + const response = await request({ input: task.input, labels: task.labels }); + expect(response.status).toBe(502); }); test("fast refuses large work before fetching; multi counts each label", () => { @@ -167,7 +184,6 @@ test("fast refuses large work before fetching; multi counts each label", () => { expect(() => planLaya([{ ...task, input: "x".repeat(2001) }], "bulk")).toThrow("short text"); const plan = planLaya(Array.from({ length: 100 }, () => ({ ...task, multi: true })), "bulk"); expect(plan.cost).toBe(200); - expect(plan.batches.map(b => b.length)).toEqual([32, 32, 32, 4]); }); test("POST selects fast or bulk and returns familiar output", async () => { @@ -176,8 +192,8 @@ test("POST selects fast or bulk and returns familiar output", async () => { expect(response.status).toBe(200); expect(response.headers.get("ratelimit-limit")).toBe("60"); expect(response.headers.get("x-classifier-processing")).toBe("fast"); - expect(await response.json()).toMatchObject({ results: [{ label: "billing", model: "laya-0.3.4-english-fast" }] }); - expect(calls[0].url).toContain("fast.example"); + expect(await response.json()).toMatchObject({ results: [{ label: "billing", model: "jev/laya" }] }); + expect(calls[0].url).toBe(BEAM_URL); const bulk = await request({ inputs: ["a", "b"], labels: task.labels, processing: "bulk", multi: true }); expect(bulk.status).toBe(200); expect(await bulk.json()).toMatchObject({ results: [{ labels: ["billing", "tech"] }, { labels: ["billing", "tech"] }] }); @@ -201,14 +217,14 @@ test.each([ const response = await request(body); expect(response.status).toBe(200); expect(response.headers.get("x-classifier-processing")).toBe("bulk"); - expect(calls.every(call => call.url === "https://bulk.example/predict")).toBe(true); + expect(calls.every(call => call.url === BEAM_URL)).toBe(true); }); test.each(["fast", "bulk"])("processing %s implies Laya when model is omitted", async processing => { const calls = mockLaya(); const response = await request({ model: undefined, input: task.input, labels: task.labels, processing }); expect(response.status).toBe(200); - expect(calls[0].url).toBe(`https://${processing}.example/predict`); + expect(calls[0].url).toBe(BEAM_URL); }); test("explicit model and lane choices are not silently overridden", async () => { @@ -230,7 +246,7 @@ test("chat's MCP batch succeeds with omitted model or omitted lane", async () => expect(body.result.isError).not.toBe(true); expect(body.result.structuredContent.results).toHaveLength(30); } - expect(calls.every(call => call.url === "https://bulk.example/predict")).toBe(true); + expect(calls.every(call => call.url === BEAM_URL)).toBe(true); }); test.each([undefined, "fast", "bulk"])("explicit Jev batches accept processing hint %s without selecting Laya", async processing => { @@ -297,11 +313,10 @@ test("anonymous Laya success reports only numeric critical-path timings", async const rows = header.split(", "); expect(rows.every(row => /^[a-z_]+;dur=\d+\.\d{2}$/.test(row))).toBe(true); const timing = Object.fromEntries(rows.map(row => { const [name, duration] = row.split(";dur="); return [name, Number(duration)]; })); - expect(timing.backend).toBe(7.25); expect(timing.quota_regular).toBeGreaterThanOrEqual(10); expect(timing.quota_laya).toBeGreaterThanOrEqual(10); - expect(timing.modal_fetch).toBeGreaterThanOrEqual(10); - expect(timing.laya_run).toBeGreaterThanOrEqual(timing.modal_fetch); + expect(timing.beam_fetch).toBeGreaterThanOrEqual(10); + expect(timing.laya_run).toBeGreaterThanOrEqual(timing.beam_fetch); expect(timing.worker_total).toBeGreaterThanOrEqual(timing.quota_regular + timing.quota_laya + timing.laya_run - .1); expect(header).not.toContain(task.input); }); @@ -351,12 +366,11 @@ test("normal batch and tier rejection cannot debit Laya quota", async () => { expect(calls).toHaveLength(0); }); -test("trial billing explicitly prices both Laya lanes at zero without making reviews free", () => { +test("trial billing explicitly prices both Beam models at zero without making reviews free", () => { const card = parseTokenRateCard(JSON.stringify(rates))!; - for (const lane of ["fast", "bulk"]) { - const model = `laya-0.3.4-routed-${lane}`; - expect(providerCallBound(card, "modal", model, 0)).toBe(0); - expect(priceTokens(card, [{ provider: "modal", model, calls: 1, inputTokens: 1000, outputTokens: 0, cachedInputTokens: 0 }])?.nanodollars).toBe(0n); + for (const model of ["jev/laya", "jev/kev"]) { + expect(providerCallBound(card, "beam", model, 0)).toBe(0); + expect(priceTokens(card, [{ provider: "beam", model, calls: 1, inputTokens: 1000, outputTokens: 0, cachedInputTokens: 0 }])?.nanodollars).toBe(0n); } expect(providerCallBound(card, "openrouter", "google/gemini-3.8-flash", 2000)).toBeGreaterThan(0); }); diff --git a/tests/app-http.test.ts b/tests/app-http.test.ts index c3ca378..7c8772f 100644 --- a/tests/app-http.test.ts +++ b/tests/app-http.test.ts @@ -497,21 +497,24 @@ test.each([ ])("account Laya %j settles free inference without a Jev credential", async selection => { const { processing } = selection; delete env.TYPESAFE_API_KEY; - Object.assign(env, { LAYA_ENABLED: "true", LAYA_FAST_URL: "https://fast.example", LAYA_BULK_URL: "https://bulk.example", - LAYA_MODAL_KEY: "fixture", LAYA_MODAL_SECRET: "fixture", LIMITER: { + Object.assign(env, { LAYA_ENABLED: "true", BEAM_API_KEY: "fixture", LIMITER: { idFromName: (name: string) => name, get: () => ({ fetch: async () => Response.json({ limited: false, remaining: 59 }) }), } }); // A free inference must also work when the account has no paid or free credits. await env.APP_DB.prepare("UPDATE app_accounts SET balance=0,paid_balance=0 WHERE id='local-demo'").run(); globalThis.fetch = (async (url, init) => { - expect(String(url)).toBe(`https://${processing}.example/predict`); - const { batch } = JSON.parse(String(init?.body)); - return Response.json({ results: batch.map((row: { questions: Record }) => ({ - answers: Object.fromEntries(Object.keys(row.questions).map(id => [id, { + // Both lanes reach the same Beam endpoint; the lane is a quota, not a host. + expect(String(url)).toBe("https://app.beam.cloud/v1/systemone"); + const body = JSON.parse(String(init?.body)) as { model: string; questions: Record }; + expect(body.model).toBe("jev/laya"); + return Response.json({ + model: "jev/laya", + answers: Object.fromEntries(Object.keys(body.questions).map(id => [id, { choice: "yes", confidence: .99, probabilities: { yes: .99, no: .01 }, - }])), usage: { input_tokens: 41 }, routing: { model: "english" }, - })) }); + }])), + usage: { input_tokens: 41, output_tokens: 0 }, + }); }) as typeof fetch; const response = await accountClassification(request({ input: "Hello", labels: ["yes", "no"], ...selection }), env); expect(response?.status).toBe(200); @@ -534,21 +537,23 @@ test.each(["fast", "smart"])("unconfigured %s Jev inference fails before spendin }); test("account Laya Smart charges only the actual review tokens", async () => { - Object.assign(env, { LAYA_ENABLED: "true", LAYA_FAST_URL: "https://fast.example", LAYA_MODAL_KEY: "fixture", - LAYA_MODAL_SECRET: "fixture", OPENROUTER_API_KEY: "fixture", LIMITER: { + Object.assign(env, { LAYA_ENABLED: "true", BEAM_API_KEY: "fixture", OPENROUTER_API_KEY: "fixture", LIMITER: { idFromName: (name: string) => name, get: () => ({ fetch: async () => Response.json({ limited: false, remaining: 59 }) }), } }); let calls = 0; globalThis.fetch = (async (url, init) => { calls++; - if (String(url) === "https://fast.example/predict") { - const { batch } = JSON.parse(String(init?.body)); - return Response.json({ results: batch.map((row: { questions: Record }) => ({ - answers: Object.fromEntries(Object.keys(row.questions).map(id => [id, { + if (String(url) === "https://app.beam.cloud/v1/systemone") { + const body = JSON.parse(String(init?.body)) as { questions: Record }; + // Low confidence, so the smart tier escalates this answer to a review. + return Response.json({ + model: "jev/laya", + answers: Object.fromEntries(Object.keys(body.questions).map(id => [id, { choice: "yes", confidence: .6, probabilities: { yes: .6, no: .4 }, - }])), usage: { input_tokens: 41 }, routing: { model: "english" }, - })) }); + }])), + usage: { input_tokens: 41, output_tokens: 0 }, + }); } expect(String(url)).toBe("https://openrouter.ai/api/v1/chat/completions"); return Response.json({ model: "google/gemini-3.8-flash", choices: [{ message: { content: "B" } }], diff --git a/tests/plans-prices.test.ts b/tests/plans-prices.test.ts index 58cc184..1baa093 100644 --- a/tests/plans-prices.test.ts +++ b/tests/plans-prices.test.ts @@ -6,12 +6,12 @@ import { getSnapshot } from "../src/server/accounts"; import { provisionTestAccount } from "./support/account"; import { database } from "./support/postgres"; -test("account plans identify Laya trial rates without advertising free Gemini reviews", async () => { +test("account plans identify the Beam trial rates without advertising free Gemini reviews", async () => { const env = { APP_DB: database(), APP_ACCOUNTS_ENABLED: "true" }; await provisionTestAccount(new Request("http://localhost/login"), env); const html = renderToStaticMarkup(createElement(Plans, { snapshot: await getSnapshot("local-demo", env), navigate: () => {} })); const names = Array.from(html.matchAll(/]*>(.*?)<\/th>/g), match => match[1]); - expect(names).toEqual(["Jev", "Gemini escalation", "laya-0.3.4-routed-fast", "laya-0.3.4-routed-bulk"]); + expect(names).toEqual(["Jev", "Gemini escalation", "jev/laya", "jev/kev"]); expect(html).toContain("Laya inference is free during the trial"); expect(html).toContain("Smart reviews are still billed"); }); diff --git a/tests/public-pricing.test.ts b/tests/public-pricing.test.ts index 9859f80..0936d7a 100644 --- a/tests/public-pricing.test.ts +++ b/tests/public-pricing.test.ts @@ -83,8 +83,8 @@ test("pricing renders workspace plans and Pro rate limits", () => { expect(html).toContain('class="plan-action" href="/auth/sign-up"'); expect(html).toContain(`$${BILLING_PLANS.pro.priceCents / 100}`); expect(html).toContain(formatCreditsUsd(BILLING_PLANS.pro.includedCredits)); - expect(html).toContain("laya-0.3.4-routed-fast"); - expect(html).toContain("laya-0.3.4-routed-bulk"); + expect(html).toContain("jev/laya"); + expect(html).toContain("jev/kev"); expect(html).toContain('href="/auth/sign-up?returnTo=/app/plans"'); expect(html).not.toContain("classifier_pro_"); expect(html).toContain("30,000/min · 200,000/day"); diff --git a/wrangler.example.toml b/wrangler.example.toml index f3f2080..2f93b9b 100644 --- a/wrangler.example.toml +++ b/wrangler.example.toml @@ -55,10 +55,10 @@ namespace_id = "910001" [vars] # Transfer-aware legacy objects were deployed before enabling admission here. QUOTA_COORDINATOR_ENABLED = "true" +# Laya and Kev are hosted by Beam on its shared endpoints; there is no +# deployment of our own to point at. BEAM_API_KEY is a Worker secret, never a +# plaintext var, and it is the only credential either model needs. LAYA_ENABLED = "true" -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" # Gateway free-tier 429s observed in production. Re-enable only after capacity is verified.