// `archilyzer seed` — the home seeder of last resort (release 21 D4b), and // `archilyzer tracker`, the pilot's self-hosted tracker. // // THE SEEDER. WebTorrent 3 in node: TCP for desktop clients, WebRTC (through // node-datachannel, which WebTorrent's peer library loads) for browsers. It // seeds the playable torrents (controller/preparePlayable.ts) of every channel // of the sites named in `settings.seeder.sites` — and, per torrent, only while // no other seeder has it: lib/lastResortSeeder.ts decides, this module // observes (each tracker's scrape, its own wires) and applies. STANDBY is the // torrent removed from the client with its data kept (`stopped` is announced, // every peer closed); SEEDING is it added back, trusting the copy on disk // (`skipVerify`: the manifest's size was checked, and every peer verifies each // piece it takes). // // No DHT, no local discovery, no UPnP/NAT-PMP, no web seeds: the seeder is // found through the trackers it is given and nowhere else, and it opens no // port on a router. `bindInterface`, when set, must exist (the VPN tunnel): // the seeder waits for it to start, and drops every torrent the moment it is // gone — the docker profile's namespace is the real kill switch // (docker-compose.seeder.yml); this is the same rule in the app. // // THE TRACKER. bittorrent-tracker's server, HTTP + WebSocket (no UDP), bound // to the address given (127.0.0.1 by default), tracking ONLY the infohashes // in the playable manifests of the seeder's sites — an open tracker on a home // connection would be anyone's. // // SERVER-ONLY. /// import os from "node:os"; import path from "node:path"; import { readFile, stat } from "node:fs/promises"; import WebTorrent, { type Torrent } from "webtorrent"; import { Client as TrackerClient, Server as TrackerServer, type ScrapeData } from "bittorrent-tracker"; import parseTorrent from "parse-torrent"; import type { Paths } from "../lib/paths"; import { getSite, siteChannelSlugs } from "../lib/site"; import { readChannelConfig } from "./channels"; import { defaultPlayableRoot, loadPlayableManifest } from "./preparePlayable"; import { initialSeederMemory, stepSeeder, type SeederMemory, type SeederPolicy, type SwarmObservation, } from "../lib/lastResortSeeder"; import { seederConfigured, type SeederSettings } from "../lib/seederSettings"; export type SeedTarget = { slug: string; id: string; infoHash: string; torrentPath: string; // The torrent's one file (`.`), and the directory holding it. file: string; dir: string; bytes: number; }; // Every playable torrent of the sites' channels, deduplicated by infohash. // A site that does not load, or a channel with no manifest, contributes none. export async function listSeedTargets( paths: Paths, settings: Pick, onLog: (line: string) => void = () => {}, ): Promise { const out = new Map(); const slugs = new Set(); for (const siteId of settings.sites) { try { for (const slug of siteChannelSlugs(getSite(siteId, paths))) slugs.add(slug); } catch (err) { onLog(`site ${siteId}: ${(err as Error).message} — skipped`); } } for (const slug of [...slugs].sort()) { const config = await readChannelConfig(paths, slug).catch(() => null); const root = await defaultPlayableRoot(paths, config ?? null); const manifest = await loadPlayableManifest(root, slug); for (const [id, e] of Object.entries(manifest.entries)) { if (out.has(e.infoHash)) continue; out.set(e.infoHash, { slug, id, infoHash: e.infoHash, torrentPath: path.join(root, slug, e.torrent), file: path.join(root, slug, e.file), dir: path.dirname(path.join(root, slug, e.file)), bytes: e.bytes, }); } } return [...out.values()]; } export function interfacePresent( name: string, interfaces: NodeJS.Dict = os.networkInterfaces(), ): boolean { return (interfaces[name] ?? []).some((a) => !a.internal); } export type ScrapeFn = ( announce: string, infoHashes: string[], ) => Promise>; // One scrape request per tracker, for every infohash at once. export const scrapeTracker: ScrapeFn = (announce, infoHashes) => new Promise((resolve, reject) => { if (infoHashes.length === 0) return resolve(new Map()); TrackerClient.scrape({ announce, infoHash: infoHashes }, (err, data) => { if (err) return reject(err); const rows: ScrapeData[] = infoHashes.length === 1 ? [data as ScrapeData] : Object.values(data as Record); resolve(new Map(rows.map((r) => [r.infoHash, { complete: r.complete, incomplete: r.incomplete }]))); }); }); // Several trackers' answers for one torrent → one: the most seeders and the // most leechers any tracker reported (the swarms overlap; a sum would count // one browser on two trackers twice). Null when no tracker answered. export function mergeScrapes( answers: (Map | null)[], infoHash: string, ): { complete: number; incomplete: number } | null { let out: { complete: number; incomplete: number } | null = null; for (const a of answers) { const row = a?.get(infoHash); if (!row) continue; out = out ? { complete: Math.max(out.complete, row.complete), incomplete: Math.max(out.incomplete, row.incomplete) } : { ...row }; } return out; } export type RunSeederOptions = { paths: Paths; settings: SeederSettings; onLog: (line: string) => void; signal?: AbortSignal; deps?: { now?: () => number; sleep?: (ms: number, signal?: AbortSignal) => Promise; scrape?: ScrapeFn; interfaces?: () => NodeJS.Dict; targets?: () => Promise; // Overrides for a test's clock: the poll and the policy's windows, in ms. pollMs?: number; policy?: SeederPolicy; // Stop after this many polls (tests). maxPolls?: number; onTransition?: (t: { infoHash: string; to: string; reason: string }) => void; client?: WebTorrent; }; }; const defaultSleep = (ms: number, signal?: AbortSignal) => new Promise((resolve) => { const t = setTimeout(resolve, ms); signal?.addEventListener("abort", () => { clearTimeout(t); resolve(); }, { once: true }); }); function addTorrent(client: WebTorrent, buf: Buffer, target: SeedTarget, trackers: string[]): Promise { return new Promise((resolve, reject) => { const t = client.add( buf, { path: target.dir, skipVerify: true, ...(trackers.length ? { announce: trackers } : {}) }, (torrent) => resolve(torrent), ); t.once("error", reject); }); } function removeTorrent(client: WebTorrent, torrent: Torrent): Promise { return new Promise((resolve) => { client.remove(torrent, { destroyStore: false }, () => resolve()).catch(() => resolve()); }); } export async function runSeeder(opts: RunSeederOptions): Promise { const { settings, onLog } = opts; const d = opts.deps ?? {}; const now = d.now ?? Date.now; const sleep = d.sleep ?? defaultSleep; const scrape = d.scrape ?? scrapeTracker; const ifaces = d.interfaces ?? (() => os.networkInterfaces()); const pollMs = d.pollMs ?? settings.pollSeconds * 1000; const policy: SeederPolicy = d.policy ?? { standbyAfterSeconds: settings.standbyAfterSeconds, resumeAfterSeconds: settings.standbyAfterSeconds, }; if (!seederConfigured(settings)) { throw new Error("the seeder is not configured: settings.seeder.sites is empty"); } const client = d.client ?? new WebTorrent({ maxConns: settings.maxConnections, uploadLimit: settings.maxUploadKiBps > 0 ? settings.maxUploadKiBps * 1024 : -1, dht: false, lsd: false, natUpnp: false, natPmp: false, utp: false, webSeeds: false, }); onLog(`seeder: WebRTC ${WebTorrent.WEBRTC_SUPPORT ? "available" : "NOT available (TCP peers only)"}; sites ${settings.sites.join(", ")}`); type Entry = { target: SeedTarget; memory: SeederMemory; torrent: Torrent | null; trackers: string[]; buf: Buffer }; const entries = new Map(); let tunnelUp = true; let polls = 0; try { while (!opts.signal?.aborted) { polls += 1; // The tunnel, first: no tunnel, no torrents. if (settings.bindInterface) { const up = interfacePresent(settings.bindInterface, ifaces()); if (!up && tunnelUp) { onLog(`seeder: ${settings.bindInterface} is gone — every torrent stopped until it is back`); for (const e of entries.values()) { if (e.torrent) await removeTorrent(client, e.torrent); e.torrent = null; } } else if (up && !tunnelUp) { onLog(`seeder: ${settings.bindInterface} is back`); } tunnelUp = up; } if (tunnelUp) { // The targets, re-read each poll: a new prepare adds torrents. const targets = await (d.targets ?? (() => listSeedTargets(opts.paths, settings, onLog)))(); const wanted = new Set(targets.map((t) => t.infoHash)); for (const [hash, e] of entries) { if (!wanted.has(hash)) { if (e.torrent) await removeTorrent(client, e.torrent); entries.delete(hash); onLog(`${e.target.id}: no longer in a playable manifest — dropped`); } } for (const target of targets) { if (entries.has(target.infoHash)) continue; // skipVerify trusts the copy, so the copy must at least be whole. const st = await stat(target.file).catch(() => null); if (!st || st.size !== target.bytes) { onLog(`${target.id}: ${target.file} is ${st ? `${st.size} bytes, not ${target.bytes}` : "not there"} — skipped`); continue; } const buf = await readFile(target.torrentPath).catch(() => null); if (!buf) { onLog(`${target.id}: ${target.torrentPath} is not there — skipped`); continue; } const own = (await parseTorrent(buf)).announce ?? []; const trackers = settings.trackers.length ? settings.trackers : own; entries.set(target.infoHash, { target, memory: initialSeederMemory(now()), torrent: null, trackers, buf }); } // Every torrent the machine says to seed and is not seeding, added. for (const e of entries.values()) { if (e.memory.state === "seeding" && !e.torrent) { try { e.torrent = await addTorrent(client, e.buf, e.target, e.trackers); onLog(`${e.target.id}: seeding (${e.target.infoHash})`); } catch (err) { onLog(`${e.target.id}: could not seed — ${(err as Error).message}`); } } } // Observe and decide. const byTracker = new Map(); for (const e of entries.values()) { for (const tr of e.trackers) byTracker.set(tr, [...(byTracker.get(tr) ?? []), e.target.infoHash]); } const answers = new Map | null>(); for (const [tr, hashes] of byTracker) { answers.set(tr, await scrape(tr, hashes).catch((err: Error) => { onLog(`scrape ${tr}: ${err.message}`); return null; })); } for (const e of entries.values()) { const wires = e.torrent?.wires ?? []; // By peer, not by wire: two peers that dialled each other at once // hold two wires. const peers = (seeds: boolean) => new Set(wires.filter((w) => w.isSeeder === seeds).map((w, i) => w.peerId ?? w.remoteAddress ?? `wire${i}`)).size; const o: SwarmObservation = { at: now(), scrape: mergeScrapes(e.trackers.map((tr) => answers.get(tr) ?? null), e.target.infoHash), selfCounted: e.memory.state === "seeding" && !!e.torrent, connectedSeeds: peers(true), connectedLeechers: peers(false), }; const step = stepSeeder(e.memory, o, policy); e.memory = step.memory; if (!step.transition) continue; onLog(`${e.target.id}: ${step.transition.from} → ${step.transition.to} — ${step.transition.reason} (${step.seen})`); d.onTransition?.({ infoHash: e.target.infoHash, to: step.transition.to, reason: step.transition.reason }); if (step.transition.to === "standby" && e.torrent) { await removeTorrent(client, e.torrent); e.torrent = null; } else if (step.transition.to === "seeding" && !e.torrent) { try { e.torrent = await addTorrent(client, e.buf, e.target, e.trackers); } catch (err) { onLog(`${e.target.id}: could not resume — ${(err as Error).message}`); } } } } if (d.maxPolls && polls >= d.maxPolls) break; await sleep(pollMs, opts.signal); } } finally { if (!d.client) await new Promise((resolve) => client.destroy(() => resolve())); } } // --- the tracker ------------------------------------------------------------ export type RunningTracker = { port: number; host: string; // What the tracker itself knows of one swarm (null: never announced). swarm: (infoHash: string) => { complete: number; incomplete: number } | null; close: () => Promise; }; export async function startTracker(opts: { host: string; port: number; // The infohashes it tracks; anything else is refused. Re-read on every // announce through this function, so a new prepare needs no restart. allowed: () => Promise> | Set; onLog?: (line: string) => void; }): Promise { const server = new TrackerServer({ udp: false, http: true, ws: true, stats: false, trustProxy: false, filter: (infoHash, _params, cb) => { Promise.resolve(opts.allowed()) .then((set) => cb(set.has(infoHash) ? null : new Error("this tracker does not track that torrent"))) .catch((err: Error) => cb(err)); }, }); server.on("error", (err: Error) => opts.onLog?.(`tracker: ${err.message}`)); server.on("warning", (err: Error) => opts.onLog?.(`tracker warning: ${err.message}`)); await new Promise((resolve) => server.listen(opts.port, opts.host, () => resolve())); const address = server.http?.address(); const port = address?.port ?? opts.port; opts.onLog?.(`tracker: http://${opts.host}:${port}/announce and ws://${opts.host}:${port}`); return { port, host: opts.host, swarm: (infoHash) => { const t = server.torrents[infoHash]; return t ? { complete: t.complete, incomplete: t.incomplete } : null; }, close: () => new Promise((resolve) => server.close(() => resolve())), }; }