Archilyzer · Source

archilyzer

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

commit a39afd0001f7f59e691ac1143c72bddfe452351e
parent e67b93cc417f615789bc6d0ed3140cfafef696ee
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sun, 30 Aug 2026 11:26:55 -0400

jobs: the scheduler's view rides the live payload, and heals its drift

The queue-diagnostics page was one screen doing two things: reporting the
scheduler's slots against the registry's statuses, and — as a side effect of
being looked at — freeing a slot whose record had gone terminal or been evicted.
It was also the only place that did either, so a wedged queue stayed wedged
until someone thought to open a page most operators never knew about.

The reconciliation is a pure function over the queues snapshot now, and it
rides the live payload. Every surface that draws that payload — /jobs, the
dashboard, the monitor widget, /api/jobs/active — frees a running slot whose
record is finished or gone, after drawing it once so the operator sees what was
fixed. The soft "possibly stalled" case is surfaced and never healed: it may be
a genuinely long operation, and that is the operator's call.

A stuck slot is a row like any other: its reason as a badge, its held-for time,
the last line of its log, Force-release on the row and Reap stuck on the health
line above the table (active queues · running · queued · stuck · workers). A
slot the registry has forgotten gets a row too — the phantom, the one row no
registry record can answer for.

The last log line is now read per stuck row rather than per slot per poll; a
transcription log runs to megabytes and there was one line to show.

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

Diffstat:
Deditor/app/api/jobs/queue/route.ts | 12------------
Meditor/app/api/test/stuck-job/route.ts | 17+++++++++--------
Meditor/app/jobs/actions.ts | 13++++++-------
Meditor/app/jobs/active/buildActiveJobs.ts | 179++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Meditor/app/jobs/components/JobsTable.tsx | 92+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Aeditor/app/jobs/components/ReapStuckButton.tsx | 37+++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/jobRowView.ts | 4++--
Meditor/app/jobs/jobRows.test.ts | 116+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/jobRows.ts | 97++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Deditor/app/jobs/queue/buildQueueView.ts | 248-------------------------------------------------------------------------------
Deditor/app/jobs/queue/components/QueueViewLive.tsx | 237-------------------------------------------------------------------------------
Deditor/app/jobs/queue/page.tsx | 31-------------------------------
Meditor/app/widget/components/MonitorWidget.tsx | 4++--
Meditor/e2e/queue.spec.ts | 52+++++++++++++++++++++++++++-------------------------
14 files changed, 544 insertions(+), 595 deletions(-)

