Archilyzer · Source

archilyzer

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

commit a303efc1a364c232654338db9bc4f6809eaa45d9
parent 935569938878d3ecc5ae8d3c90286f239ec367cf
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 22 Jun 2026 22:15:25 -0400

Auto-queue: per-platform rate-limit backoff + subtitle progress

The auto-download runner discarded each download's outcome and counted
every result as success, so an Odysee 429 was marked done and the next
video on that platform was hit immediately — it never downloaded and
re-stormed the source on every restart.

- Add per-platform exponential backoff (rate_limit/network): pause that
  platform with a 1min→x2→30min (jittered) cooldown, keep other platforms
  flowing, and leave the video uncompleted so it retries after cooldown.
  A success clears the platform's backoff. Cooldowns persist across
  restarts in .auto-queue state (lapsed entries pruned on boot).
  New common/jobs/platformBackoff.ts, wired into autoRunner + autoQueueState.

- Classify download failures against the FULL yt-dlp stderr tail (surfaced
  as DownloadOutcomeRecord.failureClass) instead of the truncated per-attempt
  error, so a 429 logged as a WARNING before a different final error still
  triggers the backoff. Used by both the runner and the batch path.

- Per-video progress bar now steps once per subtitle track for subtitle-only
  (youtube-handling, --skip-download) downloads, which report no byte total;
  real media downloads keep their byte-based bar.

Unit tests added for the backoff math, state persistence, and subtitle
progress stepping.

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

