diff --git a/CHANGELOG.md b/CHANGELOG.md index 2614faa..10f8714 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,73 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [8.2.0] - 2026-09-26 + +### Added + +- **The client reads the rate-limit headers, and acts on them.** Both meters, + the per-minute REST one and the separate hourly history budget, on + `client.rateLimit`: + + ```js + tp.rateLimit.rest.remaining; + tp.rateLimit.history.remaining; + secondsUntilReset(tp.rateLimit.rest); + ``` + + Every field can be null, and null means the server did not say rather than + "nothing left". Use `isExhausted`, true only when the server said zero. The + per-minute figures ride most responses; the hourly history ones are withheld + from anything a shared cache may store, because they are per-caller; an + unmetered plan advertises nothing. A response served from a cache is ignored + entirely, because its figures belong to whoever populated the entry. `reset` is a relative + countdown frozen when it was read, so `secondsUntilReset` ages it rather + than returning a stale number. + + A window the server says is spent is now waited out instead of walked into, + since that request is a certain 429 that also spends budget being refused. + `retry: { respectRemaining: false }` opts out. + + The hourly history budget is new on the wire; before it there was nothing to + read. + +### Changed + +- **Calls may now block before sending.** When the server has said your window + is spent, or has issued a 429 that is still in force, the client waits rather + than sending a request certain to be refused. A call that used to return in + 200ms can now take up to `retry.maxRetryAfterMs` (120000) first. That is a + TOTAL across the call, not per wait: the shared 429 gate and the + spent-window wait stack, and before the budget existed a 429 carrying both a + `Retry-After` and a spent window blocked for 180 seconds under a 120 second + cap. Turn the two halves off with `retry: { respectRemaining: false }` and + `retry: { on429: false }`. + +### Fixed + +- **A paged history call crashed on the default configuration.** `#private` + fields on the transport failed their brand check through the caching Proxy, + so `history.days()` threw `TypeError: Receiver must be an instance of class +Transport` on page two for anyone who had not passed `cache: false`. Every + test passed `cache: false`, so none of them saw it. + +- **`on429: false` did not opt out.** It threw the error the caller asked for + and then held their NEXT call for the full `Retry-After` anyway, because the + shared gate was closed regardless of the setting. + +- **The gate timed off the wall clock.** A backward NTP step turned a + five-second wait into however far the clock moved, unbounded, because the cap + is applied when the gate is armed and not when it is served. It uses a + monotonic clock now, as the Python sibling always did. + +- **A 429 was waited out once per in-flight request.** The wait belongs to the + caller, not to whichever request met it, so ten concurrent requests each + slept their own `Retry-After` and then retried at the same instant, + re-tripping the limit together. It is taken once now, on a gate shared by the + whole client, with a little jitter so waiters do not wake in unison. A + shorter wait arriving while a longer one is in force no longer brings the + gate forward. + ## [8.1.0] - 2026-09-23 ### Added diff --git a/README.md b/README.md index 6ad1d20..86081a9 100644 --- a/README.md +++ b/README.md @@ -49,15 +49,15 @@ Jungle Cruise 40 min `new ThemeParks(options)` takes the following keyword options: -| Option | Type | Default | Purpose | -| ----------- | ----------------------------------- | -------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `baseUrl` | `string` | `https://api.themeparks.wiki/v1` | API base URL (point at a mock / staging if you need to). | -| `userAgent` | `string` | `themeparks-sdk-js/` | Sent as the `User-Agent` header. Set this to identify your app. | -| `apiKey` | `string` | none | API key from api.themeparks.wiki, sent as `X-API-Key`. Optional; a key raises the limits. | -| `fetch` | `typeof fetch` | `globalThis.fetch` | Custom fetch implementation. Useful for logging, mocking, or older runtimes. | -| `timeoutMs` | `number` | `10000` | Per-request timeout in milliseconds. | -| `retry` | `Partial` | `{ max: 3, on429: true, maxRetryAfterMs: 120000 }` | Retry/backoff behavior. `max` counts retries **beyond** the initial attempt (so `3` = up to 4 total). `maxRetryAfterMs` is the longest `Retry-After` the client will sleep through; past it you get `RateLimitError` instead of a silent wait. | -| `cache` | `Cache \| false \| { maxEntries? }` | in-memory LRU | See [Caching](#caching) below. `false` disables caching entirely. | +| Option | Type | Default | Purpose | +| ----------- | ----------------------------------- | -------------------------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `baseUrl` | `string` | `https://api.themeparks.wiki/v1` | API base URL (point at a mock / staging if you need to). | +| `userAgent` | `string` | `themeparks-sdk-js/` | Sent as the `User-Agent` header. Set this to identify your app. | +| `apiKey` | `string` | none | API key from api.themeparks.wiki, sent as `X-API-Key`. Optional; a key raises the limits. | +| `fetch` | `typeof fetch` | `globalThis.fetch` | Custom fetch implementation. Useful for logging, mocking, or older runtimes. | +| `timeoutMs` | `number` | `10000` | Per-request timeout in milliseconds. | +| `retry` | `Partial` | `{ max: 3, on429: true, maxRetryAfterMs: 120000, respectRemaining: true }` | Retry/backoff behavior. `max` counts retries **beyond** the initial attempt (so `3` = up to 4 total). `maxRetryAfterMs` is the longest `Retry-After` the client will sleep through; past it you get `RateLimitError` instead of a silent wait. | +| `cache` | `Cache \| false \| { maxEntries? }` | in-memory LRU | See [Caching](#caching) below. `false` disables caching entirely. | Example: @@ -161,6 +161,49 @@ const entries = await tp.entity(mk).schedule.range(new Date('2026-05-01'), new D console.log(`${entries.length} schedule entries`); ``` +## Rate limits + +The API meters requests per minute, and history requests again per hour. Both +are read off every response that carries them: + +```js +import { ThemeParks, isExhausted, secondsUntilReset } from 'themeparks'; + +const tp = new ThemeParks({ apiKey: KEY }); +await tp.entity(parkId).live(); + +tp.rateLimit.rest.remaining; // 299 +secondsUntilReset(tp.rateLimit.rest); // 40 +tp.rateLimit.history.remaining; // on a history call +``` + +**`null` means the server did not say, never "nothing left".** Use +`isExhausted`, which is true only when the server actually said zero. + +Which figures you get depends on the response: + +- The **per-minute** figures ride most responses, anonymous ones included. +- The **hourly history** figures are withheld from anything a shared cache may + store, because they are per-caller and a cache would hand one caller's budget + to another. In practice you get them on calls made with a key. +- An **unmetered plan** advertises nothing at all. + +A response served from a cache is ignored entirely. Its figures belong to +whoever populated the entry and its countdown is already wrong: a cached +`remaining: 0` would otherwise make the client sleep out someone else's +window. + +The client acts on what it reads. When a response says the window is spent, the +next request waits for the advertised reset rather than sending one that is +certain to be refused, and to cost a unit of budget being refused. Opt out with +`retry: { respectRemaining: false }`. + +**A 429 is held once for the whole client.** The wait belongs to the caller, +not to whichever request met it, so it goes on a shared gate with a little +jitter. Without that, ten concurrent requests each sleep their own copy of +`Retry-After` and then all retry at the same instant, re-tripping the limit +together. + ## History Three endpoints answer what an entity did in the past. Days are park-local; the diff --git a/package.json b/package.json index e27822b..524bf7a 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "themeparks", - "version": "8.1.0", + "version": "8.2.0", "description": "Official SDK for the ThemeParks.wiki API", "license": "MIT", "repository": "github:ThemeParks/ThemeParks_JavaScript", diff --git a/src/client.ts b/src/client.ts index daef7e4..21ef1b0 100644 --- a/src/client.ts +++ b/src/client.ts @@ -8,9 +8,10 @@ import { type FetchLike, type RetryConfig, } from './transport'; +import type { RateLimits } from './ratelimit'; const DEFAULT_BASE_URL = 'https://api.themeparks.wiki/v1'; -const PACKAGE_VERSION = '8.1.0'; +const PACKAGE_VERSION = '8.2.0'; const DEFAULT_USER_AGENT = `themeparks-sdk-js/${PACKAGE_VERSION}`; export interface ThemeParksOptions { @@ -51,6 +52,7 @@ export class ThemeParks { max: options.retry?.max ?? 3, on429: options.retry?.on429 ?? true, maxRetryAfterMs: options.retry?.maxRetryAfterMs ?? DEFAULT_MAX_RETRY_AFTER_MS, + respectRemaining: options.retry?.respectRemaining ?? true, }, fetch: fetchFn, }); @@ -60,6 +62,26 @@ export class ThemeParks { this.destinations = new DestinationsApi(this.raw); } + /** + * What the server last said about your two budgets. + * + * `rateLimit.rest` is the per-minute REST meter; `rateLimit.history` is the + * separate hourly history budget. Every field can be null, because every + * field can legitimately be absent: an unmetered plan advertises nothing, + * and neither does a publicly cacheable response, since the figures are + * per-caller and a shared cache would hand one caller's to another. + * + * null therefore means "the server did not say", never "nothing left". Use + * `isExhausted`, which is true only when it said zero. + * + * Read from the inner transport rather than a copy, so a cache HIT -- which + * sends no request and so learns nothing -- correctly leaves the last known + * figures standing. A hit spent no budget either. + */ + get rateLimit(): RateLimits { + return this.transport.rateLimit; + } + entity(id: string): EntityHandle { return new EntityHandle(this.raw, id); } diff --git a/src/index.ts b/src/index.ts index 71bc6f5..fc2cbbc 100644 --- a/src/index.ts +++ b/src/index.ts @@ -39,4 +39,5 @@ export { type LiveDataEntry, type LiveQueue, } from './ergonomic/live'; +export { isExhausted, secondsUntilReset, type RateLimit, type RateLimits } from './ratelimit'; export { parseApiDateTime } from './dates'; diff --git a/src/ratelimit.ts b/src/ratelimit.ts new file mode 100644 index 0000000..ee67ddb --- /dev/null +++ b/src/ratelimit.ts @@ -0,0 +1,212 @@ +/** + * What the server says about your budget, and how to stay inside it. + * + * TWO BUDGETS, SEPARATELY METERED. The API meters requests per minute, and + * history requests again per hour. Different windows over different counters, + * so the server advertises them in two sets of headers: + * + * RateLimit-Limit / -Policy / -Remaining / -Reset per minute + * RateLimit-History-Limit / -Policy / -Remaining / -Reset per hour + * + * Both were being thrown away. The SDK read only `Retry-After`, and only after + * a 429 had already happened, so it could tell you that you had run out and + * never that you were about to. + * + * ABSENCE IS NOT ZERO. `null` here means "the server did not say", which is a + * different thing from "nothing left". Reading an unknown as zero would stall + * a client permanently. + * + * A CACHED RESPONSE'S FIGURES ARE NOT YOURS. The server withholds the HISTORY + * figures from anything a shared cache may store, but the per-minute REST ones + * ride those responses, so an anonymous call can come back from a CDN carrying + * another caller's numbers frozen at whatever they were when the entry was + * populated. Measured against production: three consecutive calls returning + * `age: 9` and an unmoving `remaining: 285`. A cached `remaining: 0` would make + * the client sleep out a window belonging to someone else, so a response with a + * non-zero `Age` is treated as saying nothing at all. + */ + +/** One meter's state, as of the last response that mentioned it. */ +export interface RateLimit { + /** Requests allowed per window, or null if the server did not say. */ + readonly limit: number | null; + /** Requests left in the current window, or null if the server did not say. */ + readonly remaining: number | null; + /** Seconds until the window resets, as of `observedAt`. */ + readonly reset: number | null; + /** The raw policy string, e.g. `"300;w=60"`. */ + readonly policy: string | null; + /** + * A monotonic reading (milliseconds) from when this was read, so `reset` + * can be aged. NOT an epoch timestamp: it is not meaningful as a date, and + * `new Date(observedAt)` is nonsense. The Python sibling's `observed_at` is + * the same idea in seconds. + */ + readonly observedAt: number | null; +} + +/** Both meters. Reached as `client.rateLimit`. */ +export interface RateLimits { + readonly rest: RateLimit; + readonly history: RateLimit; +} + +export const UNKNOWN_RATE_LIMIT: RateLimit = Object.freeze({ + limit: null, + remaining: null, + reset: null, + policy: null, + observedAt: null, +}); + +export const UNKNOWN_RATE_LIMITS: RateLimits = Object.freeze({ + rest: UNKNOWN_RATE_LIMIT, + history: UNKNOWN_RATE_LIMIT, +}); + +/** + * True only when the server SAID there is nothing left. + * + * An unknown remaining is not exhaustion. Treating it as such would make an + * anonymous caller, whose responses never carry figures, wait forever. + */ +export function isExhausted(meter: RateLimit): boolean { + return meter.remaining === 0; +} + +/** + * How much of the window is left, counted down from when we read it. + * + * `reset` is relative and frozen at `observedAt`; using it later without + * ageing it is how a client waits far longer than it needs to. + */ +export function secondsUntilReset(meter: RateLimit, now = now_()): number | null { + if (meter.reset === null || meter.observedAt === null) return null; + const elapsed = (now - meter.observedAt) / 1000; + return Math.max(0, meter.reset - elapsed); +} + +/** + * A clock that only moves forward. + * + * `Date.now()` is wall-clock: it steps backwards on an NTP correction, a + * manual clock change or a container resync. The Gate stores a deadline and + * subtracts the clock from it, so a backward step turns a five-second wait + * into however far the clock moved -- measured at 60 minutes for a one-hour + * step, silently, before a request the caller thinks is in flight. The cap is + * applied when the gate is ARMED, never when it is served, so nothing bounds + * it. + * + * `performance.now()` is monotonic, which is what the Python sibling gets from + * `time.monotonic()`. It excludes time the host spent suspended, so a resumed + * laptop waits out a remainder it already slept through: conservative and + * bounded, which is the right side to err on. + */ +function now_(): number { + return typeof performance !== 'undefined' && typeof performance.now === 'function' + ? performance.now() + : Date.now(); +} + +/** + * These fields are integers per the draft spec, so anything else is a header + * we do not understand and the honest answer is "unknown". + * + * `Number()` was too generous and differed from the Python sibling on the + * same input: it read "0.4" as 0, which makes `isExhausted` true and sleeps + * out a window the caller has not spent, and it accepted "0x10" as 16 and + * "1e3" as 1000. The server only ever sends a non-negative integer, so none + * of that fires today; a pair of libraries whose selling point is parity + * should not disagree on it regardless. + */ +function intOrNull(raw: string | null): number | null { + if (raw === null) return null; + const trimmed = raw.trim(); + if (!/^\d+$/.test(trimmed)) return null; + const value = Number(trimmed); + return Number.isSafeInteger(value) ? value : null; +} + +interface HeaderBag { + get(name: string): string | null; +} + +function readOne(headers: HeaderBag, prefix: string, now: number): RateLimit { + // Exact names, never a prefix scan: `ratelimit-history-limit` also begins + // with `ratelimit-`, and a scan would read the hourly figure as the + // per-minute one. + const limit = intOrNull(headers.get(`${prefix}-limit`)); + const remaining = intOrNull(headers.get(`${prefix}-remaining`)); + const reset = intOrNull(headers.get(`${prefix}-reset`)); + const policy = headers.get(`${prefix}-policy`); + if (limit === null && remaining === null && reset === null && policy === null) { + return UNKNOWN_RATE_LIMIT; + } + return { limit, remaining, reset, policy, observedAt: now }; +} + +/** + * Merge whatever this response said into what we already knew. + * + * A response mentioning neither meter leaves both alone, and so does a + * response served from a cache. Most responses + * mention only one -- the history headers appear on history routes, and a + * cacheable response carries neither -- so overwriting with blanks would mean + * the last cacheable response erased everything the client had learned. + */ +export function readRateLimits( + headers: HeaderBag, + previous: RateLimits = UNKNOWN_RATE_LIMITS, +): RateLimits { + // A cache HIT carries the figures of whoever populated the entry, frozen at + // that moment. They are not ours and the countdown is already wrong, so the + // honest reading is that this response said nothing. + const age = intOrNull(headers.get('age')); + if (age !== null && age > 0) return previous; + + const now = now_(); + const rest = readOne(headers, 'ratelimit', now); + const history = readOne(headers, 'ratelimit-history', now); + return { + rest: rest.observedAt !== null ? rest : previous.rest, + history: history.observedAt !== null ? history : previous.history, + }; +} + +/** + * One shared "not before" instant for a whole client. + * + * WHY SHARED. A 429 applies to the CALLER, not to the request that happened to + * meet it. With a per-request backoff, ten concurrent requests each sleep + * their own Retry-After and then all retry at the same instant, re-tripping + * the limit together: a thundering herd the client inflicts on itself, and on + * us. One gate means the wait is taken once. + * + * Each waiter adds its own small jitter on the way out, because waking + * together is the other half of the same problem. + */ +export class Gate { + #until = 0; + + constructor(private readonly jitterMs = 250) {} + + /** When the gate opens, on the monotonic clock. 0 if it is open. */ + get deadline(): number { + return this.#until; + } + + /** Hold every request on this client for at least `ms`. */ + closeFor(ms: number): void { + // Never bring the gate forward: a shorter Retry-After arriving while a + // longer one is in force would release the herd early. + this.#until = Math.max(this.#until, now_() + Math.max(0, ms)); + } + + /** How long this caller should hold off, jitter included. 0 if open. */ + waitMs(): number { + const remaining = this.#until - now_(); + if (!Number.isFinite(remaining) || remaining <= 0) return 0; + const jitter = Number.isFinite(this.jitterMs) ? this.jitterMs : 0; + return remaining + Math.random() * jitter; + } +} diff --git a/src/transport.ts b/src/transport.ts index e4617ae..a20ce44 100644 --- a/src/transport.ts +++ b/src/transport.ts @@ -1,4 +1,12 @@ import { ApiError, NetworkError, RateLimitError, TimeoutError } from './errors'; +import { + Gate, + isExhausted, + readRateLimits, + secondsUntilReset, + UNKNOWN_RATE_LIMITS, + type RateLimits, +} from './ratelimit'; /** * Minimal fetch-like function signature covering only what the SDK uses. @@ -49,6 +57,17 @@ export interface RetryConfig { max: number; /** If true, HTTP 429 responses are retried (honouring `Retry-After`). */ on429: boolean; + /** + * Wait out a window the server has already said is spent. + * + * When a response reports `remaining: 0`, the next request is a guaranteed + * 429 that also costs a unit of the caller's budget to be refused. Waiting + * for the reset it advertised is strictly better than sending it. Setting + * this false turns the client back into a purely reactive one. + * + * Defaults to true. + */ + respectRemaining?: boolean; /** * Longest `Retry-After` this client will sleep through, in milliseconds. * Defaults to {@link DEFAULT_MAX_RETRY_AFTER_MS}. @@ -81,15 +100,41 @@ export interface TransportOptions { /** Two minutes: longer than any REST 429 asks for, far short of an hourly budget. */ export const DEFAULT_MAX_RETRY_AFTER_MS = 120_000; +/** Below this, a computed wait is floating-point residue rather than a wait. */ +const MIN_SLEEP_MS = 1; + +/** Spread applied to a synchronised release, so waiters do not wake as one. */ +const SPREAD_MS = 250; + const defaultSleep = (ms: number): Promise => new Promise((r) => setTimeout(r, ms)); +/** + * Milliseconds to wait, or null when the header gives us nothing usable. + * + * NULL AND ZERO ARE DIFFERENT ANSWERS, and conflating them turned the client + * into a hammer. Only `null` reaches the exponential backoff, so a header that + * parsed to 0 -- `Retry-After: 0`, which RFC 9110 permits, or a negative, or + * an already-past date -- meant no wait at all. Measured in the Python + * sibling, which had the same shape: four requests in 3ms against a server + * that had just said 429, and 204 a second across ten threads. + * + * `Number()` was also far too generous for a `delta-seconds = 1*DIGIT` field: + * it read ' ' as 0 and spun, and '0x10' as 16 and slept 48 seconds. + */ function parseRetryAfter(header: string | null): number | null { - if (header === null || header === '') return null; - const asInt = Number(header); - if (Number.isFinite(asInt)) return Math.max(0, asInt * 1000); - const asDate = Date.parse(header); - if (Number.isFinite(asDate)) return Math.max(0, asDate - Date.now()); - return null; + if (header === null) return null; + const trimmed = header.trim(); + if (trimmed === '') return null; + + let ms: number; + if (/^\d+(\.\d+)?$/.test(trimmed)) { + ms = Number(trimmed) * 1000; + } else { + const asDate = Date.parse(trimmed); + if (!Number.isFinite(asDate)) return null; + ms = asDate - Date.now(); + } + return ms > 0 ? ms : null; } function backoff(attempt: number): number { @@ -134,8 +179,80 @@ function formatBodyExcerpt(body: unknown): string | undefined { } export class Transport { + /** What the server last said about the two budgets. */ + rateLimit: RateLimits = UNKNOWN_RATE_LIMITS; + /** + * TypeScript `private`, deliberately NOT an ECMAScript `#` field. + * + * client.ts wraps this transport in a Proxy to add caching, and a Proxy + * forwards methods with `this` bound to the PROXY. A `#` member's brand + * check then fails with "Receiver must be an instance of class Transport", + * which crashed every paged history call on the default configuration -- + * `history.days()` threw on page two for anyone who had not passed + * `cache: false`. Every test passed `cache: false`, so 130 of them missed + * it. A TS `private` compiles to a plain property and forwards fine. + */ + private readonly gate = new Gate(); + constructor(private readonly opts: TransportOptions) {} + /** + * Wait before sending, if we already know this request would fail. + * + * Two reasons, and they are different. The GATE is a 429 the server has + * already issued to this caller: the wait belongs to them, not to whichever + * request met it, so it is shared and taken once. The REMAINING check is a + * window the server told us is spent, where sending is a certain 429 that + * also spends budget being refused. + * + * A remaining we were never told is not a spent one. Anonymous responses + * carry no figures at all, so an unknown must never hold. + */ + private async hold(sleep: (ms: number) => Promise, budget: number): Promise { + // Returns how long it slept, so the caller can keep a running total. The + // TOTAL is what maxRetryAfterMs bounds, not each leg: the gate wait and + // the spent-window wait are both self-initiated holds and they stack. + // Measured before this budget existed: a 429 carrying Retry-After 5 and + // RateLimit-Reset 60 slept 5 then 55, three times over -- 180 seconds + // inside one call whose cap was 120. Each leg was under the cap, so the + // per-leg check never fired and the promise was reachable around. + let spent = 0; + // Re-read the gate after each wait. It slept once and returned, so a + // waiter that woke while someone else's 429 had pushed the gate further + // out sent anyway: measured waking at 584ms with the gate shut for + // another two seconds. Bounded by the same budget, so it cannot spin. + for (;;) { + const before = this.gate.deadline; + const gated = Math.min(this.gate.waitMs(), budget - spent); + if (gated <= MIN_SLEEP_MS) break; + await sleep(gated); + spent += gated; + // Only wait again if someone genuinely pushed the gate further out + // while we slept. Re-reading unconditionally loops against any clock + // that does not advance, which is every test harness and, briefly, a + // suspended host. + if (this.gate.deadline <= before) break; + } + if (this.opts.retry.respectRemaining === false) return spent; + const cap = this.opts.retry.maxRetryAfterMs ?? DEFAULT_MAX_RETRY_AFTER_MS; + for (const meter of [this.rateLimit.rest, this.rateLimit.history]) { + if (!isExhausted(meter)) continue; + const left = secondsUntilReset(meter); + // Past the cap we do not sit on it: the caller gets the 429 and its + // Retry-After and can decide. Same rule the retry path follows. + if (left === null || left <= 0 || left * 1000 > cap) continue; + // Jittered like the gate. Without it every waiter derived `left` from + // the same observedAt and woke at the same absolute instant: measured + // spread 0ms across ten waiters, the tightest burst in the client, on + // the very branch that exists to avoid a 429. + const wait = Math.min(left * 1000 + Math.random() * SPREAD_MS, budget - spent); + if (wait <= MIN_SLEEP_MS) break; + await sleep(wait); + spent += wait; + } + return spent; + } + async get(path: string): Promise { return this.request('GET', this.opts.baseUrl.replace(/\/$/, '') + path); } @@ -156,7 +273,10 @@ export class Transport { const sleep = this.opts.sleep ?? defaultSleep; let attempt = 0; + // One budget for the whole call, because that is what the cap promises. + let budget = this.opts.retry.maxRetryAfterMs ?? DEFAULT_MAX_RETRY_AFTER_MS; while (true) { + budget -= await this.hold(sleep, budget); const controller = new AbortController(); const timer = setTimeout(() => { controller.abort(); @@ -189,6 +309,7 @@ export class Transport { throw new NetworkError(`network error calling ${url}`, { cause }); } clearTimeout(timer); + this.rateLimit = readRateLimits(response.headers, this.rateLimit); if (response.ok) { return (await response.json()) as T; @@ -200,13 +321,33 @@ export class Transport { const retryAfterMs = parseRetryAfter(response.headers.get('retry-after')); const cap = this.opts.retry.maxRetryAfterMs ?? DEFAULT_MAX_RETRY_AFTER_MS; const waitTooLong = retryAfterMs !== null && retryAfterMs > cap; + if ( + response.status === 429 && + this.opts.retry.on429 && + retryAfterMs !== null && + !waitTooLong + ) { + // The wait belongs to the CALLER, not to whichever request met it, so + // it goes on the shared gate and hold() serves it once. Sleeping here + // as well would pay it twice, and ten concurrent requests would each + // pay their own and then retry in unison. + // + // on429: false means "do not wait on a 429", so it must gate the gate + // too. Without this the caller got the error they asked for and then + // their NEXT call silently blocked: an opt-out that does not opt out. + // + // Past the cap the gate is left OPEN on purpose: we throw instead, and + // blocking the caller's next call for most of an hour is the opposite + // of letting them checkpoint and resume. + this.gate.closeFor(retryAfterMs); + } if ( response.status === 429 && this.opts.retry.on429 && attempt < this.opts.retry.max && !waitTooLong ) { - await sleep(retryAfterMs ?? backoff(attempt)); + if (retryAfterMs === null) await sleep(backoff(attempt)); attempt++; continue; } diff --git a/test/unit/history-paging.test.ts b/test/unit/history-paging.test.ts index 0b1b151..8d29e51 100644 --- a/test/unit/history-paging.test.ts +++ b/test/unit/history-paging.test.ts @@ -40,6 +40,31 @@ async function collect(iterable: AsyncIterable): Promise { return out; } +describe('days() paging on the DEFAULT configuration', () => { + // Every other test in this file passes `cache: false`, so the caching + // Proxy in client.ts is never in the call path for getUrl -- and that is + // exactly where paging goes. A `#private` field on Transport made the + // Proxy's forwarded `this` fail its brand check, so page two threw + // "Receiver must be an instance of class Transport" for every real caller + // while 130 tests stayed green. Tests that all take the same non-default + // path cover a shape the users do not have. + it('follows next with the cache ON', async () => { + const page = await loadFixture('mk_history_daily.json'); + let call = 0; + const fetchFn = vi.fn(() => { + call++; + return Promise.resolve( + json({ ...page, next: call === 1 ? 'https://api.themeparks.wiki/v1/p2' : null }), + ); + }); + // No `cache: false`: this is the configuration every real caller gets. + const tp = new ThemeParks({ fetch: fetchFn as unknown as FetchLike }); + const rows = await collect(tp.entity('park-1').history.days()); + expect(fetchFn).toHaveBeenCalledTimes(2); + expect(rows.length).toBeGreaterThan(0); + }); +}); + describe('days() paging', () => { it('follows the server next URL verbatim and stops at null', async () => { const page = await loadFixture('mk_history_daily.json'); @@ -184,7 +209,15 @@ describe('the hourly history budget', () => { expect(error).toBeInstanceOf(RateLimitError); expect(error).not.toBeInstanceOf(BudgetExhaustedError); - expect(slept).toEqual([2000, 2000, 2000]); + // Jittered: the 429 wait is now taken once on a gate shared by the whole + // client, and without a little spread every waiter would wake at the same + // instant and re-trip the limit together. One wait per retry, never + // doubled by the retry path paying it as well. + expect(slept).toHaveLength(3); + for (const ms of slept) { + expect(ms).toBeGreaterThanOrEqual(2000); + expect(ms).toBeLessThan(2300); + } expect(fetchFn).toHaveBeenCalledTimes(4); }); diff --git a/test/unit/ratelimit.test.ts b/test/unit/ratelimit.test.ts new file mode 100644 index 0000000..82b767a --- /dev/null +++ b/test/unit/ratelimit.test.ts @@ -0,0 +1,517 @@ +/** + * Reading the two budgets, and staying inside them. + * + * The SDK read exactly one header, `Retry-After`, and only after a 429 had + * already happened: it could say you had run out, never that you were about + * to. These cover the three things that changed. The figures are read, + * absence is not confused with zero, and the wait a 429 imposes is taken ONCE + * for the whole client rather than once per in-flight request. + */ + +import { describe, it, expect, vi } from 'vitest'; +import { ThemeParks } from '../../src/client'; +import { RateLimitError } from '../../src/errors'; +import { + Gate, + isExhausted, + readRateLimits, + secondsUntilReset, + UNKNOWN_RATE_LIMIT, + UNKNOWN_RATE_LIMITS, +} from '../../src/ratelimit'; +import type { FetchLike } from '../../src/transport'; + +const REST = { + 'RateLimit-Limit': '300', + 'RateLimit-Policy': '300;w=60', + 'RateLimit-Remaining': '299', + 'RateLimit-Reset': '60', +}; +const HISTORY = { + 'RateLimit-History-Limit': '600', + 'RateLimit-History-Policy': '600;w=3600', + 'RateLimit-History-Remaining': '599', + 'RateLimit-History-Reset': '3412', +}; + +function bag(headers: Record) { + return new Headers(headers); +} + +function client(headers: Record, options: Record = {}) { + const fetchFn = vi.fn(() => + Promise.resolve( + new Response(JSON.stringify({ destinations: [] }), { + headers: { 'content-type': 'application/json', ...headers }, + }), + ), + ); + const tp = new ThemeParks({ fetch: fetchFn as unknown as FetchLike, ...options }); + return { tp, fetchFn }; +} + +describe('reading the headers', () => { + it('reads the REST meter', () => { + const { rest } = readRateLimits(bag(REST)); + expect([rest.limit, rest.remaining, rest.reset]).toEqual([300, 299, 60]); + expect(rest.policy).toBe('300;w=60'); + }); + + it('reads the history meter', () => { + const { history } = readRateLimits(bag(HISTORY)); + expect([history.limit, history.remaining, history.reset]).toEqual([600, 599, 3412]); + }); + + it('keeps the two meters apart', () => { + // `ratelimit-history-limit` also begins with `ratelimit-`, so a prefix + // scan reads the hourly figure as the per-minute one and the client paces + // itself against the wrong window. + const out = readRateLimits(bag({ ...REST, ...HISTORY })); + expect(out.rest.limit).toBe(300); + expect(out.history.limit).toBe(600); + expect(out.rest.reset).toBe(60); + expect(out.history.reset).toBe(3412); + }); + + it('does not invent a REST meter from history headers alone', () => { + expect(readRateLimits(bag(HISTORY)).rest.limit).toBeNull(); + }); + + it('leaves what we knew alone when a response mentions neither', () => { + // Most responses mention one meter, or -- if publicly cacheable -- + // neither. Overwriting with blanks would let the last cacheable response + // erase everything the client had learned. + const known = readRateLimits(bag({ ...REST, ...HISTORY })); + const after = readRateLimits(bag({ 'content-type': 'application/json' }), known); + expect(after.rest.remaining).toBe(299); + expect(after.history.remaining).toBe(599); + }); + + it('treats a malformed value as unknown rather than zero', () => { + // A failed parse must not become 0, or the client holds forever waiting + // on a window it invented. + const { rest } = readRateLimits(bag({ ...REST, 'RateLimit-Remaining': 'lots' })); + expect(rest.remaining).toBeNull(); + expect(isExhausted(rest)).toBe(false); + }); +}); + +describe('a cached response says nothing', () => { + // Its figures belong to whoever populated the entry. The server withholds + // the HISTORY figures from anything a shared cache may store, but the + // per-minute ones ride those responses. Measured against production: three + // consecutive calls returning `age: 9` and an unmoving `remaining: 285`. A + // cached `remaining: 0` would make the client sleep out someone else's + // window. + it('ignores a cache hit', () => { + const known = readRateLimits(bag(REST)); + const after = readRateLimits(bag({ ...REST, 'RateLimit-Remaining': '0', Age: '1713' }), known); + expect(after.rest.remaining).toBe(299); + }); + + it('records a fresh response', () => { + // A cache MISS carries no Age at all, which is the path that matters: + // confirmed against production, a MISS returns the figures. + expect(readRateLimits(bag(REST)).rest.remaining).toBe(299); + }); + + it('treats Age: 0 as fresh', () => { + expect(readRateLimits(bag({ ...REST, Age: '0' })).rest.remaining).toBe(299); + }); + + it('does not erase what we knew', () => { + const known = readRateLimits(bag({ ...REST, ...HISTORY })); + const after = readRateLimits(bag({ Age: '60' }), known); + expect(after.rest.remaining).toBe(299); + expect(after.history.remaining).toBe(599); + }); +}); + +describe('absence is not zero', () => { + it('an unknown remaining is not exhausted', () => { + // Anonymous responses carry no figures, because they are publicly + // cacheable and the numbers are per-caller. Reading that as "nothing + // left" would stall every anonymous client permanently. + expect(isExhausted(UNKNOWN_RATE_LIMIT)).toBe(false); + }); + + it('a zero remaining is exhausted', () => { + expect(isExhausted({ ...UNKNOWN_RATE_LIMIT, remaining: 0 })).toBe(true); + }); + + // Both of these pass an explicit `now` rather than reading a clock. The + // function takes one precisely so these can be exact; asserting a range + // around real wall-clock time makes a gate test flaky under CPU load, and + // one of these did flake when the two suites ran concurrently. + it('counts the reset down from when it was read', () => { + const meter = { ...UNKNOWN_RATE_LIMIT, reset: 60, observedAt: 1_000_000 }; + expect(secondsUntilReset(meter, 1_050_000)).toBe(10); + }); + + it('never reports a negative countdown', () => { + const meter = { ...UNKNOWN_RATE_LIMIT, reset: 5, observedAt: 1_000_000 }; + expect(secondsUntilReset(meter, 1_100_000)).toBe(0); + }); + + it('observedAt is monotonic, not an epoch timestamp', () => { + // Date.now() in the Gate made a backward clock step turn a 5s wait into + // an hour, silently, with nothing bounding it: the cap is applied when + // the gate is armed, never when it is served. + const { rest } = readRateLimits(bag(REST)); + expect(rest.observedAt).not.toBeNull(); + // An epoch reading would be ~1.8e12; a monotonic one is process uptime. + expect(rest.observedAt!).toBeLessThan(1e11); + }); + + it('has no countdown without a reset', () => { + expect(secondsUntilReset(UNKNOWN_RATE_LIMITS.rest)).toBeNull(); + }); +}); + +describe('the gate', () => { + it('costs nothing while open', () => { + expect(new Gate().waitMs()).toBe(0); + }); + + it('holds every caller once closed', () => { + const gate = new Gate(); + gate.closeFor(5000); + expect(gate.waitMs()).toBeGreaterThan(4000); + expect(gate.waitMs()).toBeGreaterThan(4000); + }); + + it('jitters waiters so they do not wake together', () => { + // Waking in unison is the other half of the thundering herd: the wait is + // shared, then everyone retries at the same instant and re-trips it. + // + // Measured against a FROZEN clock. With a live one this asserted only + // that time passes between calls -- it stayed green with the jitter term + // deleted, which is a test that cannot fail. + let clock = 1_000_000; + vi.spyOn(performance, 'now').mockImplementation(() => clock); + try { + const gate = new Gate(); + gate.closeFor(5000); + const waits = new Set(Array.from({ length: 20 }, () => gate.waitMs())); + expect(waits.size).toBeGreaterThan(1); + // And the spread is bounded, so it cannot be mistaken for the wait. + for (const w of waits) { + expect(w).toBeGreaterThanOrEqual(5000); + expect(w).toBeLessThan(5300); + } + } finally { + vi.mocked(performance.now).mockRestore(); + void clock; + } + }); + + it('is never brought forward by a shorter wait', () => { + // A 2s Retry-After arriving while a 60s one is in force would otherwise + // release the herd early. + const gate = new Gate(); + gate.closeFor(60_000); + gate.closeFor(2000); + expect(gate.waitMs()).toBeGreaterThan(55_000); + }); +}); + +describe('the opt-outs actually opt out', () => { + // An advertised switch that switches nothing is worse than no switch. + // `on429: false` threw the error the caller asked for and then closed the + // shared gate anyway, so their NEXT call blocked for the full Retry-After. + // The setting says "do not wait on a 429"; the gate is a wait on a 429. + function limited(retry: Record) { + const slept: number[] = []; + const fetchFn = vi.fn(() => + Promise.resolve( + new Response(JSON.stringify({}), { + status: 429, + headers: { 'content-type': 'application/json', 'retry-after': '45' }, + }), + ), + ); + const tp = new ThemeParks({ fetch: fetchFn as unknown as FetchLike, cache: false, retry }); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + return Promise.resolve(); + }; + return { tp, slept }; + } + + it('on429 false never sleeps, even on a later call', async () => { + const { tp, slept } = limited({ on429: false }); + for (let i = 0; i < 3; i++) { + await expect(tp.destinations.list()).rejects.toBeInstanceOf(RateLimitError); + } + expect(slept).toEqual([]); + }); + + it('on429 true still holds the gate', async () => { + // The opt-out must not have disabled the feature for everyone else. + const { tp, slept } = limited({ max: 0 }); + await expect(tp.destinations.list()).rejects.toBeInstanceOf(RateLimitError); + await expect(tp.destinations.list()).rejects.toBeInstanceOf(RateLimitError); + expect(slept.length).toBeGreaterThan(0); + }); +}); + +describe('through the client', () => { + it('records both meters from a real call', async () => { + const { tp } = client({ ...REST, ...HISTORY }); + await tp.destinations.list(); + expect(tp.rateLimit.rest.remaining).toBe(299); + expect(tp.rateLimit.history.remaining).toBe(599); + }); + + it('knows nothing before the first call', () => { + const { tp } = client(REST); + expect(tp.rateLimit.rest.limit).toBeNull(); + expect(isExhausted(tp.rateLimit.rest)).toBe(false); + }); + + it('survives the cache wrapper, which is the default path', async () => { + // Caching is ON by default and wraps the transport. In the Python sibling + // every test used cache:false and the first real call against production + // threw, because the wrapper had no rateLimit to forward. + const { tp, fetchFn } = client(REST); // cache default: on + await tp.destinations.list(); + await tp.destinations.list(); + expect(fetchFn).toHaveBeenCalledOnce(); // second was a cache hit + // A hit sends no request and so learns nothing, which is right: it spent + // no budget either, so the previous figures still stand. + expect(tp.rateLimit.rest.remaining).toBe(299); + }); + + it('waits out a window the server said is spent', async () => { + const slept: number[] = []; + const { tp, fetchFn } = client( + { ...REST, 'RateLimit-Remaining': '0', 'RateLimit-Reset': '7' }, + { cache: false }, + ); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + return Promise.resolve(); + }; + await tp.destinations.list(); // learns remaining 0 + await tp.destinations.list(); // must hold first + // Both bounds. Upper-only let `sleep(left)` through in place of + // `sleep(left * 1000)` -- 7 milliseconds instead of 7 seconds, a + // thousandfold too short, passing green. + expect(slept.length).toBeGreaterThan(0); + expect(slept[0]).toBeGreaterThan(6000); + // Upper bound allows the spread now applied to this path: without it all + // waiters woke inside the same millisecond. + expect(slept[0]).toBeLessThanOrEqual(7000 + 250); + expect(fetchFn).toHaveBeenCalledTimes(2); + }); + + it('never holds on an unknown remaining', async () => { + const slept: number[] = []; + const { tp } = client({ 'content-type': 'application/json' }, { cache: false }); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + return Promise.resolve(); + }; + await tp.destinations.list(); + await tp.destinations.list(); + expect(slept).toEqual([]); + }); + + it('can be told not to respect remaining', async () => { + const slept: number[] = []; + const { tp } = client( + { ...REST, 'RateLimit-Remaining': '0', 'RateLimit-Reset': '7' }, + { cache: false, retry: { respectRemaining: false } }, + ); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + return Promise.resolve(); + }; + await tp.destinations.list(); + await tp.destinations.list(); + expect(slept).toEqual([]); + }); +}); + +describe('the cap bounds the whole call', () => { + // Two self-initiated holds exist -- the shared gate and the spent-window + // wait -- and they stack. A 429 carrying BOTH a Retry-After and + // RateLimit-Remaining: 0 slept 5s at the gate then 55s for the window, + // three times over: 180 seconds inside one call whose cap was 120. Each leg + // was under the cap, so the per-leg check never fired. + // + // Nothing caught it because no test sent a 429 carrying rate-limit headers, + // and a fake sleep that does not advance the clock cannot show a cumulative + // total at all. This advances one, which is what the real world does. + const REAL_429 = { + 'retry-after': '5', + 'RateLimit-Limit': '300', + 'RateLimit-Remaining': '0', + 'RateLimit-Reset': '60', + }; + + async function run(capMs: number) { + const slept: number[] = []; + let clock = 1_000_000; + const original = performance.now.bind(performance); + vi.spyOn(performance, 'now').mockImplementation(() => clock); + try { + const fetchFn = vi.fn(() => + Promise.resolve( + new Response(JSON.stringify({}), { + status: 429, + headers: { 'content-type': 'application/json', ...REAL_429 }, + }), + ), + ); + const tp = new ThemeParks({ + fetch: fetchFn as unknown as FetchLike, + cache: false, + retry: { maxRetryAfterMs: capMs }, + }); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + clock += ms; + return Promise.resolve(); + }; + await expect(tp.destinations.list()).rejects.toBeInstanceOf(RateLimitError); + return slept; + } finally { + vi.mocked(performance.now).mockRestore(); + void original; + } + } + + it('never exceeds the cap in total', async () => { + const slept = await run(120_000); + const total = slept.reduce((a, b) => a + b, 0); + expect(total).toBeLessThanOrEqual(120_000 + 10); + }); + + it('binds harder with a smaller cap', async () => { + const slept = await run(30_000); + const total = slept.reduce((a, b) => a + b, 0); + expect(total).toBeLessThanOrEqual(30_000 + 10); + }); + + it('still waits when there is budget', async () => { + // The cap must bound the feature, not disable it. + const slept = await run(120_000); + const total = slept.reduce((a, b) => a + b, 0); + expect(total).toBeGreaterThan(5000); + }); + + it('emits no meaningless micro-sleeps', async () => { + const slept = await run(120_000); + for (const ms of slept) expect(ms).toBeGreaterThan(1); + }); +}); + +describe('the waits are actually awaited', () => { + // The suite could prove a sleep was REQUESTED and never that it was waited + // on: dropping the `await` in front of it passed every test, because a fake + // sleep that resolves synchronously advances the world either way. A + // dropped await would make the whole feature a no-op that still looks busy. + // + // So this sleep resolves on a later macrotask and counts itself as pending. + // A request that starts while a sleep is pending means the hold was not + // awaited. + it('no request is sent while a hold is still pending', async () => { + let pending = 0; + const violations: string[] = []; + + const fetchFn = vi.fn(() => { + if (pending > 0) violations.push('fetch started with a sleep in flight'); + return Promise.resolve( + new Response(JSON.stringify({}), { + status: 429, + headers: { + 'content-type': 'application/json', + 'retry-after': '5', + 'RateLimit-Limit': '300', + 'RateLimit-Remaining': '0', + 'RateLimit-Reset': '60', + }, + }), + ); + }); + + const tp = new ThemeParks({ fetch: fetchFn as unknown as FetchLike, cache: false }); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = () => { + pending++; + return new Promise((resolve) => { + setTimeout(() => { + pending--; + resolve(); + }, 0); + }); + }; + + await expect(tp.destinations.list()).rejects.toBeInstanceOf(RateLimitError); + expect(violations).toEqual([]); + expect(fetchFn.mock.calls.length).toBeGreaterThan(1); + }); +}); + +describe('waiters do not wake as one', () => { + // The gate exists for concurrency and nothing measured concurrency. Both + // release paths are covered, because the spent-window one had NO spread: + // ten waiters derived the same deadline from the same observedAt and left + // inside the same millisecond -- the tightest burst in the client, on the + // branch that exists to avoid a 429. + it('spreads the spent-window release', async () => { + // Against a FROZEN clock. With a live one `left` varies by itself as time + // passes, so the set is distinct with or without spread and the test + // measures nothing. + let clock = 1_000_000; + vi.spyOn(performance, 'now').mockImplementation(() => clock); + try { + const seen = new Set(); + for (let i = 0; i < 20; i++) { + const slept: number[] = []; + const fetchFn = vi.fn(() => + Promise.resolve( + new Response(JSON.stringify({ destinations: [] }), { + headers: { + 'content-type': 'application/json', + ...REST, + 'RateLimit-Remaining': '0', + 'RateLimit-Reset': '7', + }, + }), + ), + ); + const tp = new ThemeParks({ fetch: fetchFn as unknown as FetchLike, cache: false }); + ( + tp as unknown as { transport: { opts: { sleep: (ms: number) => Promise } } } + ).transport.opts.sleep = (ms: number) => { + slept.push(ms); + return Promise.resolve(); + }; + await tp.destinations.list(); + await tp.destinations.list(); + if (slept.length > 0) seen.add(Math.round(slept[0]!)); + } + expect(seen.size).toBeGreaterThan(1); + for (const ms of seen) { + expect(ms).toBeGreaterThanOrEqual(7000); + expect(ms).toBeLessThan(7300); + } + } finally { + vi.mocked(performance.now).mockRestore(); + void clock; + } + }); +});