Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit a1a5dea0f932cc5b0fa49215979bf5567f5bc409
parent d75d4da88abf351e00396bc3b8a0fd5034856715
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon,  5 Oct 2026 15:55:05 -0400

sources: archive.org files over BitTorrent when possible, else a direct download — no yt-dlp

An archive.org record is no longer a yt-dlp download (its page scraper fails
on some items: "opening play-av tag not found"). downloadOneManaged routes an
archive.org URL to controller/archiveOrgDownload.ts, which works from the
item's metadata API:

- the one media file the record is (an mp4/mp3 transcode of an .avi/.flac
  original, as the old format selector chose);
- BitTorrent first: aria2c fetches only that file from the item's
  <identifier>_archive.torrent (archive.org is its web seed), then seeds
  (settings.archiveOrg: seedMinutes 10 / seedRatio 1, whichever first; stall
  timeout 5 min; peers and rate caps). aria2c runs in its own process group,
  killed by pgid on a cancel, and stops itself if the editor dies;
- a direct download from archive.org/download/ otherwise (no aria2c, no
  torrent, the file not in it, a stall), resumed with Range, backing off on
  429/503;
- verified against the item's sha1/md5 (size at least); a mismatch falls back
  once to a direct download, a second fails the record;
- the record: archiveorg.json provenance (mirror original kept),
  metadata.info.json synthesised from the item (extractor_key ArchiveOrg,
  canonical id, the file's page, ffprobe duration, playable formats),
  audio.<fmt> (an audio file already in the format is the audio; anything
  else through the app's extraction, a video persisted per the plan), the
  media tier hook, the archive line, download-outcome.json.

Pure bencode decoder and torrent reader; finalizeAppExtraction split out of
downloadOneManaged for reuse; ARIA2C_BIN, the doctor row and aria2 in the
runtime images. SETTINGS.md, ENVIRONMENT.md and settings.json.example
regenerated.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

Diffstat:
MDockerfile | 6++++--
MENVIRONMENT.md | 1+
MSETTINGS.md | 31+++++++++++++++++++++++++++++++
Mcommon/bin/doctor.test.ts | 1+
Mcommon/bin/doctor.ts | 3++-
Acommon/controller/archiveOrgDownload.ts | 545+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/archiveOrg-server.ts | 61++++++++++++++++++++++++++++++++++++++++---------------------
Mcommon/lib/archiveOrg.ts | 159+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/archiveOrgClient.ts | 180+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Acommon/lib/archiveOrgTorrent-server.ts | 293+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/archiveOrgTorrent.ts | 300+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/bencode.ts | 133+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/downloadOutcome.ts | 8+++++++-
Mcommon/lib/envVars.ts | 1+
Mcommon/lib/metadataHistory.ts | 4++++
Mcommon/lib/paths.ts | 5+++++
Mcommon/lib/settingsDocs.ts | 8++++++++
Mcommon/lib/settingsSchema.test.ts | 6+++++-
Mcommon/lib/settingsSchema.ts | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Dcommon/ytdlp/archiveOrgProvenanceHook.test.ts | 80-------------------------------------------------------------------------------
Mcommon/ytdlp/downloadOneManaged.ts | 147+++++--------------------------------------------------------------------------
Acommon/ytdlp/finalizeAppExtraction.ts | 129+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msettings.json.example | 9+++++++++
23 files changed, 1916 insertions(+), 252 deletions(-)

diff --git a/Dockerfile b/Dockerfile @@ -301,9 +301,11 @@ FROM ${RUNTIME_IMAGE} AS runtime-base # with it present. # procps: `ps`, so "is the transcriber actually running in there?" is answerable # from a `docker compose exec` shell. debian:slim ships without it. +# aria2: archive.org files over BitTorrent (archive.org as the web seed), then +# seeded for a while; without it they are downloaded directly. RUN apt-get update \ && apt-get install -y --no-install-recommends \ - ffmpeg zip unzip tar xz-utils gzip rsync curl ca-certificates git procps \ + ffmpeg aria2 zip unzip tar xz-utils gzip rsync curl ca-certificates git procps \ && rm -rf /var/lib/apt/lists/* # yt-dlp as the standalone release binary (it bundles its own python), installed @@ -411,7 +413,7 @@ FROM ${CUDA_RUNTIME_IMAGE} AS runtime-cuda ARG NODE_MAJOR=20 RUN apt-get update \ && apt-get install -y --no-install-recommends \ - ffmpeg zip unzip tar xz-utils gzip rsync curl ca-certificates git gnupg procps \ + ffmpeg aria2 zip unzip tar xz-utils gzip rsync curl ca-certificates git gnupg procps \ && curl -fsSL "https://deb.nodesource.com/setup_${NODE_MAJOR}.x" | bash - \ && apt-get install -y --no-install-recommends nodejs \ && rm -rf /var/lib/apt/lists/* diff --git a/ENVIRONMENT.md b/ENVIRONMENT.md @@ -37,6 +37,7 @@ The one override surface for where things live and which binary runs. Every one | `FINDMNT_BIN` | `findmnt` on PATH | The read-only volume-identity probe behind storage locations. Optional. | common/lib/paths.ts (getPaths) | | `UDISKSCTL_BIN` | `udisksctl` on PATH | Mounts an attached volume from `/storage`. Optional. | common/lib/paths.ts (getPaths) | | `GALLERY_DL_BIN` | `gallery-dl` on PATH | The X/Twitter post fetcher, for social channels. | common/lib/paths.ts (getPaths) | +| `ARIA2C_BIN` | `aria2c` on PATH | Fetches an archive.org file over BitTorrent (the item's torrent, archive.org as web seed), then seeds it for a while. Optional: without it archive.org files are downloaded directly. See `archiveOrg` in [SETTINGS.md](SETTINGS.md). | common/lib/paths.ts (getPaths) | | `OLLAMA_URL` | `http://127.0.0.1:11434` | The local ollama server, the local digest and attribution engine. | common/lib/paths.ts (getPaths) | | `CLAUDE_BIN` | `claude` on PATH | The `claude` CLI, driving the opt-in metered digest lane. | common/lib/paths.ts (getPaths) | | `ARCHILYZER_CONFIG_DIR` | `~/.config/archilyzer` | The operator's private config dir, outside the repo: the two inputs of `archilyzer source publish` below. Never committed. | common/lib/paths.ts (getPaths) | diff --git a/SETTINGS.md b/SETTINGS.md @@ -46,6 +46,7 @@ A copied example PINS every default it spells — including each lane's `autoQue | [`diarization`](#diarization) | object — see below | | [`backfill`](#backfill) | object — see below | | [`attribution`](#attribution) | object — see below | +| [`archiveOrg`](#archiveorg) | object — see below | ## `adminTitle` @@ -790,3 +791,33 @@ Default: "promptVersion": 1 } ``` + +## `archiveOrg` + +How archive.org files are fetched (controller/archiveOrgDownload.ts). Over BitTorrent with aria2c when the item's torrent carries the file — archive.org is the torrent's web seed, so the swarm takes load off archive.org — then seeded for a while; otherwise, or when the torrent stalls, a direct download from archive.org. Either way the file is verified against archive.org's sha1/md5. No yt-dlp. + +#### `archiveOrg` + +| Key | Default | Description | +|---|---|---| +| `torrent` | `true` | Fetch an archive.org file over BitTorrent when the item's `<identifier>_archive.torrent` carries it and aria2c is installed (ARIA2C_BIN). archive.org is the torrent's web seed, so what other peers give never touches archive.org. False = always the direct download. | +| `seedMinutes` | `10` | Minutes to seed the file after it is complete, as a good swarm citizen. The import holds archive.org's queue while it seeds. 0 = do not seed. Clamped to [0, 1440]; default 10. | +| `seedRatio` | `1` | Stop seeding sooner once this much has been uploaded relative to the file's size (aria2c --seed-ratio), whichever of the two comes first. 0 = no ratio limit. Clamped to [0, 100]; default 1. | +| `stallMinutes` | `5` | No download progress for this long stops aria2c and the file is downloaded directly instead. Clamped to [1, 120]; default 5. | +| `maxPeers` | `30` | Most peers per torrent (aria2c --bt-max-peers). Clamped to [1, 500]; default 30. | +| `maxDownloadKiBps` | `0` | Download rate cap for a torrent fetch, KiB/s (aria2c --max-overall-download-limit). 0 = unlimited. | +| `maxUploadKiBps` | `0` | Upload rate cap while downloading and seeding, KiB/s (aria2c --max-overall-upload-limit). 0 = unlimited. | + +Default: + +```json +{ + "torrent": true, + "seedMinutes": 10, + "seedRatio": 1, + "stallMinutes": 5, + "maxPeers": 30, + "maxDownloadKiBps": 0, + "maxUploadKiBps": 0 +} +``` diff --git a/common/bin/doctor.test.ts b/common/bin/doctor.test.ts @@ -56,6 +56,7 @@ function checkout(): { root: string; bin: string; paths: Paths } { ffmpegBin: "ffmpeg", ffprobeBin: "ffprobe", galleryDlBin: "gallery-dl", + aria2cBin: "aria2c", rsyncBin: "rsync", findmntBin: "findmnt", udisksctlBin: "udisksctl", diff --git a/common/bin/doctor.ts b/common/bin/doctor.ts @@ -289,6 +289,7 @@ export async function collectDoctorReport(deps: DoctorDeps): Promise<DoctorRepor { id: "ffmpeg", bin: paths.ffmpegBin, args: ["-version"], neededBy: ["audio extraction", "diarization", "parakeet"] }, { id: "ffprobe", bin: paths.ffprobeBin, args: ["-version"], neededBy: ["duration checks"] }, { id: "gallery-dl", bin: paths.galleryDlBin, neededBy: ["X/Twitter post fetches"] }, + { id: "aria2c", bin: paths.aria2cBin, neededBy: ["archive.org files over BitTorrent (optional; else a direct download)"] }, { id: "rsync", bin: paths.rsyncBin, neededBy: ["saved-video backup"] }, { id: "findmnt", bin: paths.findmntBin, neededBy: ["storage-location identity (optional)"] }, // Looked up, not run: udisksctl has no version flag. @@ -339,7 +340,7 @@ export async function collectDoctorReport(deps: DoctorDeps): Promise<DoctorRepor ); const envNameFor: Record<string, string> = { "yt-dlp": "YTDLP_BIN", ffmpeg: "FFMPEG_BIN", ffprobe: "FFPROBE_BIN", "gallery-dl": "GALLERY_DL_BIN", - rsync: "RSYNC_BIN", findmnt: "FINDMNT_BIN", udisksctl: "UDISKSCTL_BIN", claude: "CLAUDE_BIN", diarize: "DIARIZE_BIN", + aria2c: "ARIA2C_BIN", rsync: "RSYNC_BIN", findmnt: "FINDMNT_BIN", udisksctl: "UDISKSCTL_BIN", claude: "CLAUDE_BIN", diarize: "DIARIZE_BIN", }; const reports = await Promise.all(specs.map((s) => probe(s))); for (const r of reports) { diff --git a/common/controller/archiveOrgDownload.ts b/common/controller/archiveOrgDownload.ts @@ -0,0 +1,545 @@ +// ONE archive.org RECORD DOWNLOADED — without yt-dlp. +// +// Every managed download of an archive.org URL lands here +// (ytdlp/downloadOneManaged.ts routes it before any yt-dlp spawn): the +// single-URL import, each file of a bulk import (controller/archiveOrgImport.ts), +// a re-download, "Persist source video". yt-dlp's ArchiveOrg extractor scrapes +// the item's page, and fails on items whose page it does not understand ("opening +// play-av tag not found" on an audio item); the item's metadata API, which the +// import has already asked for (lib/archiveOrgClient.ts caches it), names every +// file, its size and its checksums. So the record is built from that. +// +// THE LADDER, for the one media file the record is (lib/archiveOrg.ts +// archiveOrgRecordFile → archiveOrgFetchFile): +// +// 1. BITTORRENT (the operator: "use torrents when possible to be extra polite +// to archive.org") — when settings.archiveOrg.torrent is on, aria2c is +// installed, the item lists `<identifier>_archive.torrent`, and the +// torrent carries the file. aria2c fetches only that file from the swarm +// and archive.org's web seed, then seeds it (lib/archiveOrgTorrent-server.ts). +// A stall, an aria2c failure, or a torrent that cannot be used falls back: +// "fell back to direct download: <reason>". +// 2. DIRECT — https://archive.org/download/<identifier>/<file>, one stream, +// resumed with a Range request, backing off on 429/503 +// (ArchiveOrgClient.downloadFile). +// +// Whichever fetched it, the file is VERIFIED against the item's sha1 (else +// md5, else size). A mismatch deletes it and falls back once to a direct +// download; a second mismatch fails the record. +// +// THEN THE RECORD, in data/<id>/ as any managed download leaves it: +// - archiveorg.json, the provenance (kept when already there; the mirror's +// uploaded info.json read for the original's title and date); +// - metadata.info.json, synthesised (lib/archiveOrg.ts archiveOrgInfoJson) +// with the duration ffprobe measures, inside the metadata history; +// - audio.<fmt>: an audio file already in the channel's format IS the audio; +// anything else is source-media.<ext>, extracted by the app's own +// extraction (ytdlp/finalizeAppExtraction.ts) — a video kept in the +// saved-video store when the persistence plan says so, else removed; an +// audio original is never kept as a "source video"; +// - the media tier's hook, the archive line `archiveorg <id>`, and +// download-outcome.json with the attempts ("archiveorg-torrent", +// "archiveorg-direct"). +// +// A CANCEL kills aria2c by its process group and leaves the partial in the +// staging dir (`.archiveorg-fetch/` under data/<id>/), so the next run resumes. + +import path from "node:path"; +import { createHash } from "node:crypto"; +import { createReadStream } from "node:fs"; +import { appendFile, mkdir, readdir, rename, rm, stat } from "node:fs/promises"; +import type { ManagedDownloadOpts } from "../ytdlp/downloadOneManaged"; +import type { AudioFormat } from "../lib/channelConfig"; +import type { + DownloadAttempt, + DownloadAttemptKind, + DownloadOutcomeRecord, + DownloadOutcomeStatus, +} from "../lib/downloadOutcome"; +import type { DownloadFailureClass } from "../lib/availability"; +import { writeDownloadOutcome } from "../lib/downloadOutcome-server"; +import { + archiveOrgFetchFile, + archiveOrgInfoJson, + archiveOrgRecordFile, + buildArchiveOrgProvenance, + isArchiveOrgVideoFileName, + listArchiveOrgMediaFiles, + type ArchiveOrgItemFile, + type ArchiveOrgItemMetadata, + type ArchiveOrgProvenance, +} from "../lib/archiveOrg"; +import { archiveOrgVideoId, parseArchiveOrgUrl, type ArchiveOrgRef } from "../lib/archiveOrgId"; +import { + ARCHIVE_ORG_USER_AGENT, + ArchiveOrgRequestError, + archiveOrgClient, + type ArchiveOrgClient, +} from "../lib/archiveOrgClient"; +import { ensureArchiveOrgProvenanceSidecar } from "../lib/archiveOrg-server"; +import { + DEFAULT_ARCHIVE_ORG_FETCH_SETTINGS, + findTorrentFile, + parseTorrent, + type ArchiveOrgFetchSettings, +} from "../lib/archiveOrgTorrent"; +import { + aria2cAvailable as defaultAria2cAvailable, + fetchFileByTorrent, + type TorrentFetchOpts, + type TorrentFetchResult, +} from "../lib/archiveOrgTorrent-server"; +import { AUDIO_EXTS, MEDIA_EXTS } from "../lib/mediaFiles"; +import { tierMediaFile, tierVideoDir } from "../lib/mediaTier-server"; +import { withMetadataHistory } from "../lib/metadataHistory-server"; +import { writeJsonAtomic } from "../lib/jsonFile-server"; +import { getSettings } from "../lib/settings"; +import { formatBytes } from "../lib/format"; +import { isDoNotClean } from "../lib/doNotClean-server"; +import { isInKeepWindow, uploadKeyFor } from "./keptVideos"; +import { resolvePersistenceDecision } from "../ytdlp/persistencePlan"; +import { finalizeAppExtraction } from "../ytdlp/finalizeAppExtraction"; +import { probeMediaDurationSec } from "../ytdlp/ffprobeDuration"; +import { transcribeWithWorker } from "./transcribeOne"; + +export const ARCHIVE_ORG_STAGING_DIR = ".archiveorg-fetch"; + +export type ArchiveOrgVerifyResult = + | { ok: true; by: "sha1" | "md5" | "size" | "nothing" } + | { ok: false; reason: string }; + +export type ArchiveOrgDownloadDeps = { + client?: ArchiveOrgClient; + settings?: ArchiveOrgFetchSettings; + aria2cAvailable?: (bin: string) => Promise<boolean>; + fetchByTorrent?: (o: TorrentFetchOpts) => Promise<TorrentFetchResult>; + // The plain download: the client's, by default. + fetchDirect?: (o: { + identifier: string; + file: string; + dest: string; + expectedSize?: number; + signal: AbortSignal; + onLog: (s: string) => void; + }) => Promise<void>; + verify?: (file: string, entry: ArchiveOrgItemFile) => Promise<ArchiveOrgVerifyResult>; + probeDuration?: (file: string) => Promise<number | null>; + finalize?: typeof finalizeAppExtraction; + transcribe?: typeof transcribeWithWorker; +}; + +// ─── Verification ─── + +// The file against archive.org's checksum: sha1 when the item lists one, else +// md5, else the size alone, else nothing to check against. One read. +export async function verifyArchiveOrgFile( + file: string, + entry: ArchiveOrgItemFile, +): Promise<ArchiveOrgVerifyResult> { + const size = Number(entry.size); + const st = await stat(file).catch(() => null); + if (!st) return { ok: false, reason: "the file is missing" }; + if (Number.isFinite(size) && size > 0 && st.size !== size) { + return { ok: false, reason: `size ${st.size} bytes, archive.org lists ${size}` }; + } + const want = entry.sha1 ? { algo: "sha1" as const, hex: entry.sha1 } : entry.md5 ? { algo: "md5" as const, hex: entry.md5 } : null; + if (!want) return { ok: true, by: Number.isFinite(size) && size > 0 ? "size" : "nothing" }; + const hash = createHash(want.algo); + await new Promise<void>((resolve, reject) => { + createReadStream(file) + .on("data", (c) => hash.update(c)) + .on("end", () => resolve()) + .on("error", reject); + }); + const got = hash.digest("hex"); + if (got.toLowerCase() !== want.hex.trim().toLowerCase()) { + return { ok: false, reason: `${want.algo} ${got}, archive.org lists ${want.hex}` }; + } + return { ok: true, by: want.algo }; +} + +// ─── Small helpers ─── + +function extOf(name: string): string { + const m = /\.([A-Za-z0-9]{1,8})$/.exec(name); + return m ? m[1].toLowerCase() : ""; +} + +async function hasTranscriptOnDisk(videoDir: string): Promise<boolean> { + const entries = await readdir(videoDir).catch(() => [] as string[]); + return entries.some((e) => { + if (e === "transcript.json") return true; + const m = /^transcript\.([^.]+)\.(?:vtt|json|json3|srv1|srv2|srv3)$/.exec(e); + return !!m && m[1] !== "live_chat"; + }); +} + +function failureClassOf(err: unknown): DownloadFailureClass { + if (err instanceof ArchiveOrgRequestError) { + if (err.rateLimited) return "rate_limit"; + if (err.status === 403 || err.status === 404 || err.status === 410) return "per_video"; + return err.status === null ? "network" : "unknown"; + } + return "network"; +} + +function attempt(n: number, kind: DownloadAttemptKind, handling: DownloadAttempt["handling"], error?: string): DownloadAttempt { + return { + n, + kind, + handling, + usedCookies: false, + // No yt-dlp ran: 0 for a fetch that delivered a verified file, 1 for one + // that did not, so every reader of the field keeps meaning "succeeded". + ytdlpExitCode: error ? 1 : 0, + ...(error ? { error } : {}), + }; +} + +// ─── The download ─── + +export async function downloadArchiveOrgManaged( + opts: ManagedDownloadOpts, + deps: ArchiveOrgDownloadDeps = {}, +): Promise<DownloadOutcomeRecord> { + const startedAt = new Date().toISOString(); + const log = (s: string) => opts.onLog(s.endsWith("\n") ? s : `${s}\n`); + const client = deps.client ?? archiveOrgClient; + const handling = opts.channelConfig.handling; + const channelDir = path.join(opts.paths.channelsDir, opts.channelSlug); + const attempts: DownloadAttempt[] = []; + + const ref: ArchiveOrgRef | null = parseArchiveOrgUrl(opts.videoUrl); + if (!ref) { + log(`Not an archive.org item URL: ${opts.videoUrl}`); + return { + videoId: "unknown", + webpageUrl: opts.videoUrl, + status: "failed", + startedAt, + finishedAt: new Date().toISOString(), + attempts: [attempt(1, "archiveorg-direct", handling, "not an archive.org item URL")], + failureClass: "per_video", + }; + } + const videoId = archiveOrgVideoId(ref); + const videoDir = path.join(channelDir, "data", videoId); + const staging = path.join(videoDir, ARCHIVE_ORG_STAGING_DIR); + await mkdir(videoDir, { recursive: true }); + + const finish = async (o: { + status: DownloadOutcomeStatus; + failureClass?: DownloadFailureClass; + fellBackToTranscribe?: boolean; + }): Promise<DownloadOutcomeRecord> => { + const record: DownloadOutcomeRecord = { + videoId, + webpageUrl: opts.videoUrl, + status: o.status, + startedAt, + finishedAt: new Date().toISOString(), + attempts, + ...(o.failureClass ? { failureClass: o.failureClass } : {}), + ...(o.fellBackToTranscribe ? { fellBackToTranscribe: true } : {}), + }; + try { + // THE MEDIA TIER'S HOOK (release 17): what this run finalised moves into + // channels/<slug>/media when the channel has one. Never throws. + await tierVideoDir(videoDir, { onLog: opts.onLog }); + await writeDownloadOutcome(videoDir, record); + } catch (err) { + log(`Failed to write download-outcome.json: ${(err as Error).message}`); + } + return record; + }; + const fail = (kind: DownloadAttemptKind, error: string, failureClass: DownloadFailureClass) => { + attempts.push(attempt(attempts.length + 1, kind, handling, error)); + return finish({ status: "failed", failureClass }); + }; + + // ---------- The item ---------- + let item: ArchiveOrgItemMetadata; + try { + item = await client.itemMetadata(ref.identifier, opts.signal); + } catch (err) { + log(`archive.org metadata for ${ref.identifier} not fetched: ${(err as Error).message}`); + return fail("archiveorg-direct", (err as Error).message, failureClassOf(err)); + } + const identifier = item.metadata.identifier || ref.identifier; + const file = archiveOrgRecordFile(item, ref.file); + if (!file) { + const media = listArchiveOrgMediaFiles(item).length; + const why = ref.file + ? `archive.org item "${identifier}" has no file "${ref.file}"` + : media === 0 + ? `archive.org item "${identifier}" has no media files` + : `archive.org item "${identifier}" holds ${media} media files; import one by its file URL`; + log(why); + return fail("archiveorg-direct", why, "per_video"); + } + const fetchEntry = archiveOrgFetchFile(item, file)!; + const fetchSize = Number(fetchEntry.size); + log( + `archive.org: ${identifier} / ${file}` + + (fetchEntry.name !== file ? ` — fetching archive.org's ${fetchEntry.format ?? extOf(fetchEntry.name)} of it, ${fetchEntry.name}` : "") + + (Number.isFinite(fetchSize) && fetchSize > 0 ? ` (${formatBytes(fetchSize)})` : "") + + ".", + ); + + // ---------- Provenance (before the bytes: it is what the record says) ---------- + let prov: ArchiveOrgProvenance | null; + try { + prov = await ensureArchiveOrgProvenanceSidecar( + videoDir, + { identifier, ...(ref.file ? { file: ref.file } : {}) }, + { client, signal: opts.signal, onLog: opts.onLog }, + ); + } catch { + return fail("archiveorg-direct", "cancelled", "unknown"); + } + if (!prov) { + // archive.org answered the metadata a moment ago but not this: the record + // is built from the item alone (the next download fills the sidecar in). + prov = buildArchiveOrgProvenance({ + ref: { identifier, ...(ref.file ? { file: ref.file } : {}) }, + item, + fetchedAt: new Date().toISOString(), + }); + } + + // ---------- The bytes ---------- + const settings = deps.settings ?? (getSettings().archiveOrg ?? DEFAULT_ARCHIVE_ORG_FETCH_SETTINGS); + const verify = deps.verify ?? verifyArchiveOrgFile; + const aria2cBin = opts.paths.aria2cBin ?? "aria2c"; + let got: string | null = null; + let mismatches = 0; + + if (settings.torrent) { + const why = await (async (): Promise<string | null> => { + if (!(await (deps.aria2cAvailable ?? defaultAria2cAvailable)(aria2cBin))) { + return "torrents need aria2c (not found on PATH; ARIA2C_BIN names another binary)"; + } + const torrentName = `${identifier}_archive.torrent`; + if (!item.files.some((f) => f.name === torrentName)) return "the item lists no torrent"; + let buf: Buffer; + try { + buf = await client.itemTorrent(identifier, opts.signal); + } catch (err) { + if (opts.signal.aborted) return "cancelled"; + return `the torrent was not fetched (${(err as Error).message})`; + } + const parsed = parseTorrent(buf); + if (!parsed) return "the torrent could not be read"; + const entry = findTorrentFile(parsed, fetchEntry.name); + if (!entry) return `the item's torrent does not carry ${fetchEntry.name} (older than the file?)`; + const res = await (deps.fetchByTorrent ?? fetchFileByTorrent)({ + aria2cBin, + torrent: buf, + parsed, + entry, + stagingDir: path.join(staging, "torrent"), + settings, + userAgent: ARCHIVE_ORG_USER_AGENT, + onLog: opts.onLog, + signal: opts.signal, + }); + if (!res.ok) { + attempts.push(attempt(attempts.length + 1, "archiveorg-torrent", handling, res.reason)); + return res.cancelled ? "cancelled" : res.reason; + } + const v = await verify(res.file, fetchEntry); + if (!v.ok) { + mismatches++; + attempts.push(attempt(attempts.length + 1, "archiveorg-torrent", handling, `checksum mismatch: ${v.reason}`)); + return `checksum mismatch (${v.reason})`; + } + attempts.push(attempt(attempts.length + 1, "archiveorg-torrent", handling)); + log(`Verified ${fetchEntry.name} (${v.by}).`); + got = res.file; + return null; + })(); + if (why === "cancelled" || opts.signal.aborted) { + log("Cancelled; the partial download is kept for the next run."); + if (attempts.at(-1)?.error !== "cancelled") attempts.push(attempt(attempts.length + 1, "archiveorg-torrent", handling, "cancelled")); + return finish({ status: "failed" }); + } + if (why) { + log(`fell back to direct download: ${why}`); + await rm(path.join(staging, "torrent"), { recursive: true, force: true }); + } + } else { + log("archiveOrg.torrent is off; downloading directly."); + } + + while (!got) { + const dest = path.join(staging, "direct", fetchEntry.name.split("/").pop() || "file"); + await mkdir(path.dirname(dest), { recursive: true }); + try { + await (deps.fetchDirect ?? defaultFetchDirect(client))({ + identifier, + file: fetchEntry.name, + dest, + ...(Number.isFinite(fetchSize) && fetchSize > 0 ? { expectedSize: fetchSize } : {}), + signal: opts.signal, + onLog: opts.onLog, + }); + } catch (err) { + if (opts.signal.aborted) { + attempts.push(attempt(attempts.length + 1, "archiveorg-direct", handling, "cancelled")); + log("Cancelled; the partial download is kept for the next run."); + return finish({ status: "failed" }); + } + log(`Direct download failed: ${(err as Error).message}`); + return fail("archiveorg-direct", (err as Error).message, failureClassOf(err)); + } + const v = await verify(dest, fetchEntry); + if (v.ok) { + attempts.push(attempt(attempts.length + 1, "archiveorg-direct", handling)); + log(`Verified ${fetchEntry.name} (${v.by}).`); + got = dest; + break; + } + mismatches++; + await rm(dest, { force: true }); + if (mismatches >= 2) { + log(`Checksum mismatch again (${v.reason}); the record fails.`); + return fail("archiveorg-direct", `checksum mismatch: ${v.reason}`, "per_video"); + } + attempts.push(attempt(attempts.length + 1, "archiveorg-direct", handling, `checksum mismatch: ${v.reason}`)); + log(`fell back to direct download: checksum mismatch (${v.reason}); downloading once more.`); + } + + // ---------- The record ---------- + const fetchedPath: string = got; + const durationSec = + (await (deps.probeDuration ?? + ((f: string) => + probeMediaDurationSec({ ffprobeBin: opts.paths.ffprobeBin, file: f, signal: opts.signal, onLog: () => {} })))( + fetchedPath, + )) ?? undefined; + const info = archiveOrgInfoJson({ prov, item, file, fetched: fetchEntry, durationSec }); + try { + await withMetadataHistory(videoDir, { by: "archiveorg-import", onLog: opts.onLog }, () => + writeJsonAtomic(path.join(videoDir, "metadata.info.json"), info, { indent: 0, newline: false }), + ); + } catch (err) { + return fail("archiveorg-direct", `metadata.info.json not written: ${(err as Error).message}`, "unknown"); + } + + const fmt: AudioFormat = opts.audioFormatOverride ?? opts.channelConfig.audioFormat ?? "mp3"; + const ext = extOf(fetchEntry.name); + const isVideo = isArchiveOrgVideoFileName(fetchEntry.name); + const isAudio = !isVideo && AUDIO_EXTS.includes(ext); + const hasTranscript = await hasTranscriptOnDisk(videoDir); + // "Persist source video" over a record that already has a transcript keeps + // the container and extracts nothing (the managed download's keepTranscript). + const keepTranscript = opts.forceMedia === true && hasTranscript; + const plan = resolvePersistenceDecision({ + inWindow: isInKeepWindow(uploadKeyFor(info.upload_date as string | undefined, videoId), opts.keepWindow), + pinned: await isDoNotClean(videoDir), + channelExtractionMode: opts.channelConfig.extractionMode, + overrides: { + keepSourceVideoOverride: opts.keepSourceVideoOverride, + extractImmediately: opts.extractImmediately, + }, + }); + const persist = isVideo && plan.persist; + if (isVideo) { + log(`Persistence: ${persist ? "keep source video" : "audio-only"} (${plan.reason}).`); + } + + if (isAudio && ext === fmt && !keepTranscript) { + // Already the channel's audio format: it IS the audio. + const name = `audio.${fmt}`; + await rename(fetchedPath, path.join(videoDir, name)); + await tierMediaFile(videoDir, name, { onLog: opts.onLog }); + log(`Wrote ${name} (${fetchEntry.name}, no transcode).`); + } else if (keepTranscript && !persist) { + log( + `forceMedia: a transcript is on disk and ${isVideo ? `this download keeps no source video (${plan.reason})` : "an audio original is not kept as a source video"}; nothing kept.`, + ); + await rm(fetchedPath, { force: true }); + } else { + const sourceName = `source-media.${MEDIA_EXTS.includes(ext) ? ext : ext || "bin"}`; + await rename(fetchedPath, path.join(videoDir, sourceName)); + await (deps.finalize ?? finalizeAppExtraction)({ + paths: opts.paths, + channelSlug: opts.channelSlug, + channelConfig: opts.channelConfig, + videoDir, + videoId, + fmt, + persist, + category: plan.category, + origin: opts.persistOrigin, + extractAudio: !keepTranscript, + sourceFilename: sourceName, + onLog: opts.onLog, + signal: opts.signal, + }); + // An audio original whose extraction failed is still the only copy of the + // recording: finalize keeps it, and the record fails below for want of + // audio. A video that was not persisted has been removed by finalize. + } + await rm(staging, { recursive: true, force: true }); + + const entries = await readdir(videoDir).catch(() => [] as string[]); + const haveAudio = entries.includes(`audio.${fmt}`); + if (!keepTranscript && !haveAudio) { + // Whatever finalize kept (the container) stays for a retry to use. + return fail("archiveorg-direct", `no audio.${fmt} was produced from ${fetchEntry.name}`, "unknown"); + } + + // ---------- Transcription hand-off ---------- + // archive.org has no captions: a youtube-handling channel's record is the + // no-subs fallback's, a transcribe-handling channel's is a plain download. + let status: DownloadOutcomeStatus = "ok"; + const fellBackToTranscribe = handling === "youtube" && !keepTranscript; + if (fellBackToTranscribe && opts.inlineTranscribeOnFallback && !opts.signal.aborted) { + try { + await (deps.transcribe ?? transcribeWithWorker)({ + paths: opts.paths, + videoDir, + videoId, + audioFilename: `audio.${fmt}`, + onLog: opts.onLog, + signal: opts.signal, + }); + status = "ok-auto-transcribed"; + } catch (err) { + log(`Whisper failed after the archive.org download: ${(err as Error).message}`); + status = "failed"; + } + } else if (!keepTranscript) { + log(`Audio downloaded for ${videoId}; it is transcribed by the next "Transcribe missing" pass.`); + } + + if (status !== "failed" && opts.appendArchive !== false) { + try { + await mkdir(channelDir, { recursive: true }); + await appendFile(path.join(channelDir, "archive"), `archiveorg ${videoId}\n`); + } catch (err) { + log(`Archive append failed (archiveorg ${videoId}): ${(err as Error).message}`); + } + } + return finish({ status, fellBackToTranscribe }); +} + +function defaultFetchDirect(client: ArchiveOrgClient): NonNullable<ArchiveOrgDownloadDeps["fetchDirect"]> { + return async ({ identifier, file, dest, expectedSize, signal, onLog }) => { + let lastLogAt = 0; + onLog(`direct: https://archive.org/download/${identifier}/${file}\n`); + await client.downloadFile(identifier, file, dest, { + signal, + ...(expectedSize ? { expectedSize } : {}), + onProgress: ({ bytes, total }) => { + const t = Date.now(); + if (t - lastLogAt < 15_000) return; + lastLogAt = t; + onLog( + `direct: ${file} ${formatBytes(bytes)}${total ? ` of ${formatBytes(total)} (${Math.floor((bytes / total) * 100)}%)` : ""}\n`, + ); + }, + }); + }; +} diff --git a/common/lib/archiveOrg-server.ts b/common/lib/archiveOrg-server.ts @@ -1,8 +1,12 @@ -// archive.org PROVENANCE ON DISK — the `archiveorg.json` sidecar, and the step -// every managed download of an archive.org record runs after yt-dlp writes its -// metadata (ytdlp/downloadOneManaged.ts, right after the prefetch). +// archive.org PROVENANCE ON DISK — the `archiveorg.json` sidecar. // -// THE STEP, `ensureArchiveOrgProvenance`: +// `ensureArchiveOrgProvenanceSidecar` is what every archive.org download runs +// (controller/archiveOrgDownload.ts), before it writes the record's +// metadata.info.json from the item. `ensureArchiveOrgProvenance` is the older, +// whole step for a record yt-dlp wrote (archive.org records were fetched by +// yt-dlp before the torrent slice): the sidecar, then that record corrected. +// +// THE WHOLE STEP, `ensureArchiveOrgProvenance`: // // 1. the sidecar: kept when it is already there for this item and file; // otherwise built from the item's metadata API (one cached request per @@ -93,29 +97,44 @@ export type EnsureArchiveOrgProvenanceOpts = { now?: () => Date; }; -export async function ensureArchiveOrgProvenance( - opts: EnsureArchiveOrgProvenanceOpts, +// The sidecar for `ref`: the one on disk when it is for this item and file, +// else fetched (the item's metadata is cached; a mirror's info.json is one more +// request) and written. Null when archive.org could not be asked — logged, +// never thrown, unless the job was cancelled. +export async function ensureArchiveOrgProvenanceSidecar( + videoDir: string, + ref: ArchiveOrgRef, + opts: { client?: ArchiveOrgClient; signal?: AbortSignal; now?: () => Date; onLog?: (line: string) => void } = {}, ): Promise<ArchiveOrgProvenance | null> { const log = opts.onLog ?? (() => {}); - const ref = parseArchiveOrgUrl(opts.videoUrl); - if (!ref) return null; - let prov = await loadArchiveOrgProvenance(opts.videoDir); + let prov = await loadArchiveOrgProvenance(videoDir); if (prov && (prov.identifier !== ref.identifier || (prov.file ?? "") !== (ref.file ?? ""))) { prov = null; } - if (!prov) { - try { - prov = await fetchArchiveOrgProvenance(ref, opts); - await writeArchiveOrgProvenance(opts.videoDir, prov); - log( - `archive.org provenance: ${prov.identifier}${prov.file ? ` / ${prov.file}` : ""}` + - (prov.mirror ? ` — mirror of YouTube ${prov.mirror.id} (from ${prov.mirror.from})` : "") + - `.\n`, - ); - } catch (err) { - log(`archive.org provenance not fetched (${(err as Error).message}); the record keeps yt-dlp's fields.\n`); - } + if (prov) return prov; + try { + prov = await fetchArchiveOrgProvenance(ref, opts); + await writeArchiveOrgProvenance(videoDir, prov); + log( + `archive.org provenance: ${prov.identifier}${prov.file ? ` / ${prov.file}` : ""}` + + (prov.mirror ? ` — mirror of YouTube ${prov.mirror.id} (from ${prov.mirror.from})` : "") + + `.\n`, + ); + return prov; + } catch (err) { + if (opts.signal?.aborted) throw err; + log(`archive.org provenance not fetched (${(err as Error).message}).\n`); + return null; } +} + +export async function ensureArchiveOrgProvenance( + opts: EnsureArchiveOrgProvenanceOpts, +): Promise<ArchiveOrgProvenance | null> { + const log = opts.onLog ?? (() => {}); + const ref = parseArchiveOrgUrl(opts.videoUrl); + if (!ref) return null; + const prov = await ensureArchiveOrgProvenanceSidecar(opts.videoDir, ref, opts); const info = await readInfo(opts.videoDir); if (!info) return prov; // Without a sidecar only the page is corrected — it is the one field the diff --git a/common/lib/archiveOrg.ts b/common/lib/archiveOrg.ts @@ -22,6 +22,7 @@ import { archiveOrgDetailsUrl, archiveOrgDownloadUrl, archiveOrgTorrentUrl, + archiveOrgVideoId, parseArchiveOrgUrl, youtubeIdFromFileName, youtubeIdFromIdentifier, @@ -46,6 +47,10 @@ export type ArchiveOrgItemFile = { title?: string; // On a derivative: the original it was made from. original?: string; + // archive.org's checksums of the file, hex. What a fetched file is verified + // against (controller/archiveOrgDownload.ts). + md5?: string; + sha1?: string; }; export type ArchiveOrgItemMetadata = { @@ -59,6 +64,7 @@ export type ArchiveOrgItemMetadata = { collection?: string | string[]; mediatype?: string; description?: string | string[]; + subject?: string | string[]; }; files: ArchiveOrgItemFile[]; }; @@ -444,3 +450,156 @@ export function archiveOrgPlayableUrl(meta: { } return undefined; } + +// ─── What a record is fetched from, and the record it becomes ─── +// +// An archive.org record is not fetched by yt-dlp (controller/ +// archiveOrgDownload.ts): one FILE of the item comes down, over BitTorrent or +// as a plain download, and the record's metadata.info.json is written here +// from the item's metadata. These are the pure rules for both. + +// The media file a record IS: the file the URL named, or an item's one media +// original. Null for an item URL of an item with none, or several. +export function archiveOrgRecordFile( + item: ArchiveOrgItemMetadata, + file: string | undefined, +): string | null { + if (file) return item.files.some((f) => f.name === file) ? file : null; + const media = listArchiveOrgMediaFiles(item); + return media.length === 1 ? media[0].name : null; +} + +// Containers fetched as they are; an original in any other (.avi, .mpeg, +// .flac, .wav, …) is fetched as archive.org's own transcode of it when there +// is one — the h.264 mp4 of a video, the MP3 of an audio file — a fraction of +// the bytes for the same transcript. The rule yt-dlp's selector used to apply +// (ytdlp/downloadFormat.ts ARCHIVE_ORG_AUTO_FORMAT_SELECTOR), kept. +const FETCH_AS_IS = new Set(["mp4", "m4v", "mkv", "webm", "mp3", "m4a", "ogg", "oga", "opus"]); +const VIDEO_EXTS = new Set(["mp4", "m4v", "mkv", "webm", "mov", "avi", "mpeg", "mpg", "ogv", "flv", "wmv", "3gp", "ts"]); + +export function isArchiveOrgVideoFileName(name: string): boolean { + return VIDEO_EXTS.has(extOf(name)); +} + +export function archiveOrgFetchFile( + item: ArchiveOrgItemMetadata, + file: string, +): ArchiveOrgItemFile | null { + const own = item.files.find((f) => f.name === file); + if (!own) return null; + if (FETCH_AS_IS.has(extOf(own.name))) return own; + const derivs = item.files.filter((f) => f.source === "derivative" && f.original === file); + const wanted = isArchiveOrgVideoFileName(own.name) ? ["mp4", "webm"] : ["mp3", "ogg", "m4a"]; + for (const ext of wanted) { + const d = derivs.find((f) => extOf(f.name) === ext); + if (d) return d; + } + return own; +} + +// archive.org's `length`: seconds ("123.45") or a clock ("1:02:03"). +export function parseArchiveOrgLength(v: string | undefined): number | undefined { + if (!v) return undefined; + const t = v.trim(); + if (/^\d+(?:\.\d+)?$/.test(t)) return Number(t); + const m = /^(?:(\d+):)?(\d{1,2}):(\d{1,2}(?:\.\d+)?)$/.exec(t); + if (!m) return undefined; + return Number(m[1] ?? 0) * 3600 + Number(m[2]) * 60 + Number(m[3]); +} + +// "2019-05-04", "2019-05-04 12:00:00", "2019-05" or "2019" → YYYYMMDD, only +// for a whole day. +function ymdOf(v: string | undefined): string | undefined { + const m = v ? /^(\d{4})-(\d{2})-(\d{2})/.exec(v.trim()) : null; + return m ? `${m[1]}${m[2]}${m[3]}` : undefined; +} + +function descriptionOf(v: unknown): string | undefined { + if (typeof v === "string") return v.trim() || undefined; + if (Array.isArray(v)) { + const parts = v.filter((x): x is string => typeof x === "string" && !!x.trim()); + return parts.length > 0 ? parts.join("\n\n") : undefined; + } + return undefined; +} + +// THE RECORD'S metadata.info.json, in the shape every reader expects of +// yt-dlp's (lib/transcripts-server.ts RawMetadata): `extractor_key` +// "ArchiveOrg" (so platformFromMetadata says archiveorg), `id` the canonical +// video id, `webpage_url` the page of what was imported (the file's page for +// one file of many — the snapshot names the directory after it), `duration` +// measured off the fetched file, and `formats` the file and archive.org's +// playable transcodes of it, which is what the file player picks from +// (archiveOrgPlayableUrl). +// +// The provenance decides title, date, uploader and description exactly as it +// corrected yt-dlp's record before (archiveOrgMetadataPatch): a mirror's +// original, else the file's own title, else the item's. The uploading +// account's e-mail address (`uploader` in the item's metadata) is never read. +// +// KEY ORDER IS NOT COSMETIC: two readers scan the file's bytes, the title from +// a 16 KB head (videoTitles.ts) and upload_date from an 8 KB tail +// (recencyIndex.ts), so the title is near the front and the date is last, +// after the description and the formats. +export function archiveOrgInfoJson(opts: { + prov: ArchiveOrgProvenance; + item: ArchiveOrgItemMetadata; + // The record's media file (the original), and the one actually fetched. + file: string; + fetched: ArchiveOrgItemFile; + durationSec?: number; +}): Record<string, unknown> { + const { prov, item, file, fetched } = opts; + const m = item.metadata; + const identifier = prov.identifier; + const ref = prov.file ? { identifier, file: prov.file } : { identifier }; + const page = archiveOrgDetailsUrl(ref); + const original = item.files.find((f) => f.name === file); + const formats = item.files + .filter((f) => f.name === file || (f.source === "derivative" && f.original === file)) + .filter((f) => isMediaFileName(f.name)) + .map((f) => { + const size = Number(f.size); + return { + format_id: f.name, + url: archiveOrgDownloadUrl(identifier, f.name), + ext: extOf(f.name), + format_note: f.source === "original" ? "original" : "derivative", + ...(f.format ? { format: f.format } : {}), + ...(Number.isFinite(size) && size > 0 ? { filesize: size } : {}), + }; + }); + const duration = + opts.durationSec ?? parseArchiveOrgLength(fetched.length) ?? parseArchiveOrgLength(original?.length); + const subjects = stringList(m.subject).flatMap((s) => s.split(/\s*;\s*/)).filter(Boolean); + const base: Record<string, unknown> = { + id: archiveOrgVideoId(ref), + title: firstString(m.title) ?? identifier, + webpage_url: page, + original_url: page, + extractor: "archive.org", + extractor_key: "ArchiveOrg", + ...(firstString(m.creator) ? { uploader: firstString(m.creator) } : {}), + ...(typeof duration === "number" && Number.isFinite(duration) ? { duration: Math.round(duration * 100) / 100 } : {}), + ext: extOf(fetched.name), + ...(subjects.length > 0 ? { tags: subjects } : {}), + ...(descriptionOf(m.description) ? { description: descriptionOf(m.description) } : {}), + formats, + ...(ymdOf(firstString(m.date)) ?? ymdOf(firstString(m.publicdate)) + ? { upload_date: ymdOf(firstString(m.date)) ?? ymdOf(firstString(m.publicdate)) } + : {}), + }; + const patch = archiveOrgMetadataPatch(prov, base); + const out: Record<string, unknown> = { ...base }; + for (const [k, v] of Object.entries(patch)) { + if (v === null) delete out[k]; + else out[k] = v; + } + // A patched date is still the tail's last key. + if ("upload_date" in out) { + const d = out.upload_date; + delete out.upload_date; + out.upload_date = d; + } + return out; +} diff --git a/common/lib/archiveOrgClient.ts b/common/lib/archiveOrgClient.ts @@ -1,5 +1,7 @@ // THE POLITE archive.org CLIENT — every request this app makes to archive.org -// that is not a yt-dlp spawn (the item metadata API, a mirror's info.json). +// 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: @@ -14,20 +16,31 @@ // 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 is cached in memory for `cacheTtlMs` -// (6 h): a bulk import of 161 files of one item asks for it -// once, and so does every per-file provenance step after it. +// 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. // -// The downloads themselves are yt-dlp's (one plain HTTP stream per file, paced -// by PLATFORM_ARGS.archiveorg) and run on the `platform:archiveorg` job queue, -// one at a time. +// 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 } from "./archiveOrgId"; +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})`; @@ -56,6 +69,19 @@ export type ArchiveOrgClientOpts = { 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]); @@ -91,6 +117,7 @@ export class ArchiveOrgClient { private chain: Promise<unknown> = Promise.resolve(); private lastDoneAt = -Infinity; private readonly cache = new Map<string, { at: number; value: ArchiveOrgItemMetadata }>(); + private readonly torrentCache = new Map<string, { at: number; value: Buffer }>(); // Requests actually sent (each attempt), for tests and logs. requests = 0; @@ -108,6 +135,7 @@ export class ArchiveOrgClient { maxBackoffMs: opts.maxBackoffMs ?? 120_000, cacheTtlMs: opts.cacheTtlMs ?? 6 * 60 * 60 * 1000, timeoutMs: opts.timeoutMs ?? 60_000, + idleTimeoutMs: opts.idleTimeoutMs ?? 120_000, }; } @@ -190,6 +218,142 @@ export class ArchiveOrgClient { itemJsonFile(identifier: string, file: string, signal?: AbortSignal): Promise<unknown> { return this.getJson(archiveOrgDownloadUrl(identifier, file), signal); } + + // The item's `<identifier>_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<Buffer> { + 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/<id>/<file>. + // Written to `<dest>.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<void>((r) => out.once("drain", () => r())); + opts.onProgress?.({ bytes, total }); + } + } finally { + await new Promise<void>((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. diff --git a/common/lib/archiveOrgTorrent-server.ts b/common/lib/archiveOrgTorrent-server.ts @@ -0,0 +1,293 @@ +// 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: <file> (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<string, Promise<boolean>>(); + +// `aria2c --version` once per binary per process. ENOENT (or any failure to +// run) is "not installed". +export function aria2cAvailable(bin: string): Promise<boolean> { + let p = probes.get(bin); + if (!p) { + p = new Promise<boolean>((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; + 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<boolean> { + return stat(p).then( + () => true, + () => false, + ); +} + +export async function fetchFileByTorrent(opts: TorrentFetchOpts): Promise<TorrentFetchResult> { + 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 = 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 <gid> <number of files> <path>`. 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, + }); + 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 } : {}) }; + const tail = lastLines.filter((l) => /ERROR|errorCode|Exception/.test(l)).slice(-2).join(" | "); + return { + ok: false, + reason: `aria2c exited ${exit.code ?? exit.signal}${tail ? `: ${tail}` : ""}`, + }; +} diff --git a/common/lib/archiveOrgTorrent.ts b/common/lib/archiveOrgTorrent.ts @@ -0,0 +1,300 @@ +// archive.org OVER BITTORRENT — the pure half: what an item's torrent holds, +// which of its files is the one wanted, the aria2c command line that fetches +// only that file, and what aria2c's progress lines say. +// +// WHY A TORRENT (the operator: "use torrents when possible to be extra polite +// to archive.org"). archive.org derives `<identifier>_archive.torrent` for every +// item, and lists archive.org itself as a WEB SEED in it (`url-list`): a +// torrent client takes pieces from any peer that has them and from archive.org +// for the rest, so what other peers can give never touches archive.org's disks, +// and the file is then SEEDED back for a while. One file of a 160-file item is +// fetched alone (`--select-file`). +// +// The process side — spawning aria2c, watching it stall, seeding, killing it +// by its process group — is lib/archiveOrgTorrent-server.ts; the ladder that +// falls back to a plain download is controller/archiveOrgDownload.ts. + +import { bdecode, bInt, bList, bString, isDict, type BDict, type BValue } from "./bencode"; + +// ─── The torrent ─── + +export type TorrentFileEntry = { + // 1-based, in the torrent's own order: what aria2c's --select-file takes. + index: number; + // The path inside the torrent, `/`-joined (an archive.org torrent's name is + // the identifier and its paths are the item's file names). + path: string; + length: number; + // Byte offset of the file in the torrent's concatenated payload. + offset: number; +}; + +export type ParsedTorrent = { + name: string; + pieceLength: number; + pieceCount: number; + // A single-file torrent has one entry whose path is `name`. + multiFile: boolean; + files: TorrentFileEntry[]; + // BEP 19 web seeds (`url-list`). archive.org's name archive.org itself. + webSeeds: string[]; + trackers: string[]; +}; + +function pathOf(file: BDict): string | null { + // `path.utf-8` (a BitComet extension archive.org also writes) wins over the + // raw `path` when present. + const parts = bList(file["path.utf-8"] ?? file.path) + .map((p) => bString(p)) + .filter((p): p is string => typeof p === "string"); + return parts.length > 0 ? parts.join("/") : null; +} + +function urlList(v: BValue | undefined): string[] { + const one = bString(v); + if (one) return [one]; + return bList(v) + .map((x) => bString(x)) + .filter((x): x is string => !!x); +} + +// Null for anything that is not a torrent with an info dictionary. +export function parseTorrent(buf: Buffer): ParsedTorrent | null { + let root: BValue; + try { + root = bdecode(buf); + } catch { + return null; + } + if (!isDict(root) || !isDict(root.info)) return null; + const info = root.info; + const name = bString(info["name.utf-8"] ?? info.name) ?? ""; + const pieceLength = bInt(info["piece length"]) ?? 0; + const pieces = info.pieces; + const pieceCount = Buffer.isBuffer(pieces) ? Math.floor(pieces.length / 20) : 0; + if (!name || pieceLength <= 0) return null; + const files: TorrentFileEntry[] = []; + let multiFile = false; + if (Array.isArray(info.files)) { + multiFile = true; + let offset = 0; + let index = 0; + for (const f of info.files) { + index++; + if (!isDict(f)) continue; + const length = bInt(f.length) ?? 0; + const p = pathOf(f); + if (p !== null) files.push({ index, path: p, length, offset }); + offset += length; + } + } else { + files.push({ index: 1, path: name, length: bInt(info.length) ?? 0, offset: 0 }); + } + const trackers = [ + ...urlList(root.announce), + ...bList(root["announce-list"]).flatMap((tier) => urlList(tier)), + ]; + return { + name, + pieceLength, + pieceCount, + multiFile, + files, + webSeeds: urlList(root["url-list"]), + trackers: [...new Set(trackers)], + }; +} + +// The torrent's entry for one file of the item, or null when the torrent does +// not carry it (archive.org regenerates an item's torrent when the item +// changes, but a file added since, or one archive.org leaves out, is not in it). +export function findTorrentFile(t: ParsedTorrent, file: string): TorrentFileEntry | null { + return t.files.find((f) => f.path === file) ?? null; +} + +// Where aria2c writes the file under its --dir: `<dir>/<name>/<path>` for a +// multi-file torrent, `<dir>/<name>` for a single-file one. As path segments, +// for the caller to join. +export function torrentFileSegments(t: ParsedTorrent, f: TorrentFileEntry): string[] { + return t.multiFile ? [t.name, ...f.path.split("/")] : [t.name]; +} + +// The pieces the file spans: a file shares its first and last piece with its +// neighbours, which is why aria2c writes a sliver of each (and +// --bt-remove-unselected-file deletes them when the file is complete). +export function torrentFilePieces( + t: ParsedTorrent, + f: TorrentFileEntry, +): { first: number; last: number; count: number } { + const first = Math.floor(f.offset / t.pieceLength); + const last = Math.max(first, Math.floor((f.offset + Math.max(0, f.length - 1)) / t.pieceLength)); + return { first, last, count: last - first + 1 }; +} + +// ─── Politeness ─── + +// settings.json `archiveOrg` (lib/settingsSchema.ts documents each key). +export type ArchiveOrgFetchSettings = { + // Fetch over BitTorrent when the item's torrent carries the file and aria2c + // is installed. False = always the plain download. + torrent: boolean; + // Seed the file for this long after it is complete (0 = do not seed). + seedMinutes: number; + // ...or until this much has been uploaded relative to the file's size, + // whichever comes first (0 = no ratio limit; the time alone ends it). + seedRatio: number; + // No progress for this long while downloading: aria2c is stopped and the + // file is downloaded directly. + stallMinutes: number; + maxPeers: number; + // 0 = unlimited. + maxDownloadKiBps: number; + maxUploadKiBps: number; +}; + +export const DEFAULT_ARCHIVE_ORG_FETCH_SETTINGS: ArchiveOrgFetchSettings = { + torrent: true, + seedMinutes: 10, + seedRatio: 1, + stallMinutes: 5, + maxPeers: 30, + maxDownloadKiBps: 0, + maxUploadKiBps: 0, +}; + +// ─── The aria2c command line ─── + +export type Aria2cSpec = { + torrentPath: string; + // Where aria2c writes (the record's staging dir, on the corpus disk). + dir: string; + fileIndex: number; + settings: ArchiveOrgFetchSettings; + // Called by aria2c when the selected file is complete, BEFORE seeding (its + // documented `--on-bt-download-complete` contract): how the runner knows the + // download is over while aria2c keeps running to seed. + onCompleteHook: string; + userAgent: string; + // aria2c stops on its own if this process dies, so a crashed editor never + // leaves a seeder behind. + parentPid?: number; + summaryIntervalSec?: number; +}; + +// Every flag is spelled `--name=value`, one argv entry each, so nothing here +// is ever word-split. +export function aria2cArgs(spec: Aria2cSpec): string[] { + const s = spec.settings; + const args = [ + `--dir=${spec.dir}`, + `--select-file=${spec.fileIndex}`, + // The slivers of neighbouring files a shared piece leaves are removed when + // the selected file is complete. + "--bt-remove-unselected-file=true", + `--seed-time=${Math.max(0, s.seedMinutes)}`, + `--seed-ratio=${Math.max(0, s.seedRatio).toFixed(1)}`, + `--bt-max-peers=${Math.max(1, Math.floor(s.maxPeers))}`, + `--max-overall-download-limit=${s.maxDownloadKiBps > 0 ? `${Math.floor(s.maxDownloadKiBps)}K` : "0"}`, + `--max-overall-upload-limit=${s.maxUploadKiBps > 0 ? `${Math.floor(s.maxUploadKiBps)}K` : "0"}`, + // One connection to each web seed: archive.org is the web seed. + "--max-connection-per-server=1", + `--user-agent=${spec.userAgent}`, + "--follow-torrent=mem", + "--file-allocation=none", + // A re-run resumes: the partial and its .aria2 control file are kept on a + // cancel, and a complete file (a run killed while seeding) is hash-checked + // and seeded rather than refused or fetched again. + "--continue=true", + "--check-integrity=true", + "--auto-file-renaming=false", + "--bt-save-metadata=false", + `--on-bt-download-complete=${spec.onCompleteHook}`, + `--summary-interval=${spec.summaryIntervalSec ?? 10}`, + "--show-console-readout=false", + "--console-log-level=notice", + "--enable-color=false", + ...(spec.parentPid ? [`--stop-with-process=${spec.parentPid}`] : []), + `--torrent-file=${spec.torrentPath}`, + ]; + return args; +} + +// ─── aria2c's progress lines ─── + +// The bracket line of a progress summary: +// [#09cba8 608KiB/2.8MiB(20%) CN:1 SD:3 DL:302KiB ETA:7s] +// [#09cba8 SEED(0.4) CN:2 SD:0 UL:12KiB(4.9MiB)] +export type Aria2cStatus = + | { + seeding: false; + completedBytes: number; + totalBytes: number; + percent: number; + connections: number; + seeders?: number; + } + | { seeding: true; ratio: number; connections: number; uploadedBytes?: number }; + +const UNITS: Record<string, number> = { B: 1, KiB: 1024, MiB: 1024 ** 2, GiB: 1024 ** 3, TiB: 1024 ** 4 }; + +export function parseAria2cSize(s: string): number | null { + const m = /^([\d.]+)(B|KiB|MiB|GiB|TiB)$/.exec(s.trim()); + if (!m) return null; + const n = Number(m[1]); + return Number.isFinite(n) ? Math.round(n * UNITS[m[2]]) : null; +} + +export function parseAria2cStatusLine(line: string): Aria2cStatus | null { + const m = /^\s*\[#[0-9a-f]+ (.*)\]\s*$/.exec(line); + if (!m) return null; + const body = m[1]; + const field = (k: string) => { + const f = new RegExp(`(?:^| )${k}:([^ \\]]+)`).exec(body); + return f ? f[1] : undefined; + }; + const cn = Number(field("CN") ?? 0) || 0; + const seed = /^SEED\(([\d.]+|inf)\)/.exec(body); + if (seed) { + const ul = /UL:[^(]*\(([^)]+)\)/.exec(body); + const uploaded = ul ? parseAria2cSize(ul[1]) : null; + return { + seeding: true, + ratio: seed[1] === "inf" ? Infinity : Number(seed[1]), + connections: cn, + ...(uploaded !== null ? { uploadedBytes: uploaded } : {}), + }; + } + const prog = /^([\d.]+(?:B|KiB|MiB|GiB|TiB))\/([\d.]+(?:B|KiB|MiB|GiB|TiB))\((\d+)%\)/.exec(body); + if (!prog) return null; + const done = parseAria2cSize(prog[1]); + const total = parseAria2cSize(prog[2]); + if (done === null || total === null) return null; + const sd = field("SD"); + return { + seeding: false, + completedBytes: done, + totalBytes: total, + percent: Number(prog[3]), + connections: cn, + ...(sd !== undefined && Number.isFinite(Number(sd)) ? { seeders: Number(sd) } : {}), + }; +} + +// The line a human follows: "torrent: <file> (n of m pieces, peers p, web +// seed yes)". The piece count is the file's own span; n is read off the bytes +// aria2c reports, so it is the pieces' worth done, not a bitfield. +export function torrentProgressLine(opts: { + file: string; + status: Extract<Aria2cStatus, { seeding: false }>; + pieces: number; + webSeed: boolean; +}): string { + const { status, pieces } = opts; + const frac = status.totalBytes > 0 ? status.completedBytes / status.totalBytes : 0; + const done = Math.min(pieces, Math.floor(frac * pieces)); + return ( + `torrent: ${opts.file} (${done} of ${pieces} pieces, ` + + `peers ${status.connections}${status.seeders !== undefined ? ` (${status.seeders} seeding)` : ""}, ` + + `web seed ${opts.webSeed ? "yes" : "no"})` + ); +} diff --git a/common/lib/bencode.ts b/common/lib/bencode.ts @@ -0,0 +1,133 @@ +// BENCODE — the encoding of a .torrent file (BEP 3), decoded. +// +// Pure and dependency-free: the archive.org torrent fetch reads one file's +// index and byte range out of an item's `<identifier>_archive.torrent` +// (lib/archiveOrgTorrent.ts), and that is all this is for. A byte string stays +// a Buffer (a torrent's `pieces` is raw SHA-1 bytes, and a path may not be +// UTF-8); a dictionary's keys are read as UTF-8 strings, which every key a +// torrent defines is. +// +// STRICT where it is cheap: a truncated or malformed input throws +// BencodeError naming the offset, rather than returning half a value. There is +// no encoder — nothing here writes a torrent. + +export type BValue = number | Buffer | BValue[] | BDict; +export type BDict = { [key: string]: BValue }; + +export class BencodeError extends Error { + constructor(message: string, readonly offset: number) { + super(`${message} at byte ${offset}`); + this.name = "BencodeError"; + } +} + +const CH_I = 0x69; // i +const CH_L = 0x6c; // l +const CH_D = 0x64; // d +const CH_E = 0x65; // e +const CH_COLON = 0x3a; +const CH_MINUS = 0x2d; +const CH_0 = 0x30; +const CH_9 = 0x39; + +// Nesting deeper than this is not a torrent; it is an attempt to blow the stack. +const MAX_DEPTH = 64; + +export function bdecode(buf: Buffer): BValue { + let pos = 0; + + function readInt(end: number): number { + const start = pos; + let neg = false; + if (buf[pos] === CH_MINUS) { + neg = true; + pos++; + } + if (pos >= end) throw new BencodeError("empty integer", start); + let n = 0; + for (; pos < end; pos++) { + const c = buf[pos]; + if (c < CH_0 || c > CH_9) throw new BencodeError("bad digit", pos); + n = n * 10 + (c - CH_0); + } + if (!Number.isSafeInteger(n)) throw new BencodeError("integer too large", start); + return neg ? -n : n; + } + + function indexOf(byte: number, from: number): number { + const i = buf.indexOf(byte, from); + if (i < 0) throw new BencodeError("unterminated value", from); + return i; + } + + function readBytes(): Buffer { + const colon = indexOf(CH_COLON, pos); + const len = readInt(colon); + if (len < 0) throw new BencodeError("negative length", pos); + pos = colon + 1; + if (pos + len > buf.length) throw new BencodeError("string past end of input", pos); + const out = buf.subarray(pos, pos + len); + pos += len; + return out; + } + + function readValue(depth: number): BValue { + if (depth > MAX_DEPTH) throw new BencodeError("nested too deeply", pos); + if (pos >= buf.length) throw new BencodeError("unexpected end of input", pos); + const c = buf[pos]; + if (c === CH_I) { + pos++; + const end = indexOf(CH_E, pos); + const n = readInt(end); + pos = end + 1; + return n; + } + if (c === CH_L) { + pos++; + const list: BValue[] = []; + while (true) { + if (pos >= buf.length) throw new BencodeError("unterminated list", pos); + if (buf[pos] === CH_E) break; + list.push(readValue(depth + 1)); + } + pos++; + return list; + } + if (c === CH_D) { + pos++; + const dict: BDict = {}; + while (true) { + if (pos >= buf.length) throw new BencodeError("unterminated dictionary", pos); + if (buf[pos] === CH_E) break; + const key = readBytes().toString("utf8"); + dict[key] = readValue(depth + 1); + } + pos++; + return dict; + } + if (c >= CH_0 && c <= CH_9) return readBytes(); + throw new BencodeError(`unexpected byte 0x${c.toString(16)}`, pos); + } + + const value = readValue(0); + if (pos !== buf.length) throw new BencodeError("trailing bytes", pos); + return value; +} + +// ─── Reading a decoded value ─── + +export function isDict(v: BValue | undefined): v is BDict { + return !!v && typeof v === "object" && !Array.isArray(v) && !Buffer.isBuffer(v); +} + +export function bString(v: BValue | undefined): string | undefined { + return Buffer.isBuffer(v) ? v.toString("utf8") : undefined; +} + +export function bInt(v: BValue | undefined): number | undefined { + return typeof v === "number" ? v : undefined; +} + +export function bList(v: BValue | undefined): BValue[] { + return Array.isArray(v) ? v : []; +} diff --git a/common/lib/downloadOutcome.ts b/common/lib/downloadOutcome.ts @@ -71,7 +71,13 @@ export type DownloadAttemptKind = // failed ONLY on its subtitle fetch (a 429): the media is what the download // is for, and the subtitles are deferred (release 17, slice RL). Recorded // with n: 1 beside the primary it re-runs. - | "primary-without-subs"; + | "primary-without-subs" + // An archive.org record's file fetched without yt-dlp + // (controller/archiveOrgDownload.ts): over BitTorrent with aria2c, or as a + // plain download from archive.org. `ytdlpExitCode` is 0 for a fetch that + // delivered a verified file and 1 for one that did not. + | "archiveorg-torrent" + | "archiveorg-direct"; export type AudioCheckProbeVerdict = "clean" | "partial" | "malformed"; diff --git a/common/lib/envVars.ts b/common/lib/envVars.ts @@ -78,6 +78,7 @@ const DECLARED: EnvVarDecl[] = [ paths("FINDMNT_BIN", "`findmnt` on PATH", "The read-only volume-identity probe behind storage locations. Optional."), paths("UDISKSCTL_BIN", "`udisksctl` on PATH", "Mounts an attached volume from `/storage`. Optional."), paths("GALLERY_DL_BIN", "`gallery-dl` on PATH", "The X/Twitter post fetcher, for social channels."), + paths("ARIA2C_BIN", "`aria2c` on PATH", "Fetches an archive.org file over BitTorrent (the item's torrent, archive.org as web seed), then seeds it for a while. Optional: without it archive.org files are downloaded directly. See `archiveOrg` in [SETTINGS.md](SETTINGS.md)."), paths("OLLAMA_URL", "`http://127.0.0.1:11434`", "The local ollama server, the local digest and attribution engine."), paths("CLAUDE_BIN", "`claude` on PATH", "The `claude` CLI, driving the opt-in metered digest lane."), paths("ARCHILYZER_CONFIG_DIR", "`~/.config/archilyzer`", "The operator's private config dir, outside the repo: the two inputs of `archilyzer source publish` below. Never committed."), diff --git a/common/lib/metadataHistory.ts b/common/lib/metadataHistory.ts @@ -64,6 +64,10 @@ export const METADATA_HISTORY_WRITERS = [ // archiveOrgMetadataPatch): the file's page and title instead of the item's, // and a mirror's original title, date and uploader. "archiveorg-provenance", + // An archive.org record's metadata.info.json written whole from the item's + // metadata API (controller/archiveOrgDownload.ts) — archive.org records are + // not fetched by yt-dlp. + "archiveorg-import", ] as const; export type MetadataHistoryWriter = (typeof METADATA_HISTORY_WRITERS)[number]; diff --git a/common/lib/paths.ts b/common/lib/paths.ts @@ -124,6 +124,10 @@ export type Paths = { // (Phase 4 of the video-persistence feature). See // common/controller/backupSavedVideos.ts. rsyncBin: string; + // aria2c, which fetches an archive.org file over BitTorrent + // (lib/archiveOrgTorrent-server.ts). Optional: without it the file is + // downloaded directly. + aria2cBin: string; // findmnt (util-linux): the read-only identity probe behind storage // locations — which volume a root is on, and where a UUID is mounted now. // See common/lib/storageVolumes.ts. Never required: every call fails open to @@ -280,6 +284,7 @@ export function getPaths(): Paths { ffmpegBin: process.env.FFMPEG_BIN ?? "ffmpeg", ffprobeBin: process.env.FFPROBE_BIN ?? "ffprobe", rsyncBin: process.env.RSYNC_BIN ?? "rsync", + aria2cBin: process.env.ARIA2C_BIN ?? "aria2c", findmntBin: process.env.FINDMNT_BIN ?? "findmnt", udisksctlBin: process.env.UDISKSCTL_BIN ?? "udisksctl", galleryDlBin: process.env.GALLERY_DL_BIN ?? "gallery-dl", diff --git a/common/lib/settingsDocs.ts b/common/lib/settingsDocs.ts @@ -11,6 +11,7 @@ // a schema edit followed by regenerating, never an edit here. import { + ARCHIVE_ORG_SETTINGS_FIELD_DOCS, ARCHIVE_STORAGE_SETTINGS_FIELD_DOCS, ATTRIBUTION_SETTINGS_FIELD_DOCS, BACKFILL_SETTINGS_FIELD_DOCS, @@ -236,6 +237,13 @@ export function blockTables(d: SiteSettings): Partial<Record<keyof SiteSettings, defaults: fromObject(d.attribution), }, ], + archiveOrg: [ + { + path: "archiveOrg", + docs: ARCHIVE_ORG_SETTINGS_FIELD_DOCS, + defaults: fromObject(d.archiveOrg), + }, + ], }; } diff --git a/common/lib/settingsSchema.test.ts b/common/lib/settingsSchema.test.ts @@ -29,6 +29,7 @@ import type { SourceVideoQuality, } from "../ytdlp/downloadFormat"; import type { StorageSettings } from "./storageLocations"; +import type { ArchiveOrgFetchSettings } from "./archiveOrgTorrent"; import type { AttributionSettings, BackfillSettings, @@ -97,13 +98,15 @@ type PreSchemaSiteSettings = { diarization: DiarizationSettings; backfill: BackfillSettings; attribution: AttributionSettings; + // How archive.org files are fetched: BitTorrent with seeding, else direct. + archiveOrg: ArchiveOrgFetchSettings; }; // Bracketed so the conditional does not distribute (see commit 8c43231). type Same<A, B> = [A] extends [B] ? ([B] extends [A] ? true : false) : false; const shapeUnchanged: Same<SiteSettings, PreSchemaSiteSettings> = true; -test("SiteSettings keeps its 34 fields, in file order", () => { +test("SiteSettings keeps its 35 fields, in file order", () => { assert.equal(shapeUnchanged, true); assert.deepEqual(Object.keys(siteSettingsSchema.shape), [ "adminTitle", @@ -140,6 +143,7 @@ test("SiteSettings keeps its 34 fields, in file order", () => { "diarization", "backfill", "attribution", + "archiveOrg", ]); // A parsed object carries every key, in that order — writeSettings writes // exactly this, so the order is the on-disk order. diff --git a/common/lib/settingsSchema.ts b/common/lib/settingsSchema.ts @@ -73,6 +73,10 @@ import { } from "./diarization"; import { ATTRIBUTION_PROMPT_VERSION } from "./attribution"; import { + DEFAULT_ARCHIVE_ORG_FETCH_SETTINGS, + type ArchiveOrgFetchSettings, +} from "./archiveOrgTorrent"; +import { DEFAULT_COOKIE_MODE, isCookieMode, type CookieMode, @@ -481,6 +485,57 @@ export const BUILD_PIPELINE_SETTINGS_FIELD_DOCS: FieldDocs<BuildPipelineSettings "image.", }; +// Each field is documented in ARCHIVE_ORG_SETTINGS_FIELD_DOCS below (rendered +// into SETTINGS.md). The type and defaults live with the torrent code +// (lib/archiveOrgTorrent.ts), which is pure. +export type { ArchiveOrgFetchSettings }; + +export const ARCHIVE_ORG_SETTINGS_FIELD_DOCS: FieldDocs<ArchiveOrgFetchSettings> = { + torrent: + "Fetch an archive.org file over BitTorrent when the item's " + + "`<identifier>_archive.torrent` carries it and aria2c is installed " + + "(ARIA2C_BIN). archive.org is the torrent's web seed, so what other peers " + + "give never touches archive.org. False = always the direct download.", + seedMinutes: + "Minutes to seed the file after it is complete, as a good swarm citizen. " + + "The import holds archive.org's queue while it seeds. 0 = do not seed. " + + "Clamped to [0, 1440]; default 10.", + seedRatio: + "Stop seeding sooner once this much has been uploaded relative to the " + + "file's size (aria2c --seed-ratio), whichever of the two comes first. " + + "0 = no ratio limit. Clamped to [0, 100]; default 1.", + stallMinutes: + "No download progress for this long stops aria2c and the file is " + + "downloaded directly instead. Clamped to [1, 120]; default 5.", + maxPeers: "Most peers per torrent (aria2c --bt-max-peers). Clamped to [1, 500]; default 30.", + maxDownloadKiBps: + "Download rate cap for a torrent fetch, KiB/s (aria2c " + + "--max-overall-download-limit). 0 = unlimited.", + maxUploadKiBps: + "Upload rate cap while downloading and seeding, KiB/s (aria2c " + + "--max-overall-upload-limit). 0 = unlimited.", +}; + +function clampNumber(value: unknown, fallback: number, min: number, max: number): number { + const n = typeof value === "number" && Number.isFinite(value) ? value : fallback; + return Math.min(max, Math.max(min, n)); +} + +// Coerce a raw settings.archiveOrg value into clean ArchiveOrgFetchSettings. +export function sanitizeArchiveOrg(value: unknown): ArchiveOrgFetchSettings { + const d = DEFAULT_ARCHIVE_ORG_FETCH_SETTINGS; + const r = (value && typeof value === "object" ? value : {}) as Record<string, unknown>; + return { + torrent: r.torrent !== false, + seedMinutes: clampNumber(r.seedMinutes, d.seedMinutes, 0, 1440), + seedRatio: clampNumber(r.seedRatio, d.seedRatio, 0, 100), + stallMinutes: clampNumber(r.stallMinutes, d.stallMinutes, 1, 120), + maxPeers: Math.floor(clampNumber(r.maxPeers, d.maxPeers, 1, 500)), + maxDownloadKiBps: Math.floor(clampNumber(r.maxDownloadKiBps, d.maxDownloadKiBps, 0, 10_000_000)), + maxUploadKiBps: Math.floor(clampNumber(r.maxUploadKiBps, d.maxUploadKiBps, 0, 10_000_000)), + }; +} + // Each field is documented in SAVED_VIDEO_BACKUP_SETTINGS_FIELD_DOCS below (rendered into SETTINGS.md). export type SavedVideoBackupSettings = { enabled: boolean; @@ -1647,6 +1702,9 @@ export const siteSettingsSchema = z.object({ attribution: settingsField((v): AttributionSettings => sanitizeAttribution(v)).describe( "Naming the speakers diarization found (or reconstructing them from the transcript when it found none). OFF by default. See AttributionSettings.", ), + archiveOrg: settingsField((v): ArchiveOrgFetchSettings => sanitizeArchiveOrg(v)).describe( + "How archive.org files are fetched (controller/archiveOrgDownload.ts). Over BitTorrent with aria2c when the item's torrent carries the file — archive.org is the torrent's web seed, so the swarm takes load off archive.org — then seeded for a while; otherwise, or when the torrent stalls, a direct download from archive.org. Either way the file is verified against archive.org's sha1/md5. No yt-dlp.", + ), }); export type SiteSettings = z.infer<typeof siteSettingsSchema>; diff --git a/common/ytdlp/archiveOrgProvenanceHook.test.ts b/common/ytdlp/archiveOrgProvenanceHook.test.ts @@ -1,80 +0,0 @@ -import { test } from "node:test"; -import assert from "node:assert/strict"; -import { chmod, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; -import path from "node:path"; -import type { Paths } from "../lib/paths"; -import type { ChannelConfig } from "../lib/channelConfig"; -import { downloadOneManaged } from "./downloadOneManaged"; -import { archiveOrgDetailsUrl, archiveOrgVideoId } from "../lib/archiveOrgId"; - -// Run with: -// pnpm --filter yt-dlp-transcript-common exec tsx --test ytdlp/archiveOrgProvenanceHook.test.ts -// -// THE archive.org PROVENANCE STEP RUNS RIGHT AFTER A PREFETCH THAT SUCCEEDED, -// and never after one that failed. The yt-dlp here is a temp script: the -// prefetch (`--skip-download`) writes a metadata.info.json shaped like the -// ArchiveOrg extractor's (the ITEM's page for one file of many) and exits 0 or -// 1; the download pass fails. The step itself is a spy. Every name is invented. - -const ITEM = "example-item"; -const FILE = "Clip One.mp4"; -const URL = archiveOrgDetailsUrl({ identifier: ITEM, file: FILE }); -const ID = archiveOrgVideoId({ identifier: ITEM, file: FILE }); - -async function runWith(prefetchExit: 0 | 1) { - const root = await mkdtemp(path.join(tmpdir(), "archiveorg-hook-")); - try { - const bin = path.join(root, "fake-ytdlp.mjs"); - await writeFile( - bin, - `#!/usr/bin/env node -import { mkdirSync, writeFileSync } from "node:fs"; -import path from "node:path"; -const args = process.argv.slice(2); -if (!args.includes("--skip-download")) { console.error("ERROR: no download in this test"); process.exit(1); } -const o = args.find((a) => a.startsWith("infojson:")); -const file = o.slice("infojson:".length) + ".info.json"; -mkdirSync(path.dirname(file), { recursive: true }); -writeFileSync(file, JSON.stringify({ id: ${JSON.stringify(`${ITEM}/${FILE}`)}, extractor_key: "ArchiveOrg", title: "Example Archive", webpage_url: "https://archive.org/details/${ITEM}" })); -process.exit(${prefetchExit}); -`, - ); - await chmod(bin, 0o755); - const paths = { channelsDir: path.join(root, "channels"), ytdlpBin: bin } as Paths; - const calls: { videoDir: string; videoUrl: string; infoAtCall: string }[] = []; - await downloadOneManaged({ - channelSlug: "c", - channelConfig: { handling: "transcribe", platform: "archiveorg" } as ChannelConfig, - paths, - videoUrl: URL, - onLog: () => {}, - signal: new AbortController().signal, - archiveOrgProvenance: async (o) => { - calls.push({ - videoDir: o.videoDir, - videoUrl: o.videoUrl, - infoAtCall: await readFile(path.join(o.videoDir, "metadata.info.json"), "utf8").catch(() => ""), - }); - return null; - }, - }); - return { calls, dataDir: path.join(paths.channelsDir, "c", "data") }; - } finally { - await rm(root, { recursive: true, force: true }); - } -} - -test("after a prefetch that succeeded, the step runs once on the record's own dir", async () => { - const { calls, dataDir } = await runWith(0); - assert.equal(calls.length, 1); - assert.equal(calls[0].videoDir, path.join(dataDir, ID)); - assert.equal(calls[0].videoUrl, URL); - // The prefetch's record is already on disk when it runs. - assert.match(calls[0].infoAtCall, /"extractor_key":"ArchiveOrg"/); -}); - -test("after a prefetch that failed, it does not run", async () => { - const { calls } = await runWith(1); - assert.equal(calls.length, 0); -}); diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts @@ -64,7 +64,8 @@ import { type MetadataScanEntry, } from "../controller/metadataScanStore"; import { detectPlatform, type Platform } from "../lib/platform"; -import { ensureArchiveOrgProvenance } from "../lib/archiveOrg-server"; +import { downloadArchiveOrgManaged } from "../controller/archiveOrgDownload"; +import { finalizeAppExtraction } from "./finalizeAppExtraction"; import { probeMediaDurationSec } from "./ffprobeDuration"; import { isShortAudio, @@ -185,9 +186,6 @@ export type ManagedDownloadOpts = { // "video_720" = the ≤720p H.264 selector (downloadFormat.ts). Never touches // the audio-only selector above. persistFormatPreset?: SourceVideoQuality; - // Test seam: the archive.org provenance step (lib/archiveOrg-server.ts). - // Every production caller passes none. - archiveOrgProvenance?: typeof ensureArchiveOrgProvenance; }; // When `reuseInfoJson` is true, the real download reuses the metadata the @@ -316,111 +314,6 @@ function transcribeMediaArgs( ]; } -// After an app-mode transcribe download, produce audio.<fmt> from the downloaded -// source-media.<ext> container via ffmpeg. When the video is persisted the -// container is MOVED into the saved-video store (Phase 3) and a pointer is left -// in the data dir; otherwise the container is removed (the extract-now / -// save-disk path). Best-effort throughout: a failed extraction leaves the -// container in place (so a later pass or the transcribe-from-container fallback -// can still recover), and a failed store-move leaves the container in the data -// dir as source-media.<ext> (still persisted, just not relocated). Neither fails -// the whole download. -async function finalizeAppExtraction(opts: { - paths: Paths; - channelSlug: string; - channelConfig: ChannelConfig; - videoDir: string; - videoId: string; - fmt: AudioFormat; - persist: boolean; - // The persistence cause, recorded on the saved-video pointer (governs the - // retention prune). Only meaningful when persist is true. - category: PersistenceDecision["category"]; - // Threaded straight onto the pointer; see ManagedDownloadOpts.persistOrigin. - origin?: SavedVideoOrigin; - // PERSIST ONLY when false: move the container, extract nothing. A forced - // media download of a video that already has a transcript (attempt 3, - // `keepTranscript`) wants the source kept, and an audio.<fmt> beside an - // existing transcript is bytes nothing will read. Default true. - extractAudio?: boolean; - // What the persist asked for and what yt-dlp took (the pass's FORMAT_MARKER - // line), recorded on the pointer. Absent = not reported. - format?: SavedVideoFormat | null; - onLog: (s: string) => void; - signal: AbortSignal; -}): Promise<void> { - const entries = await readdir(opts.videoDir).catch(() => [] as string[]); - const source = findSourceMedia(entries); - if (!source) { - opts.onLog( - `App extraction: no source-media container in ${opts.videoDir}; skipping.\n`, - ); - return; - } - if (opts.extractAudio === false) { - opts.onLog( - `Persist only: a transcript is on disk, so no audio.${opts.fmt} is extracted from ${source}.\n`, - ); - if (!opts.persist) { - // Unreachable from attempt 3 today (it does not download when the plan - // persists nothing), kept so the option means one thing on its own. - opts.onLog(`Not persisting ${source}; it stays in the data dir.\n`); - return; - } - } else { - try { - await transcodeAudio({ - paths: opts.paths, - videoDir: opts.videoDir, - sourceFilename: source, - targetFormat: opts.fmt, - onLog: opts.onLog, - signal: opts.signal, - }); - } catch (err) { - opts.onLog( - `App extraction failed (${source} -> audio.${opts.fmt}): ${(err as Error).message}. Keeping the source container.\n`, - ); - return; - } - if (!opts.persist) { - await removeMediaFile(opts.videoDir, source); - opts.onLog(`Discarded source container ${source} (audio-only).\n`); - return; - } - } - // Persist: move the container into the saved-video store + write a pointer. - const storeDir = savedVideoDir( - opts.paths, - opts.channelConfig, - opts.channelSlug, - opts.videoId, - ); - try { - const pointer = await persistSourceVideo({ - videoDir: opts.videoDir, - sourceFilename: source, - storeDir, - keepReason: opts.category === "none" ? undefined : opts.category, - origin: opts.origin, - ...(opts.format ? { format: opts.format } : {}), - // The store-in-transition guard. A move of the saved-video store is a - // multi-hour copy, and a container landing in the middle of it is lost - // three different ways — see lib/savedVideoStore.ts. The catch below is - // already the right handling: the container stays in the data dir with a - // line in the log, and the next persist (after the move) picks it up. - paths: opts.paths, - }); - opts.onLog( - `Persisted source video to ${path.join(pointer.dir, pointer.file)} (${pointer.bytes} bytes).\n`, - ); - } catch (err) { - opts.onLog( - `Failed to move source container into the saved-video store (${(err as Error).message}); leaving ${source} in the data dir.\n`, - ); - } -} - // The source a download attempt reads from: either the prefetched info json // (no positional URL needed) or the URL itself. Mirrors how the no-subs // fallback has always fed yt-dlp a --load-info-json instead of a URL. @@ -500,9 +393,6 @@ const PREFETCH_OWN_FILES: ReadonlySet<string> = new Set([ // one a previous rejection wrote before this rule existed. It records what // happened, never what is on disk, so it is ours to drop with the rest. "download-outcome.json", - // An archive.org record's provenance, written by this pass right after the - // prefetch (lib/archiveOrg-server.ts) and fetchable again. - "archiveorg.json", // NOT metadata.history.json, deliberately: a history means an earlier // metadata.info.json was here before this pass, and it is the one record of // what the source used to say (lib/metadataHistory-server.ts). @@ -753,6 +643,13 @@ export async function downloadOneManaged( } try { + // AN archive.org RECORD IS NOT A yt-dlp DOWNLOAD: its file comes over + // BitTorrent or straight from archive.org, verified against the item's + // checksums, and the record is built from the item's metadata + // (controller/archiveOrgDownload.ts). Same log, same outcome sidecar. + if (detectPlatform(opts.videoUrl) === "archiveorg") { + return await downloadArchiveOrgManaged(opts); + } return await runManagedDownload(opts, channelDir, startedAt, canonicalId); } finally { logStream?.end(); @@ -933,22 +830,6 @@ async function runManagedDownload( } const metaPath = path.join(videoDir, "metadata.info.json"); - // AN archive.org RECORD IS CORRECTED BEFORE ANYTHING READS IT: its - // provenance sidecar written (one cached metadata request per item), and - // metadata.info.json given the file's own page and title — yt-dlp writes - // the ITEM's for one file of a multi-file item (lib/archiveOrg-server.ts). - // Only after a prefetch that succeeded: a refused one is not a record. - if ( - detectPlatform(opts.videoUrl) === "archiveorg" && - attemptSucceeded(attempts.at(-1)?.ytdlpExitCode ?? null) - ) { - await (opts.archiveOrgProvenance ?? ensureArchiveOrgProvenance)({ - videoDir, - videoUrl: opts.videoUrl, - onLog: opts.onLog, - signal: opts.signal, - }); - } const metadata = await loadRawMetadata(metaPath); // Only wire --load-info-json into the real attempts when we actually have // the metadata file; a failed prefetch falls through to the legacy path so @@ -1263,16 +1144,6 @@ async function runManagedDownload( runAudioCheck, ) : await runAudioCheck(); - // The audio-checked pass re-extracted, so the archive.org correction above - // is re-applied (from the sidecar; no request). - if (canonicalId && detectPlatform(opts.videoUrl) === "archiveorg") { - await (opts.archiveOrgProvenance ?? ensureArchiveOrgProvenance)({ - videoDir: path.join(channelDir, "data", canonicalId), - videoUrl: opts.videoUrl, - onLog: opts.onLog, - signal: opts.signal, - }); - } primaryRes = { exitCode: audioOutcome.ytdlpExitCode, stderrTail: audioOutcome.stderrTail, diff --git a/common/ytdlp/finalizeAppExtraction.ts b/common/ytdlp/finalizeAppExtraction.ts @@ -0,0 +1,129 @@ +// THE APP'S AUDIO EXTRACTION, after a download that left a source container. +// Split out of ytdlp/downloadOneManaged.ts so the archive.org fetch +// (controller/archiveOrgDownload.ts), which downloads no-yt-dlp files into the +// same data/<id>/, finalises them through the same code: audio.<fmt> by +// ffmpeg, the container kept in the saved-video store or removed. + +import path from "node:path"; +import { readdir } from "node:fs/promises"; +import { removeMediaFile } from "../lib/mediaTier-server"; +import type { AudioFormat, ChannelConfig } from "../lib/channelConfig"; +import type { Paths } from "../lib/paths"; +import type { PersistenceDecision } from "./persistencePlan"; +import { transcodeAudio } from "../controller/transcode"; +import { findSourceMedia } from "../lib/videoStatus"; +import { + savedVideoDir, + type SavedVideoFormat, + type SavedVideoOrigin, +} from "../lib/savedVideo"; +import { persistSourceVideo } from "../lib/savedVideo-server"; + +// After an app-mode transcribe download, produce audio.<fmt> from the downloaded +// source-media.<ext> container via ffmpeg. When the video is persisted the +// container is MOVED into the saved-video store (Phase 3) and a pointer is left +// in the data dir; otherwise the container is removed (the extract-now / +// save-disk path). Best-effort throughout: a failed extraction leaves the +// container in place (so a later pass or the transcribe-from-container fallback +// can still recover), and a failed store-move leaves the container in the data +// dir as source-media.<ext> (still persisted, just not relocated). Neither fails +// the whole download. +export async function finalizeAppExtraction(opts: { + paths: Paths; + channelSlug: string; + channelConfig: ChannelConfig; + videoDir: string; + videoId: string; + fmt: AudioFormat; + persist: boolean; + // The persistence cause, recorded on the saved-video pointer (governs the + // retention prune). Only meaningful when persist is true. + category: PersistenceDecision["category"]; + // Threaded straight onto the pointer; see ManagedDownloadOpts.persistOrigin. + origin?: SavedVideoOrigin; + // PERSIST ONLY when false: move the container, extract nothing. A forced + // media download of a video that already has a transcript (attempt 3, + // `keepTranscript`) wants the source kept, and an audio.<fmt> beside an + // existing transcript is bytes nothing will read. Default true. + extractAudio?: boolean; + // What the persist asked for and what yt-dlp took (the pass's FORMAT_MARKER + // line), recorded on the pointer. Absent = not reported. + format?: SavedVideoFormat | null; + // The container's name in videoDir, when the caller knows it (the + // archive.org fetch). Absent: the sorted first source-media.<ext> there. + sourceFilename?: string; + onLog: (s: string) => void; + signal: AbortSignal; +}): Promise<void> { + const source = + opts.sourceFilename ?? + findSourceMedia(await readdir(opts.videoDir).catch(() => [] as string[])); + if (!source) { + opts.onLog( + `App extraction: no source-media container in ${opts.videoDir}; skipping.\n`, + ); + return; + } + if (opts.extractAudio === false) { + opts.onLog( + `Persist only: a transcript is on disk, so no audio.${opts.fmt} is extracted from ${source}.\n`, + ); + if (!opts.persist) { + // Unreachable from attempt 3 today (it does not download when the plan + // persists nothing), kept so the option means one thing on its own. + opts.onLog(`Not persisting ${source}; it stays in the data dir.\n`); + return; + } + } else { + try { + await transcodeAudio({ + paths: opts.paths, + videoDir: opts.videoDir, + sourceFilename: source, + targetFormat: opts.fmt, + onLog: opts.onLog, + signal: opts.signal, + }); + } catch (err) { + opts.onLog( + `App extraction failed (${source} -> audio.${opts.fmt}): ${(err as Error).message}. Keeping the source container.\n`, + ); + return; + } + if (!opts.persist) { + await removeMediaFile(opts.videoDir, source); + opts.onLog(`Discarded source container ${source} (audio-only).\n`); + return; + } + } + // Persist: move the container into the saved-video store + write a pointer. + const storeDir = savedVideoDir( + opts.paths, + opts.channelConfig, + opts.channelSlug, + opts.videoId, + ); + try { + const pointer = await persistSourceVideo({ + videoDir: opts.videoDir, + sourceFilename: source, + storeDir, + keepReason: opts.category === "none" ? undefined : opts.category, + origin: opts.origin, + ...(opts.format ? { format: opts.format } : {}), + // The store-in-transition guard. A move of the saved-video store is a + // multi-hour copy, and a container landing in the middle of it is lost + // three different ways — see lib/savedVideoStore.ts. The catch below is + // already the right handling: the container stays in the data dir with a + // line in the log, and the next persist (after the move) picks it up. + paths: opts.paths, + }); + opts.onLog( + `Persisted source video to ${path.join(pointer.dir, pointer.file)} (${pointer.bytes} bytes).\n`, + ); + } catch (err) { + opts.onLog( + `Failed to move source container into the saved-video store (${(err as Error).message}); leaving ${source} in the data dir.\n`, + ); + } +} diff --git a/settings.json.example b/settings.json.example @@ -187,5 +187,14 @@ "diarizedEnabled": false, "textOnlyEnabled": false, "promptVersion": 1 + }, + "archiveOrg": { + "torrent": true, + "seedMinutes": 10, + "seedRatio": 1, + "stallMinutes": 5, + "maxPeers": 30, + "maxDownloadKiBps": 0, + "maxUploadKiBps": 0 } }