Diffstat:
Mcommon/controller/autoRunner.ts | 125++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Acommon/jobs/autoQueueState.test.ts | 71+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/autoQueueState.ts | 17+++++++++++++++--
Acommon/jobs/platformBackoff.test.ts | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/platformBackoff.ts | 94+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/progressParsers.test.ts | 61+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/progressParsers.ts | 54++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/downloadOutcome.ts | 9++++++++-
Mcommon/ytdlp/downloadOneManaged.ts | 17+++++++++++++++++
Mcommon/ytdlp/runYtdlp.ts | 12++++++++----
Meditor/CHANGELOG.md | 2++
11 files changed, 518 insertions(+), 37 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -25,6 +25,13 @@ import { recordPick, writeAutoQueueState, } from "../jobs/autoQueueState"; +import { + clearBackoff, + isCoolingDown, + nextBackoff, + pruneExpired, +} from "../jobs/platformBackoff"; +import { type DownloadFailureClass } from "../lib/availability"; import { readChannelConfig } from "./channels"; import { readChannelSnapshot } from "./channelSnapshot"; import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; @@ -219,6 +226,9 @@ async function runLoop( const state = await readAutoQueueState(paths); const kindState = state[kind]; const runtime = kindState.runtime; + // Drop long-lapsed platform cooldowns on boot; entries still in (or recently + // out of) their window are kept so an Odysee 429 cooldown survives a restart. + pruneExpired(kindState.platformBackoff, Date.now()); const completed = new Set<string>(); const completedOrder: string[] = []; let counters = { transcribed: 0, failed: 0, skipped: 0 }; @@ -343,17 +353,21 @@ async function runLoop( const exclude = new Set<string>([...live.inFlight.keys(), ...completed]); removeIds(pending, exclude); // Download only: drop videos whose platform already has an in-flight - // download, so a busy platform yields to the next-priority free one. + // download (busy) OR is in a rate-limit/network backoff window (cooling + // down), so a busy/throttled platform yields to the next-priority free one. if (kind === "download") { - const busy = new Set<string>(); + const skip = new Set<string>(); for (const [pf, n] of platformInFlight) { - if (n >= PER_PLATFORM_CAP) busy.add(pf); + if (n >= PER_PLATFORM_CAP) skip.add(pf); + } + for (const pf of Object.keys(kindState.platformBackoff)) { + if (isCoolingDown(kindState.platformBackoff, pf, now)) skip.add(pf); } - if (busy.size > 0) { + if (skip.size > 0) { for (const leafId of Object.keys(pending)) { pending[leafId] = pending[leafId].filter((id) => { const slug = owner.get(id); - return !(slug && busy.has(platformKey(slug, slugToPlatform))); + return !(slug && skip.has(platformKey(slug, slugToPlatform))); }); } } @@ -373,14 +387,15 @@ async function runLoop( continue; } + const unitPlatform = + kind === "download" ? platformKey(channelSlug, slugToPlatform) : null; + let result: UnitResult = { outcome: "failed" }; launchUnit({ kind, paths, pick, channelSlug, tracker, onLog, signal }) - .then((outcome) => { - if (outcome === "transcribed") counters.transcribed++; - else if (outcome === "failed") counters.failed++; - else counters.skipped++; + .then((r) => { + result = r; }) .catch((err) => { - counters.failed++; + result = { outcome: "failed" }; onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); }) .finally(() => { @@ -388,11 +403,44 @@ async function runLoop( for (const id of pick.path) { live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1); } - if (kind === "download") { - const pf = platformKey(channelSlug, slugToPlatform); - platformInFlight.set(pf, Math.max(0, (platformInFlight.get(pf) ?? 0) - 1)); + if (unitPlatform) { + platformInFlight.set( + unitPlatform, + Math.max(0, (platformInFlight.get(unitPlatform) ?? 0) - 1), + ); + } + + if (result.outcome === "transcribed") counters.transcribed++; + else if (result.outcome === "failed") counters.failed++; + else counters.skipped++; + + // Per-platform backoff: a rate-limit / network failure pauses the whole + // platform with an exponential cooldown, and the video is deliberately + // NOT marked completed so it's retried once the cooldown lapses. Any + // other outcome marks the video done-for-session; a success also clears + // the platform's cooldown. + const backoffHit = + unitPlatform !== null && + (result.failureClass === "rate_limit" || + result.failureClass === "network"); + if (backoffHit && unitPlatform) { + const entry = nextBackoff( + kindState.platformBackoff[unitPlatform], + Date.now(), + ); + kindState.platformBackoff[unitPlatform] = entry; + const secs = Math.round((entry.until - Date.now()) / 1000); + onLog( + `Auto-download: ${unitPlatform} ${result.failureClass} — backing off ${secs}s (attempt ${entry.fails}). ${pick.videoId} will retry after cooldown.`, + ); + } else { + if (unitPlatform && result.outcome === "transcribed") { + clearBackoff(kindState.platformBackoff, unitPlatform); + } + markCompleted(pick.videoId); } - markCompleted(pick.videoId); + // Persist so the cooldown (and reset) survive a restart. + persist(); // The runner job has no channelSlug, so the registry's own // snapshot-on-finish hook doesn't fire — request a regen here so the // channel's pending counts refresh and a restart sees fresh buckets. @@ -409,9 +457,11 @@ async function runLoop( startedAt: Date.now(), }); for (const id of pick.path) live.active[id] = (live.active[id] ?? 0) + 1; - if (kind === "download") { - const pf = platformKey(channelSlug, slugToPlatform); - platformInFlight.set(pf, (platformInFlight.get(pf) ?? 0) + 1); + if (unitPlatform) { + platformInFlight.set( + unitPlatform, + (platformInFlight.get(unitPlatform) ?? 0) + 1, + ); } recordPick(kindState, { at: Date.now(), @@ -441,12 +491,19 @@ type LaunchArgs = { signal: AbortSignal; }; -// Run a single unit of work. Returns the outcome for counters. Hard cancel -// (signal) aborts an in-flight unit; drain is handled by the loop (it stops -// launching new units), so the drain signal is intentionally NOT forwarded here. -async function launchUnit(args: LaunchArgs): Promise< - "transcribed" | "skipped" | "failed" -> { +// The outcome of one unit. `failureClass` is set only for download failures +// (the runner uses rate_limit/network to drive per-platform backoff); it's +// derived from the returned outcome record via classifyDownloadFailure. +type UnitResult = { + outcome: "transcribed" | "skipped" | "failed"; + failureClass?: DownloadFailureClass; +}; + +// Run a single unit of work. Returns the outcome for counters/backoff. Hard +// cancel (signal) aborts an in-flight unit; drain is handled by the loop (it +// stops launching new units), so the drain signal is intentionally NOT +// forwarded here. +async function launchUnit(args: LaunchArgs): Promise<UnitResult> { if (args.kind === "transcription") { const res = await transcribeOneFromQueue({ paths: args.paths, @@ -461,7 +518,7 @@ async function launchUnit(args: LaunchArgs): Promise< onLog: args.onLog, signal: args.signal, }); - return res.outcome; + return { outcome: res.outcome }; } // Download: resolve the video's source URL and run the same managed // single-video download the import/download-missing paths use, so archive @@ -469,7 +526,7 @@ async function launchUnit(args: LaunchArgs): Promise< // per-platform in-flight gate in the loop keeps this to one download per // platform at a time (matching the per-platform serial queue's politeness). const config = await readChannelConfig(args.paths, args.channelSlug); - if (!config) return "skipped"; + if (!config) return { outcome: "skipped" }; const url = await findVideoSourceUrl( args.paths, args.channelSlug, @@ -478,7 +535,7 @@ async function launchUnit(args: LaunchArgs): Promise< ); if (!url) { args.onLog(`Auto-download: no source URL for ${args.pick.videoId} — skipped`); - return "skipped"; + return { outcome: "skipped" }; } const settings = getSettings(); const task = args.tracker.start({ @@ -487,7 +544,11 @@ async function launchUnit(args: LaunchArgs): Promise< kind: "download", }); try { - await downloadOneManaged({ + // downloadOneManaged does NOT throw on a 429/network failure — it returns + // a record whose last attempt carries the stderr tail + availability class. + // Classify it the same way the batch path does (runYtdlp.ts) so the runner + // can back off per-platform instead of immediately re-hammering the source. + const record = await downloadOneManaged({ channelSlug: args.channelSlug, channelConfig: config, paths: args.paths, @@ -499,13 +560,17 @@ async function launchUnit(args: LaunchArgs): Promise< globalSkipLiveDownloads: settings.skipLiveDownloads, appendArchive: true, }); - return "transcribed"; // generic "success" — feeds the done counter + if (record.status === "skipped-filtered") return { outcome: "skipped" }; + if (record.status.startsWith("ok")) return { outcome: "transcribed" }; + // failureClass is computed inside downloadOneManaged against the full stderr + // tail (catches a rate-limit WARNING masked by a later error line). + return { outcome: "failed", failureClass: record.failureClass }; } catch (err) { if (args.signal.aborted || (err as Error)?.name === "AbortError") { - return "skipped"; + return { outcome: "skipped" }; } args.onLog(`Auto-download failed for ${args.pick.videoId}: ${String(err)}`); - return "failed"; + return { outcome: "failed" }; } finally { task.end(); } diff --git a/common/jobs/autoQueueState.test.ts b/common/jobs/autoQueueState.test.ts @@ -0,0 +1,71 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { + emptyAutoQueueState, + readAutoQueueState, + writeAutoQueueState, +} from "./autoQueueState"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/autoQueueState.test.ts + +// Only autoQueueStateFile is touched by these functions, so a minimal Paths is fine. +async function withPaths( + fn: (paths: Paths) => Promise<void>, +): Promise<void> { + const dir = await mkdtemp(path.join(tmpdir(), "ttb-aqstate-")); + const paths = { + autoQueueStateFile: path.join(dir, ".auto-queue", "state.json"), + } as Paths; + try { + await fn(paths); + } finally { + await rm(dir, { recursive: true, force: true }); + } +} + +test("platformBackoff round-trips through write/read", async () => { + await withPaths(async (paths) => { + const state = emptyAutoQueueState(); + state.download.platformBackoff = { + odysee: { until: 123456, fails: 3 }, + youtube: { until: 999, fails: 1 }, + }; + await writeAutoQueueState(paths, state); + + const back = await readAutoQueueState(paths); + assert.deepEqual(back.download.platformBackoff, { + odysee: { until: 123456, fails: 3 }, + youtube: { until: 999, fails: 1 }, + }); + assert.deepEqual(back.transcription.platformBackoff, {}); + }); +}); + +test("missing file yields empty state with empty platformBackoff", async () => { + await withPaths(async (paths) => { + const back = await readAutoQueueState(paths); + assert.deepEqual(back.download.platformBackoff, {}); + assert.deepEqual(back.transcription.platformBackoff, {}); + }); +}); + +test("corrupt platformBackoff field coerces to empty without throwing", async () => { + await withPaths(async (paths) => { + // Write a valid state first so the .auto-queue dir exists, then overwrite + // the file with a hand-corrupted shape. + await writeAutoQueueState(paths, emptyAutoQueueState()); + await writeFile( + paths.autoQueueStateFile, + JSON.stringify({ + download: { platformBackoff: { odysee: { until: "nope" }, junk: 5 } }, + }), + ); + + const back = await readAutoQueueState(paths); + assert.deepEqual(back.download.platformBackoff, {}); + }); +}); diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts @@ -5,6 +5,10 @@ import { type AutoQueueRuntime, emptyAutoQueueRuntime, } from "./autoQueuePolicy"; +import { + type PlatformBackoffState, + coercePlatformBackoff, +} from "./platformBackoff"; // Persistent fairness state for the auto-queue runners. Unlike in-flight worker // counts (which are zero by definition after a restart and live only on the @@ -30,6 +34,10 @@ export type AutoQueueKindState = { runtime: AutoQueueRuntime; // Newest-first, bounded to AUTO_QUEUE_PICK_LOG_LIMIT entries. picks: AutoQueuePick[]; + // Per-platform download cooldowns (rate-limit/network backoff). Only used by + // the "download" kind; empty for transcription. Persisted so an Odysee 429 + // cooldown survives a server restart instead of re-storming on boot. + platformBackoff: PlatformBackoffState; }; export type AutoQueueState = { @@ -40,7 +48,7 @@ export type AutoQueueState = { export const AUTO_QUEUE_PICK_LOG_LIMIT = 50; export function emptyAutoQueueKindState(): AutoQueueKindState { - return { runtime: emptyAutoQueueRuntime(), picks: [] }; + return { runtime: emptyAutoQueueRuntime(), picks: [], platformBackoff: {} }; } export function emptyAutoQueueState(): AutoQueueState { @@ -93,7 +101,11 @@ function coerceKindState(value: unknown): AutoQueueKindState { .filter((p): p is AutoQueuePick => p !== null) .slice(0, AUTO_QUEUE_PICK_LOG_LIMIT) : []; - return { runtime: coerceRuntime(r.runtime), picks }; + return { + runtime: coerceRuntime(r.runtime), + picks, + platformBackoff: coercePlatformBackoff(r.platformBackoff), + }; } // Read the state file, tolerating a missing/corrupt file by returning empty @@ -123,6 +135,7 @@ export async function writeAutoQueueState( const trim = (k: AutoQueueKindState): AutoQueueKindState => ({ runtime: k.runtime, picks: k.picks.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT), + platformBackoff: k.platformBackoff ?? {}, }); const out: AutoQueueState = { transcription: trim(state.transcription), diff --git a/common/jobs/platformBackoff.test.ts b/common/jobs/platformBackoff.test.ts @@ -0,0 +1,93 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + BACKOFF_BASE_MS, + BACKOFF_MAX_MS, + type PlatformBackoffState, + clearBackoff, + coercePlatformBackoff, + isCoolingDown, + nextBackoff, + pruneExpired, +} from "./platformBackoff"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/platformBackoff.test.ts + +// Deterministic jitter: no jitter (returns 0.5 -> factor 1.0). +const noJitter = () => 0.5; + +test("nextBackoff doubles the cooldown across consecutive failures", () => { + const now = 1_000_000; + const first = nextBackoff(undefined, now, noJitter); + assert.equal(first.fails, 1); + assert.equal(first.until - now, BACKOFF_BASE_MS); // 60s + + const second = nextBackoff(first, now, noJitter); + assert.equal(second.fails, 2); + assert.equal(second.until - now, BACKOFF_BASE_MS * 2); // 120s + + const third = nextBackoff(second, now, noJitter); + assert.equal(third.fails, 3); + assert.equal(third.until - now, BACKOFF_BASE_MS * 4); // 240s +}); + +test("nextBackoff caps the cooldown at BACKOFF_MAX_MS", () => { + const now = 0; + let entry = nextBackoff(undefined, now, noJitter); + for (let i = 0; i < 20; i++) entry = nextBackoff(entry, now, noJitter); + assert.equal(entry.until - now, BACKOFF_MAX_MS); + assert.ok(entry.fails > 1); +}); + +test("nextBackoff applies bounded jitter (±10%)", () => { + const now = 0; + const low = nextBackoff(undefined, now, () => 0); // factor 0.9 + const high = nextBackoff(undefined, now, () => 1); // factor 1.1 + assert.equal(low.until, Math.round(BACKOFF_BASE_MS * 0.9)); + assert.equal(high.until, Math.round(BACKOFF_BASE_MS * 1.1)); +}); + +test("isCoolingDown reflects the until window", () => { + const state: PlatformBackoffState = { odysee: { until: 500, fails: 1 } }; + assert.equal(isCoolingDown(state, "odysee", 499), true); + assert.equal(isCoolingDown(state, "odysee", 500), false); // boundary: until is exclusive + assert.equal(isCoolingDown(state, "odysee", 501), false); + assert.equal(isCoolingDown(state, "youtube", 0), false); // absent platform +}); + +test("clearBackoff removes a platform (success resets escalation)", () => { + const state: PlatformBackoffState = { odysee: { until: 999, fails: 3 } }; + clearBackoff(state, "odysee"); + assert.deepEqual(state, {}); + // Next failure starts over at attempt 1. + const fresh = nextBackoff(state["odysee"], 0, noJitter); + assert.equal(fresh.fails, 1); +}); + +test("pruneExpired drops only long-lapsed entries, keeps active/recent ones", () => { + const now = 10_000_000; + const state: PlatformBackoffState = { + active: { until: now + 1000, fails: 1 }, // still cooling down + recent: { until: now - 1000, fails: 2 }, // lapsed but within retain window + old: { until: now - BACKOFF_MAX_MS - 1, fails: 5 }, // long lapsed + }; + pruneExpired(state, now); + assert.ok(state.active, "active kept"); + assert.ok(state.recent, "recently-lapsed kept for escalation memory"); + assert.equal(state.old, undefined, "long-lapsed dropped"); +}); + +test("coercePlatformBackoff tolerates corrupt/missing shapes", () => { + assert.deepEqual(coercePlatformBackoff(undefined), {}); + assert.deepEqual(coercePlatformBackoff("nope"), {}); + assert.deepEqual(coercePlatformBackoff({ odysee: { until: 5 } }), {}); // missing fails + assert.deepEqual(coercePlatformBackoff({ odysee: { until: "x", fails: 1 } }), {}); + assert.deepEqual( + coercePlatformBackoff({ + odysee: { until: 5, fails: 2 }, + junk: 42, + youtube: { until: 7, fails: 1, extra: "ignored" }, + }), + { odysee: { until: 5, fails: 2 }, youtube: { until: 7, fails: 1 } }, + ); +}); diff --git a/common/jobs/platformBackoff.ts b/common/jobs/platformBackoff.ts @@ -0,0 +1,94 @@ +// Per-platform exponential backoff for the auto-download runner. Pure logic + +// types, no I/O — so it's unit-testable and can be persisted by +// autoQueueState.ts and driven by autoRunner.ts. +// +// When a download fails with a rate-limit (HTTP 429 / throttle) or network +// error, the runner must stop launching downloads for THAT platform for a +// while instead of immediately re-picking the next video and hammering the +// source again (the original Odysee 429 storm). Each consecutive failure +// doubles the cooldown up to a cap; a successful download resets it. + +export type PlatformBackoffEntry = { + // Epoch ms: downloads for this platform are paused until this instant. + until: number; + // Consecutive backoff-triggering failures, drives the exponential delay. + fails: number; +}; + +// Keyed by platform name (e.g. "odysee", "youtube", or "unknown"). +export type PlatformBackoffState = Record<string, PlatformBackoffEntry>; + +// First cooldown after a single failure; doubles from here. +export const BACKOFF_BASE_MS = 60_000; // 1 minute +// Ceiling on the cooldown — long rate limits shouldn't park a platform forever. +export const BACKOFF_MAX_MS = 30 * 60_000; // 30 minutes +// ±10% jitter so multiple platforms recovering at once don't sync up. +const JITTER = 0.1; + +// Compute the next backoff entry for a platform given its previous one (or +// undefined for the first failure). `now` and `rand` are injected for +// deterministic tests; `rand` defaults to Math.random and returns [0, 1). +export function nextBackoff( + prev: PlatformBackoffEntry | undefined, + now: number, + rand: () => number = Math.random, +): PlatformBackoffEntry { + const fails = (prev?.fails ?? 0) + 1; + const base = Math.min(BACKOFF_BASE_MS * 2 ** (fails - 1), BACKOFF_MAX_MS); + const jitter = 1 + (rand() * 2 - 1) * JITTER; + return { until: now + Math.round(base * jitter), fails }; +} + +// True when the platform is currently in a cooldown window. +export function isCoolingDown( + state: PlatformBackoffState, + platform: string, + now: number, +): boolean { + const entry = state[platform]; + return entry !== undefined && entry.until > now; +} + +// Clear a platform's backoff (called on a successful download). Mutates in place. +export function clearBackoff( + state: PlatformBackoffState, + platform: string, +): void { + delete state[platform]; +} + +// Drop entries whose cooldown lapsed long enough ago that keeping their `fails` +// count for escalation memory no longer matters. Mutates in place, keeping the +// map bounded across a long-running server. We retain a lapsed entry for one +// full BACKOFF_MAX_MS window so a platform that recovers and immediately +// re-fails keeps escalating rather than restarting at the base delay. +export function pruneExpired( + state: PlatformBackoffState, + now: number, + retainMs: number = BACKOFF_MAX_MS, +): void { + for (const [platform, entry] of Object.entries(state)) { + if (entry.until + retainMs < now) delete state[platform]; + } +} + +// Defensive coercion for the persisted shape — any stored value that doesn't +// match is dropped, so a hand-edited/corrupt state file can't crash the runner. +// Mirrors coerceRuntime/coercePick in autoQueueState.ts. +export function coercePlatformBackoff(value: unknown): PlatformBackoffState { + const out: PlatformBackoffState = {}; + if (!value || typeof value !== "object") return out; + for (const [platform, raw] of Object.entries(value as Record<string, unknown>)) { + if (!raw || typeof raw !== "object") continue; + const r = raw as Record<string, unknown>; + if ( + typeof r.until === "number" && + Number.isFinite(r.until) && + typeof r.fails === "number" && + Number.isFinite(r.fails) + ) { + out[platform] = { until: r.until, fails: r.fails }; + } + } + return out; +} diff --git a/common/jobs/progressParsers.test.ts b/common/jobs/progressParsers.test.ts @@ -0,0 +1,61 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { createDownloadProgressParser } from "./progressParsers"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/progressParsers.test.ts + +// Subtitle-only (youtube --skip-download) progress lines, as seen in real job +// logs: total=NA while downloading, a real total only on the finished line. +const DL_NA = "DLOM_PROGRESS status=downloading downloaded=7168 total=NA frag=NA/NA speed=3442662 eta=NA"; +const DL_FIN = "DLOM_PROGRESS status=finished downloaded=177273 total=177273 frag=NA/NA speed=745217 eta=NA"; + +test("subtitle-only download advances one step per completed track", () => { + const p = createDownloadProgressParser(); + // Without a total, byte fraction can't move — but once the track count is + // known the bar steps per finished track. + assert.equal(p.feed("[info] vid: Downloading subtitles: en-orig, en"), null); + + let u = p.feed(DL_NA); // first track, downloading + assert.equal(u?.fraction, 0); + assert.match(u?.detail ?? "", /subs 0\/2/); + + u = p.feed(DL_FIN); // first track done + assert.equal(u?.fraction, 0.5); + assert.match(u?.detail ?? "", /subs 1\/2/); + + u = p.feed(DL_NA); // second track, downloading + assert.equal(u?.fraction, 0.5); + + u = p.feed(DL_FIN); // second track done + assert.equal(u?.fraction, 1); + assert.match(u?.detail ?? "", /subs 2\/2/); +}); + +test("without the subtitle announcement, behaviour is unchanged (jumps at finish)", () => { + const p = createDownloadProgressParser(); + const downloading = p.feed(DL_NA); + // total=NA, no track count → no fraction, just whatever detail it can build. + assert.equal(downloading?.fraction, undefined); + const finished = p.feed(DL_FIN); + assert.equal(finished?.fraction, 1); // finished line still reports 100% +}); + +test("a real media download is byte-based, not stepped, even after a subs announce", () => { + const p = createDownloadProgressParser(); + // youtube no-subs fallback: subs announced, then it switches to an audio + // download that DOES report a byte total — subtitle-step mode must release. + p.feed("[info] vid: Downloading subtitles: en-orig, en"); + const u = p.feed( + "DLOM_PROGRESS status=downloading downloaded=500 total=1000 frag=NA/NA speed=1000 eta=1", + ); + assert.equal(u?.fraction, 0.5); // 500/1000, not a subtitle step + assert.doesNotMatch(u?.detail ?? "", /subs /); +}); + +test("media download fraction is byte-based from the start", () => { + const p = createDownloadProgressParser(); + const u = p.feed( + "DLOM_PROGRESS status=downloading downloaded=250 total=1000 frag=NA/NA speed=1000 eta=3", + ); + assert.equal(u?.fraction, 0.25); +}); diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts @@ -190,6 +190,16 @@ export function createDownloadProgressParser(): { let lastProbeMs: number | undefined; let diffMs = 0; // smoothed per-probe duration increase let intervalSec: number | undefined; + // Subtitle-only downloads (youtube handling runs yt-dlp with --skip-download) + // fetch one or more tiny .vtt files. yt-dlp can't report a byte total for + // them, so DLOM_PROGRESS carries total=NA while downloading and the per-byte + // fraction never moves. Instead, count the announced subtitle tracks and + // advance the bar one step per completed track. `subsTotal > 0` means we're in + // this mode; it's cleared the moment a real media download starts (a + // "downloading" line with a concrete byte total — e.g. a youtube no-subs + // fallback that switches to an audio download). + let subsTotal = 0; + let subsDone = 0; const adjustEtaSeconds = (etaSeconds: number): number => { if (lastProbeMs === undefined || !intervalSec || intervalSec <= 0) { @@ -242,6 +252,50 @@ export function createDownloadProgressParser(): { } return null; } + + // yt-dlp announces the subtitle tracks it's about to fetch, e.g. + // [info] <id>: Downloading subtitles: en-orig, en + // Use the count as the denominator for the step-progress bar below. + const subsAnnounce = line.match(/Downloading subtitles:\s*(.+?)\s*$/); + if (subsAnnounce) { + const tracks = subsAnnounce[1] + .split(",") + .map((t) => t.trim()) + .filter(Boolean); + if (tracks.length > 0) { + subsTotal = tracks.length; + subsDone = 0; + } + return null; + } + + if (line.includes(DOWNLOAD_PROGRESS_MARKER)) { + const field = (key: string): string | undefined => + line.match(new RegExp(`\\b${key}=(\\S+)`))?.[1]; + const status = field("status"); + const total = num(field("total")); + // A "downloading" line with a real byte total means actual media is + // streaming (not a subtitle file) — leave subtitle-step mode. + if (status === "downloading" && total !== undefined && total > 0) { + subsTotal = 0; + } + if (subsTotal > 0) { + if (status === "finished" && subsDone < subsTotal) subsDone += 1; + const structured = parseStructuredDownloadProgress( + line, + adjustEtaSeconds, + ); + const out: ProgressUpdate = { + fraction: clamp01(subsDone / subsTotal), + }; + const parts: string[] = []; + if (structured?.detail) parts.push(structured.detail); + parts.push(`subs ${subsDone}/${subsTotal}`); + out.detail = parts.join(" · "); + return out; + } + } + const structured = parseStructuredDownloadProgress(line, adjustEtaSeconds); if (structured) return structured; // Legacy "[download] …%" lines (non-structured) carry no numeric ETA we diff --git a/common/lib/downloadOutcome.ts b/common/lib/downloadOutcome.ts @@ -1,7 +1,7 @@ // Client-safe types and constants for the per-video download outcome sidecar. // Mirrors availability.ts: server-only I/O lives in downloadOutcome-server.ts. -import type { Availability } from "./availability"; +import type { Availability, DownloadFailureClass } from "./availability"; import type { ChannelHandling } from "./channelConfig"; export type DownloadOutcomeStatus = @@ -63,6 +63,13 @@ export type DownloadOutcomeRecord = { startedAt: string; finishedAt: string; attempts: DownloadAttempt[]; + // Set on a failed download: the failure classified against the FULL stderr + // tail of the last attempt (not the truncated per-attempt `error`), so a + // rate-limit that yt-dlp logs as a WARNING before failing with a different + // final line (e.g. Odysee "HTTP Error 429" → "No video formats found") is + // still recognised. Drives the auto-runner's per-platform backoff. Undefined + // on success / skipped-filtered. + failureClass?: DownloadFailureClass; fellBackToTranscribe?: boolean; // Set when status is "skipped-filtered": which app-level filter declined the // download and why. Recorded so the UI/log can explain the skip. diff --git a/common/ytdlp/downloadOneManaged.ts b/common/ytdlp/downloadOneManaged.ts @@ -3,6 +3,7 @@ import { appendFile, mkdir, readdir, readFile } from "node:fs/promises"; import { createWriteStream, type WriteStream } from "node:fs"; import { execa } from "execa"; import { + classifyDownloadFailure, parseUnavailableFromStderr, type Availability, } from "../lib/availability"; @@ -317,6 +318,10 @@ async function runManagedDownload( let status: DownloadOutcomeStatus = "failed"; let fellBackToTranscribe = false; let lastArchiveLine: string | null = null; + // Full (untruncated) stderr tail of the most recent attempt, so a failed + // download can be classified (rate_limit/network) against everything yt-dlp + // printed — not just the last 3 lines stored on the attempt record. + let lastFullTail = ""; // ---------- Attempt 0: metadata prefetch + app-level filters ---------- // Split the per-video download into a cheap metadata-only pass followed by @@ -345,6 +350,7 @@ async function runManagedDownload( opts.videoUrl, ]; const prefetchRes = await runOneYtdlp(opts, channelDir, prefetchArgs); + lastFullTail = prefetchRes.stderrTail; const prefetchAvail = attemptSucceeded(prefetchRes.exitCode) ? undefined : parseUnavailableFromStderr(prefetchRes.stderrTail); @@ -472,6 +478,7 @@ async function runManagedDownload( primaryRes = await runOneYtdlp(opts, channelDir, primaryArgs); } + lastFullTail = primaryRes.stderrTail; const primaryAvail = attemptSucceeded(primaryRes.exitCode) ? undefined : parseUnavailableFromStderr(primaryRes.stderrTail); @@ -520,6 +527,7 @@ async function runManagedDownload( ...sourceArgs(opts.videoUrl, infoJsonPath), ]; const retryRes = await runOneYtdlp(opts, channelDir, retryArgs); + lastFullTail = retryRes.stderrTail; const retryAvail = attemptSucceeded(retryRes.exitCode) ? undefined : parseUnavailableFromStderr(retryRes.stderrTail); @@ -587,6 +595,7 @@ async function runManagedDownload( ), ]; const fallbackRes = await runOneYtdlp(opts, channelDir, fallbackArgs); + lastFullTail = fallbackRes.stderrTail; const fallbackAvail = attemptSucceeded(fallbackRes.exitCode) ? undefined : parseUnavailableFromStderr(fallbackRes.stderrTail); @@ -654,6 +663,13 @@ async function runManagedDownload( // ---------- Sidecar ---------- const finishedAt = new Date().toISOString(); + // Classify a failed download against the FULL stderr tail of the last attempt + // so a rate-limit logged as a WARNING (then masked by a different final error) + // is still caught. Skipped for successes and filter-skips. + const isFailure = status === "failed" || status === "failed-corrupt-source"; + const failureClass = isFailure + ? classifyDownloadFailure(lastFullTail, attempts.at(-1)?.availabilityClass) + : undefined; const record: DownloadOutcomeRecord = { videoId, webpageUrl: opts.videoUrl, @@ -661,6 +677,7 @@ async function runManagedDownload( startedAt, finishedAt, attempts, + ...(failureClass ? { failureClass } : {}), ...(fellBackToTranscribe ? { fellBackToTranscribe: true } : {}), }; // Only write the sidecar if we know which dir to put it in. If the very first diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -594,11 +594,15 @@ async function runManagedDownloads( } else if (outcome.status === "failed") { failedCount++; if (abortOnError && !firstFailure) { + // downloadOneManaged classifies against the full stderr tail; fall + // back to the truncated attempt error only if it's somehow absent. const lastAttempt = outcome.attempts[outcome.attempts.length - 1]; - const failureClass = classifyDownloadFailure( - lastAttempt?.error ?? "", - lastAttempt?.availabilityClass, - ); + const failureClass = + outcome.failureClass ?? + classifyDownloadFailure( + lastAttempt?.error ?? "", + lastAttempt?.availabilityClass, + ); // Per-video failures (private, members-only, age-gated, deleted) // are properties of that video, not signals the batch should // stop. Only abort on rate-limit / network / unclassified-fatal. diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,8 @@ # Changelog ## [Unreleased] +- **Auto-queue downloads now back off per-platform on rate limits instead of hammering the source.** When a managed auto-download hits an HTTP 429 / "too many requests" or a network error (most often on Odysee), the runner pauses *that platform* for an exponential cooldown (1 min, doubling up to 30 min, with jitter) while other platforms keep flowing, and the affected video is retried after the cooldown rather than being burned as a false success. Previously the runner discarded the download outcome, counted the rate-limited video as done, and immediately re-hit the same platform — so it never actually downloaded and re-stormed the source on every restart. Failures are now classified against the full yt-dlp stderr (not just the last few lines), so a 429 that yt-dlp logs as a WARNING before failing with a different final error still triggers the backoff. Cooldowns persist across restarts (`.auto-queue` state), and a successful download clears the platform's backoff. To stop a runner immediately, use Cancel in Active Jobs (Drain still waits for the in-flight video to finish, by design). See `common/jobs/platformBackoff.ts` and `common/controller/autoRunner.ts`. +- **Per-video progress bar now advances for subtitle-only (YouTube-handling) downloads.** YouTube-handling channels fetch only subtitles (`--skip-download`), which yt-dlp reports with no byte total — so the Active Jobs progress bar sat empty and jumped straight to 100%. It now steps once per subtitle track (e.g. `subs 1/2`) using the track list yt-dlp announces, while real media downloads (Odysee/transcribe) keep their byte-based bar. See `createDownloadProgressParser` in `common/jobs/progressParsers.ts`. - **New Archilyzer homepage (hub): a standalone marketing/docs site, configurable here.** A fourth workspace package, `homepage`, builds a single instance-level static site (`output: "export"`) that sits above the per-content export sites — for info/docs pages and cross-site charts that don't belong on any one content site. It's managed from the new **Manage → Homepage** page: edit branding (title/header/description/tagline/public URL/Cloudflare project) and author **Markdown pages** (a slug + nav label + body, with a live markdown-to-jsx preview), which the hub renders server-side at build (the `index` page is the home body; every other slug gets a `/<slug>` route). Page content lives in the data dir (`sites/_homepage/`, a reserved id `listSiteIds()` ignores), so copy changes need no code deploy. The hub's **/stats** dashboard charts downloads/transcriptions completed across **all** content sites, leading with a per-site breakdown — the chart engine gains a **Site** grouping option (`groupBy: "site"`) that fans each video out to every site exposing its channel, resolved through a channel→sites map the hub supplies; combined whole-pool totals remain one series. Build with `pnpm build:homepage` (a `compose-homepage` step stages whole-pool stats + the channel→sites map ahead of `next build`). See `common/lib/{homepage,homepagePages}.ts`, `common/bin/compose-homepage.ts`, `common/components/charts/channelSites.tsx`, the `homepage/` package, and `editor/app/homepage/*`. - **"Cut release" can now create the release commit for you.** After turning `## [Unreleased]` into a dated semver heading, cutting a release used to leave the changelog edit sitting in your working tree to `git commit` by hand. A **Commit changelog** checkbox now sits next to the **Cut release** button (on `/changelog` for the editor and `/deploy` for the export), checked by default — leave it on and the cut is followed by a path-limited `git commit` of just that one CHANGELOG.md, with the message `Release <workspace> <version>` (e.g. `Release export 0.4.1`). To keep the release commit clean it commits *only* the changelog: if the working tree has any *other* uncommitted change, the cut is refused up front (nothing is written) with an error telling you to commit or stash those first — a dirty changelog itself is fine, so uncommitted `[Unreleased]` bullets get folded into the release commit. Uncheck the box to cut without committing, exactly as before. See `editor/app/deploy/cutReleaseAction.ts`, `editor/app/deploy/components/CutReleaseForm.tsx`, and the new `common/lib/git.ts`. - **"Stop & keep progress" no longer mislabels the paused video as a failed transcription.** Using **Stop & keep progress** on a busy parakeet worker (or any partial-capable engine) sends the engine a graceful SIGTERM so it stops after the current window and the video resumes next run. But if the engine took longer than execa's 5-second force-kill window to exit — which a parakeet window routinely does, since finishing/stitching one ~480s window outlasts 5s — execa force-SIGKILLed it and the resulting "Command was killed with SIGTERM … forcefully terminated after 5000 milliseconds" error escaped the pause handling: it was treated as a genuine transcription failure and the video was written to the channel's `failed-transcriptions` file *permanently* (so even though its completed windows were cached for resume, it was skipped as "failed" on every later run). The transcribe path now recognizes that a force-killed **requested pause** is still a pause, not a failure — it returns the `paused` outcome (a skip, not a failure), so nothing lands in `failed-transcriptions` and the next "Transcribe missing" resumes it from the cached windows. Hard **Cancel** and **Drain** were never affected (their abort signal already classifies the kill as a skip). A video wrongly blacklisted by the old behavior won't auto-prune (it has real audio) — clear it with the channel's **Clear failed transcriptions** action to retry. See `common/controller/transcribeOne.ts`.