// archive.org OVER BITTORRENT — the process side: one aria2c, run for ONE file // of an item, watched until the file is complete and then for as long as it // seeds (lib/archiveOrgTorrent.ts has the command line and the line parser). // // THE LIFE OF ONE RUN // 1. the item's torrent is written beside the staging dir, and a two-line // `--on-bt-download-complete` hook that touches `.complete` there; // 2. aria2c starts in its OWN PROCESS GROUP (detached), so every child it // has is killed with it, by pgid, never by name; // 3. DOWNLOADING: its progress summary is read every few seconds; a human // line ("torrent: (n of m pieces, peers p, web seed yes)") is // logged every `progressLogMs`; no growth in the completed bytes for // `stallMinutes` stops it — the caller then downloads directly; // 4. COMPLETE: the hook fired (or aria2c says SEED, or it exited 0): the // run is a success from here on whatever seeding does; // 5. SEEDING: aria2c's own --seed-time / --seed-ratio end it ("seeding // 10 min…"); a cancel ends it sooner. A seeding aria2c that does not exit // within its seed time plus a grace period is stopped. // // A CANCEL (the job's signal) kills the group: SIGTERM, then SIGKILL after // `killGraceMs`. A partial and its .aria2 control file stay in the staging dir, // so the next run resumes. import { spawn, execFile, type ChildProcess } from "node:child_process"; import path from "node:path"; import { chmod, mkdir, rm, stat, writeFile } from "node:fs/promises"; import { aria2cArgs, parseAria2cStatusLine, torrentFileSegments, torrentFilePieces, torrentProgressLine, type ArchiveOrgFetchSettings, type ParsedTorrent, type TorrentFileEntry, } from "./archiveOrgTorrent"; // ─── Is aria2c here? ─── const probes = new Map>(); // `aria2c --version` once per binary per process. ENOENT (or any failure to // run) is "not installed". export function aria2cAvailable(bin: string): Promise { let p = probes.get(bin); if (!p) { p = new Promise((resolve) => { execFile(bin, ["--version"], { timeout: 5_000 }, (err, stdout) => { resolve(!err && /aria2/i.test(String(stdout))); }); }); probes.set(bin, p); } return p; } // Tests only: forget a probe. export function __resetAria2cProbeForTest(): void { probes.clear(); } // ─── One run ─── export type TorrentFetchResult = | { ok: true; file: string; seeded: boolean } | { ok: false; reason: string; stalled?: boolean; cancelled?: boolean }; export type TorrentFetchOpts = { aria2cBin: string; torrent: Buffer; parsed: ParsedTorrent; entry: TorrentFileEntry; // aria2c's --dir. Created; the caller removes it when it is done with it. stagingDir: string; settings: ArchiveOrgFetchSettings; userAgent: string; onLog: (line: string) => void; signal: AbortSignal; // Test seams. tickMs?: number; progressLogMs?: number; summaryIntervalSec?: number; killGraceMs?: number; // How long past --seed-time a seeding aria2c may run before it is stopped. seedGraceMs?: number; // Overrides settings.stallMinutes (tests). stallMs?: number; now?: () => number; }; const COMPLETE_MARKER = ".complete"; const HOOK_NAME = "on-complete.sh"; const TORRENT_NAME = "item.torrent"; function killGroup(child: ChildProcess, sig: NodeJS.Signals): void { if (child.pid === undefined || child.exitCode !== null || child.signalCode !== null) return; try { // Negative pid: the whole process group aria2c leads (spawned detached). process.kill(-child.pid, sig); } catch { try { child.kill(sig); } catch { /* already gone */ } } } async function exists(p: string): Promise { return stat(p).then( () => true, () => false, ); } export async function fetchFileByTorrent(opts: TorrentFetchOpts): Promise { const now = opts.now ?? (() => Date.now()); const tickMs = opts.tickMs ?? 2_000; const progressLogMs = opts.progressLogMs ?? 30_000; const killGraceMs = opts.killGraceMs ?? 10_000; const stallMs = opts.stallMs ?? Math.max(1, opts.settings.stallMinutes) * 60_000; const seedMs = Math.max(0, opts.settings.seedMinutes) * 60_000; const seedGraceMs = opts.seedGraceMs ?? 5 * 60_000; const log = (s: string) => opts.onLog(s.endsWith("\n") ? s : `${s}\n`); if (opts.signal.aborted) return { ok: false, reason: "cancelled", cancelled: true }; await mkdir(opts.stagingDir, { recursive: true }); const torrentPath = path.join(opts.stagingDir, TORRENT_NAME); const marker = path.join(opts.stagingDir, COMPLETE_MARKER); const hook = path.join(opts.stagingDir, HOOK_NAME); await writeFile(torrentPath, opts.torrent); await rm(marker, { force: true }); // aria2c runs the hook as `hook `. It only has // to say "complete"; the file's path is known already. await writeFile(hook, `#!/bin/sh\n: > "$(dirname "$0")/${COMPLETE_MARKER}"\n`); await chmod(hook, 0o755); const target = path.join(opts.stagingDir, ...torrentFileSegments(opts.parsed, opts.entry)); const pieces = torrentFilePieces(opts.parsed, opts.entry).count; const webSeed = opts.parsed.webSeeds.length > 0; const label = opts.entry.path; const args = aria2cArgs({ torrentPath, dir: opts.stagingDir, fileIndex: opts.entry.index, settings: opts.settings, onCompleteHook: hook, userAgent: opts.userAgent, parentPid: process.pid, summaryIntervalSec: opts.summaryIntervalSec ?? 5, resume: (await exists(target)) || (await exists(`${target}.aria2`)), }); log( `torrent: fetching ${label} (${pieces} pieces of ${opts.parsed.pieceLength} bytes, ` + `web seed ${webSeed ? "yes" : "no"}, ${opts.parsed.trackers.length} trackers) with aria2c.`, ); log(`$ ${opts.aria2cBin} ${args.join(" ")}`); const child = spawn(opts.aria2cBin, args, { cwd: opts.stagingDir, detached: true, stdio: ["ignore", "pipe", "pipe"], }); let phase = "downloading" as "downloading" | "seeding"; let completedBytes = -1; let lastProgressAt = now(); let lastLogAt = 0; let seedingSince = 0; let lastLines: string[] = []; type StopReason = { reason: string; stalled?: boolean; cancelled?: boolean }; // A holder, not a `let`: it is set from callbacks, where TypeScript cannot // see the assignment. const stopped: { why: StopReason | null } = { why: null }; let killTimer: NodeJS.Timeout | null = null; const enterSeeding = () => { if (phase === "seeding") return; phase = "seeding"; seedingSince = now(); if (seedMs > 0) { log( `torrent: ${label} complete; seeding ${opts.settings.seedMinutes} min` + `${opts.settings.seedRatio > 0 ? ` or to ratio ${opts.settings.seedRatio}` : ""}, whichever comes first…`, ); } else { log(`torrent: ${label} complete; not seeding (archiveOrg.seedMinutes is 0).`); } }; const stop = (why: StopReason | null) => { if (why && !stopped.why) stopped.why = why; killGroup(child, "SIGTERM"); if (!killTimer) { killTimer = setTimeout(() => killGroup(child, "SIGKILL"), killGraceMs); killTimer.unref?.(); } }; const onLine = (raw: string) => { const line = raw.trimEnd(); if (!line.trim()) return; lastLines = [...lastLines.slice(-11), line]; const st = parseAria2cStatusLine(line); if (st) { if (st.seeding) { enterSeeding(); return; } if (st.completedBytes > completedBytes) { completedBytes = st.completedBytes; lastProgressAt = now(); } if (phase === "downloading" && now() - lastLogAt >= progressLogMs) { lastLogAt = now(); log(torrentProgressLine({ file: label, status: st, pieces, webSeed })); } return; } // The summary's frame (rules, FILE:, the heading) is noise; aria2c's own // notices and errors are not. if (/\[(NOTICE|WARN|ERROR)\]/.test(line) || /^Exception|errorCode=/.test(line)) { log(`aria2c: ${line.replace(/^\d\d\/\d\d \d\d:\d\d:\d\d /, "")}`); } }; const reader = (stream: NodeJS.ReadableStream | null) => { let buf = ""; stream?.setEncoding?.("utf8"); stream?.on("data", (chunk: string) => { buf += chunk; let i: number; while ((i = buf.search(/\r?\n|\r/)) >= 0) { const m = /\r?\n|\r/.exec(buf.slice(i))!; onLine(buf.slice(0, i)); buf = buf.slice(i + m[0].length); } }); stream?.on("end", () => { if (buf) onLine(buf); }); }; reader(child.stdout); reader(child.stderr); const onAbort = () => stop({ reason: "cancelled", cancelled: true }); opts.signal.addEventListener("abort", onAbort, { once: true }); const ticker = setInterval(() => { void (async () => { if (phase === "downloading" && (await exists(marker))) enterSeeding(); const t = now(); if (phase === "downloading" && t - lastProgressAt >= stallMs) { stop({ reason: `stalled: no progress in ${opts.settings.stallMinutes} min` + (completedBytes >= 0 ? ` (at ${completedBytes} of ${opts.entry.length} bytes)` : ""), stalled: true, }); } if (phase === "seeding" && t - seedingSince >= seedMs + seedGraceMs) { log(`torrent: seeding outlasted ${opts.settings.seedMinutes} min; stopping aria2c.`); stop(null); } })(); }, tickMs); const exit = await new Promise<{ code: number | null; signal: NodeJS.Signals | null; error?: Error }>( (resolve) => { child.once("error", (error) => resolve({ code: null, signal: null, error })); child.once("close", (code, signal) => resolve({ code, signal })); }, ); clearInterval(ticker); if (killTimer) clearTimeout(killTimer); opts.signal.removeEventListener("abort", onAbort); if (exit.error) return { ok: false, reason: `aria2c did not start: ${exit.error.message}` }; if (phase === "downloading" && (await exists(marker))) phase = "seeding"; const why = stopped.why; const complete = phase === "seeding" || (exit.code === 0 && !why && (await exists(target))); if (why?.cancelled) { // A cancel is a cancel even mid-seed: the caller stops the whole download. return { ok: false, reason: why.reason, cancelled: true }; } if (complete && (await exists(target))) { const seeded = seedMs > 0 && seedingSince > 0; if (seeded) log(`torrent: seeding over after ${Math.round((now() - seedingSince) / 60_000)} min.`); return { ok: true, file: target, seeded }; } if (why) return { ok: false, reason: why.reason, ...(why.stalled ? { stalled: true } : {}) }; // aria2c's DHT routing-table complaint (an empty or first-run dht.dat) is // noise, never the reason a run failed. const tail = lastLines .filter((l) => /ERROR|errorCode|Exception/.test(l) && !/DHT/.test(l)) .slice(-2) .join(" | "); return { ok: false, reason: `aria2c exited ${exit.code ?? exit.signal}${tail ? `: ${tail}` : ""}`, }; }