import type { ChildProcess } from "node:child_process"; import type { JobSpec } from "./jobSpec"; import { getScheduler } from "./scheduler"; import { getJobKind, type SchedulerTier } from "./jobKinds"; import { ulid } from "./ulid"; export type JobStatus = | "queued" | "running" | "done" | "failed" | "cancelled"; // Extended for the digest sweep. NOTE: this union has one re-spelled copy in // editor jobRowView.ts's JobRowProgress — kept as an IMPORT there, so // TypeScript actually flags the next member added here. export type JobProgressMetric = | "downloads" | "transcripts" | "digests" | "backfills" // The metadata scan. Its progress CANNOT be re-counted from disk — the scan // deliberately writes no video directory — so its runner always sets // `current` itself. See ytdlp/metadataScan.ts. | "scans" // Clip windows of a fetch-windows batch (controller/fetchWindows.ts). Spans // channels, so there is no channel count to re-count; the runner always sets // `current`. | "clips"; export type JobProgress = { metric: JobProgressMetric; initial: number; target: number; // Progress as counted BY THE RUNNER, when the runner knows better than the // disk does. // // The default is to re-count `current` from on-disk channel stats, which is // right for downloads and transcripts: a file appears, the count goes up. // It is WRONG for digests under regeneration — a regenerated digest is // rewritten in place, so the file count never moves and the bar sits at 0% // for the whole job (observed: {initial:39, current:39, target:72, pct:0}). // Cosmetic for a five-minute job; over an 81-day sweep it makes working work // look wedged. // // Only set it where the runner has a genuinely better number. Absent → the // disk re-count, unchanged. current?: number; // Audio-seconds still to process. Paired with completedTaskAudioSeconds it // gives an ETA in the unit the work is actually priced in; without it the // estimate falls back to averaging TASKS, which for digests is wrong by more // than an order of magnitude between a VOD channel and a shorts channel. remainingAudioSeconds?: number; }; // "relocate" is one CHANNEL'S media move, not one video: the relocate job has // exactly one sub-operation and its progress is bytes copied by rsync. It is a // task rather than a JobProgress metric because JobProgress counts artifacts // re-countable from disk, and a copy in flight is neither. export type JobTaskKind = | "download" | "transcribe" | "digest" | "backfill" | "relocate"; // A single in-flight sub-operation within a job (one video download or one // transcription). Only currently-running tasks are kept on the record — they // are added on start and pruned on completion — so memory stays bounded by the // job's concurrency regardless of how many videos a batch spans. export type JobTask = { id: string; // videoId (or URL-derived id) label: string; // human label kind: JobTaskKind; fraction?: number; // 0..1, undefined until first parseable progress detail?: string; // e.g. "2.31MiB/s ETA 00:03:12" or "01:50:02 / 03:24:47" startedAt: number; // For transcribe tasks: the worker running it (id into the worker pool), so // the Workers page can show which engine each task is on. workerId?: string; // Audio-integrity phase of a download task. "probing" while yt-dlp is paused // for an ffmpeg integrity probe (the determinate "scanning" bar), otherwise // "downloading" (or undefined for tasks that never probe). phase?: "downloading" | "probing"; // Wall-clock start of the current probe, stamped when phase flips to // "probing"; the UI fills the scanning bar against it. Cleared on resume. probeStartedAt?: number; // Estimated probe duration (seconds) for the scanning bar's denominator. probeEtaSeconds?: number; }; export type JobRecord = { id: string; kind: string; queueKey: string; channelSlug?: string; videoId?: string; // A serializable replay descriptor, set for replayable job kinds. Its // presence is what offers Retry on a job in the UI. See jobSpec.ts. spec?: JobSpec; // Background jobs queue BEHIND foreground (default) jobs on the same queueKey // (see enqueue). The auto-download runner marks its per-video units background // so a manually-clicked Sync on that platform preempts the queued units. background?: boolean; status: JobStatus; queuedAt: number; startedAt?: number; endedAt?: number; exitCode?: number; logPath: string; child?: ChildProcess; abortController?: AbortController; // Separate from abortController: a drain (soft-cancel) aborts this signal so // a batch stops starting NEW sub-operations but lets in-flight ones finish, // then completes normally (releasing the queue). Hard cancel aborts both. drainController?: AbortController; draining?: boolean; progress?: JobProgress; tasks?: JobTask[]; // Running totals over completed sub-operations, used to estimate time // remaining (remaining tasks × average measured task duration). Accumulated // as tasks finish — unlike `tasks`, these are NOT cleared, so the average // reflects the whole batch. Both undefined until the first task completes. completedTaskCount?: number; completedTaskMs?: number; // Audio-seconds covered by completed sub-operations. Only digest/transcribe // tasks report it, because only they consume work proportional to a video's // LENGTH — a digest of a 3-hour VOD is not one task's worth of anything. // // It is what makes a sweep ETA expressible: `completedTaskMs` over this gives // seconds-per-audio-hour, the unit the whole backfill is estimated in, and // the only unit in which "how long is this going to take" has an answer. A // task-count average cannot express it — the corpus is 77k videos and 77k // audio-hours, and those distribute completely differently. completedTaskAudioSeconds?: number; // When the job last MOVED (epoch ms): its progress snapshot changed // (`noteProgress`) or a sub-operation finished (`recordTaskDuration`). // Undefined until the first. What the /jobs stall heuristic measures quiet // time from — a slow scan at 10 videos/min that holds a slot for an hour with // no in-flight tasks is still moving (release 9, B4a). progressAt?: number; }; // Store a progress snapshot on the record, stamping `progressAt` only when the // numbers changed — a runner that re-reports the same count is not moving. export function noteProgress( record: JobRecord, snap: JobProgress, now: number = Date.now(), ): void { const prev = record.progress; if ( !prev || prev.current !== snap.current || prev.target !== snap.target || prev.initial !== snap.initial || prev.remainingAudioSeconds !== snap.remainingAudioSeconds ) { record.progressAt = now; } record.progress = snap; } export type QueueSnapshot = { name: string; running?: JobRecord; queued: JobRecord[]; }; type StartFn = () => void; type CancelFn = () => void; // The terminal transition itself: status, endedAt, exitCode — once. A record // already terminal keeps the state it ended in (a cancel racing a child's exit // stays "cancelled"). function markTerminal( job: JobRecord, status: "done" | "failed" | "cancelled", exitCode?: number, ): void { if (job.status !== "queued" && job.status !== "running") return; job.status = status; job.endedAt = Date.now(); if (typeof exitCode === "number") job.exitCode = exitCode; } // The registry owns job LIFECYCLE/STATE (the JobRecord, meta sidecars, tasks, // terminal transitions). Queue ORDERING — which job runs vs. waits, and the // foreground-before-background priority — is delegated to the shared Scheduler // (scheduler.ts), the single source of truth that the worker pool also shares. // Every non-"" job is submitted with concurrency 1 (strict per-queueKey serial), // preserving the long-standing one-running-job-per-queue behavior. class JobRegistry { private jobs = new Map(); // Set by the graceful-shutdown reaper (shutdownCancel.ts) before it cancels // every live job; see cancel(). private shuttingDown = false; // The server is going down: the cancels that follow are nobody's decision. beginShutdown(): void { this.shuttingDown = true; } register(record: JobRecord): void { this.jobs.set(record.id, record); // Bound the registry. Drop oldest finished records past 100, but never // evict records that are queued or running. if (this.jobs.size > 100) { const finished = Array.from(this.jobs.values()) .filter((j) => j.status !== "running" && j.status !== "queued") .sort((a, b) => (a.endedAt ?? 0) - (b.endedAt ?? 0)); for (const drop of finished.slice(0, this.jobs.size - 100)) { this.jobs.delete(drop.id); } } } get(id: string): JobRecord | undefined { return this.jobs.get(id); } list(): JobRecord[] { return Array.from(this.jobs.values()).sort( (a, b) => b.queuedAt - a.queuedAt, ); } // Submit a job for scheduling. The scheduler decides whether it runs now (its // queueKey has a free slot) or waits, and invokes `start` when it becomes head // of its queue — immediately, or later when an earlier job finalizes. // // Priority: a foreground (default) job runs ahead of any queued background // jobs but never displaces a running one — so a manually clicked Sync jumps // ahead of the auto-download runner's queued units without interrupting one // that's already mid-download. This is the scheduler's compareTier ordering // (background -> "background" tier, everything else -> "foreground"). // // queueKey === "" is special: the job bypasses serialization entirely and runs // immediately, parallel to any other work, untracked by the scheduler. // Cancellation still works (running -> SIGTERM/abort). enqueue( record: JobRecord, callbacks: { start: StartFn; onCancel: CancelFn }, ): { willRunNow: boolean; position: number } { // Setting status/startedAt is coupled with actually launching, so wrap the // caller's start: the scheduler runs this exactly when the job starts. const start: StartFn = () => { record.status = "running"; record.startedAt = Date.now(); callbacks.start(); }; const tier: SchedulerTier = record.background ? "background" : getJobKind(record.kind)?.defaultTier ?? "foreground"; return getScheduler().submit({ id: record.id, queueKey: record.queueKey, tier, concurrency: 1, start, onCancel: callbacks.onCancel, }); } // Idempotent: marks the job terminal (if not already), then tells the // scheduler it is done so the next queued job on that queueKey is promoted. finalize( id: string, status: "done" | "failed" | "cancelled", exitCode?: number, ): void { const job = this.jobs.get(id); if (!job) return; markTerminal(job, status, exitCode); // A RUNNING job's cancel sets `cancelled` at once and leaves the record // without `endedAt` while its function winds down (a transcriber finishing // its window); markTerminal then skips it, because the status is already // terminal. This is where the job has actually stopped, so it is stamped // here — and "cancelled with no endedAt" means "still stopping" to anyone // asking whether the job is still writing (controller/channelWriters.ts). if (job.endedAt === undefined && job.status === "cancelled") { job.endedAt = Date.now(); } // Release per-task and drain references on terminal jobs. job.tasks = []; job.drainController = undefined; getScheduler().complete(id); } cancel(id: string): boolean { const job = this.jobs.get(id); if (!job) return false; if (job.status === "queued") { // Terminal FIRST, then out of the scheduler — which fires the queued // job's onCancel, and a managed job persists its sidecar there // (streamCommand.ts), so the record it writes must already say // cancelled. It used to be marked after, and every job cancelled while // queued kept "queued" in its .meta.json (release 13 slice W1). Not // finalize() first: its scheduler.complete would take the entry out of // the queue WITHOUT firing onCancel, and the job's stream and `done` // would never settle. finalize() after is idempotent — it releases the // task references, and scheduler.complete is a no-op by then. // // EXCEPT AT SHUTDOWN. The reaper cancels a queued job only so the exit // cannot promote it into a child; nobody cancelled it. Its sidecar must // stay `queued`, which is what the boot pass (bootQueuedJobs.ts) settles // or re-queues on the next start — so it reaches onCancel still queued, // and onCancel persists only a terminal record. if (!this.shuttingDown) markTerminal(job, "cancelled"); getScheduler().cancel(id); this.finalize(id, "cancelled"); return true; } if (job.status === "running") { // Set status synchronously so the streamCommand exit handler sees // "cancelled" and skips its own status update. SIGTERM the child; // streamCommand's .finally calls finalize(cancelled) which is // idempotent and will start the next queued job. job.status = "cancelled"; job.abortController?.abort(); // Also abort the drain signal so any code awaiting a drain unblocks. job.drainController?.abort(); if (job.child) { job.child.kill("SIGTERM"); setTimeout(() => { if (job.child && !job.child.killed) job.child.kill("SIGKILL"); }, 5_000).unref?.(); } return true; } return false; } // Escape hatch for a WEDGED slot: unconditionally free the scheduler slot for // this id, even when the registry record is already terminal (or evicted). This // deliberately BYPASSES finalize's status guard — cancel() marks a job // "cancelled" synchronously but the slot is only freed later when the child // promise settles → finalize() → scheduler.complete(). If that never settles // (a child that ignored SIGTERM, a crashed finalizer), the registry says // "cancelled" while the scheduler still marks the slot "running", blocking every // job behind it on that queueKey. This kills the child hard, marks any still- // active record terminal, and calls scheduler.complete unconditionally to free // the slot and promote the next queued job. forceRelease(id: string): boolean { const job = this.jobs.get(id); if (job) { job.abortController?.abort(); job.drainController?.abort(); if (job.child && !job.child.killed) job.child.kill("SIGKILL"); if (job.status === "queued" || job.status === "running") { job.status = "cancelled"; job.endedAt = Date.now(); } else if (job.status === "cancelled" && job.endedAt === undefined) { // Cancelled and never finalized (the wedged case this exists for): the // operator has declared it stopped, so it stops counting as a writer. job.endedAt = Date.now(); } job.tasks = []; job.drainController = undefined; } // Free the slot even if the record is terminal or evicted. scheduler.complete // is idempotent, a no-op for unknown ids, and promotes the next queued job. getScheduler().complete(id); return true; } // Soft-cancel ("drain"): stop the job from starting NEW sub-operations but // let in-flight ones run to completion. The batch's fn returns normally, // which finalizes it "done" and promotes the next queued job. A queued job // has nothing in flight, so draining it is equivalent to a plain cancel. requestDrain(id: string): boolean { const job = this.jobs.get(id); if (!job) return false; if (job.status === "queued") return this.cancel(id); if (job.status === "running") { job.draining = true; job.drainController?.abort(); return true; } return false; } // Per-task tracking. No-ops if the job is gone. tasks is created lazily and // only ever holds currently-running sub-operations. addTask(jobId: string, task: JobTask): void { const job = this.jobs.get(jobId); if (!job) return; if (!job.tasks) job.tasks = []; job.tasks.push(task); } updateTask( jobId: string, taskId: string, patch: { fraction?: number; detail?: string; phase?: "downloading" | "probing"; probeEtaSeconds?: number; }, ): void { const job = this.jobs.get(jobId); const task = job?.tasks?.find((t) => t.id === taskId); if (!task) return; if (patch.fraction !== undefined) task.fraction = patch.fraction; if (patch.detail !== undefined) task.detail = patch.detail; if (patch.phase !== undefined && patch.phase !== task.phase) { task.phase = patch.phase; if (patch.phase === "probing") { task.probeStartedAt = Date.now(); task.probeEtaSeconds = patch.probeEtaSeconds; } else { task.probeStartedAt = undefined; task.probeEtaSeconds = undefined; } } } removeTask(jobId: string, taskId: string): void { const job = this.jobs.get(jobId); if (!job?.tasks) return; const idx = job.tasks.findIndex((t) => t.id === taskId); if (idx >= 0) job.tasks.splice(idx, 1); } // Fold one finished sub-operation's wall-clock duration into the job's // running totals (see JobRecord.completedTaskCount/Ms). No-op if the job is // gone or the duration is nonsensical. recordTaskDuration( jobId: string, durationMs: number, audioSeconds?: number, ): void { const job = this.jobs.get(jobId); if (!job || durationMs < 0) return; job.completedTaskCount = (job.completedTaskCount ?? 0) + 1; job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs; job.progressAt = Date.now(); if (typeof audioSeconds === "number" && audioSeconds > 0) { job.completedTaskAudioSeconds = (job.completedTaskAudioSeconds ?? 0) + audioSeconds; } } // Snapshot of queues for UI, resolved from the scheduler's id views back to // JobRecords. Already sorted by queue name (the scheduler sorts). Running/ // queued ids always resolve: the 100-job eviction never drops non-terminal // jobs, but filter defensively anyway. listQueues(): QueueSnapshot[] { return getScheduler() .queues() .map((v) => ({ name: v.name, running: v.running .map((id) => this.jobs.get(id)) .find((j): j is JobRecord => Boolean(j)), queued: v.queued .map((id) => this.jobs.get(id)) .filter((j): j is JobRecord => Boolean(j)), })); } // Names of queues that currently have any non-terminal jobs. Used by the // UI to populate the QueuePicker dropdown. activeQueueNames(): string[] { return getScheduler() .queues() .map((v) => v.name); } // Position of the job in its queue (0 = currently running). Returns -1 if the // job is no longer in any queue (terminal or unserialized ""). positionInQueue(id: string): number { return getScheduler().positionInQueue(id); } // Granular manual control: move a QUEUED job one slot toward (-1) or away from // (+1) the head, or to the front of the queued section. Never reorders a // running job. Returns false if the move isn't possible. Delegates to the // scheduler, which owns ordering. reorder(id: string, dir: -1 | 1): boolean { return getScheduler().reorder(id, dir); } promote(id: string): boolean { return getScheduler().promote(id); } } declare global { // eslint-disable-next-line no-var var __yttJobRegistry__: JobRegistry | undefined; } export function getRegistry(): JobRegistry { if (!globalThis.__yttJobRegistry__) { globalThis.__yttJobRegistry__ = new JobRegistry(); } return globalThis.__yttJobRegistry__; } export function newJobId(): string { // ULIDs: lexicographically time-sortable and time-decodable, so the /jobs // list can paginate over `.jobs` by filename without stat-ing every file. // See common/jobs/ulid.ts (jobIdTime handles both this and the legacy scheme). return ulid(); } // Which kinds honor the drain signal now lives in the job-kind metadata table // (jobKinds.ts), the single source of truth. Re-exported here so existing // importers of `isDrainableKind` from the registry keep resolving unchanged. export { isDrainableKind } from "./jobKinds";