import { fromEntry, fromRecord, fromSlot, orderLiveRows, reconcileSlots, RECENT_MS, } from "./jobRows"; // Directive-free and value-import-free by construction (laneState.ts imports // only tone maps and a type), so a server module can read the same derivation // the client rail does. import { deriveLaneState } from "./laneState"; import { autoRunnerJobKind } from "../controller/autoRunner"; import { isGateHeld } from "../lib/pauseGates"; import { LANES, type AutoQueueKind } from "../lib/autoQueueTypes"; import type { AutoRunnerIdleReason, AutoRunnerStatus, } from "../controller/autoRunner"; import type { JobRecord } from "../jobs/registry"; import type { QueueView } from "../jobs/scheduler"; import type { JobMeta } from "../jobs/jobMeta"; import type { JobsPage } from "../jobs/listJobs"; import type { ChannelStat } from "../controller/channels"; import type { DiskGateReason, DiskGateStatus } from "../lib/diskSpace"; import type { JobRowView } from "./jobRowView"; import type { LiveInputs } from "./inputs"; // THE READERS ARRIVE AS FUNCTIONS, not as data. Every other view in this // directory can be handed the rows it folds, because its inputs are known // before it runs. This one cannot: which channels to stat is not knowable // until `registry.list()` has been read, and which job ids need a meta sidecar // is not knowable until the scheduler's slots have been compared against it. // So the disk stays out of here the only other way — the caller supplies the // reader and the view decides what to ask for. export type LiveJobRowsInputs = LiveInputs & { channelStat: (slug: string) => Promise; }; // The full payload also needs the three things a per-page row list must NOT // have: the log tail (read per stuck row), the disk gate (already sampled by // the shell, in "observe" mode) and the runner statuses behind the lane strip. export type ActiveJobsInputs = LiveJobRowsInputs & { jobMeta: (id: string) => Promise; tailLog: (id: string) => Promise; disk: DiskGateStatus; autoRunnerStatus: (kind: AutoQueueKind) => AutoRunnerStatus; }; export type DiskStatusView = { // Whether the low-disk gate is configured (minFreeDiskGB > 0). When false the // UI hides the indicator entirely. enabled: boolean; freeBytes: number; thresholdBytes: number; // True when the gate is enabled and free space is at/below the floor — i.e. // downloads are currently being blocked. low: boolean; // The bar a stopped pipeline has to clear to resume (floor + resume margin). // Equal to thresholdBytes when the margin is 0. resumeBytes: number; // Why, in one word, so a surface can say "disk" rather than leaving the // operator to read a red "downloads paused" as the manual toggle. reason: DiskGateReason; // Log-ready explanation; empty when not low. message: string; }; // A RUNNER IS A LANE, NOT A JOB. // // The auto-queue runners are channel-less jobs, so this screen grouped them by // kind and gave each group the same bordered card a real channel gets — a whole // card, per always-on daemon, to say "running". With four pipelines converging // on this model that is four cards stacked above the work that actually has // progress bars. // // So a runner ships as one of these instead: a line on a strip, in the lane // vocabulary the dashboard already uses (Running / Holding / Idle / Off), with // its reason in words. The JOB is still in `jobs` and still carries its Drain // and Cancel controls — this is the same fact, said in one line instead of a // card. export type ActiveLaneView = { // The job kind, so the client can find the runner's own job row. kind: string; label: string; state: "running" | "holding" | "idle" | "unavailable"; // Why it is not working, in words. Null when it IS working. note: string | null; inFlight: number; }; export type ActiveJobsPayload = { // Running and queued work, plus (commit 3) the scheduler slots with no live // record behind them. jobs: JobRowView[]; // Registry records that went terminal within the last RECENT_MS, newest // first — so a job that finishes between two polls is shown finishing rather // than vanishing. The /jobs head keeps it; the widget and the dashboard // ignore this field on purpose (idle must flip when work ends). recent: JobRowView[]; // The server's clock at build time. The /jobs table renders whichever of the // SSR prop and its own poll is newer, so a page served from the router cache // cannot show a finished job as running again. builtAt: number; // The scheduler's own health, for the line above the table. `running` and // `queued` count SLOTS, not registry records: a parallel ("") job holds no // slot at all (registry.ts), which is why the dashboard's running · queued // chip keeps counting rows instead of reading these. summary: { queues: number; running: number; queued: number; stuck: number; workersPaused: boolean; workerCount: number; }; disk: DiskStatusView; // Every lane, whether or not it has a job right now. A stopped runner still // gets a line: "not running" is exactly the state a screen full of channel // cards used to hide. lanes: ActiveLaneView[]; }; // THE ONE BUILDER FOR EVERY JOB LIST. // // Shared by the /jobs page (initial state), the dashboard, the widget and the // /api/jobs/active poll route. Five pages hand-rolled a six-field copy of the // row this returns and so drew the same job with progress bars on Active Jobs // and none on its channel page; liveJobRows() below is what they call instead. async function rowsForRecords( i: LiveJobRowsInputs, records: JobRecord[], views: QueueView[], now: number, ): Promise { const channelSlugs = Array.from( new Set( records.map((j) => j.channelSlug).filter((s): s is string => Boolean(s)), ), ); const statResults = await Promise.all( channelSlugs.map((slug) => i.channelStat(slug)), ); const channelStats = new Map(); for (let n = 0; n < channelSlugs.length; n++) { const stat = statResults[n]; if (stat) channelStats.set(channelSlugs[n], stat); } // Map worker id → display name so each transcribe task can show which worker // it's running on. const workerNames = new Map(); for (const w of i.pool.summary()) workerNames.set(w.id, w.name); // Queue positions, to derive each queued job's reorder bounds, read off the // ONE scheduler snapshot rather than a positionInQueue() call per record: a // slot completing mid-build would otherwise move a position under us, and two // rows would disagree about who is next (the queue view made the same point). // position 0 is the running head; queued jobs occupy 1..max. A parallel ("") // job is in no queue, has position -1 and is never reorderable. const positions = new Map(); const queueMaxPos = new Map(); for (const v of views) { const ids = [...v.running, ...v.queued]; for (const [n, id] of ids.entries()) positions.set(id, n); queueMaxPos.set(v.name, ids.length - 1); } return records.map((j) => fromRecord(j, { now, stat: j.channelSlug ? channelStats.get(j.channelSlug) : undefined, workerNames, position: positions.get(j.id) ?? -1, queueMax: queueMaxPos.get(j.queueKey) ?? 0, }), ); } // THE ONE BUILDER FOR EVERY JOB LIST, for a page that wants only its own jobs. // No lanes, no disk, no reconciliation and NO HEAL here — a channel page // reading its own jobs must not free a scheduler slot as a side effect. export async function liveJobRows( i: LiveJobRowsInputs, filter: (j: JobRecord) => boolean, ): Promise { const records = i.registry .list() .filter((j) => j.status === "running" || j.status === "queued") .filter(filter); const rows = await rowsForRecords(i, records, i.scheduler.queues(), i.now()); return orderLiveRows(rows); } // The tail's adapter, beside the head's: one page of the on-disk listing as // rows. Paging, ordering and totals are listAllJobs's, unchanged — the read // itself is the shell's, so what is left here is the fold over the page. export function rowsFromJobsPage(page: JobsPage): { rows: JobRowView[]; hasMore: boolean; total: number; } { return { rows: page.entries.map(fromEntry), hasMore: page.hasMore, total: page.total, }; } // EVERY CONSUMER OF THIS PAYLOAD HEALS THE SCHEDULER'S DRIFT WHEN IT OBSERVES // IT. The queue-diagnostics page used to be the one place that compared the // scheduler's slots to the registry's statuses and freed a slot whose record // was terminal or evicted; the Active Jobs page hid exactly those rows. Now // the /jobs page, the dashboard, the widget and /api/jobs/active all build this, // so any of them observing a terminal record holding a running slot frees it — // AFTER the rows are built, so the stuck row is drawn once. The decision is the // operator's (plans/editor-ia-slice-8c.md): a running slot with a terminal // record is never right, and the surface that sees it is the surface that // fixes it. The soft "possibly-stalled" case is surfaced, never auto-healed. // /api/pulse does NOT build this and must not: it observes, it never constructs. // // Builds the live payload from the in-memory registry, the scheduler's slots // and on-disk channel stats. Shared by the /jobs page (initial state), the // dashboard, the widget and the /api/jobs/active poll route. export async function buildActiveJobsPayload( i: ActiveJobsInputs, ): Promise { const { registry, scheduler, pool } = i; const now = i.now(); // Snapshot the scheduler ONCE — each queues() call returns fresh arrays, so // completing slots below cannot mutate it under us. const views = scheduler.queues(); const facts = reconcileSlots(views, (id) => registry.get(id), now); const factById = new Map(facts.map((f) => [f.id, f])); const all = registry.list(); const holdsSlot = new Set(facts.map((f) => f.id)); // Running and queued records, PLUS any record the scheduler still has a slot // for even though its status went terminal — that row is the drift, and the // page that hid it is the reason this fold exists. const records = all .filter( (j) => j.status === "running" || j.status === "queued" || holdsSlot.has(j.id), ) .sort((a, b) => { if (a.status === b.status) return 0; return a.status === "running" ? -1 : 1; }); const recordRows = await rowsForRecords(i, records, views, now); for (const row of recordRows) { const f = factById.get(row.id); if (!f) continue; row.heldMs = f.heldMs; row.stuck = f.stuck; row.pid = registry.get(row.id)?.child?.pid; } // The phantoms: a slot whose record the registry has forgotten. Its kind // comes from the on-disk meta sidecar. const phantoms = await Promise.all( facts .filter((f) => !registry.get(f.id)) .map(async (f) => { const meta = await i.jobMeta(f.id).catch(() => null); return fromSlot(f, meta?.kind ?? "unknown"); }), ); const jobs = orderLiveRows([...recordRows, ...phantoms]); // The last log line is read PER STUCK ROW, not per slot per poll: a // transcription log runs to megabytes and there is one line to show and one // reason to show it. await Promise.all( jobs .filter((r) => r.stuck) .map(async (r) => { r.lastLogLine = await i.tailLog(r.id); }), ); // The ~30 s finished window, built from the same registry read. const live = new Set(jobs.map((r) => r.id)); const recentRecords = all.filter( (j) => !live.has(j.id) && typeof j.endedAt === "number" && j.endedAt >= now - RECENT_MS, ); const recent = (await rowsForRecords(i, recentRecords, views, now)).sort( (a, b) => (b.endedAt ?? 0) - (a.endedAt ?? 0), ); // Auto-heal the hard cases now that they have been captured into the view: a // terminal/evicted record's slot MUST be freed. Idempotent and safe; the // freed slot (and any promoted successor) shows correctly on the next poll. for (const f of facts) if (f.heal) scheduler.complete(f.id); // "observe" — a UI poll, several times a minute. It reports the state the // runners are actually in (latch included, so a pipeline held for the resume // margin reads as stopped rather than green) without being the thing that // moves that latch. const diskStatus = i.disk; const disk: DiskStatusView = { enabled: diskStatus.enabled, freeBytes: diskStatus.freeBytes, thresholdBytes: diskStatus.thresholdBytes, resumeBytes: diskStatus.resumeBytes, low: diskStatus.enabled && !diskStatus.ok, reason: diskStatus.reason, message: diskStatus.message, }; return { jobs, recent, builtAt: now, summary: { queues: views.length, running: views.reduce((n, v) => n + v.running.length, 0), queued: views.reduce((n, v) => n + v.queued.length, 0), stuck: facts.filter((f) => f.stuck).length, workersPaused: pool.isPaused(), workerCount: pool.summary().length, }, disk, lanes: buildLanes(i), }; } // The ids the live payload currently reports stuck. `reapStuckJobsAction` force- // releases each: the hard cases were already healed by the build that listed // them, and forceRelease is a no-op on a freed slot (registry.forceRelease). export async function stuckJobIds(i: ActiveJobsInputs): Promise { const payload = await buildActiveJobsPayload(i); return payload.jobs.filter((r) => r.stuck).map((r) => r.id); } // The four lanes, from state this process already holds — a runner's status is // an in-memory read, so the strip costs nothing on a 1-second poll. // // ONE LOOP FOR FOUR LANES since slice 1.3. Two of these rows used to be SWEEP // rows, derived from an armed flag and a registry lookup for the orchestrator // job, while the two above them were runner rows — the same strip describing // two different mechanisms in one vocabulary. Every lane has a runner now. // // The wording comes from the SAME helpers the console uses. A runner's reason // is its own idleReason, mapped to words here rather than on the client so // /jobs, the widget and the dashboard cannot describe one state three ways. function buildLanes(i: ActiveJobsInputs): ActiveLaneView[] { const settings = i.settings; return LANES.map((kind) => { const status = i.autoRunnerStatus(kind); const inFlight = status.inFlight.length; return { kind: autoRunnerJobKind(kind), label: LANE_LABEL[kind], // A STOPPED runner is "unavailable", not "idle". It will never pick // anything up, and an idle-looking lane reads as "all caught up". state: !status.running ? "unavailable" : deriveLaneState({ gateHeld: isGateHeld(settings, kind), activeCount: inFlight, }), note: status.running ? autoIdleNote(status.idleReason, kind) : "not running", inFlight, }; }); } // The strip's name for each lane. The two runner lanes keep the labels the // strip has always shown; the two operation lanes keep the ones their sweep // rows had, so nothing on the dashboard is renamed by this slice. const LANE_LABEL: Record = { transcription: "Auto-transcribe", download: "Auto-download", digest: "Digest", backfill: "Backfill", }; // The runner idle reasons, in words. A copy of the console's idleReasonText, // kept server-side because this payload is consumed by three clients and the // sentence must be the same in all three. function autoIdleNote( reason: AutoRunnerIdleReason | null, kind: AutoQueueKind, ): string | null { switch (reason) { case "no-pending": return "nothing pending"; case "capped": return "every route to the work is at a worker cap"; case "cooldown": return "every pending platform is in a rate-limit cooldown"; case "deferred": return "every pending video was rate-limited recently and is deferred"; case "held": return "every pending platform is held after repeated rate limits — one probe at a time"; case "paced": return "every pending platform is pausing between downloads"; case "no-workers": return "no enabled worker"; case "workers-paused": return "transcriptions paused globally"; case "disk-gate": return "disk gate closed"; case "downloads-paused": return "downloads paused globally"; case "lane-held": return "lane held"; case "engine-unreachable": return "engine unreachable"; case "snoozed": return "snoozed"; case "disabled": return `the ${LANE_LABEL[kind].toLowerCase()} lane is switched off`; // Exhaustive (no default), so a new idle reason is a compile error here // as it is in the console's idleReasonText. case "stopped": case null: return null; } }