// THE POLITE archive.org CLIENT — every request this app makes to archive.org // itself: the item metadata API, a mirror's info.json, the item's torrent, and // a file's bytes when it is not fetched over BitTorrent (the torrent's own web // seeding is aria2c's, lib/archiveOrgTorrent-server.ts). // // archive.org is a non-profit serving files off its own disks; the rules here // are the operator's "be polite to archive.org", made mechanical: // // ONE AT A TIME requests are chained: a second caller waits for the first // to finish, and then for the gap. // A GAP at least `minGapMs` (2 s) between the end of one request // and the start of the next, from this process. // IDENTIFIED a User-Agent naming the project and its URL. // BACKS OFF a 429 or 503 (and 502/504, and a dropped connection) is // retried after the server's Retry-After when it sends one, // else after an exponential wait (5 s, 10 s, 20 s… capped at // 2 min, ±20 % jitter), at most `maxAttempts` times in all; // then it stops with ArchiveOrgRequestError, `rateLimited`. // ASKS ONCE an item's metadata, and its torrent, are cached in // memory for `cacheTtlMs` (6 h): a bulk import of 161 files // of one item asks for each once, and so does every // per-file provenance step after it. // // A FILE DOWNLOAD (`downloadFile`, the fallback when a torrent cannot be used) // is one plain HTTP stream: it waits out the same gap before it starts, but it // does not hold the chain for its minutes-long body — a metadata request from // another caller (the editor resolving an import URL) is not queued behind it. // The downloads already go one at a time, on the `platform:archiveorg` job // queue. A dropped stream resumes with a Range request from the bytes on disk; // a 429/503 waits as above. import { PROJECT_NAME, PROJECT_URL } from "./project"; import { parseArchiveOrgItemMetadata, type ArchiveOrgItemMetadata, } from "./archiveOrg"; import { archiveOrgDownloadUrl, archiveOrgMetadataUrl, archiveOrgTorrentUrl, } from "./archiveOrgId"; import { createWriteStream } from "node:fs"; import { rename, rm, stat } from "node:fs/promises"; export const ARCHIVE_ORG_USER_AGENT = `${PROJECT_NAME} archive.org import (+${PROJECT_URL})`; export class ArchiveOrgRequestError extends Error { readonly status: number | null; readonly rateLimited: boolean; constructor(message: string, status: number | null, rateLimited: boolean) { super(message); this.name = "ArchiveOrgRequestError"; this.status = status; this.rateLimited = rateLimited; } } export type ArchiveOrgClientDeps = { fetch?: (url: string, init: RequestInit) => Promise; now?: () => number; sleep?: (ms: number, signal?: AbortSignal) => Promise; random?: () => number; }; export type ArchiveOrgClientOpts = { minGapMs?: number; maxAttempts?: number; baseBackoffMs?: number; maxBackoffMs?: number; cacheTtlMs?: number; timeoutMs?: number; // A file download that receives no bytes for this long is dropped (and // retried, resuming). idleTimeoutMs?: number; }; export type DownloadFileProgress = { bytes: number; total: number | null }; export type DownloadFileOpts = { signal?: AbortSignal; // The size archive.org lists for the file: a partial of exactly this size is // already complete. expectedSize?: number; onProgress?: (p: DownloadFileProgress) => void; }; const RETRYABLE = new Set([429, 502, 503, 504]); function defaultSleep(ms: number, signal?: AbortSignal): Promise { return new Promise((resolve, reject) => { if (signal?.aborted) return reject(signal.reason ?? new Error("aborted")); const t = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); const onAbort = () => { clearTimeout(t); reject(signal?.reason ?? new Error("aborted")); }; signal?.addEventListener("abort", onAbort, { once: true }); }); } // Retry-After: delta-seconds or an HTTP date. Null when absent or unreadable. export function parseRetryAfterMs(value: string | null, nowMs: number): number | null { if (!value) return null; const v = value.trim(); if (/^\d+$/.test(v)) return Number(v) * 1000; const at = Date.parse(v); if (Number.isFinite(at)) return Math.max(0, at - nowMs); return null; } export class ArchiveOrgClient { private readonly deps: Required; private readonly opts: Required; private chain: Promise = Promise.resolve(); private lastDoneAt = -Infinity; private readonly cache = new Map(); private readonly torrentCache = new Map(); // Requests actually sent (each attempt), for tests and logs. requests = 0; constructor(deps: ArchiveOrgClientDeps = {}, opts: ArchiveOrgClientOpts = {}) { this.deps = { fetch: deps.fetch ?? ((url, init) => fetch(url, init)), now: deps.now ?? (() => Date.now()), sleep: deps.sleep ?? defaultSleep, random: deps.random ?? Math.random, }; this.opts = { minGapMs: opts.minGapMs ?? 2_000, maxAttempts: opts.maxAttempts ?? 4, baseBackoffMs: opts.baseBackoffMs ?? 5_000, maxBackoffMs: opts.maxBackoffMs ?? 120_000, cacheTtlMs: opts.cacheTtlMs ?? 6 * 60 * 60 * 1000, timeoutMs: opts.timeoutMs ?? 60_000, idleTimeoutMs: opts.idleTimeoutMs ?? 120_000, }; } // Run `fn` after every earlier request (and its gap) has finished. private serial(fn: () => Promise): Promise { const run = this.chain.then(fn, fn); this.chain = run.catch(() => {}); return run; } private backoffMs(attempt: number): number { const exp = Math.min(this.opts.maxBackoffMs, this.opts.baseBackoffMs * 2 ** (attempt - 1)); const jitter = 1 + (this.deps.random() * 0.4 - 0.2); return Math.round(exp * jitter); } // One GET, with the gap before it and the retry policy around it. private async getOnce(url: string, accept: string, signal?: AbortSignal): Promise { let lastError = ""; for (let attempt = 1; attempt <= this.opts.maxAttempts; attempt++) { const wait = this.lastDoneAt + this.opts.minGapMs - this.deps.now(); if (wait > 0) await this.deps.sleep(wait, signal); this.requests++; let res: Response | null = null; try { const timeout = AbortSignal.timeout(this.opts.timeoutMs); res = await this.deps.fetch(url, { signal: signal ? AbortSignal.any([signal, timeout]) : timeout, redirect: "follow", headers: { accept, "user-agent": ARCHIVE_ORG_USER_AGENT }, }); } catch (err) { if (signal?.aborted) throw err; lastError = (err as Error).message; } finally { this.lastDoneAt = this.deps.now(); } if (res && res.ok) return res; if (res && !RETRYABLE.has(res.status)) { throw new ArchiveOrgRequestError( `archive.org answered HTTP ${res.status} for ${url}`, res.status, false, ); } if (res) lastError = `HTTP ${res.status}`; if (attempt === this.opts.maxAttempts) break; const retryAfter = res ? parseRetryAfterMs(res.headers.get("retry-after"), this.deps.now()) : null; const delay = Math.min(this.opts.maxBackoffMs, retryAfter ?? this.backoffMs(attempt)); await this.deps.sleep(delay, signal); } throw new ArchiveOrgRequestError( `archive.org did not answer ${url} after ${this.opts.maxAttempts} attempts (${lastError}); stopping`, null, true, ); } getJson(url: string, signal?: AbortSignal): Promise { return this.serial(async () => { const res = await this.getOnce(url, "application/json", signal); return res.json(); }); } // The item's metadata, once per `cacheTtlMs`. async itemMetadata(identifier: string, signal?: AbortSignal): Promise { const hit = this.cache.get(identifier); if (hit && this.deps.now() - hit.at < this.opts.cacheTtlMs) return hit.value; const raw = await this.getJson(archiveOrgMetadataUrl(identifier), signal); const item = parseArchiveOrgItemMetadata(raw); if (!item) { throw new ArchiveOrgRequestError(`archive.org has no item "${identifier}"`, 404, false); } this.cache.set(identifier, { at: this.deps.now(), value: item }); return item; } // A small JSON file inside an item (a mirror's `.info.json`). itemJsonFile(identifier: string, file: string, signal?: AbortSignal): Promise { return this.getJson(archiveOrgDownloadUrl(identifier, file), signal); } // The item's `_archive.torrent`, once per `cacheTtlMs`. A torrent // is a few kilobytes per file; a 160-file item's is still well under a // megabyte. async itemTorrent(identifier: string, signal?: AbortSignal): Promise { const hit = this.torrentCache.get(identifier); if (hit && this.deps.now() - hit.at < this.opts.cacheTtlMs) return hit.value; const buf = await this.serial(async () => { const res = await this.getOnce(archiveOrgTorrentUrl(identifier), "application/x-bittorrent", signal); return Buffer.from(await res.arrayBuffer()); }); this.torrentCache.set(identifier, { at: this.deps.now(), value: buf }); return buf; } // ONE FILE'S BYTES into `dest`, by https://archive.org/download//. // Written to `.part` and renamed when whole; a `.part` already there is // resumed with a Range request (a server that ignores the Range sends the // whole file, which then replaces it). Retries as `getOnce` does — a 429/503 // after its Retry-After, a dropped connection or an idle stream after the // exponential wait, each attempt resuming — and gives up with // ArchiveOrgRequestError (`rateLimited` when archive.org kept refusing). async downloadFile( identifier: string, file: string, dest: string, opts: DownloadFileOpts = {}, ): Promise<{ bytes: number }> { const url = archiveOrgDownloadUrl(identifier, file); const part = `${dest}.part`; const signal = opts.signal; let lastError = ""; // Whether the LAST failed attempt was archive.org refusing (429/5xx), as // opposed to a dropped stream: only that is a rate limit. let refused = false; for (let attempt = 1; attempt <= this.opts.maxAttempts; attempt++) { const have = await stat(part).then((s) => s.size, () => 0); if (opts.expectedSize !== undefined && have === opts.expectedSize && have > 0) { await rename(part, dest); return { bytes: have }; } const wait = this.lastDoneAt + this.opts.minGapMs - this.deps.now(); if (wait > 0) await this.deps.sleep(wait, signal); this.requests++; const ctl = new AbortController(); const onAbort = () => ctl.abort(signal?.reason); signal?.addEventListener("abort", onAbort, { once: true }); let idle: NodeJS.Timeout | null = null; const arm = () => { if (idle) clearTimeout(idle); idle = setTimeout( () => ctl.abort(new Error(`no data for ${Math.round(this.opts.idleTimeoutMs / 1000)}s`)), this.opts.idleTimeoutMs, ); }; let res: Response | null = null; try { arm(); res = await this.deps.fetch(url, { signal: ctl.signal, redirect: "follow", headers: { accept: "*/*", "user-agent": ARCHIVE_ORG_USER_AGENT, ...(have > 0 ? { range: `bytes=${have}-` } : {}), }, }); if (res.status === 416 && have > 0) { // The partial is not a prefix of what archive.org has now: start over. await rm(part, { force: true }); lastError = "HTTP 416 (partial discarded)"; continue; } if (!res.ok) { if (!RETRYABLE.has(res.status)) { throw new ArchiveOrgRequestError( `archive.org answered HTTP ${res.status} for ${url}`, res.status, false, ); } refused = true; lastError = `HTTP ${res.status}`; } else { const append = res.status === 206 && have > 0; const lengthHeader = Number(res.headers.get("content-length")); const total = Number.isFinite(lengthHeader) && lengthHeader > 0 ? (append ? have : 0) + lengthHeader : (opts.expectedSize ?? null); let bytes = append ? have : 0; const out = createWriteStream(part, { flags: append ? "a" : "w" }); try { const reader = res.body?.getReader(); if (!reader) throw new Error("empty response body"); while (true) { const { done, value } = await reader.read(); if (done) break; arm(); bytes += value.byteLength; if (!out.write(value)) await new Promise((r) => out.once("drain", () => r())); opts.onProgress?.({ bytes, total }); } } finally { await new Promise((resolve, reject) => out.end((err?: Error | null) => (err ? reject(err) : resolve())), ); } if (total !== null && bytes < total) { refused = false; lastError = `stream ended at ${bytes} of ${total} bytes`; } else { await rename(part, dest); return { bytes }; } } } catch (err) { if (signal?.aborted) throw err; if (err instanceof ArchiveOrgRequestError) throw err; refused = false; lastError = (err as Error).message; } finally { if (idle) clearTimeout(idle); signal?.removeEventListener("abort", onAbort); this.lastDoneAt = this.deps.now(); } if (attempt === this.opts.maxAttempts) break; const retryAfter = res && !res.ok ? parseRetryAfterMs(res.headers.get("retry-after"), this.deps.now()) : null; const delay = Math.min(this.opts.maxBackoffMs, retryAfter ?? this.backoffMs(attempt)); await this.deps.sleep(delay, signal); } throw new ArchiveOrgRequestError( `archive.org did not deliver ${url} after ${this.opts.maxAttempts} attempts (${lastError}); stopping`, null, refused, ); } } // THE process's client: one chain, one gap, one cache for every caller. export const archiveOrgClient = new ArchiveOrgClient();