Archilyzer · Source

archilyzer

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

commit 0baa0fa1e03f3442c4e848ae3e20c2b63f84ea09
parent 8ccd3dc569453445937df267bcc0f76c0a1ba98a
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sat, 10 Oct 2026 02:22:49 -0400

ops: archive.org imports take many items or a search; Odysee/BitChute remote listings (release 19 A6)

import-archive-org takes {item} as before, {items: [...]} or {query, limit?}
(archive.org's advanced search, one request through the polite client) as one
job: the download gap and failure count shared across the batch, a jittered
pause between items, records already held (on disk or in the saved-video
store) skipped, files archive.org marks not for download listed as
RESTRICTED by a dry run and skipped by an import, a 401/403 skipping the rest
of its item and three such items in a row stopping the job with the platform
backed off. The log ends with a summary: line.

remote-listing {slug} (and `pnpm ops get remote-listing <slug>`, waited on,
the result on stdout) lists an Odysee or BitChute channel with one
flat-playlist read on the platform's queue, refused while the platform is
held or cooling down, after its import floor, backing it off on a 429, and
answers what it lists that is not held and what is held but not listed.
Nothing is written.

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

Diffstat:
MCOMMANDS.md | 4+++-
MOPERATING.md | 12++++++++++--
Mcommon/controller/archiveOrgImport.test.ts | 187+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/archiveOrgImport.ts | 503++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Acommon/controller/remoteListing.test.ts | 87+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/remoteListing.ts | 133+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/jobKinds.ts | 13+++++++++++++
Aeditor/app/api/ops/import-archive-org/route.test.ts | 52++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/api/ops/import-archive-org/route.ts | 98+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Aeditor/app/api/ops/remote-listing/route.test.ts | 74++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/ops/remote-listing/route.ts | 16++++++++++++++++
Meditor/app/channels/[slug]/pipelineActions.ts | 211+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mscripts/archilyzer-ops.mjs | 68+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mscripts/archilyzer-ops.test.mjs | 26++++++++++++++++++++++++++
14 files changed, 1326 insertions(+), 158 deletions(-)

diff --git a/COMMANDS.md b/COMMANDS.md @@ -86,7 +86,8 @@ Drives a running editor over HTTP (`/api/ops/*`, the same actions its pages run) | `metadata-scan` | — | `pnpm ops metadata-scan --json '{"slug":"the-quartering"}'` | | `refresh-metadata` | refresh-metadata re-reads ONE video's metadata.info.json from its source (no subtitles, no media) on the platform's queue: {"slug", "id"}. The job's log ends with what the source now says — live\_status, formats, audio-only formats and whether any is non-fragmented, English captions, the keys that changed. An id with no data/&lt;id&gt;/ is refused (a refresh re-reads a video already archived), as are archive.org and Wayback records. | `pnpm ops refresh-metadata --json '{"slug":"the-quartering","id":"<videoId>"}' --wait` | | `import-video` | — | `pnpm ops import-video --json '{"slug":"demo-archive","url":"https://archive.org/details/example-item"}'` | -| `import-archive-org` | — | `pnpm ops import-archive-org --json '{"slug":"demo-archive","item":"example-item","match":"\\.mp4$"}' --wait` | +| `import-archive-org` | import-archive-org imports archive.org media into a channel, as one job on archive.org's queue: {"slug", "item", "files": \[...\] \| "match": "&lt;regex&gt;"} for one item; {"slug", "items": \["&lt;id&gt;" \| {"item", "files"? \| "match"?}, ...\]} for many (a bare id takes "match", else every media original); or {"slug", "query": "&lt;archive.org search&gt;", "limit"?: 100} for the first items a search finds (at most 500). One file at a time, a jittered pause between files and between items; a record already held (on disk, or in the saved-video store) is skipped. "dryRun": true lists each file as held, RESTRICTED (archive.org marks it not for download: a fetch would answer 401/403) or would get, and fetches nothing. A 401/403 skips the rest of its item; three such items in a row, a 429 or three failures in a row stop the job. The log ends with a summary: line; a re-run resumes. | `pnpm ops import-archive-org --json '{"slug":"demo-archive","item":"example-item","match":"\\.mp4$"}' --wait`<br>`pnpm ops import-archive-org --json '{"slug":"demo-archive","query":"collection:example-collection","dryRun":true}' --wait` | +| `remote-listing` | remote-listing lists an Odysee or BitChute channel upstream and diffs it against what it holds: {"slug"}. One flat-playlist read on the platform's queue (refused while it is held or cooling down; a 429 backs it off), nothing written. `get remote-listing <slug>` waits for the job and prints {listed, held, notHeld: \[{id, url}\], heldNotListed: \[id\], ...} on stdout. | | | `attach-media` | attach-media copies each held video's file out of a LOCAL archive into the saved-video store, as its source container (nothing fetched): {"slug", "source"} — an absolute path to a directory, a .zip (read in place) or a .7z. The id is the folder's trailing "(&lt;id&gt;)", else the file's yt-dlp suffix; "items": \[{"id", "path"}\] names exact files (path inside the source), "match" narrows the folders by regex. "createRecords": true writes a record for a video the channel does not hold; "replace": true re-attaches over a saved container. "dryRun": true logs each folder's class (attach, not-held, already-attached, lost, no-media, unmatched, ambiguous) and the held videos with no media in it, and writes nothing. The log ends with a summary: line; a re-run resumes. | | | `prepare-playable` | prepare-playable remuxes each of a channel's saved containers, losslessly (-c copy), into a browser-playable mp4 (+faststart) or webm, and makes one single-file torrent per copy, under playable/ beside the saved-video store: {"slug"}. "ids" narrows it; "trackers": \[...\] is each .torrent's announce list (none by default; the infohash does not depend on it); "root" names another playable root. A video prepared from the same source (by sha256) is skipped; "dryRun": true logs each decision. | | | `feed-metadata` | — | `pnpm ops feed-metadata --json '{"slug":"demo-podcast","dryRun":true}' --wait` | @@ -141,6 +142,7 @@ Usage: pnpm ops <action> [--json '<body>' | --file <path>] [--wait] pnpm ops get settings [<key>] pnpm ops get storage | sites | workers | auto-queue | scheduler pnpm ops get cleanup <slug> + pnpm ops get remote-listing <slug> [--wait-timeout <seconds>] pnpm ops list --wait follows the job's log and survives a poll that fails (a busy diff --git a/OPERATING.md b/OPERATING.md @@ -60,13 +60,21 @@ pnpm ops persist-videos --json '{"items":[{"slug":"example","id":"<videoId>"}] ```sh pnpm ops import-archive-org --json '{"slug":"example-archive","item":"<item>","match":"\\.mp4$","dryRun":true}' --wait pnpm ops import-archive-org --json '{"slug":"example-archive","item":"<item>","match":"\\.mp4$"}' --wait +pnpm ops import-archive-org --json '{"slug":"example-archive","items":["<item>","<item>"],"dryRun":true}' --wait +pnpm ops import-archive-org --json '{"slug":"example-archive","query":"collection:<collection>","limit":50,"dryRun":true}' --wait +pnpm ops get remote-listing example-odysee pnpm ops import-video --json '{"slug":"example","url":"https://www.bitchute.com/video/<id>/"}' --wait pnpm ops import-video --json '{"slug":"example","url":"https://odysee.com/@example:0/<video>:0"}' --wait pnpm ops import-video --json '{"slug":"example","url":"https://web.archive.org/web/<timestamp>/<original-url>"}' --wait ``` -- `import-archive-org` takes one item's files (`files` exact, or `match` a regex), one at a time on archive.org's - queue; files already held are skipped. archive.org is fetched over BitTorrent when it can be, never by yt-dlp. +- `import-archive-org` takes one item's files (`files` exact, or `match` a regex), many items (`items`), or the + items an archive.org search finds (`query`, the first `limit`), one file at a time on archive.org's queue with a + pause between files and between items; records already held (on disk, or in the saved-video store) are skipped. + A dry run lists each file as held, RESTRICTED (archive.org will answer 401/403) or would-get. archive.org is + fetched over BitTorrent when it can be, never by yt-dlp. +- `get remote-listing <slug>` lists an Odysee or BitChute channel upstream (one paced read on the platform's queue) + and prints what it lists that the channel does not hold, and what it holds that is no longer listed. - A one-off Odysee or BitChute import is paced like that platform's own downloads. A whole Odysee or BitChute channel is a channel with that URL, then `sync`. - A Wayback capture is named by what it copies; `wayback.json` beside it records the capture. diff --git a/common/controller/archiveOrgImport.test.ts b/common/controller/archiveOrgImport.test.ts @@ -24,10 +24,16 @@ writeFileSync( mkdirSync(path.join(ROOT, "channels", "c", "data"), { recursive: true }); const { + ARCHIVE_ORG_ITEM_GAP_SECONDS, ARCHIVE_ORG_MIN_GAP_SECONDS, archiveOrgGapMs, + archiveOrgRestriction, + archiveOrgSearchUrl, resolveArchiveOrgImportUrl, + runArchiveOrgBatchImport, runArchiveOrgImport, + searchArchiveOrgItems, + summarizeArchiveOrgBatch, } = await import("./archiveOrgImport"); const { ArchiveOrgClient } = await import("../lib/archiveOrgClient"); const { getPaths } = await import("../lib/paths"); @@ -221,3 +227,184 @@ test("bulk: a dry run fetches nothing", async () => { assert.equal(n, 0); assert.equal(result.planned, 3); }); + +// ─── Many items (release 19 A6) ─── + +// A scripted archive.org that answers by URL: each item's metadata, a search, +// and `{}` (no such item) for anything else. Every request is recorded. +function routedClient(items: Record<string, unknown>, search?: unknown) { + const asked: string[] = []; + const c = new ArchiveOrgClient( + { + sleep: async () => {}, + fetch: async (url: string) => { + asked.push(url); + if (url.includes("advancedsearch.php")) return new Response(JSON.stringify(search ?? {}), { status: 200 }); + const id = /\/metadata\/([^/?]+)/.exec(url)?.[1] ?? ""; + return new Response(JSON.stringify(items[decodeURIComponent(id)] ?? {}), { status: 200 }); + }, + }, + { minGapMs: 0 }, + ); + return { client: c, asked }; +} + +const one = (identifier: string, extra: Record<string, unknown> = {}, file: Record<string, unknown> = {}) => ({ + metadata: { identifier, title: identifier, ...extra }, + files: [{ name: `${identifier}.mp4`, source: "original", ...file }], +}); + +test("restriction: an access-restricted item restricts every file; else a private file only", () => { + const whole = archiveOrgRestriction(one("r1", { "access-restricted-item": "true" }) as never); + assert.equal(whole.item, true); + assert.deepEqual([...whole.files], ["r1.mp4"]); + const priv = archiveOrgRestriction({ + metadata: { identifier: "r2" }, + files: [ + { name: "a.mp4", source: "original", private: "true" }, + { name: "b.mp4", source: "original" }, + ], + } as never); + assert.equal(priv.item, false); + assert.deepEqual([...priv.files], ["a.mp4"]); +}); + +test("search: one request through the client, identifiers in order; the URL sorts and caps", async () => { + const url = archiveOrgSearchUrl("collection:example AND mediatype:movies", 50); + assert.match(url, /^https:\/\/archive\.org\/advancedsearch\.php\?/); + const q = new URL(url).searchParams; + assert.equal(q.get("q"), "collection:example AND mediatype:movies"); + assert.deepEqual(q.getAll("fl[]"), ["identifier", "title"]); + assert.equal(q.get("sort[]"), "identifier asc"); + assert.equal(q.get("rows"), "50"); + assert.equal(q.get("output"), "json"); + const { client: c, asked } = routedClient( + {}, + { response: { numFound: 7, docs: [{ identifier: "a-1", title: ["A"] }, { identifier: "b-2" }, { title: "no id" }] } }, + ); + const r = await searchArchiveOrgItems("x", { rows: 9999, client: c }); + assert.equal(asked.length, 1); + assert.equal(new URL(asked[0]).searchParams.get("rows"), "500"); + assert.equal(r.found, 7); + assert.deepEqual(r.items, [{ identifier: "a-1", title: "A" }, { identifier: "b-2" }]); +}); + +test("batch dry run: held (on disk or saved), restricted and missing items are told apart; nothing fetched", async () => { + const dataDir = path.join(ROOT, "channels", "c", "data"); + // h-disk holds a transcript; h-saved has only its saved-video pointer (the + // saved-container tier) — both are held. + mkdirSync(path.join(dataDir, "h-disk"), { recursive: true }); + writeFileSync(path.join(dataDir, "h-disk", "transcript.json"), "{}"); + mkdirSync(path.join(dataDir, "h-saved"), { recursive: true }); + writeFileSync( + path.join(dataDir, "h-saved", "saved-video.json"), + JSON.stringify({ storedAt: "2026-01-01T00:00:00Z", dir: "/store/c/h-saved", file: "source.mp4", bytes: 10 }), + ); + const { client: c } = routedClient({ + "h-disk": one("h-disk"), + "h-saved": one("h-saved"), + "r-item": one("r-item", { "access-restricted-item": true }), + "r-file": one("r-file", {}, { private: "true" }), + fresh: one("fresh"), + }); + const sleeps: number[] = []; + const lines: string[] = []; + let n = 0; + const b = await runArchiveOrgBatchImport({ + paths: getPaths(), + slug: "c", + channelConfig: CONFIG, + items: ["h-disk", "h-saved", "r-item", "r-file", "gone", "fresh"].map((identifier) => ({ + identifier, + selection: { match: "." }, + })), + onLog: (l) => lines.push(l), + signal: new AbortController().signal, + dryRun: true, + deps: { + client: c, + downloadOne: async () => { + n++; + return outcome("ok"); + }, + sleep: async (ms) => { + sleeps.push(ms); + }, + random: () => 0, + }, + }); + assert.equal(n, 0); + // The inter-item gap before every item after the first, and no download gap. + assert.deepEqual(sleeps, Array(5).fill(ARCHIVE_ORG_ITEM_GAP_SECONDS * 1000)); + const s = summarizeArchiveOrgBatch(b); + assert.equal(s.dryRun, true); + assert.equal(s.items, 5); + assert.equal(s.held, 2); + assert.deepEqual(s.restricted, [ + { identifier: "r-item", files: ["r-item.mp4"] }, + { identifier: "r-file", files: ["r-file.mp4"] }, + ]); + assert.deepEqual(s.missing, ["gone"]); + assert.match(lines.join(""), /RESTRICTED r-item/); + assert.match(lines.join(""), /would get fresh/); +}); + +test("batch import: a 401 skips the rest of its item; three refused items in a row stop the batch", async () => { + const two = (id: string) => ({ + metadata: { identifier: id }, + files: [`${id}-a.mp4`, `${id}-b.mp4`].map((name) => ({ name, source: "original" })), + }); + const { client: c } = routedClient({ p1: two("p1"), p2: two("p2"), p3: two("p3"), p4: two("p4") }); + const urls: string[] = []; + const b = await runArchiveOrgBatchImport({ + paths: getPaths(), + slug: "c", + channelConfig: CONFIG, + items: ["p1", "p2", "p3", "p4"].map((identifier) => ({ identifier, selection: { match: "." } })), + onLog: () => {}, + signal: new AbortController().signal, + deps: { + client: c, + downloadOne: async (o) => { + urls.push(o.videoUrl); + return { + ...outcome("failed"), + attempts: [{ n: 1, error: `archive.org answered HTTP 401 for ${o.videoUrl}` }], + } as unknown as DownloadOutcomeRecord; + }, + sleep: async () => {}, + }, + }); + // One attempt per item (its second file skipped as restricted), three items, then the stop. + assert.equal(urls.length, 3); + assert.equal(b.refusedStorm, true); + assert.match(b.stopped ?? "", /refused 3 items in a row/); + assert.deepEqual(b.notReached, ["p4"]); + assert.deepEqual(b.items[0].restricted, ["p1-a.mp4", "p1-b.mp4"]); + assert.equal(b.items[0].failed.length, 0); +}); + +test("batch import: a rate limit stops everything; the rest are not reached", async () => { + const { client: c } = routedClient({ q1: one("q1"), q2: one("q2"), q3: one("q3") }); + let n = 0; + const b = await runArchiveOrgBatchImport({ + paths: getPaths(), + slug: "c", + channelConfig: CONFIG, + items: ["q1", "q2", "q3"].map((identifier) => ({ identifier, selection: { match: "." } })), + onLog: () => {}, + signal: new AbortController().signal, + deps: { + client: c, + downloadOne: async () => { + n++; + return n === 1 ? outcome("ok") : outcome("failed", "rate_limit"); + }, + sleep: async () => {}, + }, + }); + assert.equal(n, 2); + assert.equal(b.rateLimited, true); + assert.deepEqual(b.notReached, ["q3"]); + assert.deepEqual(summarizeArchiveOrgBatch(b).imported, 1); +}); diff --git a/common/controller/archiveOrgImport.ts b/common/controller/archiveOrgImport.ts @@ -27,6 +27,15 @@ // backoff takes over), and ARCHIVE_ORG_MAX_CONSECUTIVE_FAILURES failures in // a row stop it too. Re-running the same command resumes: what landed is // skipped. +// +// MANY ITEMS (release 19 A6): `{items: [...]}` or `{query}` (archive.org's +// advanced search, one request) is one job over many items, the same polite +// rules shared across the batch plus a jittered pause between items. A file +// archive.org marks as not for download (`access-restricted-item` on the item, +// `private` on the file) is RESTRICTED: a dry run lists it, an import skips it; +// a download answered 401/403 marks the rest of its item restricted, and three +// such items in a row stop the batch (the platform backs off). A record already +// HELD — on disk, or its media in the saved-video store — is never fetched. import path from "node:path"; import type { ChannelConfig } from "../lib/channelConfig"; @@ -43,7 +52,12 @@ import { archiveOrgVideoId, parseArchiveOrgUrl, } from "../lib/archiveOrgId"; -import { archiveOrgClient, type ArchiveOrgClient } from "../lib/archiveOrgClient"; +import { + ArchiveOrgRequestError, + archiveOrgClient, + type ArchiveOrgClient, +} from "../lib/archiveOrgClient"; +import { isSavedVideo } from "../lib/savedVideo-server"; import { getSettings } from "../lib/settings"; import { diskGate } from "../lib/diskSpace"; import { resolveCookiePolicy } from "../lib/cookiePolicy"; @@ -123,8 +137,12 @@ export type ArchiveOrgImportPlanEntry = { file: string; url: string; id: string; - // Already downloaded: skipped, never fetched again. + // HELD: already downloaded, or its media kept in the saved-video store (the + // saved-container tier) — skipped, never fetched again. onDisk: boolean; + // archive.org marks the file (or the whole item) as not for download: a + // fetch answers 401/403. Listed by a dry run, skipped by an import. + restricted: boolean; }; export type ArchiveOrgImportPlan = { @@ -134,8 +152,41 @@ export type ArchiveOrgImportPlan = { entries: ArchiveOrgImportPlanEntry[]; // Named in `files` but not a media original of the item. unknown: string[]; + // The whole item is access-restricted (archive.org's lending flag). + restrictedItem: boolean; }; +const truthy = (v: unknown): boolean => v === true || v === "true" || v === "1"; + +// What archive.org will not hand out, read off the item's own metadata: an +// item carrying `access-restricted-item` (every file), else each file marked +// `private`. Either answers a download with 401/403; the metadata says so +// before anything is asked for. +export function archiveOrgRestriction(item: ArchiveOrgItemMetadata): { + item: boolean; + files: Set<string>; +} { + const meta = item.metadata as Record<string, unknown>; + const whole = truthy(meta["access-restricted-item"]); + const files = new Set<string>(); + for (const f of item.files) { + if (whole || truthy((f as Record<string, unknown>).private)) files.add(f.name); + } + return { item: whole, files }; +} + +// HELD: the record has its transcript or (on a transcribe channel) its audio — +// `destinationExists`, what every download path asks — or its media lives in +// the saved-video store, which a download must never bypass. +export async function isHeldArchiveOrgRecord( + dataDir: string, + id: string, + handling: ChannelConfig["handling"], +): Promise<boolean> { + if (await destinationExists(dataDir, id, handling)) return true; + return isSavedVideo(path.join(dataDir, id)); +} + export async function planArchiveOrgImport(opts: { identifier: string; selection: ArchiveOrgFileSelection; @@ -148,6 +199,7 @@ export async function planArchiveOrgImport(opts: { const item = await client.itemMetadata(opts.identifier, opts.signal); const identifier = item.metadata.identifier || opts.identifier; const media = listArchiveOrgMediaFiles(item); + const restriction = archiveOrgRestriction(item); const { picked, unknown } = pickArchiveOrgFiles(item, opts.selection); const entries: ArchiveOrgImportPlanEntry[] = []; for (const file of picked) { @@ -159,7 +211,8 @@ export async function planArchiveOrgImport(opts: { file, url, id, - onDisk: await destinationExists(opts.dataDir, id, opts.handling), + onDisk: await isHeldArchiveOrgRecord(opts.dataDir, id, opts.handling), + restricted: restriction.files.has(file), }); } const title = item.metadata.title; @@ -169,9 +222,51 @@ export async function planArchiveOrgImport(opts: { mediaFiles: media.length, entries, unknown, + restrictedItem: restriction.item, }; } +// ─── A search ─── + +export const ARCHIVE_ORG_SEARCH_DEFAULT_ROWS = 100; +export const ARCHIVE_ORG_SEARCH_MAX_ROWS = 500; + +// archive.org's advanced search, ONE request through the polite client (its +// chain, its gap, its backoff): the identifiers of the first `rows` items the +// query matches, in identifier order so a re-run walks the same list. +export function archiveOrgSearchUrl(query: string, rows: number): string { + const q = new URLSearchParams(); + q.set("q", query); + q.append("fl[]", "identifier"); + q.append("fl[]", "title"); + q.append("sort[]", "identifier asc"); + q.set("rows", String(rows)); + q.set("page", "1"); + q.set("output", "json"); + return `https://archive.org/advancedsearch.php?${q.toString()}`; +} + +export async function searchArchiveOrgItems( + query: string, + opts: { rows?: number; client?: ArchiveOrgClient; signal?: AbortSignal } = {}, +): Promise<{ found: number; items: { identifier: string; title?: string }[] }> { + const rows = Math.min(ARCHIVE_ORG_SEARCH_MAX_ROWS, Math.max(1, opts.rows ?? ARCHIVE_ORG_SEARCH_DEFAULT_ROWS)); + const client = opts.client ?? archiveOrgClient; + const raw = (await client.getJson(archiveOrgSearchUrl(query, rows), opts.signal)) as { + response?: { numFound?: unknown; docs?: unknown }; + } | null; + const docs = Array.isArray(raw?.response?.docs) ? (raw!.response!.docs as unknown[]) : []; + const items: { identifier: string; title?: string }[] = []; + for (const d of docs) { + const r = d as { identifier?: unknown; title?: unknown }; + if (typeof r?.identifier !== "string" || !r.identifier) continue; + const title = Array.isArray(r.title) ? r.title.find((t) => typeof t === "string") : r.title; + items.push({ identifier: r.identifier, ...(typeof title === "string" ? { title } : {}) }); + } + const found = typeof raw?.response?.numFound === "number" ? raw.response.numFound : items.length; + return { found, items }; +} + // The gap before the next file: the configured pause, floored, plus up to // half again at random so a batch never settles into a fixed beat. export function archiveOrgGapMs(sleepBetweenDownloadsSeconds: number, random: number): number { @@ -179,15 +274,39 @@ export function archiveOrgGapMs(sleepBetweenDownloadsSeconds: number, random: nu return Math.round(base * (1 + 0.5 * Math.min(1, Math.max(0, random))) * 1000); } +// THE INTER-ITEM GAP: before the next item of a batch is asked about, a +// jittered pause on top of the client's own 2 s between requests, so a dry run +// over a search of hundreds of items reads as a person paging, not a crawl. +// (A download after the first is paced by archiveOrgGapMs as well.) +export const ARCHIVE_ORG_ITEM_GAP_SECONDS = 3; +export function archiveOrgItemGapMs(random: number): number { + return Math.round(ARCHIVE_ORG_ITEM_GAP_SECONDS * (1 + 0.5 * Math.min(1, Math.max(0, random))) * 1000); +} + +// Items in a row whose downloads archive.org refused with 401/403 before the +// batch stops: one restricted item is that item; three in a row is archive.org +// refusing us, and the platform backs off. +export const ARCHIVE_ORG_MAX_REFUSED_ITEMS = 3; + function isOk(rec: DownloadOutcomeRecord): boolean { return rec.status.startsWith("ok"); } +// A download error that is archive.org saying "not for you" (lib/archiveOrgClient +// words every non-retryable answer "archive.org answered HTTP <n> for <url>"). +export function isRefusalError(error: string): boolean { + return /\bHTTP 40[13]\b/.test(error); +} + export type ArchiveOrgImportResult = { identifier: string; planned: number; imported: string[]; + // Held: on disk or in the saved-video store. skipped: string[]; + // Not for download (archive.org's metadata says so, or a fetch answered + // 401/403): never fetched. + restricted: string[]; failed: { file: string; error: string }[]; unknown: string[]; // Why the batch ended before its last file, when it did. @@ -219,18 +338,47 @@ function abortableSleep(ms: number, signal: AbortSignal): Promise<void> { }); } -export async function runArchiveOrgImport(opts: { +export type ArchiveOrgItemRequest = { + identifier: string; + selection: ArchiveOrgFileSelection; +}; + +export type ArchiveOrgBatchResult = { + // One per item reached, in order. + items: ArchiveOrgImportResult[]; + // Items archive.org has no record of (or whose metadata it would not give). + missing: { identifier: string; error: string }[]; + // Items never reached: the batch stopped first. + notReached: string[]; + stopped?: string; + rateLimited?: boolean; + // Three items in a row refused with 401/403: the caller backs the platform off. + refusedStorm?: boolean; + dryRun: boolean; +}; + +type BatchOpts = { paths: Paths; slug: string; channelConfig: ChannelConfig; - identifier: string; - selection: ArchiveOrgFileSelection; onLog: (line: string) => void; signal: AbortSignal; drainSignal?: AbortSignal; dryRun?: boolean; deps?: ArchiveOrgImportDeps; -}): Promise<ArchiveOrgImportResult> { +}; + +// MANY ITEMS, ONE JOB (`import-archive-org {items | query}`). Each item is +// planned (one cached metadata request) and imported as a single-item import +// is, with the state that makes it polite shared across the batch: the gap +// before every download after the first, the consecutive-failure count, the +// rate-limit stop. Between items, the inter-item gap. An item archive.org has +// no record of is listed and passed over; a rate limit stops everything (the +// rest are `notReached`, and re-running the same body resumes — held records +// are skipped); three items in a row refused with 401/403 stop it too. +export async function runArchiveOrgBatchImport( + opts: BatchOpts & { items: ArchiveOrgItemRequest[] }, +): Promise<ArchiveOrgBatchResult> { const deps = opts.deps ?? {}; const downloadOne = deps.downloadOne ?? downloadOneManaged; const sleep = deps.sleep ?? abortableSleep; @@ -238,124 +386,269 @@ export async function runArchiveOrgImport(opts: { const settings = getSettings(); const dataDir = path.join(opts.paths.channelsDir, opts.slug, "data"); const log = opts.onLog; - - const plan = await planArchiveOrgImport({ - identifier: opts.identifier, - selection: opts.selection, - dataDir, - handling: opts.channelConfig.handling, - client: deps.client, - signal: opts.signal, - }); - const result: ArchiveOrgImportResult = { - identifier: plan.identifier, - planned: plan.entries.length, - imported: [], - skipped: [], - failed: [], - unknown: plan.unknown, - }; - log( - `archive.org item ${plan.identifier}${plan.title ? ` ("${plan.title}")` : ""}: ` + - `${plan.mediaFiles} media files, ${plan.entries.length} chosen, ` + - `${plan.entries.filter((e) => e.onDisk).length} already downloaded.\n`, - ); - if (plan.unknown.length > 0) { - log(`Not media files of the item (ignored): ${plan.unknown.map((n) => JSON.stringify(n)).join(", ")}\n`); - } - if (opts.dryRun) { - for (const e of plan.entries) log(` ${e.onDisk ? "on disk " : "would get"} ${e.id} ${e.file}\n`); - result.skipped = plan.entries.filter((e) => e.onDisk).map((e) => e.file); - return result; - } - const sleepSeconds = opts.channelConfig.sleepBetweenDownloadsSeconds ?? settings.sleepBetweenDownloadsSeconds; + const batch: ArchiveOrgBatchResult = { items: [], missing: [], notReached: [], dryRun: opts.dryRun === true }; let fetched = 0; let consecutiveFailures = 0; - for (const entry of plan.entries) { + let refusedItemsInARow = 0; + const stop = (why: string) => { + batch.stopped = why; + }; + + for (let i = 0; i < opts.items.length; i++) { + const req = opts.items[i]; + if (batch.stopped) { + batch.notReached.push(req.identifier); + continue; + } if (opts.signal.aborted) { - result.stopped = "cancelled"; - break; + stop("cancelled"); + batch.notReached.push(req.identifier); + continue; } if (opts.drainSignal?.aborted) { - result.stopped = "drained"; - break; - } - if (entry.onDisk || (await destinationExists(dataDir, entry.id, opts.channelConfig.handling))) { - result.skipped.push(entry.file); + stop("drained"); + batch.notReached.push(req.identifier); continue; } - const gate = await diskGate(opts.paths, settings, { dir: dataDir }); - if (!gate.ok) { - result.stopped = gate.message; - log(`Stopping: ${gate.message}.\n`); - break; - } - if (fetched > 0) { - const gap = archiveOrgGapMs(sleepSeconds, random()); - log(`Waiting ${(gap / 1000).toFixed(1)}s before the next file (archive.org pacing)...\n`); - await sleep(gap, opts.signal); + if (i > 0) { + await sleep(archiveOrgItemGapMs(random()), opts.signal); if (opts.signal.aborted) { - result.stopped = "cancelled"; - break; + stop("cancelled"); + batch.notReached.push(req.identifier); + continue; } } - fetched++; - log(`[${fetched}] ${entry.file} → data/${entry.id}/\n`); - let rec: DownloadOutcomeRecord | null = null; - let error = ""; + + let plan: ArchiveOrgImportPlan; try { - rec = await downloadOne({ - channelSlug: opts.slug, - channelConfig: opts.channelConfig, - paths: opts.paths, - videoUrl: entry.url, - onLog: log, + plan = await planArchiveOrgImport({ + identifier: req.identifier, + selection: req.selection, + dataDir, + handling: opts.channelConfig.handling, + client: deps.client, signal: opts.signal, - cookiePolicy: resolveCookiePolicy(settings, opts.channelConfig), - inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, - globalSkipLiveDownloads: settings.skipLiveDownloads, - appendArchive: true, }); } catch (err) { - error = (err as Error).message; + const e = err as ArchiveOrgRequestError; + if (e instanceof ArchiveOrgRequestError && e.rateLimited) { + batch.rateLimited = true; + stop(`archive.org did not answer for item ${req.identifier}; stopping (re-run later — held records are skipped)`); + log(`${batch.stopped}.\n`); + batch.notReached.push(req.identifier); + continue; + } + if (opts.signal.aborted) { + stop("cancelled"); + batch.notReached.push(req.identifier); + continue; + } + batch.missing.push({ identifier: req.identifier, error: e.message }); + log(`archive.org item ${req.identifier}: ${e.message} — passed over.\n`); + continue; } - if (rec && isOk(rec)) { - consecutiveFailures = 0; - result.imported.push(entry.file); - await mergeRosterFile( - opts.paths, - opts.slug, - [{ id: entry.id, url: entry.url }], - new Date().toISOString(), - "import", - ).catch(() => { - /* the download succeeded; a roster write failure must not fail it */ - }); - deps.onImported?.(entry.id); + + const result: ArchiveOrgImportResult = { + identifier: plan.identifier, + planned: plan.entries.length, + imported: [], + skipped: [], + restricted: [], + failed: [], + unknown: plan.unknown, + }; + batch.items.push(result); + const held = plan.entries.filter((e) => e.onDisk).length; + const restricted = plan.entries.filter((e) => e.restricted && !e.onDisk).length; + log( + `archive.org item ${plan.identifier}${plan.title ? ` ("${plan.title}")` : ""}: ` + + `${plan.mediaFiles} media files, ${plan.entries.length} chosen, ` + + `${held} already held` + + `${restricted ? `, ${restricted} restricted${plan.restrictedItem ? " (the item is access-restricted)" : ""}` : ""}.\n`, + ); + if (plan.unknown.length > 0) { + log(`Not media files of the item (ignored): ${plan.unknown.map((n) => JSON.stringify(n)).join(", ")}\n`); + } + if (opts.dryRun) { + for (const e of plan.entries) { + const state = e.onDisk ? "held " : e.restricted ? "RESTRICTED" : "would get"; + log(` ${state} ${e.id} ${e.file}\n`); + } + result.skipped = plan.entries.filter((e) => e.onDisk).map((e) => e.file); + result.restricted = plan.entries.filter((e) => !e.onDisk && e.restricted).map((e) => e.file); continue; } - consecutiveFailures++; - const why = error || rec?.attempts.at(-1)?.error || rec?.status || "failed"; - result.failed.push({ file: entry.file, error: why }); - log(` failed: ${why}\n`); - if (rec?.failureClass === "rate_limit") { - result.rateLimited = true; - result.stopped = "archive.org rate-limited the download; stopping (re-run later — files on disk are skipped)"; - log(`${result.stopped}.\n`); - break; + + let refusedHere = false; + for (const entry of plan.entries) { + if (opts.signal.aborted) { + stop("cancelled"); + break; + } + if (opts.drainSignal?.aborted) { + stop("drained"); + break; + } + if (entry.onDisk || (await isHeldArchiveOrgRecord(dataDir, entry.id, opts.channelConfig.handling))) { + result.skipped.push(entry.file); + continue; + } + if (entry.restricted || refusedHere) { + result.restricted.push(entry.file); + continue; + } + const gate = await diskGate(opts.paths, settings, { dir: dataDir }); + if (!gate.ok) { + stop(gate.message); + log(`Stopping: ${gate.message}.\n`); + break; + } + if (fetched > 0) { + const gap = archiveOrgGapMs(sleepSeconds, random()); + log(`Waiting ${(gap / 1000).toFixed(1)}s before the next file (archive.org pacing)...\n`); + await sleep(gap, opts.signal); + if (opts.signal.aborted) { + stop("cancelled"); + break; + } + } + fetched++; + log(`[${fetched}] ${entry.file} → data/${entry.id}/\n`); + let rec: DownloadOutcomeRecord | null = null; + let error = ""; + try { + rec = await downloadOne({ + channelSlug: opts.slug, + channelConfig: opts.channelConfig, + paths: opts.paths, + videoUrl: entry.url, + onLog: log, + signal: opts.signal, + cookiePolicy: resolveCookiePolicy(settings, opts.channelConfig), + inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, + globalSkipLiveDownloads: settings.skipLiveDownloads, + appendArchive: true, + }); + } catch (err) { + error = (err as Error).message; + } + if (rec && isOk(rec)) { + consecutiveFailures = 0; + refusedItemsInARow = 0; + result.imported.push(entry.file); + await mergeRosterFile( + opts.paths, + opts.slug, + [{ id: entry.id, url: entry.url }], + new Date().toISOString(), + "import", + ).catch(() => { + /* the download succeeded; a roster write failure must not fail it */ + }); + deps.onImported?.(entry.id); + continue; + } + const why = error || rec?.attempts.at(-1)?.error || rec?.status || "failed"; + if (isRefusalError(why)) { + // archive.org will not hand this item's files to us: the rest of the + // item is skipped as restricted, and it is not a failure streak. + refusedHere = true; + result.restricted.push(entry.file); + log(` refused (${why}) — the rest of ${plan.identifier} is skipped as restricted\n`); + continue; + } + consecutiveFailures++; + result.failed.push({ file: entry.file, error: why }); + log(` failed: ${why}\n`); + if (rec?.failureClass === "rate_limit") { + result.rateLimited = true; + batch.rateLimited = true; + stop("archive.org rate-limited the download; stopping (re-run later — held records are skipped)"); + log(`${batch.stopped}.\n`); + break; + } + if (consecutiveFailures >= ARCHIVE_ORG_MAX_CONSECUTIVE_FAILURES) { + stop(`${consecutiveFailures} failures in a row; stopping`); + log(`${batch.stopped}.\n`); + break; + } + } + if (batch.stopped) result.stopped = batch.stopped; + if (refusedHere) { + refusedItemsInARow++; + if (refusedItemsInARow >= ARCHIVE_ORG_MAX_REFUSED_ITEMS && !batch.stopped) { + batch.refusedStorm = true; + stop(`archive.org refused ${refusedItemsInARow} items in a row (401/403); stopping`); + result.stopped = batch.stopped; + log(`${batch.stopped}.\n`); + } } - if (consecutiveFailures >= ARCHIVE_ORG_MAX_CONSECUTIVE_FAILURES) { - result.stopped = `${consecutiveFailures} failures in a row; stopping`; - log(`${result.stopped}.\n`); - break; + log( + `archive.org import of ${plan.identifier}: ${result.imported.length} imported, ` + + `${result.skipped.length} already held, ${result.restricted.length} restricted, ${result.failed.length} failed` + + `${result.stopped ? ` — stopped: ${result.stopped}` : ""}.\n`, + ); + } + return batch; +} + +// The batch's totals, for its last log line (`summary: {…}`) and the action's +// verdict. +export function summarizeArchiveOrgBatch(b: ArchiveOrgBatchResult): { + dryRun: boolean; + items: number; + planned: number; + imported: number; + held: number; + restricted: { identifier: string; files: string[] }[]; + failed: number; + missing: string[]; + notReached: string[]; + stopped?: string; +} { + const sum = (f: (r: ArchiveOrgImportResult) => number) => b.items.reduce((n, r) => n + f(r), 0); + return { + dryRun: b.dryRun, + items: b.items.length, + planned: sum((r) => r.planned), + imported: sum((r) => r.imported.length), + held: sum((r) => r.skipped.length), + restricted: b.items + .filter((r) => r.restricted.length > 0) + .map((r) => ({ identifier: r.identifier, files: r.restricted })), + failed: sum((r) => r.failed.length), + missing: b.missing.map((m) => m.identifier), + notReached: b.notReached, + ...(b.stopped ? { stopped: b.stopped } : {}), + }; +} + +// ONE ITEM — the original `import-archive-org {item}` — as a batch of one. +export async function runArchiveOrgImport( + opts: BatchOpts & { identifier: string; selection: ArchiveOrgFileSelection }, +): Promise<ArchiveOrgImportResult> { + const { identifier, selection, ...rest } = opts; + const batch = await runArchiveOrgBatchImport({ ...rest, items: [{ identifier, selection }] }); + // A single item's metadata failure is the import's error, as it always was. + if (batch.missing.length > 0) throw new Error(batch.missing[0].error); + const result = batch.items[0]; + if (!result) { + if (batch.stopped && batch.stopped !== "cancelled" && batch.stopped !== "drained") { + throw new Error(batch.stopped); } + return { + identifier, + planned: 0, + imported: [], + skipped: [], + restricted: [], + failed: [], + unknown: [], + ...(batch.stopped ? { stopped: batch.stopped } : {}), + ...(batch.rateLimited ? { rateLimited: true } : {}), + }; } - log( - `archive.org import of ${plan.identifier}: ${result.imported.length} imported, ` + - `${result.skipped.length} already on disk, ${result.failed.length} failed` + - `${result.stopped ? ` — stopped: ${result.stopped}` : ""}.\n`, - ); - return result; + return { ...result, ...(batch.rateLimited ? { rateLimited: true } : {}) }; } diff --git a/common/controller/remoteListing.test.ts b/common/controller/remoteListing.test.ts @@ -0,0 +1,87 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdirSync, mkdtempSync, writeFileSync } from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import type { ChannelConfig } from "../lib/channelConfig"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test controller/remoteListing.test.ts +// +// The diff of a listing against the held ids, and one run with the listing +// injected. Every id here is invented. + +const ROOT = mkdtempSync(path.join(os.tmpdir(), "remote-listing-")); +process.env.TRANSCRIPTS_DIR = ROOT; +process.env.SETTINGS_FILE = path.join(ROOT, "settings.json"); +writeFileSync(process.env.SETTINGS_FILE, "{}\n"); + +const { diffRemoteListing, heldVideoIds, isRemoteListingPlatform, runRemoteListing } = await import( + "./remoteListing" +); +const { getPaths } = await import("../lib/paths"); + +test("only Odysee and BitChute are listed this way", () => { + assert.equal(isRemoteListingPlatform("odysee"), true); + assert.equal(isRemoteListingPlatform("bitchute"), true); + assert.equal(isRemoteListingPlatform("youtube"), false); + assert.equal(isRemoteListingPlatform(null), false); +}); + +test("the diff keeps listing order, collapses duplicates, and counts entries with no id", () => { + const d = diffRemoteListing( + [ + "https://www.bitchute.com/video/AbC123/", + "https://www.bitchute.com/video/held01/", + "https://www.bitchute.com/channel/somebody/", + "https://www.bitchute.com/embed/AbC123/", + "https://odysee.com/@demo:0/a-video:7", + ], + new Set(["held01", "gone99"]), + ); + assert.equal(d.listed, 5); + assert.equal(d.unparsed, 1); + assert.equal(d.held, 1); + assert.deepEqual(d.notHeld, [ + { id: "AbC123", url: "https://www.bitchute.com/video/AbC123/" }, + { id: "7", url: "https://odysee.com/@demo:0/a-video:7" }, + ]); + assert.deepEqual(d.heldNotListed, ["gone99"]); +}); + +test("held ids are data/ directories; an absent data/ holds nothing", async () => { + assert.deepEqual([...(await heldVideoIds(path.join(ROOT, "nowhere")))], []); + const dir = path.join(ROOT, "h", "data"); + mkdirSync(path.join(dir, "v1"), { recursive: true }); + mkdirSync(path.join(dir, ".hidden"), { recursive: true }); + writeFileSync(path.join(dir, "stray.txt"), ""); + assert.deepEqual([...(await heldVideoIds(dir))], ["v1"]); +}); + +test("a run lists through the injected reader and answers the diff", async () => { + const slug = "demo-bitchute"; + mkdirSync(path.join(ROOT, "channels", slug, "data", "held01"), { recursive: true }); + const config = { handling: "transcribe", platform: "bitchute", url: "https://www.bitchute.com/channel/demo/" } as ChannelConfig; + writeFileSync(path.join(ROOT, "channels", slug, "config.json"), JSON.stringify(config)); + const lines: string[] = []; + let asked = 0; + const r = await runRemoteListing({ + paths: getPaths(), + slug, + channelConfig: config, + platform: "bitchute", + onLog: (l) => lines.push(l), + signal: new AbortController().signal, + deps: { + list: async () => { + asked++; + return ["https://www.bitchute.com/video/held01/", "https://www.bitchute.com/video/new001/"]; + }, + now: () => new Date("2026-10-10T00:00:00Z"), + }, + }); + assert.equal(asked, 1); + assert.equal(r.listedAt, "2026-10-10T00:00:00.000Z"); + assert.deepEqual(r.notHeld, [{ id: "new001", url: "https://www.bitchute.com/video/new001/" }]); + assert.match(lines.join(""), /2 listed, 1 held, 1 not held, 0 held but not listed/); +}); diff --git a/common/controller/remoteListing.ts b/common/controller/remoteListing.ts @@ -0,0 +1,133 @@ +// AN ODYSEE OR BITCHUTE CHANNEL'S REMOTE LISTING, DIFFED AGAINST WHAT IS HELD +// (release 19 A6, `pnpm ops get remote-listing <slug>`). +// +// One flat-playlist enumeration of the channel's own URL — the same one-spawn +// read the quick availability check and a sync's full sweep make, with the +// platform's paced args (`configArgs`) — on the platform's own queue, so it is +// ONE STREAM with every other request to that platform. Nothing is written: no +// roster merge, no maybe-missing, no download. The answer is the diff: +// +// notHeld listed upstream, no data/<id>/ here — what a download would get +// heldNotListed held here, absent from the listing (deleted, unlisted, or a +// record imported from elsewhere) +// +// POLITE: the editor refuses it while the platform is held or cooling down, +// waits out the platform's import floor before the read (jobs/platformGap.ts), +// and a rate-limited read backs the platform off (the caller records it). +// Ids are the URL-derived canonical ids, which ARE the data-dir names on both +// platforms (lib/videoId.ts). + +import path from "node:path"; +import { readdir } from "node:fs/promises"; +import type { ChannelConfig } from "../lib/channelConfig"; +import type { Paths } from "../lib/paths"; +import { assertChannelTextReadable } from "../lib/channelMedia"; +import { extractVideoId } from "../lib/videoId"; +import { fetchFlatPlaylistUrls } from "../ytdlp/runYtdlp"; + +export const REMOTE_LISTING_PLATFORMS = ["odysee", "bitchute"] as const; +export type RemoteListingPlatform = (typeof REMOTE_LISTING_PLATFORMS)[number]; + +export function isRemoteListingPlatform(p: string | null | undefined): p is RemoteListingPlatform { + return (REMOTE_LISTING_PLATFORMS as readonly string[]).includes(p ?? ""); +} + +// The log line a remote-listing job ends with: this marker, then the result as +// compact JSON (`pnpm ops get remote-listing` prints it on stdout). The same +// string is REMOTE_LISTING_RESULT_MARKER in scripts/archilyzer-ops.mjs. +export const REMOTE_LISTING_RESULT_MARKER = "@@remote-listing "; + +export type RemoteListingResult = { + slug: string; + platform: RemoteListingPlatform; + url: string; + listedAt: string; + // Entries the enumeration printed, and how many named no video id. + listed: number; + unparsed: number; + held: number; + notHeld: { id: string; url: string }[]; + heldNotListed: string[]; +}; + +// The pure half: a listing's URLs against the held ids. Listing order is kept +// (newest first, as the platform lists); duplicates collapse. +export function diffRemoteListing( + urls: string[], + heldIds: ReadonlySet<string>, +): Pick<RemoteListingResult, "listed" | "unparsed" | "held" | "notHeld" | "heldNotListed"> { + const seen = new Set<string>(); + const notHeld: { id: string; url: string }[] = []; + let unparsed = 0; + let held = 0; + for (const url of urls) { + const id = extractVideoId(url); + if (!id) { + unparsed++; + continue; + } + if (seen.has(id)) continue; + seen.add(id); + if (heldIds.has(id)) held++; + else notHeld.push({ id, url }); + } + const heldNotListed = [...heldIds].filter((id) => !seen.has(id)).sort(); + return { listed: urls.length, unparsed, held, notHeld, heldNotListed }; +} + +// The channel's held ids: its data/<id>/ directories. An absent data/ is a +// channel with nothing held; any other read error is thrown, never read as +// "nothing held" (the text guard runs first). +export async function heldVideoIds(dataDir: string): Promise<Set<string>> { + try { + const entries = await readdir(dataDir, { withFileTypes: true }); + return new Set(entries.filter((d) => d.isDirectory() && !d.name.startsWith(".")).map((d) => d.name)); + } catch (err) { + if ((err as NodeJS.ErrnoException).code === "ENOENT") return new Set(); + throw err; + } +} + +export type RemoteListingDeps = { + list?: typeof fetchFlatPlaylistUrls; + now?: () => Date; +}; + +export async function runRemoteListing(opts: { + paths: Paths; + slug: string; + channelConfig: ChannelConfig; + platform: RemoteListingPlatform; + onLog: (line: string) => void; + signal: AbortSignal; + deps?: RemoteListingDeps; +}): Promise<RemoteListingResult> { + const list = opts.deps?.list ?? fetchFlatPlaylistUrls; + const now = opts.deps?.now ?? (() => new Date()); + const url = opts.channelConfig.url; + if (!url) throw new Error(`Channel "${opts.slug}" has no URL configured`); + // The text guard: an unmounted or moving drive is not "nothing held". + await assertChannelTextReadable(opts.paths, opts.slug, opts.channelConfig); + const urls = await list({ + channelConfig: opts.channelConfig, + paths: opts.paths, + channelSlug: opts.slug, + onLog: opts.onLog, + signal: opts.signal, + }); + const held = await heldVideoIds(path.join(opts.paths.channelsDir, opts.slug, "data")); + const diff = diffRemoteListing(urls, held); + const result: RemoteListingResult = { + slug: opts.slug, + platform: opts.platform, + url, + listedAt: now().toISOString(), + ...diff, + }; + opts.onLog( + `Remote listing of ${opts.slug} (${opts.platform}): ${diff.listed} listed, ${diff.held} held, ` + + `${diff.notHeld.length} not held, ${diff.heldNotListed.length} held but not listed` + + `${diff.unparsed ? `, ${diff.unparsed} entries with no video id` : ""}.\n`, + ); + return result; +} diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -293,6 +293,19 @@ const JOB_KINDS: Record<string, JobKindMeta> = { queueKeyStrategy: "platform", needsMedia: true, }, + // AN ODYSEE OR BITCHUTE CHANNEL'S LISTING, diffed against what is held + // (release 19 A6, controller/remoteListing.ts): one flat-playlist read on the + // platform's queue, nothing written. Reads data/ (the text tier) for the + // held ids, so the text guard covers it. + "remote-listing": { + kind: "remote-listing", + label: "Remote listing", + drainable: false, + replayable: false, + queueKeyStrategy: "platform", + needsMedia: false, + needsText: true, + }, // ONE WINDOW of a video's source media, fetched into data/<id>/clips/ for a // tool that asked for it by name (umtool's clip bench). On its platform's // CLIP queue (`clips:<platform>`, lib/queueKeys.ts clipWindowQueueKey — diff --git a/editor/app/api/ops/import-archive-org/route.test.ts b/editor/app/api/ops/import-archive-org/route.test.ts @@ -0,0 +1,52 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { setupOpsCorpus, callPost } from "../_testCorpus"; + +// Run with: +// pnpm -C editor exec tsx --test "app/api/ops/import-archive-org/route.test.ts" +// +// The three body shapes (item / items / query) and every refusal that comes +// before a job — none of them asks archive.org anything, so no network is +// touched and no job is queued. The batch itself (held, restricted, the +// inter-item gap, the 401 storm) is common/controller/archiveOrgImport.test.ts. + +const corpus = await setupOpsCorpus(null); +await corpus.writeChannelConfig("demo-archive", { handling: "transcribe", platform: "archiveorg" }); +const { POST } = await import("./route"); +test.after(() => corpus.cleanup()); + +test("refusals before any job: keys, modes, items, limit, regex, identifiers, channel", async () => { + const many = Array.from({ length: 501 }, (_, i) => `item-${i}`); + const cases: [Record<string, unknown>, RegExp][] = [ + [{ slug: "demo-archive", items: ["a"], force: true }, /unknown key\(s\): force/], + [{ slug: "demo-archive" }, /exactly one of "item", "items" \(a list\) or "query"/], + [{ slug: "demo-archive", item: "a", query: "x" }, /exactly one of "item", "items"/], + [{ slug: "demo-archive", items: ["a"], query: "x" }, /exactly one of "item", "items"/], + [{ slug: "demo-archive", item: "a" }, /exactly one of "files" \(a list\) or "match"/], + [{ slug: "demo-archive", item: "a", files: ["x.mp4"], match: "." }, /exactly one of "files"/], + [{ slug: "demo-archive", item: "../a", match: "." }, /not an archive.org identifier/], + [{ slug: "demo-archive", items: [] }, /"items" must be a non-empty array/], + [{ slug: "demo-archive", items: many }, /holds 501 items; at most 500/], + [{ slug: "demo-archive", items: [{ item: "a", size: 1 }] }, /"items\[0\]" has unknown key\(s\): size/], + [{ slug: "demo-archive", items: [{ files: ["x.mp4"] }] }, /"items\[0\].item" is required/], + [{ slug: "demo-archive", items: [{ item: "a", files: [] }] }, /"items\[0\].files" must be a non-empty array/], + [{ slug: "demo-archive", items: [{ item: "a", files: ["x"], match: "." }] }, /item "a": "files" or "match", not both/], + [{ slug: "demo-archive", items: ["ok-1", "bad id"] }, /"bad id" is not an archive.org identifier/], + [{ slug: "demo-archive", items: [{ item: "a", match: "(" }] }, /item "a": "match" is not a valid regex/], + [{ slug: "demo-archive", items: ["a"], files: ["x.mp4"] }, /with "items", give each its own/], + [{ slug: "demo-archive", items: ["a"], limit: 5 }, /"limit" caps a "query"/], + [{ slug: "demo-archive", query: "x", limit: 501 }, /"limit" is at most 500/], + [{ slug: "demo-archive", query: "x", limit: 0 }, /"limit" must be a whole number above zero/], + [{ slug: "demo-archive", query: " " }, /"query" must be a non-empty string/], + [{ slug: "demo-archive", query: "x", match: "(" }, /"match" is not a valid regex/], + [{ slug: "demo-archive", query: "x", dryRun: "yes" }, /"dryRun" must be a boolean/], + [{ slug: "no-such", items: ["a"] }, /Channel "no-such" not found/], + ]; + for (const [body, re] of cases) { + const r = await callPost(POST, body); + assert.equal(r.status, 400, JSON.stringify(body)); + assert.match(r.body.error ?? "", re, JSON.stringify(body)); + assert.equal(r.body.jobId, undefined); + } + assert.deepEqual(await corpus.listJobIds(), []); +}); diff --git a/editor/app/api/ops/import-archive-org/route.ts b/editor/app/api/ops/import-archive-org/route.ts @@ -1,44 +1,102 @@ -import { importArchiveOrgAction } from "../../../channels/[slug]/pipelineActions"; +import { + importArchiveOrgAction, + type ArchiveOrgImportEntry, +} from "../../../channels/[slug]/pipelineActions"; import { OpsInputError, jobResponse, ops, optBool, + optPositiveInt, optString, reqSlug, - reqString, type OpsBody, } from "../_lib"; export const dynamic = "force-dynamic"; // POST { slug, item, files?: string[], match?: string, dryRun? } +// | { slug, items: (string | { item, files? | match? })[], match?, dryRun? } +// | { slug, query: string, limit?: number, match?, dryRun? } // -> { ok: true, jobId } // -// Import chosen media files of ONE archive.org item into an existing channel, -// as one job on archive.org's own queue: one file at a time, a jittered pause -// between them, files already downloaded skipped, a rate limit or three -// failures in a row ending it (controller/archiveOrgImport.ts). Exactly one of -// `files` (exact paths in the item, untrimmed) and `match` (a case-insensitive -// regex over them). `dryRun` logs what would be fetched and fetches nothing. -function optFiles(body: OpsBody): string[] | undefined { - const v = body.files; +// Import media files from archive.org into an existing channel, as ONE job on +// archive.org's own queue: one file at a time, a jittered pause between them +// and between items, records already held (on disk or in the saved-video +// store) skipped, a rate limit or three failures in a row ending it, three +// items in a row refused with 401/403 ending it too (controller/ +// archiveOrgImport.ts). `item` is one item and needs exactly one of `files` +// (exact paths in the item, untrimmed) and `match` (a case-insensitive regex +// over them); `items` is many (a bare identifier takes `match`, else every +// media original); `query` is an archive.org search, its first `limit` (100, +// at most 500) items. `dryRun` logs each file as held, RESTRICTED (archive.org +// marks it not for download — a fetch would answer 401/403) or would-get, and +// fetches nothing. The log ends with `summary: {…}`. +const MAX_ARCHIVE_ORG_ITEMS = 500; + +function optFiles(v: unknown, key: string): string[] | undefined { if (v === undefined) return undefined; if (!Array.isArray(v) || v.length === 0 || v.some((s) => typeof s !== "string" || !s)) { - throw new OpsInputError('"files" must be a non-empty array of file names'); + throw new OpsInputError(`"${key}" must be a non-empty array of file names`); } return v as string[]; } +function optItems(body: OpsBody): ArchiveOrgImportEntry[] | undefined { + const v = body.items; + if (v === undefined) return undefined; + if (!Array.isArray(v) || v.length === 0) { + throw new OpsInputError('"items" must be a non-empty array of identifiers or { "item", "files"? | "match"? }'); + } + if (v.length > MAX_ARCHIVE_ORG_ITEMS) { + throw new OpsInputError(`"items" holds ${v.length} items; at most ${MAX_ARCHIVE_ORG_ITEMS} per job`); + } + return v.map((entry, i) => { + if (typeof entry === "string" && entry.trim()) return entry.trim(); + if (typeof entry !== "object" || entry === null || Array.isArray(entry)) { + throw new OpsInputError(`"items[${i}]" must be an identifier or { "item", "files"? | "match"? }`); + } + const e = entry as Record<string, unknown>; + const stray = Object.keys(e).filter((k) => !["item", "files", "match"].includes(k)); + if (stray.length) { + throw new OpsInputError(`"items[${i}]" has unknown key(s): ${stray.join(", ")} — accepted: item, files, match`); + } + if (typeof e.item !== "string" || !e.item.trim()) { + throw new OpsInputError(`"items[${i}].item" is required and must be a non-empty string`); + } + if (e.match !== undefined && typeof e.match !== "string") { + throw new OpsInputError(`"items[${i}].match" must be a string`); + } + return { + item: e.item.trim(), + ...(e.files !== undefined ? { files: optFiles(e.files, `items[${i}].files`) } : {}), + ...(e.match !== undefined ? { match: e.match as string } : {}), + }; + }); +} + export async function POST(request: Request) { - return ops(request, ["slug", "item", "files", "match", "dryRun"], async (body) => - jobResponse( - await importArchiveOrgAction(reqSlug(body, "slug"), { - item: reqString(body, "item"), - files: optFiles(body), - match: optString(body, "match"), - dryRun: optBool(body, "dryRun"), - }), - ), + return ops( + request, + ["slug", "item", "items", "query", "limit", "files", "match", "dryRun"], + async (body) => { + const query = optString(body, "query"); + if (query !== undefined && !query.trim()) throw new OpsInputError('"query" must be a non-empty string'); + const limit = optPositiveInt(body, "limit"); + if (limit !== undefined && limit > 500) throw new OpsInputError('"limit" is at most 500'); + const item = optString(body, "item"); + if (item !== undefined && !item.trim()) throw new OpsInputError('"item" must be a non-empty string'); + return jobResponse( + await importArchiveOrgAction(reqSlug(body, "slug"), { + item: item?.trim(), + items: optItems(body), + query: query?.trim(), + limit, + files: optFiles(body.files, "files"), + match: optString(body, "match"), + dryRun: optBool(body, "dryRun"), + }), + ); + }, ); } diff --git a/editor/app/api/ops/remote-listing/route.test.ts b/editor/app/api/ops/remote-listing/route.test.ts @@ -0,0 +1,74 @@ +import test from "node:test"; +import assert from "node:assert/strict"; +import { mkdir, readFile } from "node:fs/promises"; +import path from "node:path"; +import { setupOpsCorpus, callPost } from "../_testCorpus"; + +// Run with: +// pnpm -C editor exec tsx --test "app/api/ops/remote-listing/route.test.ts" +// +// The refusals before any job, and one listing run end to end against the e2e +// fake yt-dlp (it lists fake00000001…05 for any channel URL): the job's last +// line carries the diff, and nothing is written into the channel. + +const corpus = await setupOpsCorpus(null); +await corpus.writeChannelConfig("demo-odysee", { + handling: "transcribe", + platform: "odysee", + url: "https://odysee.com/@demo:0", +}); +await corpus.writeChannelConfig("demo-yt", { url: "https://www.youtube.com/@demo" }); +await corpus.writeChannelConfig("demo-nourl", {}); +const DATA = path.join(corpus.transcripts, "channels", "demo-odysee", "data"); +// Held: one id the fake lists, one it does not. +await mkdir(path.join(DATA, "fake00000002"), { recursive: true }); +await mkdir(path.join(DATA, "imported-elsewhere"), { recursive: true }); +const { POST } = await import("./route"); +const { getRegistry } = await import("yt-dlp-transcript-common/jobs/registry"); +const { getPaths } = await import("yt-dlp-transcript-common/lib/paths"); +test.after(() => corpus.cleanup()); + +test("refusals before any job: keys, slug, channel, platform, url", async () => { + const cases: [Record<string, unknown>, RegExp][] = [ + [{}, /"slug" is required/], + [{ slug: "demo-odysee", full: true }, /unknown key\(s\): full/], + [{ slug: "../x" }, /not a valid channel slug/], + [{ slug: "no-such" }, /Channel "no-such" not found/], + [{ slug: "demo-yt" }, /reads an Odysee or BitChute channel; "demo-yt" is on youtube/], + [{ slug: "demo-nourl" }, /no `url` configured/], + ]; + for (const [body, re] of cases) { + const r = await callPost(POST, body); + assert.equal(r.status, 400, JSON.stringify(body)); + assert.match(r.body.error ?? "", re, JSON.stringify(body)); + } + assert.deepEqual(await corpus.listJobIds(), []); +}); + +test("a listing is a job whose last line is the diff against what is held", async () => { + const r = await callPost(POST, { slug: "demo-odysee" }); + assert.equal(r.status, 200, JSON.stringify(r.body)); + const jobId = r.body.jobId as string; + for (let i = 0; i < 200; i++) { + const s = getRegistry().get(jobId)?.status; + if (s === "done" || s === "failed") break; + await new Promise((res) => setTimeout(res, 50)); + } + const log = await readFile(path.join(getPaths().jobsDir, `${jobId}.log`), "utf8"); + assert.equal(getRegistry().get(jobId)?.status, "done", log); + const line = log.split("\n").find((l) => l.startsWith("@@remote-listing ")); + assert.ok(line, log); + const result = JSON.parse(line!.slice("@@remote-listing ".length)); + assert.equal(result.platform, "odysee"); + assert.equal(result.listed, 5); + assert.equal(result.held, 1); + assert.deepEqual( + result.notHeld.map((n: { id: string }) => n.id), + ["fake00000001", "fake00000003", "fake00000004", "fake00000005"], + ); + assert.deepEqual(result.heldNotListed, ["imported-elsewhere"]); + // Read-only: no roster, no playlist, no maybe-missing written. + const { readdir } = await import("node:fs/promises"); + const top = await readdir(path.join(corpus.transcripts, "channels", "demo-odysee")); + assert.deepEqual(top.filter((f) => !["config.json", "data", "fake-ytdlp.invocations"].includes(f)), []); +}); diff --git a/editor/app/api/ops/remote-listing/route.ts b/editor/app/api/ops/remote-listing/route.ts @@ -0,0 +1,16 @@ +import { remoteListingAction } from "../../../channels/[slug]/pipelineActions"; +import { jobResponse, ops, reqSlug } from "../_lib"; + +export const dynamic = "force-dynamic"; + +// POST { slug } -> { ok: true, jobId } +// +// An Odysee or BitChute channel's listing, diffed against what it holds +// (controller/remoteListing.ts): one flat-playlist read on the platform's own +// queue, paced, nothing written. The job's log ends with +// `@@remote-listing {slug, platform, url, listedAt, listed, unparsed, held, +// notHeld: [{id, url}], heldNotListed: [id]}`; `pnpm ops get remote-listing +// <slug>` waits for it and prints that JSON on stdout. +export async function POST(request: Request) { + return ops(request, ["slug"], async (body) => jobResponse(await remoteListingAction(reqSlug(body, "slug")))); +} diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -17,10 +17,20 @@ import { platformQueueKey, } from "yt-dlp-transcript-common/lib/platform"; import { + ARCHIVE_ORG_SEARCH_DEFAULT_ROWS, resolveArchiveOrgImportUrl, - runArchiveOrgImport, + runArchiveOrgBatchImport, + searchArchiveOrgItems, + summarizeArchiveOrgBatch, + type ArchiveOrgItemRequest, } from "yt-dlp-transcript-common/controller/archiveOrgImport"; import { + REMOTE_LISTING_RESULT_MARKER, + isRemoteListingPlatform, + runRemoteListing, +} from "yt-dlp-transcript-common/controller/remoteListing"; +import { EnumerationIncompleteError } from "yt-dlp-transcript-common/ytdlp/runYtdlp"; +import { importPlatformSignal, resolveBitchuteImportUrl, } from "yt-dlp-transcript-common/controller/bitchuteImport"; @@ -721,42 +731,105 @@ function archiveOrgRefusal(paths: Paths, what: string): Promise<string | null> { return platformRefusal(paths, "archiveorg", "archive.org", what); } -// IMPORT CHOSEN FILES OF ONE archive.org ITEM (`pnpm ops import-archive-org`): -// one job on archive.org's own queue that imports the files one at a time, -// with a jittered pause between them, skipping any already downloaded, and -// stopping on a rate limit or three failures in a row -// (controller/archiveOrgImport.ts). `files` names exact paths in the item; -// `match` is a case-insensitive regex over them. `dryRun` lists what would be -// fetched and fetches nothing. +// IMPORT FROM archive.org (`pnpm ops import-archive-org`): one job on +// archive.org's own queue that imports files one at a time, with a jittered +// pause between them, skipping any already held, and stopping on a rate limit +// or three failures in a row (controller/archiveOrgImport.ts). +// +// item ONE item; exactly one of `files` (exact paths in the item) and +// `match` (a case-insensitive regex over them). +// items MANY: identifiers, or `{item, files? | match?}`; a bare identifier +// takes the top-level `match`, else every media original of it. +// query archive.org's advanced search: the first `limit` (100) items it +// matches, each as a bare identifier in `items`. +// +// `dryRun` lists each file as held, RESTRICTED (archive.org will answer +// 401/403) or "would get", and fetches nothing. The log ends with +// `summary: {…}`. +export type ArchiveOrgImportEntry = string | { item: string; files?: string[]; match?: string }; + +const ARCHIVE_ORG_ID_RE = /^[A-Za-z0-9][A-Za-z0-9._-]*$/; + export async function importArchiveOrgAction( slug: string, - opts: { item: string; files?: string[]; match?: string; dryRun?: boolean }, + opts: { + item?: string; + items?: ArchiveOrgImportEntry[]; + query?: string; + limit?: number; + files?: string[]; + match?: string; + dryRun?: boolean; + }, ): Promise<StreamActionResult> { const paths = getPaths(); const channelConfig = await readChannelConfig(paths, slug); if (!channelConfig) { return { ok: false, error: `Channel "${slug}" not found` }; } - const item = opts.item.trim(); - if (!/^[A-Za-z0-9][A-Za-z0-9._-]*$/.test(item)) { - return { ok: false, error: `"${item}" is not an archive.org identifier` }; - } - if ((opts.files === undefined) === (opts.match === undefined)) { - return { ok: false, error: 'Name the files: exactly one of "files" (a list) or "match" (a regex)' }; + const modes = [opts.item !== undefined, opts.items !== undefined, opts.query !== undefined].filter(Boolean).length; + if (modes !== 1) { + return { ok: false, error: 'Name what to import: exactly one of "item", "items" (a list) or "query" (an archive.org search)' }; } - if (opts.match !== undefined) { + const badRegex = (m: string | undefined): string | null => { + if (m === undefined) return null; try { - new RegExp(opts.match, "i"); + new RegExp(m, "i"); + return null; } catch (e) { - return { ok: false, error: `"match" is not a valid regex: ${(e as Error).message}` }; + return `"match" is not a valid regex: ${(e as Error).message}`; + } + }; + const top = badRegex(opts.match); + if (top) return { ok: false, error: top }; + if (opts.item === undefined && opts.files !== undefined) { + return { ok: false, error: '"files" names files of one item — with "items", give each its own: {"item", "files"}' }; + } + if (opts.query === undefined && opts.limit !== undefined) { + return { ok: false, error: '"limit" caps a "query" — it has nothing to cap here' }; + } + // Every item, validated before any job: one bad identifier is a refusal + // naming it, never a batch that stops part-way. + const requests: ArchiveOrgItemRequest[] = []; + const every = { match: opts.match ?? "." }; + if (opts.item !== undefined) { + const item = opts.item.trim(); + if (!ARCHIVE_ORG_ID_RE.test(item)) { + return { ok: false, error: `"${item}" is not an archive.org identifier` }; + } + if ((opts.files === undefined) === (opts.match === undefined)) { + return { ok: false, error: 'Name the files: exactly one of "files" (a list) or "match" (a regex)' }; + } + requests.push({ identifier: item, selection: opts.files !== undefined ? { files: opts.files } : { match: opts.match! } }); + } else if (opts.items !== undefined) { + const seen = new Set<string>(); + for (const entry of opts.items) { + const e = typeof entry === "string" ? { item: entry } : entry; + const item = e.item.trim(); + if (!ARCHIVE_ORG_ID_RE.test(item)) { + return { ok: false, error: `"${item}" is not an archive.org identifier` }; + } + if (seen.has(item)) continue; + seen.add(item); + if (e.files !== undefined && e.match !== undefined) { + return { ok: false, error: `item "${item}": "files" or "match", not both` }; + } + const own = badRegex(e.match); + if (own) return { ok: false, error: `item "${item}": ${own}` }; + requests.push({ + identifier: item, + selection: e.files !== undefined ? { files: e.files } : e.match !== undefined ? { match: e.match } : every, + }); } } if (!opts.dryRun) { const err = await lowDiskError(paths, slug); if (err) return err; - const refused = await archiveOrgRefusal(paths, "The archive.org import"); - if (refused) return { ok: false, info: true, error: refused }; } + // A dry run asks archive.org too (metadata, a search), so a hold or a + // cooldown refuses it as well. + const refused = await archiveOrgRefusal(paths, "The archive.org import"); + if (refused) return { ok: false, info: true, error: refused }; const downloadConfig = channelConfig.platform === "archiveorg" ? channelConfig @@ -767,12 +840,26 @@ export async function importArchiveOrgAction( paths, channelSlug: slug, fn: async (onLog, signal, _setProgress, ctx) => { - const result = await runArchiveOrgImport({ + let items = requests; + if (opts.query !== undefined) { + const rows = opts.limit ?? ARCHIVE_ORG_SEARCH_DEFAULT_ROWS; + const found = await searchArchiveOrgItems(opts.query, { rows, signal }).catch(async (err) => { + if ((err as { rateLimited?: boolean }).rateLimited) { + await recordDownloadBackoff("archiveorg", paths, "rate_limit"); + } + throw err; + }); + onLog( + `archive.org search ${JSON.stringify(opts.query)}: ${found.found} item(s) match, ` + + `taking the first ${found.items.length} (identifier order).\n`, + ); + items = found.items.map((it) => ({ identifier: it.identifier, selection: every })); + } + const batch = await runArchiveOrgBatchImport({ paths, slug, channelConfig: downloadConfig, - identifier: item, - selection: opts.files !== undefined ? { files: opts.files } : { match: opts.match! }, + items, onLog, signal, drainSignal: ctx.drainSignal, @@ -781,16 +868,86 @@ export async function importArchiveOrgAction( onImported: (id) => safeRevalidate([`/channels/${slug}/videos/${id}`]), }, }); + const summary = summarizeArchiveOrgBatch(batch); + onLog(`summary: ${JSON.stringify(summary)}\n`); safeRevalidate([`/channels/${slug}`, "/channels"]); // The platform's shared pacing state learns what archive.org said, so // the next import (and any other archive.org job) backs off or settles. - if (result.rateLimited) { + // Three items in a row refused with 401/403 back it off as a rate limit + // does: archive.org is refusing us, not one item. + if (batch.rateLimited || batch.refusedStorm) { await recordDownloadBackoff("archiveorg", paths, "rate_limit"); - } else if (result.imported.length > 0 && result.failed.length === 0) { + } else if (summary.imported > 0 && summary.failed === 0) { await recordPlatformClean("archiveorg", paths); } - if (result.failed.length > 0 && result.imported.length === 0 && !opts.dryRun) { - throw new Error(result.stopped ?? `${result.failed.length} file(s) failed`); + if (!opts.dryRun && summary.failed > 0 && summary.imported === 0) { + throw new Error(batch.stopped ?? `${summary.failed} file(s) failed`); + } + if (opts.item !== undefined && batch.missing.length > 0) { + throw new Error(batch.missing[0].error); + } + }, + }); +} + +// AN ODYSEE OR BITCHUTE CHANNEL'S REMOTE LISTING, diffed against what it +// holds (`pnpm ops get remote-listing <slug>`, controller/remoteListing.ts): +// one flat-playlist read of the channel's URL on the platform's own queue — +// one stream with every other request there — refused while the platform is +// held or cooling down, after the platform's import floor, and a rate-limited +// read backs the platform off. Nothing is written. The job's log ends with +// the result as `@@remote-listing {…}`. +export async function remoteListingAction(slug: string): Promise<StreamActionResult> { + const paths = getPaths(); + const channelConfig = await readChannelConfig(paths, slug); + if (!channelConfig) { + return { ok: false, error: `Channel "${slug}" not found` }; + } + if (!channelConfig.url) { + return { ok: false, error: "Channel has no `url` configured" }; + } + const platform = detectPlatform(channelConfig.url); + if (!isRemoteListingPlatform(platform)) { + return { + ok: false, + error: `A remote listing reads an Odysee or BitChute channel; "${slug}" is on ${platform ?? "an unknown platform"} — its sync lists it`, + }; + } + const label = PACED_LABELS[platform] ?? platform; + const refused = await platformRefusal(paths, platform, label, "The remote listing"); + if (refused) return { ok: false, info: true, error: refused }; + const settings = getSettings(); + return runManagedFunction({ + kind: "remote-listing", + queueKey: platformQueueKey(platform), + paths, + channelSlug: slug, + fn: async (onLog, signal) => { + await waitForPlatformGap({ + label, + remainingMs: async () => + Math.max( + platformGapRemainingMs(platform), + await platformCooldownRemainingMs(platform, paths).catch(() => 0), + ), + signal, + onLog, + }); + try { + const result = await runRemoteListing({ paths, slug, channelConfig, platform, onLog, signal }); + onLog(`${REMOTE_LISTING_RESULT_MARKER}${JSON.stringify(result)}\n`); + } catch (err) { + if (err instanceof EnumerationIncompleteError) { + await recordDownloadBackoff(platform, paths, "rate_limit").catch(() => {}); + } + throw err; + } finally { + notePlatformGap( + platform, + downloadGapMs(settings.sleepBetweenDownloadsSeconds, 0, 0, { + minSeconds: platformImportMinGapSeconds(platform), + }), + ); } }, }); diff --git a/scripts/archilyzer-ops.mjs b/scripts/archilyzer-ops.mjs @@ -15,6 +15,7 @@ // pnpm ops get tags [<tagId>] // pnpm ops get job <id> [--tail [<lines>]] // pnpm ops get jobs [--active | --failed] [--kind <kind>] [--slug <slug>] [--limit <n>] +// pnpm ops get remote-listing <slug> // pnpm ops job <cancel|drain|promote|force-release|retry> <id>... [--wait] // pnpm ops job retry-failed | job wait <id>... // pnpm ops list @@ -41,6 +42,8 @@ // pnpm ops feed-metadata --json '{"slug":"demo-podcast","dryRun":true}' --wait // pnpm ops import-video --json '{"slug":"demo-archive","url":"https://archive.org/details/example-item"}' // pnpm ops import-archive-org --json '{"slug":"demo-archive","item":"example-item","match":"\\.mp4$"}' --wait +// pnpm ops import-archive-org --json '{"slug":"demo-archive","query":"collection:example-collection","dryRun":true}' --wait +// pnpm ops get remote-listing demo-odysee // pnpm ops channel-config --json '{"slug":"x","patch":{"downloadFilterExclude":"rerun"}}' // pnpm ops channel-config --json '{"slug":"x","sites":[{"siteId":"anilyzer"}]}' // pnpm ops channel-config --json '{"slug":"x","sites":[],"excludeFromBuild":true}' @@ -162,6 +165,20 @@ const GETTERS = { cleanup: (slug) => `/api/ops/cleanup/${encodeURIComponent(slug)}`, }; +// READS THAT ARE JOBS: the answer needs a request upstream, so the editor +// runs it on the platform's queue and the job's log carries the result. `get` +// POSTs the action, waits, and prints the result on stdout (--wait-timeout +// bounds the wait). +const GET_JOBS = { + // An Odysee or BitChute channel's listing, diffed against what it holds + // (release 19 A6): {notHeld: [{id, url}], heldNotListed: [id], …}. + "remote-listing": (slug) => ({ + path: "/api/ops/remote-listing", + body: { slug }, + resultMarker: REMOTE_LISTING_RESULT_MARKER, + }), +}; + // Nouns whose read takes no argument. `get channel` without a slug is a // mistake; `get tags` without one is the whole vocabulary. const GET_ARG_OPTIONAL = new Set([ @@ -210,9 +227,14 @@ const ACTIONS = [ // subtitles, no media; the rewrite lands in metadata.history.json. "refresh-metadata", "import-video", - // Chosen media files of ONE archive.org item ({slug, item, files: [...] | - // match: "<regex>", dryRun?}): one job, one file at a time, paced. + // archive.org media into a channel: one item ({slug, item, files: [...] | + // match: "<regex>"}), many ({items: [...]}) or a search ({query, limit?}); + // dryRun? lists held / RESTRICTED / would-get. One job, one file at a time, + // paced, a pause between items. "import-archive-org", + // An Odysee or BitChute channel's listing diffed against what it holds + // ({slug}); `get remote-listing <slug>` is the same, waited on. + "remote-listing", // A channel's held videos given their media from a LOCAL archive — a // directory, a zip read in place, a 7z ({slug, source, items?, match?, // createRecords?, replace?, dryRun?}): one job, nothing fetched. @@ -306,6 +328,9 @@ const ACTIONS = [ // compact JSON. The same string is TRANSCRIBE_RESULT_MARKER in // common/controller/transcribeFile.ts. export const TRANSCRIBE_RESULT_MARKER = "@@transcribe-result "; +// The same for a remote listing (REMOTE_LISTING_RESULT_MARKER in +// common/controller/remoteListing.ts). +export const REMOTE_LISTING_RESULT_MARKER = "@@remote-listing "; // The provenance a tag write from this CLI carries. Everything else ignores it. function agentSource() { @@ -429,9 +454,25 @@ export function parseArgs(argv) { } if (positional[0] === "get") { const noun = positional[1]; + if (noun && GET_JOBS[noun]) { + if (!positional[2]) return { error: `get ${noun}: needs an argument` }; + if (used.size) { + return { error: `get ${noun} does not take ${[...used].map((f) => `--${f}`).join(", ")}` }; + } + const job = GET_JOBS[noun](positional[2]); + return { + method: "POST", + path: job.path, + body: job.body, + resultMarker: job.resultMarker, + wait: true, + quiet, + waitTimeout, + }; + } if (!noun || !GETTERS[noun]) { return { - error: `get: unknown noun "${noun ?? ""}" — known: ${Object.keys(GETTERS).join(", ")}`, + error: `get: unknown noun "${noun ?? ""}" — known: ${[...Object.keys(GETTERS), ...Object.keys(GET_JOBS)].join(", ")}`, }; } if (!positional[2] && !GET_ARG_OPTIONAL.has(noun)) { @@ -511,6 +552,7 @@ export function parseArgs(argv) { ...(action === "tag-videos" ? { defaultSource: agentSource() } : {}), // A job whose log carries a RESULT, which --wait prints on stdout. ...(action === "transcribe" ? { resultMarker: TRANSCRIBE_RESULT_MARKER } : {}), + ...(action === "remote-listing" ? { resultMarker: REMOTE_LISTING_RESULT_MARKER } : {}), wait, quiet, waitTimeout, @@ -593,6 +635,7 @@ export function usage() { " pnpm ops get settings [<key>]", " pnpm ops get storage | sites | workers | auto-queue | scheduler", " pnpm ops get cleanup <slug>", + " pnpm ops get remote-listing <slug> [--wait-timeout <seconds>]", " pnpm ops list", "", `Actions: ${ACTIONS.join(", ")}`, @@ -631,6 +674,12 @@ export function usage() { " reclaim (they overlap — never add them), what holds the rest, and the", " failed-transcriptions count. measured: false means unknown, not zero.", "", + 'remote-listing lists an Odysee or BitChute channel upstream and diffs it', + ' against what it holds: {"slug"}. One flat-playlist read on the platform\'s', + " queue (refused while it is held or cooling down; a 429 backs it off),", + " nothing written. `get remote-listing <slug>` waits for the job and prints", + " {listed, held, notHeld: [{id, url}], heldNotListed: [id], ...} on stdout.", + "", 'publish runs publish stages on the editor\'s publish queue, one at a time,', ' under one run id: {"verb": …}. "index" updates the index; "build" builds', ' "siteId"/"siteIds" (forced; the index first when stale; "runner":', @@ -773,6 +822,19 @@ export function usage() { " queue; a low disk or a rate limit stops it, and running the same body", " again resumes — saved videos are skipped.", "", + 'import-archive-org imports archive.org media into a channel, as one job on', + ' archive.org\'s queue: {"slug", "item", "files": [...] | "match": "<regex>"}', + ' for one item; {"slug", "items": ["<id>" | {"item", "files"? | "match"?},', + ' ...]} for many (a bare id takes "match", else every media original); or', + ' {"slug", "query": "<archive.org search>", "limit"?: 100} for the first', + ' items a search finds (at most 500). One file at a time, a jittered pause', + ' between files and between items; a record already held (on disk, or in', + ' the saved-video store) is skipped. "dryRun": true lists each file as', + ' held, RESTRICTED (archive.org marks it not for download: a fetch would', + ' answer 401/403) or would get, and fetches nothing. A 401/403 skips the', + ' rest of its item; three such items in a row, a 429 or three failures in', + ' a row stop the job. The log ends with a summary: line; a re-run resumes.', + "", 'attach-media copies each held video\'s file out of a LOCAL archive into the', ' saved-video store, as its source container (nothing fetched): {"slug",', ' "source"} — an absolute path to a directory, a .zip (read in place) or a', diff --git a/scripts/archilyzer-ops.test.mjs b/scripts/archilyzer-ops.test.mjs @@ -6,6 +6,7 @@ import assert from "node:assert/strict"; import test from "node:test"; import { + REMOTE_LISTING_RESULT_MARKER, TRANSCRIBE_RESULT_MARKER, editorEnvFiles, followJob, @@ -783,3 +784,28 @@ test("the archival writes are POSTs to their routes, each named in the usage", ( assert.match(usage(), /"dryRun": true answers each channel's preview/); assert.match(usage(), /"publish", "enabled"\?, "held"\?, "action"\?/); }); + +test("get remote-listing posts the job, waits, and carries the result marker", () => { + const r = parseArgs(["get", "remote-listing", "demo-odysee"]); + assert.equal(r.method, "POST"); + assert.equal(r.path, "/api/ops/remote-listing"); + assert.deepEqual(r.body, { slug: "demo-odysee" }); + assert.equal(r.wait, true); + assert.equal(r.resultMarker, REMOTE_LISTING_RESULT_MARKER); + assert.equal(REMOTE_LISTING_RESULT_MARKER, "@@remote-listing "); + assert.match(parseArgs(["get", "remote-listing"]).error, /needs an argument/); + assert.match(parseArgs(["get", "remote-listing", "x", "--counts"]).error, /does not take --counts/); + assert.equal(parseArgs(["get", "remote-listing", "x", "--wait-timeout", "60"]).waitTimeout, 60); + // The action form, waited on, prints the result too. + const a = parseArgs(["remote-listing", "--json", '{"slug":"x"}', "--wait"]); + assert.equal(a.path, "/api/ops/remote-listing"); + assert.equal(a.resultMarker, REMOTE_LISTING_RESULT_MARKER); +}); + +test("usage documents import-archive-org's three shapes and the remote listing", () => { + const u = usage(); + assert.match(u, /import-archive-org imports archive\.org media/); + assert.match(u, /"query": "<archive\.org search>"/); + assert.match(u, /RESTRICTED/); + assert.match(u, /remote-listing lists an Odysee or BitChute channel/); +});