Archilyzer · Source

archilyzer

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

commit 90028d915e239c83fe25c84ec6a5ae03c7461c78
parent a071ee9ff7e410f37604dcb7abf91cbabb78f2d1
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 22 Jun 2026 15:09:24 -0400

Add automatic dynamic-priority queue (auto-transcribe / auto-download)

Two always-on runners that pick the next video to process across all
channels by a configurable priority policy, instead of running one
per-channel batch at a time.

- Pure policy tree (common/jobs/autoQueuePolicy.ts): strict / round-robin
  / weighted-fair (smooth weighted round-robin) groups; channel / platform
  / all leaves narrowable to a snapshot bucket; ancestor-aware max-workers
  caps; first-match-wins leaf assignment. Covered by tsx --test units.
- Runners (common/controller/autoRunner.ts): fill-to-capacity loop with
  non-destructive preemption (claim next freed slot). Transcription is
  bounded by the worker pool's eligible slots; downloads gate to one per
  platform to preserve per-platform politeness. Each runner is a drainable
  queueKey:"" registry job; fairness cursors persist in
  transcripts/.auto-queue/state.json.
- Shared per-video gating extracted to transcribeOneFromQueue.ts and reused
  by the existing whisper batch (behavior-preserving).
- Settings: autoQueue block in settings.json (sanitized like syncScheduler).
- UI: Pool -> Auto-queue page with a recursive policy-tree editor, live
  per-rule pending counts + recent-pick log, Start/Stop; /api/auto-queue
  status + control routes; instrumentation starts enabled runners at boot.

Verified: 16 resolver unit tests; 5 auto-queue e2e (strict fallback,
round-robin, UI round-trip, status, download ordering); regression green
(batch/drain, workers, whisper, settings).

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

Diffstat:
Acommon/controller/autoRunner.ts | 565+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/transcribeOneFromQueue.ts | 142+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/whisperBatch.ts | 107+++++++++++++++----------------------------------------------------------------
Acommon/jobs/autoQueuePolicy.test.ts | 338+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/autoQueuePolicy.ts | 363+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/jobs/autoQueueState.ts | 142+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/registry.ts | 4++++
Mcommon/lib/paths.ts | 5+++++
Mcommon/lib/settings.ts | 15+++++++++++++++
Meditor/CHANGELOG.md | 1+
Aeditor/app/api/auto-queue/control/route.ts | 44++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/api/auto-queue/status/route.ts | 11+++++++++++
Meditor/app/api/test/invalidate-cache/route.ts | 6++++++
Aeditor/app/auto-queue/actions.ts | 70++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/auto-queue/components/AutoQueueView.tsx | 246+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/auto-queue/components/PolicyTreeEditor.tsx | 494+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/auto-queue/page.tsx | 60++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/auto-queue/status.ts | 49+++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/layout.tsx | 1+
Meditor/app/settings/actions.ts | 5+++++
Aeditor/e2e/auto-queue.spec.ts | 409+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/instrumentation.ts | 13+++++++++++++
Mpnpm-workspace.yaml | 7+++++++
23 files changed, 3010 insertions(+), 87 deletions(-)

