diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index fbb1eff..616981d 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -13,32 +13,31 @@ Callers import from `@corbits/embedding`. The barrel is the public surface: | `embedTexts` | Batch texts, post sequentially, return `number[][]` in input order | | `probeEmbedDims` | Embed one probe string; report the vector length that came back | | `EmbedConfig` / `EmbedConfigSchema` | Endpoint, model, and optional knobs | -| `EmbedOptions` | Harness deps plus optional retry, `Retry-After` override, abort | -| `runJSONRequest` | Classified, retried JSON POST (sibling clients reuse this) | -| `extractRetryAfterMs` | Default `Retry-After` reader (seconds or HTTP-date) | -| `ModelRequestError` | Transport, HTTP, or non-JSON 200 — `reason` is the classified `InferenceError` | +| `EmbedOptions` | Optional harness deps, retry, `Retry-After` override, abort | +| `RequestDependencies` / `RetryAfterExtractor` | Types of the `EmbedOptions` fields | +| `EmbeddingRequestError` | Every request failure — `reason` is the classified `InferenceError` | + +The JSON transport (`runJSONRequest`) is private. Quickstart shape (matches the README): ```ts -import { createDefaultScheduler } from "@intx/inference"; -import { embedTexts, probeEmbedDims } from "@corbits/embedding"; - -const deps = { fetch, scheduler: createDefaultScheduler() }; -const [vector] = await embedTexts( - ["hello world"], - { baseURL: "http://localhost:11434/v1", model: "nomic-embed-text" }, - { deps }, -); +import { embedTexts } from "@corbits/embedding"; + +const [vector] = await embedTexts(["hello world"], { + baseURL: "http://localhost:11434/v1", + model: "nomic-embed-text", +}); ``` -`probeEmbedDims(config, { deps })` is the same config and deps; it is not a +`probeEmbedDims(config)` is the same config and options; it is not a separate protocol. ## Request dependencies The transport reads only `fetch` and `scheduler` from the harness -(`RequestDependencies`). That is deliberately narrower than Interchange +(`RequestDependencies`), defaulting to global `fetch` and +`createDefaultScheduler()`. That is deliberately narrower than Interchange `Dependencies`, whose `adapters` registry a caller should not have to assemble in order to embed a string. A real `Dependencies` still satisfies the pick structurally. @@ -88,20 +87,22 @@ and `createDefaultRetryPolicy` (back off retryables, abort the rest). A 429 is therefore classified and backed off as a chat 429 would be; a `credential_failure` aborts immediately. -`runJSONRequest` is duplicated, modulo comments, in `@corbits/reranking`. It -is not shared infrastructure yet — the intended home is `@intx/inference` as -a non-streaming sibling of `runInference`. Until then, `ModelRequestError` is -a distinct class in each package: `instanceof` does not hold across the two. -Catch on `error.name === "ModelRequestError"`. +Following Interchange, the `InferenceError` is carried as data on a thrown +`Error` subclass (`EmbeddingRequestError`), never thrown itself. + +`runJSONRequest` is duplicated, modulo comments, in `@corbits/reranking`, +which still throws its own `ModelRequestError`; code using both packages +checks two error classes. The intended home for the transport is +`@intx/inference`, as a non-streaming sibling of `runInference`. ## Failure modes - Transport, HTTP status, or a 200 whose body is not JSON → - `ModelRequestError` with classified `reason` and the URL. -- Malformed embeddings JSON, wrong count, bad index → thrown `Error` (not - `ModelRequestError`); the protocol succeeded, the payload did not. -- `batchSize` not an integer `>= 1` → `RangeError` before any request - (`i += size` of 0 would never advance). + `EmbeddingRequestError` with classified `reason` and the URL. +- Malformed embeddings JSON, wrong count, bad index → + `EmbeddingRequestError` with a `classifyProtocolMismatch` reason. +- Config failing `EmbedConfigSchema` (e.g. `batchSize` not an integer + `>= 1`) → arktype `TraversalError` before any request. - Empty `texts` → `[]`, no request. - Caller `signal` or per-attempt timeout (default 30s) abort the attempt; aborting mid-delay wakes immediately and the next attempt fails its diff --git a/IMPLEMENTATION.md b/IMPLEMENTATION.md index c50c2c2..4556a07 100644 --- a/IMPLEMENTATION.md +++ b/IMPLEMENTATION.md @@ -80,8 +80,8 @@ Reply: `{ data: { index: number, embedding: number[] | string }[] }`. Base64 strings decode with `atob` → `Uint8Array` → little-endian `Float32Array` → `number[]`. -`embedTexts` also range-checks `batchSize` at runtime so a value that -bypassed the schema cannot stall `batches` (`i += size`). +`embedTexts` asserts `config` against `EmbedConfigSchema` first, so a +`batchSize` below 1 cannot stall `batches` (`i += size`). ## Errors and retry @@ -92,16 +92,13 @@ bypassed the schema cannot stall `batches` (`i += size`). 3. network throw → `classifyNetworkError` or `classifyAbortError` 4. JSON parse fail → `classifyProtocolMismatch` 5. `retryPolicy` (default `createDefaultRetryPolicy`) → `abort` throws - `ModelRequestError`, else sleep `delayMs` on `deps.scheduler` + `EmbeddingRequestError`, else sleep `delayMs` on `deps.scheduler` `extractRetryAfterMs` reads `Retry-After` as seconds or HTTP-date, clamped at zero. Callers override via `EmbedOptions.extractRetryAfterMs` (same signature as Interchange's unexported extractor; re-declared so this package depends only on the published surface). -`ModelRequestError.name` is `"ModelRequestError"`. Discriminate on the -name when catching both this package and `@corbits/reranking`. - ## Tests `bun run test` is `bun test ./src`. Coverage is the client contract: empty diff --git a/README.md b/README.md index 07a4c52..03e37b9 100644 --- a/README.md +++ b/README.md @@ -73,16 +73,7 @@ scheduler), with optional `retryPolicy`, `extractRetryAfterMs`, and `signal`. ## Errors -Every failure mode — transport, HTTP status, or a 200 with an unexpected body — raises `ModelRequestError`, carrying the classified `InferenceError` as `reason` plus the request URL. The embedding and reranking packages each carry their own copy of this class while the shared transport is upstreamed, so code catching both discriminates on `error.name === "ModelRequestError"`. - -## Lower-level: transport - -The barrel also re-exports the one-shot JSON transport `embedTexts` is built on, for sibling clients that want the same classified, retried request path: - -- `runJSONRequest` — one JSON POST with error classification and retry. -- `extractRetryAfterMs` — default `Retry-After` reader, overridable per call. -- `ModelRequestError` — error for every failure mode, with `reason` and URL. -- Types `RunRequestOptions` and `RetryAfterExtractor`. +Every request failure — transport, HTTP status, or a 200 with an unexpected body — throws `EmbeddingRequestError` (`extends Error`), carrying the classified `InferenceError` as `reason` and the request `url`. A config that fails `EmbedConfigSchema` throws arktype's `TraversalError` before any request. ## Interchange diff --git a/src/embed.test.ts b/src/embed.test.ts index 22cb770..a3bd7d0 100644 --- a/src/embed.test.ts +++ b/src/embed.test.ts @@ -8,7 +8,7 @@ import { EmbedConfigSchema, type EmbedConfig, } from "./embed"; -import { extractRetryAfterMs, ModelRequestError } from "./request"; +import { EmbeddingRequestError } from "./request"; const CONFIG: EmbedConfig = { baseURL: "https://embed.example/v1", @@ -134,12 +134,12 @@ describe("embedTexts", () => { it("rejects a batchSize that is not an integer >= 1 instead of looping forever", async () => { // `batches` advances by `i += size`, so a size of 0 never advances and // grows the batch list until OOM. Fractional sizes would also stall or - // slice wrongly. Schema and runtime both require integer >= 1. + // slice wrongly. The schema requires integer >= 1. for (const batchSize of [0, -1, 0.5, Number.NaN, 1.5]) { const { deps: d, calls } = deps([embedding(3)]); await expect( embedTexts(["a"], { ...CONFIG, batchSize }, { deps: d }), - ).rejects.toBeInstanceOf(RangeError); + ).rejects.toThrow(/batchSize/); expect(calls).toHaveLength(0); } }); @@ -148,9 +148,9 @@ describe("embedTexts", () => { const { deps: d } = deps([ Response.json({ data: [{ index: 0, embedding: true }] }), ]); - await expect(embedTexts(["a"], CONFIG, { deps: d })).rejects.toThrow( - /malformed embeddings response/, - ); + const rejection = expect(embedTexts(["a"], CONFIG, { deps: d })).rejects; + await rejection.toBeInstanceOf(EmbeddingRequestError); + await rejection.toThrow(/malformed embeddings response/); }); }); @@ -195,7 +195,7 @@ describe("retry behaviour inherited from the inference policy", () => { new Response("unauthorized", { status: 401 }), ]); await expect(embedTexts(["a"], CONFIG, { deps: d })).rejects.toBeInstanceOf( - ModelRequestError, + EmbeddingRequestError, ); expect(calls).toHaveLength(1); }); @@ -232,21 +232,6 @@ describe("retry-after handling", () => { }); expect(seen).toEqual(["7"]); }); - - it("reads the seconds form and tolerates an absent header", () => { - expect(extractRetryAfterMs(new Headers({ "retry-after": "2" }))).toBe( - 2_000, - ); - expect(extractRetryAfterMs(new Headers())).toBeUndefined(); - }); - - it("never returns a negative delay for a date already past", () => { - expect( - extractRetryAfterMs( - new Headers({ "retry-after": "Wed, 21 Oct 2015 07:28:00 GMT" }), - ), - ).toBeGreaterThanOrEqual(0); - }); }); describe("reply integrity", () => { @@ -319,7 +304,7 @@ describe("transport failure modes", () => { }), ]); await expect(embedTexts(["a"], CONFIG, { deps: d })).rejects.toBeInstanceOf( - ModelRequestError, + EmbeddingRequestError, ); }); @@ -347,7 +332,7 @@ describe("transport failure modes", () => { { ...CONFIG, timeoutMs: 5 }, { deps: d, signal: controller.signal }, ), - ).rejects.toBeInstanceOf(ModelRequestError); + ).rejects.toBeInstanceOf(EmbeddingRequestError); expect(calls).toHaveLength(1); }); }); diff --git a/src/embed.ts b/src/embed.ts index 7ceccbd..2a12416 100644 --- a/src/embed.ts +++ b/src/embed.ts @@ -1,8 +1,16 @@ import { type } from "arktype"; -import type { BuiltRequest } from "@intx/inference"; +import { + classifyProtocolMismatch, + createDefaultRetryPolicy, + createDefaultScheduler, + type BuiltRequest, +} from "@intx/inference"; import type { RetryPolicy } from "@intx/types/runtime"; import { + DEFAULT_TIMEOUT_MS, + EmbeddingRequestError, + extractRetryAfterMs, runJSONRequest, type RequestDependencies, type RetryAfterExtractor, @@ -42,7 +50,7 @@ export type EmbedConfig = typeof EmbedConfigSchema.infer; */ const EmbedResponse = type({ data: type({ - index: "number", + index: "number.integer", embedding: "number[] | string", }).array(), }); @@ -51,11 +59,14 @@ const DEFAULT_BATCH_SIZE = 32; const PROBE_TEXT = "embedding dimension probe"; export type EmbedOptions = { - deps: RequestDependencies; + /** Defaults to global `fetch` and Interchange's default scheduler. */ + deps?: RequestDependencies; + /** Defaults to Interchange's policy: back off retryables, abort the rest. */ retryPolicy?: RetryPolicy; /** * Reads `Retry-After` off a failed response, for a provider that signals - * pacing its own way. Named after `ProviderAdapter.extractRetryAfterMs`. + * pacing its own way. Defaults to seconds or HTTP-date parsing. Named + * after `ProviderAdapter.extractRetryAfterMs`. * There is no adapter object here because every provider serves the one * `/v1/embeddings` shape, so this hangs off the options instead. */ @@ -77,22 +88,27 @@ function buildRequest( return { url: `${config.baseURL}/embeddings`, headers, + // Unset `dimensions` and `encoding_format` are dropped by JSON.stringify. body: JSON.stringify({ model: config.model, input, - ...(config.dimensions !== undefined - ? { dimensions: config.dimensions } - : {}), - ...(config.encodingFormat !== undefined - ? { encoding_format: config.encodingFormat } - : {}), + dimensions: config.dimensions, + encoding_format: config.encodingFormat, }), }; } /** Unpacks the `base64` encoding: a little-endian float32 buffer. */ -function decodeBase64Embedding(encoded: string): number[] { - const bytes = Uint8Array.from(atob(encoded), (char) => char.charCodeAt(0)); +function decodeBase64Embedding(encoded: string): number[] | undefined { + let binary: string; + try { + binary = atob(encoded); + } catch { + return undefined; + } + if (binary.length % Float32Array.BYTES_PER_ELEMENT !== 0) return undefined; + + const bytes = Uint8Array.from(binary, (char) => char.charCodeAt(0)); const floats = new Float32Array( bytes.buffer, bytes.byteOffset, @@ -114,32 +130,37 @@ function parseResponse( url: string, expected: number, ): number[][] { + const malformed = (detail: string): EmbeddingRequestError => + new EmbeddingRequestError(classifyProtocolMismatch(detail, body), url); + const parsed = EmbedResponse(body); if (parsed instanceof type.errors) { - throw new Error( - `${url}: malformed embeddings response — ${parsed.summary}`, - ); + throw malformed(`malformed embeddings response — ${parsed.summary}`); } if (parsed.data.length !== expected) { - throw new Error( - `${url}: expected ${expected} embeddings, got ${parsed.data.length}`, + throw malformed( + `expected ${expected} embeddings, got ${parsed.data.length}`, ); } - const vectors = new Array(expected); - for (const entry of parsed.data) { - const taken = entry.index in vectors; - if (entry.index < 0 || entry.index >= expected || taken) { - throw new Error(`${url}: bad embedding index ${entry.index}`); + const seen = new Set(); + for (const { index } of parsed.data) { + if (index < 0 || index >= expected || seen.has(index)) { + throw malformed(`bad embedding index ${index}`); } - vectors[entry.index] = - typeof entry.embedding === "string" - ? decodeBase64Embedding(entry.embedding) - : entry.embedding; + seen.add(index); } - // No slot is empty: the count matches and every index was distinct. - return vectors as number[][]; + // Every index in [0, expected) appears exactly once, so sorting restores + // request order. + return parsed.data + .toSorted((a, b) => a.index - b.index) + .map(({ embedding }) => { + if (typeof embedding !== "string") return embedding; + const decoded = decodeBase64Embedding(embedding); + if (decoded === undefined) throw malformed("invalid base64 embedding"); + return decoded; + }); } function batches( @@ -161,32 +182,28 @@ function batches( export async function embedTexts( texts: readonly string[], config: EmbedConfig, - options: EmbedOptions, + options: EmbedOptions = {}, ): Promise { + const valid = EmbedConfigSchema.assert(config); if (texts.length === 0) return []; - const batchSize = config.batchSize ?? DEFAULT_BATCH_SIZE; - if (!Number.isInteger(batchSize) || batchSize < 1) { - throw new RangeError( - `batchSize must be an integer >= 1, got ${batchSize}: a batch size below 1 would never advance through the input`, - ); - } + const deps = options.deps ?? { + fetch, + scheduler: createDefaultScheduler(), + }; + + const retryPolicy = options.retryPolicy ?? createDefaultRetryPolicy(); + const timeoutMs = valid.timeoutMs ?? DEFAULT_TIMEOUT_MS; + const extractRetryAfter = options.extractRetryAfterMs ?? extractRetryAfterMs; const vectors: number[][] = []; - for (const batch of batches(texts, batchSize)) { - const request = buildRequest(config, batch); - const body = await runJSONRequest(request, { - deps: options.deps, - ...(options.retryPolicy !== undefined - ? { retryPolicy: options.retryPolicy } - : {}), - ...(config.timeoutMs !== undefined - ? { timeoutMs: config.timeoutMs } - : {}), - ...(options.extractRetryAfterMs !== undefined - ? { extractRetryAfterMs: options.extractRetryAfterMs } - : {}), - ...(options.signal !== undefined ? { signal: options.signal } : {}), + for (const batch of batches(texts, valid.batchSize ?? DEFAULT_BATCH_SIZE)) { + const request = buildRequest(valid, batch); + const body = await runJSONRequest(request, deps, { + retryPolicy, + timeoutMs, + extractRetryAfterMs: extractRetryAfter, + signal: options.signal, }); vectors.push(...parseResponse(body, request.url, batch.length)); } @@ -202,11 +219,14 @@ export async function embedTexts( */ export async function probeEmbedDims( config: EmbedConfig, - options: EmbedOptions, + options: EmbedOptions = {}, ): Promise { const [vector] = await embedTexts([PROBE_TEXT], config, options); if (vector === undefined) { - throw new Error(`${config.baseURL}: embed probe returned no vector`); + throw new EmbeddingRequestError( + classifyProtocolMismatch("embed probe returned no vector"), + config.baseURL, + ); } return vector.length; } diff --git a/src/index.ts b/src/index.ts index 5847f42..dfbbf2d 100644 --- a/src/index.ts +++ b/src/index.ts @@ -6,9 +6,7 @@ export { type EmbedOptions, } from "./embed.js"; export { - runJSONRequest, - extractRetryAfterMs, - ModelRequestError, + EmbeddingRequestError, + type RequestDependencies, type RetryAfterExtractor, - type RunRequestOptions, } from "./request.js"; diff --git a/src/request.ts b/src/request.ts index a467881..89dd74c 100644 --- a/src/request.ts +++ b/src/request.ts @@ -3,28 +3,11 @@ import { classifyHTTPError, classifyNetworkError, classifyProtocolMismatch, - createDefaultRetryPolicy, type BuiltRequest, type Dependencies, } from "@intx/inference"; import type { InferenceError, RetryPolicy } from "@intx/types/runtime"; -/** - * One JSON request, classified and retried the way Interchange classifies and - * retries an inference call. - * - * `runInference` cannot serve this: it is turn-shaped and streams SSE, whereas - * an embedding is one-shot JSON. What carries over is everything around it — the - * single `deps.fetch` transport, the error taxonomy, and the retry policy — so - * a 429 from an embedding endpoint is classified and backed off exactly as one from - * a chat endpoint. - * - * XXX: this file is duplicated, modulo the wording of this comment, in `@corbits/reranking`. It is not - * shared infrastructure yet — the intended home is `@intx/inference` itself, as - * a non-streaming sibling of `runInference`, pending that upstream - * conversation. Fix both copies or neither. - */ - /** * What the transport actually reads from the harness. * @@ -37,19 +20,17 @@ export type RequestDependencies = Pick; /** * Raised for every failure mode: transport, HTTP status, and a 200 whose body - * is not the JSON the protocol promises. - * - * XXX: duplicated alongside this module, so `instanceof` does not hold across - * the embedding and reranking copies. Discriminate on - * `error.name === "ModelRequestError"` when catching both. + * is not the embeddings reply the protocol promises. `reason` is the + * Interchange classification, so retryability and status read the same as + * for an inference call. */ -export class ModelRequestError extends Error { +export class EmbeddingRequestError extends Error { constructor( readonly reason: InferenceError, readonly url: string, ) { super(`${url}: ${reason.message}`); - this.name = "ModelRequestError"; + this.name = "EmbeddingRequestError"; } } @@ -60,30 +41,27 @@ export class ModelRequestError extends Error { */ export type RetryAfterExtractor = (headers: Headers) => number | undefined; -export type RunRequestOptions = { - deps: RequestDependencies; - /** Defaults to Interchange's policy: back off retryables, abort the rest. */ - retryPolicy?: RetryPolicy; +type RunRequestOptions = { + retryPolicy: RetryPolicy; /** * Per-attempt ceiling, enforced alongside any caller `signal` rather than * instead of it. A cold local model can take a while to page in. */ - timeoutMs?: number; + timeoutMs: number; /** * Reads `Retry-After` off a failed response. Mirrors `ProviderAdapter`'s * member of the same name, so a provider that signals pacing its own way can - * override without touching the transport. Defaults to - * {@link extractRetryAfterMs}. + * override without touching the transport. */ - extractRetryAfterMs?: RetryAfterExtractor; - signal?: AbortSignal; + extractRetryAfterMs: RetryAfterExtractor; + signal: AbortSignal | undefined; }; type Attempt = | { ok: true; body: unknown } | { ok: false; error: InferenceError }; -const DEFAULT_TIMEOUT_MS = 30_000; +export const DEFAULT_TIMEOUT_MS = 30_000; /** * `Retry-After` in seconds or as an HTTP date; undefined when absent. @@ -198,14 +176,23 @@ function delay( }); } +/** + * One JSON request, classified and retried the way Interchange classifies and + * retries an inference call. + * + * `runInference` cannot serve this: it is turn-shaped and streams SSE, whereas + * an embedding is one-shot JSON. What carries over is everything around it — the + * single `deps.fetch` transport, the error taxonomy, and the retry policy — so + * a 429 from an embedding endpoint is classified and backed off exactly as one from + * a chat endpoint. + */ export async function runJSONRequest( request: BuiltRequest, + deps: RequestDependencies, options: RunRequestOptions, ): Promise { - const { deps, signal } = options; - const policy = options.retryPolicy ?? createDefaultRetryPolicy(); - const timeoutMs = options.timeoutMs ?? DEFAULT_TIMEOUT_MS; - const extractRetryAfter = options.extractRetryAfterMs ?? extractRetryAfterMs; + const { retryPolicy: policy, timeoutMs, signal } = options; + const extractRetryAfter = options.extractRetryAfterMs; // Time comes from the harness scheduler, not globals, so a virtual-clock // test scheduler drives the retry loop deterministically. @@ -227,7 +214,7 @@ export async function runJSONRequest( elapsedMs: deps.scheduler.now() - startedAt, }); if (decision.kind === "abort") { - throw new ModelRequestError(result.error, request.url); + throw new EmbeddingRequestError(result.error, request.url); } // Aborting mid-delay wakes immediately; the next attempt then fails its