Archilyzer · Source

archilyzer

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

commit e3178ac1a4634aedef09ad9758cc2714be12c26f
parent c033cc421678de189db5f844cf0b26075f4f8521
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  2 Oct 2026 01:09:35 -0400

common: a hold cannot outlive its lane — refused only before its probe, lifted by a clean manual run; the pace eases hourly (review H2, L1, L2, L3, L5)

A manual Sync, download or scan on a held platform is refused only while
its probe is ahead; once due, the run is the probe, and a clean one
(runYtdlp onPlatformClean, a clean metadata scan) clears the hold and the
backoff through recordPlatformClean. clearPlatformHold drops hold, backoff
and pace. A held platform stays in the status view whatever its probe
time. A probe whose media came down but whose subtitles 429'd lifts the
hold without easing the pace. The pace eases one step per hour with no
rate limit (steppedAt; effectivePaceSeconds for readers). Download missing
subs waits the lane gap between videos. A hold reached through network
failures alone is worded as failing, not rate-limited.

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

Diffstat:
Mcommon/bin/doctor.test.ts | 4++--
Mcommon/bin/doctor.ts | 11++++++++---
Mcommon/jobs/autoQueueState.test.ts | 6+++---
Mcommon/jobs/autoQueueState.ts | 5++++-
Mcommon/jobs/downloadBackoff.test.ts | 43++++++++++++++++++++++++++++++++++++++-----
Mcommon/jobs/downloadBackoff.ts | 79+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcommon/jobs/platformBackoff.test.ts | 6+++---
Mcommon/jobs/platformBackoff.ts | 93++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcommon/jobs/unitOutcome.test.ts | 63++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Mcommon/jobs/unitOutcome.ts | 23++++++++++++++++++++++-
Mcommon/views/autoQueueStatus.test.ts | 16+++++++++++++---
Mcommon/views/autoQueueStatus.ts | 18+++++++++++++-----
Mcommon/ytdlp/channelArgs.test.ts | 2+-
Mcommon/ytdlp/metadataScan.ts | 12+++++++++++-
Mcommon/ytdlp/runYtdlp.ts | 51++++++++++++++++++++++++++++++++++++++++++++++++++-
Meditor/app/channels/[slug]/pipelineActions.ts | 9+++++++--
16 files changed, 388 insertions(+), 53 deletions(-)