diff --git a/common/controller/autoRunner.ts b/common/controller/autoRunner.ts @@ -0,0 +1,565 @@ +import { readdir } from "node:fs/promises"; +import type { Paths } from "../lib/paths"; +import { getPaths } from "../lib/paths"; +import { getSettings } from "../lib/settings"; +import { detectPlatform } from "../lib/platform"; +import { getWorkerPool } from "../jobs/workerPool"; +import { getRegistry } from "../jobs/registry"; +import { runManagedFunction } from "../jobs/streamCommand"; +import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; +import { makeTaskTracker } from "../jobs/taskHooks"; +import type { JobRunContext } from "../jobs/streamCommand"; +import { + type ActiveCounts, + type AutoQueuePolicy, + type ChannelWork, + type WorkPick, + buildPendingByLeaf, + selectNextWork, + DEFAULT_TRANSCRIBE_BUCKET, + DEFAULT_DOWNLOAD_BUCKET, +} from "../jobs/autoQueuePolicy"; +import { + type AutoQueueKind, + readAutoQueueState, + recordPick, + writeAutoQueueState, +} from "../jobs/autoQueueState"; +import { readChannelConfig } from "./channels"; +import { readChannelSnapshot } from "./channelSnapshot"; +import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; +import { findVideoSourceUrl } from "./undownloadedVideos"; +import { downloadOneManaged } from "../ytdlp/downloadOneManaged"; + +// The automatic priority-queue runners. Each kind (auto-transcribe / +// auto-download) is ONE long-lived registry job (queueKey "" so it runs in +// parallel to everything, with working cancel + drain). The job's body is a +// fill-to-capacity loop: it reads pending work from channel snapshots, asks the +// pure policy (autoQueuePolicy.selectNextWork) which video to do next across all +// channels, and runs it. +// +// Priority is preserved by capping the number of OUTSTANDING units to the count +// of eligible worker slots — never dumping a backlog into the worker pool's FIFO +// park queue (which would decide order by hardware, not policy). Each time a slot +// frees, the loop re-picks from fresh policy, so a higher-priority channel's +// video claims the next freed slot (non-destructive "preempt at next slot"). + +export const AUTO_TRANSCRIBE_KIND = "auto-transcribe"; +export const AUTO_DOWNLOAD_KIND = "auto-download"; + +export function autoRunnerJobKind(kind: AutoQueueKind): string { + return kind === "transcription" ? AUTO_TRANSCRIBE_KIND : AUTO_DOWNLOAD_KIND; +} + +// How long the loop sleeps when there's no work and nothing in flight, before +// re-reading snapshots (new downloads may have appeared, or a worker re-enabled). +const IDLE_POLL_MS = 3000; +// Re-read channel list + platforms at most this often (configs change rarely). +const CHANNEL_LIST_TTL_MS = 30_000; +// Cap on the session "recently completed" set, so a long-running server doesn't +// grow it unbounded. FIFO eviction; the snapshot catches up within ~1s anyway. +const COMPLETED_CAP = 5000; + +// --- Live status (for /api/auto-queue/status), per kind -------------------- + +export type AutoRunnerInFlight = { + videoId: string; + leafId: string; + channelSlug: string; + startedAt: number; +}; + +type RunnerLive = { + jobId: string; + inFlight: Map<string, AutoRunnerInFlight>; + active: ActiveCounts; + startedAt: number; +}; + +type AutoRunnerSingleton = { runners: Map<AutoQueueKind, RunnerLive> }; + +declare global { + // eslint-disable-next-line no-var + var __yttAutoRunner__: AutoRunnerSingleton | undefined; +} + +function getSingleton(): AutoRunnerSingleton { + if (!globalThis.__yttAutoRunner__) { + globalThis.__yttAutoRunner__ = { runners: new Map() }; + } + return globalThis.__yttAutoRunner__; +} + +export type AutoRunnerStatus = { + kind: AutoQueueKind; + running: boolean; + jobId: string | null; + startedAt: number | null; + inFlight: AutoRunnerInFlight[]; + activeByNode: ActiveCounts; +}; + +export function getAutoRunnerStatus(kind: AutoQueueKind): AutoRunnerStatus { + const live = getSingleton().runners.get(kind); + const running = !!live && getRegistry().get(live.jobId)?.status === "running"; + return { + kind, + running, + jobId: live?.jobId ?? null, + startedAt: live?.startedAt ?? null, + inFlight: live ? [...live.inFlight.values()] : [], + activeByNode: live ? { ...live.active } : {}, + }; +} + +// --- Pending-work construction from snapshots ------------------------------ + +type ChannelMeta = { slug: string; platform: ReturnType<typeof detectPlatform> }; + +async function listChannelMeta(paths: Paths): Promise<ChannelMeta[]> { + let names: string[]; + try { + names = (await readdir(paths.channelsDir, { withFileTypes: true })) + .filter((e) => e.isDirectory()) + .map((e) => e.name); + } catch { + return []; + } + const out: ChannelMeta[] = []; + for (const slug of names) { + const config = await readChannelConfig(paths, slug); + if (!config) continue; + out.push({ slug, platform: detectPlatform(config.url) }); + } + return out; +} + +const BUCKET_NAMES = ["downloadedNoTranscript", "failedListed"] as const; + +// Read each channel's snapshot and project the buckets this runner kind cares +// about into ChannelWork, plus a videoId -> owning channel map (a platform/all +// leaf spans channels, so the pick needs the owner to locate the video dir). +async function buildChannelWork( + paths: Paths, + kind: AutoQueueKind, + meta: ReadonlyArray<ChannelMeta>, +): Promise<{ channels: ChannelWork[]; owner: Map<string, string> }> { + const channels: ChannelWork[] = []; + const owner = new Map<string, string>(); + for (const { slug, platform } of meta) { + const snap = await readChannelSnapshot(paths, slug); + if (!snap) continue; + const buckets: Record<string, string[]> = {}; + if (kind === "transcription") { + for (const name of BUCKET_NAMES) { + const ids = snap.buckets?.[name] ?? []; + buckets[name] = ids; + for (const id of ids) if (!owner.has(id)) owner.set(id, slug); + } + } else { + const ids = snap.undownloadedIds ?? []; + buckets[DEFAULT_DOWNLOAD_BUCKET] = ids; + for (const id of ids) if (!owner.has(id)) owner.set(id, slug); + } + channels.push({ slug, platform, buckets }); + } + return { channels, owner }; +} + +// Per-leaf count of pending (snapshot-derived) work for the status panel — what +// the operator sees as "cornbreadman: 12 pending". Uses the same matching as the +// runner so the numbers line up with what would actually be picked. +export async function computeLeafPendingCounts( + kind: AutoQueueKind, + paths: Paths = getPaths(), +): Promise<Record<string, number>> { + const policy = getSettings().autoQueue[kind]; + const defaultBucket = + kind === "transcription" ? DEFAULT_TRANSCRIBE_BUCKET : DEFAULT_DOWNLOAD_BUCKET; + const meta = await listChannelMeta(paths); + const { channels } = await buildChannelWork(paths, kind, meta); + const pending = buildPendingByLeaf(policy.root, channels, defaultBucket); + const counts: Record<string, number> = {}; + for (const [leafId, ids] of Object.entries(pending)) counts[leafId] = ids.length; + return counts; +} + +function removeIds( + pending: Record<string, string[]>, + exclude: ReadonlySet<string>, +): void { + if (exclude.size === 0) return; + for (const leafId of Object.keys(pending)) { + pending[leafId] = pending[leafId].filter((id) => !exclude.has(id)); + } +} + +// Count of worker slots that could run a transcription right now (enabled and +// not degraded — busy ones still count as capacity, just occupied). This is the +// ceiling on outstanding units so the loop never parks a backlog. +function eligibleSlots(): number { + return getWorkerPool() + .summary() + .filter((s) => s.state === "enabled" && !s.degraded).length; +} + +// --- The fill-to-capacity loop --------------------------------------------- + +async function runLoop( + kind: AutoQueueKind, + paths: Paths, + live: RunnerLive, + onLog: (line: string) => void, + signal: AbortSignal, + ctx: JobRunContext, +): Promise<void> { + const defaultBucket = + kind === "transcription" ? DEFAULT_TRANSCRIBE_BUCKET : DEFAULT_DOWNLOAD_BUCKET; + const tracker = makeTaskTracker(ctx, onLog); + const state = await readAutoQueueState(paths); + const kindState = state[kind]; + const runtime = kindState.runtime; + const completed = new Set<string>(); + const completedOrder: string[] = []; + let counters = { transcribed: 0, failed: 0, skipped: 0 }; + + // Wake coordination: a finishing unit (or an abort) wakes the loop out of its + // wait. pendingWake guards against a wake that fires between waits. + let waiter: (() => void) | null = null; + let pendingWake = false; + const wake = () => { + if (waiter) { + const w = waiter; + waiter = null; + w(); + } else { + pendingWake = true; + } + }; + const waitNext = (maxMs: number): Promise<void> => + new Promise<void>((resolve) => { + if (pendingWake || signal.aborted || ctx.drainSignal.aborted) { + pendingWake = false; + resolve(); + return; + } + let done = false; + const finish = () => { + if (done) return; + done = true; + clearTimeout(timer); + signal.removeEventListener("abort", finish); + ctx.drainSignal.removeEventListener("abort", finish); + if (waiter === finish) waiter = null; + resolve(); + }; + const timer = setTimeout(finish, maxMs); + timer.unref?.(); + waiter = finish; + signal.addEventListener("abort", finish, { once: true }); + ctx.drainSignal.addEventListener("abort", finish, { once: true }); + }); + + const markCompleted = (videoId: string) => { + if (completed.has(videoId)) return; + completed.add(videoId); + completedOrder.push(videoId); + if (completedOrder.length > COMPLETED_CAP) { + const old = completedOrder.shift(); + if (old) completed.delete(old); + } + }; + + const persist = () => { + void writeAutoQueueState(paths, state).catch(() => {}); + }; + + // Downloads have no worker pool; concurrency politeness is per-platform (one + // download per source at a time, matching the per-platform serial job queue). + // Track in-flight downloads per platform so the loop never launches a second + // download for a platform that's already busy. Keyed by platform name (or + // "unknown" for unrecognized hosts). Unused for transcription. + const platformInFlight = new Map<string, number>(); + const PER_PLATFORM_CAP = 1; + const platformKey = (slug: string, slugToPlatform: Map<string, string>) => + slugToPlatform.get(slug) ?? "unknown"; + + let metaCache: ChannelMeta[] = []; + let metaAt = 0; + + // Seed/refresh the global worker pool from current settings so eligibleSlots() + // and acquire() see the latest worker config (mirrors runWhisperBatch). After + // an e2e cache-reset the pool is recreated and re-seeds here. + getWorkerPool().reconfigure(); + + onLog(`Auto-${kind} runner started.`); + + while (!signal.aborted) { + // Stop if the job record is gone (e2e registry reset) or no longer running + // (hard cancel flips it to "cancelled"). Drain keeps it "running" and is + // handled via ctx.drainSignal below. + const rec = getRegistry().get(live.jobId); + if (!rec || rec.status !== "running") break; + + const policy: AutoQueuePolicy = getSettings().autoQueue[kind]; + if (!policy.enabled) { + onLog(`Auto-${kind} disabled — stopping runner.`); + break; + } + const drained = ctx.drainSignal.aborted; + // Transcription is throttled by the worker pool's slot count; download has no + // pool, so its real limiter is the per-platform gate below — leave the + // numeric cap at the policy max (or unbounded). + const target = drained + ? 0 + : kind === "transcription" + ? Math.min(eligibleSlots(), policy.maxWorkers ?? Number.POSITIVE_INFINITY) + : (policy.maxWorkers ?? Number.POSITIVE_INFINITY); + + if (live.inFlight.size >= target) { + // At capacity (or paused/drained). When draining/cancelling and nothing is + // left in flight, the run is complete. + if (live.inFlight.size === 0 && (drained || signal.aborted)) break; + await waitNext(IDLE_POLL_MS); + continue; + } + + // Refresh the (rarely-changing) channel list/platforms on a TTL. + const now = Date.now(); + if (now - metaAt > CHANNEL_LIST_TTL_MS || metaCache.length === 0) { + metaCache = await listChannelMeta(paths); + metaAt = now; + } + const slugToPlatform = new Map( + metaCache.map((m) => [m.slug, m.platform ?? "unknown"]), + ); + const { channels, owner } = await buildChannelWork(paths, kind, metaCache); + const pending = buildPendingByLeaf(policy.root, channels, defaultBucket); + const exclude = new Set<string>([...live.inFlight.keys(), ...completed]); + removeIds(pending, exclude); + // Download only: drop videos whose platform already has an in-flight + // download, so a busy platform yields to the next-priority free one. + if (kind === "download") { + const busy = new Set<string>(); + for (const [pf, n] of platformInFlight) { + if (n >= PER_PLATFORM_CAP) busy.add(pf); + } + if (busy.size > 0) { + for (const leafId of Object.keys(pending)) { + pending[leafId] = pending[leafId].filter((id) => { + const slug = owner.get(id); + return !(slug && busy.has(platformKey(slug, slugToPlatform))); + }); + } + } + } + + const pick = selectNextWork(policy.root, pending, runtime, live.active); + if (!pick) { + if (live.inFlight.size === 0 && (drained || signal.aborted)) break; + await waitNext(IDLE_POLL_MS); + continue; + } + + const channelSlug = owner.get(pick.videoId); + if (!channelSlug) { + // Shouldn't happen (every pending id came from a channel) — guard anyway. + markCompleted(pick.videoId); + continue; + } + + launchUnit({ kind, paths, pick, channelSlug, tracker, onLog, signal }) + .then((outcome) => { + if (outcome === "transcribed") counters.transcribed++; + else if (outcome === "failed") counters.failed++; + else counters.skipped++; + }) + .catch((err) => { + counters.failed++; + onLog(`Auto-${kind} unit error on ${pick.videoId}: ${String(err)}`); + }) + .finally(() => { + live.inFlight.delete(pick.videoId); + for (const id of pick.path) { + live.active[id] = Math.max(0, (live.active[id] ?? 0) - 1); + } + if (kind === "download") { + const pf = platformKey(channelSlug, slugToPlatform); + platformInFlight.set(pf, Math.max(0, (platformInFlight.get(pf) ?? 0) - 1)); + } + markCompleted(pick.videoId); + // The runner job has no channelSlug, so the registry's own + // snapshot-on-finish hook doesn't fire — request a regen here so the + // channel's pending counts refresh and a restart sees fresh buckets. + requestChannelSnapshot(paths, channelSlug); + wake(); + }); + + // Reserve the slot and account active counts up the path BEFORE looping, so + // the next pick sees the reservation (and ancestor caps reflect it). + live.inFlight.set(pick.videoId, { + videoId: pick.videoId, + leafId: pick.leafId, + channelSlug, + startedAt: Date.now(), + }); + for (const id of pick.path) live.active[id] = (live.active[id] ?? 0) + 1; + if (kind === "download") { + const pf = platformKey(channelSlug, slugToPlatform); + platformInFlight.set(pf, (platformInFlight.get(pf) ?? 0) + 1); + } + recordPick(kindState, { + at: Date.now(), + leafId: pick.leafId, + videoId: pick.videoId, + channelSlug, + }); + persist(); + onLog(`Auto-${kind}: ${channelSlug}/${pick.videoId} (leaf ${pick.leafId})`); + } + + // Drain remaining in-flight units (graceful stop / drain / cancel). + while (live.inFlight.size > 0) await waitNext(IDLE_POLL_MS); + persist(); + onLog( + `Auto-${kind} runner finished: ${counters.transcribed} done, ${counters.failed} failed, ${counters.skipped} skipped.`, + ); +} + +type LaunchArgs = { + kind: AutoQueueKind; + paths: Paths; + pick: WorkPick; + channelSlug: string; + tracker: ReturnType<typeof makeTaskTracker>; + onLog: (line: string) => void; + signal: AbortSignal; +}; + +// Run a single unit of work. Returns the outcome for counters. Hard cancel +// (signal) aborts an in-flight unit; drain is handled by the loop (it stops +// launching new units), so the drain signal is intentionally NOT forwarded here. +async function launchUnit(args: LaunchArgs): Promise< + "transcribed" | "skipped" | "failed" +> { + if (args.kind === "transcription") { + const res = await transcribeOneFromQueue({ + paths: args.paths, + channelSlug: args.channelSlug, + videoId: args.pick.videoId, + audioFilename: "audio.mp3", + strict: false, + tracker: args.tracker, + onLog: args.onLog, + signal: args.signal, + }); + return res.outcome; + } + // Download: resolve the video's source URL and run the same managed + // single-video download the import/download-missing paths use, so archive + // appends, auth/no-subs fallbacks and audio-integrity checks all apply. The + // per-platform in-flight gate in the loop keeps this to one download per + // platform at a time (matching the per-platform serial queue's politeness). + const config = await readChannelConfig(args.paths, args.channelSlug); + if (!config) return "skipped"; + const url = await findVideoSourceUrl( + args.paths, + args.channelSlug, + args.pick.videoId, + config, + ); + if (!url) { + args.onLog(`Auto-download: no source URL for ${args.pick.videoId} — skipped`); + return "skipped"; + } + const settings = getSettings(); + const task = args.tracker.start({ + id: args.pick.videoId, + label: args.pick.videoId, + kind: "download", + }); + try { + await downloadOneManaged({ + channelSlug: args.channelSlug, + channelConfig: config, + paths: args.paths, + videoUrl: url, + onLog: task.onLog, + signal: args.signal, + globalCookiesFromBrowser: settings.cookiesFromBrowser || undefined, + inlineTranscribeOnFallback: settings.inlineTranscribeOnFallback, + globalSkipLiveDownloads: settings.skipLiveDownloads, + appendArchive: true, + }); + return "transcribed"; // generic "success" — feeds the done counter + } catch (err) { + if (args.signal.aborted || (err as Error)?.name === "AbortError") { + return "skipped"; + } + args.onLog(`Auto-download failed for ${args.pick.videoId}: ${String(err)}`); + return "failed"; + } finally { + task.end(); + } +} + +// --- Lifecycle ------------------------------------------------------------- + +// Launch a runner job for `kind` if its policy is enabled and one isn't already +// running. Idempotent across HMR reloads / repeated calls (it checks the live +// registry). Returns the job id, or null when not started. +export async function startAutoRunner( + kind: AutoQueueKind, + paths: Paths = getPaths(), +): Promise<string | null> { + const policy = getSettings().autoQueue[kind]; + if (!policy.enabled) return null; + const singleton = getSingleton(); + const existing = singleton.runners.get(kind); + if (existing && getRegistry().get(existing.jobId)?.status === "running") { + return existing.jobId; // already running + } + + const live: RunnerLive = { + jobId: "", + inFlight: new Map(), + active: {}, + startedAt: Date.now(), + }; + + const result = await runManagedFunction({ + kind: autoRunnerJobKind(kind), + queueKey: "", // run in parallel to everything; not serialized + paths, + fn: async (onLog, signal, _setProgress, ctx) => { + live.jobId = ctx.jobId; + try { + await runLoop(kind, paths, live, onLog, signal, ctx); + } finally { + if (singleton.runners.get(kind) === live) { + singleton.runners.delete(kind); + } + } + }, + }); + if (!result.ok) return null; + live.jobId = result.jobId; + singleton.runners.set(kind, live); + return result.jobId; +} + +// Start every enabled runner. Called from the editor instrumentation hook at +// server startup. +export async function startAutoRunnersIfEnabled( + paths: Paths = getPaths(), +): Promise<void> { + await startAutoRunner("transcription", paths); + await startAutoRunner("download", paths); +} + +// Stop a runner (hard cancel its job). Returns true if one was running. +export function stopAutoRunner(kind: AutoQueueKind): boolean { + const live = getSingleton().runners.get(kind); + if (!live) return false; + return getRegistry().cancel(live.jobId); +} diff --git a/common/controller/transcribeOneFromQueue.ts b/common/controller/transcribeOneFromQueue.ts @@ -0,0 +1,142 @@ +import path from "node:path"; +import fs from "fs-extra"; +import type { Paths } from "../lib/paths"; +import { + isVideoDownloaded, + isVideoTranscribed, + readVideoFiles, +} from "../lib/videoStatus"; +import { transcribeWithWorker } from "./transcribeOne"; +import { TranscribeError } from "./transcribeError"; +import type { TaskTracker } from "../jobs/taskHooks"; + +const { pathExists, appendFile } = fs; + +// One video's gate-checks + transcription, factored out of runWhisperBatch so +// the per-channel batch (whisperBatch.ts) and the cross-channel auto-runner +// (autoRunner.ts) share IDENTICAL skip/fail rules. It acquires a worker lease +// via transcribeWithWorker (parking until a slot frees), records genuine +// failures to the channel's failed-transcriptions file, and never throws for the +// expected skip/fail cases — it reports them through the returned result. + +export type TranscribeOneOutcome = "transcribed" | "skipped" | "failed"; + +export type TranscribeOneResult = { + // True once the video cleared every gate and a transcription was attempted. + // Mirrors runWhisperBatch's `attempted++` (which counts paused/cancelled/ + // no-audio outcomes too — they increment attempted then fall into skipped). + attempted: boolean; + outcome: TranscribeOneOutcome; +}; + +export type TranscribeOneOptions = { + paths: Paths; + channelSlug: string; + videoId: string; + // Resolved audio filename (e.g. "audio.mp3"); under strict mode the video is + // skipped unless exactly this file exists. + audioFilename: string; + strict?: boolean; + // Optional in-memory failed-transcriptions set for the pre-attempt skip check + // (the batch pre-reads it once per run). When omitted, the on-disk failure + // list still governs future runs via the append on failure below. + failedSet?: ReadonlySet<string>; + tracker?: TaskTracker; + onLog?: (msg: string) => void; + signal?: AbortSignal; + drainSignal?: AbortSignal; +}; + +const skip = (): TranscribeOneResult => ({ attempted: false, outcome: "skipped" }); + +export async function transcribeOneFromQueue({ + paths, + channelSlug, + videoId, + audioFilename, + strict = false, + failedSet, + tracker, + onLog, + signal, + drainSignal, +}: TranscribeOneOptions): Promise<TranscribeOneResult> { + const log = onLog ?? ((m: string) => console.log(m)); + const channelDir = path.join(paths.channelsDir, channelSlug); + const dataDir = path.join(channelDir, "data"); + const failureListFile = path.join(channelDir, "failed-transcriptions"); + const videoPath = path.join(dataDir, videoId); + + if (signal?.aborted || drainSignal?.aborted) return skip(); + + if (failedSet?.has(videoId)) { + log(`Skipping previously failed transcription for ${videoId}`); + return skip(); + } + // A transcript already exists if either whisper has run (transcript.json) OR + // yt-dlp wrote auto/manual English subs (transcript.en.vtt, or a regional/auto + // fallback like en-US). Without the VTT branch, videos on a youtube-handling + // channel with only VTT auto-subs end up attempted -> fail with "no audio file + // found" -> get added to failed-transcriptions on every run. + const videoFiles = await readVideoFiles(videoPath); + if (isVideoTranscribed(videoFiles)) { + log(`Transcription for ${videoId} already exists`); + return skip(); + } + // No real audio file: a download/source problem, not a transcription + // candidate. Skip (don't fail) so a re-download lets it transcribe later. + if (!isVideoDownloaded(videoFiles)) { + log(`Skipping ${videoId}: not downloaded (no audio file)`); + return skip(); + } + if (strict && !(await pathExists(path.join(videoPath, audioFilename)))) { + log(`Skipping ${videoId}: no ${audioFilename} (strict mode)`); + return skip(); + } + + try { + // transcribeWithWorker acquires a pool worker (parking until one frees, or + // forever while all are disabled), retries across workers on transport + // failures, and tracks the per-video task. + const outcome = await transcribeWithWorker({ + paths, + videoDir: videoPath, + videoId, + audioFilename, + strictAudio: strict, + tracker, + taskId: videoId, + taskLabel: videoId, + onLog: log, + signal, + drainSignal, + }); + if (outcome === "paused") { + // Parakeet was paused with progress cached; leave the video untranscribed + // (not failed) so the next run resumes it. + log(`Transcribe ${videoId} paused (will resume next run)`); + return { attempted: true, outcome: "skipped" }; + } + return { attempted: true, outcome: "transcribed" }; + } catch (err) { + // Cancelled or drained while parked/in-flight: a skip, not a failure. + if ( + signal?.aborted || + drainSignal?.aborted || + (err as Error)?.name === "AbortError" + ) { + log(`Transcribe ${videoId} cancelled`); + return { attempted: true, outcome: "skipped" }; + } + // A missing audio file is a download/source problem, not a transcription + // failure — skip it (it transcribes once a good audio file exists) rather + // than permanently recording it in failed-transcriptions. + if (err instanceof TranscribeError && err.failureClass === "no-audio") { + log(`Skipping ${videoId}: no audio file (download/source problem)`); + return { attempted: true, outcome: "skipped" }; + } + log(`FAILED TO TRANSCRIBE ${videoId}: ${String(err)}`); + await appendFile(failureListFile, `${videoId}\n`); + return { attempted: true, outcome: "failed" }; + } +} diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts @@ -7,8 +7,7 @@ import { isVideoTranscribed, readVideoFiles, } from "../lib/videoStatus"; -import { transcribeWithWorker } from "./transcribeOne"; -import { TranscribeError } from "./transcribeError"; +import { transcribeOneFromQueue } from "./transcribeOneFromQueue"; import { pruneFailedTranscriptions } from "./failedTranscriptions"; import { resolveShardItems } from "./shard"; import { countNotYetTranscribed } from "./channels"; @@ -16,7 +15,7 @@ import { getWorkerPool } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; import type { JobProgress } from "../jobs/registry"; -const { pathExists, readdir, appendFile, readFile, ensureFile } = fs; +const { readdir, readFile, ensureFile } = fs; export type WhisperBatchOptions = { channelSlug: string; @@ -195,91 +194,25 @@ export async function runWhisperBatch({ // a worker slot frees. The pool's total free slots throttle the batch, and a // freed high-priority slot is handed to the oldest waiter. When all workers are // disabled, acquire() blocks indefinitely — the batch pauses (it does not fail) - // and resumes when a worker is re-enabled. + // and resumes when a worker is re-enabled. The per-video gating + transcribe is + // shared with the cross-channel auto-runner via transcribeOneFromQueue. const runOne = async (videoDir: string): Promise<void> => { - if (signal?.aborted || drainSignal?.aborted) { - skipped++; - return; - } - const videoPath = path.join(dataDir, videoDir); - if (failedSet.has(videoDir)) { - log(`Skipping previously failed transcription for ${videoDir}`); - skipped++; - return; - } - // A transcript already exists if either whisper has run (transcript.json) OR - // yt-dlp wrote auto/manual English subs (transcript.en.vtt, or a regional/auto - // fallback like en-US). Without the VTT branch, videos on a youtube-handling - // channel with only VTT auto-subs end up attempted -> fail with "no audio file - // found" -> get added to failed-transcriptions on every run. - const videoFiles = await readVideoFiles(videoPath); - if (isVideoTranscribed(videoFiles)) { - log(`Transcription for ${videoDir} already exists`); - skipped++; - return; - } - // No real audio file: a download/source problem, not a transcription - // candidate. Skip (don't fail) so a re-download lets it transcribe later. - if (!isVideoDownloaded(videoFiles)) { - log(`Skipping ${videoDir}: not downloaded (no audio file)`); - skipped++; - return; - } - if (strict && !(await pathExists(path.join(videoPath, resolvedAudioFilename)))) { - log(`Skipping ${videoDir}: no ${resolvedAudioFilename} (strict mode)`); - skipped++; - return; - } - - attempted++; - try { - // transcribeWithWorker acquires a pool worker (parking until one frees, or - // forever while all are disabled), retries across workers on transport - // failures, and tracks the per-video task. - const outcome = await transcribeWithWorker({ - paths, - videoDir: videoPath, - videoId: videoDir, - audioFilename: resolvedAudioFilename, - strictAudio: strict, - tracker, - taskId: videoDir, - taskLabel: videoDir, - onLog: log, - signal, - drainSignal, - }); - if (outcome === "paused") { - // Parakeet was paused with progress cached; leave the video untranscribed - // (not failed) so the next run resumes it. - log(`Transcribe ${videoDir} paused (will resume next run)`); - skipped++; - } else { - succeededCount++; - } - } catch (err) { - // Cancelled or drained while parked/in-flight: a skip, not a failure. - if ( - signal?.aborted || - drainSignal?.aborted || - (err as Error)?.name === "AbortError" - ) { - log(`Transcribe ${videoDir} cancelled`); - skipped++; - return; - } - // A missing audio file is a download/source problem, not a transcription - // failure — skip it (it transcribes once a good audio file exists) rather - // than permanently recording it in failed-transcriptions. - if (err instanceof TranscribeError && err.failureClass === "no-audio") { - log(`Skipping ${videoDir}: no audio file (download/source problem)`); - skipped++; - return; - } - log(`FAILED TO TRANSCRIBE ${videoDir}: ${String(err)}`); - await appendFile(failureListFile, `${videoDir}\n`); - failed++; - } + const res = await transcribeOneFromQueue({ + paths, + channelSlug, + videoId: videoDir, + audioFilename: resolvedAudioFilename, + strict, + failedSet, + tracker, + onLog: log, + signal, + drainSignal, + }); + if (res.attempted) attempted++; + if (res.outcome === "transcribed") succeededCount++; + else if (res.outcome === "failed") failed++; + else skipped++; }; await Promise.all(videoDirs.map(runOne)); diff --git a/common/jobs/autoQueuePolicy.test.ts b/common/jobs/autoQueuePolicy.test.ts @@ -0,0 +1,338 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + type AutoQueueGroup, + type AutoQueueRuntime, + type ChannelWork, + buildPendingByLeaf, + defaultAutoQueue, + emptyAutoQueueRuntime, + flattenLeaves, + sanitizeAutoQueue, + selectNextWork, +} from "./autoQueuePolicy"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test common/jobs/autoQueuePolicy.test.ts +// (or `node_modules/.bin/tsx --test common/jobs/autoQueuePolicy.test.ts` from the repo root) + +// Drain a policy: repeatedly selectNextWork, consuming the chosen video from +// the pending map, until nothing is left. Returns the leaf ids in served order. +function drain( + root: AutoQueueGroup, + pending: Record<string, string[]>, + runtime: AutoQueueRuntime = emptyAutoQueueRuntime(), + max = 1000, +): string[] { + const order: string[] = []; + for (let i = 0; i < max; i++) { + const pick = selectNextWork(root, pending, runtime); + if (!pick) break; + order.push(pick.leafId); + const ids = pending[pick.leafId]; + assert.equal(ids[0], pick.videoId, "pick should be head of leaf queue"); + ids.shift(); + } + return order; +} + +test("strict: drains highest-priority leaf first, then falls back", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { id: "corn", match: { type: "channel", value: "cornbreadman" } }, + { id: "hasan", match: { type: "channel", value: "hasanabi" } }, + ], + }; + const pending = { corn: ["c1", "c2"], hasan: ["h1", "h2", "h3"] }; + assert.deepEqual(drain(root, pending), [ + "corn", + "corn", + "hasan", + "hasan", + "hasan", + ]); +}); + +test("strict: a higher-priority arrival is served before lower (preemption-at-next-slot)", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { id: "corn", match: { type: "channel", value: "cornbreadman" } }, + { id: "hasan", match: { type: "channel", value: "hasanabi" } }, + ], + }; + const runtime = emptyAutoQueueRuntime(); + // Mid-"hasanabi batch": only hasan has work. + const pending: Record<string, string[]> = { corn: [], hasan: ["h1", "h2"] }; + let pick = selectNextWork(root, pending, runtime); + assert.equal(pick?.leafId, "hasan"); + pending.hasan.shift(); + // cornbreadman gains a video — next selection jumps to it. + pending.corn.push("c1"); + pick = selectNextWork(root, pending, runtime); + assert.equal(pick?.leafId, "corn"); +}); + +test("round-robin: alternates A B A B between non-empty children", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "round-robin", + children: [ + { id: "a", match: { type: "channel", value: "a" } }, + { id: "b", match: { type: "channel", value: "b" } }, + ], + }; + const pending = { a: ["a1", "a2", "a3"], b: ["b1", "b2", "b3"] }; + const order = drain(root, pending); + // SWRR with equal weights interleaves; first two picks cover both children. + assert.deepEqual(new Set(order.slice(0, 2)), new Set(["a", "b"])); + assert.equal(order.filter((x) => x === "a").length, 3); + assert.equal(order.filter((x) => x === "b").length, 3); + // No three-in-a-row of the same leaf while both have work. + assert.ok(!/(.)\1\1/.test(order.join(""))); +}); + +test("round-robin: skips an empty child without burning its turn", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "round-robin", + children: [ + { id: "a", match: { type: "channel", value: "a" } }, + { id: "b", match: { type: "channel", value: "b" } }, + ], + }; + const pending = { a: ["a1", "a2"], b: [] }; + assert.deepEqual(drain(root, pending), ["a", "a"]); +}); + +test("weighted-fair: honors the 2:1 weight ratio", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "weighted-fair", + children: [ + { id: "a", match: { type: "channel", value: "a" }, weight: 2 }, + { id: "b", match: { type: "channel", value: "b" }, weight: 1 }, + ], + }; + const pending = { + a: Array.from({ length: 20 }, (_, i) => `a${i}`), + b: Array.from({ length: 20 }, (_, i) => `b${i}`), + }; + const order = drain(root, pending, emptyAutoQueueRuntime(), 30); + // Over the first 30 picks, A should be served ~twice as often as B. + const a = order.slice(0, 30).filter((x) => x === "a").length; + const b = order.slice(0, 30).filter((x) => x === "b").length; + assert.equal(a + b, 30); + assert.ok(Math.abs(a - 20) <= 1, `expected ~20 A picks, got ${a}`); + assert.ok(Math.abs(b - 10) <= 1, `expected ~10 B picks, got ${b}`); +}); + +test("maxWorkers: a saturated leaf falls through to the next-priority sibling", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { id: "corn", match: { type: "channel", value: "corn" }, maxWorkers: 1 }, + { id: "hasan", match: { type: "channel", value: "hasan" } }, + ], + }; + const pending = { corn: ["c1", "c2"], hasan: ["h1"] }; + // corn already has 1 in-flight worker -> saturated -> hasan wins. + const pick = selectNextWork(root, pending, emptyAutoQueueRuntime(), { corn: 1 }); + assert.equal(pick?.leafId, "hasan"); +}); + +test("maxWorkers: a capped ancestor group blocks its whole subtree", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { + id: "vips", + mode: "strict", + maxWorkers: 2, + children: [ + { id: "corn", match: { type: "channel", value: "corn" } }, + { id: "hasan", match: { type: "channel", value: "hasan" } }, + ], + }, + { id: "rest", match: { type: "all" } }, + ], + }; + const pending = { corn: ["c1"], hasan: ["h1"], rest: ["r1"] }; + // The vips group has 2 in-flight (at its cap) -> skip it -> rest wins. + const pick = selectNextWork(root, pending, emptyAutoQueueRuntime(), { vips: 2 }); + assert.equal(pick?.leafId, "rest"); + // Path includes every ancestor id for the chosen leaf. + assert.deepEqual(pick?.path, ["root", "rest"]); +}); + +test("path: returns root..leaf ids for nested selection", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { + id: "g", + mode: "round-robin", + children: [{ id: "corn", match: { type: "channel", value: "corn" } }], + }, + ], + }; + const pick = selectNextWork(root, { corn: ["c1"] }, emptyAutoQueueRuntime()); + assert.deepEqual(pick?.path, ["root", "g", "corn"]); +}); + +test("selectNextWork: returns null when everything is empty", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [{ id: "a", match: { type: "channel", value: "a" } }], + }; + assert.equal(selectNextWork(root, { a: [] }, emptyAutoQueueRuntime()), null); +}); + +// --- buildPendingByLeaf ----------------------------------------------------- + +const CHANNELS: ChannelWork[] = [ + { + slug: "cornbreadman", + platform: "youtube", + buckets: { downloadedNoTranscript: ["c1", "c2"], failedListed: ["cf1"] }, + }, + { + slug: "hasanabi", + platform: "twitch", + buckets: { downloadedNoTranscript: ["h1"], failedListed: [] }, + }, +]; + +test("buildPendingByLeaf: channel + platform + all matchers", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { id: "corn", match: { type: "channel", value: "cornbreadman" } }, + { id: "tw", match: { type: "platform", value: "twitch" } }, + { id: "rest", match: { type: "all" } }, + ], + }; + const pending = buildPendingByLeaf(root, CHANNELS, "downloadedNoTranscript"); + assert.deepEqual(pending.corn, ["c1", "c2"]); + assert.deepEqual(pending.tw, ["h1"]); // hasanabi is twitch + assert.deepEqual(pending.rest, []); // both channels already claimed +}); + +test("buildPendingByLeaf: first-match-wins prevents double-claiming a video", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { id: "all1", match: { type: "all" } }, + { id: "corn", match: { type: "channel", value: "cornbreadman" } }, + ], + }; + const pending = buildPendingByLeaf(root, CHANNELS, "downloadedNoTranscript"); + assert.deepEqual(pending.all1.sort(), ["c1", "c2", "h1"]); + assert.deepEqual(pending.corn, []); // already claimed by the catch-all above it +}); + +test("buildPendingByLeaf: a bucket leaf draws from its named bucket", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { + id: "retry", + match: { type: "channel", value: "cornbreadman", bucket: "failedListed" }, + }, + { id: "corn", match: { type: "channel", value: "cornbreadman" } }, + ], + }; + const pending = buildPendingByLeaf(root, CHANNELS, "downloadedNoTranscript"); + assert.deepEqual(pending.retry, ["cf1"]); + assert.deepEqual(pending.corn, ["c1", "c2"]); +}); + +test("flattenLeaves: pre-order priority order", () => { + const root: AutoQueueGroup = { + id: "root", + mode: "strict", + children: [ + { + id: "g", + mode: "strict", + children: [ + { id: "a", match: { type: "all" } }, + { id: "b", match: { type: "all" } }, + ], + }, + { id: "c", match: { type: "all" } }, + ], + }; + assert.deepEqual( + flattenLeaves(root).map((l) => l.id), + ["a", "b", "c"], + ); +}); + +// --- sanitize --------------------------------------------------------------- + +test("sanitizeAutoQueue: empty/garbage -> defaults", () => { + assert.deepEqual(sanitizeAutoQueue(undefined), defaultAutoQueue()); + assert.deepEqual(sanitizeAutoQueue("nope"), defaultAutoQueue()); + const def = defaultAutoQueue(); + assert.equal(def.transcription.enabled, false); + assert.equal(def.transcription.root.children.length, 0); +}); + +test("sanitizeAutoQueue: coerces a hand-written tree, assigns missing ids", () => { + const raw = { + transcription: { + enabled: true, + maxWorkers: 3, + root: { + id: "root", + mode: "round-robin", + children: [ + { match: { type: "channel", value: "corn" }, weight: 2 }, + { mode: "bogus", children: [{ match: { type: "all" } }] }, + ], + }, + }, + }; + const s = sanitizeAutoQueue(raw); + assert.equal(s.transcription.enabled, true); + assert.equal(s.transcription.maxWorkers, 3); + assert.equal(s.transcription.root.mode, "round-robin"); + const [leaf, group] = s.transcription.root.children; + assert.ok(leaf.id, "missing id was assigned"); + assert.equal((leaf as { weight?: number }).weight, 2); + // An invalid mode falls back to strict. + assert.equal((group as AutoQueueGroup).mode, "strict"); + // download policy defaults in when absent. + assert.equal(s.download.enabled, false); +}); + +test("sanitizeAutoQueue: duplicate ids are de-duplicated", () => { + const raw = { + transcription: { + root: { + id: "dup", + mode: "strict", + children: [ + { id: "dup", match: { type: "all" } }, + { id: "dup", match: { type: "all" } }, + ], + }, + }, + }; + const s = sanitizeAutoQueue(raw); + const ids = [ + s.transcription.root.id, + ...s.transcription.root.children.map((c) => c.id), + ]; + assert.equal(new Set(ids).size, ids.length, "all ids unique after sanitize"); +}); diff --git a/common/jobs/autoQueuePolicy.ts b/common/jobs/autoQueuePolicy.ts @@ -0,0 +1,363 @@ +import type { Platform } from "../lib/platform"; + +// Pure, side-effect-free policy engine for the automatic priority queue. It +// decides WHICH pending video to process next, across all channels, from a +// configurable tree. The runner (common/controller/autoRunner.ts) does the I/O +// (reads snapshots to build the pending sets, acquires worker leases, persists +// fairness state). Keeping this pure makes the selection rules unit-testable +// without a server — exactly how syncScheduler.ts splits selectDueChannels from +// the editor tick route. +// +// The tree is a direct analog of Linux HTB / cgroup hierarchies: internal nodes +// (groups) carry a `mode` governing how their children compete; leaves carry a +// `match` describing which videos they own. A capped subtree (its in-flight +// worker count at its `maxWorkers` ceiling) reads as "no work" and the algorithm +// falls through to the next-priority sibling — exactly like HTB's class ceil. + +// --- Tree types ------------------------------------------------------------- + +export type AutoQueueMode = "strict" | "round-robin" | "weighted-fair"; + +export const AUTO_QUEUE_MODES: ReadonlyArray<AutoQueueMode> = [ + "strict", + "round-robin", + "weighted-fair", +]; + +export type AutoQueueMatchType = "channel" | "platform" | "all"; + +export type AutoQueueMatch = { + type: AutoQueueMatchType; + // Channel slug (type=channel) or platform name (type=platform). Ignored for + // type=all. A type=channel leaf with no value matches nothing. + value?: string; + // Optional snapshot bucket this leaf draws from, narrowing the default for the + // runner kind (transcription → downloadedNoTranscript, download → + // undownloadedIds). E.g. bucket="failedListed" prioritizes retries. + bucket?: string; +}; + +export type AutoQueueLeaf = { + id: string; + match: AutoQueueMatch; + // Relative share under a weighted-fair parent. Default 1. Ignored otherwise. + weight?: number; + // Optional ceiling on concurrent in-flight workers drawn from this leaf. + maxWorkers?: number | null; +}; + +export type AutoQueueGroup = { + id: string; + mode: AutoQueueMode; + children: AutoQueueNode[]; + weight?: number; + maxWorkers?: number | null; +}; + +export type AutoQueueNode = AutoQueueLeaf | AutoQueueGroup; + +export function isGroup(node: AutoQueueNode): node is AutoQueueGroup { + return Array.isArray((node as AutoQueueGroup).children); +} + +// --- Settings (persisted in settings.json under `autoQueue`) ---------------- + +export type AutoQueuePolicy = { + // Master switch for this runner (transcription / download independently). + enabled: boolean; + // Overall ceiling on concurrent in-flight workers for this runner. null = no + // runner-level cap (the worker pool / platform queues are the real throttle). + maxWorkers: number | null; + root: AutoQueueGroup; +}; + +export type AutoQueueSettings = { + transcription: AutoQueuePolicy; + download: AutoQueuePolicy; +}; + +// Default snapshot bucket each runner kind draws from when a leaf doesn't +// specify one. +export const DEFAULT_TRANSCRIBE_BUCKET = "downloadedNoTranscript"; +export const DEFAULT_DOWNLOAD_BUCKET = "undownloadedIds"; + +export const AUTO_QUEUE_MAX_WORKERS_MAX = 64; + +// --- Runtime fairness state (persisted best-effort by autoQueueState.ts) ---- + +// Smooth Weighted Round-Robin (SWRR, nginx's algorithm) current-weight per node +// id. Round-robin mode is SWRR with every weight forced to 1; weighted-fair uses +// each child's `weight`. Mutated in place by selectNextWork when a group makes a +// choice, so fairness advances only on an actual grant. +export type AutoQueueRuntime = { + currentWeights: Record<string, number>; +}; + +export function emptyAutoQueueRuntime(): AutoQueueRuntime { + return { currentWeights: {} }; +} + +// In-flight worker counts keyed by node id (a leaf and every ancestor on the +// path to it). The runner increments each id on grant and decrements on +// release; the resolver reads them to enforce maxWorkers ceilings. +export type ActiveCounts = Record<string, number>; + +// --- Pending-set construction (pure) --------------------------------------- + +// One channel's available work, grouped by snapshot bucket. The runner builds +// these from each channel's snapshot.json (already filtered to genuinely +// actionable ids), then hands them here so leaf matching stays pure/testable. +export type ChannelWork = { + slug: string; + platform: Platform | null; + // bucket name -> available video ids in that bucket + buckets: Record<string, string[]>; +}; + +// Pre-order (priority-order) flatten of every leaf in the tree. +export function flattenLeaves(node: AutoQueueNode): AutoQueueLeaf[] { + if (!isGroup(node)) return [node]; + const out: AutoQueueLeaf[] = []; + for (const child of node.children) out.push(...flattenLeaves(child)); + return out; +} + +function matchesChannel(match: AutoQueueMatch, ch: ChannelWork): boolean { + if (match.type === "all") return true; + if (match.type === "channel") return !!match.value && ch.slug === match.value; + if (match.type === "platform") { + return !!match.value && ch.platform === match.value; + } + return false; +} + +// Assign every available video to exactly ONE leaf: the first (highest-priority, +// pre-order) leaf whose match covers the video's channel and whose bucket +// contains it. First-match-wins prevents the same video being claimed — and thus +// double-processed — by two overlapping leaves (e.g. a channel leaf and an `all` +// catch-all). Returns leafId -> available video ids (in their snapshot order). +export function buildPendingByLeaf( + root: AutoQueueNode, + channels: ReadonlyArray<ChannelWork>, + defaultBucket: string, +): Record<string, string[]> { + const leaves = flattenLeaves(root); + const pending: Record<string, string[]> = {}; + for (const leaf of leaves) pending[leaf.id] = []; + const claimed = new Set<string>(); + for (const leaf of leaves) { + const bucket = leaf.match.bucket ?? defaultBucket; + for (const ch of channels) { + if (!matchesChannel(leaf.match, ch)) continue; + const ids = ch.buckets[bucket]; + if (!ids) continue; + for (const id of ids) { + if (claimed.has(id)) continue; + claimed.add(id); + pending[leaf.id].push(id); + } + } + } + return pending; +} + +// --- Core selection (pure) -------------------------------------------------- + +export type WorkPick = { + leafId: string; + videoId: string; + // Node ids from root..leaf. The runner increments active counts for every id + // on grant and decrements them on release, so ancestor caps are enforced. + path: string[]; +}; + +function capSaturated(node: AutoQueueNode, active: ActiveCounts): boolean { + const cap = node.maxWorkers; + if (cap == null) return false; + return (active[node.id] ?? 0) >= cap; +} + +// Whether a subtree has at least one available, non-cap-saturated video. +function hasWork( + node: AutoQueueNode, + pending: Record<string, string[]>, + active: ActiveCounts, +): boolean { + if (capSaturated(node, active)) return false; + if (!isGroup(node)) return (pending[node.id]?.length ?? 0) > 0; + return node.children.some((c) => hasWork(c, pending, active)); +} + +function weightOf(node: AutoQueueNode, mode: AutoQueueMode): number { + if (mode !== "weighted-fair") return 1; + const w = node.weight; + return typeof w === "number" && w > 0 ? w : 1; +} + +// Pick one leaf within a group via Smooth Weighted Round-Robin over only its +// children that currently have work. Mutates runtime.currentWeights. The chosen +// child is guaranteed to have work, so the recursive descent never dead-ends. +function pickSWRR( + group: AutoQueueGroup, + eligible: AutoQueueNode[], + runtime: AutoQueueRuntime, +): AutoQueueNode { + const cw = runtime.currentWeights; + let total = 0; + let best: AutoQueueNode | null = null; + for (const child of eligible) { + const w = weightOf(child, group.mode); + total += w; + cw[child.id] = (cw[child.id] ?? 0) + w; + if (!best || cw[child.id] > cw[best.id]) best = child; + } + best = best as AutoQueueNode; + cw[best.id] -= total; + return best; +} + +function pick( + node: AutoQueueNode, + pending: Record<string, string[]>, + runtime: AutoQueueRuntime, + active: ActiveCounts, + parents: string[], +): WorkPick | null { + if (capSaturated(node, active)) return null; + const path = [...parents, node.id]; + if (!isGroup(node)) { + const ids = pending[node.id]; + if (!ids || ids.length === 0) return null; + return { leafId: node.id, videoId: ids[0], path }; + } + const eligible = node.children.filter((c) => hasWork(c, pending, active)); + if (eligible.length === 0) return null; + if (node.mode === "strict") { + // Highest-priority child with work = first in declared order. + return pick(eligible[0], pending, runtime, active, path); + } + // round-robin (all weights 1) and weighted-fair share one SWRR path. + const chosen = pickSWRR(node, eligible, runtime); + return pick(chosen, pending, runtime, active, path); +} + +// Select the next video to process, or null when nothing is available (every +// leaf empty or every path cap-saturated). Mutates `runtime` (advancing fairness +// only on an actual grant). The caller increments `active` for each id in the +// returned path on start and decrements on release. +export function selectNextWork( + root: AutoQueueNode, + pending: Record<string, string[]>, + runtime: AutoQueueRuntime, + active: ActiveCounts = {}, +): WorkPick | null { + return pick(root, pending, runtime, active, []); +} + +// --- Defaults + sanitization (defensive, like sanitizeSyncScheduler) -------- + +function clampMaxWorkers(value: unknown): number | null { + if (value == null) return null; + if (typeof value !== "number" || !Number.isFinite(value)) return null; + const n = Math.floor(value); + if (n < 1) return null; + return Math.min(n, AUTO_QUEUE_MAX_WORKERS_MAX); +} + +function clampWeight(value: unknown): number { + if (typeof value !== "number" || !Number.isFinite(value)) return 1; + const n = Math.floor(value); + return n < 1 ? 1 : Math.min(n, AUTO_QUEUE_MAX_WORKERS_MAX); +} + +function sanitizeMatch(value: unknown): AutoQueueMatch { + const r = (value ?? {}) as Record<string, unknown>; + const type: AutoQueueMatchType = + r.type === "channel" || r.type === "platform" || r.type === "all" + ? r.type + : "all"; + const out: AutoQueueMatch = { type }; + if (typeof r.value === "string" && r.value.trim()) out.value = r.value.trim(); + if (typeof r.bucket === "string" && r.bucket.trim()) { + out.bucket = r.bucket.trim(); + } + return out; +} + +// Coerce a raw node, assigning a unique id (provided id preserved when valid and +// not already taken, so persisted fairness state survives an unrelated edit). +function sanitizeNode(value: unknown, seen: Set<string>): AutoQueueNode { + const r = (value ?? {}) as Record<string, unknown>; + const id = takeId(r.id, seen); + const weight = clampWeight(r.weight); + const maxWorkers = clampMaxWorkers(r.maxWorkers); + if (Array.isArray(r.children)) { + const mode: AutoQueueMode = AUTO_QUEUE_MODES.includes(r.mode as AutoQueueMode) + ? (r.mode as AutoQueueMode) + : "strict"; + return { + id, + mode, + weight, + maxWorkers, + children: r.children.map((c) => sanitizeNode(c, seen)), + }; + } + return { id, match: sanitizeMatch(r.match), weight, maxWorkers }; +} + +let idCounter = 0; +function takeId(raw: unknown, seen: Set<string>): string { + let id = typeof raw === "string" && raw.trim() ? raw.trim() : ""; + if (!id || seen.has(id)) { + do { + id = `node-${++idCounter}`; + } while (seen.has(id)); + } + seen.add(id); + return id; +} + +function sanitizeRoot(value: unknown, seen: Set<string>): AutoQueueGroup { + const node = sanitizeNode( + value && typeof value === "object" ? value : { mode: "strict", children: [] }, + seen, + ); + if (isGroup(node)) return node; + // A root that deserialized as a leaf is meaningless — wrap into an empty group. + return { id: node.id, mode: "strict", weight: 1, maxWorkers: null, children: [] }; +} + +export function defaultAutoQueuePolicy(): AutoQueuePolicy { + return { + enabled: false, + maxWorkers: null, + root: { id: "root", mode: "strict", weight: 1, maxWorkers: null, children: [] }, + }; +} + +export function defaultAutoQueue(): AutoQueueSettings { + return { + transcription: defaultAutoQueuePolicy(), + download: defaultAutoQueuePolicy(), + }; +} + +function sanitizePolicy(value: unknown): AutoQueuePolicy { + const r = (value ?? {}) as Record<string, unknown>; + const seen = new Set<string>(); + return { + enabled: r.enabled === true, + maxWorkers: clampMaxWorkers(r.maxWorkers), + root: sanitizeRoot(r.root, seen), + }; +} + +export function sanitizeAutoQueue(value: unknown): AutoQueueSettings { + if (!value || typeof value !== "object") return defaultAutoQueue(); + const r = value as Record<string, unknown>; + return { + transcription: sanitizePolicy(r.transcription), + download: sanitizePolicy(r.download), + }; +} diff --git a/common/jobs/autoQueueState.ts b/common/jobs/autoQueueState.ts @@ -0,0 +1,142 @@ +import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { + type AutoQueueRuntime, + emptyAutoQueueRuntime, +} from "./autoQueuePolicy"; + +// Persistent fairness state for the auto-queue runners. Unlike in-flight worker +// counts (which are zero by definition after a restart and live only on the +// runner singleton), the SWRR current-weights are "fairness memory": persisting +// them avoids always restarting round-robin from the first child after a reboot. +// +// Best-effort and small: one runtime per kind plus a bounded pick log for the +// observability panel. Written atomically (tmp + rename). This module is pure +// I/O + types — the selection logic lives in autoQueuePolicy.ts and the +// orchestration in autoRunner.ts. Mirrors syncSchedulerState.ts. + +export type AutoQueueKind = "transcription" | "download"; + +// One recorded grant, newest-first, for the status panel. +export type AutoQueuePick = { + at: number; + leafId: string; + videoId: string; + channelSlug: string; +}; + +export type AutoQueueKindState = { + runtime: AutoQueueRuntime; + // Newest-first, bounded to AUTO_QUEUE_PICK_LOG_LIMIT entries. + picks: AutoQueuePick[]; +}; + +export type AutoQueueState = { + transcription: AutoQueueKindState; + download: AutoQueueKindState; +}; + +export const AUTO_QUEUE_PICK_LOG_LIMIT = 50; + +export function emptyAutoQueueKindState(): AutoQueueKindState { + return { runtime: emptyAutoQueueRuntime(), picks: [] }; +} + +export function emptyAutoQueueState(): AutoQueueState { + return { + transcription: emptyAutoQueueKindState(), + download: emptyAutoQueueKindState(), + }; +} + +function coerceRuntime(value: unknown): AutoQueueRuntime { + const out = emptyAutoQueueRuntime(); + if (!value || typeof value !== "object") return out; + const r = value as Record<string, unknown>; + if (r.currentWeights && typeof r.currentWeights === "object") { + for (const [id, w] of Object.entries( + r.currentWeights as Record<string, unknown>, + )) { + if (typeof w === "number" && Number.isFinite(w)) { + out.currentWeights[id] = w; + } + } + } + return out; +} + +function coercePick(value: unknown): AutoQueuePick | null { + if (!value || typeof value !== "object") return null; + const r = value as Record<string, unknown>; + if ( + typeof r.leafId !== "string" || + typeof r.videoId !== "string" || + typeof r.channelSlug !== "string" + ) { + return null; + } + return { + at: typeof r.at === "number" && Number.isFinite(r.at) ? r.at : 0, + leafId: r.leafId, + videoId: r.videoId, + channelSlug: r.channelSlug, + }; +} + +function coerceKindState(value: unknown): AutoQueueKindState { + if (!value || typeof value !== "object") return emptyAutoQueueKindState(); + const r = value as Record<string, unknown>; + const picks = Array.isArray(r.picks) + ? r.picks + .map(coercePick) + .filter((p): p is AutoQueuePick => p !== null) + .slice(0, AUTO_QUEUE_PICK_LOG_LIMIT) + : []; + return { runtime: coerceRuntime(r.runtime), picks }; +} + +// Read the state file, tolerating a missing/corrupt file by returning empty +// state. Defensive: any stored shape that doesn't match is coerced so a +// hand-edited file can't crash the runner. +export async function readAutoQueueState(paths: Paths): Promise<AutoQueueState> { + let raw: unknown; + try { + raw = JSON.parse(await readFile(paths.autoQueueStateFile, "utf8")); + } catch { + return emptyAutoQueueState(); + } + if (!raw || typeof raw !== "object") return emptyAutoQueueState(); + const r = raw as Record<string, unknown>; + return { + transcription: coerceKindState(r.transcription), + download: coerceKindState(r.download), + }; +} + +// Write the state atomically (tmp + rename), creating the .auto-queue dir on +// first use. Each kind's pick log is trimmed to the limit on the way out. +export async function writeAutoQueueState( + paths: Paths, + state: AutoQueueState, +): Promise<void> { + const trim = (k: AutoQueueKindState): AutoQueueKindState => ({ + runtime: k.runtime, + picks: k.picks.slice(0, AUTO_QUEUE_PICK_LOG_LIMIT), + }); + const out: AutoQueueState = { + transcription: trim(state.transcription), + download: trim(state.download), + }; + await mkdir(path.dirname(paths.autoQueueStateFile), { recursive: true }); + const tmp = `${paths.autoQueueStateFile}.tmp-${process.pid}`; + await writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); + await rename(tmp, paths.autoQueueStateFile); +} + +export function recordPick(state: AutoQueueKindState, pick: AutoQueuePick): void { + state.picks.unshift(pick); + if (state.picks.length > AUTO_QUEUE_PICK_LOG_LIMIT) { + state.picks.length = AUTO_QUEUE_PICK_LOG_LIMIT; + } +} diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -336,6 +336,10 @@ const DRAINABLE_KINDS = new Set<string>([ "download-missing-subs", "sync", "retry-bucket", + // The auto-queue runners: drain stops launching new units, lets in-flight + // finish, then the loop returns and the job completes. See autoRunner.ts. + "auto-transcribe", + "auto-download", ]); export function isDrainableKind(kind: string): boolean { diff --git a/common/lib/paths.ts b/common/lib/paths.ts @@ -19,6 +19,10 @@ export type Paths = { // a rolling run log). Survives restarts, unlike the in-memory job registry. // See common/jobs/syncSchedulerState.ts. schedulerStateFile: string; + // Persistent fairness state for the automatic priority-queue runners + // (round-robin / weighted-fair cursors + a rolling pick log). Best-effort: + // drift across a restart is cosmetic. See common/jobs/autoQueueState.ts. + autoQueueStateFile: string; // Persisted "default" worker arrangement: the set of worker ids that should // start enabled on the next server launch (workers not listed start // disabled). Written by the Workers page "Set as default" button; applied by @@ -87,6 +91,7 @@ export function getPaths(): Paths { jobsDir: path.join(transcriptsDir, ".jobs"), workerScratchDir: path.join(transcriptsDir, ".worker-scratch"), schedulerStateFile: path.join(transcriptsDir, ".scheduler", "state.json"), + autoQueueStateFile: path.join(transcriptsDir, ".auto-queue", "state.json"), workerDefaultsFile: path.join(transcriptsDir, ".workers", "defaults.json"), bookmarksFile: path.join(transcriptsDir, ".bookmarks", "bookmarks.json"), lmdbPath: path.join(transcriptsDir, "index.mdb"), diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -15,8 +15,14 @@ import { sanitizeWorkers, validateWorkers, } from "./workers"; +import { + type AutoQueueSettings, + defaultAutoQueue, + sanitizeAutoQueue, +} from "../jobs/autoQueuePolicy"; export type { Worker } from "./workers"; +export type { AutoQueueSettings } from "../jobs/autoQueuePolicy"; // Transcribe placeholder/arg helpers now live with the whisper-cpp app in // transcriptionApps.ts. Re-exported here so existing import sites keep working. @@ -96,6 +102,12 @@ export type SiteSettings = { // block holds the defaults and guard rails the scheduler applies across all // channels. See common/jobs/syncScheduler.ts. syncScheduler: SyncSchedulerSettings; + // Configuration for the automatic priority-queue runners (auto-transcribe / + // auto-download). Each holds a tree policy that decides which channel's video + // to process next, cross-channel, by priority/round-robin/weighted-fair rules. + // Independent of syncScheduler (which decides staleness, not work order). See + // common/jobs/autoQueuePolicy.ts. + autoQueue: AutoQueueSettings; // 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 @@ -297,6 +309,7 @@ function defaults(): SiteSettings { reportDebouncePreset: DEFAULT_REPORT_DEBOUNCE_PRESET, autoRefreshIntervalSeconds: AUTO_REFRESH_INTERVAL_DEFAULT_SECONDS, syncScheduler: defaultSyncScheduler(), + autoQueue: defaultAutoQueue(), socialLinks: [], }; } @@ -458,6 +471,7 @@ export function getSettings(): SiteSettings { merged.autoRefreshIntervalSeconds, ); merged.syncScheduler = sanitizeSyncScheduler(merged.syncScheduler); + merged.autoQueue = sanitizeAutoQueue(merged.autoQueue); 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 @@ -627,6 +641,7 @@ export async function writeSettings(next: SiteSettings): Promise<void> { next.autoRefreshIntervalSeconds, ), syncScheduler: sanitizeSyncScheduler(next.syncScheduler), + autoQueue: sanitizeAutoQueue(next.autoQueue), socialLinks, }; const tmp = `${file}.tmp-${process.pid}`; diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **New Auto-queue: automatically transcribe (and download) across all channels by a configurable priority policy, instead of running one channel batch at a time.** Previously the only way to process pending work was to manually fire a per-channel batch (e.g. *Transcribe missing* on one channel), and since every transcription job serialized on a single queue, a batch ran to completion before any other channel got a turn — there was no way to say "do cornbreadman first, then fall back to hasanabi." The new **Auto-queue** page (under Pool → Auto-queue) adds two always-on runners, **auto-transcribe** and **auto-download**, each driven by a **policy tree**: order rules top-to-bottom for **strict** priority, or wrap rules in a group set to **round-robin** or **weighted-fair** (smooth weighted round-robin) to *alternate* between rulesets. A rule (leaf) matches a **channel**, a whole **platform**, or **all** channels, optionally narrowed to a snapshot **bucket** (e.g. prioritize `failedListed` retries over fresh `downloadedNoTranscript`), and any rule or group can carry a **max-workers** cap (a saturated subtree falls through to the next-priority sibling, like an HTB ceil). The highest-priority channel with available work claims the **next freed worker slot** — non-destructive, so a higher-priority video never kills an in-flight transcription, it just wins the next slot; when a channel's work runs out the runner falls back automatically. Transcription concurrency is bounded by the worker pool's eligible slots (so policy decides *which* video runs, the pool decides *how many*); downloads have no pool, so the runner gates to **one download per platform at a time**, matching the per-platform serial queue's politeness. Each runner is a real, drainable/cancellable job (visible on the Jobs pages), and the Auto-queue page shows live per-rule pending counts and a recent-pick log. Independent of the sync **Schedule** (which only decides *when* to re-fetch a channel) — manual batches keep working alongside it. Policies live in `settings.json` under `autoQueue` (defensively sanitized like `syncScheduler`); fairness cursors persist in `transcripts/.auto-queue/state.json`. See `common/jobs/autoQueuePolicy.ts` (pure selection engine + unit tests), `common/controller/autoRunner.ts`, `common/controller/transcribeOneFromQueue.ts` (shared per-video gating, also used by the existing whisper batch), and `editor/app/auto-queue/*`. - **Charts can now track content *added to the sites* over time, not just when creators uploaded it.** The chart engine previously only binned the time axis on a video's **upload date**. Two acquisition dates are now recorded per video — when *we* downloaded it and when *we* transcribed it — and the chart editor's X-axis gains a **Date field** selector (Uploaded / Downloaded / Transcribed) alongside the bin. A new **Content added** preset group ships three ready charts (cumulative *Library growth (added)*, cumulative *Transcribed over time*, and *Added per month* stacked by channel), and the default dashboard now includes the cumulative library-growth-by-acquisition chart so the progress view is present out of the box. Everything reuses the existing charts UI, so per-channel filtering, cumulative curves, CSV/PNG export, and shareable URLs all work unchanged. Acquisition dates come from the per-video `download-outcome.json` (`finishedAt`) and a new `transcribe-outcome.json` sidecar written when a transcript is finalized; the stats build falls back to file mtimes for content added before the sidecars existed. Requires a one-time data rebuild (`build:index` + `build:stats`) on the bumped `STATS_SCHEMA_VERSION`. See `common/lib/{stats,chartConfig,chartAggregate,chartShare,transcribeOutcome}.ts`, `common/controller/{buildStats,transcribeOne}.ts`, and `common/components/charts/ChartConfigEditor.tsx`. - **The Schedule page is now a one-stop editor for per-channel sync cadence.** The `/scheduler` page used to be read-only — you could see each channel's interval, last sync, and next-due time, but to *change* a cadence you had to open that channel's editor (Source → Auto-sync), one channel at a time, and the headline global toggles lived only in Settings. Now each row's **Interval** cell is an inline editor: pick a preset (Default / Off / Every 10–30m / Hourly / 6h / 12h / Daily / Weekly) **or** choose **Custom (minutes)…** and type an exact minute count, then **Save** — writing just `syncIntervalMinutes` to that channel's `config.json` and leaving every other field untouched (it does *not* go through the full channel-form merge). The page also gained a **Global controls** block to toggle the master **enable**, the **default interval**, and the **internal heartbeat** right there (the advanced knobs — concurrency, quiet hours, backoff — still link out to Settings). The underlying due logic is unchanged: a channel auto-syncs on the next heartbeat once `now − lastSyncedAt ≥ its interval`. The live status columns keep polling every 5s, but each row's editor holds its own state seeded once from the stored value, so a refresh can't clobber an in-progress edit. The preset list is shared with the channel editor (`editor/app/scheduler/intervalPresets.ts`), and `GET /api/scheduler/status` now carries each channel's raw `configuredIntervalMinutes` so the editor can tell *inherit-default* from *explicit-off* from *explicit-minutes*. See `editor/app/scheduler/actions.ts` and `editor/app/scheduler/components/{ChannelIntervalEditor,SchedulerSettingsForm}.tsx`. - **Scheduled sync can now run without an external cron job.** The sync scheduler previously only fired when an OS cron entry POSTed to `/api/scheduler/tick` (via `pnpm sync:tick`) — fine on a server, but a chore to set up just to call a function the editor already hosts in-process. The editor can now drive its own heartbeat through a **Next.js instrumentation hook** (`editor/instrumentation.ts`): on server startup it arms a single in-process timer that calls `runSchedulerTick()` directly — no HTTP, no cron, no token. Turn it on with **Settings → Sync scheduler → Internal heartbeat (seconds)**: `0` = off (keep using an external cron heartbeat), any positive value is clamped to `[15, 3600]`s and is the cadence the editor ticks itself at; the `SYNC_HEARTBEAT_SECONDS` env var overrides the setting at runtime. The timer is a self-rescheduling, `unref`'d `setTimeout` loop (so it never holds the process open and never overlaps a tick), re-reading the cadence each fire so a change takes effect on the next tick — though turning it on *from 0* needs a restart, since the timer is armed once at boot. It's modeled on the existing snapshot-scheduler timer and reuses the already-overlap-guarded `runSchedulerTick()`, so internal and external heartbeats are interchangeable and may even coexist. The **Schedule** page header now reports how ticks are driven ("internal heartbeat every N" vs. "external heartbeat (cron)"), and `GET /api/scheduler/status` carries the effective `heartbeatSeconds`. Defaults to off, so dev/test and existing cron installs are unchanged. One caveat for multi-instance deployments: the timer runs once *per server instance*, so a cluster against one data dir should set `SYNC_HEARTBEAT_SECONDS=0` on all but one — see `SCHEDULED_SYNC.md`. diff --git a/editor/app/api/auto-queue/control/route.ts b/editor/app/api/auto-queue/control/route.ts @@ -0,0 +1,44 @@ +import { NextResponse } from "next/server"; +import { + startAutoRunner, + stopAutoRunner, +} from "yt-dlp-transcript-common/controller/autoRunner"; +import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState"; + +export const dynamic = "force-dynamic"; + +// Start or stop an auto-queue runner without a server restart. The runner only +// auto-starts at boot (editor/instrumentation.ts) when enabled, so the save flow +// and this endpoint are how an operator (or an e2e test) brings a runner up after +// flipping the enable toggle. The editor admin surface is otherwise +// unauthenticated (trusted self-host), consistent with the rest of the app. +// +// Body: { kind: "transcription" | "download", action: "start" | "stop" }. +export async function POST(req: Request) { + let body: { kind?: unknown; action?: unknown }; + try { + body = (await req.json()) as typeof body; + } catch { + return NextResponse.json({ ok: false, error: "invalid JSON" }, { status: 400 }); + } + const kind = body.kind; + if (kind !== "transcription" && kind !== "download") { + return NextResponse.json( + { ok: false, error: "kind must be 'transcription' or 'download'" }, + { status: 400 }, + ); + } + const action = body.action; + if (action === "start") { + const jobId = await startAutoRunner(kind as AutoQueueKind); + return NextResponse.json({ ok: true, started: jobId !== null, jobId }); + } + if (action === "stop") { + const stopped = stopAutoRunner(kind as AutoQueueKind); + return NextResponse.json({ ok: true, stopped }); + } + return NextResponse.json( + { ok: false, error: "action must be 'start' or 'stop'" }, + { status: 400 }, + ); +} diff --git a/editor/app/api/auto-queue/status/route.ts b/editor/app/api/auto-queue/status/route.ts @@ -0,0 +1,11 @@ +import { NextResponse } from "next/server"; +import { buildAutoQueueStatusPayload } from "../../../auto-queue/status"; + +export const dynamic = "force-dynamic"; + +// Read-only view for the Auto-Queue panel: per-kind runner status, effective +// policy, recent picks, and snapshot-derived pending counts. Backs a passive UI +// poll, like /api/scheduler/status and /api/jobs/active. +export async function GET() { + return NextResponse.json(await buildAutoQueueStatusPayload()); +} diff --git a/editor/app/api/test/invalidate-cache/route.ts b/editor/app/api/test/invalidate-cache/route.ts @@ -30,6 +30,12 @@ function invalidate() { // use. (In production the pool legitimately persists; this is e2e-only.) // eslint-disable-next-line @typescript-eslint/no-explicit-any (globalThis as any).__yttWorkerPool__ = undefined; + // Reset the auto-queue runner singleton too: its live entries reference job + // ids in the (now-wiped) registry. A still-executing runner loop from a prior + // spec detects its job is gone (registry.get -> undefined) and exits on its + // next iteration; clearing this lets the next spec start fresh runners. + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (globalThis as any).__yttAutoRunner__ = undefined; revalidatePath("/", "layout"); return NextResponse.json({ ok: true }); } diff --git a/editor/app/auto-queue/actions.ts b/editor/app/auto-queue/actions.ts @@ -0,0 +1,70 @@ +"use server"; + +import { revalidatePath } from "next/cache"; +import { + getSettings, + writeSettings, + type SiteSettings, +} from "yt-dlp-transcript-common/lib/settings"; +import { + startAutoRunner, + stopAutoRunner, +} from "yt-dlp-transcript-common/controller/autoRunner"; +import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState"; +import type { AutoQueueGroup } from "yt-dlp-transcript-common/jobs/autoQueuePolicy"; + +export type SaveResult = { ok: true } | { ok: false; error: string }; + +// Persist one runner kind's policy and bring its runner up/down to match, +// WITHOUT a server restart. writeSettings sanitizes the tree (sanitizeAutoQueue), +// so a slightly-off client payload is coerced rather than trusted. Enabling +// starts the runner immediately; disabling is picked up by the running loop on +// its next iteration (getSettings reads from disk), so it stops on its own. +export async function saveAutoQueueAction( + kind: AutoQueueKind, + input: { enabled: boolean; maxWorkers: number | null; root: AutoQueueGroup }, +): Promise<SaveResult> { + const current = getSettings(); + const next: SiteSettings = { + ...current, + autoQueue: { + ...current.autoQueue, + [kind]: { + enabled: input.enabled, + maxWorkers: input.maxWorkers, + root: input.root, + }, + }, + }; + try { + await writeSettings(next); + } catch (e) { + return { ok: false, error: (e as Error).message }; + } + if (input.enabled) { + await startAutoRunner(kind); + } + revalidatePath("/auto-queue"); + return { ok: true }; +} + +// Explicit start/stop for the Start/Stop buttons (also reachable as +// /api/auto-queue/control for e2e). Start is a no-op when the policy is disabled. +export async function startAutoQueueAction( + kind: AutoQueueKind, +): Promise<SaveResult> { + const jobId = await startAutoRunner(kind); + if (jobId === null && !getSettings().autoQueue[kind].enabled) { + return { ok: false, error: "Enable the policy before starting the runner." }; + } + revalidatePath("/auto-queue"); + return { ok: true }; +} + +export async function stopAutoQueueAction( + kind: AutoQueueKind, +): Promise<SaveResult> { + stopAutoRunner(kind); + revalidatePath("/auto-queue"); + return { ok: true }; +} diff --git a/editor/app/auto-queue/components/AutoQueueView.tsx b/editor/app/auto-queue/components/AutoQueueView.tsx @@ -0,0 +1,246 @@ +"use client"; + +import { useCallback, useEffect, useRef, useState } from "react"; +import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState"; +import type { AutoQueueStatusPayload, AutoQueueKindStatus } from "../status"; +import { PolicyTreeEditor } from "./PolicyTreeEditor"; + +// Live panel for the auto-queue runners: one section per kind (auto-transcribe / +// auto-download), each with a status header + Start/Stop, the policy-tree editor, +// snapshot-derived pending counts per rule, and the recent pick log. Polls +// /api/auto-queue/status like the scheduler panel. + +type Channel = { slug: string; name: string | null }; + +export function AutoQueueView({ + initial, + channels, + platforms, + bucketsByKind, +}: { + initial: AutoQueueStatusPayload; + channels: Channel[]; + platforms: string[]; + bucketsByKind: Record<AutoQueueKind, string[]>; +}) { + const [data, setData] = useState<AutoQueueStatusPayload>(initial); + const mounted = useRef(true); + + const refresh = useCallback(async () => { + try { + const res = await fetch("/api/auto-queue/status", { cache: "no-store" }); + if (!res.ok) return; + const next = (await res.json()) as AutoQueueStatusPayload; + if (mounted.current) setData(next); + } catch { + /* transient; next poll retries */ + } + }, []); + + useEffect(() => { + mounted.current = true; + const id = setInterval(refresh, 3000); + return () => { + mounted.current = false; + clearInterval(id); + }; + }, [refresh]); + + return ( + <div className="flex flex-col gap-8"> + <KindPanel + kind="transcription" + title="Auto-transcribe" + status={data.transcription} + channels={channels} + platforms={platforms} + buckets={bucketsByKind.transcription} + onRefresh={refresh} + /> + <KindPanel + kind="download" + title="Auto-download" + status={data.download} + channels={channels} + platforms={platforms} + buckets={bucketsByKind.download} + onRefresh={refresh} + /> + </div> + ); +} + +function KindPanel({ + kind, + title, + status, + channels, + platforms, + buckets, + onRefresh, +}: { + kind: AutoQueueKind; + title: string; + status: AutoQueueKindStatus; + channels: Channel[]; + platforms: string[]; + buckets: string[]; + onRefresh: () => Promise<void>; +}) { + const [busy, setBusy] = useState(false); + const running = status.runner.running; + + const control = useCallback( + async (action: "start" | "stop") => { + setBusy(true); + try { + await fetch("/api/auto-queue/control", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ kind, action }), + }); + await onRefresh(); + } finally { + setBusy(false); + } + }, + [kind, onRefresh], + ); + + const leafName = (leafId: string) => leafLabel(status.policy.root, leafId); + const totalPending = Object.values(status.pendingByLeaf).reduce( + (a, b) => a + b, + 0, + ); + + return ( + <section className="flex flex-col gap-4"> + <div className="flex flex-wrap items-center gap-3"> + <h2 className="text-lg font-semibold">{title}</h2> + <span + className={`text-sm rounded-full px-3 py-1 border ${ + running + ? "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" + }`} + > + {running ? "Runner running" : "Runner stopped"} + </span> + <span className="text-sm text-zinc-500"> + {status.runner.inFlight.length} in flight · {totalPending} pending + </span> + <span className="ml-auto flex gap-2"> + <button + type="button" + aria-label={`Start ${title}`} + onClick={() => control("start")} + disabled={busy || running} + className="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" + > + Start + </button> + <button + type="button" + aria-label={`Stop ${title}`} + onClick={() => control("stop")} + disabled={busy || !running} + className="px-3 py-1.5 rounded-md border border-zinc-300 dark:border-zinc-700 text-sm font-medium hover:bg-zinc-100 dark:hover:bg-zinc-800 disabled:opacity-50" + > + Stop + </button> + </span> + </div> + + <PolicyTreeEditor + kind={kind} + initialEnabled={status.policy.enabled} + initialMaxWorkers={status.policy.maxWorkers} + initialRoot={status.policy.root} + channels={channels} + platforms={platforms} + buckets={buckets} + /> + + <div className="grid gap-5 md:grid-cols-2"> + <div className="flex flex-col gap-2"> + <h3 className="text-sm font-semibold">Pending per rule</h3> + {Object.keys(status.pendingByLeaf).length === 0 ? ( + <p className="text-sm text-zinc-500">No rules configured.</p> + ) : ( + <ul className="flex flex-col gap-1 text-sm"> + {Object.entries(status.pendingByLeaf).map(([leafId, count]) => ( + <li + key={leafId} + className="flex justify-between gap-3 text-zinc-600 dark:text-zinc-300" + > + <span>{leafName(leafId)}</span> + <span className="tabular-nums">{count}</span> + </li> + ))} + </ul> + )} + </div> + + <div className="flex flex-col gap-2"> + <h3 className="text-sm font-semibold">Recent picks</h3> + {status.picks.length === 0 ? ( + <p className="text-sm text-zinc-500">Nothing picked yet.</p> + ) : ( + <ul className="flex flex-col gap-1 text-sm"> + {status.picks.slice(0, 15).map((p, i) => ( + <li + key={`${p.at}-${i}`} + className="flex flex-wrap gap-x-3 text-zinc-600 dark:text-zinc-300" + > + <span className="tabular-nums text-zinc-500"> + {formatClock(p.at)} + </span> + <span> + {p.channelSlug}/{p.videoId} + </span> + <span className="text-zinc-400">{leafName(p.leafId)}</span> + </li> + ))} + </ul> + )} + </div> + </div> + </section> + ); +} + +// Human label for a leaf id from the current policy tree (its match), falling +// back to the raw id. +function leafLabel( + node: AutoQueueKindStatus["policy"]["root"], + leafId: string, +): string { + type N = + | typeof node + | { id: string; match?: { type: string; value?: string; bucket?: string } }; + const walk = (n: N): string | null => { + if (n.id === leafId && "match" in n && n.match) { + const m = n.match; + const base = + m.type === "all" ? "all channels" : `${m.type}: ${m.value ?? "?"}`; + return m.bucket ? `${base} [${m.bucket}]` : base; + } + const children = (n as { children?: N[] }).children; + if (children) { + for (const c of children) { + const found = walk(c); + if (found) return found; + } + } + return null; + }; + return walk(node) ?? leafId; +} + +function formatClock(ms: number): string { + if (!ms) return "—"; + return new Date(ms).toLocaleTimeString([], { + hour: "2-digit", + minute: "2-digit", + }); +} diff --git a/editor/app/auto-queue/components/PolicyTreeEditor.tsx b/editor/app/auto-queue/components/PolicyTreeEditor.tsx @@ -0,0 +1,494 @@ +"use client"; + +import { useState } from "react"; +import { + type AutoQueueGroup, + type AutoQueueLeaf, + type AutoQueueMode, + type AutoQueueNode, + AUTO_QUEUE_MODES, + isGroup, +} from "yt-dlp-transcript-common/jobs/autoQueuePolicy"; +import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState"; +import { saveAutoQueueAction, type SaveResult } from "../actions"; + +// A nested, controlled editor for one runner's policy tree. Groups carry a mode +// (strict / round-robin / weighted-fair) and optional worker cap; leaves match a +// channel / platform / everything, optionally narrowed to a snapshot bucket. The +// whole tree is held in React state and saved through a typed server action +// (writeSettings re-sanitizes it server-side). + +const inputClass = + "rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-1 text-sm"; +const btnClass = + "px-2 py-0.5 rounded border border-zinc-300 dark:border-zinc-700 text-xs hover:bg-zinc-100 dark:hover:bg-zinc-800"; + +let idSeq = 0; +function newId(): string { + idSeq += 1; + const rnd = + typeof crypto !== "undefined" && "randomUUID" in crypto + ? crypto.randomUUID().slice(0, 8) + : `${idSeq}`; + return `n-${rnd}-${idSeq}`; +} + +// --- Immutable tree ops (by node id) --------------------------------------- + +function mapNode( + node: AutoQueueNode, + id: string, + fn: (n: AutoQueueNode) => AutoQueueNode, +): AutoQueueNode { + if (node.id === id) return fn(node); + if (!isGroup(node)) return node; + return { ...node, children: node.children.map((c) => mapNode(c, id, fn)) }; +} + +function removeFrom(node: AutoQueueGroup, id: string): AutoQueueGroup { + return { + ...node, + children: node.children + .filter((c) => c.id !== id) + .map((c) => (isGroup(c) ? removeFrom(c, id) : c)), + }; +} + +function addChildTo( + node: AutoQueueNode, + parentId: string, + child: AutoQueueNode, +): AutoQueueNode { + if (!isGroup(node)) return node; + if (node.id === parentId) { + return { ...node, children: [...node.children, child] }; + } + return { + ...node, + children: node.children.map((c) => addChildTo(c, parentId, child)), + }; +} + +function moveChildIn( + node: AutoQueueNode, + parentId: string, + index: number, + dir: -1 | 1, +): AutoQueueNode { + if (!isGroup(node)) return node; + if (node.id === parentId) { + const next = [...node.children]; + const j = index + dir; + if (j < 0 || j >= next.length) return node; + [next[index], next[j]] = [next[j], next[index]]; + return { ...node, children: next }; + } + return { + ...node, + children: node.children.map((c) => moveChildIn(c, parentId, index, dir)), + }; +} + +function makeLeaf(type: "channel" | "platform" | "all"): AutoQueueLeaf { + return { id: newId(), match: { type }, weight: 1, maxWorkers: null }; +} +function makeGroup(): AutoQueueGroup { + return { + id: newId(), + mode: "strict", + weight: 1, + maxWorkers: null, + children: [], + }; +} + +// --- Component -------------------------------------------------------------- + +export function PolicyTreeEditor({ + kind, + initialEnabled, + initialMaxWorkers, + initialRoot, + channels, + platforms, + buckets, +}: { + kind: AutoQueueKind; + initialEnabled: boolean; + initialMaxWorkers: number | null; + initialRoot: AutoQueueGroup; + channels: { slug: string; name: string | null }[]; + platforms: string[]; + buckets: string[]; +}) { + const [root, setRoot] = useState<AutoQueueGroup>(initialRoot); + const [enabled, setEnabled] = useState(initialEnabled); + const [maxWorkers, setMaxWorkers] = useState<number | null>(initialMaxWorkers); + const [saving, setSaving] = useState(false); + const [result, setResult] = useState<SaveResult | null>(null); + + const update = (id: string, fn: (n: AutoQueueNode) => AutoQueueNode) => + setRoot((r) => mapNode(r, id, fn) as AutoQueueGroup); + const remove = (id: string) => setRoot((r) => removeFrom(r, id)); + const addChild = (parentId: string, child: AutoQueueNode) => + setRoot((r) => addChildTo(r, parentId, child) as AutoQueueGroup); + const move = (parentId: string, index: number, dir: -1 | 1) => + setRoot((r) => moveChildIn(r, parentId, index, dir) as AutoQueueGroup); + + const save = async () => { + setSaving(true); + setResult(null); + try { + setResult(await saveAutoQueueAction(kind, { enabled, maxWorkers, root })); + } catch (e) { + setResult({ ok: false, error: (e as Error).message }); + } finally { + setSaving(false); + } + }; + + return ( + <div className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded p-3"> + <div className="flex flex-wrap items-center gap-4"> + <label className="flex items-center gap-2 text-sm font-medium"> + <input + type="checkbox" + checked={enabled} + onChange={(e) => setEnabled(e.target.checked)} + /> + Enable auto-{kind === "transcription" ? "transcribe" : "download"} + </label> + <label className="flex items-center gap-2 text-sm"> + <span>Runner max workers</span> + <input + type="number" + min={1} + value={maxWorkers ?? ""} + placeholder="∞" + onChange={(e) => + setMaxWorkers( + e.target.value.trim() === "" + ? null + : Math.max(1, Number.parseInt(e.target.value, 10) || 1), + ) + } + className={`${inputClass} w-20`} + /> + </label> + </div> + + <NodeEditor + node={root} + depth={0} + parentId={null} + parentMode={null} + index={0} + siblingCount={1} + channels={channels} + platforms={platforms} + buckets={buckets} + update={update} + remove={remove} + addChild={addChild} + move={move} + /> + + <div className="flex items-center gap-3"> + <button + type="button" + onClick={save} + disabled={saving} + className="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" + > + {saving ? "Saving…" : "Save policy"} + </button> + {result?.ok === true && ( + <span role="status" className="text-sm text-green-700 dark:text-green-300"> + Saved. + </span> + )} + {result?.ok === false && ( + <span role="alert" className="text-sm text-red-700 dark:text-red-300"> + {result.error} + </span> + )} + </div> + </div> + ); +} + +type EditorProps = { + node: AutoQueueNode; + depth: number; + parentId: string | null; + parentMode: AutoQueueMode | null; + index: number; + siblingCount: number; + channels: { slug: string; name: string | null }[]; + platforms: string[]; + buckets: string[]; + update: (id: string, fn: (n: AutoQueueNode) => AutoQueueNode) => void; + remove: (id: string) => void; + addChild: (parentId: string, child: AutoQueueNode) => void; + move: (parentId: string, index: number, dir: -1 | 1) => void; +}; + +function NodeEditor(props: EditorProps) { + const { node, depth, parentId, parentMode, index, siblingCount } = props; + const group = isGroup(node); + return ( + <div + className="flex flex-col gap-2 rounded border border-zinc-200 dark:border-zinc-800 p-2" + style={{ marginLeft: depth > 0 ? 12 : 0 }} + > + <div className="flex flex-wrap items-center gap-2"> + <span className="text-xs font-semibold text-zinc-500"> + {group ? "Group" : "Rule"} + </span> + + {group ? ( + <ModeSelect node={node} update={props.update} /> + ) : ( + <LeafControls {...props} leaf={node} /> + )} + + {/* weight, shown when the PARENT is weighted-fair */} + {parentMode === "weighted-fair" && ( + <label className="flex items-center gap-1 text-xs text-zinc-500"> + weight + <input + type="number" + min={1} + value={node.weight ?? 1} + onChange={(e) => + props.update(node.id, (n) => ({ + ...n, + weight: Math.max(1, Number.parseInt(e.target.value, 10) || 1), + })) + } + className={`${inputClass} w-16`} + /> + </label> + )} + + <label className="flex items-center gap-1 text-xs text-zinc-500"> + max + <input + type="number" + min={1} + placeholder="∞" + value={node.maxWorkers ?? ""} + onChange={(e) => + props.update(node.id, (n) => ({ + ...n, + maxWorkers: + e.target.value.trim() === "" + ? null + : Math.max(1, Number.parseInt(e.target.value, 10) || 1), + })) + } + className={`${inputClass} w-16`} + /> + </label> + + <span className="ml-auto flex items-center gap-1"> + {parentId && siblingCount > 1 && ( + <> + <button + type="button" + className={btnClass} + disabled={index === 0} + onClick={() => props.move(parentId, index, -1)} + > + ↑ + </button> + <button + type="button" + className={btnClass} + disabled={index === siblingCount - 1} + onClick={() => props.move(parentId, index, 1)} + > + ↓ + </button> + </> + )} + {parentId && ( + <button + type="button" + className={`${btnClass} text-red-700 dark:text-red-300`} + onClick={() => props.remove(node.id)} + > + Remove + </button> + )} + </span> + </div> + + {group && ( + <> + <div className="flex flex-col gap-2"> + {(node as AutoQueueGroup).children.map((child, i) => ( + <NodeEditor + key={child.id} + {...props} + node={child} + depth={depth + 1} + parentId={node.id} + parentMode={(node as AutoQueueGroup).mode} + index={i} + siblingCount={(node as AutoQueueGroup).children.length} + /> + ))} + {(node as AutoQueueGroup).children.length === 0 && ( + <p className="text-xs text-zinc-400 pl-2"> + Empty group — add a rule or nested group below. + </p> + )} + </div> + <div className="flex flex-wrap gap-2"> + <button + type="button" + className={btnClass} + onClick={() => props.addChild(node.id, makeLeaf("channel"))} + > + + Channel rule + </button> + <button + type="button" + className={btnClass} + onClick={() => props.addChild(node.id, makeLeaf("platform"))} + > + + Platform rule + </button> + <button + type="button" + className={btnClass} + onClick={() => props.addChild(node.id, makeLeaf("all"))} + > + + Catch-all rule + </button> + <button + type="button" + className={btnClass} + onClick={() => props.addChild(node.id, makeGroup())} + > + + Group + </button> + </div> + </> + )} + </div> + ); +} + +function ModeSelect({ + node, + update, +}: { + node: AutoQueueNode; + update: (id: string, fn: (n: AutoQueueNode) => AutoQueueNode) => void; +}) { + return ( + <select + value={(node as AutoQueueGroup).mode} + onChange={(e) => + update(node.id, (n) => ({ ...n, mode: e.target.value as AutoQueueMode })) + } + className={inputClass} + > + {AUTO_QUEUE_MODES.map((m) => ( + <option key={m} value={m}> + {m} + </option> + ))} + </select> + ); +} + +function LeafControls(props: EditorProps & { leaf: AutoQueueLeaf }) { + const { leaf, channels, platforms, buckets, update } = props; + const type = leaf.match.type; + return ( + <> + <select + value={type} + onChange={(e) => + update(leaf.id, (n) => ({ + ...n, + match: { type: e.target.value as AutoQueueLeaf["match"]["type"] }, + })) + } + className={inputClass} + > + <option value="channel">channel</option> + <option value="platform">platform</option> + <option value="all">all</option> + </select> + + {type === "channel" && ( + <select + aria-label="channel rule value" + value={leaf.match.value ?? ""} + onChange={(e) => + update(leaf.id, (n) => ({ + ...n, + match: { ...(n as AutoQueueLeaf).match, value: e.target.value }, + })) + } + className={inputClass} + > + <option value="">— pick channel —</option> + {channels.map((c) => ( + <option key={c.slug} value={c.slug}> + {c.name ?? c.slug} + </option> + ))} + </select> + )} + + {type === "platform" && ( + <select + aria-label="platform rule value" + value={leaf.match.value ?? ""} + onChange={(e) => + update(leaf.id, (n) => ({ + ...n, + match: { ...(n as AutoQueueLeaf).match, value: e.target.value }, + })) + } + className={inputClass} + > + <option value="">— pick platform —</option> + {platforms.map((p) => ( + <option key={p} value={p}> + {p} + </option> + ))} + </select> + )} + + <label className="flex items-center gap-1 text-xs text-zinc-500"> + bucket + <select + value={leaf.match.bucket ?? ""} + onChange={(e) => + update(leaf.id, (n) => { + const match = { ...(n as AutoQueueLeaf).match }; + if (e.target.value) match.bucket = e.target.value; + else delete match.bucket; + return { ...n, match }; + }) + } + className={inputClass} + > + <option value="">default</option> + {buckets.map((b) => ( + <option key={b} value={b}> + {b} + </option> + ))} + </select> + </label> + </> + ); +} diff --git a/editor/app/auto-queue/page.tsx b/editor/app/auto-queue/page.tsx @@ -0,0 +1,60 @@ +import type { Metadata } from "next"; +import Link from "next/link"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { listChannels } from "yt-dlp-transcript-common/controller/channels"; +import { PLATFORM_VALUES } from "yt-dlp-transcript-common/lib/platform"; +import { buildAutoQueueStatusPayload } from "./status"; +import { AutoQueueView } from "./components/AutoQueueView"; + +export const dynamic = "force-dynamic"; + +export const metadata: Metadata = { title: "Auto-queue" }; + +// Buckets each runner can draw from. Transcription: fresh downloaded-no- +// transcript plus retry-the-failed. Download: the channel's undownloaded ids. +const BUCKETS_BY_KIND = { + transcription: ["downloadedNoTranscript", "failedListed"], + download: ["undownloadedIds"], +}; + +export default async function AutoQueuePage() { + const [initial, channels] = await Promise.all([ + buildAutoQueueStatusPayload(), + listChannels(getPaths()), + ]); + const channelOptions = channels.map((c) => ({ + slug: c.slug, + name: c.config.name ?? null, + })); + + return ( + <div className="flex flex-col gap-4"> + <div className="flex items-center justify-between"> + <h1 className="text-2xl font-semibold">Auto-queue</h1> + <Link href="/scheduler" className="text-sm underline"> + Sync schedule + </Link> + </div> + <p className="text-sm text-zinc-500"> + Automatically pick the next transcription across all channels by a + priority policy, instead of running one channel batch at a time. Order the + rules top-to-bottom for strict priority, or wrap rules in a group set to{" "} + <em>round-robin</em> / <em>weighted-fair</em> to alternate between them. + The highest-priority channel with available work claims the next freed + worker slot; when its work runs out the runner falls back to the next + rule automatically. This is independent of the{" "} + <Link href="/scheduler" className="underline"> + sync schedule + </Link>{" "} + (which only decides when to re-fetch each channel) and manual batches keep + working alongside it. + </p> + <AutoQueueView + initial={initial} + channels={channelOptions} + platforms={[...PLATFORM_VALUES]} + bucketsByKind={BUCKETS_BY_KIND} + /> + </div> + ); +} diff --git a/editor/app/auto-queue/status.ts b/editor/app/auto-queue/status.ts @@ -0,0 +1,49 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { getSettings } from "yt-dlp-transcript-common/lib/settings"; +import { + type AutoRunnerStatus, + computeLeafPendingCounts, + getAutoRunnerStatus, +} from "yt-dlp-transcript-common/controller/autoRunner"; +import { + type AutoQueueKind, + type AutoQueuePick, + readAutoQueueState, +} from "yt-dlp-transcript-common/jobs/autoQueueState"; +import type { AutoQueuePolicy } from "yt-dlp-transcript-common/jobs/autoQueuePolicy"; + +// Read-only payload for the Auto-Queue panel: per-kind runner status (running?, +// what's in flight, in-flight counts per tree node), the effective policy, the +// recent pick log, and snapshot-derived pending counts per leaf — so the UI can +// show "cornbreadman: 12 pending" and which leaf is currently being serviced. +// Backs a passive UI poll, like /api/scheduler/status. + +export type AutoQueueKindStatus = { + kind: AutoQueueKind; + policy: AutoQueuePolicy; + runner: AutoRunnerStatus; + pendingByLeaf: Record<string, number>; + picks: AutoQueuePick[]; +}; + +export type AutoQueueStatusPayload = { + transcription: AutoQueueKindStatus; + download: AutoQueueKindStatus; +}; + +async function buildKind(kind: AutoQueueKind): Promise<AutoQueueKindStatus> { + const paths = getPaths(); + const policy = getSettings().autoQueue[kind]; + const runner = getAutoRunnerStatus(kind); + const state = await readAutoQueueState(paths); + const pendingByLeaf = await computeLeafPendingCounts(kind, paths); + return { kind, policy, runner, pendingByLeaf, picks: state[kind].picks }; +} + +export async function buildAutoQueueStatusPayload(): Promise<AutoQueueStatusPayload> { + const [transcription, download] = await Promise.all([ + buildKind("transcription"), + buildKind("download"), + ]); + return { transcription, download }; +} diff --git a/editor/app/layout.tsx b/editor/app/layout.tsx @@ -56,6 +56,7 @@ const NAV_GROUPS: NavGroup[] = [ { href: "/workers", label: "Workers" }, { href: "/widget/builder", label: "Monitor" }, { href: "/scheduler", label: "Schedule" }, + { href: "/auto-queue", label: "Auto-queue" }, { href: "/build", label: "Build" }, { href: "/actionable", label: "Actionable" }, ], diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -6,6 +6,7 @@ import { AUTO_REFRESH_INTERVAL_MIN_SECONDS, DEFAULT_REPORT_DEBOUNCE_PRESET, isReportDebouncePreset, + getSettings, MIN_FREE_DISK_GB_MAX, normalizeSocialSvg, parseSocialLinks, @@ -189,6 +190,10 @@ export async function saveSettingsAction( reportDebouncePreset, autoRefreshIntervalSeconds: autoRefreshParsed, syncScheduler, + // Preserve the existing auto-queue policy on an unrelated settings save + // (this form doesn't edit it; the Auto-queue page does). writeSettings + // re-sanitizes it regardless. + autoQueue: getSettings().autoQueue, socialLinks, }; try { diff --git a/editor/e2e/auto-queue.spec.ts b/editor/e2e/auto-queue.spec.ts @@ -0,0 +1,409 @@ +import { mkdir, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { + resetData, + resolvePath, + pathExists, + readJson, + writeSettings, +} from "./helpers"; +import { baseUrl } from "./baseUrl"; + +// End-to-end coverage for the automatic priority-queue runner +// (common/controller/autoRunner.ts). Uses the fake-whisper fixture (wired via +// WHISPER_BIN in dev:test) so transcriptions complete near-instantly, a single +// worker + maxWorkers:1 so picks are fully serialized and deterministic, and the +// runner's pick log (/api/auto-queue/status) to assert the cross-channel order +// the policy produced. The pure selection rules are unit-tested in +// common/jobs/autoQueuePolicy.test.ts; this proves the runner actually drives +// transcriptions in that order across channels. + +type Leaf = { + id: string; + match: { type: string; value?: string; bucket?: string }; + weight?: number; + maxWorkers?: number | null; +}; +type Group = { + id: string; + mode: string; + children: (Group | Leaf)[]; + maxWorkers?: number | null; +}; + +// One enabled local worker using the fake whisper engine. +const ONE_WORKER = [ + { + id: "w1", + name: "W1", + kind: "local", + enabled: true, + priority: 0, + appId: "whisper-cpp", + config: {}, + }, +]; + +async function makeChannel(slug: string, ids: string[]) { + const root = resolvePath(`test-transcripts/channels/${slug}`); + await mkdir(`${root}/data`, { recursive: true }); + await writeFile( + `${root}/config.json`, + JSON.stringify({ + handling: "transcribe", + name: slug, + url: `https://example.com/@${slug}`, + audioFormat: "mp3", + }), + ); + for (const id of ids) { + await mkdir(`${root}/data/${id}`, { recursive: true }); + await writeFile(`${root}/data/${id}/audio.mp3`, `fake audio ${id}\n`); + } + // Seed a static snapshot so the runner sees these as downloaded-no-transcript. + await writeFile( + `${root}/snapshot.json`, + JSON.stringify({ + generatedAt: "2026-06-01T00:00:00.000Z", + totals: { videos: ids.length, transcribed: 0, downloaded: ids.length }, + buckets: { + noTranscript: [], + downloadedNoTranscript: ids, + untranscoded: [], + multipleAudioFormats: [], + transcribedWithAudio: [], + untranscribable: [], + noMetadata: [], + failedListed: [], + missingFromArchive: [], + duplicateDirs: [], + partialDownloads: [], + corruptSource: [], + nonStandardVtt: [], + skippedByFilter: [], + }, + undownloadedIds: [], + }), + ); +} + +function transcriptionAutoQueue(root: Group, maxWorkers: number | null = 1) { + return { transcription: { enabled: true, maxWorkers, root }, download: {} }; +} + +function downloadAutoQueue(root: Group, maxWorkers: number | null = null) { + return { transcription: {}, download: { enabled: true, maxWorkers, root } }; +} + +// A YouTube channel with undownloaded videos: a playlist of watch URLs whose ids +// match, and a snapshot listing them as undownloadedIds. Same platform for both +// channels so the runner's per-platform gate serializes downloads → a +// deterministic cross-channel order. +async function makeDownloadChannel(slug: string, ids: string[]) { + const root = resolvePath(`test-transcripts/channels/${slug}`); + await mkdir(`${root}/data`, { recursive: true }); + await writeFile( + `${root}/config.json`, + JSON.stringify({ + handling: "youtube", + name: slug, + url: `https://www.youtube.com/@${slug}/videos`, + }), + ); + await writeFile( + `${root}/playlist`, + ids.map((id) => `https://www.youtube.com/watch?v=${id}`).join("\n") + "\n", + ); + await writeFile( + `${root}/snapshot.json`, + JSON.stringify({ + generatedAt: "2026-06-01T00:00:00.000Z", + totals: { videos: ids.length, transcribed: 0, downloaded: 0 }, + buckets: { + noTranscript: [], + downloadedNoTranscript: [], + untranscoded: [], + multipleAudioFormats: [], + transcribedWithAudio: [], + untranscribable: [], + noMetadata: [], + failedListed: [], + missingFromArchive: [], + duplicateDirs: [], + partialDownloads: [], + corruptSource: [], + nonStandardVtt: [], + skippedByFilter: [], + }, + undownloadedIds: ids, + }), + ); +} + +async function startRunner( + request: import("@playwright/test").APIRequestContext, + kind: "transcription" | "download" = "transcription", +) { + const res = await request.post(`${baseUrl}/api/auto-queue/control`, { + data: { kind, action: "start" }, + }); + expect(res.ok()).toBeTruthy(); + return res.json(); +} + +async function stopRunner( + request: import("@playwright/test").APIRequestContext, + kind: "transcription" | "download" = "transcription", +) { + await request.post(`${baseUrl}/api/auto-queue/control`, { + data: { kind, action: "stop" }, + }); +} + +type KindStatus = { + runner: { running: boolean; jobId: string | null }; + pendingByLeaf: Record<string, number>; + picks: { leafId: string; videoId: string; channelSlug: string }[]; +}; +type Status = { transcription: KindStatus; download: KindStatus }; + +async function getStatus( + request: import("@playwright/test").APIRequestContext, +): Promise<Status> { + const res = await request.get(`${baseUrl}/api/auto-queue/status`); + expect(res.ok()).toBeTruthy(); + return res.json(); +} + +// Chronological pick order (the status log is newest-first). +function pickOrder( + status: Status, + kind: "transcription" | "download" = "transcription", +): string[] { + return status[kind].picks.map((p) => p.videoId).reverse(); +} + +async function allTranscribed(slug: string, ids: string[]): Promise<boolean> { + for (const id of ids) { + if ( + !(await pathExists(`test-transcripts/channels/${slug}/data/${id}/transcript.json`)) + ) { + return false; + } + } + return true; +} + +test.afterEach(async ({ request }) => { + // Stop any runner left from a test so its loop can't bleed into the next spec. + await stopRunner(request, "transcription"); + await stopRunner(request, "download"); +}); + +test("strict priority: drains the high-priority channel first, then falls back", async ({ + request, +}) => { + await resetData(null); + await makeChannel("alpha", ["a1", "a2"]); + await makeChannel("beta", ["b1", "b2"]); + + const root: Group = { + id: "root", + mode: "strict", + children: [ + { id: "leaf-alpha", match: { type: "channel", value: "alpha" } }, + { id: "leaf-beta", match: { type: "channel", value: "beta" } }, + ], + }; + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: ONE_WORKER, + autoQueue: transcriptionAutoQueue(root), + }); + + await startRunner(request); + + await expect + .poll(async () => (await getStatus(request)).transcription.picks.length, { + timeout: 60_000, + }) + .toBe(4); + + // Every alpha video is picked before any beta video, and in-channel order holds. + expect(pickOrder(await getStatus(request))).toEqual(["a1", "a2", "b1", "b2"]); + expect(await allTranscribed("alpha", ["a1", "a2"])).toBe(true); + expect(await allTranscribed("beta", ["b1", "b2"])).toBe(true); +}); + +test("round-robin: alternates between the two channels", async ({ request }) => { + await resetData(null); + await makeChannel("alpha", ["a1", "a2"]); + await makeChannel("beta", ["b1", "b2"]); + + const root: Group = { + id: "root", + mode: "round-robin", + children: [ + { id: "leaf-alpha", match: { type: "channel", value: "alpha" } }, + { id: "leaf-beta", match: { type: "channel", value: "beta" } }, + ], + }; + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: ONE_WORKER, + autoQueue: transcriptionAutoQueue(root), + }); + + await startRunner(request); + + await expect + .poll(async () => (await getStatus(request)).transcription.picks.length, { + timeout: 60_000, + }) + .toBe(4); + + // SWRR with equal weights interleaves the channels: a, b, a, b. + const order = pickOrder(await getStatus(request)); + expect(order).toEqual(["a1", "b1", "a2", "b2"]); +}); + +test("UI: build a policy in the editor, save it, and start the runner", async ({ + page, +}) => { + await resetData(null); + await makeChannel("alpha", ["a1", "a2"]); + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: ONE_WORKER, + autoQueue: { + transcription: { + enabled: false, + maxWorkers: 1, + root: { id: "root", mode: "strict", children: [] }, + }, + download: {}, + }, + }); + + await page.goto("/auto-queue"); + await expect(page.getByRole("heading", { name: "Auto-queue" })).toBeVisible(); + // Scope to the Auto-transcribe section (the page also has an Auto-download one). + const section = page.locator("section", { + has: page.getByRole("heading", { name: "Auto-transcribe" }), + }); + await expect(section.getByText("Runner stopped")).toBeVisible(); + + // Add a channel rule and point it at alpha. + await section.getByRole("button", { name: "+ Channel rule" }).click(); + await section.getByLabel("channel rule value").selectOption("alpha"); + // Enable and save. + await section.getByLabel(/Enable auto-transcribe/).check(); + await section.getByRole("button", { name: "Save policy" }).click(); + await expect(section.getByText("Saved.")).toBeVisible(); + + // The policy persisted to settings. + const settings = await readJson<{ + autoQueue: { transcription: { enabled: boolean; root: { children: { match?: { value?: string } }[] } } }; + }>("test-settings.json"); + expect(settings.autoQueue.transcription.enabled).toBe(true); + expect( + settings.autoQueue.transcription.root.children.some( + (c) => c.match?.value === "alpha", + ), + ).toBe(true); + + // Start the runner from the UI; it should transcribe alpha's videos. + await section.getByRole("button", { name: "Start Auto-transcribe" }).click(); + await expect + .poll(async () => allTranscribed("alpha", ["a1", "a2"]), { timeout: 60_000 }) + .toBe(true); +}); + +test("status: reports the runner running and snapshot-derived pending counts", async ({ + request, +}) => { + await resetData(null); + await makeChannel("alpha", ["a1", "a2", "a3"]); + + const root: Group = { + id: "root", + mode: "strict", + children: [{ id: "leaf-alpha", match: { type: "channel", value: "alpha" } }], + }; + // No workers enabled -> the runner has zero eligible slots, so it parks and + // launches nothing, leaving the pending count observable and stable. + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: [], + autoQueue: transcriptionAutoQueue(root, null), + }); + + await startRunner(request); + + await expect + .poll(async () => (await getStatus(request)).transcription.runner.running, { + timeout: 15_000, + }) + .toBe(true); + + const status = await getStatus(request); + expect(status.transcription.pendingByLeaf["leaf-alpha"]).toBe(3); + // With no eligible workers nothing should have been picked. + expect(status.transcription.picks.length).toBe(0); +}); + +test("download: prioritizes channels across the per-platform queue", async ({ + request, +}) => { + await resetData(null); + // Same platform (youtube) for both, so the per-platform gate serializes the + // downloads and the cross-channel pick order is deterministic. + await makeDownloadChannel("alpha", ["a1", "a2"]); + await makeDownloadChannel("beta", ["b1", "b2"]); + + const root: Group = { + id: "root", + mode: "strict", + children: [ + { id: "leaf-alpha", match: { type: "channel", value: "alpha" } }, + { id: "leaf-beta", match: { type: "channel", value: "beta" } }, + ], + }; + await writeSettings({ + adminTitle: "Test Admin", + maxTranscriptPageBytes: 8388608, + sleepBetweenDownloadsSeconds: 0, + minFreeDiskGB: 0, + workers: ONE_WORKER, + autoQueue: downloadAutoQueue(root), + }); + + await startRunner(request, "download"); + + await expect + .poll(async () => (await getStatus(request)).download.picks.length, { + timeout: 60_000, + }) + .toBe(4); + + // alpha (higher priority) downloads before beta; same-platform serialization + // keeps in-channel order too. + expect(pickOrder(await getStatus(request), "download")).toEqual([ + "a1", + "a2", + "b1", + "b2", + ]); +}); diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -14,4 +14,17 @@ export async function register() { // blocks the server from becoming ready. const { startSyncHeartbeat } = await import("./app/scheduler/heartbeat"); startSyncHeartbeat(); + + // Start the automatic priority-queue runners (auto-transcribe / auto-download) + // if their policies are enabled. Each is a self-managed registry job; this only + // kicks them off and returns. Best-effort — a failure here must not stop the + // server from starting. + try { + const { startAutoRunnersIfEnabled } = await import( + "yt-dlp-transcript-common/controller/autoRunner" + ); + await startAutoRunnersIfEnabled(); + } catch { + /* a runner that fails to start must not block server readiness */ + } } diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml @@ -3,6 +3,13 @@ packages: - editor - export +allowBuilds: + esbuild: true + lmdb: true + msgpackr-extract: true + sharp: true + unrs-resolver: true + ignoredBuiltDependencies: - sharp - unrs-resolver