From 9f16b3329ff5f1f2dc8e4a0db7b9695d07aa184b Mon Sep 17 00:00:00 2001 From: Michael Ryaboy Date: Wed, 23 Sep 2026 16:52:07 -0700 Subject: [PATCH] Accept whole 10M-token documents in one upload --- .github/workflows/check.yml | 8 ++ .github/workflows/deploy.yml | 3 + AGENTS.md | 2 +- README.md | 24 ++-- cli/README.md | 14 ++ cli/classify.js | 77 +++++++++- e2e/document-admission.ts | 60 ++++++++ e2e/long-context-billing.ts | 2 +- e2e/whole-document.live.mjs | 57 ++++++++ e2e/whole-document.mjs | 269 +++++++++++++++++++++++++++++++++++ package-lock.json | 25 ++++ package.json | 4 +- src/docs.ts | 98 +++++++------ src/document-upload.ts | 178 +++++++++++++++++++++++ src/http/account-api.ts | 2 +- src/http/classification.ts | 6 + src/http/document.ts | 57 ++++++++ src/index.ts | 2 +- src/long-context-job.ts | 235 +++++++++++++++++++++++++----- src/openapi.ts | 62 +++++--- src/pages.ts | 42 ++++-- src/pricingui.ts | 4 +- src/server/usage.ts | 40 +++++- wrangler.example.toml | 3 + 24 files changed, 1141 insertions(+), 133 deletions(-) create mode 100644 e2e/document-admission.ts create mode 100644 e2e/whole-document.live.mjs create mode 100644 e2e/whole-document.mjs create mode 100644 src/document-upload.ts create mode 100644 src/http/document.ts diff --git a/.github/workflows/check.yml b/.github/workflows/check.yml index 585bf6b..66563cd 100644 --- a/.github/workflows/check.yml +++ b/.github/workflows/check.yml @@ -34,6 +34,10 @@ jobs: run: bun test tests/spending.e2e.test.ts env: POSTGRES_TEST_URL: postgres://postgres:spending-test@localhost:5432/postgres + - name: Verify document admission under concurrent PostgreSQL requests + run: bun e2e/document-admission.ts + env: + POSTGRES_TEST_URL: postgres://postgres:spending-test@localhost:5432/postgres - name: Render secret-free build configuration env: CLOUDFLARE_ACCOUNT_ID: "00000000000000000000000000000000" @@ -52,6 +56,8 @@ jobs: run: npm run test:e2e:long-context-billing - name: Verify ten-million-token job and settlement run: npm run test:e2e:long-context-job + - name: Verify one whole ten-million-token document upload + run: npm run test:e2e:whole-document - name: Save long-context verification evidence uses: actions/upload-artifact@v4 with: @@ -60,6 +66,8 @@ jobs: captures/long-context-e2e.json captures/long-context-billing.json captures/long-context-job.json + captures/whole-document.json + captures/document-admission.json - name: Save classification usage evidence uses: actions/upload-artifact@v4 with: diff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml index 7c6c625..e9c9362 100644 --- a/.github/workflows/deploy.yml +++ b/.github/workflows/deploy.yml @@ -104,6 +104,8 @@ jobs: run: bun e2e/complimentary-pro.ts - name: Verify ten-million-token job and settlement run: npm run test:e2e:long-context-job + - name: Verify one whole ten-million-token document upload + run: npm run test:e2e:whole-document - name: Save long-context verification evidence uses: actions/upload-artifact@v4 with: @@ -112,6 +114,7 @@ jobs: captures/long-context-e2e.json captures/long-context-billing.json captures/long-context-job.json + captures/whole-document.json - name: Save complimentary Pro verification evidence uses: actions/upload-artifact@v4 with: diff --git a/AGENTS.md b/AGENTS.md index 9eb0a1e..e0cacd2 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -24,7 +24,7 @@ the plain text (`curl classifier.dev`), the HTML and the Markdown never drift. - Docs are plain text with UPPERCASE headings (`src/docs.ts`, `src/pages.ts`); `renderDoc` turns them into HTML and `toMarkdown` into Markdown. - Discovery files (`/.well-known/*`, sitemap, robots, auth.md) are generated in `src/wellknown.ts` from `SITE` and the MCP tool table — edit the source, never a served file. - The MCP servers (`src/mcp.ts`) are stateless Streamable HTTP; tools call the API through `worker.fetch` so limits and logging are shared. -- Long-context jobs in `src/long-context-job.ts` accept ordered bounded parts through a SQLite Durable Object. It holds a maximum-price workspace reservation, retains only selected evidence until completion/cancellation/24-hour expiry, and settles the original part-token total once. Keep `wrangler.example.toml`, the generated docs and the job E2E in sync with this contract. +- Whole-document uploads use `POST /v1/classify`: `src/document-upload.ts` stream-parses JSON or UTF-8 text up to 10M tokens/100 MB, preserving exact cl100k counts across internal fragment boundaries. `src/long-context-job.ts` stores source privately while queued, deletes it as SQLite Durable Object alarms screen it, then judges and settles automatically. Reserve the actual document price after upload. Delete remaining source/evidence on failure, cancellation or 24-hour expiry. The old manual-part API remains supported. Keep generated docs and `e2e/whole-document.mjs` in sync; test through the built Worker, real Durable Objects and PostgreSQL ledger. - Never commit secrets; `.secrets.env`, `.dev.vars` are ignored. `eval/data/` is ignored except the summary copied to `src/vs-jev.json`. - Measured numbers on the site come from `eval/`; do not type numbers in by hand. - Jev is asked through Vercel's AI Gateway first when `AI_GATEWAY_API_KEY` diff --git a/README.md b/README.md index e2066dd..9b6cdfa 100644 --- a/README.md +++ b/README.md @@ -226,7 +226,7 @@ provider context check. Eligible chunks can be omitted when that budget fills; This does not guarantee full-document final reading or universal accuracy. Use Fast with a workspace key backed by paid balance or an active paid -subscription; anonymous access and signup credit do not qualify. Per request: +subscription; anonymous access and signup credit do not qualify. Synchronous limits: 250,000 original context tokens summed across inputs, 20 documents, 32 decisions, and a 1 MB body. Retail is $0.084/M original `cl100k_base` context tokens, counted once across inputs regardless of dimensions and actual inference usage. @@ -236,15 +236,19 @@ tokenization or chunking, bounding tokenizer work on pathological inputs. No eligible evidence returns `422 long_context_no_evidence` without charge. Explicit `model: "chunklaya"` keeps its separate legacy opt-in behavior. -For up to 10 million input tokens, funded workspaces use a long-context job: -create a UUID job with a `max_tokens` ceiling, upload ordered parts of at most -50,000 tokens and 1 MB each, then call `finish`. The server screens every part -without retaining the full source and keeps a bounded set of selected excerpts -until completion or 24-hour expiry. The final Jev call reads at most 20,000 -selected evidence tokens per document. `max_tokens` reserves workspace credit; -successful completion charges the actual part-token sum at the same $0.084/M -rate. Cancellation, expiry and no eligible evidence refund the reservation. -The full API sequence and resume/status route are in the generated docs. +Send one whole document of up to 10 million input tokens (100 MB) to +`POST /v1/classify` with a funded workspace key. Use the normal JSON input and +labels, or send `text/plain` with a `labels` query parameter. Large uploads +return `202` with a `status_url`; `Prefer: respond-async` also works for smaller +documents. The server splits, screens, retries and finishes automatically. +The repository CLI handles upload and polling with +`node cli/classify.js a,b --document document.txt --json`. + +The exact whole-document token count sets the reservation and $0.084/M charge. +Source text is stored privately while queued and deleted as it is screened. +Final Jev reads at most 20,000 selected evidence tokens. Completion, failure, +cancellation and 24-hour expiry delete remaining source and evidence; failed +or canceled work is refunded. Results remain available until expiry. ## Analytics diff --git a/cli/README.md b/cli/README.md index 74773c8..d7d6697 100644 --- a/cli/README.md +++ b/cli/README.md @@ -23,6 +23,19 @@ One line per input, in input order: `label`, `confidence`, `text`, tab-separated everything the API returns; `--quiet` gives labels only; `--count` gives a histogram. +## One whole document + +With a funded workspace key in `CLASSIFY_API_KEY`, upload a UTF-8 document of +up to 10 million tokens (100 MB): + + classify renewal,cancellation --document agreement.txt --json + +The CLI streams the file once and waits for the result. The server handles +chunking, screening, retries and final judgment. The price is $0.084 per +million original tokens; 10M tokens costs $0.84. Failed jobs are refunded. +`CLASSIFY_JOB_TIMEOUT` controls how long to poll (seconds, default 86400). +Interrupting the CLI does not cancel accepted work; stderr includes its job ID. + ## Built for agents The confidence is calibrated (on a six-way emotion set, answers at ≥ 0.9 were @@ -50,6 +63,7 @@ resumes; a daily quota stops immediately. Errors go to stderr with exit code 1. -k, --max at most n labels (implies --multi) -s, --smart re-ask uncertain answers of a reasoning model -i, --instructions extra criteria + --document one whole UTF-8 document (funded workspace) -r, --review print only inputs with confidence below t -c, --count label histogram instead of rows -j, --json NDJSON output diff --git a/cli/classify.js b/cli/classify.js index f376cec..c1472aa 100755 --- a/cli/classify.js +++ b/cli/classify.js @@ -5,7 +5,8 @@ // tab-separated and greppable, or NDJSON with --json. Errors go to stderr with // exit 1; nothing else ever does. -import { readFileSync, writeFileSync, mkdirSync, statSync, realpathSync } from "node:fs"; +import { createReadStream, readFileSync, writeFileSync, mkdirSync, statSync, realpathSync } from "node:fs"; +import { randomUUID } from "node:crypto"; import { homedir } from "node:os"; import { join } from "node:path"; import { fileURLToPath } from "node:url"; @@ -23,6 +24,7 @@ const HELP = `classify ${PKG.version} — sort text into your own labels, with a USAGE classify "" one text, prints the label + classify --document file.txt one whole document, uploads and waits classify < items.txt one line per input, in order cat items.jsonl | classify --field title @@ -42,9 +44,10 @@ OPTIONS -k, --max at most n labels (implies --multi) -s, --smart re-ask uncertain answers of a reasoning model (slower) --model laya opt into the automatically routed Laya trial (default: jev) - --model chunklaya the long-document model; also automatic past 32,000 characters + --model chunklaya opt into the legacy long-document model --processing bulk Laya bulk lane; default fast is one decision per call -i, --instructions extra criteria: "judge only the service, ignore the food" + --document upload one UTF-8 document, up to 10M tokens (paid workspace) -r, --review print only inputs with confidence below t -c, --count print a label histogram instead of rows -j, --json NDJSON output with every field the API returns @@ -86,7 +89,7 @@ Docs: https://classifier.dev Agent skill: npx skills add https://classifier.de function parseArgs(argv) { const o = { labels: null, text: null, multi: false, max: null, smart: false, instructions: "", review: null, count: false, json: false, quiet: false, field: "text", id: null, - endpoint: ENDPOINT, model: "jev", processing: "fast", apiKey: process.env.CLASSIFY_API_KEY || process.env.CLASSIFIER_API_KEY || "", help: false, version: false }; + endpoint: ENDPOINT, model: "jev", processing: "fast", document: null, apiKey: process.env.CLASSIFY_API_KEY || process.env.CLASSIFIER_API_KEY || "", help: false, version: false }; const positional = []; for (let i = 0; i < argv.length; i++) { const a = argv[i]; @@ -109,6 +112,7 @@ function parseArgs(argv) { else if (a === "--id") o.id = next(); else if (a === "--endpoint") o.endpoint = next(); else if (a === "--api-key") o.apiKey = next(); + else if (a === "--document") o.document = next(); else if (a === "--model") o.model = next(); else if (a === "--processing") o.processing = next(); else if (a.startsWith("--") && a.includes("=")) { argv.splice(i + 1, 0, a.slice(a.indexOf("=") + 1)); argv[i] = a.slice(0, a.indexOf("=")); i--; } @@ -205,6 +209,7 @@ async function post(o, inputs) { continue; } const payload = await res.json().catch(() => ({})); + if (res.status === 202) return checkResults(o, await waitDocument(o, payload), inputs.length); if (res.ok) return checkResults(o, payload, inputs.length); last = payload.error || `HTTP ${res.status}`; // A daily quota cannot recover during a normal CLI run. Keep the API's @@ -224,6 +229,62 @@ async function post(o, inputs) { throw new Error(last); } +async function waitDocument(o, accepted) { + const url = new URL(accepted.status_url, o.endpoint); + if (url.origin !== new URL(o.endpoint).origin) throw new Error("The API returned an invalid job URL."); + if (!process.env.CLASSIFY_NO_PROGRESS) process.stderr.write(`classify: document accepted; waiting for ${accepted.id}\n`); + const deadline = Date.now() + (Number(process.env.CLASSIFY_JOB_TIMEOUT) || 86400) * 1000; + let failures = 0; + while (Date.now() < deadline) { + let response; + try { + response = await fetch(url, { headers: { authorization: `Bearer ${o.apiKey}` }, signal: AbortSignal.timeout(TIMEOUT_MS) }); + } catch { + if (++failures >= 5) throw new Error(`Polling interrupted. Your job continues at ${url}`); + await sleep(2000); continue; + } + const job = await response.json().catch(() => ({})); + if (response.ok && job.status === "finished") return job.result; + if (job.status === "failed") throw new Error(job.error?.message || "Document classification failed; the reservation was refunded."); + if (!response.ok && response.status !== 429 && response.status < 500) + throw new Error(job.error || `Unable to read job: HTTP ${response.status}`); + if (!response.ok && ++failures >= 5) throw new Error(`Polling interrupted. Your job continues at ${url}`); + if (response.ok) failures = 0; + await sleep(Math.min(30, Math.max(1, Number(response.headers.get("retry-after")) || 2)) * 1000); + } + throw new Error(`Stopped waiting. Check your job at ${url}`); +} + +async function classifyDocument(o) { + if (!o.apiKey) throw new Error("--document requires a funded workspace API key (CLASSIFY_API_KEY)."); + if (o.text !== null || o.smart || o.model !== "jev" || o.max !== null) + throw new Error("--document accepts one file with Jev fast, optional --multi and --instructions."); + if (!statSync(o.document).isFile()) throw new Error("--document must name a UTF-8 text file."); + const url = new URL(o.endpoint); + for (const label of o.labels) url.searchParams.append("label", label); + if (o.instructions) url.searchParams.set("instructions", o.instructions); + if (o.multi) url.searchParams.set("multi", "true"); + const id = randomUUID(); + for (let attempt = 0; attempt < ATTEMPTS; attempt++) { + const stream = createReadStream(o.document); + try { + const response = await fetch(url, { method: "POST", duplex: "half", body: stream, + headers: { authorization: `Bearer ${o.apiKey}`, "content-type": "text/plain; charset=utf-8", + "idempotency-key": id, prefer: "respond-async" }, signal: AbortSignal.timeout(TIMEOUT_MS) }); + const payload = await response.json().catch(() => ({})); + if (response.status === 202) return await waitDocument(o, payload); + if (response.status < 500 && response.status !== 409 && response.status !== 429) + throw new Error(payload.error || `Document upload failed: HTTP ${response.status}`); + if (attempt + 1 === ATTEMPTS) throw new Error(payload.error || "Document upload is unavailable."); + } catch (error) { + if (error instanceof TypeError || error.name === "TimeoutError") { + if (attempt + 1 === ATTEMPTS) throw new Error(`Upload interrupted. Retry with Idempotency-Key ${id}, or check /v1/long-context/jobs/${id}/status.`); + } else throw error; + } finally { stream.destroy(); } + await sleep(1000 * 2 ** attempt); + } +} + const ATTEMPTS = 5; let hinted = false; @@ -439,6 +500,16 @@ export async function main(argv) { if (o.labels.length > 100) fail("at most 100 labels"); try { batchSize(o); } catch (e) { fail(e.message); } + if (o.document) { + const payload = await classifyDocument(o); + const results = checkResults(o, payload, 1); + if (o.count) process.stdout.write(formatCount(o, results).join("\n") + "\n"); + else if (kept(o, results[0])) process.stdout.write((o.json + ? JSON.stringify({ file: o.document, ...results[0], usage: payload.usage, pricing: payload.pricing }) + : labelOf(o, results[0])) + "\n"); + return 0; + } + const hint = updateHint(); let items; diff --git a/e2e/document-admission.ts b/e2e/document-admission.ts new file mode 100644 index 0000000..20bdde7 --- /dev/null +++ b/e2e/document-admission.ts @@ -0,0 +1,60 @@ +import assert from "node:assert/strict"; +import { SQL } from "bun"; +import { mkdir, readdir, readFile, writeFile } from "node:fs/promises"; +import { postgresDatabase, hashToken, type AppEnv } from "../src/server/db"; +import { authorizeAndReserve } from "../src/server/usage"; +import { refundTokenReservation } from "../src/server/token-ledger"; + +// Failure contract: different concurrent document IDs must not overdraw; the +// same admission ID must debit once even when recovery reads race. Refunds +// restore one hold. Real PostgreSQL connections, migrations and production SQL. +assert.ok(process.env.POSTGRES_TEST_URL, "Set POSTGRES_TEST_URL to a local test database."); +const sql = new SQL(process.env.POSTGRES_TEST_URL!, { max: 20 }); +const schema = `document_admission_${crypto.randomUUID().replaceAll("-", "")}`; +const token = "classifier_agent_document_admission_fixture"; +const env = { APP_ACCOUNTS_ENABLED: "true", APP_DB: postgresDatabase(async (queries, transaction) => + sql.begin(transaction ? "ISOLATION LEVEL SERIALIZABLE" : "ISOLATION LEVEL READ COMMITTED", async connection => { + await connection.unsafe(`SET LOCAL search_path TO ${schema}`); + const results = []; + for (const query of queries) { + const rows = await connection.unsafe(query.sql, query.params); + results.push({ results: Array.from(rows), meta: { changes: rows.count } }); + } + return results; + })) } as AppEnv; +const report: unknown[] = []; +try { + await sql.unsafe(`CREATE SCHEMA ${schema}`); + await sql.begin(async connection => { + await connection.unsafe(`SET LOCAL search_path TO ${schema}`); + for (const file of (await readdir("migrations/postgres")).filter(p => p.endsWith(".sql")).sort()) + await connection.unsafe(await readFile(`migrations/postgres/${file}`, "utf8")).simple(); + }); + await env.APP_DB.prepare("INSERT INTO app_accounts(id,email,name,balance,paid_balance,reset_at,created_at,period_start) VALUES('fixture','fixture@example.com','Fixture',100,100,now()+interval '1 day',now(),now())").run(); + await env.APP_DB.prepare("INSERT INTO app_agents(id,account_id,name,client,status,token_hash,prefix,created_at) VALUES('fixture','fixture','Fixture','API','connected',?,'fixture',now())") + .bind(await hashToken(token)).run(); + const reserve = (id: string) => authorizeAndReserve(new Request("http://fixture/", { + headers: { authorization: `Bearer ${token}` }, + }), env, 60, 1, { meteringMode: "tokens", reservationId: id }); + for (const sameId of [false, true]) { + const shared = crypto.randomUUID(); + const attempts = await Promise.allSettled(Array.from({ length: 20 }, () => reserve(sameId ? shared : crypto.randomUUID()))); + const accepted = attempts.flatMap(a => a.status === "fulfilled" && a.value ? [a.value] : []); + assert.equal(new Set(accepted.map(a => a.id)).size, 1); + if (sameId) assert.equal(accepted.length, 20, JSON.stringify(attempts.filter(a => a.status === "rejected").map(a => String(a.reason)))); + else assert.equal(accepted.length, 1); + const account = await env.APP_DB.prepare("SELECT balance::integer AS balance,paid_balance::integer AS paid FROM app_accounts WHERE id='fixture'").first(); + assert.deepEqual(account, { balance: 40, paid: 40 }); + const agent = await env.APP_DB.prepare("SELECT used::integer AS used FROM app_agents WHERE id='fixture'").first(); + assert.deepEqual(agent, { used: 60 }); + await Promise.all(Array.from({ length: 5 }, () => refundTokenReservation(env.APP_DB, accepted[0].id))); + assert.deepEqual(await env.APP_DB.prepare("SELECT balance::integer AS balance FROM app_accounts WHERE id='fixture'").first(), { balance: 100 }); + report.push({ sameId, attempts: 20, accepted: accepted.length, distinctHolds: 1, account, agent, refundedBalance: 100 }); + } + await mkdir("captures", { recursive: true }); + await writeFile("captures/document-admission.json", JSON.stringify({ runtime: "PostgreSQL concurrent connections", results: report }, null, 2)); + console.log("Document admission concurrency passed; captures/document-admission.json"); +} finally { + await sql.unsafe(`DROP SCHEMA IF EXISTS ${schema} CASCADE`); + await sql.close(); +} diff --git a/e2e/long-context-billing.ts b/e2e/long-context-billing.ts index ece6045..a254f46 100644 --- a/e2e/long-context-billing.ts +++ b/e2e/long-context-billing.ts @@ -265,7 +265,7 @@ try { [{ inputs: [document, ...Array(20).fill("short")], labels: body.labels }, "too_many_inputs"], [{ input: document, dimensions: Object.fromEntries(Array.from({ length: 33 }, (_, i) => [`field${i}`, ["yes", "no"]])) }, "too_many_decisions"], [{ ...body, tier: "smart" }, "bad_tier"], - [{ input: " x".repeat(250001), labels: body.labels }, "long_context_too_large"], + [{ inputs: [" x".repeat(125001), " x".repeat(125001)], labels: body.labels }, "long_context_too_large"], [{ inputs: [document, ""], labels: body.labels }, "long_context_input"], ]; for (const [invalid, code] of limits) { diff --git a/e2e/whole-document.live.mjs b/e2e/whole-document.live.mjs new file mode 100644 index 0000000..bd3da83 --- /dev/null +++ b/e2e/whole-document.live.mjs @@ -0,0 +1,57 @@ +import assert from 'node:assert/strict'; +import { randomUUID } from 'node:crypto'; +import { mkdir, writeFile } from 'node:fs/promises'; +import { countTokens } from 'gpt-tokenizer/encoding/cl100k_base'; + +// Opt-in paid smoke: real upload, alarms, Jev and settlement. Synthetic filler +// stresses document size, not classification accuracy. Never writes the key. +const key = process.env.CLASSIFIER_API_KEY ?? process.env.CLASSIFY_API_KEY; +assert.ok(key, 'Set CLASSIFIER_API_KEY to a funded workspace key.'); +const base = process.env.CLASSIFIER_BASE_URL ?? 'https://classifier.dev'; +const tokens = process.argv.includes('--full') ? 10_000_000 : 600_000; +const id = process.env.DOCUMENT_JOB_ID ?? randomUUID(); +const headers = { authorization: `Bearer ${key}` }; +const prefix = 'The agreement automatically renews for another year. No cancellation was sent. The remaining text is archival filler:'; +const started = Date.now(); +if (!process.env.DOCUMENT_JOB_ID) { + let left = tokens - countTokens(prefix), phase = 0; + const encoder = new TextEncoder(); + const body = new ReadableStream({ pull(controller) { + if (phase++ === 0) controller.enqueue(encoder.encode(JSON.stringify({ labels: ['renewal', 'cancellation'], + instructions: 'Classify the agreement outcome; archival filler is irrelevant.', input: prefix }).slice(0, -2))); + else if (left) { const n = Math.min(32768, left); left -= n; controller.enqueue(encoder.encode(' x'.repeat(n))); } + else { controller.enqueue(encoder.encode('"}')); controller.close(); } + } }); + const response = await fetch(new URL('/v1/classify', base), { method: 'POST', duplex: 'half', body, + headers: { ...headers, 'content-type': 'application/json', 'idempotency-key': id }, signal: AbortSignal.timeout(300000) }); + const accepted = await response.json(); + console.log(JSON.stringify({ uploadStatus: response.status, ...accepted })); + assert.equal(response.status, 202); + assert.equal(accepted.context_tokens, tokens); +} +const url = new URL(`/v1/long-context/jobs/${id}/status`, base); +console.log(`Resume with DOCUMENT_JOB_ID=${id}`); +let previous = -1; +for (;;) { + const response = await fetch(url, { headers, signal: AbortSignal.timeout(30000) }); + assert.equal(response.status, 200); + const job = await response.json(); + if (job.processed_tokens !== previous) { + previous = job.processed_tokens; + console.log(JSON.stringify({ id, status: job.status, processedTokens: previous, tokens: job.context_tokens })); + } + if (job.status === 'finished' || job.status === 'failed') { + await mkdir('captures', { recursive: true }); + await writeFile(`captures/whole-document-live-${id}.json`, JSON.stringify({ endpoint: base, synthetic: true, + elapsedMs: Date.now() - started, job }, null, 2)); + assert.equal(job.status, 'finished', JSON.stringify(job.error)); + assert.equal(job.result.results[0].label, 'renewal'); + assert.equal(job.result.pricing.input_tokens, tokens); + assert.equal(job.result.pricing.total_usd, tokens * 84 / 1e9); + assert.equal(job.result.usage.long_context.screened_chunks, job.result.usage.long_context.chunks); + console.log(JSON.stringify({ id, result: job.result })); + break; + } + assert.ok(Date.now() - started < 3600000, `Still processing: resume ${id}`); + await new Promise(resolve => setTimeout(resolve, 10000)); +} diff --git a/e2e/whole-document.mjs b/e2e/whole-document.mjs new file mode 100644 index 0000000..634f302 --- /dev/null +++ b/e2e/whole-document.mjs @@ -0,0 +1,269 @@ +import assert from 'node:assert/strict'; +import { createHash, randomUUID } from 'node:crypto'; +import { execFile } from 'node:child_process'; +import { promisify } from 'node:util'; +import { mkdir, readFile, readdir, writeFile } from 'node:fs/promises'; +import { PGlite } from '@electric-sql/pglite'; +import { Miniflare, convertV4MiniflareOptions, Response as WorkerResponse } from 'miniflare'; +import { countTokens } from 'gpt-tokenizer/encoding/cl100k_base'; + +// Failure contract: one upload must not buffer the document, lose Unicode or +// escaped JSON, change token counts at transport boundaries, require client +// chunking/finish calls, double-charge retries, or lose work after a restart. +// Reject unfunded, malformed, empty and oversized inputs before inference; +// protect ownership and refund cancelled/no-evidence jobs. Real workerd, SQLite +// Durable Objects, alarms and PostgreSQL ledger; only upstream HTTP is a fixture. +const report = { runtime: 'built Worker + SQLite Durable Objects/alarms + PostgreSQL ledger', results: [] }; +const pg = new PGlite(); +for (const file of (await readdir('migrations/postgres')).filter(p => p.endsWith('.sql')).sort()) + await pg.exec(await readFile(`migrations/postgres/${file}`, 'utf8')); +const key = 'classifier_agent_whole_document_fixture'; +const freeKey = 'classifier_agent_free_document_fixture'; +for (const [id, token, paid] of [['whole', key, 500000], ['free', freeKey, 0]]) { + await pg.query(`INSERT INTO app_accounts(id,email,name,balance,paid_balance,reset_at,created_at,period_start) + VALUES($1,$2,'Fixture',500000,$3,$4,$5,$5)`, [id, `${id}@example.com`, paid, + new Date(Date.now() + 30 * 86400000).toISOString(), new Date().toISOString()]); + await pg.query(`INSERT INTO app_agents(id,account_id,name,client,status,credit_limit,token_hash,prefix,created_at) + VALUES($1,$1,'Fixture','test','connected',500000,$2,'fixture',$3)`, + [id, createHash('sha256').update(token).digest('hex'), new Date().toISOString()]); +} +let calls = 0, screenedCharacters = 0, mode = 'normal', failScreen = 0, finalEvidence = ''; +let loseAdmission = false, loseRefund = false, staleAdmissionRead = false; +const modules = (await readdir('dist/server', { recursive: true })).filter(p => /\.(js|wasm)$/.test(p)) + .sort((a, b) => a === 'index.js' ? -1 : b === 'index.js' ? 1 : a.localeCompare(b)) + .map(p => ({ type: p.endsWith('.wasm') ? 'CompiledWasm' : 'ESModule', path: `dist/server/${p}` })); +const options = convertV4MiniflareOptions({ workers: [{ name: 'whole-document', modules, modulesRoot: 'dist/server', + compatibilityDate: '2026-08-01', compatibilityFlags: ['nodejs_compat'], + durableObjects: { LONG_CONTEXT_JOBS: { className: 'LongContextJob', useSQLite: true } }, kvNamespaces: ['STATS'], + bindings: { DATABASE_URL: 'postgres://fixture:fixture@fixture.neon.tech/fixture', APP_ACCOUNTS_ENABLED: 'true', + SPENDING_ENABLED: 'true', TYPESAFE_API_KEY: 'fixture', AI_GATEWAY_DISABLED: 'true' }, + outboundService: async request => { + const body = await request.json(); + if (request.url.includes('neon.tech')) { + if ((body.queries ?? [body]).some(q => q.query.includes('WITH admitted AS'))) + assert.equal(request.headers.get('neon-batch-isolation-level'), 'Serializable'); + const execute = async tx => { + const results = []; + for (const q of body.queries ?? [body]) { + const r = await tx.query(q.query, q.params); + results.push({ fields: r.fields, rowCount: r.affectedRows ?? r.rows.length, + rows: r.rows.map(row => r.fields.map(field => { + const value = row[field.name]; + return value === null ? null : typeof value === 'boolean' ? value ? 't' : 'f' : String(value); + })) }); + } + return results; + }; + try { + const results = await pg.transaction(execute); + if (staleAdmissionRead && body.query?.startsWith('SELECT u.account_id,u.agent_id')) { + staleAdmissionRead = false; results[0].rows = []; results[0].rowCount = 0; + } + if ((loseAdmission && (body.queries ?? [body]).some(q => q.query.includes('INSERT INTO app_usage'))) || + (loseRefund && body.query?.includes('refund_token_reservation'))) { + loseAdmission = false; loseRefund = false; + return WorkerResponse.json({ message: 'Fixture: response lost after commit' }, { status: 503 }); + } + return WorkerResponse.json(body.queries ? { results } : results[0]); + } catch (error) { return WorkerResponse.json({ message: error.message, code: error.code }, { status: 400 }); } + } + assert.ok(request.url.includes('typesafe.ai'), request.url); + calls++; + const screening = Object.values(body.questions).some(q => Object.hasOwn(q.criteria ?? {}, 'irrelevant')); + if (screening && failScreen > 0) { failScreen--; return WorkerResponse.json({ detail: { error_type: 'invalid_request' } }, { status: 400 }); } + if (screening && mode === 'slow') await new Promise(resolve => setTimeout(resolve, 200)); + if (screening) screenedCharacters += body.state.reduce((sum, state) => sum + state.text.length, 0); + else finalEvidence = body.state.map(state => state.text).join(''); + return WorkerResponse.json({ model: 'jev-1.13.0', usage: { input_tokens: 1000, output_tokens: 0 }, + answers: Object.fromEntries(Object.entries(body.questions).map(([id, q]) => { + const state = body.state.find(state => state.id === id) ?? body.state[0]; + const labels = Object.keys(q.criteria); + const choice = screening ? mode === 'none' ? 'irrelevant' : state.text.includes('DECISIVE') ? 'relevant' : 'irrelevant' : labels[0]; + return [id, { choice, confidence: 0.99, + probabilities: Object.fromEntries(labels.map(label => [label, label === choice ? 0.99 : 0.005])) }]; + })) }); + }, +}] }); +const mf = new Miniflare(options); +const headers = { authorization: `Bearer ${key}`, 'content-type': 'application/json', prefer: 'respond-async' }; +const upload = (body, extra = {}) => mf.dispatchFetch('https://classifier.dev/v1/classify', { + method: 'POST', duplex: 'half', headers: { ...headers, ...extra }, body: typeof body === 'string' || body instanceof ReadableStream ? body : JSON.stringify(body), +}); +async function waitForJob(url, failure = false) { + for (let i = 0; i < 1800; i++) { + const response = await mf.dispatchFetch(url, { headers }); + assert.equal(response.status, 200, await response.clone().text()); + const job = await response.json(); + if (job.status === 'finished') return job; + if (failure && job.status === 'failed') return job; + assert.notEqual(job.status, 'failed', JSON.stringify(job)); + await new Promise(resolve => setTimeout(resolve, 100)); + } + throw new Error('Job did not finish within three minutes.'); +} +try { + await mf.ready; + const text = 'DECISIVE START: This agreement renews.\n\n' + 'Meeting notes, café, 日本語 😀.\n\n'.repeat(6000) + 'DECISIVE END: Renewal is confirmed.'; + const payload = { input: text, labels: ['renewal', 'cancellation'], instructions: 'Does the agreement renew?' }; + const accepted = await upload(payload); + assert.equal(accepted.status, 202, await accepted.clone().text()); + const job = await accepted.json(); + assert.ok(job.id && job.status_url); + const done = await waitForJob(job.status_url); + assert.equal(done.result.results[0].label, 'renewal'); + assert.equal(done.result.usage.long_context.context_tokens, countTokens(text)); + assert.equal(screenedCharacters, text.length); + assert.ok(finalEvidence.includes('DECISIVE START') && finalEvidence.includes('DECISIVE END')); + report.results.push({ name: 'one whole JSON document, automatic processing and exact token count', result: done.result }); + + const repeatCalls = calls; + const duplicate = await upload(payload, { 'idempotency-key': job.id }); + assert.equal(duplicate.status, 202); + assert.equal((await duplicate.json()).id, job.id); + assert.equal(calls, repeatCalls); + assert.equal(Number((await pg.query('SELECT count(*) AS n FROM app_usage')).rows[0].n), 1); + report.results.push({ name: 'retrying a lost upload response reuses one job and one charge' }); + + const previousCalls = calls; + const free = await upload(payload, { authorization: `Bearer ${freeKey}` }); + assert.equal(free.status, 402); + for (const bad of ['{"input":"unfinished', '{"input":"a","input":"b","labels":["x","y"]}', + JSON.stringify({ input: '', labels: ['x', 'y'] })]) { + const response = await upload(bad); + assert.equal(response.status, 400, await response.clone().text()); + } + assert.equal(calls, previousCalls); + const invalidUtf8 = await mf.dispatchFetch('https://classifier.dev/v1/classify?labels=a,b', { + method: 'POST', headers: { ...headers, 'content-type': 'text/plain' }, body: new Uint8Array([0xff]), + }); + assert.equal(invalidUtf8.status, 400); + const hidden = await mf.dispatchFetch(job.status_url, { headers: { authorization: `Bearer ${freeKey}` } }); + assert.equal(hidden.status, 404); + report.results.push({ name: 'free, malformed, duplicate, empty and cross-account requests rejected', providerCalls: 0 }); + + const escapedText = 'DECISIVE: renewal. "Quoted" \\ \t café 😀 日本語\n'.repeat(4000); + const escapedJson = JSON.stringify({ input: escapedText, labels: ['renewal', 'cancellation'] }) + .replace(/[\u007f-\uffff]/g, c => '\\u' + c.charCodeAt(0).toString(16).padStart(4, '0')); + let position = 0; + const encoder = new TextEncoder(); + const fragmented = new ReadableStream({ pull(controller) { + if (position === escapedJson.length) { controller.close(); return; } + controller.enqueue(encoder.encode(escapedJson.slice(position, position + 997))); + position = Math.min(escapedJson.length, position + 997); + } }); + const escaped = await upload(fragmented); + assert.equal(escaped.status, 202, await escaped.clone().text()); + const escapedDone = await waitForJob((await escaped.json()).status_url); + assert.equal(escapedDone.result.pricing.input_tokens, countTokens(escapedText)); + report.results.push({ name: 'escaped Unicode and JSON across arbitrary transport boundaries', tokens: countTokens(escapedText) }); + + const tail = 'DECISIVE renewal. ' + 'x '.repeat(37000) + ' \t'.repeat(4000); + const tailResponse = await upload({ input: tail, labels: ['renewal', 'cancellation'] }); + assert.equal(tailResponse.status, 202); + const tailDone = await waitForJob((await tailResponse.json()).status_url); + assert.equal(tailDone.result.pricing.input_tokens, countTokens(tail)); + report.results.push({ name: 'whitespace-only internal tail retains exact original token count' }); + + const recoveryId = randomUUID(); + const beforeHolds = Number((await pg.query('SELECT count(*) AS n FROM app_usage')).rows[0].n); + const exactCredits = Math.ceil(countTokens(text) * 84 / 10000); + await pg.query("UPDATE app_accounts SET balance=$1,paid_balance=$1 WHERE id='whole'", [exactCredits]); + loseAdmission = true; + const lost = await upload(payload, { 'idempotency-key': recoveryId }); + assert.equal(lost.status, 503); + staleAdmissionRead = true; + const recovered = await upload(payload, { 'idempotency-key': recoveryId }); + assert.equal(recovered.status, 202, await recovered.clone().text()); + await waitForJob((await recovered.json()).status_url); + assert.equal(Number((await pg.query('SELECT count(*) AS n FROM app_usage')).rows[0].n), beforeHolds + 1); + const recoveredBalance = Number((await pg.query("SELECT balance FROM app_accounts WHERE id='whole'")).rows[0].balance); + assert.ok(recoveredBalance >= 0 && recoveredBalance <= 1, 'Only one debit; settlement may return one rounded credit.'); + await pg.query("UPDATE app_accounts SET balance=500000,paid_balance=500000 WHERE id='whole'"); + report.results.push({ name: 'lost reservation response and stale retry read debit once, even with only enough balance for one hold' }); + + mode = 'slow'; + const cancelStart = calls; + const cancelUpload = await upload(payload); + assert.equal(cancelUpload.status, 202); + const cancelJob = await cancelUpload.json(); + while (calls === cancelStart) await new Promise(resolve => setTimeout(resolve, 10)); + const cancellation = await mf.dispatchFetch(cancelJob.status_url.replace(/status$/, 'cancel'), { method: 'POST', headers }); + assert.ok([200, 202].includes(cancellation.status)); + const canceled = await waitForJob(cancelJob.status_url, true); + assert.equal(canceled.status, 'failed'); + assert.match(canceled.error.message, /canceled/); + assert.equal((await pg.query('SELECT status FROM app_usage ORDER BY created_at DESC LIMIT 1')).rows[0].status, 'refunded'); + mode = 'none'; + loseRefund = true; + const emptyEvidence = await upload({ input: 'No meaningful evidence.', labels: ['renewal', 'cancellation'] }); + const emptyJob = await waitForJob((await emptyEvidence.json()).status_url, true); + assert.equal(emptyJob.error.code, 'long_context_no_evidence'); + assert.equal((await pg.query('SELECT status FROM app_usage ORDER BY created_at DESC LIMIT 1')).rows[0].status, 'refunded'); + report.results.push({ name: 'cancellation and no evidence refund, including lost refund acknowledgement' }); + + mode = 'normal'; failScreen = 1; + const restart = await upload(payload); + assert.equal(restart.status, 202); + const restartJob = await restart.json(); + await mf.setOptions({ ...options, workers: options.workers.map(worker => ({ ...worker, + config: { ...worker.config, env: { ...worker.config.env, RESTART_PROBE: { type: 'text', value: 'reloaded' } } } })) }); + const restarted = await waitForJob(restartJob.status_url); + assert.equal(restarted.result.pricing.input_tokens, countTokens(text)); + report.results.push({ name: 'background job survives Worker reload and provider refusal', id: restarted.id }); + + await mkdir('captures', { recursive: true }); + await writeFile('captures/document-fixture.txt', text); + const cli = await promisify(execFile)(process.execPath, ['cli/classify.js', 'renewal,cancellation', + '--document', 'captures/document-fixture.txt', '--json', '--endpoint', String(await mf.ready)], + { env: { ...process.env, CLASSIFY_API_KEY: key, CLASSIFY_NO_UPDATE_CHECK: '1', CLASSIFY_NO_PROGRESS: '1' }, timeout: 30000 }); + const cliResult = JSON.parse(cli.stdout); + assert.equal(cliResult.label, 'renewal'); + assert.equal(cliResult.usage.long_context.context_tokens, countTokens(text)); + report.results.push({ name: 'CLI streams one file and polls automatically', tokens: cliResult.usage.long_context.context_tokens }); + + if (process.argv.includes('--full')) { + // Exactly 10M tokens, sent as one 20MB JSON upload. No caller tokenization, + // part calls or finish request. The HTTP stream uses ordinary 64KB frames. + let remaining = 10_000_000 - countTokens('DECISIVE renewal.'); + const encoder = new TextEncoder(); + let phase = 0; + const body = new ReadableStream({ pull(controller) { + if (phase === 0) { controller.enqueue(encoder.encode('{"labels":["renewal","cancellation"],"input":"DECISIVE renewal.')); phase = 1; } + else if (remaining) { const n = Math.min(32768, remaining); remaining -= n; controller.enqueue(encoder.encode(' x'.repeat(n))); } + else { controller.enqueue(encoder.encode('"}')); controller.close(); } + } }); + const started = Date.now(); + const response = await upload(body, { prefer: '' }); + assert.equal(response.status, 202, await response.clone().text()); + const pending = await response.json(); + const completed = await waitForJob(pending.status_url); + assert.equal(completed.result.usage.long_context.context_tokens, 10_000_000); + assert.equal(completed.result.pricing.total_usd, 0.84); + assert.equal(completed.result.usage.long_context.screened_chunks, completed.result.usage.long_context.chunks); + const ledger = await pg.query("SELECT status,actual_nano::text AS nano FROM app_usage ORDER BY created_at DESC LIMIT 1"); + assert.equal(ledger.rows[0].nano, '840000000'); + assert.equal(ledger.rows[0].status, 'completed'); + report.results.push({ name: '10M tokens in one streamed upload through real workerd and SQLite alarms', + elapsedMs: Date.now() - started, usage: completed.result.usage, ledger: ledger.rows[0] }); + + const beforeOversize = calls; + let left = 10_000_001; + const oversized = new ReadableStream({ pull(controller) { + if (!left) { controller.close(); return; } + const n = Math.min(32768, left); left -= n; + controller.enqueue(encoder.encode(' x'.repeat(n))); + } }); + const rejected = await mf.dispatchFetch('https://classifier.dev/v1/classify?labels=a,b', { + method: 'POST', duplex: 'half', headers: { ...headers, 'content-type': 'text/plain' }, body: oversized, + }); + assert.equal(rejected.status, 413, await rejected.clone().text()); + assert.equal(calls, beforeOversize); + report.results.push({ name: '10M+1 tokens rejected before inference or reservation' }); + } +} finally { + await mkdir('captures', { recursive: true }); + await writeFile('captures/whole-document.json', JSON.stringify(report, null, 2)); + await mf.dispose(); + await pg.close(); +} +console.log(`${report.results.length} whole-document E2E scenarios passed; captures/whole-document.json`); diff --git a/package-lock.json b/package-lock.json index 9db75ea..122033b 100644 --- a/package-lock.json +++ b/package-lock.json @@ -29,6 +29,7 @@ "react": "^19.3.0", "react-dom": "^19.3.0", "recharts": "^3.10.1", + "stream-json": "^3.7.0", "svix": "^2.5.0", "tailwind-merge": "^3.7.0", "tailwindcss": "^4.3.3", @@ -6785,6 +6786,30 @@ "fast-sha256": "^1.3.0" } }, + "node_modules/stream-chain": { + "version": "4.2.5", + "resolved": "https://registry.npmjs.org/stream-chain/-/stream-chain-4.2.5.tgz", + "integrity": "sha512-Wtyq3bNE3ggLR0v2vftqvuhltym3WbZAkZpfIrkr5F/6vpeUmWmwTgXa16zD87gpahwJ/Qulq3zVfUlgIc0J2A==", + "license": "BSD-3-Clause", + "engines": { + "node": ">=22" + }, + "funding": { + "url": "https://github.com/sponsors/uhop" + } + }, + "node_modules/stream-json": { + "version": "3.7.0", + "resolved": "https://registry.npmjs.org/stream-json/-/stream-json-3.7.0.tgz", + "integrity": "sha512-rCSBdcBP/bPk6T8QFcxAj1MSzAuc5i49cYW6IE7sYObNEPccBJIUiL6fU9c9BWt40aJK0siVynLX7lWaW8alXw==", + "license": "BSD-3-Clause", + "dependencies": { + "stream-chain": "^4.2.5" + }, + "funding": { + "url": "https://github.com/sponsors/uhop" + } + }, "node_modules/supports-color": { "version": "10.2.2", "resolved": "https://registry.npmjs.org/supports-color/-/supports-color-10.2.2.tgz", diff --git a/package.json b/package.json index 6b680a6..f5bd835 100644 --- a/package.json +++ b/package.json @@ -25,7 +25,8 @@ "test:e2e:usage": "node e2e/usage.mjs", "test:e2e:long-context": "node e2e/long-context.mjs", "test:e2e:long-context-billing": "bun e2e/long-context-billing.ts", - "test:e2e:long-context-job": "bun e2e/long-context-job.ts --full" + "test:e2e:long-context-job": "bun e2e/long-context-job.ts --full", + "test:e2e:whole-document": "node e2e/whole-document.mjs --full" }, "keywords": [], "author": "Michael Ryaboy", @@ -68,6 +69,7 @@ "react": "^19.3.0", "react-dom": "^19.3.0", "recharts": "^3.10.1", + "stream-json": "^3.7.0", "svix": "^2.5.0", "tailwind-merge": "^3.7.0", "tailwindcss": "^4.3.3", diff --git a/src/docs.ts b/src/docs.ts index 062f573..6220936 100644 --- a/src/docs.ts +++ b/src/docs.ts @@ -1,7 +1,7 @@ import { vsJevText } from "./vsjev"; import { roadmapDoc } from "./newsletter"; import { INPUT_PRICE_PER_MILLION, LONG_CONTEXT_PRICING } from "./lib/classification-pricing"; -import { LONG_CONTEXT_JOB_MAX_TOKENS, LONG_CONTEXT_PART_MAX_TOKENS, LONG_CONTEXT_MAX_PARTS } from "./long-context"; +import { LONG_CONTEXT_JOB_MAX_TOKENS } from "./long-context"; export const SPENDING_LIMITS = `SPENDING LIMITS @@ -14,8 +14,10 @@ export const SPENDING_LIMITS = `SPENDING LIMITS Anonymous proxy traffic requires a funded key. Unfunded workspace keys share the free limits. Funded work uses its workspace balance, outside the shared free budget, with a default $10 maximum request allowance. - Request bodies are limited to 1 MB. A supplied Idempotency-Key prevents - re-execution: repeated keys receive 409, not a cached response. Free keys + Synchronous request bodies are limited to 1 MB; whole-document jobs allow + 100 MB (see TEN-MILLION-TOKEN JOBS). A supplied Idempotency-Key prevents + re-execution: synchronous repeated keys receive 409, not a cached response; + whole-document UUID keys return the same job. Free keys are scoped to the IP and UTC day; workspace keys are scoped to the account.`; export const DOCS = `classifier.dev @@ -172,7 +174,7 @@ LONG DOCUMENTS Anonymous access and free signup credit do not qualify. Only tier: "fast" is supported; tier: "smart" is refused with bad_tier. - Limits per request: 250,000 original context tokens in total, counted with + Synchronous limits: 250,000 original context tokens in total, counted with cl100k_base across all inputs; 20 documents; 32 decisions (documents times dimensions, or documents times labels in multi-label mode); and a 1 MB request body. Instructions allow 4,000 characters and labels 200 each. @@ -223,46 +225,54 @@ LONG DOCUMENTS TEN-MILLION-TOKEN JOBS - Paid workspaces can classify up to ${LONG_CONTEXT_JOB_MAX_TOKENS.toLocaleString("en-US")} original input tokens as one - job. The single-request JSON API keeps its 250,000-token and 1 MB limits. - Jobs accept ordered parts so no request or Worker needs to hold the whole - document. Create a UUID yourself and reuse it if the create response is lost: - - PUT /v1/long-context/jobs/{uuid}/create - {"max_tokens":${LONG_CONTEXT_JOB_MAX_TOKENS},"documents":1,"labels":["renewing","not-renewing"]} - - PUT /v1/long-context/jobs/{uuid}/parts/0 - {"document":0,"text":"first excerpt, including its original spacing"} - - PUT /v1/long-context/jobs/{uuid}/parts/1 - {"document":0,"text":"next excerpt"} - - POST /v1/long-context/jobs/{uuid}/finish - - Use the same Authorization: Bearer classifier_agent_... key on every call. - Each part may contain at most ${LONG_CONTEXT_PART_MAX_TOKENS.toLocaleString("en-US")} cl100k_base tokens and a 1 MB JSON body; - a job accepts at most ${LONG_CONTEXT_MAX_PARTS.toLocaleString("en-US")} parts. - Send document indexes from 0 in nondecreasing order and part numbers from 0 - without gaps. Split at sentence or paragraph boundaries where possible; - concatenate the parts to recover each document exactly. Jobs allow 20 - documents and 32 decisions (documents × labels in multi-label mode). They - accept labels and optional multi/instructions; dimension maps use the - single-request API. Only the fast tier is available. - - max_tokens is a ceiling, not a charge: the workspace holds enough credits - for that ceiling at creation, then settles at $0.084 per million original - input tokens actually uploaded. cl100k_base counts each part separately; - the full source text is not retained. Screening runs on every part. The job - keeps only selected evidence, up to 30,000 tokens per document while open; - final Jev sees at most 20,000 evidence tokens per document. Omitted evidence - is reported in usage.long_context. Final judgment can take many requests - and minutes for a 10M-token job; poll GET .../{uuid}/status to resume. - - POST .../{uuid}/cancel refunds the hold. A job expires after 24 hours; - unfinished jobs are refunded and selected evidence is deleted. Completion - deletes selected evidence, retains only the result until expiry, and charges - the exact original-part token total once. Failed final judgment can be - retried before expiry. No usable evidence returns 422 and refunds the job. + Send the whole document once to POST /v1/classify. A funded workspace can + upload up to ${LONG_CONTEXT_JOB_MAX_TOKENS.toLocaleString("en-US")} original cl100k_base tokens in a 100 MB request. + Large single-document requests automatically return 202 with a status_url; + add Prefer: respond-async to use that flow for a smaller document too. + Splitting, screening, retries and final judgment happen on the server. + + POST /v1/classify + Authorization: Bearer classifier_agent_... + Content-Type: application/json + Prefer: respond-async + + {"input":"","labels":["renewing","not-renewing"]} + + Or upload a UTF-8 text file directly: + + curl 'https://classifier.dev/v1/classify?labels=renewing,not-renewing' \\ + -H "Authorization: Bearer $CLASSIFY_API_KEY" \\ + -H 'Content-Type: text/plain' --data-binary @document.txt + + JSON accepts one input string (or a one-element inputs/items array), labels, + optional instructions and multi, and fast/jev only. Multi-label jobs allow + up to 32 labels. Raw text accepts repeated label query parameters or + comma-separated labels, plus optional instructions and multi=true. + Labels and instructions may appear anywhere in the JSON object. + + Poll the returned status_url with the same workspace key. Status progresses + from queued to processing to finished; finished includes result with the + normal results, usage and pricing fields. A failed job includes an error + and refunds its reservation. Network/provider failures retry automatically. + Use an optional UUID Idempotency-Key to safely retry a lost upload response. + + The repository CLI uploads and waits in one command: + + node cli/classify.js renewing,not-renewing --document document.txt --json + + Tokens are counted exactly as in the original whole document, independent + of network boundaries. The workspace reserves that actual token price after + upload and settles it once on success: $0.084/M, or $0.84 for 10M tokens. + Source text is temporarily stored privately while queued, then deleted as + it is screened. Final Jev reads up to 20,000 selected evidence tokens; + usage.long_context discloses any omitted eligible chunks. + + Cancel unfinished work and refund its hold with + POST /v1/long-context/jobs/{id}/cancel. + + Jobs expire after 24 hours. Source and selected evidence are + deleted on completion, cancellation, failure or expiry; results remain + available until expiry. Signup credit alone does not enable this feature. LEGACY CHUNKLAYA diff --git a/src/document-upload.ts b/src/document-upload.ts new file mode 100644 index 0000000..7fe801f --- /dev/null +++ b/src/document-upload.ts @@ -0,0 +1,178 @@ +import { parser, type Token } from "stream-json/web/parser.js"; +import { CL100K_TOKEN_SPLIT_REGEX } from "gpt-tokenizer/encodingParams/constants"; +import { countContextTokens, LONG_CONTEXT_JOB_MAX_TOKENS, LONG_CONTEXT_MAX_RUN_CHARS, LongContextError } from "./long-context"; + +export const DOCUMENT_UPLOAD_MAX_BYTES = 100_000_000; +const invalid = (message: string) => new LongContextError(message, 400, "long_context_input"); +const tooLarge = () => new LongContextError("A document supports at most 10,000,000 tokens and a 100 MB upload.", 413, "long_context_too_large"); +type Save = (index: number, text: string, tokens: number) => Promise; + +/** BPE merges stay inside cl100k pre-tokenizer matches. Ending on a complete, + * non-whitespace match also prevents its end-of-input whitespace rule changing + * the count. Thus stored fragments sum to the whole document's exact tokens. */ +class DocumentBuffer { + private text = ""; + private runLength = 0; + private whitespace = false; + private nonempty = false; + tokens = 0; + parts = 0; + constructor(private readonly save: Save) {} + + async append(value: string) { + for (const [run] of value.matchAll(/\s+|\S+/g)) { + const whitespace = /^\s/.test(run); + this.runLength = whitespace === this.whitespace ? this.runLength + run.length : run.length; + this.whitespace = whitespace; + this.nonempty ||= !whitespace; + if (this.runLength > LONG_CONTEXT_MAX_RUN_CHARS) + throw invalid("Continuous whitespace or non-whitespace runs must not exceed 8192 UTF-16 code units."); + } + this.text += value; + if (this.text.length >= 131_072) await this.flush(false); + } + + private boundary(target: number) { + let end = 0; + for (const match of this.text.matchAll(new RegExp(CL100K_TOKEN_SPLIT_REGEX))) { + const next = match.index + match[0].length; + if (next > target) break; + if (/\S$/.test(match[0])) end = next; + } + if (!end) throw invalid("The document contains an unsupported continuous text run."); + return end; + } + + private async flush(final: boolean) { + while (this.text.length && (final || this.text.length >= 131_072)) { + let end = final && this.text.length < 131_072 ? this.text.length : this.boundary(98_304); + let tokens = countContextTokens(this.text.slice(0, end)); + while (tokens > 40_000) { + end = this.boundary(Math.floor(end * 38_000 / tokens)); + tokens = countContextTokens(this.text.slice(0, end)); + } + if (this.tokens + tokens > LONG_CONTEXT_JOB_MAX_TOKENS) throw tooLarge(); + await this.save(this.parts++, this.text.slice(0, end), tokens); + this.tokens += tokens; + this.text = this.text.slice(end); + } + } + + async end() { + if (!this.nonempty) throw invalid("Provide a nonempty input document."); + await this.flush(true); + return { tokens: this.tokens, parts: this.parts }; + } +} + +/** Decode bounded network frames, even when a runtime delivers a larger buffer. */ +function documentText(request: Request): ReadableStream { + if (!request.body) throw invalid("Provide an input document."); + if (Number(request.headers.get("content-length")) > DOCUMENT_UPLOAD_MAX_BYTES) throw tooLarge(); + let bytes = 0; + const decoder = new TextDecoder("utf-8", { fatal: true }); + const reader = request.body.getReader(); + let frame = new Uint8Array(0), offset = 0; + return new ReadableStream({ + async pull(controller) { + if (offset === frame.length) { + const next = await reader.read(); + if (next.done) { controller.enqueue(decoder.decode()); controller.close(); return; } + frame = next.value; offset = 0; + bytes += frame.byteLength; + if (bytes > DOCUMENT_UPLOAD_MAX_BYTES) throw tooLarge(); + } + const end = Math.min(offset + 16_384, frame.length); + controller.enqueue(decoder.decode(frame.subarray(offset, end), { stream: true })); + offset = end; + }, + cancel(reason) { return reader.cancel(reason); }, + }); +} + +/** Parse one JSON input string without ever materializing that string. Labels + * and the rubric can appear before or after it, just as with ordinary JSON. */ +export async function receiveDocument(request: Request, save: Save) { + const document = new DocumentBuffer(save); + const body: Record = {}; + const source = documentText(request); + const contentType = request.headers.get("content-type")?.split(";")[0].trim(); + if (contentType === "text/plain") { + const query = new URL(request.url).searchParams; + body.labels = query.getAll("label"); + if (!(body.labels as string[]).length) body.labels = (query.get("labels") ?? "").split(","); + if (query.has("instructions")) body.instructions = query.get("instructions"); + if (query.has("multi")) { + if (!["true", "false"].includes(query.get("multi")!)) throw invalid("multi must be true or false."); + body.multi = query.get("multi") === "true"; + } + const reader = source.getReader(); + try { + for (;;) { const { done, value } = await reader.read(); if (done) break; await document.append(value); } + } catch (error) { + if (error instanceof LongContextError) throw error; + throw invalid("Send a valid UTF-8 document."); + } finally { await reader.cancel().catch(() => {}); } + } else { + if (contentType && contentType !== "application/json") throw invalid("Use application/json or text/plain."); + const reader = source.pipeThrough(parser.asWebStream({ packValues: false, streamValues: true })).getReader(); + const seen = new Set(); + let root = false, closed = false, field = "", array = "", string = "", value = "", input = false; + try { + for (;;) { + const next = await reader.read(); + if (next.done) break; + const token = next.value as Token; + switch (token.name) { + case "startObject": + if (root) throw invalid("Provide a JSON object with one input string and labels."); + root = true; break; + case "endObject": closed = true; break; + case "startKey": string = "key"; value = ""; break; + case "endKey": + if (seen.has(value)) throw invalid("Duplicate JSON fields are not allowed."); + if (!["input", "inputs", "items", "labels", "instructions", "multi", "tier", "model"].includes(value)) + throw invalid(`Unsupported document field: ${value}.`); + seen.add(value); field = value; string = ""; break; + case "startArray": + if (array || !["labels", "inputs", "items"].includes(field)) throw invalid("Provide one input document and a labels array."); + array = field; if (field === "labels") body.labels = []; break; + case "endArray": array = ""; field = ""; break; + case "startString": + value = ""; + if (field === "input" || array === "inputs" || array === "items") { + if (input) throw invalid("A document upload accepts exactly one input."); + input = true; string = "input"; + } else if (array === "labels") string = "label"; + else if (["instructions", "tier", "model"].includes(field)) string = field; + else throw invalid("Provide a JSON object with one input string and labels."); + break; + case "stringChunk": + if (string === "input") await document.append(token.value); + else { + value += token.value; + const limit = string === "instructions" ? 4000 : string === "label" ? 200 : 32; + if (value.length > limit) throw invalid("The labels, instructions or field name exceed their limit."); + } + break; + case "endString": + if (string === "label") { + (body.labels as string[]).push(value); + if ((body.labels as string[]).length > 100) throw invalid("Provide at most 100 labels."); + } else if (string !== "input") body[string] = value; + string = ""; if (!array) field = ""; break; + case "trueValue": case "falseValue": + if (field !== "multi" || array) throw invalid("Only multi accepts a boolean."); + body.multi = token.value; field = ""; break; + default: throw invalid("Provide one input string, labels, and optional instructions/multi."); + } + } + if (!root || !closed || !input) throw invalid("Provide a JSON object with one input string and labels."); + if (body.model !== undefined && body.model !== "jev") throw invalid("Document uploads use the Jev model."); + } catch (error) { + if (error instanceof LongContextError) throw error; + throw invalid("Send valid UTF-8 JSON with one input string and labels."); + } finally { await reader.cancel().catch(() => {}); } + } + return { body, ...await document.end() }; +} diff --git a/src/http/account-api.ts b/src/http/account-api.ts index 6c95b24..3b835b8 100644 --- a/src/http/account-api.ts +++ b/src/http/account-api.ts @@ -27,6 +27,6 @@ export async function accountApi(request: Request, env: AppEnv & Env, ctx: Execu const headers = new Headers(response.headers); headers.set("access-control-allow-origin", "*"); const exposed = headers.get("access-control-expose-headers"); - headers.set("access-control-expose-headers", [exposed, "x-request-id", "x-billing-status", "x-billed-input-tokens", "x-smart-escalations", "x-usage-cost-usd"].filter(Boolean).join(", ")); + headers.set("access-control-expose-headers", [exposed, "Location", "Retry-After", "x-request-id", "x-billing-status", "x-billed-input-tokens", "x-smart-escalations", "x-usage-cost-usd"].filter(Boolean).join(", ")); return new Response(response.body, { status: response.status, headers }); } diff --git a/src/http/classification.ts b/src/http/classification.ts index 247ab77..e2e82df 100644 --- a/src/http/classification.ts +++ b/src/http/classification.ts @@ -12,6 +12,7 @@ import { refundTokenReservation, settleTokenReservation } from "../server/token- import { writeAccountAnalytics } from "../server/analytics/write"; import { typeSafeDecisionCount } from "../typesafe-compat"; import { isLongContextRequest } from "../long-context"; +import { documentRequest } from "./document"; const card = parseTokenRateCard(JSON.stringify(rates))!; const background = { waitUntil(promise: Promise) { void promise.catch(() => {}); } } as ExecutionContext; @@ -27,6 +28,11 @@ export async function accountClassification(request: Request, env: AppEnv & Part if (!env.LONG_CONTEXT_JOBS) throw new AppError(503, "Long-context jobs are unavailable."); return env.LONG_CONTEXT_JOBS.get(env.LONG_CONTEXT_JOBS.idFromName(jobPath[1])).fetch(request); } + if (request.method === "POST" && ["/", "/v1/classify"].includes(path)) { + const document = await documentRequest(request, env); + if (document instanceof Response) return document; + request = document; + } if (request.method !== "POST" || !["/", "/v1/classify", "/v1/classify/batch", "/sandbox/classify", "/v1/sandbox/classify", "/v1/systemone"].includes(path)) return null; if (env.SPENDING_ENABLED === "true") return spendingClassification(request, env, source, ctx); const accountId = await requireApiAccount(request, env); diff --git a/src/http/document.ts b/src/http/document.ts new file mode 100644 index 0000000..ebd09c0 --- /dev/null +++ b/src/http/document.ts @@ -0,0 +1,57 @@ +import type { Env } from "../index"; +import { countContextTokens, LONG_CONTEXT_MAX_TOKENS, LONG_CONTEXT_THRESHOLD } from "../long-context"; +import { AppError } from "../server/db"; + +export const JOB_ID_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i; +const JSON_REQUEST_BYTES = 1_000_000; + +/** Keep the ordinary API synchronous. Larger single documents are forwarded + * as a stream; the caller sends the same JSON once and receives a job URL. */ +export async function documentRequest(request: Request, env: Partial): Promise { + let asynchronous = /(?:^|,)\s*respond-async(?:\s|,|$)/i.test(request.headers.get("prefer") ?? "") || + request.headers.get("content-type")?.startsWith("text/plain") || + Number(request.headers.get("content-length")) > JSON_REQUEST_BYTES; + if (!asynchronous && request.body) { + const reader = request.body.getReader(); + const prefix: Uint8Array[] = []; + let size = 0, complete = false; + while (size <= JSON_REQUEST_BYTES) { + const next = await reader.read(); + if (next.done) { complete = true; break; } + prefix.push(next.value); size += next.value.byteLength; + } + if (complete) { + const bytes = new Uint8Array(size); + let offset = 0; + for (const chunk of prefix) { bytes.set(chunk, offset); offset += chunk.byteLength; } + try { + const body = JSON.parse(new TextDecoder().decode(bytes)); + const inputs = body.input ?? body.inputs ?? body.items; + const input = typeof inputs === "string" ? inputs : Array.isArray(inputs) && inputs.length === 1 ? inputs[0] : null; + if ((body.model === undefined || body.model === "jev") && typeof input === "string" && input.length > LONG_CONTEXT_THRESHOLD) + asynchronous = countContextTokens(input) > LONG_CONTEXT_MAX_TOKENS; + } catch { /* Ordinary validation owns malformed small requests. */ } + request = new Request(request, { body: bytes }); + } else { + asynchronous = true; + let index = 0; + const stream = new ReadableStream({ + async pull(controller) { + if (index < prefix.length) { controller.enqueue(prefix[index++]); return; } + const next = await reader.read(); + if (next.done) controller.close(); else controller.enqueue(next.value); + }, + cancel(reason) { return reader.cancel(reason); }, + }); + request = new Request(request, { body: stream }); + } + } + if (!asynchronous) return request; + if (!env.LONG_CONTEXT_JOBS) throw new AppError(503, "Long-context jobs are unavailable."); + const supplied = request.headers.get("idempotency-key"); + if (supplied && !JOB_ID_PATTERN.test(supplied)) throw new AppError(400, "Use a UUID for the document Idempotency-Key."); + const id = supplied?.toLowerCase() ?? crypto.randomUUID(); + const url = new URL(request.url); + url.pathname = `/v1/long-context/jobs/${id}/upload`; + return env.LONG_CONTEXT_JOBS.get(env.LONG_CONTEXT_JOBS.idFromName(id)).fetch(new Request(url, request)); +} diff --git a/src/index.ts b/src/index.ts index 2e2d974..c9406a4 100644 --- a/src/index.ts +++ b/src/index.ts @@ -215,7 +215,7 @@ export function primaryModels(): string[] { const CORS = { "access-control-allow-origin": "*", "access-control-allow-methods": "GET, POST, PUT, OPTIONS", - "access-control-allow-headers": "Content-Type, Authorization, Accept, Idempotency-Key, If-None-Match, Mcp-Session-Id, MCP-Protocol-Version, X-TypeSafe-SDK, X-TypeSafe-Runtime, X-TypeSafe-Retry-Count", + "access-control-allow-headers": "Content-Type, Authorization, Accept, Prefer, Idempotency-Key, If-None-Match, Mcp-Session-Id, MCP-Protocol-Version, X-TypeSafe-SDK, X-TypeSafe-Runtime, X-TypeSafe-Retry-Count", "access-control-expose-headers": "RateLimit-Limit, RateLimit-Remaining, RateLimit-Policy, Retry-After, Retry-After-Ms, x-api-version, x-typesafe-request-id, Idempotency-Key", }; diff --git a/src/long-context-job.ts b/src/long-context-job.ts index cf2597e..c012b57 100644 --- a/src/long-context-job.ts +++ b/src/long-context-job.ts @@ -20,6 +20,7 @@ import { boundedRequest } from "./spending"; import { Permit } from "./spending/permit"; import { policy, SpendingError } from "./spending/policy"; import type { Env } from "./index"; +import { receiveDocument } from "./document-upload"; const EVIDENCE_RESERVOIR_TOKENS = 30_000; const JOB_TTL_MS = 24 * 60 * 60 * 1000; @@ -27,11 +28,13 @@ const JOB_TTL_MS = 24 * 60 * 60 * 1000; type Candidate = { index: number; text: string; score: number; tokens: number }; type Job = { id: string; accountId: string; agentId: string; reservationId: string; createdAt: number; - status: "open" | "settling" | "finished" | "failed"; + status: "open" | "settling" | "finished" | "failing" | "failed"; maxTokens: number; documents: number; labels: string[]; instructions?: string; multi: boolean; nextPart: number; lastDocument: number; chunkIndexes: number[]; stats: NonNullable; tokens: ModelTokenUsage[]; providerCostUsd: number | null; + automatic?: { parts: number; tokens: number; retries: number }; + error?: { code: string; message: string }; result?: Record; }; @@ -55,7 +58,7 @@ function joinUsage(into: ModelTokenUsage[], rows: ModelTokenUsage[]) { } } -/** One private job stores selected evidence only; full source parts are discarded after screening. */ +/** Private source fragments are deleted as screened; selected evidence survives until final judgment. */ export class LongContextJob implements DurableObject { private busy = false; private readonly env: Env & AppEnv; @@ -68,18 +71,36 @@ export class LongContextJob implements DurableObject { const accountId = await requireApiAccount(request, this.env); const job = await this.state.storage.get("job"); if (job && job.accountId !== accountId) return answer(404, "not_found", "Job not found."); + if (!job) { + const owner = await this.state.storage.get("upload-owner"); + if (owner && owner !== accountId) return answer(404, "not_found", "Job not found."); + } const path = new URL(request.url).pathname; - if (request.method === "PUT" && path.endsWith("/create")) return this.create(request, accountId, job); + if (request.method === "POST" && path.endsWith("/upload")) return await this.upload(request, accountId, job); + if (request.method === "PUT" && path.endsWith("/create")) return await this.create(request, accountId, job); if (!job) return answer(404, "not_found", "Job not found."); if (request.method === "GET" && path.endsWith("/status")) - return Response.json({ id: job.id, status: job.status, parts: job.nextPart, context_tokens: job.stats.contextTokens, - max_tokens: job.maxTokens, ...(job.status === "finished" ? { result: job.result } : {}) }); + return Response.json({ id: job.id, status: job.status === "open" && await this.state.storage.get("cancel-requested") + ? "canceling" : this.status(job), parts: job.nextPart, + context_tokens: job.automatic?.tokens ?? job.stats.contextTokens, processed_tokens: job.stats.contextTokens, + max_tokens: job.maxTokens, ...(job.error ? { error: job.error } : {}), + ...(job.status === "finished" ? { result: job.result } : {}) }, { headers: { "cache-control": "no-store" } }); if (request.method === "PUT" && /\/parts\/\d+$/.test(path)) - return this.part(request, job, Number(path.split("/").at(-1))); - if (request.method === "POST" && path.endsWith("/finish")) return this.finish(job); + return job.automatic ? answer(409, "job_closed", "The uploaded document is processed automatically.") + : await this.part(request, job, Number(path.split("/").at(-1))); + if (request.method === "POST" && path.endsWith("/finish")) + return job.automatic && job.status === "open" ? this.accepted(job, request.url) : await this.finish(job); if (request.method === "POST" && path.endsWith("/cancel")) { - if (job.status === "open") await this.fail(job); - return Response.json({ id: job.id, status: job.status }); + if (this.busy && job.automatic && job.status === "open") { + await this.state.storage.put("cancel-requested", true); + return Response.json({ id: job.id, status: "canceling" }, { status: 202 }); + } + if (this.busy) return answer(409, "job_busy", "The job is processing; retry cancellation shortly."); + this.busy = true; + try { + if (job.status === "open") await this.cancel(job); + return Response.json({ id: job.id, status: job.status }); + } finally { this.busy = false; } } return answer(404, "not_found", "Job route not found."); } catch (error) { @@ -93,9 +114,53 @@ export class LongContextJob implements DurableObject { } } - private async create(request: Request, accountId: string, existing?: Job): Promise { + private status(job: Job) { + return job.automatic && job.status === "open" ? job.nextPart ? "processing" : "queued" : job.status; + } + + private accepted(job: Job, requestUrl: string) { + const statusUrl = new URL(`/v1/long-context/jobs/${job.id}/status`, requestUrl).href; + return Response.json({ id: job.id, status: this.status(job), status_url: statusUrl, + context_tokens: job.automatic?.tokens ?? job.stats.contextTokens, + expires_at: new Date(job.createdAt + JOB_TTL_MS).toISOString() }, + { status: 202, headers: { location: statusUrl, "retry-after": "2", "cache-control": "no-store" } }); + } + + private async upload(request: Request, accountId: string, existing?: Job) { + if (existing) { await request.body?.cancel().catch(() => {}); return this.accepted(existing, request.url); } + if (this.busy) return answer(409, "job_busy", "The document is already uploading."); + this.busy = true; + try { + const owner = await this.state.storage.get("upload-owner"); + if (owner && owner !== accountId) return answer(404, "not_found", "Job not found."); + const funded = await this.env.APP_DB.prepare(`SELECT (paid_balance>0 OR + (billing_plan IN ('pro','max','scale') AND reset_at::timestamptz>now()) OR + EXISTS(SELECT 1 FROM app_usage WHERE id=? AND account_id=? AND status='pending')) AS funded + FROM app_accounts WHERE id=?`).bind(`document:${new URL(request.url).pathname.split("/")[4]}`, accountId, accountId).first<{ funded: boolean }>(); + if (!funded?.funded) return answer(402, "long_context_payment_required", "A funded workspace is required."); + await this.deleteSources(); + await this.state.storage.setAlarm(Date.now() + JOB_TTL_MS); + await this.state.storage.put("upload-owner", accountId); + const uploaded = await receiveDocument(request, async (index, text, tokens) => { + await this.state.storage.put(`source:${index}`, { text, tokens }); + }); + const headers = new Headers(request.headers); + headers.delete("content-length"); + headers.set("content-type", "application/json"); + const response = await this.create(new Request(request.url, { method: "PUT", headers, + body: JSON.stringify({ ...uploaded.body, documents: 1, max_tokens: uploaded.tokens }) }), accountId, undefined, + { parts: uploaded.parts, tokens: uploaded.tokens, retries: 0 }); + if (!response.ok) { await this.deleteSources(); return response; } + return this.accepted((await this.state.storage.get("job"))!, request.url); + } catch (error) { + if (!await this.state.storage.get("job")) await this.deleteSources(); + throw error; + } finally { this.busy = false; } + } + + private async create(request: Request, accountId: string, existing?: Job, automatic?: Job["automatic"]): Promise { if (existing) return Response.json({ id: existing.id, status: existing.status, max_tokens: existing.maxTokens }); - if (this.busy) return answer(409, "job_busy", "The job is processing another request."); + if (this.busy && !automatic) return answer(409, "job_busy", "The job is processing another request."); this.busy = true; try { const body = await (await boundedRequest(request)).json() as Record; @@ -117,26 +182,30 @@ export class LongContextJob implements DurableObject { if (!jevClassificationFits("", SCREEN_LABELS, longContextScreeningInstructions(labels, instructions), false) || !jevClassificationFits("", labels, longContextFinalInstructions(instructions), multi)) return answer(400, "long_context_input", "The labels and instructions exceed Jev's context budget."); + const id = new URL(request.url).pathname.split("/")[4]; + const reservationId = `document:${id}`; const funded = await this.env.APP_DB.prepare(`SELECT (paid_balance>0 OR - (billing_plan IN ('pro','max','scale') AND reset_at::timestamptz>now())) AS funded - FROM app_accounts WHERE id=?`).bind(accountId).first<{ funded: boolean }>(); + (billing_plan IN ('pro','max','scale') AND reset_at::timestamptz>now()) OR + EXISTS(SELECT 1 FROM app_usage WHERE id=? AND account_id=? AND status='pending')) AS funded + FROM app_accounts WHERE id=?`).bind(reservationId, accountId, accountId).first<{ funded: boolean }>(); if (!funded?.funded) return answer(402, "long_context_payment_required", "A funded workspace is required."); const credits = Number((longContextCharge(maxTokens as number).nanodollars + 9999n) / 10000n); + await this.state.storage.setAlarm(Date.now() + JOB_TTL_MS); + await this.state.storage.put("admission", reservationId); const reservation = await authorizeAndReserve(request, this.env, credits, documents as number, - { type: "API · Long context job", meteringMode: "tokens" }); + { type: "API · Long context job", meteringMode: "tokens", reservationId }); if (!reservation) return answer(401, "invalid_api_key", "A workspace API key is required."); - const id = new URL(request.url).pathname.split("/")[4]; const job: Job = { id, accountId, agentId: reservation.agentId, reservationId: reservation.id, createdAt: Date.now(), status: "open", maxTokens: maxTokens as number, documents: documents as number, labels, instructions, multi, nextPart: 0, lastDocument: 0, chunkIndexes: Array(documents as number).fill(0), - tokens: [], providerCostUsd: 0, + tokens: [], providerCostUsd: 0, ...(automatic ? { automatic } : {}), stats: { contextTokens: 0, documents: documents as number, chunks: 0, screenedChunks: 0, eligibleChunks: 0, selectedChunks: 0, omittedChunks: 0, screeningInputTokens: 0, finalInputTokens: 0, screeningCalls: 0, finalCalls: 0, screeningMs: 0, finalMs: 0, tokenizer: "cl100k_base" } }; try { - await this.state.storage.setAlarm(Date.now() + JOB_TTL_MS); + await this.state.storage.setAlarm(Date.now() + (automatic ? 1 : JOB_TTL_MS)); await this.state.storage.put("job", job); } catch (error) { await refundTokenReservation(this.env.APP_DB, reservation.id); @@ -144,15 +213,15 @@ export class LongContextJob implements DurableObject { } return Response.json({ id, status: "open", max_tokens: maxTokens, part_max_tokens: LONG_CONTEXT_PART_MAX_TOKENS, expires_in_seconds: JOB_TTL_MS / 1000, reservation_usd: credits / 100000 }, { status: 201 }); - } finally { this.busy = false; } + } finally { if (!automatic) this.busy = false; } } - private async part(request: Request, job: Job, sequence: number): Promise { + private async part(request: Request, job: Job, sequence: number, lockHeld = false): Promise { if (job.status !== "open") return answer(409, "job_closed", "This job is closed."); - if (this.busy) return answer(409, "job_busy", "The job is processing another request."); + if (this.busy && !lockHeld) return answer(409, "job_busy", "The job is processing another request."); if (!Number.isSafeInteger(sequence) || sequence !== job.nextPart) return answer(409, "part_sequence", `The next part is ${job.nextPart}.`); - if (job.nextPart >= LONG_CONTEXT_MAX_PARTS) + if (!job.automatic && job.nextPart >= LONG_CONTEXT_MAX_PARTS) return answer(400, "long_context_too_large", `A job accepts at most ${LONG_CONTEXT_MAX_PARTS.toLocaleString("en-US")} parts.`); this.busy = true; try { @@ -160,12 +229,12 @@ export class LongContextJob implements DurableObject { const document = body.document; const text = body.text; if (!Number.isSafeInteger(document) || (document as number) < job.lastDocument || - (document as number) >= job.documents || typeof text !== "string" || !text.trim()) + (document as number) >= job.documents || typeof text !== "string" || (!text.trim() && !job.automatic)) return answer(400, "long_context_input", "Provide nonempty text and a document index in source order."); const count = countContextTokens(text as string); if (count > LONG_CONTEXT_PART_MAX_TOKENS || job.stats.contextTokens + count > job.maxTokens) return answer(400, "long_context_too_large", "The part or job exceeds its token limit."); - const chunks = await chunkLongContextPart(text as string); + const chunks = text.trim() ? await chunkLongContextPart(text as string) : []; const screenInstructions = longContextScreeningInstructions(job.labels, job.instructions); if (chunks.some(chunk => !jevClassificationFits(chunk, SCREEN_LABELS, screenInstructions, false))) return answer(400, "long_context_input", "A chunk and its instructions exceed Jev's context budget."); @@ -175,8 +244,8 @@ export class LongContextJob implements DurableObject { const keys = jevKeys(this.env); if (!keys) return answer(503, "long_context_unavailable", "Jev is unavailable."); const started = Date.now(); - const screening = await jevClassify(keys, chunks, SCREEN_LABELS, - screenInstructions, false, meter, SCREENING_BACKEND); + const screening = chunks.length ? await jevClassify(keys, chunks, SCREEN_LABELS, + screenInstructions, false, meter, SCREENING_BACKEND) : []; const existing = await this.state.storage.get(`doc:${document}`) ?? []; const candidates = [...existing]; let eligible = 0; @@ -205,26 +274,28 @@ export class LongContextJob implements DurableObject { job.chunkIndexes[document as number] += chunks.length; job.lastDocument = document as number; job.nextPart++; + if (job.automatic) job.automatic.retries = 0; await this.state.storage.transaction(async transaction => { await transaction.put(`doc:${document}`, candidates); await transaction.put("job", job); + if (job.automatic) await transaction.delete(`source:${sequence}`); }); return Response.json({ id: job.id, status: job.status, next_part: job.nextPart, context_tokens: job.stats.contextTokens, screened_chunks: job.stats.screenedChunks, max_tokens: job.maxTokens }); - } finally { this.busy = false; } + } finally { if (!lockHeld) this.busy = false; } } - private async finish(job: Job): Promise { + private async finish(job: Job, lockHeld = false): Promise { if (job.status === "finished") return Response.json(job.result); if (job.status === "settling") { - if (this.busy) return answer(409, "job_busy", "The job is settling."); + if (this.busy && !lockHeld) return answer(409, "job_busy", "The job is settling."); this.busy = true; try { return await this.settle(job); } - finally { this.busy = false; } + finally { if (!lockHeld) this.busy = false; } } if (job.status !== "open") return answer(409, "job_closed", "This job is closed."); - if (this.busy) return answer(409, "job_busy", "The job is processing another request."); + if (this.busy && !lockHeld) return answer(409, "job_busy", "The job is processing another request."); if (job.stats.contextTokens === 0 || job.chunkIndexes.some(index => index === 0)) return answer(400, "long_context_input", "Every document needs at least one part."); this.busy = true; @@ -245,6 +316,7 @@ export class LongContextJob implements DurableObject { } } if (!selected.length) { + job.error = { code: "long_context_no_evidence", message: "No usable evidence was selected. The reservation was refunded." }; await this.fail(job); return answer(422, "long_context_no_evidence", "No usable evidence was selected for a document."); } @@ -272,6 +344,10 @@ export class LongContextJob implements DurableObject { joinUsage(job.tokens, meter.tokens); job.providerCostUsd = job.providerCostUsd === null || meter.tokens.some(row => row.inputTokens === null) ? null : job.providerCostUsd + meter.usd; + if (job.automatic && await this.state.storage.get("cancel-requested")) { + await this.cancel(job); + return answer(409, "job_closed", "The document was canceled and its reservation refunded."); + } const usageMeter = newMeter(); usageMeter.longContext = job.stats; usageMeter.tokens = job.tokens; @@ -284,21 +360,27 @@ export class LongContextJob implements DurableObject { job.result = payload; job.status = "settling"; await this.state.storage.put("job", job); - return this.settle(job); + return await this.settle(job); } catch (error) { // A provider failure is retryable. The reservation remains held until retry, // cancellation or the expiry alarm; a failed final call never charges. return answer(503, "long_context_unavailable", "Final judgment is temporarily unavailable; retry before job expiry."); - } finally { this.busy = false; } + } finally { if (!lockHeld) this.busy = false; } } private async settle(job: Job): Promise { const charge = longContextCharge(job.stats.contextTokens); const tokens = job.tokens.every(row => row.inputTokens !== null) ? job.tokens.reduce((sum, row) => sum + row.inputTokens!, 0) : null; - await settleTokenReservation(this.env.APP_DB, job.reservationId, charge, + const settlement = await settleTokenReservation(this.env.APP_DB, job.reservationId, charge, { inputTokens: tokens, outputTokens: job.tokens.every(row => row.outputTokens !== null) ? job.tokens.reduce((sum, row) => sum + row.outputTokens!, 0) : null }); + if (settlement.status === "refunded") { + job.error = { code: "job_closed", message: "The reservation was refunded; start a new document job." }; + await this.fail(job); + return answer(409, "job_closed", job.error.message); + } + if (settlement.status !== "completed") throw new Error("Settlement has not completed."); job.status = "finished"; await this.state.storage.put("job", job); for (let doc = 0; doc < job.documents; doc++) @@ -311,14 +393,31 @@ export class LongContextJob implements DurableObject { } private async fail(job: Job) { + await this.state.storage.setAlarm(Date.now() + 2000); + job.status = "failing"; + await this.state.storage.put("job", job); await refundTokenReservation(this.env.APP_DB, job.reservationId); + for (let doc = 0; doc < job.documents; doc++) await this.state.storage.delete(`doc:${doc}`); + if (job.automatic) await this.deleteSources(); job.status = "failed"; await this.state.storage.put("job", job); - for (let doc = 0; doc < job.documents; doc++) await this.state.storage.delete(`doc:${doc}`); recordLongContext(this.env, job.stats, "error"); this.recordAccount(job, "error", 0); } + private async deleteSources() { + for (;;) { + const sources = await this.state.storage.list({ prefix: "source:", limit: 16 }); + if (!sources.size) return; + await this.state.storage.delete([...sources.keys()]); + } + } + + private async cancel(job: Job) { + job.error = { code: "job_closed", message: "The document was canceled and its reservation refunded." }; + await this.fail(job); + } + private recordAccount(job: Job, status: "success" | "error", retailCostUsd: number) { const sum = (field: "inputTokens" | "outputTokens") => job.tokens.every(row => row[field] !== null) ? job.tokens.reduce((total, row) => total + row[field]!, 0) : null; @@ -331,15 +430,79 @@ export class LongContextJob implements DurableObject { } async alarm() { - if (this.busy) { await this.state.storage.setAlarm(Date.now() + 300_000); return; } + if (this.busy) { await this.state.storage.setAlarm(Date.now() + 1000); return; } this.busy = true; try { const job = await this.state.storage.get("job"); + if (!job) { + const admission = await this.state.storage.get("admission"); + if (admission) { + await this.state.storage.setAlarm(Date.now() + 300_000); + if (await this.env.APP_DB.prepare("SELECT id FROM app_usage WHERE id=?").bind(admission).first()) + await refundTokenReservation(this.env.APP_DB, admission); + } + await this.state.storage.deleteAll(); + return; + } + if (job.status === "failing") { + await this.state.storage.setAlarm(Date.now() + 2000); + await this.fail(job); + await this.state.storage.setAlarm(job.createdAt + JOB_TTL_MS); + return; + } + if (job?.automatic && job.status === "open" && Date.now() < job.createdAt + JOB_TTL_MS) { + // Persist a wake-up before work so an eviction cannot strand a paid job. + await this.state.storage.setAlarm(Date.now() + 60_000); + try { + const started = Date.now(); + for (let n = 0; n < 4 && job.nextPart < job.automatic.parts && Date.now() - started < 20_000; n++) { + if (await this.state.storage.get("cancel-requested")) { + await this.cancel(job); + await this.state.storage.setAlarm(job.createdAt + JOB_TTL_MS); + return; + } + const sequence = job.nextPart; + const source = await this.state.storage.get<{ text: string; tokens: number }>(`source:${sequence}`); + if (!source) throw new Error("The uploaded document is unavailable."); + const response = await this.part(new Request("https://job.internal/part", { method: "PUT", + body: JSON.stringify({ document: 0, text: source.text }) }), job, sequence, true); + if (!response.ok) throw new Error("Document screening failed."); + } + if (job.nextPart === job.automatic.parts) { + const response = await this.finish(job, true); + if (!response.ok && response.status !== 422) throw new Error("Document judgment failed."); + } + await this.state.storage.put("job", job); + await this.state.storage.setAlarm(job.status === "open" ? Date.now() + 1 : job.createdAt + JOB_TTL_MS); + } catch { + // Reload durable progress: a response may fail after its transaction + // committed. Never replay a completed source fragment or its billing. + const current = (await this.state.storage.get("job"))!; + if (current.status === "failed" || current.status === "finished") { + await this.state.storage.setAlarm(current.createdAt + JOB_TTL_MS); return; + } + if (current.status === "settling") { await this.state.storage.setAlarm(Date.now() + 2000); return; } + current.automatic!.retries++; + if (current.automatic!.retries >= 5) { + current.error = { code: "long_context_unavailable", message: "Processing failed after retries. The reservation was refunded." }; + await this.fail(current); + await this.state.storage.setAlarm(current.createdAt + JOB_TTL_MS); + } else { + await this.state.storage.put("job", current); + await this.state.storage.setAlarm(Date.now() + 1000 * 2 ** current.automatic!.retries); + } + } + return; + } if (job?.status === "settling") { - try { await this.settle(job); } + try { await this.settle(job); await this.state.storage.setAlarm(Math.max(Date.now() + 1, job.createdAt + JOB_TTL_MS)); } catch { await this.state.storage.setAlarm(Date.now() + 300_000); } return; } + if (job?.automatic && Date.now() < job.createdAt + JOB_TTL_MS) { + await this.state.storage.setAlarm(job.createdAt + JOB_TTL_MS); + return; + } if (job?.status === "open") { try { await refundTokenReservation(this.env.APP_DB, job.reservationId); } catch { await this.state.storage.setAlarm(Date.now() + 300_000); return; } diff --git a/src/openapi.ts b/src/openapi.ts index b95209d..1521964 100644 --- a/src/openapi.ts +++ b/src/openapi.ts @@ -72,10 +72,27 @@ const ACCOUNT_BILLING_HEADERS = { const JOB_ID = { name: "id", in: "path", required: true, schema: { type: "string", format: "uuid" }, description: "A UUID chosen by the caller; reuse it to retry creation or resume the same job." }; const JOB_RESPONSE = { "200": { description: "Job progress or result.", content: { "application/json": { schema: { type: "object" } } } }, + "202": { description: "Cancellation requested; poll status until refunded." }, "400": err("Invalid job metadata or part."), "401": err("A valid workspace API key is required."), "402": err("A funded workspace and sufficient balance are required."), "409": err("The job is closed, busy, or the next part sequence is different."), "503": err("Jev or job processing is unavailable; retry the same part or finish call.") }; +const DOCUMENT_PARAMETERS = [ + { $ref: "#/components/parameters/IdempotencyKey" }, + { name: "Prefer", in: "header", schema: { type: "string", const: "respond-async" }, description: "Upload one complete document and receive a background job, even below the synchronous limits." }, + { name: "labels", in: "query", schema: { type: "string" }, description: "For text/plain uploads: comma-separated labels." }, + { name: "label", in: "query", style: "form", explode: true, schema: { type: "array", items: { type: "string" } }, description: "For text/plain uploads: repeated label parameters, allowing commas inside labels." }, + { name: "instructions", in: "query", schema: { type: "string", maxLength: 4000 }, description: "For text/plain uploads: classification criteria." }, + { name: "multi", in: "query", schema: { type: "boolean" }, description: "For text/plain uploads: return every applicable label." }, +]; +const DOCUMENT_ACCEPTED = { + description: "Whole document uploaded. Poll status_url with the same workspace key. No client chunking or finish call is needed. Requires funded access. Queued source is stored privately and deleted as screened; all source/evidence is deleted on completion, cancellation, failure or 24-hour expiry.", + headers: { Location: { schema: { type: "string" } }, "Retry-After": { schema: { type: "integer" } } }, + content: { "application/json": { schema: { type: "object", required: ["id", "status", "status_url", "context_tokens", "expires_at"], properties: { + id: { type: "string", format: "uuid" }, status: { type: "string" }, status_url: { type: "string", format: "uri" }, + context_tokens: { type: "integer", maximum: LONG_CONTEXT_JOB_MAX_TOKENS }, expires_at: { type: "string", format: "date-time" }, + } } } }, +}; const TYPESAFE_ENTRY = { anyOf: [ { type: "string" }, @@ -166,8 +183,8 @@ export const OPENAPI = { "(RFC 9745 / RFC 8594) for at least six months before removal.\n\n" + "Rate limits: RateLimit-Limit and RateLimit-Policy (IETF draft-ietf-httpapi-ratelimit-headers) on every classification " + "response, RateLimit-Remaining once the limiter has been consulted (every 200 and 429), Retry-After on 429s. " + - "Idempotency: classification has no side effects; an Idempotency-Key header is accepted and " + - "echoed so generic retry logic keeps working.\n\n" + + "Idempotency: synchronous paid requests reject replayed keys. Whole-document jobs accept an optional UUID " + + "Idempotency-Key and return the same job on retry without charging again.\n\n" + "Errors: classification failures return {error, code}; workspace authorization and billing errors return {error}. " + "See components.schemas.Error for classification codes. The two GET forms answer " + "plain text (`error:`, `usage:`, `try:` lines) unless ?verbose=1 or Accept: application/json asks for the JSON object.\n\n" + @@ -606,7 +623,7 @@ export const OPENAPI = { "/v1/long-context/jobs/{id}/create": { put: { operationId: "createLongContextJob", tags: ["classify"], - summary: "Create a paid long-context job (up to 10M tokens)", + summary: "Legacy manual-upload job creation; prefer one POST /v1/classify", description: "Choose a UUID once and reuse it on retries. Creation holds credits for max_tokens at $0.084/M original input tokens; final judgment settles the actual uploaded part-token count. Requires a funded workspace. No source text is stored at creation.", security: [{ accountKey: [] }], parameters: [JOB_ID], requestBody: { required: true, content: { "application/json": { schema: { type: "object", required: ["max_tokens", "documents", "labels"], properties: { @@ -620,7 +637,7 @@ export const OPENAPI = { "/v1/long-context/jobs/{id}/parts/{sequence}": { put: { operationId: "uploadLongContextPart", tags: ["classify"], - summary: "Screen one bounded text part", + summary: "Legacy manual part upload; not needed for whole-document jobs", description: `Upload source text in document order. Sequence starts at zero with no gaps, up to ${LONG_CONTEXT_MAX_PARTS.toLocaleString("en-US")} parts. Each part fits 1 MB JSON and ${LONG_CONTEXT_PART_MAX_TOKENS.toLocaleString("en-US")} cl100k_base tokens. The full source part is discarded after screening; only selected evidence remains until finish, cancel, or 24-hour expiry.`, security: [{ accountKey: [] }], parameters: [JOB_ID, { name: "sequence", in: "path", required: true, schema: { type: "integer", minimum: 0 } }], @@ -631,6 +648,7 @@ export const OPENAPI = { }, "/v1/long-context/jobs/{id}/status": { get: { operationId: "getLongContextJob", tags: ["classify"], summary: "Read job progress or a completed result", + description: "Automatic jobs progress from queued to processing to finished. context_tokens is the full original document; processed_tokens tracks screening. A finished job includes result with results, usage and pricing. A failed job includes error and refunds its hold. Results expire after 24 hours.", security: [{ accountKey: [] }], parameters: [JOB_ID], responses: JOB_RESPONSE, } }, "/v1/long-context/jobs/{id}/finish": { post: { @@ -647,12 +665,12 @@ export const OPENAPI = { post: { operationId: "classifyV1", summary: "Classify texts (v1). Identical to POST /.", - description: "The versioned address of the classification endpoint. Same request body, same response, same limits as POST /.", + description: "Identical to POST /. One complete Jev document supports 10M original cl100k_base tokens and a 100 MB upload. Prefer: respond-async, text/plain, a body above 1 MB or a single Jev input above 250k tokens returns 202 with a background job. The server screens every chunk then runs final Jev over selected evidence. Price: $0.084/M original tokens, reserved after upload; free/signup credit alone cannot enable jobs. Automatic jobs accept input (or one-element inputs/items), labels, instructions, multi, tier:fast and model:jev only. Other batches remain synchronous.", tags: ["classify"], - parameters: [{ $ref: "#/components/parameters/IdempotencyKey" }], + parameters: DOCUMENT_PARAMETERS, requestBody: { required: true, - content: { "application/json": { schema: { $ref: "#/components/schemas/ClassifyRequest" } } }, + content: { "application/json": { schema: { $ref: "#/components/schemas/ClassifyRequest" } }, "text/plain": { schema: { type: "string", description: "The entire UTF-8 document; supply labels in the query." } } }, }, responses: { "200": { @@ -660,6 +678,7 @@ export const OPENAPI = { headers: RATE_LIMIT_HEADERS, content: { "application/json": { schema: { $ref: "#/components/schemas/ClassifyResponse" } } }, }, + "202": DOCUMENT_ACCEPTED, ...ERRORS, }, }, @@ -667,8 +686,8 @@ export const OPENAPI = { "/v1/classify/batch": { post: { operationId: "classifyBatchV1", - summary: "Classify up to 1,000 texts in one request (alias of POST /v1/classify).", - description: "The batch operation is the normal operation: `inputs` takes up to 1,000 texts and the results come back in the same order. This path exists for callers that look for a batch endpoint by name; it behaves exactly like POST /v1/classify.", + summary: "Classify up to 1,000 texts synchronously.", + description: "inputs takes up to 1,000 texts and results come back in the same order. This path retains synchronous limits; for one whole document up to 10M tokens/100 MB, use POST /v1/classify instead.", tags: ["classify"], parameters: [{ $ref: "#/components/parameters/IdempotencyKey" }], requestBody: { @@ -731,9 +750,12 @@ export const OPENAPI = { post: { operationId: "classify", summary: "Classify one or many texts into one of the supplied labels.", + description: "Same synchronous and whole-document behavior as POST /v1/classify. One paid document supports 10M tokens in a 100 MB upload; automatic jobs return 202.", + parameters: DOCUMENT_PARAMETERS, requestBody: { required: true, content: { + "text/plain": { schema: { type: "string", description: "The entire UTF-8 document; supply labels in the query." } }, "application/json": { schema: { $ref: "#/components/schemas/ClassifyRequest" }, examples: { @@ -757,6 +779,7 @@ export const OPENAPI = { }, }, responses: { + "202": DOCUMENT_ACCEPTED, ...ERRORS, "200": { description: "Classification results", @@ -1041,16 +1064,16 @@ 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", "kev", "chunklaya"], description: "Defaults to Jev unless processing implies Laya. Default/explicit jev inputs over 32,000 characters use paid Fast-only Jev long context: 600-token Chonkie chunks, parallel evidence screening and final Jev over whole eligible chunks in source order, bounded by 20,000 cl100k_base tokens and a conservative provider estimate. Eligible evidence may be omitted; usage.long_context discloses selection. Requires paid workspace balance or active paid subscription, not signup credit. Limits per request: 250,000 original cl100k_base context tokens, 20 documents, 32 decisions, 1 MB body. Price: $0.084/M original context tokens summed once across inputs, independent of dimensions and actual inference usage. Explicit chunklaya retains the legacy opt-in (4,000,000 characters/input, 20 inputs, subject to 1 MB body; chunklaya/multilingual results; no Smart). Laya and Kev are experimental Beam models with 512-token and 8,192-token contexts respectively, 2–16 short labels, text ≤2,000 characters and instructions ≤400 characters. Results use jev/laya or jev/kev; Jev calibration claims do not apply." }, + model: { type: "string", enum: ["jev", "laya", "kev", "chunklaya"], description: "Defaults to Jev unless processing implies Laya. Default/explicit jev inputs over 32,000 characters use paid Fast-only Jev long context: 600-token Chonkie chunks, parallel evidence screening and final Jev over whole eligible chunks in source order, bounded by 20,000 cl100k_base tokens and a conservative provider estimate. Eligible evidence may be omitted; usage.long_context discloses selection. Requires paid workspace balance or active paid subscription, not signup credit. Synchronous limits: 250,000 original cl100k_base tokens, 20 documents, 32 decisions, 1 MB body. One whole document on POST / or /v1/classify supports 10M tokens and 100 MB as a background job. Price: $0.084/M original context tokens summed once across inputs, independent of dimensions and actual inference usage. Explicit chunklaya retains the legacy opt-in (4,000,000 characters/input, 20 inputs, subject to 1 MB body; chunklaya/multilingual results; no Smart). Laya and Kev are experimental Beam models with 512-token and 8,192-token contexts respectively, 2–16 short labels, text ≤2,000 characters and instructions ≤400 characters. Results use jev/laya or jev/kev; Jev calibration claims do not apply." }, 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: 4000000 }, description: "Alias for inputs in dimensions mode. Default/jev inputs above 32,000 characters use paid long context: at most 20 documents, 32 decisions and 250,000 original context tokens total, within a 1 MB body. Explicit chunklaya allows up to 4,000,000 characters/input within the body limit. Do not combine with input or inputs." }, - input: { type: "string", description: "A single text. Provide this or inputs; a string under `inputs` is read as one text too." }, + items: { type: "array", minItems: 1, maxItems: 1000, items: { type: "string", minLength: 1 }, description: "Alias for inputs. One-element arrays can use whole-document jobs (10M tokens/100 MB, labels only). Synchronous dimensions mode: Default/jev inputs above 32,000 characters use paid long context: at most 20 documents, 32 decisions and 250,000 original context tokens total, within a 1 MB body. Explicit chunklaya allows up to 4,000,000 characters/input within the body limit. Do not combine with input or inputs." }, + input: { type: "string", description: "One complete document. POST / or /v1/classify automatically returns a background job above synchronous limits: up to 10M original tokens and 100 MB. No manual parts needed. Use Prefer: respond-async to request a job at any size. Provide this or inputs." }, inputs: { type: "array", items: { type: "string" }, maxItems: 1000, - description: "Up to 1,000 texts classified in one call, results in the same order. Long context allows 20 documents, 32 decisions and 250,000 original context tokens total within a 1 MB request. Public smart requests accept at most 200 so the batch fits its per-minute quota.", + description: "Up to 1,000 texts classified in one call, results in the same order. Synchronous long context allows 20 documents, 32 decisions and 250,000 original context tokens total within 1 MB. One-element arrays on POST / or /v1/classify support automatic jobs up to 10M tokens/100 MB. Public smart requests accept at most 200 so the batch fits its per-minute quota.", }, labels: { type: "array", @@ -1348,7 +1371,7 @@ export const OPENAPI = { in: "header", required: false, schema: { type: "string", maxLength: 200 }, - description: "Optional replay protection. A previously admitted key returns 409 without executing again; responses are not cached. Free keys are scoped to the IP network and UTC day. Workspace keys are scoped to the workspace. Changing the request body does not make a used key reusable.", + description: "Optional replay protection. Whole-document jobs require a UUID and return 202 with the same job on retry, without executing or charging again; reuse only for the same document. Otherwise a previously admitted key returns 409 without executing again; responses are not cached. Free keys are scoped to the IP network and UTC day. Workspace keys are scoped to the workspace.", }, }, securitySchemes: { @@ -1491,13 +1514,20 @@ $100/day across all free traffic, with four concurrent requests per IP. IPv6 addresses share a /64. Smart requests must fit the request allowance. Funded workspace keys use their balance and bypass the shared subsidy and proxy check; the default maximum request allowance is $10. Request bodies -are limited to 1 MB. Billing settles asynchronously after the response. +are limited to 1 MB for synchronous classification. One whole Jev document can +be uploaded to POST /v1/classify as JSON or UTF-8 text/plain: up to 10M tokens +and 100 MB. Large uploads return 202 with a status_url; poll with the same key. +Prefer: respond-async requests this behavior at any size. No client splitting +or finish call is required. Jobs reserve the exact original token price after +upload, finish automatically, and refund on failure. Source is temporarily +stored privately and deleted as screened, or on failure, cancellation or expiry. +Results expire after 24 hours. Synchronous billing settles after the response. Free, per IP, counted in classifications: 3,000/minute and 20,000/day on the fast tier, 200/minute and 2,000/day on the smart tier. Default/explicit Jev inputs over 32,000 characters use paid Fast-only long context. Requires paid workspace balance or active paid subscription; anonymous access and signup credit do not -qualify. Limits: 250,000 original cl100k_base context tokens summed across inputs, +qualify. Synchronous limits: 250,000 original cl100k_base context tokens summed across inputs, 20 documents, 32 decisions and 1 MB body. Retail: $0.084/M original context tokens, counted once regardless of dimensions or actual screening/final usage. Final Jev reads selected whole chunks in source order; eligible evidence can be diff --git a/src/pages.ts b/src/pages.ts index 3016f93..1509cdb 100644 --- a/src/pages.ts +++ b/src/pages.ts @@ -13,7 +13,7 @@ import { SITE, SITE_UPDATED } from "./wellknown"; import { codeLang } from "./ui"; import { BILLING_PLANS, formatCreditsUsd } from "./lib/billing"; import { INPUT_PRICE_PER_MILLION, ESCALATION_PRICE_PER_THOUSAND, LONG_CONTEXT_PRICING } from "./lib/classification-pricing"; -import { LONG_CONTEXT_JOB_MAX_TOKENS, LONG_CONTEXT_PART_MAX_TOKENS } from "./long-context"; +import { LONG_CONTEXT_JOB_MAX_TOKENS } from "./long-context"; export const MCP_SETUP = `classifier.dev MCP @@ -279,7 +279,7 @@ LONG CONTEXT 32,000 characters through chunking, parallel Jev evidence screening and a final Jev call. POST with a workspace key backed by paid balance or an active paid subscription; anonymous access and free signup credit do not qualify. - Fast only. Limits: 250,000 original cl100k_base context tokens summed across + Fast only. Synchronous limits: 250,000 original cl100k_base context tokens summed across inputs, 20 documents, 32 decisions and a 1 MB request body. Decisions count documents × dimensions, or documents × labels in multi-label mode. Existing dedicated enterprise/operator access remains supported. @@ -297,12 +297,20 @@ LONG CONTEXT https://classifier.dev for the full field reference and limitations. Explicit model: "chunklaya" remains a separate legacy opt-in. - For up to ${LONG_CONTEXT_JOB_MAX_TOKENS.toLocaleString("en-US")} tokens, a funded workspace can create a long-context - job and upload ordered parts (at most ${LONG_CONTEXT_PART_MAX_TOKENS.toLocaleString("en-US")} tokens and 1 MB each). Screening - runs on every part; final Jev judges selected evidence. The same $0.084/M - original-token rate applies to the parts actually uploaded. Jobs hold the - maximum quoted charge, settle once at completion, and refund on cancellation - or 24-hour expiry. See TEN-MILLION-TOKEN JOBS in the full docs. + A funded workspace can upload one whole document of up to + ${LONG_CONTEXT_JOB_MAX_TOKENS.toLocaleString("en-US")} tokens (100 MB) in one POST /v1/classify request. + Send JSON with input and labels, + or text/plain with labels in the query. Large documents return 202 and a + status_url. For the same flow with smaller inputs, send the header + Prefer: respond-async. + Splitting, screening and final judgment happen automatically. + The original whole-document token count sets the hold and final charge at + $0.084/M. Failures, cancellation and 24-hour expiry refund unfinished work. + See TEN-MILLION-TOKEN JOBS in the full docs. + + The repository CLI uploads and waits: + + node cli/classify.js a,b --document document.txt --json AUTHENTICATION @@ -480,7 +488,8 @@ FREE must fit the same allowance. Large inputs or batches need a funded key. Free access pauses when the shared pool or verification capacity is spent; anonymous proxy networks require a funded key. Signup credit uses these - same free limits. Every request body is limited to 1 MB. + same free limits. Synchronous request bodies are limited to 1 MB; + funded whole-document uploads allow 10M tokens and 100 MB. PRO @@ -548,7 +557,11 @@ USAGE PRICES $${(250000 * LONG_CONTEXT_PRICING.inputNanodollars / 1e9).toFixed(3)}. No eligible evidence returns 422 long_context_no_evidence, no charge. Requires paid balance or an active paid subscription; free signup credit and anonymous access do not qualify. Fast only, up to 20 documents and 32 - decisions, within the 250,000-token and 1 MB request limits. Final Jev uses + decisions, within the synchronous 250,000-token and 1 MB limits. For one + whole document, POST /v1/classify accepts up to 10M tokens and 100 MB in one + upload and returns 202 with a status_url. No client splitting is needed. + The exact document price is reserved after upload; 10M tokens costs $0.84. + Failed or canceled jobs are refunded. Final Jev uses selected evidence; eligible chunks can be omitted when its budget fills. @@ -745,10 +758,11 @@ WORKSPACES, API KEYS AND ACTIVITY redacted, and content is size-limited; redaction cannot detect every kind of sensitive text. Public, keyless requests are not included in this dataset. - Long-context jobs keep only selected evidence excerpts in private temporary - job storage while a paid job is open. The full uploaded document is not - retained. Selected excerpts are deleted on completion, cancellation or - expiry (24 hours after creation). The result remains until that expiry; + Whole-document jobs temporarily keep uploaded source text in private job + storage while queued. Source fragments are deleted as they are screened; + selected evidence is deleted on completion. Cancellation, failure and + expiry (24 hours after creation) delete remaining source and evidence. + The result remains until that expiry; aggregate usage records follow normal workspace retention rules. Ask ${SITE.email} about access to or deletion of your account data. diff --git a/src/pricingui.ts b/src/pricingui.ts index decb6ee..d037247 100644 --- a/src/pricingui.ts +++ b/src/pricingui.ts @@ -71,8 +71,8 @@ export function pricingHtml(signedIn = false) {