diff --git a/common/bin/doctor.test.ts b/common/bin/doctor.test.ts @@ -605,8 +605,8 @@ test("download pacing: a held platform and a raised pace warn, and nothing is wr youtube: { until: now + 42 * 60_000, fails: 8 }, rumble: { until: now + 90_000, fails: 2 }, }, - platformHolds: { youtube: { since: now - 60_000, probeAt: now + 42 * 60_000 } }, - platformPace: { youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 1 } }, + platformHolds: { youtube: { since: now - 60_000, probeAt: now + 42 * 60_000, rateLimited: true } }, + platformPace: { youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 1, steppedAt: now } }, }, }), ); diff --git a/common/bin/doctor.ts b/common/bin/doctor.ts @@ -201,17 +201,22 @@ export async function collectDoctorReport(deps: DoctorDeps): Promise<DoctorRepor if (hold) { quiet = false; add(DP, pf, "warn", - `held since ${stamp(new Date(hold.since))} after ${e.fails} rate-limited/network failures in a row; ` + - `next probe ${e.until > nowMs ? `in ${mins(e.until - nowMs)}` : "due now"} — manual fetches on it are refused until a probe comes back clean`); + `held since ${stamp(new Date(hold.since))} after ${e.fails} ${hold.rateLimited ? "rate-limited/network" : "network"} failures in a row; ` + + (e.until > nowMs + ? `next probe in ${mins(e.until - nowMs)} — manual fetches on it are refused until then` + : "probe overdue — a clean manual Sync lifts it, or Clear hold on /operations/download")); } else if (e.until > nowMs) { quiet = false; add(DP, pf, "warn", `in a rate-limit cooldown for ${mins(e.until - nowMs)} more (attempt ${e.fails})`); } } + const { effectivePaceSeconds } = await import("../jobs/platformBackoff"); for (const [pf, p] of Object.entries(st.platformPace)) { + const v = effectivePaceSeconds(p, nowMs); + if (v <= p.baseSeconds) continue; quiet = false; add(DP, `${pf} pace`, "warn", - `${p.sleepRequestsSeconds}s between requests (base ${p.baseSeconds}s) — a rate limit raised it; ${p.cleanUnits} clean unit(s) toward the next easing`); + `${v}s between requests (base ${p.baseSeconds}s) — a rate limit raised it; it eases one step per hour with no rate limit, and one step per run of clean units (${p.cleanUnits} so far)`); } const subs = Object.values(st.subtitleDeferrals).filter((d) => d.until > nowMs); if (subs.length > 0) { diff --git a/common/jobs/autoQueueState.test.ts b/common/jobs/autoQueueState.test.ts @@ -175,9 +175,9 @@ test("the pacing memory round-trips; a file written before it coerces to {} (rel await withPaths(async (paths) => { const state = emptyAutoQueueState(); state.download.platformPace = { - youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2 }, + youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2, steppedAt: 7 }, }; - state.download.platformHolds = { youtube: { since: 10, probeAt: 20 } }; + state.download.platformHolds = { youtube: { since: 10, probeAt: 20, rateLimited: false } }; state.download.subtitleDeferrals = { v1: { count: 2, lastAt: 5, until: 6, channelSlug: "alpha" }, }; @@ -200,7 +200,7 @@ test("the pacing memory round-trips; a file written before it coerces to {} (rel ); const old = await readAutoQueueState(paths); assert.deepEqual(old.download.platformPace, { - rumble: { sleepRequestsSeconds: 2, baseSeconds: 0, cleanUnits: 0 }, + rumble: { sleepRequestsSeconds: 2, baseSeconds: 0, cleanUnits: 0, steppedAt: 0 }, }); assert.deepEqual(old.download.platformHolds, {}); assert.deepEqual(old.download.subtitleDeferrals, {}); diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts @@ -16,6 +16,7 @@ import { coercePlatformPace, coerceSubtitleDeferrals, coerceVideoDeferrals, + effectivePaceSeconds, } from "./platformBackoff"; // Persistent fairness state for the auto-queue runners. Unlike in-flight worker @@ -248,7 +249,9 @@ function notePace(state: AutoQueueState): void { export function livePlatformPaceSeconds(platform: string): number | undefined { const holder = getHolder(); const pace = holder.state?.download?.platformPace ?? holder.lastPace; - return pace?.[platform]?.sleepRequestsSeconds; + const entry = pace?.[platform]; + // With the hourly easing applied (platformBackoff.ts, PACE_TIME_DECAY_MS). + return entry ? effectivePaceSeconds(entry, Date.now()) : undefined; } // THE PERSISTED STATE IS ONE OBJECT FOR EVERY LANE'S RUNNER, and it has to be. diff --git a/common/jobs/downloadBackoff.test.ts b/common/jobs/downloadBackoff.test.ts @@ -13,6 +13,7 @@ import { recordSubtitleDeferral, } from "./downloadBackoff"; import { PACING_DEFAULTS } from "./platformBackoff"; +import { clearPlatformHold, recordPlatformClean } from "./downloadBackoff"; import { BACKOFF_BASE_MS } from "./platformBackoff"; import { emptyAutoQueueState, @@ -165,11 +166,9 @@ test("a manual rate limit doubles the pace; a network failure does not", async ( await recordDownloadBackoff("youtube", paths, "network", PACING_DEFAULTS); const onDisk = await readAutoQueueState(paths); assert.equal(onDisk.download.platformBackoff.youtube.fails, 2); - assert.deepEqual(onDisk.download.platformPace.youtube, { - sleepRequestsSeconds: 2, - baseSeconds: 1, - cleanUnits: 0, - }); + const { steppedAt, ...pace } = onDisk.download.platformPace.youtube; + assert.deepEqual(pace, { sleepRequestsSeconds: 2, baseSeconds: 1, cleanUnits: 0 }); + assert.ok(Math.abs(steppedAt - Date.now()) < 60_000); }); }); @@ -216,3 +215,37 @@ test("with a runner live, a subtitle deferral is written through the shared obje assert.equal((await readAutoQueueState(paths)).download.subtitleDeferrals.v9.count, 1); }); }); + +test("an overdue hold no longer refuses; a clean manual run lifts it (review H2)", async () => { + await withTempPaths(async (paths) => { + const state = emptyAutoQueueState(); + const now = Date.now(); + state.download.platformBackoff.youtube = { until: now - 1_000, fails: 9 }; + state.download.platformHolds.youtube = { since: now - 3_600_000, probeAt: now - 1_000, rateLimited: true }; + await writeAutoQueueState(paths, state); + assert.equal(await heldPlatformRefusal("youtube", "Sync", paths), null); + const line = await recordPlatformClean("youtube", paths, PACING_DEFAULTS); + assert.match(line ?? "", /^youtube answered cleanly: its hold and backoff are cleared/); + const after = (await readAutoQueueState(paths)).download; + assert.deepEqual(after.platformHolds, {}); + assert.deepEqual(after.platformBackoff, {}); + // Nothing to settle: no line and no write. + assert.equal(await recordPlatformClean("youtube", paths, PACING_DEFAULTS), null); + }); +}); + +test("Clear hold drops the hold, the backoff and the pace, and says so", async () => { + await withTempPaths(async (paths) => { + const state = emptyAutoQueueState(); + const now = Date.now(); + state.download.platformBackoff.youtube = { until: now + 3_600_000, fails: 9 }; + state.download.platformHolds.youtube = { since: now - 60_000, probeAt: now + 3_600_000, rateLimited: true }; + state.download.platformPace.youtube = { sleepRequestsSeconds: 8, baseSeconds: 1, cleanUnits: 0, steppedAt: now }; + await writeAutoQueueState(paths, state); + const line = await clearPlatformHold("youtube", paths); + assert.match(line ?? "", /^Cleared by hand for youtube: the hold \(since .* UTC\), the backoff \(9 failures\), the pace \(8s → base 1s\)\./); + const after = (await readAutoQueueState(paths)).download; + assert.deepEqual([after.platformHolds, after.platformBackoff, after.platformPace], [{}, {}, {}]); + assert.equal(await clearPlatformHold("youtube", paths), null); + }); +}); diff --git a/common/jobs/downloadBackoff.ts b/common/jobs/downloadBackoff.ts @@ -31,12 +31,15 @@ import { type PlatformHoldEntry, type SubtitleDeferral, type SubtitleDeferralState, + clearPlatformPacing, deferSubtitles, + effectivePaceSeconds, escalatePlatform, heldPlatformSentence, mergeBackoffEntry, prunePlatformPacing, pruneSubtitleDeferrals, + settlePlatformClean, } from "./platformBackoff"; import { staticSleepRequestsSeconds } from "../ytdlp/platformArgs.mjs"; @@ -163,20 +166,84 @@ export async function platformHold( return { hold, fails: state.download.platformBackoff[platform]?.fails ?? 0 }; } -// The refusal for a manual Sync / download on a HELD platform, or null when it -// is not held. `what` names the refused thing ("Sync", "This download"). +// The refusal for a manual Sync / download on a HELD platform, or null. Only +// while it is held AND its probe is still in the future (review H2): once the +// probe is due, a manual run IS the probe — refusing it too would leave a +// platform with the lane off, or nothing pending on it, held for ever. `what` +// names the refused thing ("Sync", "This download"). export async function heldPlatformRefusal( platform: string, what: string, paths: Paths = getPaths(), now: number = Date.now(), ): Promise<string | null> { - const held = await platformHold(platform, paths); - if (!held) return null; + const state = + (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); + const hold = state.download.platformHolds[platform]; + const entry = state.download.platformBackoff[platform]; + if (!hold || !entry || entry.until <= now) return null; const probeMinutes = - pacingSettings()?.holdProbeMinutes ?? Math.round((held.hold.probeAt - held.hold.since) / 60_000); - return heldPlatformSentence(platform, held.hold, held.fails, now, { + pacingSettings()?.holdProbeMinutes ?? Math.round((hold.probeAt - hold.since) / 60_000); + return heldPlatformSentence(platform, hold, entry.fails, now, { probeMinutes, what, }); } + +// A MANUAL RUN THAT CAME BACK CLEAN (review H2) — a Sync, a download, a +// metadata scan that the source answered without a rate limit or a network +// failure — settles the platform exactly as a clean lane unit does: the +// backoff and any hold clear, and it counts toward the pace's easing. Returns +// the line for the job's log, or null when there was nothing to settle (and +// then writes nothing). +export async function recordPlatformClean( + platform: string, + paths: Paths = getPaths(), + pacing: PacingSettings | undefined = pacingSettings(), +): Promise<string | null> { + const peek = + (await liveAutoQueueState(paths)) ?? (await readAutoQueueState(paths)); + const d = peek.download; + if (!d.platformBackoff[platform] && !d.platformHolds[platform] && !d.platformPace[platform]) { + return null; + } + return mutateDownloadState(paths, (state) => { + const now = Date.now(); + prunePlatformPacing(state.download, now); + const hadBackoff = Boolean(state.download.platformBackoff[platform]); + const fx = settlePlatformClean(state.download, platform, { + baseSeconds: staticSleepRequestsSeconds(platform), + now, + ...(pacing ? { pacing } : {}), + }); + if (fx.releasedHold) { + return `${platform} answered cleanly: its hold and backoff are cleared (as a clean probe would).\n`; + } + if (hadBackoff) return `${platform} answered cleanly: its rate-limit backoff is cleared.\n`; + if (fx.decayed) return `${platform} pace eased to ${fx.paceSeconds}s between requests.\n`; + return null; + }); +} + +// "Clear hold" (review H2): the operator's word beats the machine, as with an +// auto-pause. The hold, the backoff and the raised pace all go. Returns the +// sentence for the job log, or null when the platform had none of them. +export async function clearPlatformHold( + platform: string, + paths: Paths = getPaths(), +): Promise<string | null> { + return mutateDownloadState(paths, (state) => { + const now = Date.now(); + const pace = state.download.platformPace[platform]; + const was = clearPlatformPacing(state.download, platform); + if (!was.hold && was.fails === 0 && was.paceSeconds === null) return null; + const parts = [ + was.hold + ? `the hold (since ${new Date(was.hold.since).toISOString().slice(0, 16).replace("T", " ")} UTC)` + : null, + was.fails > 0 ? `the backoff (${was.fails} failures)` : null, + pace ? `the pace (${effectivePaceSeconds(pace, now)}s → base ${pace.baseSeconds}s)` : null, + ].filter(Boolean); + return `Cleared by hand for ${platform}: ${parts.join(", ")}. The next failure starts from the bottom of the backoff.\n`; + }); +} diff --git a/common/jobs/platformBackoff.test.ts b/common/jobs/platformBackoff.test.ts @@ -207,8 +207,8 @@ test("a held platform's backoff survives the prune; a hold with no backoff is dr } as PlatformBackoffState, platformPace: {}, platformHolds: { - youtube: { since: 0, probeAt: 0 }, - odysee: { since: 0, probeAt: 0 }, + youtube: { since: 0, probeAt: 0, rateLimited: true }, + odysee: { since: 0, probeAt: 0, rateLimited: true }, }, }; prunePlatformPacing(state, now); @@ -219,7 +219,7 @@ test("a held platform's backoff survives the prune; a hold with no backoff is dr test("the held sentence names the hold and the next probe", () => { const since = Date.UTC(2026, 9, 1, 14, 2); assert.equal( - heldPlatformSentence("youtube", { since, probeAt: since + 42 * 60_000 }, 8, since, { + heldPlatformSentence("youtube", { since, probeAt: since + 42 * 60_000, rateLimited: true }, 8, since, { probeMinutes: 60, what: "Sync", }), diff --git a/common/jobs/platformBackoff.ts b/common/jobs/platformBackoff.ts @@ -229,8 +229,17 @@ export type PlatformPaceEntry = { baseSeconds: number; // Clean units since the last doubling or decay step. cleanUnits: number; + // Epoch ms of the last doubling or decay step. THE ONE TIME-BASED RULE + // (review L2): the pace also eases one step per PACE_TIME_DECAY_MS with no + // rate limit, so a raised pace cannot outlive its cause while the lane is + // off or has nothing pending on the platform. 0 for an entry written before + // the field existed (it eases at once). + steppedAt: number; }; +// One easing step per hour with no rate limit (see `steppedAt`). +export const PACE_TIME_DECAY_MS = 60 * 60_000; + export type PlatformPaceState = Record<string, PlatformPaceEntry>; // A platform in a HOLD. Its backoff entry's `until` is the next probe. @@ -239,6 +248,10 @@ export type PlatformHoldEntry = { since: number; // Epoch ms of the next probe (mirrors the backoff entry's `until`). probeAt: number; + // Whether a RATE LIMIT is among the failures that held it — false for a + // hold reached through network failures alone, which is worded "failing", + // not "rate-limited" (review L5). + rateLimited: boolean; }; export type PlatformHoldState = Record<string, PlatformHoldEntry>; @@ -335,6 +348,7 @@ export function escalatePlatform( sleepRequestsSeconds: Math.min(cap, Math.max(1, cur * 2)), baseSeconds: opts.baseSeconds, cleanUnits: 0, + steppedAt: now, }; } const wasHeld = isPlatformHeld(state.platformHolds, platform); @@ -344,6 +358,9 @@ export function escalatePlatform( state.platformHolds[platform] = { since: state.platformHolds[platform]?.since ?? now, probeAt: entry.until, + rateLimited: + (state.platformHolds[platform]?.rateLimited ?? false) || + failureClass === "rate_limit", }; } state.platformBackoff[platform] = entry; @@ -369,7 +386,14 @@ export type CleanEffect = { export function settlePlatformClean( state: PlatformPacingState, platform: string, - opts: { pacing?: PacingSettings; baseSeconds: number }, + opts: { + pacing?: PacingSettings; + baseSeconds: number; + // False for a probe whose media came down but whose subtitles 429'd: it + // lifts a hold, and does not count toward the pace's easing (review L1). + countForDecay?: boolean; + now?: number; + }, ): CleanEffect { const pacing = opts.pacing ?? PACING_DEFAULTS; const releasedHold = isPlatformHeld(state.platformHolds, platform); @@ -377,13 +401,19 @@ export function settlePlatformClean( delete state.platformHolds[platform]; const pace = state.platformPace[platform]; let decayed = false; - if (pace) { + if (pace && opts.countForDecay !== false) { pace.cleanUnits += 1; if (pace.cleanUnits >= Math.max(1, pacing.decayAfterCleanUnits)) { decayed = true; const next = pace.sleepRequestsSeconds / 2; if (next <= pace.baseSeconds) delete state.platformPace[platform]; - else state.platformPace[platform] = { ...pace, sleepRequestsSeconds: next, cleanUnits: 0 }; + else + state.platformPace[platform] = { + ...pace, + sleepRequestsSeconds: next, + cleanUnits: 0, + steppedAt: opts.now ?? pace.steppedAt, + }; } } return { @@ -398,6 +428,7 @@ export function settlePlatformClean( // the escalation memory that keeps it a hold. A hold whose backoff entry is // gone (a hand-edited file) is dropped with it. Mutates. export function prunePlatformPacing(state: PlatformPacingState, now: number): void { + decayPaceByTime(state.platformPace, now); const before = { ...state.platformBackoff }; pruneExpired(state.platformBackoff, now); for (const pf of Object.keys(state.platformHolds)) { @@ -406,6 +437,48 @@ export function prunePlatformPacing(state: PlatformPacingState, now: number): vo } } +// The pace an entry stands at NOW once the hourly easing is applied: one +// halving per whole PACE_TIME_DECAY_MS since `steppedAt`, never below its base. +// Pure — the readers (the args builder, the view, the doctor) use it so a +// stale entry on disk never paces a spawn above what the rule allows. +export function effectivePaceSeconds(entry: PlatformPaceEntry, now: number): number { + const steps = Math.max(0, Math.floor((now - entry.steppedAt) / PACE_TIME_DECAY_MS)); + let v = entry.sleepRequestsSeconds; + for (let i = 0; i < steps && v > entry.baseSeconds; i++) v = v / 2; + return Math.max(v, entry.baseSeconds); +} + +// Apply the hourly easing to the stored map (an entry at its base is +// deleted). Mutates. Called where the state is pruned and on every unit. +export function decayPaceByTime(pace: PlatformPaceState, now: number): void { + for (const [pf, e] of Object.entries(pace)) { + const steps = Math.max(0, Math.floor((now - e.steppedAt) / PACE_TIME_DECAY_MS)); + if (steps === 0) continue; + const v = effectivePaceSeconds(e, now); + if (v <= e.baseSeconds) delete pace[pf]; + else pace[pf] = { ...e, sleepRequestsSeconds: v, steppedAt: e.steppedAt + steps * PACE_TIME_DECAY_MS }; + } +} + +// THE OPERATOR'S WORD (review H2): "Clear hold" on the lane page. The hold, +// the backoff and the raised pace all go, as if the platform had never been +// refused; the next failure starts the escalation from the bottom. Mutates; +// returns what there was, for the job log. +export function clearPlatformPacing( + state: PlatformPacingState, + platform: string, +): { hold: PlatformHoldEntry | null; fails: number; paceSeconds: number | null } { + const out = { + hold: state.platformHolds[platform] ?? null, + fails: state.platformBackoff[platform]?.fails ?? 0, + paceSeconds: state.platformPace[platform]?.sleepRequestsSeconds ?? null, + }; + delete state.platformHolds[platform]; + clearBackoff(state.platformBackoff, platform); + delete state.platformPace[platform]; + return out; +} + // Record one subtitle 429 against a video: the count climbs, and the batch // subtitle fetch leaves it alone for 6 h — or for 7 days from the third. A // manual fetch never reads this. Mutates; returns the entry. @@ -458,6 +531,7 @@ export function coercePlatformPace(value: unknown): PlatformPaceState { sleepRequestsSeconds: r.sleepRequestsSeconds, baseSeconds: finite(r.baseSeconds) && r.baseSeconds >= 0 ? r.baseSeconds : 0, cleanUnits: finite(r.cleanUnits) && r.cleanUnits >= 0 ? Math.floor(r.cleanUnits) : 0, + steppedAt: finite(r.steppedAt) ? r.steppedAt : 0, }; } } @@ -470,7 +544,13 @@ export function coercePlatformHolds(value: unknown): PlatformHoldState { for (const [pf, raw] of Object.entries(value as Record<string, unknown>)) { if (!raw || typeof raw !== "object") continue; const r = raw as Record<string, unknown>; - if (finite(r.since) && finite(r.probeAt)) out[pf] = { since: r.since, probeAt: r.probeAt }; + if (finite(r.since) && finite(r.probeAt)) { + out[pf] = { + since: r.since, + probeAt: r.probeAt, + rateLimited: typeof r.rateLimited === "boolean" ? r.rateLimited : true, + }; + } } return out; } @@ -519,8 +599,11 @@ export function heldPlatformSentence( opts: { probeMinutes: number; what?: string }, ): string { const what = opts.what ?? "This"; + const why = hold.rateLimited + ? "its rate limit outlasted the cooldown cap" + : "it kept failing (network errors) past the cooldown cap"; return ( - `${platform} is held: its rate limit outlasted the cooldown cap ` + + `${platform} is held: ${why} ` + `(${fails} failures in a row, held since ${clockOf(hold.since)}). ` + `Auto-download probes it once every ${opts.probeMinutes} min — the next probe ` + `is in ${minutesText(hold.probeAt - now)}. ${what} will run once a probe comes back clean.` diff --git a/common/jobs/unitOutcome.test.ts b/common/jobs/unitOutcome.test.ts @@ -4,8 +4,11 @@ import { BACKOFF_BASE_MS, FAILS_TO_REACH_CAP, PACING_DEFAULTS, + PACE_TIME_DECAY_MS, SUBTITLE_HOLD_MS, VIDEO_RATE_LIMIT_DEFER_MS, + effectivePaceSeconds, + heldPlatformSentence, isVideoDeferred, } from "./platformBackoff"; import { type UnitOutcomeState, applyUnitOutcome } from "./unitOutcome"; @@ -118,7 +121,7 @@ test("fails climbs only across DISTINCT ids: the deferred video is not re-picked test("subs_rate_limit retires the video, defers its subtitles, and touches no platform state", () => { const s = empty(); s.platformBackoff.youtube = { until: 10, fails: 3 }; - s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2 }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 2, steppedAt: 1_000_000 }; const before = structuredClone({ platformBackoff: s.platformBackoff, platformPace: s.platformPace, @@ -189,13 +192,13 @@ test("the pace doubles on each rate limit, never on network, and stops at the ca assert.equal(s.platformPace.youtube.sleepRequestsSeconds, 8); // A platform with no static pace starts from nothing and goes to 1 s. applyUnitOutcome(s, unit("o1", { platform: "odysee", failureClass: "rate_limit" }), now, noJitter, pacing); - assert.deepEqual(s.platformPace.odysee, { sleepRequestsSeconds: 1, baseSeconds: 0, cleanUnits: 0 }); + assert.deepEqual(s.platformPace.odysee, { sleepRequestsSeconds: 1, baseSeconds: 0, cleanUnits: 0, steppedAt: now }); }); test("the pace eases one step per N clean units, and an entry back at its base is gone", () => { const s = empty(); const pacing = { ...PACING_DEFAULTS, decayAfterCleanUnits: 2 }; - s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0 }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0, steppedAt: 0 }; const lines: (string | null)[] = []; for (let i = 0; i < 4; i++) { lines.push(applyUnitOutcome(s, unit(`c${i}`, { outcome: "transcribed" }), 0, noJitter, pacing).line); @@ -208,7 +211,7 @@ test("the pace eases one step per N clean units, and an entry back at its base i ]); assert.deepEqual(s.platformPace, {}); // A subtitle-429 unit is not clean: it does not count toward the decay. - s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 1 }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 1, steppedAt: 0 }; applyUnitOutcome(s, unit("x", { outcome: "transcribed", failureClass: "subs_rate_limit" }), 0, noJitter, pacing); assert.equal(s.platformPace.youtube.cleanUnits, 1); }); @@ -227,14 +230,14 @@ test("a backoff failing at the cap N times holds the platform: one probe per hol assert.equal(s.platformBackoff.youtube.fails, FAILS_TO_REACH_CAP + 1); // The third at the cap: held, and the next try is a probe an hour away. const r = applyUnitOutcome(s, unit("h1", { failureClass: "rate_limit" }), now, noJitter, pacing); - assert.deepEqual(s.platformHolds.youtube, { since: now, probeAt: now + 60 * 60_000 }); + assert.deepEqual(s.platformHolds.youtube, { since: now, probeAt: now + 60 * 60_000, rateLimited: true }); assert.equal(s.platformBackoff.youtube.until, now + 60 * 60_000); assert.match(r.line ?? "", /^Auto-download: youtube rate_limit — held after 8 failures in a row; next probe in 3600s\./); // A failed probe keeps the hold (and its `since`) and re-arms the hour. const since = now; now = s.platformBackoff.youtube.until; const r2 = applyUnitOutcome(s, unit("h2", { failureClass: "network" }), now, noJitter, pacing); - assert.deepEqual(s.platformHolds.youtube, { since, probeAt: now + 60 * 60_000 }); + assert.deepEqual(s.platformHolds.youtube, { since, probeAt: now + 60 * 60_000, rateLimited: true }); assert.match(r2.line ?? "", /still held/); // A clean probe clears the hold AND the backoff. const r3 = applyUnitOutcome(s, unit("h3", { outcome: "transcribed" }), now + 60 * 60_000, noJitter, pacing); @@ -246,16 +249,54 @@ test("a backoff failing at the cap N times holds the platform: one probe per hol ); }); -test("a subtitle 429 on a probe neither lifts nor deepens the hold", () => { +test("a probe whose media came down but whose subtitles 429'd lifts the hold, without easing the pace (review L1)", () => { const s = empty(); s.platformBackoff.youtube = { until: 100, fails: 9 }; - s.platformHolds.youtube = { since: 0, probeAt: 100 }; - applyUnitOutcome( + s.platformHolds.youtube = { since: 0, probeAt: 100, rateLimited: true }; + s.platformPace.youtube = { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0, steppedAt: 100 }; + const r = applyUnitOutcome( s, unit("p", { outcome: "transcribed", failureClass: "subs_rate_limit" }), 100, noJitter, ); - assert.deepEqual(s.platformHolds.youtube, { since: 0, probeAt: 100 }); - assert.deepEqual(s.platformBackoff.youtube, { until: 100, fails: 9 }); + assert.deepEqual(s.platformHolds, {}); + assert.equal(s.platformBackoff.youtube, undefined); + assert.equal(s.platformPace.youtube.cleanUnits, 0); + assert.equal(s.subtitleDeferrals.p.count, 1); + assert.match(r.line ?? "", /The youtube probe's media came down: the hold and the backoff are cleared\.$/); +}); + +test("the pace eases one step per hour with no rate limit (the one time-based rule, review L2)", () => { + const s = empty(); + const t0 = 10 * PACE_TIME_DECAY_MS; + s.platformPace.youtube = { sleepRequestsSeconds: 8, baseSeconds: 1, cleanUnits: 0, steppedAt: t0 }; + applyUnitOutcome(s, unit("x", { outcome: "skipped" }), t0 + PACE_TIME_DECAY_MS - 1, noJitter); + assert.equal(s.platformPace.youtube.sleepRequestsSeconds, 8); + applyUnitOutcome(s, unit("x", { outcome: "skipped" }), t0 + 2 * PACE_TIME_DECAY_MS, noJitter); + assert.deepEqual(s.platformPace.youtube, { + sleepRequestsSeconds: 2, + baseSeconds: 1, + cleanUnits: 0, + steppedAt: t0 + 2 * PACE_TIME_DECAY_MS, + }); + // The readers see the same easing without a write. + assert.equal(effectivePaceSeconds(s.platformPace.youtube, t0 + 3 * PACE_TIME_DECAY_MS), 1); + applyUnitOutcome(s, unit("x", { outcome: "skipped" }), t0 + 9 * PACE_TIME_DECAY_MS, noJitter); + assert.deepEqual(s.platformPace, {}); +}); + +test("a hold reached through network failures alone is worded as failing, not rate-limited (review L5)", () => { + const s = empty(); + const pacing = { ...PACING_DEFAULTS, holdAfterFailsAtCap: 1 }; + s.platformBackoff.youtube = { until: 0, fails: 5 }; + applyUnitOutcome(s, unit("n", { failureClass: "network" }), 10, noJitter, pacing); + assert.equal(s.platformHolds.youtube.rateLimited, false); + assert.match( + heldPlatformSentence("youtube", s.platformHolds.youtube, 6, 10, { probeMinutes: 60 }), + /^youtube is held: it kept failing \(network errors\) past the cooldown cap/, + ); + // A rate limit while held makes it a rate-limited hold. + applyUnitOutcome(s, unit("r", { failureClass: "rate_limit" }), 20, noJitter, pacing); + assert.equal(s.platformHolds.youtube.rateLimited, true); }); diff --git a/common/jobs/unitOutcome.ts b/common/jobs/unitOutcome.ts @@ -24,7 +24,9 @@ // its next try is a probe `pacing.holdProbeMinutes` away. // - `subs_rate_limit` (the subtitle fetch alone answered 429 and the media came // down) retires the video like any success and defers its SUBTITLES — the -// platform's backoff, hold and pace are not touched, either way. +// platform's backoff and pace are not touched; on a HELD platform it is the +// clean probe that lifts the hold (review L1). +// - the pace eases one step per hour with no rate limit (decayPaceByTime). // - a clean unit (`transcribed` with no subtitle 429) clears the backoff and // any hold — a clean probe is what lifts one — and counts toward the pace's // decay. @@ -38,8 +40,10 @@ import { type SubtitleDeferralState, type VideoDeferralState, SUBTITLE_HOLD_AFTER, + decayPaceByTime, deferSubtitles, deferVideo, + isPlatformHeld, escalatePlatform, pruneDeferred, pruneSubtitleDeferrals, @@ -81,6 +85,7 @@ export function applyUnitOutcome( ): UnitOutcomeEffect { pruneDeferred(state.videoDeferrals, now); pruneSubtitleDeferrals(state.subtitleDeferrals, now); + decayPaceByTime(state.platformPace, now); const pf = unit.platform; if ( pf !== null && @@ -118,6 +123,21 @@ export function applyUnitOutcome( d.count >= SUBTITLE_HOLD_AFTER ? "left alone for 7 days" : "deferred 6h for download-missing-subs"; + // A HELD platform's probe whose media came down answered on every + // platform-level request: that is the clean probe the hold waits for + // (review L1). It does not count toward the pace's easing. + if (pf !== null && isPlatformHeld(state.platformHolds, pf)) { + settlePlatformClean(state, pf, { + baseSeconds: staticSleepRequestsSeconds(pf), + countForDecay: false, + now, + ...(pacing ? { pacing } : {}), + }); + return { + markCompleted: true, + line: `Auto-download: ${unit.videoId} downloaded; its subtitles were rate-limited (${d.count}×) — ${left}. The ${pf} probe's media came down: the hold and the backoff are cleared.`, + }; + } return { markCompleted: true, line: `Auto-download: ${unit.videoId} downloaded; its subtitles were rate-limited (${d.count}×) — ${left}. ${pf ?? "the platform"} is not backed off.`, @@ -126,6 +146,7 @@ export function applyUnitOutcome( if (pf !== null && unit.outcome === "transcribed") { const fx = settlePlatformClean(state, pf, { baseSeconds: staticSleepRequestsSeconds(pf), + now, ...(pacing ? { pacing } : {}), }); if (fx.releasedHold) { diff --git a/common/views/autoQueueStatus.test.ts b/common/views/autoQueueStatus.test.ts @@ -103,9 +103,9 @@ test("a cooldown is filtered by the injected clock, newest first", () => { test("a hold, the pace and the subtitle deferrals reach the download lane (release 17, RL)", () => { const state = emptyAutoQueueState(); state.download.platformBackoff = { youtube: { until: NOW + 3_600_000, fails: 8 } }; - state.download.platformHolds = { youtube: { since: NOW - 60_000, probeAt: NOW + 3_600_000 } }; + state.download.platformHolds = { youtube: { since: NOW - 60_000, probeAt: NOW + 3_600_000, rateLimited: true } }; state.download.platformPace = { - youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0 }, + youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0, steppedAt: NOW }, }; state.download.subtitleDeferrals = { b: { count: 3, lastAt: NOW - 1, until: NOW + 7 * 86_400_000, channelSlug: "alpha" }, @@ -114,7 +114,7 @@ test("a hold, the pace and the subtitle deferrals reach the download lane (relea }; const d = buildAutoQueueStatusPayload(inputs({ state })).download; assert.deepEqual(d.cooldowns, [ - { platform: "youtube", untilMs: NOW + 3_600_000, fails: 8, hold: { sinceMs: NOW - 60_000 } }, + { platform: "youtube", untilMs: NOW + 3_600_000, fails: 8, hold: { sinceMs: NOW - 60_000, rateLimited: true } }, ]); assert.deepEqual(d.pace, [{ platform: "youtube", sleepRequestsSeconds: 4, baseSeconds: 1 }]); assert.deepEqual( @@ -329,3 +329,13 @@ test("memo: a different key misses the memo and the computation in flight", asyn assert.equal(await a, "late A2"); assert.equal(await memo.get("tree C", async () => "unused"), "C"); }); + +test("a held platform stays listed after its probe time passed (review H2)", () => { + const state = emptyAutoQueueState(); + state.download.platformBackoff = { youtube: { until: NOW - 5_000, fails: 9 }, rumble: { until: NOW - 5_000, fails: 2 } }; + state.download.platformHolds = { youtube: { since: NOW - 7_200_000, probeAt: NOW - 5_000, rateLimited: false } }; + const d = buildAutoQueueStatusPayload(inputs({ state })).download; + assert.deepEqual(d.cooldowns, [ + { platform: "youtube", untilMs: NOW - 5_000, fails: 9, hold: { sinceMs: NOW - 7_200_000, rateLimited: false } }, + ]); +}); diff --git a/common/views/autoQueueStatus.ts b/common/views/autoQueueStatus.ts @@ -12,6 +12,7 @@ import type { AutoQueueState, } from "../jobs/autoQueueState"; import { LANES } from "../lib/autoQueueTypes"; +import { effectivePaceSeconds } from "../jobs/platformBackoff"; import type { AutoQueuePolicy } from "../jobs/autoQueuePolicy"; import { isGateHeld } from "../lib/pauseGates"; import type { FocusSummary } from "../lib/channelPriority"; @@ -31,9 +32,13 @@ export type PlatformCooldownView = { platform: string; untilMs: number; fails: number; - // The platform's HOLD (release 17, slice RL): its rate limit outlasted the + // The platform's HOLD (release 17, slice RL): its failures outlasted the // cooldown cap, and `untilMs` is the next probe. Null for a plain cooldown. - hold: { sinceMs: number } | null; + // A held platform is listed WHATEVER its probe time (review H2): once the + // probe is overdue (`untilMs` passed) it stays held until a clean probe, a + // clean manual run or Clear hold. `rateLimited` is false for a hold reached + // through network failures alone. + hold: { sinceMs: number; rateLimited: boolean } | null; }; // A platform whose request pace is above its static value (release 17, slice @@ -178,25 +183,28 @@ function buildKind( const cooldowns: PlatformCooldownView[] = Object.entries( inputs.state[kind].platformBackoff, ) - .filter(([, e]) => e.until > now) + .filter(([platform, e]) => e.until > now || Boolean(inputs.state[kind].platformHolds?.[platform])) .map(([platform, e]) => { const hold = inputs.state[kind].platformHolds?.[platform]; return { platform, untilMs: e.until, fails: e.fails, - hold: hold ? { sinceMs: hold.since } : null, + hold: hold ? { sinceMs: hold.since, rateLimited: hold.rateLimited } : null, }; }) .sort((a, b) => b.untilMs - a.untilMs); const pace: PlatformPaceView[] = Object.entries( inputs.state[kind].platformPace ?? {}, ) + // At the pace the hourly easing leaves it (effectivePaceSeconds); an + // entry already back at its base is not listed. .map(([platform, p]) => ({ platform, - sleepRequestsSeconds: p.sleepRequestsSeconds, + sleepRequestsSeconds: effectivePaceSeconds(p, now), baseSeconds: p.baseSeconds, })) + .filter((p) => p.sleepRequestsSeconds > p.baseSeconds) .sort((a, b) => (a.platform < b.platform ? -1 : a.platform > b.platform ? 1 : 0)); const subtitleDeferred: SubtitleDeferralView[] = Object.entries( inputs.state[kind].subtitleDeferrals ?? {}, diff --git a/common/ytdlp/channelArgs.test.ts b/common/ytdlp/channelArgs.test.ts @@ -133,7 +133,7 @@ test("the live pace (the shared state) reaches every channelExtraArgs call", () stateFile: null, loading: null, loadingFile: null, - lastPace: { youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0 } }, + lastPace: { youtube: { sleepRequestsSeconds: 4, baseSeconds: 1, cleanUnits: 0, steppedAt: Date.now() } }, }; try { assert.deepEqual( diff --git a/common/ytdlp/metadataScan.ts b/common/ytdlp/metadataScan.ts @@ -31,7 +31,8 @@ import { authRetryCookies, resolveCookiePolicy, } from "../lib/cookiePolicy"; -import { channelExtraArgs } from "./channelArgs"; +import { channelExtraArgs, pacingPlatformKey } from "./channelArgs"; +import { recordPlatformClean } from "../jobs/downloadBackoff"; import type { Paths } from "../lib/paths"; import { getSettings } from "../lib/settings"; import { extractVideoId } from "../lib/videoId"; @@ -683,6 +684,15 @@ export async function runMetadataScan( `${scannedTotal} video(s) were read and saved; the rest stay unscanned. A per-platform cooldown has been recorded — run the scan again once it lapses.\n`, ); await opts.onPlatformBackoff?.("rate_limit"); + } else if (stopped === undefined && scannedTotal > 0 && !opts.signal.aborted) { + // A scan the source answered cleanly settles the platform like a clean + // probe — its backoff and any hold clear (release 17 review H2). + try { + const line = await recordPlatformClean(pacingPlatformKey(channelConfig), paths); + if (line) opts.onLog(line); + } catch { + /* shared-state write is best-effort */ + } } const errorCount = Object.keys(errors).length; diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -179,6 +179,12 @@ export type RunYtdlpOpts = { onPlatformBackoff?: ( failureClass: "rate_limit" | "network", ) => void | Promise<void>; + // Called once when the run finished — not thrown, not cancelled — with no + // platform backoff recorded during it: the source answered cleanly, so a + // manual run settles the platform the way a clean lane probe does (its + // backoff and any hold clear; release 17 review H2). Returns a line for the + // log, or null. + onPlatformClean?: () => Promise<string | null>; }; // The batch modes whose child writes media into `data/<id>/` — yt-dlp's own @@ -222,7 +228,30 @@ export async function runYtdlp(opts: RunYtdlpOpts): Promise<void> { } } +// A CLEAN RUN SETTLES THE PLATFORM (release 17 review H2): with no rate limit +// or network failure recorded while it ran, the run is the probe a held +// platform waits for — the lane may be off, or have nothing pending on it. async function runYtdlpMode(opts: RunYtdlpOpts): Promise<void> { + if (!opts.onPlatformClean) return runYtdlpModeInner(opts); + let backedOff = false; + const inner = opts.onPlatformBackoff; + await runYtdlpModeInner({ + ...opts, + onPlatformBackoff: async (failureClass) => { + backedOff = true; + await inner?.(failureClass); + }, + }); + if (backedOff || opts.signal.aborted) return; + try { + const line = await opts.onPlatformClean(); + if (line) opts.onLog(line); + } catch { + /* shared-state write is best-effort */ + } +} + +async function runYtdlpModeInner(opts: RunYtdlpOpts): Promise<void> { switch (opts.mode) { case "store-playlist": await storePlaylist(opts); @@ -1416,11 +1445,31 @@ async function downloadMissingSubs(opts: RunYtdlpOpts): Promise<void> { let processed = 0; let failed = 0; let firstFailure: Error | null = null; + // THE GAP BETWEEN TWO VIDEOS (release 17 review L3): the same + // sleepBetweenDownloadsSeconds plus the pace above its base the lane and the + // batch downloads wait, so a channel whose subtitles are being refused is + // not asked once per video back to back. + const subsSleepSeconds = + opts.channelConfig.sleepBetweenDownloadsSeconds ?? + getSettings().sleepBetweenDownloadsSeconds; + const subsBasePace = staticSleepRequestsSeconds(channelPlatform(opts.channelConfig)); await Promise.all( - tofetch.map((url) => + tofetch.map((url, index) => limit(async () => { if (opts.signal.aborted || opts.drainSignal?.aborted) return; if (firstFailure) return; + if (index > 0) { + const gapMs = downloadGapMs( + subsSleepSeconds, + channelPaceSeconds(opts.channelConfig), + subsBasePace, + ); + if (gapMs > 0) { + opts.onLog(`Sleeping ${gapMs / 1000}s before the next video...\n`); + await abortableSleep(gapMs, opts.signal); + if (opts.signal.aborted || opts.drainSignal?.aborted) return; + } + } const task = opts.tracker?.start({ id: extractVideoId(url) ?? url, label: extractVideoId(url) ?? url, diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -17,6 +17,7 @@ import { heldPlatformRefusal, platformCooldownRemainingMs, recordDownloadBackoff, + recordPlatformClean, } from "yt-dlp-transcript-common/jobs/downloadBackoff"; import { countNotYetDownloaded, @@ -128,8 +129,9 @@ async function runPipelineAction( // "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): its rate limit outlasted the cooldown cap, and the lane is - // probing it once an hour — a click is not a probe. The sentence names the + // 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") { @@ -229,6 +231,9 @@ async function runPipelineAction( // 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