"use server"; import path from "node:path"; import { safeRevalidate } from "../../lib/safeRevalidate"; import { HANDLING_VALUES, type AudioFormat, type ChannelHandling, } from "yt-dlp-transcript-common/lib/channelConfig"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { downloadQueueKey, resolveQueueKey, } from "yt-dlp-transcript-common/lib/queueKeys"; import { detectPlatform, platformQueueKey, } from "yt-dlp-transcript-common/lib/platform"; import { ARCHIVE_ORG_SEARCH_DEFAULT_ROWS, resolveArchiveOrgImportUrl, 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"; import { heldPlatformRefusal, platformCooldownRemainingMs, recordDownloadBackoff, recordPlatformClean, } from "yt-dlp-transcript-common/jobs/downloadBackoff"; import { downloadGapMs } from "yt-dlp-transcript-common/jobs/platformBackoff"; import { notePlatformGap, platformGapRemainingMs, waitForPlatformGap, } from "yt-dlp-transcript-common/jobs/platformGap"; import { platformImportMinGapSeconds } from "yt-dlp-transcript-common/ytdlp/channelArgs"; import { countNotYetDownloaded, readChannelConfig, readChannelStat, } from "yt-dlp-transcript-common/controller/channels"; import { destinationExists, extractVideoId, runYtdlp, } from "yt-dlp-transcript-common/ytdlp/runYtdlp"; import { mergeRosterFile } from "yt-dlp-transcript-common/controller/rosterStore"; import { downloadOneManaged } from "yt-dlp-transcript-common/ytdlp/downloadOneManaged"; import { runMetadataScanJob } from "yt-dlp-transcript-common/controller/metadataScanJob"; import { runFeedMetadataJob } from "yt-dlp-transcript-common/controller/feedMetadataJob"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { isGateHeld } from "yt-dlp-transcript-common/lib/pauseGates"; import { resolveCookiePolicy } from "yt-dlp-transcript-common/lib/cookiePolicy"; import { diskGate } from "yt-dlp-transcript-common/lib/diskSpace"; import { formatBytes } from "yt-dlp-transcript-common/lib/format"; import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; import type { Paths } from "yt-dlp-transcript-common/lib/paths"; // Preflight disk-space gate: keep a download job from even starting when free // space is already below the configured floor. Returns an error result to // surface immediately, or null to proceed. The running batch separately stops // between videos via the gate in runManagedDownloads (common/ytdlp/runYtdlp.ts). async function lowDiskError( paths: Paths, slug: string, ): Promise<{ ok: false; error: string } | null> { // The operator clicked this: only the floor applies and the shared hysteresis // latch is left alone (see diskGate). runYtdlp re-checks per video mid-batch, // against this same volume — the channel's own, which for a relocated channel // is the platter and not the corpus disk. const disk = await diskGate(paths, getSettings(), { mode: "manual", dir: path.join(paths.channelsDir, slug, "data"), }); if (disk.ok) return null; return { ok: false, error: `Low disk space: ${formatBytes(disk.freeBytes)} free, ` + `${formatBytes(disk.thresholdBytes)} required. Free up space or ` + `lower the floor in Settings.`, }; } async function runPipelineAction( slug: string, mode: | "store-playlist" | "download-from-playlist" | "sync" | "download-missing" | "download-missing-subs" | "retry-bucket", kind: string, queueKey?: string, options?: { ignoreArchive?: boolean; abortOnError?: boolean; shardTotal?: number; shardIndex?: number; bucketIds?: ReadonlyArray; handlingOverride?: ChannelHandling; // retry-bucket only (the "Needs cookies" bucket): force cookie mode // "always" for this run so every yt-dlp invocation carries the configured // cookies and the defer-mode exclusion is bypassed. forceCookies?: boolean; // retry-bucket only (the "YouTube auto-captions only" bucket): let an // ASR-provenance VTT stop counting as an existing destination, so the audio // download actually runs for videos that already have auto-captions. replaceAutoSubs?: boolean; // Per-run persistence overrides (Phase 2), forwarded to runYtdlp → // downloadOneManaged for every managed download in this run. keepSourceVideoOverride?: boolean; extractImmediately?: boolean; audioFormatOverride?: AudioFormat; // sync only: force the full sweep (whole-listing re-read) regardless of the // configured cadence. Undefined leaves the decision to the cadence gate. forceFullSweep?: boolean; // Replay descriptor, forwarded onto the job record so it can be replayed (Retry). spec?: JobSpec; }, ): Promise { 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" }; } // Global downloads pause: every mode that writes media is blocked with a // friendly notice (info, not an error). store-playlist only enumerates the // video list — no media is written — so it stays allowed even while paused. if (mode !== "store-playlist" && isGateHeld(getSettings(), "download")) { return { ok: false, info: true, error: "Downloads are paused. Resume downloads from the dashboard (or Jobs) " + "to run this.", }; } // Platform key shared with the auto-download runner (detectPlatform(url) ?? // "unknown") — drives the 429 cooldown both honor. const platform = detectPlatform(channelConfig.url) ?? "unknown"; // A HELD platform refuses every manual fetch, not only a Sync (release 17, // slice RL), while its next probe is still ahead: the lane is probing it // once an hour. Once the probe is due, a manual run IS the probe, and a // clean one lifts the hold (onPlatformClean below). The sentence names the // hold and the next probe. store-playlist stays allowed, as it is under the // downloads pause. if (mode !== "store-playlist") { const held = await heldPlatformRefusal( platform, mode === "sync" ? "Sync" : "This download", paths, ); if (held) return { ok: false, info: true, error: held }; } // A manually clicked Sync RESPECTS the per-platform rate-limit cooldown: if // auto-download (or a prior sync) hit a 429, refuse rather than re-storm the // source. Neutral (info) result so the UI shows it as a notice, not an error. // Other download modes still run (and will write their own backoff on a 429). if (mode === "sync") { const remainingMs = await platformCooldownRemainingMs(platform, paths); if (remainingMs > 0) { const secs = Math.ceil(remainingMs / 1000); return { ok: false, info: true, error: `${platform} is in a rate-limit cooldown (${secs}s remaining). ` + `Auto-download or a prior sync hit HTTP 429 — sync will run once the ` + `cooldown lapses.`, }; } } // store-playlist only fetches the video list (no media written), so it is not // gated; every other mode downloads audio/subtitles into transcriptsDir. if (mode !== "store-playlist") { const err = await lowDiskError(paths, slug); if (err) return err; } return runManagedFunction({ kind, queueKey: resolveQueueKey(downloadQueueKey(channelConfig), queueKey), paths, channelSlug: slug, spec: options?.spec, fn: async (onLog, signal, setProgress, ctx) => { // Baseline (channel-wide downloadCount at run start) handed to runYtdlp so // a sharded download can refine the progress total to its own slice once // the slice is known. let progressBaseline: number | undefined; if (mode !== "store-playlist") { const stat = await readChannelStat(paths, slug); if (stat) { progressBaseline = stat.downloadCount; let target: number; if (mode === "retry-bucket") { const bucketIds = options?.bucketIds ?? []; const remaining = await countNotYetDownloaded( paths, slug, bucketIds, ); target = stat.downloadCount + remaining; } else { target = stat.playlistCount && stat.playlistCount > 0 ? stat.playlistCount : stat.videoCount; } setProgress({ metric: "downloads", initial: stat.downloadCount, target, }); } } await runYtdlp({ channelSlug: slug, mode, channelConfig, paths, onLog, signal, drainSignal: ctx.drainSignal, tracker: makeTaskTracker(ctx, onLog), ignoreArchive: options?.ignoreArchive, abortOnError: options?.abortOnError, shardTotal: options?.shardTotal, shardIndex: options?.shardIndex, bucketIds: options?.bucketIds, handlingOverride: options?.handlingOverride, forceCookies: options?.forceCookies, forceFullSweep: options?.forceFullSweep, replaceAutoSubs: options?.replaceAutoSubs, keepSourceVideoOverride: options?.keepSourceVideoOverride, extractImmediately: options?.extractImmediately, audioFormatOverride: options?.audioFormatOverride, setProgress, progressBaseline, // On a 429/network failure, record the shared per-platform cooldown so // the auto-download runner (and a later manual sync) back off too. The // class matters since release 17: only a rate limit doubles the pace. onPlatformBackoff: (failureClass) => recordDownloadBackoff(platform, paths, failureClass), // A run the source answered cleanly settles the platform — a held // one's hold included — as a clean lane probe does (release 17). onPlatformClean: () => recordPlatformClean(platform, paths), }); // The channel report (snapshot) is regenerated automatically after this // job finishes, via the global debounced scheduler hooked into // runManagedFunction's completion. See common/jobs/snapshotScheduler.ts. safeRevalidate([ `/channels/${slug}`, "/channels", ["/operations/[id]", "page"], "/cleanup", "/", ]); }, }); } export async function storePlaylistAction( slug: string, queueKey?: string, ): Promise { return runPipelineAction(slug, "store-playlist", "store-playlist", queueKey, { spec: { kind: "store-playlist", slug, params: { queueKey } }, }); } export async function downloadAction( slug: string, queueKey?: string, abortOnError?: boolean, ): Promise { return runPipelineAction( slug, "download-from-playlist", "download-from-playlist", queueKey, { abortOnError, spec: { kind: "download-from-playlist", slug, params: { queueKey, abortOnError }, }, }, ); } export async function downloadMissingAction( slug: string, queueKey?: string, ignoreArchive?: boolean, abortOnError?: boolean, shardTotal?: number, shardIndex?: number, // Per-run persistence overrides (Phase 2). The canonical use is a backfill that // extracts audio immediately to save disk (extractImmediately), but a forced // keep / audio-format override is supported too. keepSourceVideoOverride?: boolean, extractImmediately?: boolean, audioFormatOverride?: AudioFormat, ): Promise { return runPipelineAction( slug, "download-missing", "download-missing", queueKey, { ignoreArchive, abortOnError, shardTotal, shardIndex, keepSourceVideoOverride, extractImmediately, audioFormatOverride, spec: { kind: "download-missing", slug, params: { queueKey, ignoreArchive, abortOnError, shardTotal, shardIndex, keepSourceVideoOverride, extractImmediately, audioFormatOverride, }, }, }, ); } // A sync decides for itself whether to run the cheap newest-first paged walk or // the periodic full sweep (see common/jobs/deepSync.ts). Pass fullSweep to force // one now — the sweep is still a single `sync` job: one row, one log, one // Cancel, and the same per-platform serialization. export async function syncAction( slug: string, queueKey?: string, fullSweep?: boolean, ): Promise { return runPipelineAction(slug, "sync", "sync", queueKey, { forceFullSweep: fullSweep, spec: { kind: "sync", slug, params: { queueKey, fullSweep } }, }); } // THE METADATA SCAN. Not routed through runPipelineAction, and the differences // are the whole reason it is its own operation: // // - It is NOT held by the downloads pause. It fetches no media, and the // operator runs it precisely to decide what a paused download lane should // fetch when it resumes. // - It needs no disk gate, for the same reason: the only thing it writes is // one JSON file of titles. // - It DOES respect the per-platform rate-limit cooldown, like sync, and // records one of its own when the source pushes back — the scan is one // request per listed video and is the most rate-limitable thing here. export async function runMetadataScanAction( slug: string, queueKey?: string, ): Promise { 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) ?? "unknown"; const held = await heldPlatformRefusal(platform, "The metadata scan", paths); if (held) return { ok: false, info: true, error: held }; const remainingMs = await platformCooldownRemainingMs(platform, paths); if (remainingMs > 0) { const secs = Math.ceil(remainingMs / 1000); return { ok: false, info: true, error: `${platform} is in a rate-limit cooldown (${secs}s remaining). ` + `The metadata scan will run once the cooldown lapses.`, }; } // THE JOB ITSELF LIVES IN common/controller/metadataScanJob.ts, because the // download lane's runner dispatches the same one and cannot reach a server // action. What stays here is what a CLICK is owed and a runner is not: the // cooldown answered as a sentence rather than a queued job, and the three // revalidations. return runMetadataScanJob({ paths, slug, channelConfig, queueKey, afterRun: () => { safeRevalidate([ `/channels/${slug}`, "/channels", ["/operations/[id]", "page"], ]); }, }); } export async function downloadMissingSubsAction( slug: string, queueKey?: string, abortOnError?: boolean, ): Promise { return runPipelineAction( slug, "download-missing-subs", "download-missing-subs", queueKey, { abortOnError, spec: { kind: "download-missing-subs", slug, params: { queueKey, abortOnError }, }, }, ); } export async function retryBucketAction( slug: string, bucketIds: string[], queueKey?: string, abortOnError?: boolean, handlingOverride?: string, // The snapshot bucket these ids came from. Passed by the named bucket controls // (partial downloads / no-transcript) so the job can be replayed and re-run // against the CURRENT bucket; omitted by ad-hoc checkbox selections, which are // therefore not replayable. bucketKey?: ReplayBucket, // Needs-cookies bucket only: run with cookie mode forced to "always". forceCookies?: boolean, // Auto-captions-only bucket: bypass the destination-exists prefilter for // videos whose only "destination" is the YouTube ASR VTT we're replacing. // Always paired with handlingOverride "transcribe" (a youtube-handling run // would re-fetch subs instead of audio). replaceAutoSubs?: boolean, ): Promise { if (!Array.isArray(bucketIds) || bucketIds.length === 0) { return { ok: false, error: "No video IDs supplied for retry." }; } let handling: ChannelHandling | undefined; if (handlingOverride) { if (!(HANDLING_VALUES as ReadonlyArray).includes(handlingOverride)) { return { ok: false, error: `Invalid handling override: ${handlingOverride}`, }; } handling = handlingOverride as ChannelHandling; } const spec: JobSpec | undefined = bucketKey ? { kind: "retry-bucket", slug, bucket: bucketKey, params: { queueKey, abortOnError, handlingOverride, forceCookies, replaceAutoSubs, }, } : undefined; return runPipelineAction( slug, "retry-bucket", "retry-bucket", queueKey, { bucketIds, handlingOverride: handling, abortOnError, forceCookies, replaceAutoSubs, spec, }, ); } // Import a single one-off video by URL into this channel. Unlike the playlist // pipeline, the URL comes straight from user input rather than the stored // playlist. downloadOneManaged does the rest: it pins yt-dlp's output to // data//, branches on the channel's handling (youtube subs vs transcribe // audio), runs the auth/no-subs fallbacks, and appends to the archive — exactly // like downloadVideoPipelineAction, which feeds it a URL looked up from disk. export async function importVideoAction( slug: string, url: string, queueKey?: string, ): Promise { const paths = getPaths(); const channelConfig = await readChannelConfig(paths, slug); if (!channelConfig) { return { ok: false, error: `Channel "${slug}" not found` }; } let videoUrl = url.trim(); if (!/^https?:\/\//i.test(videoUrl)) { return { ok: false, error: "Enter a valid video URL" }; } // AN archive.org URL is resolved first (one cached request for the item's // metadata): an item holding several media files is refused — import a file // of it, or several with import-archive-org — and the URL becomes the // canonical page of what is imported (controller/archiveOrgImport.ts). It // runs on archive.org's own queue whatever the channel's platform — over // BitTorrent when it can, else straight from archive.org, never yt-dlp // (controller/archiveOrgDownload.ts) — and a file already on disk is not // fetched again. const urlPlatform = detectPlatform(videoUrl); const archiveOrg = urlPlatform === "archiveorg"; // A BitChute URL likewise (controller/bitchuteImport.ts): canonicalized, run // on BitChute's own queue with BitChute's args whatever the channel's // platform, refused while BitChute is held or cooling down, never fetched // again when on disk — and what BitChute answers is recorded on its pacing // state below. const bitchute = urlPlatform === "bitchute"; // ANY URL ON A PLATFORM WITH AN IMPORT FLOOR (platformImportMinGapSeconds: // BitChute, Odysee) is paced though each import is its own job: it runs // on the platform's own queue with the platform's args, is refused while the // platform is held or cooling down, waits out the floor after the last import // there (jobs/platformGap.ts) and any cooldown that began since it was // queued, and records what the platform answered. const pacedPlatform = urlPlatform && platformImportMinGapSeconds(urlPlatform) > 0 ? urlPlatform : null; const pacedLabel = pacedPlatform ? (PACED_LABELS[pacedPlatform] ?? pacedPlatform) : ""; let downloadConfig = channelConfig; if (archiveOrg) { const refused = await archiveOrgRefusal(paths, "The import"); if (refused) return { ok: false, info: true, error: refused }; const resolved = await resolveArchiveOrgImportUrl(videoUrl); if (!resolved.ok) return { ok: false, error: resolved.error }; videoUrl = resolved.url; if ( await destinationExists( path.join(paths.channelsDir, slug, "data"), resolved.id, channelConfig.handling, ) ) { return { ok: false, info: true, error: `Already downloaded: data/${resolved.id}/ — archive.org is not asked for it again.`, }; } if (channelConfig.platform !== "archiveorg") { downloadConfig = { ...channelConfig, platform: "archiveorg" }; } } if (pacedPlatform) { const refused = await platformRefusal(paths, pacedPlatform, pacedLabel, "The import"); if (refused) return { ok: false, info: true, error: refused }; if (channelConfig.platform !== pacedPlatform) { downloadConfig = { ...channelConfig, platform: pacedPlatform }; } } if (bitchute) { const resolved = resolveBitchuteImportUrl(videoUrl); if (!resolved.ok) return { ok: false, error: resolved.error }; videoUrl = resolved.url; if ( await destinationExists( path.join(paths.channelsDir, slug, "data"), resolved.id, channelConfig.handling, ) ) { return { ok: false, info: true, error: `Already downloaded: data/${resolved.id}/ — BitChute is not asked for it again.`, }; } } // Best-effort canonical id: used only for revalidation/labels. When null, // downloadOneManaged falls back to %(id)s and the reconcile pass repairs the // dir, so we don't hard-fail here. const videoId = extractVideoId(videoUrl) ?? undefined; const err = await lowDiskError(paths, slug); if (err) return err; const settings = getSettings(); // A URL on a platform with its own queue runs there. The channel page's // queue control starts at the CHANNEL's queue, which is therefore read as // "no choice made"; any other queue the operator picks still wins. const channelQueue = downloadQueueKey(channelConfig); const ownQueue = archiveOrg || pacedPlatform ? platformQueueKey(urlPlatform) : null; const override = ownQueue && queueKey !== undefined && queueKey.trim() === channelQueue ? undefined : queueKey; return runManagedFunction({ kind: "import-one", queueKey: resolveQueueKey(ownQueue ?? channelQueue, override), paths, channelSlug: slug, videoId, fn: async (onLog, signal, _setProgress, ctx) => { const task = makeTaskTracker(ctx, onLog).start({ id: videoId ?? videoUrl, label: videoId ?? videoUrl, kind: "download", }); try { if (pacedPlatform) { await waitForPlatformGap({ label: pacedLabel, remainingMs: async () => Math.max( platformGapRemainingMs(pacedPlatform), await platformCooldownRemainingMs(pacedPlatform, paths).catch(() => 0), ), signal, onLog: task.onLog, }); } const outcome = await downloadOneManaged({ channelSlug: slug, channelConfig: downloadConfig, paths, videoUrl, onLog: task.onLog, signal, cookiePolicy: resolveCookiePolicy(settings, channelConfig), inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, globalSkipLiveDownloads: settings.skipLiveDownloads, appendArchive: true, }); if (pacedPlatform) { // The platform's shared pacing state learns what it said: a 429 // backs every path on it off (the lane, a sync, the next import). // Best-effort, as the batch downloads' bookkeeping is. The next // import there waits the floor from now, whatever the answer. notePlatformGap( pacedPlatform, downloadGapMs(settings.sleepBetweenDownloadsSeconds, 0, 0, { minSeconds: platformImportMinGapSeconds(pacedPlatform), }), ); const answer = importPlatformSignal(outcome); try { if (answer === "rate_limit" || answer === "network") { await recordDownloadBackoff(pacedPlatform, paths, answer); } else if (answer === "clean") { const line = await recordPlatformClean(pacedPlatform, paths); if (line) task.onLog(line); } } catch { /* shared-state write is best-effort */ } } // ADD-ONLY: record the imported video in the roster so a one-off import // is a known member of the channel rather than an orphan dir, and so a // later retry of the same URL is possible even if it never appears in a // listing. An import is not an enumeration, so lastListedAt does not // move (see common/controller/rosterStore.ts). if (videoId) { await mergeRosterFile( paths, slug, [{ id: videoId, url: videoUrl }], new Date().toISOString(), "import", ).catch(() => { /* the download succeeded; a roster write failure must not fail it */ }); safeRevalidate([`/channels/${slug}/videos/${videoId}`]); } safeRevalidate([`/channels/${slug}`, "/channels"]); } finally { task.end(); } }, }); } // How an import names a paced platform in what it says. const PACED_LABELS: Partial> = { bitchute: "BitChute", odysee: "Odysee" }; // A platform held, or in a rate-limit cooldown: a sentence, else null. An // import asks a platform nothing while it has asked us to wait. async function platformRefusal( paths: Paths, platform: string, label: string, what: string, ): Promise { const held = await heldPlatformRefusal(platform, what, paths); if (held) return held; const remainingMs = await platformCooldownRemainingMs(platform, paths); if (remainingMs > 0) { return `${label} is in a rate-limit cooldown (${Math.ceil(remainingMs / 1000)}s remaining). ${what} can run once it lapses.`; } return null; } function archiveOrgRefusal(paths: Paths, what: string): Promise { return platformRefusal(paths, "archiveorg", "archive.org", what); } // 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; items?: ArchiveOrgImportEntry[]; query?: string; limit?: number; files?: string[]; match?: string; dryRun?: boolean; }, ): Promise { const paths = getPaths(); const channelConfig = await readChannelConfig(paths, slug); if (!channelConfig) { return { ok: false, error: `Channel "${slug}" not found` }; } 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)' }; } const badRegex = (m: string | undefined): string | null => { if (m === undefined) return null; try { new RegExp(m, "i"); return null; } catch (e) { 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(); 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; } // 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 : { ...channelConfig, platform: "archiveorg" as const }; return runManagedFunction({ kind: "import-archive-org", queueKey: platformQueueKey("archiveorg"), paths, channelSlug: slug, fn: async (onLog, signal, _setProgress, ctx) => { 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, items, onLog, signal, drainSignal: ctx.drainSignal, dryRun: opts.dryRun, deps: { 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. // 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 (summary.imported > 0 && summary.failed === 0) { await recordPlatformClean("archiveorg", paths); } 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 `, 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 { 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), }), ); } }, }); } // A PODCAST CHANNEL'S RECORDS COMPLETED FROM ITS RSS FEED (the // `feed-metadata` job, controller/feedMetadataBackfill.ts): one fetch of the // channel's url, then title, date, description and duration written into each // record that lacks them, through the metadata history. `dryRun` counts // matched / unmatched / already complete and writes nothing. // // The feed is one request to the channel's host, so a held host or one in its // rate-limit cooldown is answered with a sentence, as the metadata scan is. export async function feedMetadataBackfillAction( slug: string, dryRun?: boolean, ): Promise { 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) ?? "unknown"; const held = await heldPlatformRefusal(platform, "The feed backfill", paths); if (held) return { ok: false, info: true, error: held }; const remainingMs = await platformCooldownRemainingMs(platform, paths); if (remainingMs > 0) { const secs = Math.ceil(remainingMs / 1000); return { ok: false, info: true, error: `${platform} is in a rate-limit cooldown (${secs}s remaining). Run the feed backfill once it lapses.`, }; } return runFeedMetadataJob({ paths, slug, ...(dryRun ? { dryRun: true } : {}), requestedBy: "editor", afterRun: () => { safeRevalidate([ `/channels/${slug}`, "/channels", ["/channels/[slug]/videos/[id]", "page"], ]); }, }); }