import path from "node:path"; import { mkdir } from "node:fs/promises"; import { createWriteStream, type WriteStream } from "node:fs"; import { execa } from "execa"; import { getRegistry, newJobId, noteProgress, type JobProgress, type JobRecord, type JobStatus, type JobTask, } from "./registry"; import type { Paths } from "../lib/paths"; import { makeSafeController } from "../lib/safeStreamController"; import { requestChannelSnapshot, shouldRequestSnapshot, } from "./snapshotScheduler"; import { writeJobMeta } from "./jobMeta"; import { maybePruneJobLogs } from "./listJobs"; import type { JobSpec } from "./jobSpec"; import { kindNeedsMedia, kindNeedsText } from "./jobKinds"; import { assertChannelMediaReachable, assertChannelTextReadable, ChannelMediaUnreachableError, } from "../lib/channelMedia"; import { mediaHoldText } from "../lib/channelMediaHold"; // Mark a job's channel report dirty so the debounced scheduler regenerates the // snapshot — called both on each completed sub-operation and on the job's // terminal state. Excludes the regen job kind itself (and other read-only // kinds) to avoid an infinite loop. Skips jobs whose record is no longer in the // registry: that means the registry was reset out from under a still-running // job (e.g. the e2e cache-reset between specs), and a wiped job must not arm a // regen that would bleed into unrelated work. Synchronous and non-throwing. function requestSnapshotOnFinish(jobId: string, opts: CommonOpts): void { if (!opts.channelSlug || !shouldRequestSnapshot(opts.kind)) return; if (!getRegistry().get(jobId)) return; requestChannelSnapshot(opts.paths, opts.channelSlug); } // Resolves when a started job reaches a terminal state. Lets a server-side // caller (e.g. the auto-download runner) enqueue a job on a serialized queue // and await its completion — the queue may hold it behind other work, so the // `{ jobId, stream }` return alone can't tell the caller when it actually ran. export type JobDoneResult = { status: JobStatus; jobId: string }; export type StreamActionResult = | { ok: true; jobId: string; stream: ReadableStream; done: Promise; } // `info: true` marks a non-error outcome that started no job (e.g. a // re-derived bucket that's currently empty) so the UI can show it neutrally // rather than as a red failure. | { ok: false; error: string; info?: boolean }; type CommonOpts = { kind: string; queueKey: string; paths: Paths; channelSlug?: string; videoId?: string; // When set, recorded on the job (and its meta sidecar) so the job can be // re-launched (retried) later. See common/jobs/jobSpec.ts. spec?: JobSpec; // Background jobs queue BEHIND any foreground (default) job on the same // queueKey: a manually-triggered job (sync, manual download) jumps ahead of // queued background work without interrupting a running one. The auto-download // runner marks its per-video units background so a clicked Sync preempts them. background?: boolean; }; export type RunManagedCommandOpts = CommonOpts & { cwd: string; command: string; args: string[]; env?: Record; }; // Passed as the optional 4th argument to a managed function. Existing callers // declare three-arg arrows and silently ignore it; batch / single-op callers // opt in to per-task progress reporting and drain (soft-cancel) support. export type JobRunContext = { jobId: string; drainSignal: AbortSignal; addTask: (task: JobTask) => void; updateTask: ( taskId: string, patch: { fraction?: number; detail?: string; phase?: "downloading" | "probing"; probeEtaSeconds?: number; }, ) => void; removeTask: (taskId: string) => void; recordTaskDone: (durationMs: number, audioSeconds?: number) => void; }; export type RunManagedFunctionOpts = CommonOpts & { fn: ( onLog: (line: string) => void, signal: AbortSignal, setProgress: (snap: JobProgress) => void, ctx: JobRunContext, ) => Promise; }; async function ensureJobsDir(paths: Paths): Promise { await mkdir(paths.jobsDir, { recursive: true }); } function makeJob( kind: string, queueKey: string, paths: Paths, channelSlug?: string, videoId?: string, spec?: JobSpec, background?: boolean, ): { id: string; logPath: string; record: JobRecord } { const id = newJobId(); const logPath = path.join(paths.jobsDir, `${id}.log`); const record: JobRecord = { id, kind, queueKey, channelSlug, videoId, spec, background, status: "queued", queuedAt: Date.now(), logPath, }; getRegistry().register(record); return { id, logPath, record }; } // Best-effort error listener for the log WriteStream: writes can fail (disk // full, permissions, etc.) and would otherwise become uncaughtExceptions. function ignoreFileStreamErrors(stream: WriteStream): void { stream.on("error", () => {}); } // A one-shot completion deferred attached to a managed job. `settle` is called // from every terminal path (finalized done/failed/cancelled, or cancelled // before start); the first call wins so callers awaiting `done` see the true // terminal status exactly once. function makeDoneDeferred(jobId: string): { done: Promise; settle: (status: JobStatus) => void; } { let resolveDone!: (r: JobDoneResult) => void; const done = new Promise((res) => { resolveDone = res; }); let settled = false; const settle = (status: JobStatus) => { if (settled) return; settled = true; resolveDone({ status, jobId }); }; return { done, settle }; } // ONE WRITER PER JOB, IN ORDER — DEFENSIVE. A job's sidecar is written at // least twice (at enqueue, then with its terminal state) and each write is // async. writeJobMeta snapshots the record before its first await, so two // writes are ISSUED in order; only the fs threadpool could complete them out // of order, leaving "queued" over "cancelled" or a shorter JSON over a longer // one's tail. A 300-job probe (release 13 W1 review) never saw that happen // without the chain. Chained, each write starts when the previous one has // settled and serializes the record as it is THEN. function metaWriter(paths: Paths, record: JobRecord): () => void { return serialWriter(() => writeJobMeta(paths, record)); } // Run `write` once per call, each call starting only after the previous one // SETTLED — resolved or rejected. writeJobMeta never rejects today; a writer // that did must not silently stall every later write, hence // `then(write, write)`, nor raise an unhandled rejection from the last one, // hence the no-op catch (the writer owns its errors, as writeJobMeta does). // Exported for its unit test. export function serialWriter(write: () => Promise): () => void { let last: Promise = Promise.resolve(); return () => { last = last.then(write, write); last.catch(() => {}); }; } export async function runManagedCommand( opts: RunManagedCommandOpts, ): Promise { const registry = getRegistry(); await ensureJobsDir(opts.paths); const { id, logPath, record } = makeJob( opts.kind, opts.queueKey, opts.paths, opts.channelSlug, opts.videoId, opts.spec, opts.background, ); const { done, settle } = makeDoneDeferred(id); const persistMeta = metaWriter(opts.paths, record); const safe = makeSafeController(); let fileStream: WriteStream | null = null; let cancelledBeforeStart = false; const stream = new ReadableStream({ start(c) { safe.setController(c); }, cancel() { // Consumer disconnected (page navigation, tab close, RSC response // tear-down). Stop pushing into the controller — but keep the child // and the on-disk log going so the job runs to completion and can be // observed by another reconnecting client via /api/jobs//log. // Explicit user cancels go through `registry.cancel()` (which kills // the child), not through this source-cancel hook. safe.markClosed(); }, }); const start = () => { if (cancelledBeforeStart) { safe.safeClose(); return; } fileStream = createWriteStream(logPath); ignoreFileStreamErrors(fileStream); const child = execa(opts.command, opts.args, { cwd: opts.cwd, env: opts.env, all: true, buffer: false, reject: false, }); record.child = child; child.all?.on("data", (chunk: Buffer) => { fileStream!.write(chunk); safe.safeEnqueue(chunk.toString("utf8")); }); child.all?.on("error", () => {}); child .then((result) => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); return; } const status = result.exitCode === 0 ? "done" : "failed"; registry.finalize(id, status, result.exitCode ?? undefined); }) .catch((err) => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); return; } safe.safeEnqueue(`\n[error] ${(err as Error).message}\n`); registry.finalize(id, "failed"); }) .finally(() => { fileStream?.end(); safe.safeClose(); requestSnapshotOnFinish(id, opts); // Persist terminal state (status/endedAt/exitCode now set by finalize). persistMeta(); // Throttled retention so the .jobs directory stays bounded on its own. void maybePruneJobLogs(opts.paths); settle(record.status); }); }; // Cancelled while still queued: start() never runs, so its .finally never // writes the terminal sidecar — this does (release 13 slice W1; it used to // stay "queued", and the next boot could re-queue a job the operator had // cancelled). registry.cancel() marks the record cancelled BEFORE the // scheduler fires this; at a graceful shutdown it does not, and a record // still `queued` is left on disk as it is, for the boot pass. const onCancel = () => { cancelledBeforeStart = true; if (record.status === "cancelled") persistMeta(); safe.safeClose(); settle("cancelled"); }; registry.enqueue(record, { start, onCancel }); // Persist queued/running identity up front so a mid-run crash still leaves a // sidecar; the .finally above (or onCancel) rewrites it with the terminal // state. persistMeta(); return { ok: true, jobId: id, stream, done }; } // GUARD 1 OF FOUR (see plans/relocate-channel-media.md). runManagedFunction is // the funnel every job-shaped action goes through, and 35 of its 53 call sites // already pass a channelSlug — so one check here covers every per-channel media // action without touching any of them. A refusal happens BEFORE the job record // is made: no queued job, no log, no sidecar, just `{ ok: false }` carrying the // reason, which every caller already renders. // // The kind decides. `needsMedia` is declarative on JobKindMeta and absent means // false, so a bookkeeping kind is never refused for a drive it does not read, // and the relocate job itself — the thing that FIXES an unreachable channel — // must never declare it. // // ASKED TWICE: here, before the record exists, and again when the queue STARTS // the job (`start` below). A job can wait hours in its queue — a "Transcribe // all" behind another channel's on the one transcription queue — and a move // can begin in that wait; the answer it got at enqueue is then stale, and the // job would write into a tree being copied. At the start the refusal fails the // job with the sentence in its log. // // A MOVE IS A HOLD, and the refusal says so in the hold's words ("held: its // media is moving …") rather than calling the media unreachable. // // A TEXT KIND ASKS THE TEXT GUARD (release 17): a kind that declares // `needsText` reads only `data/`'s text, which stays on the corpus disk, so it // runs while the channel's media is moving, stalled or unmounted, and is // refused only where the text itself cannot be read (a `legacy` channel, a // `data/` that is not a directory, a tier migration in flight). async function refuseForUnreachableMedia( opts: CommonOpts, ): Promise { if (!opts.channelSlug) return null; if (!kindNeedsMedia(opts.kind)) { if (!kindNeedsText(opts.kind)) return null; try { await assertChannelTextReadable(opts.paths, opts.channelSlug); return null; } catch (err) { return (err as Error).message; } } try { await assertChannelMediaReachable(opts.paths, opts.channelSlug); return null; } catch (err) { if ( err instanceof ChannelMediaUnreachableError && err.status === "in-transition" ) { return ( `Channel "${opts.channelSlug}" is ${mediaHoldText(err.status)} — ` + `${err.location.detail ?? "a relocation marker is present"}. Its media ` + `jobs start again when the move completes, or when its marker is ` + `cleared on the channel's Storage panel.` ); } return (err as Error).message; } } export async function runManagedFunction( opts: RunManagedFunctionOpts, ): Promise { const refusal = await refuseForUnreachableMedia(opts); if (refusal) return { ok: false, error: refusal }; const registry = getRegistry(); await ensureJobsDir(opts.paths); const { id, logPath, record } = makeJob( opts.kind, opts.queueKey, opts.paths, opts.channelSlug, opts.videoId, opts.spec, opts.background, ); const { done, settle } = makeDoneDeferred(id); const persistMeta = metaWriter(opts.paths, record); const safe = makeSafeController(); let fileStream: WriteStream | null = null; let cancelledBeforeStart = false; const stream = new ReadableStream({ start(c) { safe.setController(c); }, cancel() { // Consumer disconnected (page navigation, tab close, RSC response // tear-down). Stop pushing into the controller — but let the function // run to completion, writing to disk; reconnecting clients can poll // /api/jobs//log. Explicit user cancels go through // `registry.cancel()` which aborts this controller. safe.markClosed(); }, }); const start = () => { if (cancelledBeforeStart) { safe.safeClose(); return; } fileStream = createWriteStream(logPath); ignoreFileStreamErrors(fileStream); const abort = new AbortController(); record.abortController = abort; const drain = new AbortController(); record.drainController = drain; const onLog = (line: string) => { const text = line.endsWith("\n") ? line : `${line}\n`; fileStream!.write(text); safe.safeEnqueue(text); }; const setProgress = (snap: JobProgress) => { noteProgress(record, snap); }; const ctx: JobRunContext = { jobId: id, drainSignal: drain.signal, addTask: (task) => registry.addTask(id, task), updateTask: (taskId, patch) => registry.updateTask(id, taskId, patch), removeTask: (taskId) => registry.removeTask(id, taskId), recordTaskDone: (ms, audioSeconds) => { registry.recordTaskDuration(id, ms, audioSeconds); // Refresh the report after EACH completed sub-operation (each video // downloaded/transcribed in a batch), not only when the whole batch // finishes — so a long batch updates incrementally. The global debounce // coalesces sub-operations that finish close together. requestSnapshotOnFinish(id, opts); }, }; // THE GUARD AGAIN, NOW THAT THE QUEUE HAS STARTED THE JOB (see // refuseForUnreachableMedia): a refusal here fails the job before `fn` // touches anything, with the sentence as its log's `[error]` line. const run = async () => { const refusal = await refuseForUnreachableMedia(opts); if (refusal) throw new Error(refusal); await opts.fn(onLog, abort.signal, setProgress, ctx); }; run() .then(() => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); return; } registry.finalize(id, "done"); }) .catch((err) => { if (record.status === "cancelled") { registry.finalize(id, "cancelled"); return; } onLog(`[error] ${(err as Error).message}`); registry.finalize(id, "failed"); }) .finally(() => { fileStream?.end(); safe.safeClose(); requestSnapshotOnFinish(id, opts); // Persist terminal state (status/endedAt/exitCode now set by finalize). persistMeta(); // Throttled retention so the .jobs directory stays bounded on its own. void maybePruneJobLogs(opts.paths); settle(record.status); }); }; // Cancelled while still queued: start() never runs, so its .finally never // writes the terminal sidecar — this does (release 13 slice W1; it used to // stay "queued", and the next boot could re-queue a job the operator had // cancelled). registry.cancel() marks the record cancelled BEFORE the // scheduler fires this; at a graceful shutdown it does not, and a record // still `queued` is left on disk as it is, for the boot pass. const onCancel = () => { cancelledBeforeStart = true; if (record.status === "cancelled") persistMeta(); safe.safeClose(); settle("cancelled"); }; registry.enqueue(record, { start, onCancel }); // Persist queued/running identity up front so a mid-run crash still leaves a // sidecar; the .finally above (or onCancel) rewrites it with the terminal // state. persistMeta(); return { ok: true, jobId: id, stream, done }; }