// THIS DIRECTORY IS THE SYNC OPERATION'S RUNNER — the heartbeat, the tick and // the status payload. Its PAGE is /operations/sync (rendered by // operations/[id]), and /scheduler redirects there; nothing here moved, because // instrumentation.ts, both /api/scheduler/* routes, the operations rail's sync // row and two forms import from it. import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { listChannelConfigs } from "yt-dlp-transcript-common/controller/channels"; import { getSettings } from "yt-dlp-transcript-common/lib/settings"; import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; import { activeSyncJobs } from "yt-dlp-transcript-common/jobs/syncJobs"; import type { SyncSchedulerSettings } from "yt-dlp-transcript-common/lib/settings"; import { isInQuietHours, nextEligibleAfterFailure, selectDueChannels, type ChannelEntry, } from "yt-dlp-transcript-common/jobs/syncScheduler"; import { channelState, readSchedulerState, recordRun, writeSchedulerState, type SchedulerSkip, type SchedulerState, } from "yt-dlp-transcript-common/jobs/syncSchedulerState"; import { isSocialChannel } from "yt-dlp-transcript-common/lib/channelConfig"; import { resolveFocusSlugs } from "yt-dlp-transcript-common/lib/channelPriority"; import { siteChannelIndex } from "yt-dlp-transcript-common/lib/site"; import { syncAction } from "../channels/[slug]/pipelineActions"; import { fetchPostsAction } from "../channels/[slug]/socialActions"; // STORAGE CHORES riding this heartbeat because it is the one timer the editor // has. Neither is sync (IA doc, the Sync row): the keep-latest check consumes // outputs, the backup copies them. import { checkKeptDeletedAction } from "../channels/[slug]/whisperActions"; import { backupSavedVideosAction } from "../saved-videos/backupActions"; export type SchedulerTickResult = { ok: boolean; ranAt: number; // Why nothing (or limited work) ran, when applicable: "tick already running", // "scheduler disabled", "quiet hours". Absent on a normal active tick. reason?: string; // Number of sync jobs already running/queued when the tick started. running: number; // Channels for which a sync was queued this tick. queued: string[]; // Storage chore: channels for which a keep-latest deletion check was queued // this tick. keptChecksQueued: string[]; // Storage chore: whether a saved-video backup was queued this tick. savedVideoBackupQueued: boolean; // Channels deliberately held back, with a reason. skipped: SchedulerSkip[]; }; declare global { // eslint-disable-next-line no-var var __yttSyncSchedulerTick__: boolean | undefined; } // One scheduler heartbeat. Reconciles the outcome of previously-queued syncs // (for backoff), selects the channels that are due, queues up to the // concurrency cap (most-overdue first), persists state, and returns a summary. // // Reuses the existing job machinery end to end: syncAction enqueues through // runManagedFunction with the per-channel download queue lock, so a scheduled // sync can never collide with a manual one — the lock that a standalone cron // process would have had to reinvent. export async function runSchedulerTick(): Promise { const now = Date.now(); // Synchronous overlap guard set before the first await, so two near- // simultaneous ticks can't interleave (Laravel's withoutOverlapping). if (globalThis.__yttSyncSchedulerTick__) { return { ok: false, ranAt: now, reason: "tick already running", running: 0, queued: [], keptChecksQueued: [], savedVideoBackupQueued: false, skipped: [], }; } globalThis.__yttSyncSchedulerTick__ = true; try { const paths = getPaths(); const settings = getSettings(); const scheduler = settings.syncScheduler; const state = await readSchedulerState(paths); const activeJobs = activeSyncJobs(); const activeSlugs = new Set( activeJobs.map((j) => j.channelSlug as string), ); const running = activeJobs.length; // Fold the outcomes of prior scheduler-queued syncs into the backoff state. reconcileOutcomes(state, now, scheduler); if (!scheduler.enabled) { await writeSchedulerState(paths, state); return { ok: true, ranAt: now, reason: "scheduler disabled", running, queued: [], keptChecksQueued: [], savedVideoBackupQueued: false, skipped: [], }; } const channels: ChannelEntry[] = (await listChannelConfigs(paths)).map((c) => ({ slug: c.slug, config: c.config, })); // The focus set is the only half of the priority model that costs I/O, and // ONLY a `{kind:"site"}` focus pays it: `resolveFocusSlugs` reads a site's // `channels[]` so a focus on a site tracks its membership instead of // freezing a list. `{kind:"none"}` and `{kind:"channels"}` resolve from the // document alone. No cache here — the tick runs on the heartbeat, not per // grant, so one `siteChannelIndex()` per tick is not a cost worth memoizing. // `siteChannelIndex` (lib/site.ts) is the one spelling of that read; the // runner, the /channels writer and the status payload ask the same one. const priority = settings.channelPriority; const siteChannels = priority.focus.kind === "site" ? siteChannelIndex(paths) : {}; const { due, skipped } = selectDueChannels({ channels, scheduler, state, activeSlugs, now, priority, focusSlugs: resolveFocusSlugs( priority, siteChannels, channels.map((c) => c.slug), ), }); // Concurrency cap doubles as the stagger: queue at most (cap - running) // channels; the rest roll to the next tick instead of fanning out at once. const slots = Math.max(0, scheduler.maxConcurrentSyncs - running); const toQueue = due.slice(0, slots); for (const slug of due.slice(slots)) { skipped.push({ slug, reason: "concurrency cap (deferred to next tick)" }); } const queued: string[] = []; const bySlug = new Map(channels.map((c) => [c.slug, c.config])); for (const slug of toQueue) { // The scheduler's ELIGIBILITY rules are source-agnostic (url + sync // tier + interval + lastSyncedAt), but the dispatch is not: a // social channel must run a post fetch, not a yt-dlp video sync against // its profile URL. const result = isSocialChannel(bySlug.get(slug)) ? await fetchPostsAction(slug) : await syncAction(slug); if (!result.ok) { skipped.push({ slug, reason: result.error }); continue; } const cs = channelState(state, slug); cs.lastQueuedJobId = result.jobId; cs.lastQueuedAt = now; queued.push(slug); // Nobody consumes the log stream here — cancel it so buffered chunks are // freed. The job keeps running and logging to disk. (Same as the manual // "Sync all" action.) void result.stream.cancel(); } const quiet = isInQuietHours( now, scheduler.quietHoursStart, scheduler.quietHoursEnd, ); // Storage chore: keep-latest deletion checks run on their own cadence, on // the per-channel local queue (channel:) — separate from the platform // download queue a sync uses, so the two don't serialize against each // other. Suppressed during quiet hours, capped per tick like syncs, and // skipped for any channel we just queued a sync for (avoid hitting the // source twice in one tick). const keptChecksQueued: string[] = []; if (!quiet) { const queuedThisTick = new Set(queued); const intervalMs = scheduler.keepLatestCheckIntervalMinutes * 60_000; const dueForKeptCheck = channels.filter((c) => { if (!c.config.keepLatest || c.config.keepLatest <= 0) return false; if (queuedThisTick.has(c.slug)) return false; const last = channelState(state, c.slug).lastKeptCheckAt ?? 0; return now - last >= intervalMs; }); for (const c of dueForKeptCheck.slice(0, scheduler.maxConcurrentSyncs)) { const result = await checkKeptDeletedAction(c.slug); if (!result.ok) { skipped.push({ slug: c.slug, reason: `storage chore (keep-latest check): ${result.error}`, }); continue; } channelState(state, c.slug).lastKeptCheckAt = now; keptChecksQueued.push(c.slug); void result.stream.cancel(); } } // Storage chore: the saved-video backup, a single global job on its own // cadence. Suppressed during quiet hours; runs at most once per // intervalMinutes. Independent of the per-channel concurrency cap (it // touches no source provider). let savedVideoBackupQueued = false; const backup = settings.savedVideoBackup; if (!quiet && backup.enabled && backup.dest.trim()) { const intervalMs = backup.intervalMinutes * 60_000; const last = state.lastSavedVideoBackupAt ?? 0; if (now - last >= intervalMs) { const result = await backupSavedVideosAction(); if (result.ok) { state.lastSavedVideoBackupAt = now; savedVideoBackupQueued = true; void result.stream.cancel(); } else { skipped.push({ slug: "storage chore (saved-video backup)", reason: result.error, }); } } } recordRun(state, { at: now, queued, skipped }); await writeSchedulerState(paths, state); return { ok: true, ranAt: now, reason: quiet ? "quiet hours" : undefined, running, queued, keptChecksQueued, savedVideoBackupQueued, skipped, }; } finally { globalThis.__yttSyncSchedulerTick__ = false; } } // Inspect the registry for each channel's last scheduler-queued sync and update // its backoff counters once that job reaches a terminal state. A finished job // that's already been evicted from the bounded registry is treated as neutral // (we simply stop tracking it). "cancelled" is neutral too — a user/drain // cancel shouldn't penalize the channel; only "failed" grows the backoff. function reconcileOutcomes( state: SchedulerState, now: number, scheduler: SyncSchedulerSettings, ): void { const registry = getRegistry(); for (const cs of Object.values(state.channels)) { if (!cs.lastQueuedJobId) continue; const job = registry.get(cs.lastQueuedJobId); if (!job) { cs.lastQueuedJobId = null; // evicted; can't read outcome continue; } if (job.status === "queued" || job.status === "running") continue; if (job.status === "done") { cs.consecutiveFailures = 0; cs.nextEligibleAt = null; cs.lastOutcome = "ok"; cs.lastOutcomeAt = job.endedAt ?? now; } else if (job.status === "failed") { cs.consecutiveFailures += 1; cs.nextEligibleAt = nextEligibleAfterFailure( now, cs.consecutiveFailures, scheduler, ); cs.lastOutcome = "failed"; cs.lastOutcomeAt = job.endedAt ?? now; } cs.lastQueuedJobId = null; } }