/** * @copyright Sister Software * @license AGPL-3.0 * @author Teffen Ellis, et al. * * Download and manifest utilities for `mailwoman corpus fetch `. */ import { openWriteStream, pipeline, Readable } from "@mailwoman/core/fs/streams" import { writeLocalFile } from "@mailwoman/core/fs/writers" import { sleep } from "@mailwoman/core/utils/sleep" import type { PathBuilderLike } from "path-ts" /** * Rate limited — retryable, the server is asking us to back off. */ const HTTP_TOO_MANY_REQUESTS = 429 /** * Lowest 5xx status. Server-side failures are retryable; 4xx are not. */ const HTTP_SERVER_ERROR_MIN = 500 /** * Highest 5xx status. */ const HTTP_SERVER_ERROR_MAX = 599 /** * The option base every `mailwoman corpus fetch ` module extends. */ /** * How long a failed transfer waits before the next attempt. */ export const DEFAULT_RETRY_DELAY_MS = 5000 export interface BaseFetchOptions { /** * Destination root for downloaded source data. Each source writes its own subdirectory. */ outRoot: PathBuilderLike /** * Pause between transfer retries, in milliseconds. Defaults to {@linkcode DEFAULT_RETRY_DELAY_MS}. * * A test that exercises the failure path pays this delay once per retry in real time — measured at 20.1 s for the two * failing-transfer cases in `geonames-postal.test.ts`, which is the whole cost of that file. The retry COUNT is the * behaviour under test there; the pause between attempts is not, so it is a caller's to shorten. */ retryDelayMs?: number } /** * The per-run result every fetch module returns; the command maps `failed > 0` to exit code 1. */ export interface FetchSummary { fetched: number skipped: number failed: number failedCodes: string[] } /** * The sibling `MANIFEST.json` shape the single-file fetch modules write: origin URL + fetch timestamp + byte count + * sha256, so downstream adapters can verify provenance. */ export interface SourceManifest { source_url: string downloaded_at: string filename: string sha256: string bytes: number } /** * A status worth retrying: rate limiting or a server-side failure. */ /** * An HTTP failure that CARRIES its status, so callers branch on `error.status` rather than on message prose. The prose * route shipped a real flake: a caller classified "not published upstream" with `message.includes("404")`, and the * message contains the URL — an ephemeral test-server port such as `:40453` satisfies it while the actual status is 500. * Roughly 1–2% of ephemeral ports contain the substring, which is exactly the kind of sometimes-failure that burns a CI * run and vanishes locally. */ export class HTTPStatusError extends Error { readonly status: number constructor(status: number, message: string) { super(message) this.name = "HTTPStatusError" this.status = status } } export function isTransientStatus(status: number): boolean { return status === HTTP_TOO_MANY_REQUESTS || (status >= HTTP_SERVER_ERROR_MIN && status <= HTTP_SERVER_ERROR_MAX) } export interface DownloadOptions { url: string dest: string /** * Per-attempt timeout. Default 10 minutes — these are multi-GB government dumps. */ timeoutMs?: number /** * Extra attempts after the first, taken only on transient statuses or network errors. Default 0. */ retries?: number /** * Delay between attempts. Default 5s. */ retryDelayMs?: number headers?: Record report?: (line: string) => void } /** * Download `url` to `dest` with per-attempt timeout and transient-status retry. Throws on a non-transient HTTP status * or once retries are exhausted. Returns the byte count written. */ export async function downloadToFile(options: DownloadOptions): Promise<{ bytes: number }> { const { url, dest, timeoutMs = 600_000, retries = 0, retryDelayMs = DEFAULT_RETRY_DELAY_MS, headers, report, } = options let lastError: unknown for (let attempt = 0; attempt <= retries; attempt++) { if (attempt > 0) { report?.(`retry ${attempt}/${retries} after ${retryDelayMs}ms — ${url}`) await sleep(retryDelayMs) } let res: Response try { // Raw `fetch`, deliberately: this is the shared FILE downloader and the body is piped to disk below. // `APIClient` is the repo default for API requests — small bodies, repeated calls — and buffers a // non-stream response in memory, which is the one thing a multi-gigabyte transfer must not do. res = await fetch(url, { headers, signal: AbortSignal.timeout(timeoutMs) }) } catch (error) { // AbortSignal timeouts and network-level failures are retryable. lastError = error continue } if (!res.ok) { const error = new HTTPStatusError(res.status, `HTTP ${res.status} ${res.statusText} — ${url}`) if (!isTransientStatus(res.status)) throw error lastError = error continue } try { const buffer = Buffer.from(await res.arrayBuffer()) await writeLocalFile(buffer, dest) return { bytes: buffer.byteLength } } catch (error) { // A mid-stream abort while reading the body is retryable too. lastError = error } } throw lastError instanceof Error ? lastError : new Error(String(lastError)) } export interface StreamDownloadOptions { headers?: Record timeoutMs: number retries: number retryDelayMs: number } /** * Stream an HTTP download to disk, returning the final HTTP status (0 on network error after retries). Follows * redirects (the Census and OpenAddresses endpoints both 302 to their real hosts). * * Kept separate from {@link downloadToFile} on purpose: this one streams a multi-GB body to disk (the buffered helper * reads via `arrayBuffer()`) and returns the HTTP status instead of throwing, which the per-file result collectors and * two-URL fallback ladders consume. */ export async function streamDownload(url: string, dest: string, opts: StreamDownloadOptions): Promise { for (let attempt = 0; attempt <= opts.retries; attempt++) { try { const res = await fetch(url, { headers: opts.headers ?? {}, redirect: "follow", signal: AbortSignal.timeout(opts.timeoutMs), }) if (res.ok && res.body) { await pipeline(Readable.fromWeb(res.body), openWriteStream(dest)) return res.status } if (attempt < opts.retries && isTransientStatus(res.status)) { await sleep(opts.retryDelayMs) continue } return res.status } catch { if (attempt < opts.retries) { await sleep(opts.retryDelayMs) continue } return 0 } } return 0 }