For example, 1 million input tokens with 50 Smart escalations cost $0.142.

Output tokens are free. Input usage includes the text, labels and instructions processed by the base classifier. Retries, fallback routing and Smart model tokens add no separate charges. Paid requests already admitted can finish and leave a negative balance. New requests require a positive available balance. No automatic top-ups.

Jev long context starts automatically above 32,000 characters with the default model or explicit Jev. It costs 2 × Jev's $${INPUT_PRICE_PER_MILLION.toFixed(3)} rate: $${(LONG_CONTEXT_PRICING.inputNanodollars / 1000).toFixed(3)} per million original context tokens, counted with cl100k_base once across inputs. Dimensions and actual screening or final-call usage do not multiply this price. 250,000 context tokens cost $${(250000 * LONG_CONTEXT_PRICING.inputNanodollars / 1e9).toFixed(3)}.

-

Requires paid workspace balance or an active paid subscription; anonymous access and free signup credit do not qualify. Fast only, up to 250,000 original context tokens, 20 documents and 32 decisions within a 1 MB request. Each document × dimension or multi-label category counts as a decision.

-

For larger documents, long-context jobs screen ordered uploads up to ${(LONG_CONTEXT_JOB_MAX_TOKENS / 1e6).toFixed(0)} million original tokens at the same rate. A full job costs $${(LONG_CONTEXT_JOB_MAX_TOKENS * LONG_CONTEXT_PRICING.inputNanodollars / 1e9).toFixed(2)}. The workspace reserves its stated maximum at creation, then pays only for the tokens uploaded when final judgment succeeds. Jobs expire after 24 hours; unfinished jobs are refunded.

