import { isDrainableKind } from "../jobs/jobKinds"; import { jobIdTime } from "../jobs/ulid"; import type { JobRecord } from "../jobs/registry"; import type { JobListEntry } from "../jobs/listJobs"; import { jobSpecDetail } from "../jobs/jobDetail"; import type { ChannelStat } from "../controller/channels"; import type { QueueView } from "../jobs/scheduler"; import type { AutoRunnerInFlight } from "../controller/autoRunner"; import type { JobRowView, StuckReason } from "./jobRowView"; // THE ADAPTERS, THE ORDER AND THE MERGE — pure, and directive-free by // construction so both the server builder and the "use client" table can import // them. The precedent is components/lanes/laneState.ts: its only value imports // are things with no imports of their own, so nothing drags the registry, the // scheduler or `fs` into the client bundle. Here that is isDrainableKind // (common/jobs/jobKinds.ts, zero imports, already in the client graph via // jobKindLabels.ts) and jobIdTime (common/jobs/ulid.ts, zero imports). // Everything else is `import type`. // How long a finished job stays in the LIVE head. Long enough for the 5 s pulse // to re-render the tail with the same row twice over; short enough that the head // is still "now". export const RECENT_MS = 30_000; // A running slot is treated as "possibly stalled" (the soft case) once a // running/draining record has held its slot this long with no in-flight tasks // AND no progress: the clock runs from the later of its start and its last move // (`JobRecord.progressAt`). Age alone flagged a healthy metadata scan at ~10 // videos/min as STUCK on 2026-09-25 — it never has a task in flight, and it // advances its count every few seconds (release 9, B4a). // Distinct from the hard cases (terminal-but-holding-slot / record-evicted), // which are always safe to auto-heal; the soft case is only surfaced/force- // releasable, never auto-healed. const STUCK_AGE_MS = 10 * 60 * 1000; const TERMINAL = new Set(["done", "failed", "cancelled"]); export type SlotFact = { id: string; queueKey: string; role: "running" | "queued"; position: number; heldMs: number; stuck?: { reason: StuckReason; heldMs: number }; heal: boolean; }; // THE SCHEDULER AGAINST THE REGISTRY, as facts. The scheduler is the source of // truth for SLOTS and still references ids whose record is terminal or gone; // the registry is the source of truth for STATUS and never sees a slot. Pure: // it takes the queues snapshot and a lookup, reports what it found, and does // not heal — the builder heals, AFTER the rows are built, so a stuck row is // drawn once before it vanishes. The classification is the queue view's, // unchanged: // record-evicted and terminal-but-holding-slot are hard (heal); possibly-stalled // (a live record idle past STUCK_AGE_MS) is soft — surfaced, force-releasable, // never healed, since it may be a genuinely long operation. export function reconcileSlots( views: QueueView[], get: (id: string) => JobRecord | undefined, now: number, ): SlotFact[] { const facts: SlotFact[] = []; for (const v of views) { const ids = [...v.running, ...v.queued]; for (const [position, id] of ids.entries()) { const role: "running" | "queued" = position < v.running.length ? "running" : "queued"; const rec = get(id); const heldMs = now - (rec?.startedAt ?? rec?.queuedAt ?? now); let stuck: { reason: StuckReason; heldMs: number } | undefined; let heal = false; if (role === "running") { if (!rec) { // Scheduler references an id the registry evicted while it still // holds a running slot — a stuck slot with no live record to // describe it. stuck = { reason: "record-evicted", heldMs }; heal = true; } else if (TERMINAL.has(rec.status)) { stuck = { reason: "terminal-but-holding-slot", heldMs }; heal = true; } else { const ageMs = now - (rec.startedAt ?? rec.queuedAt); const quietMs = now - Math.max(rec.startedAt ?? rec.queuedAt, rec.progressAt ?? 0); const idle = !rec.tasks || rec.tasks.length === 0; if (idle && quietMs > STUCK_AGE_MS) { // Soft case: a live running/draining record sitting idle far too // long. Force-releasable, but deliberately NOT auto-healed. stuck = { reason: "possibly-stalled", heldMs: ageMs }; } } } facts.push({ id, queueKey: v.name, role, position, heldMs, stuck, heal }); } } return facts; } // THE PHANTOM: a slot whose record the registry no longer has. It is the one // row with no id the registry can answer for, so everything it knows comes from // the slot itself plus the on-disk meta sidecar the caller resolved the kind // from. export function fromSlot( f: SlotFact, kind: string, lastLogLine?: string, ): JobRowView { return { id: f.id, kind, status: "evicted", queueKey: f.queueKey, queuedAt: jobIdTime(f.id) || undefined, position: f.position, heldMs: f.heldMs, stuck: f.stuck, lastLogLine, source: "slot", }; } // Estimate seconds remaining as: remaining tasks × average measured task // duration, converted to wall-clock via the effective concurrency observed so // far. Returns undefined until at least one task has completed (no average // yet) or when there's no work left. `now` is injected so the same snapshot is // used across all jobs in a single payload build. function computeEtaSeconds( job: JobRecord, remaining: number, now: number, remainingAudioSeconds?: number, ): number | undefined { const count = job.completedTaskCount ?? 0; const totalMs = job.completedTaskMs ?? 0; if (count < 1 || remaining <= 0 || job.startedAt === undefined) { return undefined; } // Prefer an AUDIO-HOUR estimate where the work is proportional to length. // Averaging tasks assumes every unit costs about the same, which is true for // downloads and wildly false for digests: this corpus is ~77k videos and ~77k // audio-hours, and a channel of 9-hour VODs and a channel of 10-minute clips // have the same task count and a 50x difference in cost. A task average would // therefore quote an ETA that is wrong by more than an order of magnitude at // exactly the moment an operator most needs it — the start of an 80-day run. const doneAudio = job.completedTaskAudioSeconds ?? 0; if (doneAudio > 0 && remainingAudioSeconds && remainingAudioSeconds > 0) { const secondsPerAudioSecond = totalMs / 1000 / doneAudio; const elapsedMs = Math.max(1, now - job.startedAt); const concurrency = Math.max(1, totalMs / elapsedMs); return (remainingAudioSeconds * secondsPerAudioSecond) / concurrency; } const avgProcMs = totalMs / count; // measured average per task const elapsedMs = Math.max(1, now - job.startedAt); const concurrency = Math.max(1, totalMs / elapsedMs); // effective parallelism // Round up to whole parallel waves: the final wave runs its tasks // concurrently, so e.g. 2 tasks left across 4 workers still costs ~one full // task — not "half a task". Without the ceil the tail is underestimated. const waves = Math.ceil(remaining / concurrency); return (waves * avgProcMs) / 1000; } function computeJobProgressView( job: JobRecord, stat: ChannelStat | undefined, now: number, ): JobRowView["progress"] { const snap = job.progress; // A runner-reported `current` needs no channel stat: a batch spanning // channels (fetch-windows) has none, and still knows its own count. if (!snap || (!stat && snap.current === undefined)) return undefined; // `current` is RE-COUNTED from disk (readChannelStat), never reported by the // runner — which is why the digest metric needed its own on-disk counter // (digestCount) rather than a number the batch could have just told us. // A runner-reported `current` WINS. The disk re-count below cannot see a // regeneration — a regenerated digest is rewritten in place, so the file // count never moves and the bar sits at 0% for the whole job. Only the runner // knows it did the work. Downloads and transcripts report nothing and keep // the disk re-count, unchanged. const current = snap.current ?? (snap.metric === "downloads" ? stat!.downloadCount : snap.metric === "digests" ? (stat!.digestCount ?? 0) : stat!.transcriptCount); const range = Math.max(0, snap.target - snap.initial); const advance = Math.max(0, current - snap.initial); const pct = range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100)); const remaining = Math.max(0, snap.target - current); return { metric: snap.metric, initial: snap.initial, current, target: snap.target, pct, etaSeconds: computeEtaSeconds( job, remaining, now, snap.remainingAudioSeconds, ), }; } export type FromRecordContext = { now: number; stat?: ChannelStat; workerNames: ReadonlyMap; // 0 = running head, 1.. = queued; -1 for a parallel ("") job the scheduler // never tracked. position: number; // The last index in this record's queue, for the reorder bounds. queueMax: number; }; // The live adapter: a registry record as a row. export function fromRecord(j: JobRecord, ctx: FromRecordContext): JobRowView { return { id: j.id, kind: j.kind, status: j.status, queueKey: j.queueKey, channelSlug: j.channelSlug, videoId: j.videoId, progress: j.status === "running" && (j.channelSlug || j.progress?.current !== undefined) ? computeJobProgressView(j, ctx.stat, ctx.now) : undefined, tasks: j.tasks?.map((t) => ({ id: t.id, label: t.label, kind: t.kind, fraction: t.fraction, detail: t.detail, startedAt: t.startedAt, workerId: t.workerId, workerName: t.workerId ? ctx.workerNames.get(t.workerId) : undefined, phase: t.phase, probeStartedAt: t.probeStartedAt, probeEtaSeconds: t.probeEtaSeconds, })), draining: j.draining === true, drainable: j.status === "running" && isDrainableKind(j.kind) && j.draining !== true, background: j.background === true, // A queued job can move up/promote if it isn't the first queued (position > // 1, since the running head is at 0), and down if it isn't the last in its // queue. Running and parallel ("") jobs are never reorderable. canMoveUp: j.status === "queued" && ctx.position > 1, canMoveDown: j.status === "queued" && ctx.position >= 1 && ctx.position < ctx.queueMax, queuedAt: j.queuedAt, startedAt: j.startedAt, endedAt: j.endedAt, exitCode: j.exitCode, inRegistry: true, replayable: Boolean(j.spec), position: ctx.position >= 0 ? ctx.position : undefined, detail: jobSpecDetail(j.kind, j.spec), source: "registry", }; } // The history adapter: a .log + sidecar entry as a row. Field for field — // listJobs already resolved the registry record where there was one. export function fromEntry(e: JobListEntry): JobRowView { return { id: e.id, kind: e.kind ?? "", status: e.status, queueKey: e.queueKey, channelSlug: e.channelSlug, videoId: e.videoId, queuedAt: e.queuedAt, startedAt: e.startedAt, endedAt: e.endedAt, exitCode: e.exitCode, logSize: e.logSize, inRegistry: e.inRegistry, replayable: e.replayable, detail: e.detail, ...(e.cancelReason ? { cancelReason: e.cancelReason } : {}), source: "archive", }; } function liveRank(r: JobRowView): number { if (r.stuck) return 0; if (r.status === "running") return 1; if (r.status === "queued") return 2; return 3; } // The head's order: what is wrong first, then what is working, then what is // waiting IN THE ORDER IT WILL RUN, then what just finished. // // Queued rows sort by (queue name, position) rather than by recency — finding 8 // of plans/editor-ia-slice-8c.md. That is a change: the Active Jobs screen // ordered queued jobs by the registry's queuedAt, which put the job that will // run next anywhere in the list. "2nd in line" in the Queue cell only means anything if // the rows are in line. export function orderLiveRows(rows: JobRowView[]): JobRowView[] { return [...rows].sort((a, b) => { const ra = liveRank(a); const rb = liveRank(b); if (ra !== rb) return ra - rb; if (ra === 0) return (b.heldMs ?? 0) - (a.heldMs ?? 0); // Running: registry order, newest queued first — unchanged. if (ra === 1) return (b.queuedAt ?? 0) - (a.queuedAt ?? 0); if (ra === 2) { const byQueue = (a.queueKey ?? "").localeCompare(b.queueKey ?? ""); if (byQueue !== 0) return byQueue; return (a.position ?? 0) - (b.position ?? 0); } return (b.endedAt ?? 0) - (a.endedAt ?? 0); }); } // LIVE WINS BY ID: the same job is in the registry and on disk at once, and its // live row is the one with progress. The head keeps its order; the tail follows // in the order the directory listing gave it, minus anything already above. export function mergeJobRows( live: JobRowView[], history: JobRowView[], ): JobRowView[] { const seen = new Set(live.map((r) => r.id)); return [...live, ...history.filter((r) => !seen.has(r.id))]; } // Is this row still moving? Drives the 1 s poll: a page with nothing live makes // no requests at all. export function isLive(r: JobRowView): boolean { return ( r.status === "running" || r.status === "queued" || r.status === "evicted" || Boolean(r.stuck) ); } // AN AUTO-QUEUE LANE'S IN-FLIGHT UNIT, as a job row — so the operations board // draws what the runner is doing with the same row every other surface draws // a job with. A download unit is a real registry job (its id rides on the unit // as `jobId`, set when the runner's child job starts); a transcription, digest // or backfill unit is a task on the runner's own job and has no page of its // own, so its id is synthetic and `inRegistry` is false. // // `detail` ("rule 3 · newest first from alpha") is the caller's: the leaf // sentence is an editor helper over the lane's policy, not a pure fact of the // unit. A unit with a `note` (the download lane's channel-scoped metadata // scan, whose `videoId` names no directory) says the note instead, and carries // no videoId, so nothing links to a video that is not there. export function fromInFlight( unit: AutoRunnerInFlight, ctx: { kind: string; runnerJobId?: string | null; detail?: string }, ): JobRowView { return { id: unit.jobId ?? `${ctx.runnerJobId ?? ctx.kind}:${unit.videoId}`, kind: ctx.kind, status: "running", channelSlug: unit.channelSlug, videoId: unit.note ? undefined : unit.videoId, startedAt: unit.startedAt, background: true, inRegistry: Boolean(unit.jobId), detail: unit.note ?? ctx.detail, source: "runner", }; }