import { getRegistry, type JobRecord, type JobStatus, type JobTaskKind, } from "../jobs/registry"; import { jobKindLabel, kindNeedsMedia } from "../jobs/jobKinds"; import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; import { getAutoRunnerStatus, type AutoRunnerInFlight } from "./autoRunner"; // WHO IS WRITING INTO A CHANNEL'S MEDIA RIGHT NOW — one answer, named. // // A media move must not start over a writer (release 16 slice RM). On // 2026-09-30 a move of `realcandaceo` was queued at 21:16:42 and asked the // question then, when nothing was running; a "Transcribe all" started at // 21:23:30, and the move — which had waited twenty minutes behind other moves // on the one relocation queue — started its copy at 21:36:28 without asking // again. The transcription of `v50t5yt` finished at 21:38:07, after rsync had // passed that directory, and the verify refused a 13.57 GB copy. So the move's // PREVIEW and its JOB'S FIRST STEP ask here too, not only the action that // enqueued it, and the refusal names the writer so the operator knows what to // wait for or cancel. // // A GHOST NEVER HOLDS A MOVE (release 17 slice D0). Only THIS process's // registry is read, never a `.meta.json`: a meta a dead process left // `running` (three `refresh-report`s on 2026-10-01) is not a record here and // cannot name a writer. The boot pass closes such metas as interrupted // (jobs/bootQueuedJobs.ts `settleRunningJobMetas`) so /jobs stops showing them; // nothing here needs to know. Keep it that way: a reader of metas here would // have to ask whether the writer is alive (`writerIsGone`) first. // // TWO HALVES, for the reason editor/app/channels/lib/mediaBusy.ts gives: the // job registry is half the truth. The auto-queue lanes run their per-video // units in-process and make no job record (the omnimirror incident, // 2026-09-13), so `getAutoRunnerStatus(lane).inFlight` is the other half. Both // are memory reads — no disk, no await. // // SERVER-SIDE ONLY: it reaches the runner, which imports execa transitively. export type ChannelWriter = | { source: "job"; jobId: string; kind: string; // The kind's label (jobKindLabel), or the raw kind when it has none. label: string; // "stopping": cancelled, but its function has not returned yet — a // transcriber finishing its current window still writes. The registry // stamps `endedAt` only when the job has actually stopped. status: Extract | "stopping"; // What the job is on now: its first in-flight task, else the video the // record names (a single-video job). videoId?: string; taskKind?: JobTaskKind; } | { source: "lane"; lane: AutoQueueKind; videoId: string; // A unit that is not a video (the download lane's metadata scan) says // what it is instead. note?: string; // The registry job running the unit, when it is one. jobId?: string; }; // Where the two halves are read from. Injected by the tests; production reads // the live registry and the live runners. export type ChannelWritersSource = { jobs: () => ReadonlyArray; units: () => ReadonlyArray<{ lane: AutoQueueKind; unit: AutoRunnerInFlight }>; }; export const liveChannelWritersSource: ChannelWritersSource = { jobs: () => getRegistry().list(), units: () => LANES.flatMap((lane) => getAutoRunnerStatus(lane).inFlight.map((unit) => ({ lane, unit })), ), }; export type ChannelWritersOptions = { // Count queued jobs too. The editor's courtesy check does (a queued job is // about to write); the move's own check does not — a queued media job that // starts after the marker is written refuses itself at its start // (streamCommand.ts), so refusing the move over it would trade the // operator's twenty-minute wait in the relocation queue for nothing. includeQueued?: boolean; // Kinds that are not writers for this question. The move's own job passes // `relocate-channel-media`: it is running on this channel's slug, it is the // asker, and the relocation queue runs one move at a time. ignoreKinds?: ReadonlyArray; // MEDIA WRITERS ONLY (release 17): a move of the media tier holds only the // jobs that open or write a big file (`kindNeedsMedia`) and the lanes that // do (every lane but digest). A digest, a normalize, an availability check // reads and writes the text, which never moves, so it may run during a move. mediaOnly?: boolean; source?: ChannelWritersSource; }; // Every writer on `slug`, jobs first (running before queued), then lane units. export function channelWriters( slug: string, opts: ChannelWritersOptions = {}, ): ChannelWriter[] { const source = opts.source ?? liveChannelWritersSource; const ignore = new Set(opts.ignoreKinds ?? []); const running: ChannelWriter[] = []; const queued: ChannelWriter[] = []; for (const j of source.jobs()) { if (j.channelSlug !== slug || ignore.has(j.kind)) continue; if (opts.mediaOnly && !kindNeedsMedia(j.kind)) continue; // A CANCEL IS A REQUEST, NOT AN EXIT. The registry marks a running job // `cancelled` the moment it is asked to stop, and stamps `endedAt` when // its function has returned (registry.ts `finalize`). Until then it is // still a writer, "stopping" — an operator who cancels and moves at once // would otherwise start a copy under it. const stopping = j.status === "cancelled" && j.endedAt === undefined; if ( j.status !== "running" && !stopping && !(opts.includeQueued && j.status === "queued") ) { continue; } const task = (j.tasks ?? []).find((t) => t.kind !== "relocate"); (j.status === "queued" ? queued : running).push({ source: "job", jobId: j.id, kind: j.kind, label: jobKindLabel(j.kind), status: stopping ? "stopping" : (j.status as "running" | "queued"), ...(task ? { videoId: task.id, taskKind: task.kind } : {}), ...(!task && j.videoId ? { videoId: j.videoId } : {}), }); } const jobs = [...running, ...queued]; const units: ChannelWriter[] = source .units() .filter(({ unit }) => unit.channelSlug === slug) .filter(({ lane }) => !opts.mediaOnly || lane !== "digest") .map(({ lane, unit }) => ({ source: "lane" as const, lane, videoId: unit.videoId, ...(unit.note ? { note: unit.note } : {}), ...(unit.jobId ? { jobId: unit.jobId } : {}), })); // A lane's download unit is ALSO a registry job (auto-download-unit), so it // would be named twice; the lane's wording is the one that says where it // came from. const unitJobs = new Set( units.flatMap((u) => (u.source === "lane" && u.jobId ? [u.jobId] : [])), ); return [ ...jobs.filter((w) => !(w.source === "job" && unitJobs.has(w.jobId))), ...units, ]; } const TASK_NOUN: Record = { download: "a download", transcribe: "a transcription", digest: "a digest", backfill: "a backfill", relocate: "a media move", }; const LANE_NOUN: Record = { transcription: "a transcription", download: "a download", digest: "a digest", backfill: "a backfill", }; // One writer, as a clause: "a transcription of v50t5yt is running (Transcribe // all, job 01M3…)", "Sync is queued (job 01M3…)", "a download of abc is // running (the download lane)". export function describeChannelWriter(w: ChannelWriter): string { if (w.source === "lane") { const what = w.note ?? `${LANE_NOUN[w.lane]} of ${w.videoId}`; return `${what} is running (the ${w.lane} lane)`; } if (w.taskKind && w.videoId) { return ( `${TASK_NOUN[w.taskKind]} of ${w.videoId} is ${w.status} ` + `(${w.label}, job ${w.jobId})` ); } const what = w.videoId ? `${w.label} of ${w.videoId}` : w.label; return `${what} is ${w.status} (job ${w.jobId})`; } // What the operator does about the first writer: a job is waited for or // cancelled on /jobs (one already stopping, only waited for); a lane's unit is // waited for, or its lane held. function waysOut(w: ChannelWriter): string { if (w.source === "lane") return `wait for it, or hold the ${w.lane} lane`; return w.status === "stopping" ? "wait for it to stop" : "wait for it or cancel it"; } // The first writer named, the rest counted, and what to do — or null when // nothing is writing. `doing` is the clause the refusal opens with ("move the // media of"), so the same sentence reads right for each mover. export function channelWritersRefusal( slug: string, writers: ReadonlyArray, doing = "move the media of", ): string | null { if (writers.length === 0) return null; const [first, ...rest] = writers; const more = rest.length > 0 ? `, and ${rest.length} more writer(s)` : ""; return ( `Cannot ${doing} "${slug}" now: ${describeChannelWriter(first)}${more} — ` + `${waysOut(first)}. Nothing has been touched.` ); }