+

Requires paid workspace balance or an active paid subscription; anonymous access and free signup credit do not qualify. Fast only. Synchronous requests allow up to 250,000 original context tokens, 20 documents and 32 decisions within 1 MB. Each document × dimension or multi-label category counts as a decision.

+

Upload one whole document of up to ${(LONG_CONTEXT_JOB_MAX_TOKENS / 1e6).toFixed(0)} million original tokens (100 MB) at the same rate. Splitting and screening happen automatically. A full job costs $${(LONG_CONTEXT_JOB_MAX_TOKENS * LONG_CONTEXT_PRICING.inputNanodollars / 1e9).toFixed(2)}. The workspace reserves the actual uploaded document's token price, then settles once when final judgment succeeds. Jobs expire after 24 hours; failed or canceled jobs are refunded.

Final Jev reads selected evidence. Eligible chunks can be omitted when the final budget fills; usage.long_context reports selection. No usable evidence returns 422 long_context_no_evidence without charge. Explicit chunklaya remains a separate legacy opt-in. Read the long-context limits.

Included on every plan

diff --git a/src/server/usage.ts b/src/server/usage.ts index cd95539..86e85e1 100644 --- a/src/server/usage.ts +++ b/src/server/usage.ts @@ -1,4 +1,4 @@ -import { AppError, hashToken, now, type AppEnv } from "./db"; +import { AppError, hashToken, now, parseCreditInteger, type AppEnv } from "./db"; export interface Reservation { id: string; accountId: string; @@ -12,7 +12,7 @@ export async function authorizeAndReserve( env: AppEnv, cost: number, itemCount = cost, - metadata?: { type?: string; meteringMode?: "credits" | "tokens" }, + metadata?: { type?: string; meteringMode?: "credits" | "tokens"; reservationId?: string }, ): Promise { const token = request.headers.get("Authorization")?.replace(/^Bearer\s+/i, "") ?? ""; @@ -42,7 +42,41 @@ export async function authorizeAndReserve( .bind(await hashToken(token)) .first<{ id: string; account_id: string }>(); if (!agent) throw new AppError(401, "Invalid agent credential."); - const id = crypto.randomUUID(); + const id = metadata?.reservationId ?? crypto.randomUUID(); + if (metadata?.reservationId) { + const recover = async (): Promise => { + const prior = await env.APP_DB.prepare("SELECT u.account_id,u.agent_id,u.credits,u.items,u.status,u.metering_mode,a.billing_plan FROM app_usage u JOIN app_accounts a ON a.id=u.account_id WHERE u.id=?") + .bind(id).first<{ account_id: string; agent_id: string; credits: number; items: number; status: string; metering_mode: string; billing_plan: string }>(); + if (!prior) return null; + if (prior.account_id !== agent.account_id || parseCreditInteger(String(prior.credits)) !== cost || prior.items !== itemCount || + prior.metering_mode !== meteringMode || prior.status !== "pending") + throw new AppError(409, "This reservation cannot be reused. Start a new document job."); + return { id, accountId: prior.account_id, agentId: prior.agent_id, cost, billingPlan: prior.billing_plan }; + }; + const prior = await recover(); + if (prior) return prior; + // Debit only the row inserted by this statement, never an earlier hold + // that committed after the recovery read (including a lost HTTP reply). + const admitted = await env.APP_DB.batch([env.APP_DB.prepare(`WITH admitted AS ( + INSERT INTO app_usage(id,account_id,agent_id,items,credits,status,created_at,paid_credits,usage_type,metering_mode) + SELECT ?,?,?,?,?, 'pending', ?,GREATEST(0,?-(balance-paid_balance)),?,? + FROM app_accounts WHERE id=? AND balance>=? AND NOT billing_hold + AND EXISTS(SELECT 1 FROM app_agents WHERE id=? AND account_id=? AND status IN ('pending','connected')) + ON CONFLICT(id) DO NOTHING RETURNING id,paid_credits + ), account_debit AS ( + UPDATE app_accounts SET balance=balance-?,paid_balance=paid_balance-admitted.paid_credits + FROM admitted WHERE app_accounts.id=? RETURNING billing_plan + ), agent_debit AS ( + UPDATE app_agents SET used=used+? FROM admitted WHERE app_agents.id=? RETURNING app_agents.id + ) SELECT billing_plan FROM account_debit`).bind(id, agent.account_id, agent.id, itemCount, cost, now(), cost, + metadata.type?.slice(0, 80) || "classification", meteringMode, agent.account_id, cost, agent.id, agent.account_id, + cost, agent.account_id, cost, agent.id)]); + const inserted = admitted[0].results[0] as { billing_plan: string } | undefined; + if (inserted) return { id, accountId: agent.account_id, agentId: agent.id, cost, billingPlan: inserted.billing_plan }; + const concurrent = await recover(); + if (concurrent) return concurrent; + throw new AppError(403, "Credential is paused or revoked, or the workspace balance is too low."); + } const results = await env.APP_DB.batch([ env.APP_DB.prepare( "INSERT INTO app_usage(id,account_id,agent_id,items,credits,status,created_at,paid_credits,usage_type,metering_mode) SELECT ?,?,?,?,?,?,?,GREATEST(0,?-(balance-paid_balance)),?,? FROM app_accounts WHERE id=? AND balance>=? AND NOT billing_hold AND EXISTS(SELECT 1 FROM app_agents WHERE id=? AND account_id=? AND status IN ('pending','connected'))", diff --git a/wrangler.example.toml b/wrangler.example.toml index bab5fdf..b611115 100644 --- a/wrangler.example.toml +++ b/wrangler.example.toml @@ -10,6 +10,9 @@ compatibility_flags = ["nodejs_compat"] compatibility_date = "2026-08-01" account_id = "YOUR_CLOUDFLARE_ACCOUNT_ID" +[limits] +cpu_ms = 300_000 + # Run fetch handlers near Oregon. This is a placement target, not a # data-residency guarantee; assets remain globally served. [placement]