Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit eda647537d08237a39a0e27b9a41db81af240212
parent b0a902e068f3eca8f3241abc55a95de7280f5cd6
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 14 Sep 2026 17:13:45 -0400

views: the active-jobs payload, with its readers injected

buildActiveJobs.ts constructed five singletons, stat-ed the disk, tailed
a log file and read the clock, which is why its two load-bearing rules —
liveJobRows NEVER heals a scheduler slot, buildActiveJobsPayload ALWAYS
does — were prose in a comment and nothing else. The fold moves to
common/views/activeJobs.ts and takes its live state as an argument; both
rules are now assertions against a real createScheduler().

The readers arrive as FUNCTIONS, not as data, and that is forced: which
channels to stat is unknowable until registry.list() has been read, and
which ids need a meta sidecar is unknowable until the scheduler's slots
have been compared against it. So ActiveJobsInputs carries channelStat,
jobMeta and tailLog, plus a disk gate the shell already sampled in
"observe" mode and the four runner statuses. LiveJobRowsInputs carries
only the first — a channel page listing its own jobs pays for no gate,
no sidecar and no log tail.

The 8 KB log tail (node:fs/promises) stays in the editor, in the ~45-line
shell at the old path. Every exported name, signature and value is what
it was, so all nine consumers are untouched; the payload types are
re-exported from the shell until slice 1D repoints the five client files.
listJobRows keeps its shape, with the fold over a listing page split out
as rowsFromJobsPage.

DIVERGENCE from plans/one-core-phase-3.md: the plan names the injected
readers readChannelStat / readJobMeta, but slice 1A's own textual ban
lists `readChannelStat(` and `readJobMeta(` as tokens no views/ file may
contain, and that ban is context-blind by design. Renamed to channelStat
and jobMeta rather than loosening the guard.

Ten tests, none of which needed a registry, a pool or a filesystem
before this: builtAt off the injected clock; the tail read for stuck rows
only; one stat per distinct slug; a slot with no record drawn as a
phantom from its sidecar; a terminal record holding a slot drawn once and
then freed; liveJobRows leaving that same drift alone; a stopped runner
reading unavailable rather than idle.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

Diffstat:
Acommon/views/activeJobs.test.ts | 361+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/activeJobs.ts | 417+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/active/buildActiveJobs.ts | 449++++++++++---------------------------------------------------------------------
3 files changed, 832 insertions(+), 395 deletions(-)