diff --git a/editor/app/api/jobs/queue/route.ts b/editor/app/api/jobs/queue/route.ts @@ -1,12 +0,0 @@ -import { NextResponse } from "next/server"; -import { buildQueueView } from "../../../jobs/queue/buildQueueView"; - -export const dynamic = "force-dynamic"; - -// Backs the ~2s client poll on the Queue diagnostics screen so stuck slots (and -// the auto-heal that clears terminal/evicted ones) stay live without a full RSC -// refresh. -export async function GET() { - const payload = await buildQueueView(); - return NextResponse.json(payload); -} diff --git a/editor/app/api/test/stuck-job/route.ts b/editor/app/api/test/stuck-job/route.ts @@ -10,14 +10,15 @@ import { getPaths } from "yt-dlp-transcript-common/lib/paths"; export const dynamic = "force-dynamic"; -// E2E test harness only. Fabricates the registry/scheduler DRIFT the Queue -// diagnostics page exists to surface: a job holding a running scheduler slot that -// has gone quiet (no in-flight tasks) far longer than STUCK_AGE_MS — the soft -// "possibly-stalled" case. That case is deliberately NOT auto-healed, so it -// persists across polls and the test can prove FORCE-RELEASE (not auto-heal) -// clears it. We reproduce it directly (backdating startedAt) since a genuinely -// wedged child would be racy. Mounted unconditionally, like the other /api/test -// routes — the editor is a localhost admin tool, not deployed. +// E2E test harness only. Fabricates the registry/scheduler DRIFT that /jobs +// surfaces (the live payload's reconciliation): a job holding a running +// scheduler slot that has gone quiet (no in-flight tasks) far longer than +// STUCK_AGE_MS — the soft "possibly-stalled" case. That case is deliberately +// NOT auto-healed, so it persists across polls and the test can prove +// FORCE-RELEASE (not auto-heal) clears it. We reproduce it directly +// (backdating startedAt) since a genuinely wedged child would be racy. +// Mounted unconditionally, like the other /api/test routes — the editor is a +// localhost admin tool, not deployed. export async function GET(request: Request) { const url = new URL(request.url); const queueKey = url.searchParams.get("queue") || "stuck-queue"; diff --git a/editor/app/jobs/actions.ts b/editor/app/jobs/actions.ts @@ -18,7 +18,7 @@ import { stopBackfillSweep, } from "yt-dlp-transcript-common/controller/backfillSweep"; import { runJobSpec } from "./runJobSpec"; -import { buildQueueView } from "./queue/buildQueueView"; +import { stuckJobIds } from "./active/buildActiveJobs"; export async function cancelJobAction(id: string): Promise<{ ok: boolean }> { const ok = getRegistry().cancel(id); @@ -29,27 +29,26 @@ export async function cancelJobAction(id: string): Promise<{ ok: boolean }> { // Escape hatch for a WEDGED slot (see registry.forceRelease): unconditionally // free the scheduler slot for this id — even if its record is already terminal // or evicted — SIGKILLing any still-running child. Unblocks a queue whose head -// is stuck. Revalidates the jobs list and the queue-diagnostics page. +// is stuck (the live payload re-reads the slot on the next poll). export async function forceReleaseJobAction( id: string, ): Promise<{ ok: boolean }> { const ok = getRegistry().forceRelease(id); revalidatePath("/jobs"); - revalidatePath("/jobs/queue"); return { ok }; } -// Reap every stuck slot the queue-diagnostics view found in one pass: -// force-release each id in its stuckIds set. Returns how many were released. +// Reap every stuck slot the live payload reports: force-release each. The hard +// cases were already healed by the build that listed them; forceRelease is a +// no-op on a freed slot (registry.forceRelease). Returns how many were released. export async function reapStuckJobsAction(): Promise<{ count: number }> { const registry = getRegistry(); - const { stuckIds } = await buildQueueView(); + const stuckIds = await stuckJobIds(); let count = 0; for (const id of stuckIds) { if (registry.forceRelease(id)) count++; } revalidatePath("/jobs"); - revalidatePath("/jobs/queue"); return { count }; } diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -3,6 +3,7 @@ import { 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, @@ -14,6 +15,9 @@ import { } 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 { @@ -21,7 +25,14 @@ import { type DiskGateReason, } from "yt-dlp-transcript-common/lib/diskSpace"; import type { JobRowView } from "../jobRowView"; -import { fromEntry, fromRecord, orderLiveRows, RECENT_MS } from "../jobRows"; +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. @@ -98,6 +109,18 @@ export type ActiveJobsPayload = { // 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 @@ -105,6 +128,48 @@ export type ActiveJobsPayload = { lanes: ActiveLaneView[]; }; +// 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 +// poll, to display its final line. +const TAIL_BYTES = 8 * 1024; + +// Cheap best-effort tail: read the END of the log file and return its last +// non-empty line. Swallows every error (missing file, read failure) — this is +// diagnostic sugar, never load-bearing. +// +// Reading a fixed window from the end means a log whose last 8 KB is entirely +// blank yields nothing rather than scanning back further; that is the intended +// trade for a status line, and the first (possibly truncated) line in the +// window is discarded so a partial line is never displayed as a whole one. +async function readLastLogLine( + jobsDir: string, + id: string, +): Promise<string | undefined> { + const file = path.join(jobsDir, `${id}.log`); + let handle: FileHandle | undefined; + try { + handle = await open(file, "r"); + const { size } = await handle.stat(); + const start = Math.max(0, size - TAIL_BYTES); + const length = size - start; + if (length <= 0) return undefined; + const buf = Buffer.alloc(length); + await handle.read(buf, 0, length, start); + 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(); + if (line) return line; + } + } catch { + /* best-effort */ + } finally { + await handle?.close().catch(() => {}); + } + return undefined; +} + // THE ONE BUILDER FOR EVERY JOB LIST. // // Shared by the /jobs page (initial state), the dashboard, the widget and the @@ -140,7 +205,7 @@ async function rowsForRecords( // 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 (buildQueueView made the same point). + // 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>(); @@ -191,38 +256,98 @@ export async function listJobRows( }; } -// Builds the live payload from the in-memory registry plus on-disk channel -// stats. Shared by the /jobs page (initial state), the dashboard, the widget -// and the /api/jobs/active poll route. +// 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 - // nothing below can mutate it under us. + // 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 activeRecords = all - .filter((j) => j.status === "running" || j.status === "queued") + 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 jobs = orderLiveRows(await rowsForRecords(activeRecords, views, now)); + 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) => - typeof j.endedAt === "number" && j.endedAt >= now - RECENT_MS, + !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), ); - const live = new Set(jobs.map((r) => r.id)); - const recent = (await rowsForRecords(recentRecords, views, now)) - .filter((r) => !live.has(r.id)) - .sort((a, b) => (b.endedAt ?? 0) - (a.endedAt ?? 0)); - const paths = getPaths(); + // 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 @@ -238,7 +363,29 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { message: diskStatus.message, }; - return { jobs, recent, builtAt: now, disk, lanes: buildLanes(jobs) }; + 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(jobs), + }; +} + +// 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 diff --git a/editor/app/jobs/components/JobsTable.tsx b/editor/app/jobs/components/JobsTable.tsx @@ -10,9 +10,11 @@ import { DrainJobButton } from "./DrainJobButton"; import { ForceReleaseJobButton } from "./ForceReleaseJobButton"; import { ReorderJobButtons } from "./ReorderJobButtons"; import { RetryJobButton } from "./RetryJobButton"; +import { ReapStuckButton } from "./ReapStuckButton"; import { LaneStrip } from "./LaneStrip"; import { JobProgressBar, TaskProgressBar } from "./JobProgressBars"; import { jobKindLabel } from "../jobKindLabels"; +import { formatDuration } from "yt-dlp-transcript-common/lib/format"; import { clearJobsFilters, defaultJobsFilters, @@ -51,6 +53,28 @@ function queuePlace(position: number): string { return `${position}${suffix} in line`; } +// One figure on the health line. `data-stat` is the hook; the label text is the +// contract (the suite matches "stuck" and "active queues" exactly). +function healthStat( + key: string, + label: string, + value: string, + emphasize?: boolean, +) { + return ( + <div data-stat={key} className="flex flex-col"> + <span + className={`text-lg font-semibold tabular-nums ${ + emphasize ? "text-destructive" : "text-foreground" + }`} + > + {value} + </span> + <span className="text-xs text-muted-foreground">{label}</span> + </div> + ); +} + function statusColor(status: string): string { switch (status) { case "queued": @@ -220,10 +244,41 @@ export function JobsTable({ // renders whether or not anything is running. const strip = <LaneStrip lanes={live.lanes} jobs={byKind} />; + // THE SCHEDULER'S HEALTH, in five figures, where the queue-diagnostics page + // had them. "Reap stuck" rides this line rather than the header because its + // visibility is a live fact from the poll, not something the SSR knew. + const summary = live.summary; + const health = embedded ? null : ( + <> + <div + aria-label="Queue health" + className="flex flex-wrap items-center justify-between gap-4 border border-border rounded-md p-3 bg-card" + > + <div className="flex flex-wrap items-center gap-6"> + {healthStat("queues", "active queues", String(summary.queues))} + {healthStat("running", "running", String(summary.running))} + {healthStat("queued", "queued", String(summary.queued))} + {healthStat("stuck", "stuck", String(summary.stuck), summary.stuck > 0)} + {healthStat( + "workers", + summary.workersPaused ? "workers (paused)" : "workers", + String(summary.workerCount), + summary.workersPaused, + )} + </div> + {summary.stuck > 0 && <ReapStuckButton count={summary.stuck} />} + </div> + <p className="sr-only" aria-live="polite"> + {summary.stuck} stuck {summary.stuck === 1 ? "slot" : "slots"} + </p> + </> + ); + if (!embedded && rows.length === 0 && total === 0) { return ( <div className="flex flex-col gap-3"> {strip} + {health} <p className="text-sm text-muted-foreground border border-dashed border-border rounded p-4"> No jobs have run yet. </p> @@ -234,6 +289,7 @@ export function JobsTable({ return ( <div className="flex flex-col gap-3"> {strip} + {health} {!embedded && ( <div className="flex flex-col gap-2 text-sm"> {kinds.length > 0 && ( @@ -346,6 +402,7 @@ export function JobsTable({ data-job-id={j.id} // Came from the live payload — the SSR seed or the poll. data-live={liveIds.has(j.id) ? "" : undefined} + data-stuck={j.stuck?.reason} className="border-t border-border" > <td className="px-3 py-2 font-mono text-xs align-top"> @@ -419,6 +476,17 @@ export function JobsTable({ > {j.status} </span> + {j.stuck && ( + <span + data-stuck={j.stuck.reason} + title={ + j.pid ? `pid ${j.pid} — kill -9 ${j.pid}` : undefined + } + className="self-start text-[10px] uppercase tracking-wide px-1.5 py-0.5 rounded bg-destructive-soft text-destructive" + > + stuck · {j.stuck.reason} + </span> + )} {j.progress && <JobProgressBar progress={j.progress} />} {j.tasks && j.tasks.length > 0 && ( <ul @@ -444,12 +512,23 @@ export function JobsTable({ className="px-3 py-2 text-xs text-muted-foreground whitespace-nowrap align-top" suppressHydrationWarning > - {fmtDuration(dur)} + {j.startedAt === undefined && typeof j.heldMs === "number" + ? `held for ${formatDuration(Math.round(j.heldMs / 1000)) || "0:00"}` + : fmtDuration(dur)} </td> <td className="px-3 py-2 text-xs text-muted-foreground text-right whitespace-nowrap align-top"> - {typeof j.logSize === "number" - ? `${j.logSize.toLocaleString()} B` - : "—"} + {j.stuck && j.lastLogLine ? ( + <span + className="font-mono text-xs text-muted-foreground truncate max-w-64 block" + title={j.lastLogLine} + > + {j.lastLogLine} + </span> + ) : typeof j.logSize === "number" ? ( + `${j.logSize.toLocaleString()} B` + ) : ( + "—" + )} </td> <td className="px-3 py-2 text-right align-top"> <div className="flex items-center justify-end gap-2"> @@ -469,10 +548,11 @@ export function JobsTable({ {(j.status === "running" || j.status === "queued") && ( <CancelJobButton jobId={j.id} /> )} - {j.status === "running" && ( + {(j.status === "running" || Boolean(j.stuck)) && ( // Escape hatch: force-free the scheduler slot for any // running job, for when a Cancel leaves the queue wedged - // (see registry.forceRelease). + // (see registry.forceRelease). A stuck terminal/evicted + // row has no Cancel to offer, so it gets this instead. <ForceReleaseJobButton jobId={j.id} /> )} </div> diff --git a/editor/app/jobs/components/ReapStuckButton.tsx b/editor/app/jobs/components/ReapStuckButton.tsx @@ -0,0 +1,37 @@ +"use client"; + +import { useState } from "react"; +import { reapStuckJobsAction } from "../actions"; + +// The bulk escape hatch, on the health line beside the stuck count — where the +// queue-diagnostics strip had it. It is not a header action because its +// visibility is a LIVE fact (summary.stuck > 0 from the poll), not something +// the server render knew. +export function ReapStuckButton({ count }: { count: number }) { + const [reaping, setReaping] = useState(false); + return ( + <button + type="button" + disabled={reaping} + onClick={async () => { + if ( + !window.confirm( + `Force-release all ${count} stuck ${ + count === 1 ? "slot" : "slots" + }? Any still-running children are SIGKILLed.`, + ) + ) + return; + setReaping(true); + try { + await reapStuckJobsAction(); + } finally { + setReaping(false); + } + }} + className="px-3 py-2 rounded-md bg-destructive text-destructive-foreground text-sm font-medium hover:opacity-90 disabled:opacity-50" + > + {reaping ? "Reaping…" : "Reap stuck"} + </button> + ); +} diff --git a/editor/app/jobs/jobRowView.ts b/editor/app/jobs/jobRowView.ts @@ -8,8 +8,8 @@ import type { // live JobRecord), the .jobs directory (a .log + .meta.json sidecar the // registry has forgotten), or the scheduler (a slot whose record was evicted — // a phantom, and the one row with no id the registry can answer for). Three -// surfaces used to draw three shapes; every field name RunningJobsListItem had -// is kept, so the widget and the dashboard read what they read before. +// surfaces used to draw three shapes; every field name the card renderer's row +// type had is kept, so the widget and the dashboard read what they read before. // TYPES ONLY: this module is imported from "use client" files, and the builder // that fills it (active/buildActiveJobs.ts) reaches the registry and the // filesystem. diff --git a/editor/app/jobs/jobRows.test.ts b/editor/app/jobs/jobRows.test.ts @@ -2,12 +2,15 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { newJobId, type JobRecord } from "yt-dlp-transcript-common/jobs/registry"; import type { JobListEntry } from "yt-dlp-transcript-common/jobs/listJobs"; +import { createScheduler } from "yt-dlp-transcript-common/jobs/scheduler"; import { fromEntry, fromRecord, + fromSlot, isLive, mergeJobRows, orderLiveRows, + reconcileSlots, type FromRecordContext, } from "./jobRows"; @@ -224,3 +227,116 @@ test("isLive: anything non-terminal, plus a stuck or evicted slot", () => { true, ); }); + +// --- The scheduler against the registry ------------------------------------- + +function scheduledQueue(): { + scheduler: ReturnType<typeof createScheduler>; + ids: string[]; +} { + const scheduler = createScheduler(); + const ids = [newJobId(), newJobId(), newJobId()]; + for (const id of ids) { + scheduler.submit({ + id, + queueKey: "qFacts", + tier: "foreground", + concurrency: 1, + start: () => {}, + onCancel: () => {}, + }); + } + return { scheduler, ids }; +} + +test("reconcileSlots: a terminal record holding a running slot is hard-stuck", () => { + const { scheduler, ids } = scheduledQueue(); + const [a, b, c] = ids; + const now = 10_000_000; + const records = new Map<string, JobRecord>([ + [a, rec({ id: a, queueKey: "qFacts", status: "done", startedAt: now - 5_000 })], + [b, rec({ id: b, queueKey: "qFacts", status: "queued" })], + [c, rec({ id: c, queueKey: "qFacts", status: "queued" })], + ]); + const facts = reconcileSlots(scheduler.queues(), (id) => records.get(id), now); + const byId = new Map(facts.map((f) => [f.id, f])); + assert.equal(byId.get(a)?.stuck?.reason, "terminal-but-holding-slot"); + assert.equal(byId.get(a)?.heal, true); + assert.equal(byId.get(b)?.position, 1); + assert.equal(byId.get(c)?.position, 2); + assert.equal(byId.get(b)?.stuck, undefined); + assert.equal(byId.get(c)?.stuck, undefined); + assert.equal(facts.filter((f) => f.heal).length, 1); + + // Healing the hard cases promotes the next in line. + facts.filter((f) => f.heal).forEach((f) => scheduler.complete(f.id)); + assert.deepEqual(scheduler.queues()[0].running, [b]); +}); + +test("reconcileSlots: an evicted record is hard-stuck; a quiet live one is soft", () => { + const now = 10_000_000; + + const evicted = scheduledQueue(); + const gone = reconcileSlots( + evicted.scheduler.queues(), + () => undefined, + now, + ); + assert.equal(gone[0].stuck?.reason, "record-evicted"); + assert.equal(gone[0].heal, true); + + const stalled = scheduledQueue(); + const idA = stalled.ids[0]; + const soft = reconcileSlots( + stalled.scheduler.queues(), + (id) => + id === idA + ? rec({ + id: idA, + queueKey: "qFacts", + status: "running", + startedAt: now - 11 * 60 * 1000, + tasks: [], + }) + : undefined, + now, + ); + assert.equal(soft[0].stuck?.reason, "possibly-stalled"); + // SOFT: surfaced and force-releasable, never auto-healed. + assert.equal(soft[0].heal, false); + + const busy = scheduledQueue(); + const idB = busy.ids[0]; + const fresh = reconcileSlots( + busy.scheduler.queues(), + (id) => + id === idB + ? rec({ + id: idB, + queueKey: "qFacts", + status: "running", + startedAt: now - 60 * 1000, + tasks: [], + }) + : undefined, + now, + ); + assert.equal(fresh[0].stuck, undefined); + assert.equal(fresh[0].heal, false); +}); + +test("fromSlot builds the phantom, and a stuck row sorts above a running one", () => { + const { scheduler } = scheduledQueue(); + const facts = reconcileSlots(scheduler.queues(), () => undefined, 10_000_000); + const phantom = fromSlot(facts[0], "whisper-all"); + assert.equal(phantom.status, "evicted"); + assert.equal(phantom.source, "slot"); + assert.equal(phantom.stuck?.reason, "record-evicted"); + assert.equal(phantom.queueKey, "qFacts"); + + const ordered = orderLiveRows([ + { id: "run", kind: "", status: "running", queuedAt: 9, source: "registry" }, + phantom, + ]); + assert.equal(ordered[0].id, phantom.id); +}); diff --git a/editor/app/jobs/jobRows.ts b/editor/app/jobs/jobRows.ts @@ -1,8 +1,10 @@ import { isDrainableKind } from "yt-dlp-transcript-common/jobs/jobKinds"; +import { jobIdTime } from "yt-dlp-transcript-common/jobs/ulid"; import type { JobRecord } from "yt-dlp-transcript-common/jobs/registry"; import type { JobListEntry } from "yt-dlp-transcript-common/jobs/listJobs"; import type { ChannelStat } from "yt-dlp-transcript-common/controller/channels"; -import type { JobRowView } from "./jobRowView"; +import type { QueueView } from "yt-dlp-transcript-common/jobs/scheduler"; +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 @@ -18,6 +20,99 @@ import type { JobRowView } from "./jobRowView"; // 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. +// 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 idle = !rec.tasks || rec.tasks.length === 0; + if (idle && ageMs > 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 diff --git a/editor/app/jobs/queue/buildQueueView.ts b/editor/app/jobs/queue/buildQueueView.ts @@ -1,248 +0,0 @@ -import path from "node:path"; -import { open } from "node:fs/promises"; -import type { FileHandle } from "node:fs/promises"; -import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; -import { getScheduler } from "yt-dlp-transcript-common/jobs/scheduler"; -import { getWorkerPool } from "yt-dlp-transcript-common/jobs/workerPool"; -import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; -import { getPaths } from "yt-dlp-transcript-common/lib/paths"; - -// 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. -// 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 QueueSlotView = { - id: string; - kind: string; - // Registry status ("running"/"queued"/terminal), or "evicted" when the - // scheduler references an id the registry no longer knows about. - status: string; - ageMs: number; - stuck: boolean; - stuckReason?: string; - pid?: number; - draining: boolean; - lastLogLine?: string; -}; - -export type QueueGroupView = { - name: string; - head?: QueueSlotView; - backlog: QueueSlotView[]; -}; - -export type QueueViewPayload = { - queues: QueueGroupView[]; - summary: { - queues: number; - running: number; - queued: number; - stuck: number; - workersPaused: boolean; - workerCount: number; - }; - stuckIds: string[]; -}; - -// 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 -// poll, to display its final line. -const TAIL_BYTES = 8 * 1024; - -// Cheap best-effort tail: read the END of the log file and return its last -// non-empty line. Swallows every error (missing file, read failure) — this is -// diagnostic sugar, never load-bearing. -// -// Reading a fixed window from the end means a log whose last 8 KB is entirely -// blank yields nothing rather than scanning back further; that is the intended -// trade for a status line, and the first (possibly truncated) line in the -// window is discarded so a partial line is never displayed as a whole one. -async function readLastLogLine( - jobsDir: string, - id: string, -): Promise<string | undefined> { - const file = path.join(jobsDir, `${id}.log`); - let handle: FileHandle | undefined; - try { - handle = await open(file, "r"); - const { size } = await handle.stat(); - const start = Math.max(0, size - TAIL_BYTES); - const length = size - start; - if (length <= 0) return undefined; - const buf = Buffer.alloc(length); - await handle.read(buf, 0, length, start); - 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(); - if (line) return line; - } - } catch { - /* best-effort */ - } finally { - await handle?.close().catch(() => {}); - } - return undefined; -} - -// Reconciles the scheduler (source of truth for SLOTS) against the registry -// (source of truth for job STATUS) to surface stuck jobs the registry-status -// filter on the Active Jobs page hides. The scheduler is the primary iteration -// source: it still references evicted/terminal ids whose slot was never freed, -// which is exactly the drift we're hunting. -// -// Auto-heal: for any running slot whose record is terminal OR missing, we call -// scheduler.complete(id) inline — a terminal record's slot MUST be freed, so -// this is safe and self-correcting on every poll. The soft "possibly-stalled" -// case (a live running record sitting idle past STUCK_AGE_MS) is surfaced but -// never auto-healed, since it may just be a genuinely long-running operation. -export async function buildQueueView(): Promise<QueueViewPayload> { - 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 won't mutate this snapshot). The scheduler is the - // authoritative slot view: it still references evicted/terminal ids whose slot - // was never freed — exactly the drift we surface. - const views = scheduler.queues(); - - // Terminal/evicted running slots are auto-healed, but only AFTER the view is - // built so they still render as stuck this pass (they vanish next poll). The - // soft "possibly-stalled" case is intentionally NOT collected here. - const toHeal: string[] = []; - - const stuckIds: string[] = []; - let runningCount = 0; - let queuedCount = 0; - - // buildSlot is PURE: it reports whether the slot needs healing rather than - // pushing to toHeal/stuckIds itself. Every slot is built concurrently below, - // so a side effect in here would land in completion order — and the healing - // pass must stay deterministic. - const buildSlot = async ( - id: string, - role: "running" | "queued", - ): Promise<{ view: QueueSlotView; heal: boolean }> => { - const rec = registry.get(id); - const draining = rec?.draining === true; - const lastLogLine = await readLastLogLine(paths.jobsDir, id); - let stuck = false; - let heal = false; - let stuckReason: string | undefined; - - 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 = true; - stuckReason = "record-evicted"; - heal = true; - } else if (TERMINAL.has(rec.status)) { - stuck = true; - stuckReason = "terminal-but-holding-slot"; - heal = true; - } else { - const ageMs = now - (rec.startedAt ?? rec.queuedAt); - const idle = !rec.tasks || rec.tasks.length === 0; - if (idle && ageMs > STUCK_AGE_MS) { - // Soft case: a live running/draining record sitting idle far too long. - // Force-releasable, but deliberately NOT auto-healed above. - stuck = true; - stuckReason = "possibly-stalled"; - } - } - } - - // For an evicted record fall back to the on-disk meta sidecar for a kind. - let kind = rec?.kind; - if (!kind) { - const meta = await readJobMeta(paths, id).catch(() => null); - kind = meta?.kind ?? "unknown"; - } - - const ageMs = now - (rec?.startedAt ?? rec?.queuedAt ?? now); - - return { - view: { - id, - kind, - status: rec ? rec.status : "evicted", - ageMs, - stuck, - stuckReason, - pid: rec?.child?.pid, - draining, - lastLogLine, - }, - heal, - }; - }; - - // Flatten every slot across every queue into one list, build them all at - // once, then reassemble in the original order. This used to be a serial - // `await` per slot, each of which read a whole job log. - type SlotRef = { queue: number; role: "running" | "queued"; id: string }; - const refs: SlotRef[] = []; - for (const [queue, v] of views.entries()) { - for (const id of v.running) refs.push({ queue, role: "running", id }); - for (const id of v.queued) refs.push({ queue, role: "queued", id }); - } - const built = await Promise.all( - refs.map((r) => buildSlot(r.id, r.role)), - ); - - const queues: QueueGroupView[] = []; - let cursor = 0; - for (const v of views) { - // The scheduler keeps concurrency 1 per queueKey, so there's at most one - // running head — but iterate defensively in case that ever changes. - const backlog: QueueSlotView[] = []; - let head: QueueSlotView | undefined; - for (let i = 0; i < v.running.length; i++) { - const { view } = built[cursor++]; - if (i === 0) head = view; - else backlog.push(view); - runningCount++; - } - for (let i = 0; i < v.queued.length; i++) { - backlog.push(built[cursor++].view); - queuedCount++; - } - queues.push({ name: v.name, head, backlog }); - } - - // Fold the per-slot findings in `refs` order, so the healing pass and the - // stuck list are identical run to run regardless of which read finished first. - for (const [i, r] of refs.entries()) { - const { view, heal } = built[i]; - if (view.stuck) stuckIds.push(r.id); - if (heal) toHeal.push(r.id); - } - - // Auto-heal the hard cases now that they've 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 id of toHeal) scheduler.complete(id); - - return { - queues, - summary: { - queues: queues.length, - running: runningCount, - queued: queuedCount, - stuck: stuckIds.length, - workersPaused: pool.isPaused(), - workerCount: pool.summary().length, - }, - stuckIds, - }; -} diff --git a/editor/app/jobs/queue/components/QueueViewLive.tsx b/editor/app/jobs/queue/components/QueueViewLive.tsx @@ -1,237 +0,0 @@ -"use client"; - -import Link from "next/link"; -import { useEffect, useState } from "react"; -import { formatDuration } from "yt-dlp-transcript-common/lib/format"; -import { jobKindLabel } from "../../jobKindLabels"; -import { reapStuckJobsAction } from "../../actions"; -import { ForceReleaseJobButton } from "../../components/ForceReleaseJobButton"; -import type { - QueueSlotView, - QueueViewPayload, -} from "../buildQueueView"; - -// Poll the queue-diagnostics route so stuck slots (and the server-side auto-heal -// that clears terminal/evicted ones) stay live without a full RSC refresh. -// Mirrors ActiveJobsLive; a touch slower since this is an admin/diagnostic view. -const POLL_MS = 2000; - -export function QueueViewLive({ initial }: { initial: QueueViewPayload }) { - const [payload, setPayload] = useState<QueueViewPayload>(initial); - - useEffect(() => { - let cancelled = false; - let timer: ReturnType<typeof setTimeout> | null = null; - async function tick() { - try { - const res = await fetch("/api/jobs/queue", { cache: "no-store" }); - if (res.ok && !cancelled) { - setPayload((await res.json()) as QueueViewPayload); - } - } catch { - // transient — keep polling - } finally { - if (!cancelled) timer = setTimeout(tick, POLL_MS); - } - } - timer = setTimeout(tick, POLL_MS); - return () => { - cancelled = true; - if (timer) clearTimeout(timer); - }; - }, []); - - const { queues, summary } = payload; - - return ( - <div className="flex flex-col gap-4"> - <HealthStrip payload={payload} /> - {queues.length === 0 ? ( - <p className="text-sm text-muted-foreground border border-dashed border-border rounded p-4"> - No active queues. - </p> - ) : ( - <div className="flex flex-col gap-3"> - {queues.map((q) => ( - <QueueCard key={q.name} name={q.name} head={q.head} backlog={q.backlog} /> - ))} - </div> - )} - <p className="sr-only" aria-live="polite"> - {summary.stuck} stuck {summary.stuck === 1 ? "slot" : "slots"} - </p> - </div> - ); -} - -function HealthStrip({ payload }: { payload: QueueViewPayload }) { - const [reaping, setReaping] = useState(false); - const [copied, setCopied] = useState(false); - const { summary } = payload; - - const stat = (label: string, value: string, emphasize?: boolean) => ( - <div className="flex flex-col"> - <span - className={`text-lg font-semibold tabular-nums ${ - emphasize ? "text-destructive" : "text-foreground" - }`} - > - {value} - </span> - <span className="text-xs text-muted-foreground">{label}</span> - </div> - ); - - return ( - <div className="flex flex-wrap items-center justify-between gap-4 border border-border rounded-md p-3 bg-card"> - <div className="flex flex-wrap items-center gap-6"> - {stat("active queues", String(summary.queues))} - {stat("running", String(summary.running))} - {stat("queued", String(summary.queued))} - {stat("stuck", String(summary.stuck), summary.stuck > 0)} - {stat( - summary.workersPaused ? "workers (paused)" : "workers", - String(summary.workerCount), - summary.workersPaused, - )} - </div> - <div className="flex items-center gap-2"> - {summary.stuck > 0 && ( - <button - type="button" - disabled={reaping} - onClick={async () => { - if ( - !window.confirm( - `Force-release all ${summary.stuck} stuck ${ - summary.stuck === 1 ? "slot" : "slots" - }? Any still-running children are SIGKILLed.`, - ) - ) - return; - setReaping(true); - try { - await reapStuckJobsAction(); - } finally { - setReaping(false); - } - }} - className="px-3 py-2 rounded-md bg-destructive text-destructive-foreground text-sm font-medium hover:opacity-90 disabled:opacity-50" - > - {reaping ? "Reaping…" : "Reap all stuck"} - </button> - )} - <button - type="button" - onClick={async () => { - try { - await navigator.clipboard.writeText( - JSON.stringify(payload, null, 2), - ); - setCopied(true); - setTimeout(() => setCopied(false), 1500); - } catch { - // clipboard blocked — no-op - } - }} - className="px-3 py-2 rounded-md border border-border text-sm font-medium hover:bg-muted" - > - {copied ? "Copied" : "Copy diagnostics"} - </button> - </div> - </div> - ); -} - -function QueueCard({ - name, - head, - backlog, -}: { - name: string; - head?: QueueSlotView; - backlog: QueueSlotView[]; -}) { - const warn = head?.stuck === true; - return ( - <section - aria-label={`Queue ${name}`} - className={`flex flex-col gap-2 border rounded-md p-3 ${ - warn ? "border-warning bg-warning-soft" : "border-border bg-card" - }`} - > - <div className="flex items-baseline justify-between gap-2"> - <span className="font-mono text-sm font-medium">{name || "—"}</span> - <span className="text-xs text-muted-foreground"> - {head ? "1 running" : "idle"} - {backlog.length > 0 && ` · ${backlog.length} waiting`} - </span> - </div> - {head && <SlotRow slot={head} isHead />} - {backlog.length > 0 && ( - <ul className="flex flex-col gap-1.5 pl-2 border-l-2 border-border"> - {backlog.map((s) => ( - <li key={s.id}> - <SlotRow slot={s} /> - </li> - ))} - </ul> - )} - </section> - ); -} - -function SlotRow({ slot, isHead }: { slot: QueueSlotView; isHead?: boolean }) { - const held = formatDuration(Math.round(slot.ageMs / 1000)) || "0:00"; - return ( - <div className="flex flex-col gap-1 text-sm"> - <div className="flex flex-wrap items-center gap-2"> - <span - className={`text-xs uppercase tracking-wide px-2 py-0.5 rounded ${ - isHead - ? "bg-info-soft text-info" - : "bg-muted text-muted-foreground" - }`} - > - {isHead ? "running" : slot.status} - </span> - <span className="text-xs font-medium" title={slot.kind}> - {jobKindLabel(slot.kind)} - </span> - {slot.draining && ( - <span className="text-[10px] uppercase tracking-wide px-1.5 py-0.5 rounded bg-warning-soft text-warning"> - draining - </span> - )} - <Link - href={`/jobs/${slot.id}`} - className="font-mono text-xs underline hover:text-foreground" - > - Job {slot.id} - </Link> - <span className="text-xs text-muted-foreground">held for {held}</span> - {typeof slot.pid === "number" && ( - <span - className="font-mono text-xs text-muted-foreground" - title={`kill -9 ${slot.pid}`} - > - pid {slot.pid} - </span> - )} - {slot.stuck && ( - <span className="text-[10px] uppercase tracking-wide px-1.5 py-0.5 rounded bg-destructive-soft text-destructive"> - stuck{slot.stuckReason ? ` · ${slot.stuckReason}` : ""} - </span> - )} - <div className="ml-auto"> - <ForceReleaseJobButton jobId={slot.id} /> - </div> - </div> - {slot.lastLogLine && ( - <p className="font-mono text-xs text-muted-foreground truncate"> - {slot.lastLogLine} - </p> - )} - </div> - ); -} diff --git a/editor/app/jobs/queue/page.tsx b/editor/app/jobs/queue/page.tsx @@ -1,31 +0,0 @@ -import type { Metadata } from "next"; -import Link from "next/link"; -import { buildQueueView } from "./buildQueueView"; -import { QueueViewLive } from "./components/QueueViewLive"; - -export const dynamic = "force-dynamic"; - -export const metadata: Metadata = { title: "Queue" }; - -export default async function QueuePage() { - const initial = await buildQueueView(); - return ( - <div className="flex flex-col gap-4"> - <div className="flex items-center justify-between"> - <div className="flex flex-col gap-0.5"> - <h1 className="text-2xl font-semibold">Queue</h1> - <p className="text-sm text-muted-foreground"> - Scheduler slots per queue — spot and force-release stuck jobs. - </p> - </div> - <Link - href="/jobs/active" - className="px-3 py-2 rounded-md border border-border text-sm font-medium hover:bg-muted" - > - Active jobs - </Link> - </div> - <QueueViewLive initial={initial} /> - </div> - ); -} diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx @@ -37,8 +37,8 @@ import { import { WidgetMenu } from "./WidgetMenu"; import { WidgetControls } from "./WidgetControls"; -// Read-only monitor widget. Reuses the existing ~1s poll pattern from -// ActiveJobsLive / WorkersView against the same /api endpoints, but renders a +// Read-only monitor widget. Reuses the existing ~1s poll pattern from the +// /jobs table / WorkersView against the same /api endpoints, but renders a // compact, control-free view meant to fit in a small pinned window or iframe. export function MonitorWidget({ diff --git a/editor/e2e/queue.spec.ts b/editor/e2e/queue.spec.ts @@ -1,8 +1,9 @@ -// Queue diagnostics page: a job that holds a running scheduler slot but has gone -// quiet ("possibly-stalled") is surfaced as STUCK — the Active Jobs page's -// registry-status filter would otherwise hide it — and Force-release frees the -// slot. The stuck slot is fabricated deterministically via the /api/test -// harness (see app/api/test/stuck-job) rather than racing a real wedged child. +// /jobs: a job holding a running scheduler slot that has gone quiet +// ("possibly-stalled") is surfaced as STUCK on its own row — the registry-status +// filter the live list used to apply would otherwise hide it — and Force-release +// frees the slot. The stuck slot is fabricated deterministically via the +// /api/test harness (see app/api/test/stuck-job) rather than racing a real +// wedged child. import { test, expect } from "@playwright/test"; import { resetData } from "./helpers"; @@ -16,43 +17,44 @@ async function fabricateStuckJob(queueKey: string): Promise<string> { return body.id; } -test("a stuck slot shows on /jobs/queue and Force-release clears it", async ({ +test("a stuck slot shows on /jobs and Force-release clears it", async ({ page, }) => { await resetData(null); const id = await fabricateStuckJob("qStuck"); - await page.goto("/jobs/queue"); + await page.goto("/jobs"); - // Health strip renders with the stuck count emphasized. - await expect(page.getByRole("heading", { name: "Queue" })).toBeVisible(); + // Health line renders with the stuck count emphasized. + await expect(page.getByRole("heading", { name: "Jobs", level: 1 })).toBeVisible(); await expect(page.getByText("stuck", { exact: true })).toBeVisible(); - // The fabricated slot appears in its queue card, flagged stuck. - const card = page.locator("section[aria-label='Queue qStuck']"); - await expect(card).toBeVisible(); - await expect(card.getByText(`Job ${id}`)).toBeVisible(); - await expect(card.getByText(/possibly-stalled/)).toBeVisible(); + // The fabricated slot has its own row, flagged stuck. + const row = page.locator(`tr[data-job-id="${id}"]`); + await expect(row).toBeVisible(); + await expect(row.getByText(/possibly-stalled/)).toBeVisible(); + await expect(page.locator('[data-stat="stuck"]')).toContainText("1"); - // Force-release it (accept the confirm) — the slot is freed and the card goes. + // Force-release it (accept the confirm) — the slot is freed and the badge goes. page.on("dialog", (d) => d.accept()); - await card.getByRole("button", { name: /^Force-release$/ }).click(); + await row.getByRole("button", { name: /^Force-release$/ }).click(); - await expect(page.locator("section[aria-label='Queue qStuck']")).toHaveCount( - 0, - { timeout: 10_000 }, - ); - await expect(page.getByText("No active queues.")).toBeVisible({ + // forceRelease marks the record cancelled and stamps endedAt, so the row stays + // in the live head as a freshly-finished job for another half-minute. + await expect(row.getByText(/possibly-stalled/)).toHaveCount(0, { timeout: 10_000, }); + await expect(row).toContainText("cancelled"); + await expect(page.locator('[data-stat="stuck"]')).toContainText("0"); }); -test("loads with an empty queue and renders the health strip", async ({ +test("loads with an empty queue and renders the health line", async ({ page, }) => { await resetData(null); - await page.goto("/jobs/queue"); - await expect(page.getByRole("heading", { name: "Queue" })).toBeVisible(); + await page.goto("/jobs"); + await expect(page.getByRole("heading", { name: "Jobs", level: 1 })).toBeVisible(); await expect(page.getByText("active queues", { exact: true })).toBeVisible(); - await expect(page.getByText("No active queues.")).toBeVisible(); + await expect(page.locator('[data-stat="queues"]')).toContainText("0"); + await expect(page.getByText("No jobs have run yet.")).toBeVisible(); });