Archilyzer · Source

archilyzer

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

commit 715d3ef8340b25cb7934063039556516a6722b80
parent 6eaaca3c45198b27645fd618002538259ca86092
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Wed,  5 Aug 2026 00:02:28 -0400

Let one playlist fetch drive sync, missing, and deletion detection

A channel's listing was fetched by three separate yt-dlp spawns that never
shared their work: store-playlist refreshed the `playlist` file, sync walked
50-entry pages, and the quick availability check enumerated the whole listing
to find deletions. So sync never learned about deleted videos, and never
refreshed `playlist` — "download missing" and the snapshot's undownloadedIds
kept working off whatever store-playlist last wrote, possibly months ago.

A sync now upgrades itself to a FULL SWEEP once per cadence (default daily,
per-channel overridable): one full enumeration, from which all three outputs
are derived. "Sync all" therefore surfaces upstream deletions across every
due channel with no extra clicks and no extra spawns.

The download set is deliberately unchanged. The sweep still selects
newest-first entries up to the first page containing an already-archived
entry — a faithful port of the paged walk's stopping rule applied to a list
already in memory. Downloading everything absent from the archive would drag
down the entire back catalogue on a deliberately-partial channel, including
audio the cleanup sweep has reclaimed. The sweep changes what we know, never
what we fetch.

Suspects are resolved in-line with the per-video availability probe only
while there are few enough to be worth doing unattended (default 25).

maybe-missing's record moves to a dependency-free leaf module so both
producers can write it: quickAvailabilityCheck imports runYtdlp, so runYtdlp
could not have imported back from it.

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

Diffstat:
Acommon/controller/maybeMissingStore.ts | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/quickAvailabilityCheck.ts | 62++++++++++++++------------------------------------------------
Acommon/jobs/deepSync.test.ts | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/deepSync.ts | 43+++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/channelConfig.ts | 34++++++++++++++++++++++++++++++++++
Mcommon/lib/settings.ts | 48++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/ytdlp/runYtdlp.ts | 312+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Meditor/app/channels/[slug]/pipelineActions.ts | 12+++++++++++-
Meditor/app/channels/actions.ts | 10++++++++--
Meditor/app/jobs/jobReplayRegistry.ts | 4++--
10 files changed, 602 insertions(+), 84 deletions(-)