diff --git a/common/views/activeJobs.test.ts b/common/views/activeJobs.test.ts @@ -0,0 +1,361 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { defaultSiteSettings } from "../lib/settings"; +import { createScheduler, type Scheduler } from "../jobs/scheduler"; +import type { JobRecord } from "../jobs/registry"; +import type { JobMeta } from "../jobs/jobMeta"; +import type { Paths } from "../lib/paths"; +import type { DiskGateStatus } from "../lib/diskSpace"; +import type { AutoRunnerStatus } from "../controller/autoRunner"; +import type { AutoQueueKind } from "../lib/autoQueueTypes"; +import type { PoolReader, RegistryReader } from "./inputs"; +import { + buildActiveJobsPayload, + liveJobRows, + rowsFromJobsPage, + type ActiveJobsInputs, +} from "./activeJobs"; + +// Run with: node_modules/.bin/tsx --test common/views/activeJobs.test.ts +// +// This payload used to construct five singletons, stat the disk, tail a log +// file and read the clock, which is why the two rules it lives by — "liveJobRows +// NEVER heals" and "buildActiveJobsPayload ALWAYS heals" — were prose in a +// comment and nothing else. Both are assertions below, against a real +// createScheduler(); everything else is a counted stub, so "how many times did +// this read the disk" is a number a test can state. + +const NOW = 1_700_000_000_000; +// Older than jobRows' STUCK_AGE_MS (10 minutes), which is what makes a live +// running record with no tasks read as possibly-stalled. +const LONG_AGO = NOW - 30 * 60 * 1000; + +function job(over: Partial<JobRecord> = {}): JobRecord { + return { + id: "j1", + kind: "transcribe", + queueKey: "q", + status: "running", + queuedAt: NOW - 1000, + startedAt: NOW - 1000, + logPath: "/dev/null", + ...over, + }; +} + +const OK_DISK: DiskGateStatus = { + ok: true, + enabled: true, + freeBytes: 100, + thresholdBytes: 10, + resumeBytes: 10, + reason: "ok", + message: "", +}; + +function runnerStatus(over: Partial<AutoRunnerStatus> = {}): AutoRunnerStatus { + return { + kind: "transcription", + running: true, + jobId: null, + startedAt: null, + inFlight: [], + activeByNode: {}, + idleReason: null, + ...over, + }; +} + +// A scheduler that records every complete() and otherwise IS the real thing — +// the heal contract is about a side effect on live state, so a fake that only +// counts would pin nothing. +function watched(inner: Scheduler): { + scheduler: Scheduler; + completed: string[]; +} { + const completed: string[] = []; + return { + completed, + scheduler: { + submit: (t) => inner.submit(t), + complete: (id) => { + completed.push(id); + inner.complete(id); + }, + cancel: (id) => inner.cancel(id), + reorder: (id, dir) => inner.reorder(id, dir), + promote: (id) => inner.promote(id), + positionInQueue: (id) => inner.positionInQueue(id), + queues: () => inner.queues(), + }, + }; +} + +// Put an id in a running slot on `queueKey`, first come first served. +function occupy(scheduler: Scheduler, id: string, queueKey = "q"): void { + scheduler.submit({ + id, + queueKey, + tier: "foreground", + concurrency: 1, + start: () => {}, + onCancel: () => {}, + }); +} + +type Harness = { + inputs: ActiveJobsInputs; + completed: string[]; + statCalls: string[]; + metaCalls: string[]; + tailCalls: string[]; +}; + +function harness(over: { + jobs?: JobRecord[]; + slots?: string[]; + meta?: (id: string) => JobMeta | null; + disk?: DiskGateStatus; + runner?: (kind: AutoQueueKind) => AutoRunnerStatus; +} = {}): Harness { + const jobs = over.jobs ?? []; + const { scheduler, completed } = watched(createScheduler()); + for (const id of over.slots ?? []) occupy(scheduler, id); + + const registry: RegistryReader = { + list: () => jobs, + get: (id: string) => jobs.find((j) => j.id === id), + }; + const pool: PoolReader = { + summary: () => [], + isPaused: () => false, + canStopPartial: () => false, + }; + + const statCalls: string[] = []; + const metaCalls: string[] = []; + const tailCalls: string[] = []; + + return { + completed, + statCalls, + metaCalls, + tailCalls, + inputs: { + // The view never touches paths; it is here because LiveInputs carries it. + paths: {} as Paths, + settings: defaultSiteSettings(), + registry, + scheduler, + pool, + now: () => NOW, + channelStat: async (slug) => { + statCalls.push(slug); + return null; + }, + jobMeta: async (id) => { + metaCalls.push(id); + return over.meta?.(id) ?? null; + }, + tailLog: async (id) => { + tailCalls.push(id); + return "the last line"; + }, + disk: over.disk ?? OK_DISK, + autoRunnerStatus: over.runner ?? (() => runnerStatus()), + }, + }; +} + +test("builtAt is the injected clock, not the wall clock", async () => { + const h = harness(); + const payload = await buildActiveJobsPayload(h.inputs); + assert.equal(payload.builtAt, NOW); +}); + +test("the log tail is read for stuck rows only", async () => { + // Two running records holding slots. The first has been idle since LONG_AGO + // with no tasks — possibly-stalled, the soft case — and the second started a + // second ago with work in flight. + const h = harness({ + jobs: [ + job({ id: "stale", queuedAt: LONG_AGO, startedAt: LONG_AGO }), + job({ + id: "fresh", + queueKey: "q2", + tasks: [ + { + id: "v1", + label: "v1", + kind: "transcribe", + startedAt: NOW - 500, + }, + ], + }), + ], + slots: ["stale"], + }); + occupy(h.inputs.scheduler, "fresh", "q2"); + + const payload = await buildActiveJobsPayload(h.inputs); + assert.deepEqual(h.tailCalls, ["stale"]); + assert.equal( + payload.jobs.find((r) => r.id === "stale")?.lastLogLine, + "the last line", + ); + assert.equal(payload.jobs.find((r) => r.id === "fresh")?.lastLogLine, undefined); + // Soft case: surfaced, never healed. + assert.deepEqual(h.completed, []); + assert.equal(payload.summary.stuck, 1); +}); + +test("a channel is stat-ed once however many of its jobs are running", async () => { + const h = harness({ + jobs: [ + job({ id: "a1", channelSlug: "alpha" }), + job({ id: "a2", channelSlug: "alpha", queueKey: "q2" }), + job({ id: "b1", channelSlug: "beta", queueKey: "q3" }), + // No channel at all: a corpus-wide job contributes no slug. + job({ id: "c1", queueKey: "q4" }), + ], + }); + await buildActiveJobsPayload(h.inputs); + assert.deepEqual(h.statCalls, ["alpha", "beta"]); +}); + +test("a slot whose record is gone becomes a phantom row, named by its sidecar", async () => { + const h = harness({ + slots: ["ghost"], + meta: (id) => + id === "ghost" + ? ({ id, kind: "digest-sweep", queueKey: "q" } as JobMeta) + : null, + }); + const payload = await buildActiveJobsPayload(h.inputs); + + assert.deepEqual(h.metaCalls, ["ghost"]); + const row = payload.jobs.find((r) => r.id === "ghost"); + assert.ok(row, "the phantom is in the payload"); + assert.equal(row.kind, "digest-sweep"); + assert.equal(row.status, "evicted"); + assert.equal(row.source, "slot"); + assert.equal(row.stuck?.reason, "record-evicted"); + // Hard case: the slot is freed, once, after the row was built. + assert.deepEqual(h.completed, ["ghost"]); +}); + +test("a terminal record holding a slot is drawn once, then healed", async () => { + const h = harness({ + jobs: [job({ id: "done1", status: "done", endedAt: NOW - 100 })], + slots: ["done1"], + }); + const payload = await buildActiveJobsPayload(h.inputs); + + const rows = payload.jobs.filter((r) => r.id === "done1"); + assert.equal(rows.length, 1, "the drift is drawn exactly once"); + assert.equal(rows[0].stuck?.reason, "terminal-but-holding-slot"); + // It is not ALSO in `recent`: a row in the live head is excluded there. + assert.deepEqual(payload.recent.map((r) => r.id), []); + assert.deepEqual(h.completed, ["done1"]); + // And the slot really is free now — the scheduler is a real one. + assert.deepEqual( + h.inputs.scheduler.queues().flatMap((v) => v.running), + [], + ); +}); + +test("liveJobRows never heals, whatever it sees", async () => { + // The same drift the payload builder would have freed. A channel page reading + // its own jobs must not free a scheduler slot as a side effect. + const h = harness({ + jobs: [ + job({ id: "done1", status: "done", endedAt: NOW - 100 }), + job({ id: "run1", channelSlug: "alpha" }), + ], + slots: ["done1"], + }); + const rows = await liveJobRows(h.inputs, () => true); + + assert.deepEqual(rows.map((r) => r.id), ["run1"]); + assert.deepEqual(h.completed, []); + assert.deepEqual( + h.inputs.scheduler.queues().flatMap((v) => v.running), + ["done1"], + ); + // No lanes, no disk, no meta sidecar and no log tail on this path. + assert.deepEqual(h.metaCalls, []); + assert.deepEqual(h.tailCalls, []); +}); + +test("liveJobRows applies its filter to running and queued records only", async () => { + const h = harness({ + jobs: [ + job({ id: "run1", channelSlug: "alpha" }), + job({ id: "queued1", status: "queued", channelSlug: "alpha" }), + job({ id: "run2", channelSlug: "beta", queueKey: "q2" }), + job({ id: "done1", status: "done", channelSlug: "alpha", endedAt: NOW }), + ], + }); + const rows = await liveJobRows(h.inputs, (j) => j.channelSlug === "alpha"); + assert.deepEqual(new Set(rows.map((r) => r.id)), new Set(["run1", "queued1"])); +}); + +test("a stopped runner reads unavailable, not idle", async () => { + const h = harness({ + runner: (kind) => + kind === "digest" + ? runnerStatus({ kind, running: false, idleReason: "stopped" }) + : runnerStatus({ kind, idleReason: "no-pending" }), + }); + const payload = await buildActiveJobsPayload(h.inputs); + + const byKind = new Map(payload.lanes.map((l) => [l.label, l])); + const digest = byKind.get("Digest"); + assert.ok(digest); + // A stopped runner will never pick anything up; "idle" would read as caught up. + assert.equal(digest.state, "unavailable"); + assert.equal(digest.note, "not running"); + + const transcribe = byKind.get("Auto-transcribe"); + assert.equal(transcribe?.state, "idle"); + assert.equal(transcribe?.note, "nothing pending"); + assert.equal(payload.lanes.length, 4); +}); + +test("the disk block is the gate the shell sampled, with `low` derived", async () => { + const h = harness({ + disk: { ...OK_DISK, ok: false, freeBytes: 1, reason: "below-floor", message: "full" }, + }); + const payload = await buildActiveJobsPayload(h.inputs); + assert.equal(payload.disk.low, true); + assert.equal(payload.disk.reason, "below-floor"); + assert.equal(payload.disk.message, "full"); + + // Disabled gate: never low, whatever `ok` says. + const off = harness({ disk: { ...OK_DISK, enabled: false, ok: false } }); + assert.equal((await buildActiveJobsPayload(off.inputs)).disk.low, false); +}); + +test("rowsFromJobsPage folds a listing page and nothing else", () => { + const out = rowsFromJobsPage({ + entries: [ + { + id: "j9", + kind: "download", + queueKey: "q", + status: "done", + queuedAt: NOW - 10, + endedAt: NOW, + inRegistry: false, + replayable: false, + logPath: "/dev/null", + logSize: 0, + }, + ], + hasMore: true, + total: 42, + }); + assert.equal(out.total, 42); + assert.equal(out.hasMore, true); + assert.deepEqual(out.rows.map((r) => r.id), ["j9"]); +}); diff --git a/common/views/activeJobs.ts b/common/views/activeJobs.ts @@ -0,0 +1,417 @@ +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 { 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<ChannelStat | null>; +}; + +// 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<JobMeta | null>; + tailLog: (id: string) => Promise<string | undefined>; + 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<JobRowView[]> { + 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<string, ChannelStat>(); + 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<string, string>(); + 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<string, number>(); + const queueMaxPos = new Map<string, number>(); + 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<JobRowView[]> { + 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<ActiveJobsPayload> { + 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<string[]> { + 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<AutoQueueKind, string> = { + 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: string | 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 "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`; + default: + return null; + } +} diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -1,124 +1,38 @@ -import { - getRegistry, - type JobRecord, -} from "yt-dlp-transcript-common/jobs/registry"; -import { getScheduler, type QueueView } from "yt-dlp-transcript-common/jobs/scheduler"; -import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; -import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; -import { - readChannelStat, - type ChannelStat, -} from "yt-dlp-transcript-common/controller/channels"; -import { - listAllJobs, - type JobsPage, -} from "yt-dlp-transcript-common/jobs/listJobs"; -import type { Paths } from "yt-dlp-transcript-common/lib/paths"; -import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import path from "node:path"; import { open } from "node:fs/promises"; import type { FileHandle } from "node:fs/promises"; -import { getSettings } from "yt-dlp-transcript-common/lib/settings"; -import { isGateHeld } from "yt-dlp-transcript-common/lib/pauseGates"; -import { - diskGate, - type DiskGateReason, -} from "yt-dlp-transcript-common/lib/diskSpace"; +import type { JobRecord } from "yt-dlp-transcript-common/jobs/registry"; +import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; +import { readChannelStat } from "yt-dlp-transcript-common/controller/channels"; +import { listAllJobs } from "yt-dlp-transcript-common/jobs/listJobs"; +import type { Paths } from "yt-dlp-transcript-common/lib/paths"; +import { diskGate } from "yt-dlp-transcript-common/lib/diskSpace"; +import { getAutoRunnerStatus } from "yt-dlp-transcript-common/controller/autoRunner"; import type { JobRowView } from "yt-dlp-transcript-common/views/jobRowView"; import { - fromEntry, - fromRecord, - fromSlot, - orderLiveRows, - reconcileSlots, - RECENT_MS, -} from "yt-dlp-transcript-common/views/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 "yt-dlp-transcript-common/views/laneState"; -import { - autoRunnerJobKind, - getAutoRunnerStatus, -} from "yt-dlp-transcript-common/controller/autoRunner"; -import { - LANES, - type AutoQueueKind, -} from "yt-dlp-transcript-common/lib/autoQueueTypes"; - -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. + buildActiveJobsPayload as build, + liveJobRows as rows, + rowsFromJobsPage, + stuckJobIds as stuck, + type ActiveJobsInputs, + type ActiveJobsPayload, + type LiveJobRowsInputs, +} from "yt-dlp-transcript-common/views/activeJobs"; +import { liveInputs } from "../../lib/liveInputs"; + +// THE SHELL. The payload is `common/views/activeJobs.ts`, a pure function of +// its arguments; everything here is the reading that view refuses to do — the +// log tail, the channel stats, the meta sidecars, the disk gate and the runner +// statuses. The exported names, their signatures and their values are what +// they were, so all nine consumers are untouched. // -// 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[]; - // Server Date.now() 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 type re-export keeps the five client files compiling until slice 1D +// repoints them at the view. +export type { + ActiveJobsPayload, + ActiveLaneView, + DiskStatusView, +} from "yt-dlp-transcript-common/views/activeJobs"; // How much of the log tail to read for the one line we want. A transcription // job's log runs to megabytes and this used to read all of it, per slot, per @@ -150,8 +64,8 @@ async function readLastLogLine( const lines = buf.toString("utf8").split("\n"); // Mid-file window: drop the leading fragment, which may be half a line. const from = start > 0 ? 1 : 0; - for (let i = lines.length - 1; i >= from; i--) { - const line = lines[i].trim(); + for (let n = lines.length - 1; n >= from; n--) { + const line = lines[n].trim(); if (line) return line; } } catch { @@ -162,300 +76,45 @@ async function readLastLogLine( return undefined; } -// 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( - records: JobRecord[], - views: QueueView[], - now: number, -): Promise<JobRowView[]> { - const channelSlugs = Array.from( - new Set( - records.map((j) => j.channelSlug).filter((s): s is string => Boolean(s)), - ), - ); - - const paths = getPaths(); - const statResults = await Promise.all( - channelSlugs.map((slug) => readChannelStat(paths, slug)), - ); - const channelStats = new Map<string, ChannelStat>(); - for (let i = 0; i < channelSlugs.length; i++) { - const stat = statResults[i]; - if (stat) channelStats.set(channelSlugs[i], stat); - } - - // Map worker id → display name so each transcribe task can show which worker - // it's running on. - const workerNames = new Map<string, string>(); - for (const w of getWorkerPool().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<string, number>(); - const queueMaxPos = new Map<string, number>(); - for (const v of views) { - const ids = [...v.running, ...v.queued]; - for (const [i, id] of ids.entries()) positions.set(id, i); - queueMaxPos.set(v.name, ids.length - 1); - } +// The rows half of the inputs: no disk gate, no log tail, no runner statuses. +// A page listing its own jobs must not pay for — or perform — any of those. +function rowInputs(): LiveJobRowsInputs { + const i = liveInputs(); + return { ...i, channelStat: (slug) => readChannelStat(i.paths, slug) }; +} - 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 full payload's inputs. The gate is sampled in "observe" mode — 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. +async function payloadInputs(): Promise<ActiveJobsInputs> { + const i = rowInputs(); + return { + ...i, + jobMeta: (id) => readJobMeta(i.paths, id), + tailLog: (id) => readLastLogLine(i.paths.jobsDir, id), + disk: await diskGate(i.paths, i.settings, { mode: "observe" }), + autoRunnerStatus: getAutoRunnerStatus, + }; } -// 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( filter: (j: JobRecord) => boolean, ): Promise<JobRowView[]> { - const registry = getRegistry(); - const records = registry - .list() - .filter((j) => j.status === "running" || j.status === "queued") - .filter(filter); - const rows = await rowsForRecords(records, getScheduler().queues(), Date.now()); - return orderLiveRows(rows); + return rows(rowInputs(), filter); } -// 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. export async function listJobRows( paths: Paths, opts: { limit?: number; before?: number } = {}, ): Promise<{ rows: JobRowView[]; hasMore: boolean; total: number }> { - const page: JobsPage = await listAllJobs(paths, opts); - return { - rows: page.entries.map(fromEntry), - hasMore: page.hasMore, - total: page.total, - }; + return rowsFromJobsPage(await listAllJobs(paths, opts)); } -// 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(): Promise<ActiveJobsPayload> { - const registry = getRegistry(); - const scheduler = getScheduler(); - const pool = getWorkerPool(); - const paths = getPaths(); - const now = Date.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(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 readJobMeta(paths, 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 readLastLogLine(paths.jobsDir, 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(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 = await diskGate(paths, getSettings(), { mode: "observe" }); - 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(), - }; + return build(await payloadInputs()); } -// 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(): Promise<string[]> { - const payload = await buildActiveJobsPayload(); - return payload.jobs.filter((r) => r.stuck).map((r) => r.id); -} - -// The four lanes, from state this process already holds — getAutoRunnerStatus -// 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(): ActiveLaneView[] { - const settings = getSettings(); - return LANES.map((kind) => { - const status = getAutoRunnerStatus(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<AutoQueueKind, string> = { - 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: string | 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 "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`; - default: - return null; - } + return stuck(await payloadInputs()); }