Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 15 additions & 15 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,24 +53,24 @@ console.log(vectors.length, vectors[0]?.length); // 2 768

### Config

| Field | Type | Description |
| ---------------- | ---------------------- | ------------------------------------------------------------------------------------------- |
| `baseURL` | `string` | Server root with its version path, such as `http://host:11434/v1`. |
| `model` | `string` | Model name. |
| `apiKey` | `string?` | Sent as a bearer token. |
| `dimensions` | `number?` | Requested output dimension. Sent only when set; a reply of any other dimension is rejected. |
| `encodingFormat` | `"float" \| "base64"?` | Wire format. Defaults to `float`. Results are always `number[]`. |
| `batchSize` | `number?` | Texts per request. Defaults to 32. |
| `timeoutMs` | `number?` | Per-attempt timeout; a timed-out attempt is retried. Defaults to 30000. |
| Field | Type | Description |
| ---------------- | ---------------------- | ------------------------------------------------------------------------------------------------ |
| `baseURL` | `string` | Server root with its version path, such as `http://host:11434/v1`. Trailing slashes are ignored. |
| `model` | `string` | Model name. |
| `apiKey` | `string?` | Sent as a bearer token. |
| `dimensions` | `number?` | Requested output dimension. Sent only when set; a reply of any other dimension is rejected. |
| `encodingFormat` | `"float" \| "base64"?` | Wire format. Defaults to `float`. Results are always `number[]`. |
| `batchSize` | `number?` | Texts per request. Defaults to 32. |
| `timeoutMs` | `number?` | Per-attempt timeout; a timed-out attempt is retried. Defaults to 30000. |

### Options

| Field | Description |
| --------------------- | ------------------------------------------------------------------------------- |
| `deps` | `{ fetch, scheduler }`. Defaults to global `fetch` and Interchange's scheduler. |
| `retryPolicy` | Retry policy. Defaults to Interchange's policy. |
| `extractRetryAfterMs` | Reads the server's retry delay. Defaults to parsing `Retry-After`. |
| `signal` | Aborts all pending requests. |
| Field | Description |
| --------------------- | ---------------------------------------------------------------------------------------- |
| `deps` | `{ fetch, scheduler }`. Defaults to global `fetch` and Interchange's scheduler. |
| `retryPolicy` | Retry policy. Defaults to Interchange's policy. |
| `extractRetryAfterMs` | Reads the server's retry delay. Defaults to parsing `Retry-After`, capped at 60 seconds. |
| `signal` | Aborts all pending requests. |

### Errors

Expand Down
13 changes: 13 additions & 0 deletions e2e/retry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,3 +165,16 @@ test("reports a caller abort as aborted, without retrying", async () => {
reason: { category: "aborted" },
});
});

test("caps a Retry-After at 60 seconds", async () => {
harness = setupHarness({ enableInferenceTimers: true });
reply(harness, "rate limited", 429, { "retry-after": "86400" });
reply(harness, JSON.stringify({ data: [{ index: 0, embedding: [1] }] }), 200);

const pending = embedTexts(["a"], CONFIG, { deps: harness.deps });
await harness.run();

expect(await pending).toEqual([[1]]);
expect(harness.clock.now()).toBeGreaterThanOrEqual(60_000);
expect(harness.clock.now()).toBeLessThan(61_000);
});
19 changes: 19 additions & 0 deletions e2e/wire.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,25 @@ test("sends only configured fields to {baseURL}/embeddings", async () => {
expect(await run(harness, pending)).toEqual([[1]]);
});

test("trims trailing slashes from baseURL", async () => {
harness = setupHarness();
const stream = harness.scenario.createStream();
harness.scenario.whenRequestMatches(
(req) => req.url === "https://embed.example/v1/embeddings",
stream,
);
stream.enqueueAll(
[new TextEncoder().encode('{"data":[{"index":0,"embedding":[1]}]}')],
{ startAt: 1 },
);
const pending = embedTexts(
["a"],
{ ...CONFIG, baseURL: "https://embed.example/v1//" },
{ deps: harness.deps },
);
expect(await run(harness, pending)).toEqual([[1]]);
});

test("sends a bearer token only when apiKey is set", async () => {
harness = setupHarness();
const stream = harness.scenario.createStream();
Expand Down
2 changes: 1 addition & 1 deletion src/embed.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ function buildRequest(
}

return {
url: `${config.baseURL}/embeddings`,
url: `${config.baseURL.replace(/\/+$/, "")}/embeddings`,
headers,
// Unset `dimensions` and `encoding_format` are dropped by JSON.stringify.
body: JSON.stringify({
Expand Down
16 changes: 11 additions & 5 deletions src/request.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,12 +63,18 @@ type Attempt =

export const DEFAULT_TIMEOUT_MS = 30_000;

const MAX_RETRY_AFTER_MS = 60_000;

const clampRetryAfter = (ms: number): number =>
Math.min(MAX_RETRY_AFTER_MS, Math.max(0, ms));

/**
* `Retry-After` in seconds or as an HTTP date; undefined when absent.
*
* Both branches clamp at zero. A server may name an instant that has already
* passed, and `Retry-After: -5` is not unheard of; a negative delay would flow
* into the retry policy as though it were a pacing hint.
* Both branches clamp to [0, 60s]. A server may name an instant that has
* already passed, and `Retry-After: -5` is not unheard of; a negative delay
* would flow into the retry policy as though it were a pacing hint. A day-long
* `Retry-After` would stall the caller for that long.
*
* The HTTP-date branch is wall-clock by necessity — the header names an
* absolute instant, so it cannot be resolved against the harness scheduler's
Expand All @@ -79,10 +85,10 @@ export const extractRetryAfterMs: RetryAfterExtractor = (headers) => {
if (header === undefined || header === "") return undefined;

const seconds = Number(header);
if (Number.isFinite(seconds)) return Math.max(0, seconds * 1_000);
if (Number.isFinite(seconds)) return clampRetryAfter(seconds * 1_000);

const at = Date.parse(header);
return Number.isNaN(at) ? undefined : Math.max(0, at - Date.now());
return Number.isNaN(at) ? undefined : clampRetryAfter(at - Date.now());
};

async function attemptOnce(
Expand Down
Loading