Archilyzer · Source

archilyzer

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

commit e7862a94777ed795f5360002972845d49cca2872
parent d21d6f927a6e2e851b3e053b93f25e4252fa1bf2
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 14 Jun 2026 03:15:08 -0400

Scheduled sync: per-channel auto-sync via cron heartbeat

Add a Laravel-scheduler-inspired sync system. A lightweight cron client
(pnpm sync:tick) POSTs to /api/scheduler/tick and the editor server decides
which channels are due (now - lastSyncedAt >= the channel's interval), then
queues them through the existing job queue + per-channel lock + worker pool —
so scheduled syncs can't collide with manual ones, show up live, and need no
second process or new file locks.

- ChannelConfig.syncIntervalMinutes (undefined=inherit, 0=off, else minutes)
- SiteSettings.syncScheduler: enable, default interval, max-concurrent,
  quiet hours, failure backoff
- Pure selectDueChannels + buildScheduleView + backoff math (common/jobs/
  syncScheduler.ts); durable backoff/run-log state (.scheduler/state.json)
- runSchedulerTick: overlap guard, outcome reconciliation, concurrency-cap-as-
  stagger; POST tick + GET status routes (bearer-auth via SYNC_TICK_TOKEN)
- common/bin/sync-tick.ts + pnpm sync:tick; SCHEDULED_SYNC.md cron docs
- Channel Auto-sync field, Settings controls, new /scheduler panel
- e2e scheduler.spec.ts; generalized syncAllChannelsAction via shared
  activeSyncSlugs helper

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

Diffstat:
MREADME.md | 2++
ASCHEDULED_SYNC.md | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/bin/sync-tick.ts | 55+++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/syncJobs.ts | 20++++++++++++++++++++
Acommon/jobs/syncScheduler.ts | 207+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/syncSchedulerState.ts | 161+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/channelConfig.ts | 29+++++++++++++++++++++++++++++
Mcommon/lib/paths.ts | 5+++++
Mcommon/lib/settings.ts | 108+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/CHANGELOG.md | 1+
Aeditor/app/api/scheduler/status/route.ts | 11+++++++++++
Aeditor/app/api/scheduler/tick/route.ts | 16++++++++++++++++
Meditor/app/channels/actions.ts | 14++------------
Meditor/app/channels/components/ChannelForm.tsx | 46++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/channels/components/parseChannelForm.ts | 27+++++++++++++++++++++++++++
Meditor/app/layout.tsx | 1+
Aeditor/app/scheduler/auth.ts | 13+++++++++++++
Aeditor/app/scheduler/components/SchedulerView.tsx | 253+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/scheduler/page.tsx | 33+++++++++++++++++++++++++++++++++
Aeditor/app/scheduler/runTick.ts | 186+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/scheduler/status.ts | 41+++++++++++++++++++++++++++++++++++++++++
Meditor/app/settings/actions.ts | 21+++++++++++++++++++++
Meditor/app/settings/components/SettingsForm.tsx | 81+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/e2e/scheduler.spec.ts | 99+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mpackage.json | 1+
25 files changed, 1492 insertions(+), 12 deletions(-)

diff --git a/README.md b/README.md @@ -29,6 +29,8 @@ pnpm e2e:sharded # same suite, sharded across N Docker containers in paralle - **/jobs** — recent jobs, running and archived. Click an id to tail its log. - **/settings** — site-settings form persisted to `settings.json`, plus a read-only view of `getPaths()` for env-resolution debugging. +Channels can also sync automatically on a per-channel cadence via a cron heartbeat — see [SCHEDULED_SYNC.md](SCHEDULED_SYNC.md). + ## Environment variables All paths and binaries used by the editor and CLI shims resolve through `getPaths()`. Override any of them before launching: diff --git a/SCHEDULED_SYNC.md b/SCHEDULED_SYNC.md @@ -0,0 +1,73 @@ +# Scheduled channel sync + +Channels can sync automatically on a per-channel cadence (hourly, daily, every +N minutes, …), inspired by the Laravel scheduler: a single dumb cron heartbeat +runs often, and the **server** decides which channels are actually due. + +## How it works + +``` +cron (*/5 * * * *) + └─ pnpm sync:tick external lightweight client + └─ POST /api/scheduler/tick editor server + └─ selectDueChannels() common/jobs/syncScheduler.ts + └─ syncAction(slug) existing job queue + worker pool +``` + +- A channel is **due** when `now - lastSyncedAt >= its interval`. A missed tick + (server was down, machine asleep) just runs at the next tick — no catch-up + storms. +- All work runs inside the editor server process, so scheduled syncs reuse the + per-channel queue lock (no collisions with a manual "Sync" click), show up + live in the editor UI, and feed the same transcription worker pool. + +## Configuration + +Per-channel cadence lives on the channel (channel editor → **Auto-sync**), and is +stored as `syncIntervalMinutes` in `channels/<slug>/config.json`: + +- *Default* — inherit the global default interval. +- *Off* — never auto-sync this channel (manual sync still works). +- A concrete interval (10m / 30m / hourly / 6h / daily / …). + +Global controls live in **Settings → Sync scheduler**: + +- **Enabled** — master on/off (off by default). +- **Default interval** — fallback for channels without an override. +- **Max concurrent syncs** — cap on simultaneous sync jobs. A tick queues at most + `cap - running` channels (most-overdue first); the rest roll to the next tick. + This both bounds load and staggers a large due-batch over several ticks. +- **Quiet hours** — optional local-clock window when auto-sync is suppressed. +- **Failure backoff** — after consecutive failures a channel waits + `base · 2^(n-1)` minutes (capped) before retrying, so a broken channel doesn't + retry every tick. + +Channels already marked **Exclude from sync** are never auto-synced. + +## Cron setup + +Run the heartbeat at an interval no larger than your smallest channel interval +(`*/5` comfortably serves a 10-minute minimum): + +```cron +*/5 * * * * cd /path/to/repo && pnpm sync:tick >> /var/log/ytt-sync.log 2>&1 +``` + +The client targets `http://127.0.0.1:3001/api/scheduler/tick` by default (the +editor's dev/start port). Override with env vars: + +- `SYNC_TICK_URL` — full endpoint URL (if the editor runs on another host/port). +- `SYNC_TICK_TOKEN` — bearer token. When set, the editor server must have the + **same** `SYNC_TICK_TOKEN` in its environment, and the tick route rejects + requests without a matching `Authorization: Bearer <token>` header. Recommended + if the editor is reachable beyond localhost. When unset, the route is open + (consistent with the otherwise-unauthenticated editor admin surface). + +`pnpm sync:tick` prints a one-line summary and exits non-zero on a network/HTTP +error so cron can surface failures. + +## Observability + +**Channels → Sync schedule** shows each channel's resolved interval, last sync, +next-due time, recent outcome, and any active backoff, plus a log of recent +ticks. The same data is available as JSON at `GET /api/scheduler/status`. diff --git a/common/bin/sync-tick.ts b/common/bin/sync-tick.ts @@ -0,0 +1,55 @@ +#!/usr/bin/env tsx +// Lightweight cron client for the scheduled sync system. Cron runs this on a +// fixed heartbeat (e.g. `*/5 * * * *`); it just POSTs to the editor's +// /api/scheduler/tick and the SERVER decides which channels are due and queues +// them through the existing job machinery. Keeping the heartbeat dumb is the +// whole point — see common/jobs/syncScheduler.ts and editor/app/scheduler/. +// +// Env: +// SYNC_TICK_URL full URL of the tick endpoint. Default targets the editor's +// real dev/start port (3001): +// http://127.0.0.1:3001/api/scheduler/tick +// SYNC_TICK_TOKEN optional bearer token. When set, it must match the token in +// the editor server's environment or the request is rejected. +// +// Exits non-zero on a network/HTTP error so cron surfaces failures (e.g. mails +// the output); a normal tick prints a one-line summary for the cron log. + +const DEFAULT_URL = "http://127.0.0.1:3001/api/scheduler/tick"; + +async function main(): Promise<void> { + const url = process.env.SYNC_TICK_URL ?? DEFAULT_URL; + const token = process.env.SYNC_TICK_TOKEN; + const headers: Record<string, string> = { "content-type": "application/json" }; + if (token) headers.authorization = `Bearer ${token}`; + + const res = await fetch(url, { method: "POST", headers }); + const text = await res.text(); + if (!res.ok) { + console.error(`[sync-tick] ${res.status} ${res.statusText}: ${text}`); + process.exitCode = 1; + return; + } + + try { + const r = JSON.parse(text) as { + queued?: string[]; + skipped?: unknown[]; + reason?: string; + }; + const queued = r.queued ?? []; + const skipped = r.skipped ?? []; + const reason = r.reason ? ` (${r.reason})` : ""; + const detail = queued.length ? `: ${queued.join(", ")}` : ""; + console.log( + `[sync-tick] queued ${queued.length}, skipped ${skipped.length}${reason}${detail}`, + ); + } catch { + console.log(`[sync-tick] ${text}`); + } +} + +main().catch((err: unknown) => { + console.error(`[sync-tick] request failed: ${(err as Error).message}`); + process.exit(1); +}); diff --git a/common/jobs/syncJobs.ts b/common/jobs/syncJobs.ts @@ -0,0 +1,20 @@ +import { getRegistry, type JobRecord } from "./registry"; + +// Sync jobs that are currently running or queued. Single source of truth for +// "is this channel already syncing?" — shared by the cron scheduler +// (common/jobs/syncScheduler.ts orchestration) and the manual "Sync all" action +// so the two never drift in how they dedup against in-flight work. +export function activeSyncJobs(): JobRecord[] { + return getRegistry() + .list() + .filter( + (j) => + j.kind === "sync" && + (j.status === "queued" || j.status === "running") && + !!j.channelSlug, + ); +} + +export function activeSyncSlugs(): Set<string> { + return new Set(activeSyncJobs().map((j) => j.channelSlug as string)); +} diff --git a/common/jobs/syncScheduler.ts b/common/jobs/syncScheduler.ts @@ -0,0 +1,207 @@ +import type { ChannelConfig } from "../lib/channelConfig"; +import type { SyncSchedulerSettings } from "../lib/settings"; +import type { SchedulerSkip, SchedulerState } from "./syncSchedulerState"; + +// Pure, side-effect-free scheduling logic for the cron-driven sync system. It +// decides WHICH channels are due to sync right now; the editor tick route does +// the I/O (reads settings/state/registry, queues syncs, persists state). Keeping +// this pure makes the due/skip/ordering rules unit-testable without a server. + +export type ChannelEntry = { + slug: string; + config: ChannelConfig; +}; + +export type SelectDueInput = { + channels: ReadonlyArray<ChannelEntry>; + scheduler: SyncSchedulerSettings; + state: SchedulerState; + // Slugs that already have a running or queued sync job (from the registry). + activeSlugs: ReadonlySet<string>; + now: number; +}; + +export type SelectDueResult = { + // Slugs that should be synced now, ordered most-overdue first. + due: string[]; + // Channels deliberately held back, with a human reason (for the run log). + // The common "not yet due" case is intentionally omitted to keep the log + // meaningful — only noteworthy holds (backoff, already running, etc.) appear. + skipped: SchedulerSkip[]; +}; + +// Resolve a channel's effective interval in minutes. undefined inherits the +// global default; 0 means auto-sync is disabled for the channel. +export function resolveIntervalMinutes( + config: ChannelConfig, + scheduler: SyncSchedulerSettings, +): number { + if (config.syncIntervalMinutes === undefined) { + return scheduler.defaultIntervalMinutes; + } + return config.syncIntervalMinutes; +} + +// Whether auto-sync is suppressed at `now` by the configured quiet-hours window. +// The window is [start, end) in local clock hours and may wrap past midnight +// (start=22, end=6). A null endpoint or a zero-length window means "never". +export function isInQuietHours( + now: number, + start: number | null, + end: number | null, +): boolean { + if (start === null || end === null || start === end) return false; + const hour = new Date(now).getHours(); + if (start < end) return hour >= start && hour < end; + // Wraps midnight. + return hour >= start || hour < end; +} + +// Exponential backoff (minutes) after N consecutive failures: base * 2^(N-1), +// capped at max. 0 failures => no backoff. +export function backoffMinutes( + failures: number, + scheduler: SyncSchedulerSettings, +): number { + if (failures <= 0) return 0; + const raw = scheduler.backoffBaseMinutes * 2 ** (failures - 1); + return Math.min(raw, scheduler.backoffMaxMinutes); +} + +// The epoch-ms instant a channel becomes eligible again after `failures` +// consecutive failures observed at `now`. +export function nextEligibleAfterFailure( + now: number, + failures: number, + scheduler: SyncSchedulerSettings, +): number { + return now + backoffMinutes(failures, scheduler) * 60_000; +} + +// Core selection. Evaluates each channel against the config gates, the elapsed +// interval, the backoff window and the active-job set, then orders the winners +// most-overdue first so a concurrency-capped tick services the stalest channels. +export function selectDueChannels(input: SelectDueInput): SelectDueResult { + const { channels, scheduler, state, activeSlugs, now } = input; + if (!scheduler.enabled) return { due: [], skipped: [] }; + if ( + isInQuietHours(now, scheduler.quietHoursStart, scheduler.quietHoursEnd) + ) { + return { due: [], skipped: [] }; + } + + const skipped: SchedulerSkip[] = []; + const due: { slug: string; overdueMs: number }[] = []; + + for (const { slug, config } of channels) { + if (!config.url) continue; // not auto-sync material; no noise in the log + if (config.excludeFromSync) continue; + + const interval = resolveIntervalMinutes(config, scheduler); + if (interval <= 0) continue; // per-channel disabled + + if (activeSlugs.has(slug)) { + skipped.push({ slug, reason: "already running" }); + continue; + } + + const cs = state.channels[slug]; + if (cs?.nextEligibleAt != null && now < cs.nextEligibleAt) { + skipped.push({ + slug, + reason: `backoff (${cs.consecutiveFailures} failures, until ${new Date( + cs.nextEligibleAt, + ).toISOString()})`, + }); + continue; + } + + const overdueMs = overdueAmount(config.lastSyncedAt, interval, now); + if (overdueMs === null) continue; // not yet due + due.push({ slug, overdueMs }); + } + + due.sort((a, b) => b.overdueMs - a.overdueMs); + return { due: due.map((d) => d.slug), skipped }; +} + +// A per-channel projection of the schedule for the observability panel. Pure: +// the editor status route feeds it live channels/settings/state. +export type ChannelScheduleView = { + slug: string; + name: string | null; + // True when this channel is eligible for auto-sync (scheduler on, has a url, + // not excluded, and a positive resolved interval). + autoSyncEligible: boolean; + intervalMinutes: number; // resolved; 0 = disabled + inheritsInterval: boolean; // using the global default vs a per-channel value + lastSyncedAt: string | null; + nextDueAt: number | null; // epoch ms; null when disabled or never synced + overdue: boolean; + consecutiveFailures: number; + nextEligibleAt: number | null; + lastOutcome: "ok" | "failed" | null; + lastOutcomeAt: number | null; +}; + +export function buildScheduleView(input: { + channels: ReadonlyArray<ChannelEntry>; + scheduler: SyncSchedulerSettings; + state: SchedulerState; + now: number; +}): ChannelScheduleView[] { + const { channels, scheduler, state, now } = input; + return channels.map(({ slug, config }) => { + const interval = resolveIntervalMinutes(config, scheduler); + const cs = state.channels[slug]; + const last = config.lastSyncedAt ?? null; + let nextDueAt: number | null = null; + let overdue = false; + if (interval > 0) { + if (!last) { + overdue = true; // never synced => due now + } else { + const parsed = Date.parse(last); + if (Number.isNaN(parsed)) { + overdue = true; + } else { + nextDueAt = parsed + interval * 60_000; + overdue = now >= nextDueAt; + } + } + } + return { + slug, + name: config.name ?? null, + autoSyncEligible: + scheduler.enabled && + !!config.url && + !config.excludeFromSync && + interval > 0, + intervalMinutes: interval, + inheritsInterval: config.syncIntervalMinutes === undefined, + lastSyncedAt: last, + nextDueAt, + overdue, + consecutiveFailures: cs?.consecutiveFailures ?? 0, + nextEligibleAt: cs?.nextEligibleAt ?? null, + lastOutcome: cs?.lastOutcome ?? null, + lastOutcomeAt: cs?.lastOutcomeAt ?? null, + }; + }); +} + +// How far past its interval a channel is, in ms. A never-synced channel is +// maximally overdue (Infinity) so it sorts first. Returns null when not yet due. +function overdueAmount( + lastSyncedAt: string | undefined, + intervalMinutes: number, + now: number, +): number | null { + if (!lastSyncedAt) return Number.POSITIVE_INFINITY; + const last = Date.parse(lastSyncedAt); + if (Number.isNaN(last)) return Number.POSITIVE_INFINITY; // unparseable => sync + const elapsed = now - last; + const intervalMs = intervalMinutes * 60_000; + return elapsed >= intervalMs ? elapsed - intervalMs : null; +} diff --git a/common/jobs/syncSchedulerState.ts b/common/jobs/syncSchedulerState.ts @@ -0,0 +1,161 @@ +import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; + +// Persistent state for the cron-driven sync scheduler. Unlike the in-memory job +// registry (which is wiped on every server restart), this survives restarts so +// failure backoff and the run log are durable. It is small (one entry per +// channel plus a bounded run log) and written atomically on each tick. +// +// This module is pure I/O + types — the scheduling logic lives in +// common/jobs/syncScheduler.ts and the orchestration in the editor tick route. + +export type ChannelSyncState = { + // Number of consecutive scheduled syncs that ended in failure/cancel. Reset + // to 0 on the first success. Drives the exponential backoff window. + consecutiveFailures: number; + // Epoch ms before which the channel is held back by backoff. null = eligible. + nextEligibleAt: number | null; + // The job id of the most recent scheduler-queued sync whose terminal outcome + // has not yet been folded into the counters above. null once reconciled. + lastQueuedJobId: string | null; + // Epoch ms the scheduler last queued a sync for this channel. + lastQueuedAt: number | null; + // Outcome of the last reconciled scheduled sync, for the observability panel. + lastOutcome: "ok" | "failed" | null; + lastOutcomeAt: number | null; +}; + +export type SchedulerSkip = { slug: string; reason: string }; + +export type SchedulerRun = { + at: number; + queued: string[]; + skipped: SchedulerSkip[]; +}; + +export type SchedulerState = { + channels: Record<string, ChannelSyncState>; + // Newest-first, bounded to SCHEDULER_RUN_LOG_LIMIT entries. + runs: SchedulerRun[]; +}; + +export const SCHEDULER_RUN_LOG_LIMIT = 50; + +export function emptyChannelSyncState(): ChannelSyncState { + return { + consecutiveFailures: 0, + nextEligibleAt: null, + lastQueuedJobId: null, + lastQueuedAt: null, + lastOutcome: null, + lastOutcomeAt: null, + }; +} + +export function emptySchedulerState(): SchedulerState { + return { channels: {}, runs: [] }; +} + +// Read the state file, tolerating a missing/corrupt file by returning an empty +// state. Defensive: any stored shape that doesn't match is coerced so a +// hand-edited or partially-written file can't crash a tick. +export async function readSchedulerState(paths: Paths): Promise<SchedulerState> { + let raw: unknown; + try { + raw = JSON.parse(await readFile(paths.schedulerStateFile, "utf8")); + } catch { + return emptySchedulerState(); + } + if (!raw || typeof raw !== "object") return emptySchedulerState(); + const r = raw as Record<string, unknown>; + const channels: Record<string, ChannelSyncState> = {}; + if (r.channels && typeof r.channels === "object") { + for (const [slug, value] of Object.entries( + r.channels as Record<string, unknown>, + )) { + channels[slug] = coerceChannelState(value); + } + } + const runs: SchedulerRun[] = Array.isArray(r.runs) + ? r.runs.map(coerceRun).slice(0, SCHEDULER_RUN_LOG_LIMIT) + : []; + return { channels, runs }; +} + +// Write the state atomically (tmp file + rename), creating the .scheduler dir on +// first use. The run log is trimmed to the limit on the way out. +export async function writeSchedulerState( + paths: Paths, + state: SchedulerState, +): Promise<void> { + const out: SchedulerState = { + channels: state.channels, + runs: state.runs.slice(0, SCHEDULER_RUN_LOG_LIMIT), + }; + await mkdir(path.dirname(paths.schedulerStateFile), { recursive: true }); + const tmp = `${paths.schedulerStateFile}.tmp-${process.pid}`; + await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); + await rename(tmp, paths.schedulerStateFile); +} + +// Get (or lazily create) the mutable state entry for a channel. +export function channelState( + state: SchedulerState, + slug: string, +): ChannelSyncState { + let entry = state.channels[slug]; + if (!entry) { + entry = emptyChannelSyncState(); + state.channels[slug] = entry; + } + return entry; +} + +// Prepend a run to the log and trim to the limit. +export function recordRun(state: SchedulerState, run: SchedulerRun): void { + state.runs.unshift(run); + if (state.runs.length > SCHEDULER_RUN_LOG_LIMIT) { + state.runs.length = SCHEDULER_RUN_LOG_LIMIT; + } +} + +function num(value: unknown): number | null { + return typeof value === "number" && Number.isFinite(value) ? value : null; +} + +function coerceChannelState(value: unknown): ChannelSyncState { + const base = emptyChannelSyncState(); + if (!value || typeof value !== "object") return base; + const r = value as Record<string, unknown>; + return { + consecutiveFailures: Math.max(0, num(r.consecutiveFailures) ?? 0), + nextEligibleAt: num(r.nextEligibleAt), + lastQueuedJobId: + typeof r.lastQueuedJobId === "string" ? r.lastQueuedJobId : null, + lastQueuedAt: num(r.lastQueuedAt), + lastOutcome: + r.lastOutcome === "ok" || r.lastOutcome === "failed" + ? r.lastOutcome + : null, + lastOutcomeAt: num(r.lastOutcomeAt), + }; +} + +function coerceRun(value: unknown): SchedulerRun { + const r = (value ?? {}) as Record<string, unknown>; + const queued = Array.isArray(r.queued) + ? r.queued.filter((x): x is string => typeof x === "string") + : []; + const skipped: SchedulerSkip[] = Array.isArray(r.skipped) + ? r.skipped.flatMap((x) => { + if (!x || typeof x !== "object") return []; + const s = x as Record<string, unknown>; + if (typeof s.slug !== "string" || typeof s.reason !== "string") { + return []; + } + return [{ slug: s.slug, reason: s.reason }]; + }) + : []; + return { at: num(r.at) ?? 0, queued, skipped }; +} diff --git a/common/lib/channelConfig.ts b/common/lib/channelConfig.ts @@ -24,6 +24,14 @@ export type ChannelConfig = { lastFullDownloadAt?: string; excludeFromBuild?: boolean; excludeFromSync?: boolean; + // Auto-sync cadence for the scheduled sync system (common/jobs/syncScheduler.ts). + // The cron-driven scheduler syncs this channel when + // `now - lastSyncedAt >= syncIntervalMinutes`. Semantics: + // undefined -> inherit the global SiteSettings default interval + // 0 -> auto-sync disabled for this channel (still manually syncable) + // > 0 -> sync this often (clamped to [SYNC_INTERVAL_MIN/MAX_MINUTES]) + // A missing `url` or `excludeFromSync` also disables auto-sync. + syncIntervalMinutes?: 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. @@ -42,6 +50,11 @@ export type ChannelConfig = { export const CHANNEL_SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS = 600; +// Bounds for a channel's auto-sync interval. A nonzero value is clamped into +// this window; 0 is preserved as "disabled". Max is ~31 days. +export const SYNC_INTERVAL_MIN_MINUTES = 1; +export const SYNC_INTERVAL_MAX_MINUTES = 44640; + export const AUDIO_CHECK_INTERVAL_DEFAULT_SECONDS = 60; export const AUDIO_CHECK_INTERVAL_MIN_SECONDS = 10; export const AUDIO_CHECK_INTERVAL_MAX_SECONDS = 600; @@ -113,6 +126,22 @@ export function parseChannelConfig(raw: unknown): ChannelConfig | null { if (typeof r.excludeFromSync === "boolean") { config.excludeFromSync = r.excludeFromSync; } + if ( + typeof r.syncIntervalMinutes === "number" && + Number.isFinite(r.syncIntervalMinutes) && + r.syncIntervalMinutes >= 0 + ) { + // 0 is a sentinel ("auto-sync disabled") preserved as-is; any other value + // is clamped into the supported window. + config.syncIntervalMinutes = + r.syncIntervalMinutes === 0 + ? 0 + : clampInt( + r.syncIntervalMinutes, + SYNC_INTERVAL_MIN_MINUTES, + SYNC_INTERVAL_MAX_MINUTES, + ); + } if (typeof r.skipLiveDownloads === "boolean") { config.skipLiveDownloads = r.skipLiveDownloads; } diff --git a/common/lib/paths.ts b/common/lib/paths.ts @@ -15,6 +15,10 @@ export type Paths = { // the produced transcript.json live under workerScratchDir/<remoteJobId>/ and // are cleaned up once the requesting instance pulls the result (or cancels). workerScratchDir: string; + // Persistent state for the cron-driven sync scheduler (per-channel backoff + + // a rolling run log). Survives restarts, unlike the in-memory job registry. + // See common/jobs/syncSchedulerState.ts. + schedulerStateFile: string; lmdbPath: string; exportDir: string; // The dir Next.js serves at "/" (also holds checked-in static assets). For @@ -73,6 +77,7 @@ export function getPaths(): Paths { sitesDir: process.env.SITES_DIR ?? path.join(transcriptsDir, "sites"), jobsDir: path.join(transcriptsDir, ".jobs"), workerScratchDir: path.join(transcriptsDir, ".worker-scratch"), + schedulerStateFile: path.join(transcriptsDir, ".scheduler", "state.json"), lmdbPath: path.join(transcriptsDir, "index.mdb"), exportDir, exportPublicDir, diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -1,6 +1,7 @@ import fs from "node:fs"; import path from "node:path"; import { getPaths } from "./paths"; +import { SYNC_INTERVAL_MAX_MINUTES } from "./channelConfig"; import { type AppInstanceConfig, DEFAULT_TRANSCRIBE_ARGS, @@ -85,6 +86,11 @@ export type SiteSettings = { // stay live without a manual reload. Mounted globally; pauses while the tab is // hidden. 0 disables passive refresh entirely. See AUTO_REFRESH_INTERVAL_*. autoRefreshIntervalSeconds: number; + // Global configuration for the scheduled (cron-driven) channel sync system. + // The per-channel cadence lives on ChannelConfig.syncIntervalMinutes; this + // block holds the defaults and guard rails the scheduler applies across all + // channels. See common/jobs/syncScheduler.ts. + syncScheduler: SyncSchedulerSettings; // Default social links applied to every site that doesn't define its own. // A site inherits these unless its site.json carries an explicit // `socialLinks` array — see Site.socialLinks / resolveSocialLinks in @@ -93,6 +99,27 @@ export type SiteSettings = { socialLinks: SocialLink[]; }; +export type SyncSchedulerSettings = { + // Master switch. When false, a tick selects nothing (manual sync still works). + enabled: boolean; + // Fallback cadence (minutes) for channels with no per-channel override. + defaultIntervalMinutes: number; + // Cap on sync jobs running/queued at once. A tick queues at most + // (cap - currently-active) channels; the rest roll to the next tick. This is + // also the stagger mechanism that keeps a big due-batch from hitting the + // source all at once. + maxConcurrentSyncs: number; + // Optional local-clock quiet window during which auto-sync is suppressed. + // Both null = always allowed. The window may wrap past midnight + // (e.g. start=22, end=6). Hours are [0,23]; the window is [start, end). + quietHoursStart: number | null; + quietHoursEnd: number | null; + // Failure backoff bounds. After N consecutive failed scheduled syncs a + // channel waits min(base * 2^(N-1), max) minutes before it's eligible again. + backoffBaseMinutes: number; + backoffMaxMinutes: number; +}; + export type SocialLink = { label: string; url: string; @@ -139,6 +166,84 @@ export const TRANSCRIPT_PAGE_DEFAULT_BYTES = 8 * 1024 * 1024; export const DEFAULT_ADMIN_TITLE = "Transcript Browser Admin"; +// Sync-scheduler bounds + defaults. Default cadence is daily; concurrency is +// conservative so a tick doesn't fan out into the source provider all at once. +export const SYNC_SCHEDULER_DEFAULT_INTERVAL_MINUTES = 1440; +export const SYNC_SCHEDULER_MAX_CONCURRENT_DEFAULT = 2; +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 function defaultSyncScheduler(): SyncSchedulerSettings { + return { + enabled: false, + defaultIntervalMinutes: SYNC_SCHEDULER_DEFAULT_INTERVAL_MINUTES, + maxConcurrentSyncs: SYNC_SCHEDULER_MAX_CONCURRENT_DEFAULT, + quietHoursStart: null, + quietHoursEnd: null, + backoffBaseMinutes: SYNC_SCHEDULER_BACKOFF_BASE_DEFAULT_MINUTES, + backoffMaxMinutes: SYNC_SCHEDULER_BACKOFF_MAX_DEFAULT_MINUTES, + }; +} + +function clampHourOrNull(value: unknown): number | null { + if (typeof value !== "number" || !Number.isFinite(value)) return null; + const n = Math.floor(value); + if (n < 0 || n > 23) return null; + return n; +} + +function clampPositiveInt(value: unknown, fallback: number, max: number): number { + const n = + typeof value === "number" && Number.isFinite(value) + ? Math.floor(value) + : fallback; + if (n < 1) return 1; + if (n > max) return max; + return n; +} + +// Coerce a raw settings.syncScheduler value into a clean SyncSchedulerSettings, +// falling back to defaults for missing/ill-typed fields. Quiet hours are only +// honored when BOTH endpoints are valid hours; otherwise the window is cleared. +export function sanitizeSyncScheduler(value: unknown): SyncSchedulerSettings { + const d = defaultSyncScheduler(); + if (!value || typeof value !== "object") return d; + const r = value as Record<string, unknown>; + const start = clampHourOrNull(r.quietHoursStart); + const end = clampHourOrNull(r.quietHoursEnd); + const backoffBase = clampPositiveInt( + r.backoffBaseMinutes, + d.backoffBaseMinutes, + SYNC_INTERVAL_MAX_MINUTES, + ); + return { + enabled: r.enabled === true, + defaultIntervalMinutes: clampPositiveInt( + r.defaultIntervalMinutes, + d.defaultIntervalMinutes, + SYNC_INTERVAL_MAX_MINUTES, + ), + maxConcurrentSyncs: clampPositiveInt( + r.maxConcurrentSyncs, + d.maxConcurrentSyncs, + SYNC_SCHEDULER_MAX_CONCURRENT_MAX, + ), + quietHoursStart: start !== null && end !== null ? start : null, + quietHoursEnd: start !== null && end !== null ? end : null, + backoffBaseMinutes: backoffBase, + // Cap can't sit below the base, or backoff would never grow. + backoffMaxMinutes: Math.max( + backoffBase, + clampPositiveInt( + r.backoffMaxMinutes, + d.backoffMaxMinutes, + SYNC_INTERVAL_MAX_MINUTES, + ), + ), + }; +} + function defaults(): SiteSettings { return { adminTitle: DEFAULT_ADMIN_TITLE, @@ -153,6 +258,7 @@ function defaults(): SiteSettings { skipLiveDownloads: true, reportDebouncePreset: DEFAULT_REPORT_DEBOUNCE_PRESET, autoRefreshIntervalSeconds: AUTO_REFRESH_INTERVAL_DEFAULT_SECONDS, + syncScheduler: defaultSyncScheduler(), socialLinks: [], }; } @@ -302,6 +408,7 @@ export function getSettings(): SiteSettings { merged.autoRefreshIntervalSeconds = clampAutoRefreshIntervalSeconds( merged.autoRefreshIntervalSeconds, ); + merged.syncScheduler = sanitizeSyncScheduler(merged.syncScheduler); merged.socialLinks = parseSocialLinks(merged.socialLinks); // Workers. When the file predates the worker model (no `workers` key), // synthesize a default list from the (now-settled) active app + per-app @@ -469,6 +576,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> { autoRefreshIntervalSeconds: clampAutoRefreshIntervalSeconds( next.autoRefreshIntervalSeconds, ), + syncScheduler: sanitizeSyncScheduler(next.syncScheduler), socialLinks, }; const tmp = `${file}.tmp-${process.pid}`; diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Channels can sync automatically on a schedule.** Each channel gained an **Auto-sync** setting (channel editor → Source): *Default* (inherit the global cadence), *Off*, or a concrete interval (every 10m / 30m / hourly / 6h / 12h / daily / weekly), stored as `syncIntervalMinutes` in `config.json`. Inspired by the Laravel scheduler, a single lightweight cron heartbeat (`pnpm sync:tick`, an ~30-line client) POSTs to the editor's new `/api/scheduler/tick`, and the **server** decides which channels are due — a channel is due when `now − lastSyncedAt ≥ its interval`, so a missed tick (server down, machine asleep) simply runs at the next one with no catch-up storm. All work runs **inside the editor** through the existing job queue and per-channel lock, so a scheduled sync can't collide with a manual **Sync** click, shows up live on `/jobs`, and feeds the same transcription worker pool — no second process, no new file locks. Global controls live in **Settings → Sync scheduler**: a master **enable** (off by default), a **default interval**, a **max concurrent syncs** cap (a tick queues at most `cap − running` channels, most-overdue first, rolling the rest to the next tick — which both bounds load and staggers a large due-batch so it doesn't hit the source all at once), an optional **quiet-hours** window, and **failure backoff** (after N consecutive failures a channel waits `base·2^(N-1)` minutes, capped, before retrying). Channels already marked **Exclude from sync** never auto-sync. A new **Schedule** page (`/scheduler`) shows each channel's interval, last sync, next-due time, last outcome, and any active backoff, plus a recent-ticks log and a **Run scheduler now** button; the same data is at `GET /api/scheduler/status`. The cron client targets the editor's port (3001) by default and is hardenable with a `SYNC_TICK_TOKEN` bearer token for installs that expose the editor — see `SCHEDULED_SYNC.md`. - **Sites can link to each other.** A site's form gained a **Public URL** field (the absolute URL it's served at, e.g. `https://jeralyzer.com`) and a **Related sites** section. The export footer automatically links to every *other* site that has a Public URL, so filling these in is all that's needed for cross-site links; a site left without a URL is simply omitted from the lists. The **Related sites** editor lets a site pull closely-related siblings to the front under named groups (e.g. Jeralyzer featuring Rekietalyzer under "MTG drama") — add a group, give it an optional heading, and check which sibling sites belong; everything you don't feature falls into a trailing "Other sites" group on its own. Groups reorder with ↑/↓. The picker only lists sites that actually exist, and featured ids for sites that were since deleted are dropped on save (with a heads-up note). It's a subtle, secondary feature — see the matching note in the export changelog for how it renders. - **The editor refreshes itself on a timer so its data stays live without a manual reload.** Every page now passively re-fetches its own server-rendered data on a configurable interval — so the sidebar badges (active/running job counts, changelog dot), channel reports, and any other on-screen figures keep up to date on their own. It uses Next's `router.refresh()` (the same mechanism the jobs list already used) mounted once globally in the root layout, so it covers every page and the shared sidebar with no per-page wiring. To avoid wasting work when you're not looking, it **pauses entirely while the browser tab is hidden** and does **one immediate refresh the moment you return** to the tab (rather than waiting out the interval); it also skips a tick while a previous refresh is still settling, so refreshes can't pile up. The cadence is set in **Settings → Auto-refresh interval (seconds)**: default **5s** (clamped 1–600), or **0 to disable** passive refresh completely. This replaces the jobs page's old bespoke 2.5s auto-refresh (the `/jobs/active` page keeps its faster 1s progress-bar polling, which animates per-task bars without a full re-render). - **Transcription is now driven by configurable workers instead of one global engine.** The old single **App** dropdown in **Settings → Transcription** is replaced by a **Transcription workers** list. Each worker is **one processing slot** — one transcription at a time — with its own engine (whisper.cpp / chough / parakeet) and config, and a priority given by its position in the list (top = preferred). To run several in parallel, add more workers; a **Copy** button duplicates one (e.g. point two copies at the same chough `--server` for two togglable server slots). A batch ("Transcribe missing", bucket, bulk, single-video) hands each video — per task — to the highest-priority free worker, so a fast GPU worker and a slower CPU worker (e.g. parakeet on the GPU + chough on the CPU) run side by side instead of one engine doing everything. Total parallelism is the number of enabled workers; the old per-run **Concurrency** control and the global **Parallel transcriptions** setting are gone (add/remove workers, or disable/drain one, to change load). A pre-worker `settings.json` migrates automatically to one worker per slot of the previously-selected app (the old parallel-transcriptions count becomes that many enabled copies), plus a disabled worker for any other engine you had configured, so existing installs keep their parallelism. Scheduling is a single process-wide pool, so two batches can't oversubscribe the same GPU. One-slot-per-worker also means you can disable a single slot to free *some* of a CPU/GPU while the rest keep transcribing. diff --git a/editor/app/api/scheduler/status/route.ts b/editor/app/api/scheduler/status/route.ts @@ -0,0 +1,11 @@ +import { NextResponse } from "next/server"; +import { buildSchedulerStatusPayload } from "../../../scheduler/status"; + +export const dynamic = "force-dynamic"; + +// Read-only view for the "Sync schedule" panel: the resolved per-channel +// schedule (next due / last outcome / backoff) plus the recent tick log and the +// effective scheduler settings. Backs a passive UI poll, like /api/jobs/active. +export async function GET() { + return NextResponse.json(await buildSchedulerStatusPayload()); +} diff --git a/editor/app/api/scheduler/tick/route.ts b/editor/app/api/scheduler/tick/route.ts @@ -0,0 +1,16 @@ +import { NextResponse } from "next/server"; +import { isSchedulerRequestAuthorized } from "../../../scheduler/auth"; +import { runSchedulerTick } from "../../../scheduler/runTick"; + +export const dynamic = "force-dynamic"; + +// The cron-driven heartbeat. A lightweight external client (common/bin/sync-tick.ts) +// POSTs here on a schedule; the server decides which channels are due and queues +// them through the existing job machinery. See common/jobs/syncScheduler.ts. +export async function POST(req: Request) { + if (!isSchedulerRequestAuthorized(req)) { + return NextResponse.json({ error: "unauthorized" }, { status: 401 }); + } + const result = await runSchedulerTick(); + return NextResponse.json(result); +} diff --git a/editor/app/channels/actions.ts b/editor/app/channels/actions.ts @@ -20,7 +20,7 @@ import { listSiteIds, writeSite, } from "yt-dlp-transcript-common/lib/site"; -import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { activeSyncSlugs } from "yt-dlp-transcript-common/jobs/syncJobs"; import { CHANNEL_FORM_FIELDS, parseChannelForm, @@ -166,17 +166,7 @@ export type SyncAllResult = { export async function syncAllChannelsAction(): Promise<SyncAllResult> { const paths = getPaths(); const channels = await listChannels(paths); - const active = new Set( - getRegistry() - .list() - .filter( - (j) => - j.kind === "sync" && - (j.status === "queued" || j.status === "running") && - j.channelSlug, - ) - .map((j) => j.channelSlug as string), - ); + const active = activeSyncSlugs(); const queued: string[] = []; const skipped: { slug: string; reason: string }[] = []; for (const c of channels) { diff --git a/editor/app/channels/components/ChannelForm.tsx b/editor/app/channels/components/ChannelForm.tsx @@ -119,6 +119,7 @@ export function ChannelForm({ run sequentially. </span> </label> + <SyncIntervalField value={c?.syncIntervalMinutes} /> </Section> <Section title="Transcription"> <label className="flex flex-col gap-1 text-sm"> @@ -258,6 +259,51 @@ function CollapsibleSection({ ); } +const SYNC_INTERVAL_PRESETS: { value: string; label: string }[] = [ + { value: "", label: "Default (use global)" }, + { value: "0", label: "Off (never auto-sync)" }, + { value: "10", label: "Every 10 minutes" }, + { value: "30", label: "Every 30 minutes" }, + { value: "60", label: "Hourly" }, + { value: "360", label: "Every 6 hours" }, + { value: "720", label: "Every 12 hours" }, + { value: "1440", label: "Daily" }, + { value: "10080", label: "Weekly" }, +]; + +function SyncIntervalField({ value }: { value?: number }) { + const current = value != null ? String(value) : ""; + const known = SYNC_INTERVAL_PRESETS.some((p) => p.value === current); + return ( + <label className="flex flex-col gap-1 text-sm"> + <span className="font-medium">Auto-sync</span> + <select + name="syncIntervalMinutes" + defaultValue={current} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm" + > + {SYNC_INTERVAL_PRESETS.map((p) => ( + <option key={p.value} value={p.value}> + {p.label} + </option> + ))} + {/* Preserve a hand-set non-preset interval so saving doesn't change it. */} + {!known && current !== "" && ( + <option value={current}>{`Every ${current} minutes`}</option> + )} + </select> + <span className="text-xs text-zinc-500"> + How often the scheduled sync system auto-syncs this channel, when the + global scheduler is enabled in{" "} + <a href="/settings" className="underline"> + Settings + </a> + . Requires a URL; excluded-from-sync channels never auto-sync. + </span> + </label> + ); +} + function AudioCheckFields({ config }: { config?: ChannelConfig }) { const ac = config?.audioCheck; const enabled = ac?.enabled === true; diff --git a/editor/app/channels/components/parseChannelForm.ts b/editor/app/channels/components/parseChannelForm.ts @@ -6,6 +6,8 @@ import { AUDIO_CHECK_INTERVAL_MIN_SECONDS, AUDIO_CHECK_MAX_ROLLBACKS_MAX, AUDIO_CHECK_MAX_ROLLBACKS_MIN, + SYNC_INTERVAL_MAX_MINUTES, + SYNC_INTERVAL_MIN_MINUTES, type AudioCheckConfig, } from "yt-dlp-transcript-common/lib/channelConfig"; import { @@ -30,6 +32,7 @@ export const CHANNEL_FORM_FIELDS = [ "audioFormat", "keepSourceVideo", "ytdlpExtraArgs", + "syncIntervalMinutes", "sleepBetweenDownloadsSeconds", "audioCheck", ] as const satisfies ReadonlyArray<keyof ChannelConfig>; @@ -84,6 +87,27 @@ export function parseChannelForm(formData: FormData): ParsedChannelForm { sleepBetweenDownloadsSeconds = n; } + // Auto-sync cadence. The select emits "" for "inherit global default" + // (absent from config), "0" for "off", or a positive minute count. + const syncIntervalRaw = String( + formData.get("syncIntervalMinutes") ?? "", + ).trim(); + let syncIntervalMinutes: number | undefined; + if (syncIntervalRaw) { + const n = Number.parseInt(syncIntervalRaw, 10); + if ( + !Number.isFinite(n) || + n < 0 || + (n !== 0 && + (n < SYNC_INTERVAL_MIN_MINUTES || n > SYNC_INTERVAL_MAX_MINUTES)) + ) { + throw new Error( + `Auto-sync interval must be 0 (off) or ${SYNC_INTERVAL_MIN_MINUTES}–${SYNC_INTERVAL_MAX_MINUTES} minutes`, + ); + } + syncIntervalMinutes = n; + } + const audioCheckEnabled = formData.get("audioCheckEnabled") != null; let audioCheck: AudioCheckConfig | undefined; if (audioCheckEnabled) { @@ -147,6 +171,9 @@ export function parseChannelForm(formData: FormData): ParsedChannelForm { if (audioFormat) config.audioFormat = audioFormat; if (keepSourceVideo) config.keepSourceVideo = true; if (ytdlpExtraArgs?.length) config.ytdlpExtraArgs = ytdlpExtraArgs; + if (syncIntervalMinutes != null) { + config.syncIntervalMinutes = syncIntervalMinutes; + } if (sleepBetweenDownloadsSeconds != null) { config.sleepBetweenDownloadsSeconds = sleepBetweenDownloadsSeconds; } diff --git a/editor/app/layout.tsx b/editor/app/layout.tsx @@ -53,6 +53,7 @@ const NAV_GROUPS: NavGroup[] = [ { href: "/jobs", label: "Jobs", badgeKey: "jobs" }, { href: "/jobs/active", label: "Active", badgeKey: "running" }, { href: "/workers", label: "Workers" }, + { href: "/scheduler", label: "Schedule" }, { href: "/build", label: "Build" }, { href: "/actionable", label: "Actionable" }, ], diff --git a/editor/app/scheduler/auth.ts b/editor/app/scheduler/auth.ts @@ -0,0 +1,13 @@ +// Authorization for the scheduler tick endpoint. The editor admin surface is +// otherwise unauthenticated (trusted-network self-host), so this is opt-in +// hardening for installs that expose the editor: set SYNC_TICK_TOKEN in the +// server's environment AND in the cron script's environment, and the tick route +// will require a matching `Authorization: Bearer <token>` header. When the env +// var is unset, the route is open (consistent with the rest of the editor). +export function isSchedulerRequestAuthorized(req: Request): boolean { + const token = process.env.SYNC_TICK_TOKEN; + if (!token) return true; + const header = req.headers.get("authorization") ?? ""; + const m = /^Bearer\s+(.+)$/i.exec(header.trim()); + return m?.[1] === token; +} diff --git a/editor/app/scheduler/components/SchedulerView.tsx b/editor/app/scheduler/components/SchedulerView.tsx @@ -0,0 +1,253 @@ +"use client"; + +import { useCallback, useEffect, useRef, useState } from "react"; +import type { SchedulerStatusPayload } from "../status"; + +export function SchedulerView({ + initial, +}: { + initial: SchedulerStatusPayload; +}) { + const [data, setData] = useState<SchedulerStatusPayload>(initial); + const [busy, setBusy] = useState(false); + const [message, setMessage] = useState<string | null>(null); + const mounted = useRef(true); + + const refresh = useCallback(async () => { + try { + const res = await fetch("/api/scheduler/status", { cache: "no-store" }); + if (!res.ok) return; + const next = (await res.json()) as SchedulerStatusPayload; + if (mounted.current) setData(next); + } catch { + /* transient; the next poll retries */ + } + }, []); + + useEffect(() => { + mounted.current = true; + const id = setInterval(refresh, 5000); + return () => { + mounted.current = false; + clearInterval(id); + }; + }, [refresh]); + + const runNow = useCallback(async () => { + setBusy(true); + setMessage(null); + try { + const res = await fetch("/api/scheduler/tick", { method: "POST" }); + const body = (await res.json()) as { + queued?: string[]; + skipped?: unknown[]; + reason?: string; + }; + if (!res.ok) { + setMessage(`Tick failed (${res.status}).`); + } else { + const q = body.queued?.length ?? 0; + const s = body.skipped?.length ?? 0; + setMessage( + body.reason + ? `Tick: ${body.reason} (queued ${q}, skipped ${s}).` + : `Tick queued ${q}, skipped ${s}.`, + ); + } + await refresh(); + } catch (err) { + setMessage(`Tick failed: ${(err as Error).message}`); + } finally { + if (mounted.current) setBusy(false); + } + }, [refresh]); + + const { scheduler, channels, runs, now } = data; + const eligible = channels.filter((c) => c.autoSyncEligible); + + return ( + <div className="flex flex-col gap-5"> + <div className="flex flex-wrap items-center gap-3"> + <span + className={`text-sm rounded-full px-3 py-1 border ${ + scheduler.enabled + ? "border-green-300 bg-green-50 text-green-800 dark:border-green-800 dark:bg-green-950 dark:text-green-200" + : "border-zinc-300 bg-zinc-50 text-zinc-600 dark:border-zinc-700 dark:bg-zinc-900 dark:text-zinc-400" + }`} + > + {scheduler.enabled ? "Scheduler enabled" : "Scheduler disabled"} + </span> + <span className="text-sm text-zinc-500"> + {eligible.length} channel{eligible.length === 1 ? "" : "s"} auto-syncing + · default {formatInterval(scheduler.defaultIntervalMinutes)} · max{" "} + {scheduler.maxConcurrentSyncs} concurrent + </span> + <button + type="button" + onClick={runNow} + disabled={busy} + className="ml-auto px-3 py-1.5 rounded-md bg-zinc-900 dark:bg-zinc-100 text-zinc-100 dark:text-zinc-900 text-sm font-medium hover:opacity-90 disabled:opacity-50" + > + {busy ? "Running…" : "Run scheduler now"} + </button> + </div> + {message && ( + <p role="status" className="text-sm text-zinc-600 dark:text-zinc-300"> + {message} + </p> + )} + + <div className="overflow-x-auto rounded border border-zinc-200 dark:border-zinc-800"> + <table className="w-full text-sm"> + <thead className="bg-zinc-50 dark:bg-zinc-900 text-left text-zinc-500"> + <tr> + <th className="px-3 py-2 font-medium">Channel</th> + <th className="px-3 py-2 font-medium">Interval</th> + <th className="px-3 py-2 font-medium">Last synced</th> + <th className="px-3 py-2 font-medium">Next due</th> + <th className="px-3 py-2 font-medium">Status</th> + </tr> + </thead> + <tbody> + {channels.length === 0 && ( + <tr> + <td colSpan={5} className="px-3 py-4 text-zinc-500"> + No channels. + </td> + </tr> + )} + {channels.map((c) => ( + <tr + key={c.slug} + className="border-t border-zinc-100 dark:border-zinc-800" + > + <td className="px-3 py-2"> + <a href={`/channels/${c.slug}`} className="hover:underline"> + {c.name ?? c.slug} + </a> + </td> + <td className="px-3 py-2 text-zinc-600 dark:text-zinc-300"> + {c.intervalMinutes === 0 ? ( + <span className="text-zinc-400">off</span> + ) : ( + <> + {formatInterval(c.intervalMinutes)} + {c.inheritsInterval && ( + <span className="text-zinc-400"> (default)</span> + )} + </> + )} + </td> + <td className="px-3 py-2 text-zinc-600 dark:text-zinc-300"> + {c.lastSyncedAt ? formatAgo(c.lastSyncedAt, now) : "never"} + </td> + <td className="px-3 py-2 text-zinc-600 dark:text-zinc-300"> + {nextDueLabel(c, now)} + </td> + <td className="px-3 py-2">{statusBadge(c)}</td> + </tr> + ))} + </tbody> + </table> + </div> + + <section className="flex flex-col gap-2"> + <h2 className="text-sm font-semibold">Recent ticks</h2> + {runs.length === 0 ? ( + <p className="text-sm text-zinc-500">No ticks recorded yet.</p> + ) : ( + <ul className="flex flex-col gap-1 text-sm"> + {runs.slice(0, 15).map((r, i) => ( + <li + key={`${r.at}-${i}`} + className="flex flex-wrap gap-x-3 gap-y-0.5 text-zinc-600 dark:text-zinc-300" + > + <span className="tabular-nums text-zinc-500"> + {formatClock(r.at)} + </span> + <span> + queued {r.queued.length} + {r.queued.length > 0 && ( + <span className="text-zinc-400"> ({r.queued.join(", ")})</span> + )} + </span> + {r.skipped.length > 0 && ( + <span className="text-zinc-400"> + skipped {r.skipped.length} + </span> + )} + </li> + ))} + </ul> + )} + </section> + </div> + ); +} + +type ChannelRow = SchedulerStatusPayload["channels"][number]; + +function statusBadge(c: ChannelRow) { + if (c.intervalMinutes === 0 || !c.autoSyncEligible) { + return <span className="text-zinc-400">—</span>; + } + if (c.nextEligibleAt != null && c.nextEligibleAt > Date.now()) { + return ( + <span className="text-amber-700 dark:text-amber-300"> + backoff ({c.consecutiveFailures} fail + {c.consecutiveFailures === 1 ? "" : "s"}) + </span> + ); + } + if (c.lastOutcome === "failed") { + return <span className="text-red-700 dark:text-red-300">last failed</span>; + } + if (c.overdue) { + return <span className="text-blue-700 dark:text-blue-300">due</span>; + } + return <span className="text-green-700 dark:text-green-300">ok</span>; +} + +function nextDueLabel(c: ChannelRow, now: number): string { + if (c.intervalMinutes === 0) return "—"; + if (c.overdue) return "due now"; + if (c.nextDueAt == null) return "—"; + return `in ${formatDuration(c.nextDueAt - now)}`; +} + +function formatInterval(minutes: number): string { + if (minutes <= 0) return "off"; + if (minutes % 1440 === 0) { + const d = minutes / 1440; + return d === 1 ? "daily" : `every ${d}d`; + } + if (minutes % 60 === 0) { + const h = minutes / 60; + return h === 1 ? "hourly" : `every ${h}h`; + } + return `every ${minutes}m`; +} + +function formatDuration(ms: number): string { + const sec = Math.max(0, Math.round(ms / 1000)); + if (sec < 60) return `${sec}s`; + const min = Math.round(sec / 60); + if (min < 60) return `${min}m`; + const hr = Math.round(min / 60); + if (hr < 48) return `${hr}h`; + return `${Math.round(hr / 24)}d`; +} + +function formatAgo(iso: string, now: number): string { + const t = Date.parse(iso); + if (Number.isNaN(t)) return iso; + return `${formatDuration(now - t)} ago`; +} + +function formatClock(ms: number): string { + if (!ms) return "—"; + return new Date(ms).toLocaleTimeString([], { + hour: "2-digit", + minute: "2-digit", + }); +} diff --git a/editor/app/scheduler/page.tsx b/editor/app/scheduler/page.tsx @@ -0,0 +1,33 @@ +import type { Metadata } from "next"; +import Link from "next/link"; +import { buildSchedulerStatusPayload } from "./status"; +import { SchedulerView } from "./components/SchedulerView"; + +export const dynamic = "force-dynamic"; + +export const metadata: Metadata = { title: "Schedule" }; + +export default async function SchedulerPage() { + const initial = await buildSchedulerStatusPayload(); + return ( + <div className="flex flex-col gap-4"> + <div className="flex items-center justify-between"> + <h1 className="text-2xl font-semibold">Sync schedule</h1> + <Link href="/settings" className="text-sm underline"> + Scheduler settings + </Link> + </div> + <p className="text-sm text-zinc-500"> + Per-channel auto-sync cadence, driven by a cron heartbeat + (<code>pnpm sync:tick</code>). Set a channel&apos;s cadence under its + Auto-sync field; global controls (enable, default interval, concurrency, + quiet hours, backoff) live in{" "} + <Link href="/settings" className="underline"> + Settings + </Link> + . + </p> + <SchedulerView initial={initial} /> + </div> + ); +} diff --git a/editor/app/scheduler/runTick.ts b/editor/app/scheduler/runTick.ts @@ -0,0 +1,186 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { listChannels } 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 { syncAction } from "../channels/[slug]/pipelineActions"; + +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[]; + // 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<SchedulerTickResult> { + 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: [], + 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: [], + skipped: [], + }; + } + + const channels: ChannelEntry[] = (await listChannels(paths)).map((c) => ({ + slug: c.slug, + config: c.config, + })); + const { due, skipped } = selectDueChannels({ + channels, + scheduler, + state, + activeSlugs, + now, + }); + + // 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[] = []; + for (const slug of toQueue) { + const result = 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, + ); + recordRun(state, { at: now, queued, skipped }); + await writeSchedulerState(paths, state); + + return { + ok: true, + ranAt: now, + reason: quiet ? "quiet hours" : undefined, + running, + queued, + 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; + } +} diff --git a/editor/app/scheduler/status.ts b/editor/app/scheduler/status.ts @@ -0,0 +1,41 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { listChannels } from "yt-dlp-transcript-common/controller/channels"; +import { + getSettings, + type SyncSchedulerSettings, +} from "yt-dlp-transcript-common/lib/settings"; +import { + buildScheduleView, + type ChannelScheduleView, +} from "yt-dlp-transcript-common/jobs/syncScheduler"; +import { + readSchedulerState, + type SchedulerRun, +} from "yt-dlp-transcript-common/jobs/syncSchedulerState"; + +export type SchedulerStatusPayload = { + now: number; + scheduler: SyncSchedulerSettings; + channels: ChannelScheduleView[]; + runs: SchedulerRun[]; +}; + +// Assemble the read-only "Sync schedule" view. Shared by the SSR page and the +// /api/scheduler/status poll so they never drift. +export async function buildSchedulerStatusPayload(): Promise<SchedulerStatusPayload> { + const paths = getPaths(); + const settings = getSettings(); + const state = await readSchedulerState(paths); + const now = Date.now(); + const channels = (await listChannels(paths)).map((c) => ({ + slug: c.slug, + config: c.config, + })); + const view = buildScheduleView({ + channels, + scheduler: settings.syncScheduler, + state, + now, + }); + return { now, scheduler: settings.syncScheduler, channels: view, runs: state.runs }; +} diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -111,6 +111,26 @@ export async function saveSettingsAction( }; } + // Sync-scheduler block. Values are clamped by sanitizeSyncScheduler inside + // writeSettings, so we only coerce here (NaN/blank fall back to defaults). + const intOrNaN = (key: string): number => + Number.parseInt(String(formData.get(key) ?? "").trim(), 10); + const hourOrNull = (key: string): number | null => { + const raw = String(formData.get(key) ?? "").trim(); + if (!raw) return null; + const n = Number.parseInt(raw, 10); + return Number.isFinite(n) ? n : null; + }; + const syncScheduler = { + enabled: formData.get("syncSchedulerEnabled") === "on", + defaultIntervalMinutes: intOrNaN("syncSchedulerDefaultIntervalMinutes"), + maxConcurrentSyncs: intOrNaN("syncSchedulerMaxConcurrentSyncs"), + quietHoursStart: hourOrNull("syncSchedulerQuietHoursStart"), + quietHoursEnd: hourOrNull("syncSchedulerQuietHoursEnd"), + backoffBaseMinutes: intOrNaN("syncSchedulerBackoffBaseMinutes"), + backoffMaxMinutes: intOrNaN("syncSchedulerBackoffMaxMinutes"), + }; + let socialInput: unknown; try { socialInput = JSON.parse(String(formData.get("socialLinksJson") ?? "[]")); @@ -153,6 +173,7 @@ export async function saveSettingsAction( skipLiveDownloads, reportDebouncePreset, autoRefreshIntervalSeconds: autoRefreshParsed, + syncScheduler, socialLinks, }; try { diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -144,6 +144,87 @@ export function SettingsForm({ initial, apps }: Props) { </span> </label> <fieldset className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded p-3"> + <legend className="px-1 text-sm font-medium">Sync scheduler</legend> + <p className="text-xs text-zinc-500"> + Auto-sync channels on a per-channel cadence (set under each channel&apos;s + Auto-sync). A cron heartbeat (<code>pnpm sync:tick</code>) pings the + server, which syncs whichever channels are due. See SCHEDULED_SYNC.md. + </p> + <label className="flex items-start gap-2 text-sm"> + <input + type="checkbox" + name="syncSchedulerEnabled" + defaultChecked={initial.syncScheduler.enabled} + className="mt-1" + /> + <span className="flex flex-col gap-1"> + <span className="font-medium">Enable scheduled auto-sync</span> + <span className="text-xs text-zinc-500"> + Master switch. When off, ticks are no-ops and only manual syncs + run. + </span> + </span> + </label> + <Field + label="Default interval (minutes)" + name="syncSchedulerDefaultIntervalMinutes" + defaultValue={String(initial.syncScheduler.defaultIntervalMinutes)} + type="number" + hint="Cadence for channels left on 'Default'. e.g. 60 = hourly, 1440 = daily." + /> + <Field + label="Max concurrent syncs" + name="syncSchedulerMaxConcurrentSyncs" + defaultValue={String(initial.syncScheduler.maxConcurrentSyncs)} + type="number" + hint="A tick queues at most (this − already-running) channels, most-overdue first; the rest roll to the next tick. Bounds load and staggers big batches." + /> + <div className="flex gap-3"> + <Field + label="Quiet hours start (0–23)" + name="syncSchedulerQuietHoursStart" + defaultValue={ + initial.syncScheduler.quietHoursStart != null + ? String(initial.syncScheduler.quietHoursStart) + : "" + } + type="number" + /> + <Field + label="Quiet hours end (0–23)" + name="syncSchedulerQuietHoursEnd" + defaultValue={ + initial.syncScheduler.quietHoursEnd != null + ? String(initial.syncScheduler.quietHoursEnd) + : "" + } + type="number" + /> + </div> + <p className="text-xs text-zinc-500 -mt-1"> + Optional local-clock window when auto-sync is suppressed (may wrap + past midnight, e.g. 22 → 6). Leave both blank to always allow. + </p> + <div className="flex gap-3"> + <Field + label="Backoff base (minutes)" + name="syncSchedulerBackoffBaseMinutes" + defaultValue={String(initial.syncScheduler.backoffBaseMinutes)} + type="number" + /> + <Field + label="Backoff max (minutes)" + name="syncSchedulerBackoffMaxMinutes" + defaultValue={String(initial.syncScheduler.backoffMaxMinutes)} + type="number" + /> + </div> + <p className="text-xs text-zinc-500 -mt-1"> + After N consecutive failed syncs a channel waits base·2^(N-1) minutes + (capped at max) before retrying. + </p> + </fieldset> + <fieldset className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded p-3"> <legend className="px-1 text-sm font-medium">Social links</legend> <p className="text-xs text-zinc-500"> Default social links shown in every site&apos;s footer. Each site can diff --git a/editor/e2e/scheduler.spec.ts b/editor/e2e/scheduler.spec.ts @@ -0,0 +1,99 @@ +import { test, expect } from "@playwright/test"; +import { writeFile } from "node:fs/promises"; +import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; +import { resetData, resolvePath, writeSettings } from "./helpers"; + +// End-to-end coverage for the cron-driven sync scheduler. Uses the two-slow +// channels (their mock sync hangs long enough to stay "running"), drives the +// real /api/scheduler/tick endpoint, and asserts which channels get queued. +// +// Covers the core rules of common/jobs/syncScheduler.selectDueChannels via the +// live server: elapsed-interval due selection, the not-yet-due gate, and the +// already-running dedup — plus the editor /scheduler panel. + +const SLOW_A_CONFIG = "test-transcripts/channels/slow-a/config.json"; +const SLOW_B_CONFIG = "test-transcripts/channels/slow-b/config.json"; + +async function writeChannelConfig(rel: string, config: ChannelConfig) { + await writeFile(resolvePath(rel), JSON.stringify(config, null, 2)); +} + +type TickResult = { + ok: boolean; + queued: string[]; + skipped: { slug: string; reason: string }[]; + reason?: string; +}; + +test("scheduler queues due channels, skips not-due, and dedups running ones", async ({ + page, +}) => { + await resetData("two-slow-channels"); + + // slow-a: inherits the default interval and has never synced -> due now. + await writeChannelConfig(SLOW_A_CONFIG, { + handling: "youtube", + name: "Slow A", + url: "https://www.youtube.com/@slow-a/videos", + ytdlpExtraArgs: ["--test-slow"], + }); + // slow-b: synced just now -> NOT due under the 60-minute default. + await writeChannelConfig(SLOW_B_CONFIG, { + handling: "youtube", + name: "Slow B", + url: "https://www.youtube.com/@slow-b/videos", + ytdlpExtraArgs: ["--test-slow"], + lastSyncedAt: new Date().toISOString(), + }); + + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + syncScheduler: { + enabled: true, + defaultIntervalMinutes: 60, + maxConcurrentSyncs: 2, + quietHoursStart: null, + quietHoursEnd: null, + backoffBaseMinutes: 30, + backoffMaxMinutes: 1440, + }, + }); + + // First tick: slow-a is due and gets queued; slow-b is not due. + const tick1 = await page.request.post("/api/scheduler/tick"); + expect(tick1.ok()).toBeTruthy(); + const r1 = (await tick1.json()) as TickResult; + expect(r1.queued).toEqual(["slow-a"]); + expect(r1.queued).not.toContain("slow-b"); + + // Second tick (no time passed): slow-a's sync is still running, so it's + // skipped as already-running; slow-b is still not due. Nothing new queues. + const tick2 = await page.request.post("/api/scheduler/tick"); + expect(tick2.ok()).toBeTruthy(); + const r2 = (await tick2.json()) as TickResult; + expect(r2.queued).toEqual([]); + expect(r2.skipped).toContainEqual({ + slug: "slow-a", + reason: "already running", + }); + + // The /scheduler panel reflects the enabled scheduler and lists the channels. + await page.goto("/scheduler"); + await expect(page.getByText("Scheduler enabled")).toBeVisible(); + await expect(page.getByRole("link", { name: "Slow A" })).toBeVisible(); + await expect( + page.getByRole("button", { name: "Run scheduler now" }), + ).toBeVisible(); +}); + +test("a disabled scheduler queues nothing", async ({ page }) => { + await resetData("two-slow-channels"); + // Default test settings leave syncScheduler.enabled false. + const tick = await page.request.post("/api/scheduler/tick"); + expect(tick.ok()).toBeTruthy(); + const r = (await tick.json()) as TickResult; + expect(r.queued).toEqual([]); + expect(r.reason).toBe("scheduler disabled"); +}); diff --git a/package.json b/package.json @@ -5,6 +5,7 @@ "type": "module", "scripts": { "build:index": "pnpm --filter yt-dlp-transcript-common exec tsx bin/build-index.ts", + "sync:tick": "pnpm --filter yt-dlp-transcript-common exec tsx bin/sync-tick.ts", "build:export": "pnpm --filter export run build", "build": "pnpm --filter export run build", "start:export": "node scripts/worktree.mjs run -- pnpm --filter export run start",