diff --git a/common/controller/maybeMissingStore.ts b/common/controller/maybeMissingStore.ts @@ -0,0 +1,68 @@ +import path from "node:path"; +import { readFile, rename, writeFile } from "node:fs/promises"; +import type { Paths } from "../lib/paths"; + +// Leaf module for the maybe-missing record: the set of videos we have on disk +// that no longer appear in the channel's fresh listing (deleted, made private, +// or unlisted). Deliberately dependency-free apart from the Paths type so BOTH +// producers can import it — the quick availability check (controller/) and the +// sync full sweep (ytdlp/runYtdlp.ts). quickAvailabilityCheck imports runYtdlp, +// so runYtdlp cannot import back from it; this module breaks that cycle. + +export const MAYBE_MISSING_FILENAME = "maybe-missing.json"; + +export type MaybeMissingRecord = { + checkedAt: string; + freshPlaylistCount: number; + ids: string[]; +}; + +function maybeMissingPath(paths: Paths, slug: string): string { + return path.join(paths.channelsDir, slug, MAYBE_MISSING_FILENAME); +} + +export async function loadMaybeMissing( + paths: Paths, + slug: string, +): Promise<MaybeMissingRecord | null> { + try { + const raw = await readFile(maybeMissingPath(paths, slug), "utf8"); + const parsed = JSON.parse(raw) as Partial<MaybeMissingRecord>; + if ( + typeof parsed?.checkedAt === "string" && + Array.isArray(parsed.ids) && + typeof parsed.freshPlaylistCount === "number" + ) { + return { + checkedAt: parsed.checkedAt, + freshPlaylistCount: parsed.freshPlaylistCount, + ids: parsed.ids.filter((id): id is string => typeof id === "string"), + }; + } + return null; + } catch { + return null; + } +} + +export async function writeMaybeMissing( + paths: Paths, + slug: string, + record: MaybeMissingRecord, +): Promise<void> { + const file = maybeMissingPath(paths, slug); + const tmp = `${file}.tmp-${process.pid}`; + await writeFile(tmp, JSON.stringify(record, null, 2) + "\n"); + await rename(tmp, file); +} + +// Diff a fresh listing's canonical ids against the videos we already have on +// disk. Anything known but absent from the listing is "maybe missing". Shared +// by the quick availability check and the sync full sweep so both flag the same +// set by the same rule. +export function diffMaybeMissing( + knownIds: ReadonlyArray<string>, + freshIds: ReadonlySet<string>, +): string[] { + return knownIds.filter((id) => !freshIds.has(id)).sort(); +} diff --git a/common/controller/quickAvailabilityCheck.ts b/common/controller/quickAvailabilityCheck.ts @@ -1,55 +1,21 @@ import path from "node:path"; -import { readdir, readFile, rename, writeFile } from "node:fs/promises"; +import { readdir } from "node:fs/promises"; import { readChannelConfig } from "./channels"; import { extractVideoId, fetchFlatPlaylistUrls } from "../ytdlp/runYtdlp"; +import { diffMaybeMissing, writeMaybeMissing } from "./maybeMissingStore"; import type { Paths } from "../lib/paths"; -export const MAYBE_MISSING_FILENAME = "maybe-missing.json"; - -export type MaybeMissingRecord = { - checkedAt: string; - freshPlaylistCount: number; - ids: string[]; -}; - -function maybeMissingPath(paths: Paths, slug: string): string { - return path.join(paths.channelsDir, slug, MAYBE_MISSING_FILENAME); -} - -export async function loadMaybeMissing( - paths: Paths, - slug: string, -): Promise<MaybeMissingRecord | null> { - try { - const raw = await readFile(maybeMissingPath(paths, slug), "utf8"); - const parsed = JSON.parse(raw) as Partial<MaybeMissingRecord>; - if ( - typeof parsed?.checkedAt === "string" && - Array.isArray(parsed.ids) && - typeof parsed.freshPlaylistCount === "number" - ) { - return { - checkedAt: parsed.checkedAt, - freshPlaylistCount: parsed.freshPlaylistCount, - ids: parsed.ids.filter((id): id is string => typeof id === "string"), - }; - } - return null; - } catch { - return null; - } -} - -export async function writeMaybeMissing( - paths: Paths, - slug: string, - record: MaybeMissingRecord, -): Promise<void> { - const file = maybeMissingPath(paths, slug); - const tmp = `${file}.tmp-${process.pid}`; - await writeFile(tmp, JSON.stringify(record, null, 2) + "\n"); - await rename(tmp, file); -} +// The record itself lives in the dependency-free ./maybeMissingStore so +// runYtdlp's full sweep can write it too (this module imports runYtdlp, so the +// dependency can't run the other way). Re-exported here because every existing +// caller — availabilityActions, verifyBeforeClean, channelSnapshot, buildIndex — +// already imports it from this path. +export { + MAYBE_MISSING_FILENAME, + loadMaybeMissing, + writeMaybeMissing, + type MaybeMissingRecord, +} from "./maybeMissingStore"; export type QuickAvailabilityCheckResult = { knownCount: number; @@ -109,7 +75,7 @@ export async function runQuickAvailabilityCheck({ .filter((e) => e.isDirectory()) .map((e) => e.name); - const maybeMissing = knownIds.filter((id) => !freshIds.has(id)).sort(); + const maybeMissing = diffMaybeMissing(knownIds, freshIds); await writeMaybeMissing(paths, channelSlug, { checkedAt: new Date().toISOString(), diff --git a/common/jobs/deepSync.test.ts b/common/jobs/deepSync.test.ts @@ -0,0 +1,93 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { ChannelConfig } from "../lib/channelConfig"; +import { defaultSyncScheduler, type SyncSchedulerSettings } from "../lib/settings"; +import { isFullSweepDue, resolveFullSweepIntervalMinutes } from "./deepSync"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/deepSync.test.ts + +function channel(over: Partial<ChannelConfig> = {}): ChannelConfig { + return { handling: "transcribe", url: "https://example.com/c", ...over }; +} + +function scheduler(over: Partial<SyncSchedulerSettings> = {}): SyncSchedulerSettings { + return { ...defaultSyncScheduler(), ...over }; +} + +const NOW = Date.parse("2026-08-04T12:00:00.000Z"); +const MIN = 60_000; + +test("an unset channel interval inherits the global default", () => { + assert.equal( + resolveFullSweepIntervalMinutes( + channel(), + scheduler({ fullSweepIntervalMinutes: 720 }), + ), + 720, + ); +}); + +test("a channel override beats the global default", () => { + assert.equal( + resolveFullSweepIntervalMinutes( + channel({ fullSweepIntervalMinutes: 10080 }), + scheduler({ fullSweepIntervalMinutes: 1440 }), + ), + 10080, + ); +}); + +test("a channel override of 0 beats a nonzero global (explicit off)", () => { + assert.equal( + resolveFullSweepIntervalMinutes( + channel({ fullSweepIntervalMinutes: 0 }), + scheduler({ fullSweepIntervalMinutes: 1440 }), + ), + 0, + ); +}); + +test("a never-swept channel is due immediately", () => { + assert.equal(isFullSweepDue(undefined, 1440, NOW), true); +}); + +test("an unparseable timestamp is treated as never swept", () => { + assert.equal(isFullSweepDue("not a date", 1440, NOW), true); +}); + +test("interval 0 is never due, even having never swept", () => { + assert.equal(isFullSweepDue(undefined, 0, NOW), false); + assert.equal( + isFullSweepDue(new Date(NOW - 365 * 24 * 60 * MIN).toISOString(), 0, NOW), + false, + ); +}); + +test("a negative interval is never due", () => { + assert.equal(isFullSweepDue(undefined, -5, NOW), false); +}); + +test("not due one minute short of the interval", () => { + const last = new Date(NOW - 1439 * MIN).toISOString(); + assert.equal(isFullSweepDue(last, 1440, NOW), false); +}); + +test("due at exactly the interval boundary", () => { + const last = new Date(NOW - 1440 * MIN).toISOString(); + assert.equal(isFullSweepDue(last, 1440, NOW), true); +}); + +test("due once past the interval", () => { + const last = new Date(NOW - 1441 * MIN).toISOString(); + assert.equal(isFullSweepDue(last, 1440, NOW), true); +}); + +test("a sweep that just ran is not due again", () => { + assert.equal(isFullSweepDue(new Date(NOW).toISOString(), 1440, NOW), false); +}); + +test("the shipped default is a daily sweep with a confirm cap", () => { + const d = defaultSyncScheduler(); + assert.equal(d.fullSweepIntervalMinutes, 1440); + assert.equal(d.fullSweepConfirmMaxSuspects, 25); +}); diff --git a/common/jobs/deepSync.ts b/common/jobs/deepSync.ts @@ -0,0 +1,43 @@ +import type { ChannelConfig } from "../lib/channelConfig"; +import type { SyncSchedulerSettings } from "../lib/settings"; + +// Pure cadence logic for the sync FULL SWEEP — the periodic deep pass that +// re-enumerates a channel's entire listing in one yt-dlp spawn and derives three +// outputs from it: a refreshed `playlist` file, the maybe-missing (upstream +// deletion) diff, and the same newest-first download set the paged walk would +// have picked. A full enumeration on a 5,000-VOD channel is far more expensive +// than one 50-entry page, so it runs at most once per interval per channel while +// ordinary syncs stay on the cheap paged walk. +// +// Mirrors resolveIntervalMinutes / overdueAmount in syncScheduler.ts: kept pure +// and side-effect-free so the due rules are unit-testable without a server. + +// Resolve a channel's effective full-sweep interval in minutes. undefined +// inherits the global default; 0 means the full sweep is disabled for the +// channel (every sync stays a cheap paged walk). +export function resolveFullSweepIntervalMinutes( + config: ChannelConfig, + scheduler: SyncSchedulerSettings, +): number { + if (config.fullSweepIntervalMinutes === undefined) { + return scheduler.fullSweepIntervalMinutes; + } + return config.fullSweepIntervalMinutes; +} + +// Whether a sync starting at `now` should pay for the full enumeration. A +// channel that has never swept is due immediately (that first sweep is what +// seeds `playlist` and maybe-missing.json); an unparseable timestamp is treated +// the same rather than wedging the channel out of sweeps forever. An interval of +// 0 (or less) disables the sweep outright and always wins. +export function isFullSweepDue( + lastFullSweepAt: string | undefined, + intervalMinutes: number, + now: number, +): boolean { + if (intervalMinutes <= 0) return false; + if (!lastFullSweepAt) return true; + const last = Date.parse(lastFullSweepAt); + if (Number.isNaN(last)) return true; + return now - last >= intervalMinutes * 60_000; +} diff --git a/common/lib/channelConfig.ts b/common/lib/channelConfig.ts @@ -103,6 +103,12 @@ export type ChannelConfig = { subLangs?: string; lastSyncedAt?: string; lastFullDownloadAt?: string; + // When this channel last paid for a sync FULL SWEEP — the deep pass that + // re-enumerates the whole listing to refresh `playlist` and flag videos that + // have left it. Stamped by the sweep itself; read by the cadence gate in + // common/jobs/deepSync.ts to decide whether the next sync sweeps or stays on + // the cheap newest-first paged walk. + lastFullSweepAt?: string; excludeFromBuild?: boolean; excludeFromSync?: boolean; // Opt this channel OUT of the aggregate "cleanable data" total shown on the @@ -119,6 +125,16 @@ export type ChannelConfig = { // > 0 -> sync this often (clamped to [SYNC_INTERVAL_MIN/MAX_MINUTES]) // A missing `url` or `excludeFromSync` also disables auto-sync. syncIntervalMinutes?: number; + // Full-sweep cadence for this channel (see common/jobs/deepSync.ts). A sync + // upgrades itself to a full sweep when + // `now - lastFullSweepAt >= fullSweepIntervalMinutes`. Same semantics as + // syncIntervalMinutes: + // undefined -> inherit the global syncScheduler.fullSweepIntervalMinutes + // 0 -> never sweep this channel (every sync is a paged walk) + // > 0 -> sweep this often (clamped to [SYNC_INTERVAL_MIN/MAX_MINUTES]) + // A sweep is one full enumeration, so long cadences (daily to weekly) are the + // norm — it is much more expensive than one 50-entry sync page. + fullSweepIntervalMinutes?: number; // Per-channel override for the global setting of the same name. When // omitted, the global SiteSettings value is used. 0 disables the sleep // for this channel. @@ -253,6 +269,9 @@ export function parseChannelConfig(raw: unknown): ChannelConfig | null { if (typeof r.lastFullDownloadAt === "string") { config.lastFullDownloadAt = r.lastFullDownloadAt; } + if (typeof r.lastFullSweepAt === "string") { + config.lastFullSweepAt = r.lastFullSweepAt; + } if (typeof r.excludeFromBuild === "boolean") { config.excludeFromBuild = r.excludeFromBuild; } @@ -278,6 +297,21 @@ export function parseChannelConfig(raw: unknown): ChannelConfig | null { SYNC_INTERVAL_MAX_MINUTES, ); } + if ( + typeof r.fullSweepIntervalMinutes === "number" && + Number.isFinite(r.fullSweepIntervalMinutes) && + r.fullSweepIntervalMinutes >= 0 + ) { + // Same shape as syncIntervalMinutes above: 0 is the "never sweep" sentinel. + config.fullSweepIntervalMinutes = + r.fullSweepIntervalMinutes === 0 + ? 0 + : clampInt( + r.fullSweepIntervalMinutes, + SYNC_INTERVAL_MIN_MINUTES, + SYNC_INTERVAL_MAX_MINUTES, + ); + } if (typeof r.skipLiveDownloads === "boolean") { config.skipLiveDownloads = r.skipLiveDownloads; } diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -345,6 +345,22 @@ export type SyncSchedulerSettings = { // default daily. The check shares the same concurrency cap and quiet-hours // window as scheduled syncs. See editor/app/scheduler/runTick.ts. keepLatestCheckIntervalMinutes: number; + // Default cadence (minutes) for the sync FULL SWEEP — the deep pass that + // re-enumerates a channel's whole listing in one yt-dlp spawn, refreshes the + // stored `playlist` file, and flags videos that have left the listing into + // maybe-missing.json. Ordinary syncs stay on the cheap newest-first paged + // walk; a sync only upgrades itself to a sweep when this interval has elapsed + // since the channel's lastFullSweepAt. Per-channel override: + // ChannelConfig.fullSweepIntervalMinutes. 0 = never sweep. Default daily. + // See common/jobs/deepSync.ts. + fullSweepIntervalMinutes: number; + // Upper bound on how many maybe-missing suspects a full sweep will resolve + // in-line with the per-video availability probe (deleted vs private vs + // unlisted). At or under the cap the sweep runs the targeted check itself, so + // "Sync all" surfaces upstream deletions with no extra clicks; over it, the + // suspects are flagged and left for a manual check rather than firing hundreds + // of probes inside a sync. 0 = never auto-confirm. + fullSweepConfirmMaxSuspects: number; }; export type SocialLink = { @@ -404,6 +420,13 @@ export const SYNC_SCHEDULER_MAX_CONCURRENT_MAX = 16; export const SYNC_SCHEDULER_BACKOFF_BASE_DEFAULT_MINUTES = 30; export const SYNC_SCHEDULER_BACKOFF_MAX_DEFAULT_MINUTES = 1440; export const KEEP_LATEST_CHECK_DEFAULT_INTERVAL_MINUTES = 1440; +// Full-sweep defaults. Daily: a sweep is one full enumeration of the channel, +// far more expensive than the 50-entry page an ordinary sync fetches. The +// confirm cap keeps an unattended sweep from fanning out into hundreds of +// per-video probes when a channel's listing changes wholesale. +export const FULL_SWEEP_DEFAULT_INTERVAL_MINUTES = 1440; +export const FULL_SWEEP_CONFIRM_MAX_SUSPECTS_DEFAULT = 25; +export const FULL_SWEEP_CONFIRM_MAX_SUSPECTS_MAX = 10000; export const SAVED_VIDEO_BACKUP_DEFAULT_INTERVAL_MINUTES = 1440; // Internal-heartbeat cadence bounds. 0 means "off" (use an external cron @@ -424,6 +447,8 @@ export function defaultSyncScheduler(): SyncSchedulerSettings { backoffMaxMinutes: SYNC_SCHEDULER_BACKOFF_MAX_DEFAULT_MINUTES, heartbeatSeconds: SYNC_HEARTBEAT_DEFAULT_SECONDS, keepLatestCheckIntervalMinutes: KEEP_LATEST_CHECK_DEFAULT_INTERVAL_MINUTES, + fullSweepIntervalMinutes: FULL_SWEEP_DEFAULT_INTERVAL_MINUTES, + fullSweepConfirmMaxSuspects: FULL_SWEEP_CONFIRM_MAX_SUSPECTS_DEFAULT, }; } @@ -447,6 +472,19 @@ function clampHourOrNull(value: unknown): number | null { return n; } +// Like clampPositiveInt, but 0 survives as a sentinel ("off"/"never"). Used by +// the cadences whose disabled state is expressed as a zero rather than a +// separate boolean. +function clampIntAllowZero(value: unknown, fallback: number, max: number): number { + const n = + typeof value === "number" && Number.isFinite(value) + ? Math.floor(value) + : fallback; + if (n <= 0) return 0; + if (n > max) return max; + return n; +} + function clampPositiveInt(value: unknown, fallback: number, max: number): number { const n = typeof value === "number" && Number.isFinite(value) @@ -501,6 +539,16 @@ export function sanitizeSyncScheduler(value: unknown): SyncSchedulerSettings { d.keepLatestCheckIntervalMinutes, SYNC_INTERVAL_MAX_MINUTES, ), + fullSweepIntervalMinutes: clampIntAllowZero( + r.fullSweepIntervalMinutes, + d.fullSweepIntervalMinutes, + SYNC_INTERVAL_MAX_MINUTES, + ), + fullSweepConfirmMaxSuspects: clampIntAllowZero( + r.fullSweepConfirmMaxSuspects, + d.fullSweepConfirmMaxSuspects, + FULL_SWEEP_CONFIRM_MAX_SUSPECTS_MAX, + ), }; } diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts @@ -29,7 +29,16 @@ import { } from "../lib/cookiePolicy"; import { resolveEffectiveAvailability } from "../lib/availability-server"; import { backfillAvailabilityFromMetadata } from "../controller/backfillAvailability"; +import { runAvailabilityCheck } from "../controller/checkAvailability"; +import { + diffMaybeMissing, + writeMaybeMissing, +} from "../controller/maybeMissingStore"; import { resolveShardItems } from "../controller/shard"; +import { + isFullSweepDue, + resolveFullSweepIntervalMinutes, +} from "../jobs/deepSync"; import { computeKeepWindow, type KeepWindow } from "../controller/keptVideos"; import { downloadOneManaged } from "./downloadOneManaged"; import { @@ -99,6 +108,11 @@ export type RunYtdlpOpts = { extractImmediately?: boolean; // download-one-audio only: appended after configArgs, before the URL. extraYtdlpArgs?: string[]; + // sync only: force this run to be a FULL SWEEP regardless of the configured + // cadence (see common/jobs/deepSync.ts). Undefined leaves the decision to the + // cadence gate; false forces the cheap paged walk. Backs the explicit + // "Full sweep" action and its job replay. + forceFullSweep?: boolean; // retry-bucket only: video IDs to limit the run to (matched against the // saved playlist by extracted ID). IDs not present in the playlist are // reported and skipped. @@ -420,16 +434,27 @@ export async function probeChannelMeta(opts: { return { name: lines[0] ?? null, entryCount: null }; } -async function storePlaylist(opts: RunYtdlpOpts): Promise<void> { - const root = channelRoot(opts); - await mkdir(root, { recursive: true }); +// Atomically replace the channel's stored `playlist` file (tmp + rename, so a +// crashed write never leaves a truncated list behind). Shared by store-playlist +// and the sync full sweep, which both enumerate the whole channel and are the +// only writers of this file. +async function writePlaylistFile( + root: string, + urls: ReadonlyArray<string>, + onLog: (s: string) => void, +): Promise<void> { const playlistPath = path.join(root, "playlist"); const tmpPath = `${playlistPath}.tmp-${process.pid}`; - - const urls = await enumeratePlaylistUrls(opts, root); await writeFile(tmpPath, urls.join("\n") + (urls.length ? "\n" : "")); await rename(tmpPath, playlistPath); - opts.onLog(`Wrote ${urls.length} URLs to ${playlistPath}\n`); + onLog(`Wrote ${urls.length} URLs to ${playlistPath}\n`); +} + +async function storePlaylist(opts: RunYtdlpOpts): Promise<void> { + const root = channelRoot(opts); + await mkdir(root, { recursive: true }); + const urls = await enumeratePlaylistUrls(opts, root); + await writePlaylistFile(root, urls, opts.onLog); } async function downloadFromPlaylist(opts: RunYtdlpOpts): Promise<void> { @@ -1155,7 +1180,78 @@ export async function destinationExists( // the first page that contains an already-archived entry. const SYNC_PAGE_SIZE = 50; +// A sync is one of two passes over the same channel listing: +// +// syncPaged the cheap default. Walks newest-first one 50-entry +// flat-playlist spawn at a time and stops at the first page +// containing an already-archived entry. +// syncFullSweep the periodic deep pass. ONE full enumeration, from which all +// three outputs are derived: a refreshed `playlist` file, the +// maybe-missing (upstream deletion) diff, and the very same +// download set the paged walk would have chosen. +// +// The cadence gate decides which (see common/jobs/deepSync.ts); a sweep is far +// more expensive on a large channel, so it runs at most once per interval. +// `forceFullSweep` overrides the gate in both directions. async function sync(opts: RunYtdlpOpts): Promise<void> { + const sweep = opts.forceFullSweep ?? fullSweepDue(opts); + return sweep ? syncFullSweep(opts) : syncPaged(opts); +} + +// Whether this sync should upgrade itself to a full sweep, resolved from the +// per-channel override, the global cadence and the channel's lastFullSweepAt. +function fullSweepDue(opts: RunYtdlpOpts): boolean { + const scheduler = getSettings().syncScheduler; + const interval = resolveFullSweepIntervalMinutes( + opts.channelConfig, + scheduler, + ); + return isFullSweepDue( + opts.channelConfig.lastFullSweepAt, + interval, + Date.now(), + ); +} + +// The per-page download filter, shared by both passes so they can never drift: +// entries already in the archive are counted as hits (the paged walk's stop +// signal), defer-mode needs_auth videos are held back for the Needs-cookies +// bucket, and everything else is queued for download. +async function selectDownloadableUrls( + pageUrls: ReadonlyArray<string>, + archive: Awaited<ReturnType<typeof readArchive>>, + dataDir: string, + runCookiePolicy: ResolvedCookiePolicy, +): Promise<{ + newUrls: string[]; + archivedHits: number; + deferredAuthCount: number; +}> { + const newUrls: string[] = []; + let archivedHits = 0; + let deferredAuthCount = 0; + for (const url of pageUrls) { + const archiveId = await archiveIdForUrl(url, dataDir); + if (archiveId && archive.ids.has(archiveId)) { + archivedHits++; + continue; + } + // Cookie mode "defer": don't re-attempt known auth-gated videos on + // every sync — they wait in the Needs-cookies bucket instead. + const dirId = extractVideoId(url); + if ( + dirId && + (await isDeferredAuthExcluded(path.join(dataDir, dirId), runCookiePolicy)) + ) { + deferredAuthCount++; + continue; + } + newUrls.push(url); + } + return { newUrls, archivedHits, deferredAuthCount }; +} + +async function syncPaged(opts: RunYtdlpOpts): Promise<void> { const root = channelRoot(opts); const dataDir = path.join(root, "data"); const archivePath = path.join(root, "archive"); @@ -1190,30 +1286,8 @@ async function sync(opts: RunYtdlpOpts): Promise<void> { const pageUrls = await enumeratePlaylistUrls(opts, root, { start, end }); if (pageUrls.length === 0) break; - const newUrls: string[] = []; - let archivedHits = 0; - let deferredAuthCount = 0; - for (const url of pageUrls) { - const archiveId = await archiveIdForUrl(url, dataDir); - if (archiveId && archive.ids.has(archiveId)) { - archivedHits++; - continue; - } - // Cookie mode "defer": don't re-attempt known auth-gated videos on - // every sync — they wait in the Needs-cookies bucket instead. - const dirId = extractVideoId(url); - if ( - dirId && - (await isDeferredAuthExcluded( - path.join(dataDir, dirId), - runCookiePolicy, - )) - ) { - deferredAuthCount++; - continue; - } - newUrls.push(url); - } + const { newUrls, archivedHits, deferredAuthCount } = + await selectDownloadableUrls(pageUrls, archive, dataDir, runCookiePolicy); opts.onLog( `Sync page ${page + 1}: ${pageUrls.length} entries, ${newUrls.length} new, ${archivedHits} already archived.\n`, ); @@ -1257,6 +1331,175 @@ async function sync(opts: RunYtdlpOpts): Promise<void> { await safeBackfillAvailability(opts); } +// The periodic deep pass. ONE full `--flat-playlist --print url` enumeration +// serves all three purposes that used to cost three separate yt-dlp spawns: +// refreshing the stored `playlist`, detecting videos that have left the channel +// listing, and picking up new uploads. +// +// The download set is deliberately IDENTICAL to what syncPaged would have +// fetched: newest-first entries up to the first SYNC_PAGE_SIZE window +// containing an already-archived entry, just applied to a list already in +// memory. Downloading everything absent from the archive would be a +// download-from-playlist, and on a channel deliberately kept partial (newest +// 200 of 5,000) it would drag down the entire back catalogue — including audio +// the cleanup sweep has already reclaimed. A sweep changes what we KNOW, never +// what we FETCH; the deletion signal comes from the listing diff, not downloads. +async function syncFullSweep(opts: RunYtdlpOpts): Promise<void> { + const root = channelRoot(opts); + const dataDir = path.join(root, "data"); + const archivePath = path.join(root, "archive"); + await mkdir(root, { recursive: true }); + + opts.onLog( + `Full sweep: re-reading the whole channel listing (refreshes the video list and flags videos that have gone missing).\n`, + ); + + // 1. One enumeration, no range. + const urls = await enumeratePlaylistUrls(opts, root); + if (opts.signal.aborted) return; + + // 2. Refresh the stored playlist. Everything downstream — "download missing", + // the snapshot's undownloadedIds — reads this file, and before the sweep + // only an explicit "store playlist" ever rewrote it. + await writePlaylistFile(root, urls, opts.onLog); + + // 3. Diff the fresh listing against what's on disk. Canonical, URL-derived + // ids match data-dir names on every platform; nulls are dropped because a + // null would never match and would spuriously flag every known video. + const freshIds = new Set( + urls.map((u) => extractVideoId(u)).filter((id): id is string => Boolean(id)), + ); + const entries = await readdir(dataDir, { withFileTypes: true }).catch( + () => [] as Awaited<ReturnType<typeof readdir>> & { length: 0 }, + ); + const knownIds = (entries as Array<{ isDirectory(): boolean; name: string }>) + .filter((e) => e.isDirectory()) + .map((e) => e.name); + const maybeMissing = diffMaybeMissing(knownIds, freshIds); + await writeMaybeMissing(opts.paths, opts.channelSlug, { + checkedAt: new Date().toISOString(), + freshPlaylistCount: freshIds.size, + ids: maybeMissing, + }); + opts.onLog( + `Full sweep: ${knownIds.length} known, ${freshIds.size} in fresh listing, ${maybeMissing.length} maybe-missing.\n`, + ); + + // 4. Walk the same listing in newest-first windows, applying the paged walk's + // filter and stopping rule to slices instead of spawns. + let archive = await readArchive(archivePath); + let totalNew = 0; + let totalFailed = 0; + let totalSkipped = 0; + let firstFailure: Error | null = null; + const runCookiePolicy = resolveRunCookiePolicy(opts, opts.channelConfig); + + for ( + let page = 0; + !opts.signal.aborted && !opts.drainSignal?.aborted; + page++ + ) { + const pageUrls = urls.slice( + page * SYNC_PAGE_SIZE, + (page + 1) * SYNC_PAGE_SIZE, + ); + if (pageUrls.length === 0) break; + + const { newUrls, archivedHits, deferredAuthCount } = + await selectDownloadableUrls(pageUrls, archive, dataDir, runCookiePolicy); + opts.onLog( + `Sync page ${page + 1}: ${pageUrls.length} entries, ${newUrls.length} new, ${archivedHits} already archived.\n`, + ); + if (deferredAuthCount > 0) { + opts.onLog( + `Cookie mode is "defer": needs_auth deferred=${deferredAuthCount} — run the "Needs cookies" bucket to download them with cookies.\n`, + ); + } + + if (newUrls.length > 0) { + const res = await runManagedDownloads( + opts, + newUrls, + opts.channelConfig, + runCookiePolicy, + ); + totalNew += res.okCount; + totalFailed += res.failedCount; + totalSkipped += res.skippedCount; + if (res.firstFailure) { + firstFailure = res.firstFailure; + break; + } + archive = await readArchive(archivePath); + } + + if (archivedHits > 0) break; // reached previously-synced content + if (pageUrls.length < SYNC_PAGE_SIZE) break; // last page + } + + if (firstFailure && !opts.signal.aborted) throw firstFailure; + + opts.onLog( + `Sync complete: ${totalNew} new downloaded, ${totalFailed} failed` + + (totalSkipped ? `, ${totalSkipped} skipped by filter` : "") + + `.\n`, + ); + + // 5. Capped auto-confirm. Resolving a suspect (deleted vs private vs unlisted) + // costs one yt-dlp probe each, so it only runs in-line while the suspect + // count is small — the common case of a handful of videos disappearing. + // Over the cap this is a manual decision, not something an unattended sync + // should fan out into. Runs in the same job (one streamed log, one Cancel), + // so re-check the abort between phases. + await confirmMaybeMissing(opts, maybeMissing); + + await touchLastSync(opts); + await touchLastFullSweep(opts); + await safeBackfillAvailability(opts); +} + +// Resolve the sweep's suspects with the targeted per-video availability probe, +// when there are few enough to be worth doing unattended. The option set mirrors +// checkMaybeMissingAction / verifyBeforeClean: `ignoreShard` is required because +// a saved shard-availability.json would otherwise replace `onlyIds` wholesale, +// and `skipExpectedAbsent` avoids re-probing videos already known gone. +async function confirmMaybeMissing( + opts: RunYtdlpOpts, + maybeMissing: ReadonlyArray<string>, +): Promise<void> { + if (maybeMissing.length === 0) return; + if (opts.signal.aborted) return; + const cap = getSettings().syncScheduler.fullSweepConfirmMaxSuspects; + if (cap <= 0 || maybeMissing.length > cap) { + opts.onLog( + `Full sweep: ${maybeMissing.length} maybe-missing video(s) flagged, over the auto-confirm cap of ${cap} — run "Check maybe-missing" on the channel to resolve them.\n`, + ); + return; + } + opts.onLog( + `Full sweep: confirming ${maybeMissing.length} maybe-missing video(s) upstream…\n`, + ); + try { + await runAvailabilityCheck({ + channelSlug: opts.channelSlug, + paths: opts.paths, + mode: "recheck-all", + onlyIds: maybeMissing, + skipExpectedAbsent: true, + ignoreShard: true, + concurrency: 1, + onLog: opts.onLog, + signal: opts.signal, + }); + } catch (err) { + // Best-effort, exactly like the availability backfill: the sweep's real + // product is the flagged set, which is already persisted. + opts.onLog( + `Full sweep: maybe-missing confirmation skipped: ${(err as Error).message}\n`, + ); + } +} + async function safeBackfillAvailability(opts: RunYtdlpOpts): Promise<void> { // Best-effort: write availability.json for any new video dir from this // download whose metadata.info.json has an availability field. Don't fail @@ -1312,9 +1555,16 @@ async function touchLastFullDownload(opts: RunYtdlpOpts): Promise<void> { await updateConfigField(file, "lastFullDownloadAt", new Date().toISOString()); } +// Stamps when this channel last paid for a full enumeration, which is what the +// cadence gate reads to keep the next N syncs on the cheap paged walk. +async function touchLastFullSweep(opts: RunYtdlpOpts): Promise<void> { + const file = path.join(channelRoot(opts), "config.json"); + await updateConfigField(file, "lastFullSweepAt", new Date().toISOString()); +} + async function updateConfigField( configPath: string, - field: "lastSyncedAt" | "lastFullDownloadAt", + field: "lastSyncedAt" | "lastFullDownloadAt" | "lastFullSweepAt", value: string, ): Promise<void> { let raw: string; diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -84,6 +84,9 @@ async function runPipelineAction( 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 bookmarked. spec?: JobSpec; }, @@ -188,6 +191,7 @@ async function runPipelineAction( bucketIds: options?.bucketIds, handlingOverride: options?.handlingOverride, forceCookies: options?.forceCookies, + forceFullSweep: options?.forceFullSweep, replaceAutoSubs: options?.replaceAutoSubs, keepSourceVideoOverride: options?.keepSourceVideoOverride, extractImmediately: options?.extractImmediately, @@ -284,12 +288,18 @@ export async function downloadMissingAction( ); } +// 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<StreamActionResult> { return runPipelineAction(slug, "sync", "sync", queueKey, { - spec: { kind: "sync", slug, params: { queueKey } }, + forceFullSweep: fullSweep, + spec: { kind: "sync", slug, params: { queueKey, fullSweep } }, }); } diff --git a/editor/app/channels/actions.ts b/editor/app/channels/actions.ts @@ -362,7 +362,13 @@ export type SyncAllResult = { skipped: { slug: string; reason: string }[]; }; -export async function syncAllChannelsAction(): Promise<SyncAllResult> { +// Queue a sync for every eligible channel. Each sync decides for itself whether +// it is due for a full sweep, so "Sync all" surfaces upstream deletions on +// whichever channels are due with no extra clicks — pass fullSweep to force the +// deep pass on every channel instead. +export async function syncAllChannelsAction( + opts?: { fullSweep?: boolean }, +): Promise<SyncAllResult> { const paths = getPaths(); const channels = await listChannelConfigs(paths); const active = activeSyncSlugs(); @@ -381,7 +387,7 @@ export async function syncAllChannelsAction(): Promise<SyncAllResult> { skipped.push({ slug: c.slug, reason: "already running" }); continue; } - const result = await syncAction(c.slug); + const result = await syncAction(c.slug, undefined, opts?.fullSweep); if (!result.ok) { skipped.push({ slug: c.slug, reason: result.error }); continue; diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts @@ -180,8 +180,8 @@ export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = { ); }, sync: (spec) => { - const { queueKey } = params(spec); - return syncAction(spec.slug, queueKey); + const { p, queueKey } = params(spec); + return syncAction(spec.slug, queueKey, bool(p.fullSweep)); }, "fetch-posts": (spec) => { const { p, queueKey } = params(spec);