commit 1a25532587488bbc2fee55caaa475a1b21bd70cf
parent 68e97c10e0a913ae591ff0c703044cfa3c4f20a7
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sun, 30 Aug 2026 11:13:39 -0400
jobs: one row type, and one builder for every job list
Three surfaces drew three shapes for the same thing. The directory listing had
JobListEntry, the Active Jobs screen had RunningJobsListItem, the queue
diagnostics had QueueSlotView, and none of them was a superset of the others —
so a job that was running had progress bars on one page and nothing on the next.
A row is one job now, whichever of the three places knows about it, and
jobRowView.ts owns that shape. It is types-only, which also puts the dependency
the right way round: the server builder used to import its row type from a
"use client" component. jobRows.ts is the pure part — the two adapters, the
head's order and the merge that keeps one <tr> per id with the live row winning.
The five pages that hand-rolled a six-field copy of a running job (channel,
video, build, operation) call liveJobRows instead, so their cards show the
progress bars, the drain affordance and the reorder bounds they used to drop.
Queued rows now sort by queue and position rather than by recency: "2nd in
line" only means something if the rows are in line.
The payload gains `recent` (terminal within the last 30 seconds) and `builtAt`,
both for the one list that lands next.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Diffstat:
14 files changed, 735 insertions(+), 319 deletions(-)
diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts
@@ -12,7 +12,7 @@ export type JobStatus =
| "cancelled";
// Extended for the digest sweep. NOTE: this union has one re-spelled copy in
-// RunningJobsList.tsx's RunningJobsListItem — kept as an IMPORT there now, so
+// editor jobRowView.ts's JobRowProgress — kept as an IMPORT there, so
// TypeScript actually flags the next member added here.
export type JobProgressMetric =
| "downloads"
diff --git a/editor/app/build/page.tsx b/editor/app/build/page.tsx
@@ -1,5 +1,6 @@
import type { Metadata } from "next";
import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
+import { liveJobRows } from "../jobs/active/buildActiveJobs";
import { RunningJobsList } from "../jobs/components/RunningJobsList";
import { BuildButtons } from "./components/BuildButtons";
@@ -16,24 +17,12 @@ const BUILD_KINDS = new Set([
"archive-combined-transcripts",
]);
-export default function BuildPage() {
+export default async function BuildPage() {
const registry = getRegistry();
const existingQueues = registry.activeQueueNames();
- const activeJobs = registry
- .list()
- .filter(
- (j) =>
- BUILD_KINDS.has(j.kind) &&
- (j.status === "running" || j.status === "queued"),
- )
- .map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "queued" | "running",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- }));
+ // Through the one builder, so this list has the same progress bars /jobs does
+ // (it used to drop `progress`, `tasks`, `drainable` and the reorder bounds).
+ const activeJobs = await liveJobRows((j) => BUILD_KINDS.has(j.kind));
return (
<div className="flex flex-col gap-4">
<h1 className="text-2xl font-semibold">Build</h1>
diff --git a/editor/app/channels/[slug]/lib/stageStatus.ts b/editor/app/channels/[slug]/lib/stageStatus.ts
@@ -4,7 +4,6 @@ import {
excludedDownloadIdSet,
type ChannelSnapshot,
} from "yt-dlp-transcript-common/controller/channelSnapshot";
-import type { JobRecord } from "yt-dlp-transcript-common/jobs/registry";
import {
backfillLaneEntriesOf,
operationApplies,
@@ -162,7 +161,10 @@ export type ComputeStageStatusesInput = {
failedVideoIds: string[];
failedTranscodingIds: string[];
config: ChannelConfig;
- runningJobs: JobRecord[];
+ // Only `status` and `kind` are read (which stage has work in flight), so this
+ // takes the SHAPE rather than the record: the channel page now gets its rows
+ // from the one job-row builder (liveJobRows) and no longer holds JobRecords.
+ runningJobs: ReadonlyArray<{ status: string; kind: string }>;
// Whether ANY backfill kind is switched on.
//
// The snapshot's per-kind counts are a record of what was true when it was
diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx
@@ -59,6 +59,7 @@ import type {
import { DeleteChannelForm } from "../components/DeleteChannelForm";
import { RenameChannelForm } from "../components/RenameChannelForm";
import { RunningJobsList } from "../../jobs/components/RunningJobsList";
+import { liveJobRows } from "../../jobs/active/buildActiveJobs";
import { CleanupStage } from "./components/stages/CleanupStage";
import { DiagnosticsStage } from "./components/stages/DiagnosticsStage";
import { DownloadStage } from "./components/stages/DownloadStage";
@@ -136,26 +137,13 @@ export default async function ChannelDetailPage({
]);
const availabilityRecords = Object.values(postAvailability);
const fetcher = resolveSocialFetcher(config.postFetcher, config.url);
- const socialRunningJobs = getRegistry()
- .list()
- .filter(
- (j) =>
- j.channelSlug === slug &&
- (j.status === "running" || j.status === "queued"),
- );
+ // Through the one builder, so this list has the same progress bars /jobs
+ // does (it used to drop `progress`, `tasks`, `drainable` and the reorder
+ // bounds).
+ const socialRunningJobs = await liveJobRows((j) => j.channelSlug === slug);
return (
<div className="flex flex-col gap-6">
- <RunningJobsList
- jobs={socialRunningJobs.map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "queued" | "running",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- }))}
- hideChannelSlug
- />
+ <RunningJobsList jobs={socialRunningJobs} hideChannelSlug />
<SocialChannelPanel
slug={slug}
state={{
@@ -182,13 +170,9 @@ export default async function ChannelDetailPage({
const registry = getRegistry();
const existingQueues = registry.activeQueueNames();
- const runningJobs = registry
- .list()
- .filter(
- (j) =>
- j.channelSlug === slug &&
- (j.status === "running" || j.status === "queued"),
- );
+ // Through the one builder, so this list has the same progress bars /jobs does
+ // (it used to drop `progress`, `tasks`, `drainable` and the reorder bounds).
+ const runningJobs = await liveJobRows((j) => j.channelSlug === slug);
const platformDefaultQueueKey = config.platform
? platformQueueKey(config.platform)
: queueKeyForUrl(config.url);
@@ -579,17 +563,7 @@ export default async function ChannelDetailPage({
lastSyncedAt={config.lastSyncedAt}
/>
- <RunningJobsList
- jobs={runningJobs.map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "queued" | "running",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- }))}
- hideChannelSlug
- />
+ <RunningJobsList jobs={runningJobs} hideChannelSlug />
<div className="flex flex-col gap-4">
<StageSwitcher
diff --git a/editor/app/channels/[slug]/videos/[id]/page.tsx b/editor/app/channels/[slug]/videos/[id]/page.tsx
@@ -27,6 +27,7 @@ import {
} from "yt-dlp-transcript-common/lib/platform";
import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
import { getSettings } from "yt-dlp-transcript-common/lib/settings";
+import { liveJobRows } from "../../../../jobs/active/buildActiveJobs";
import { RunningJobsList } from "../../../../jobs/components/RunningJobsList";
import { VideoPanel, type VideoFile } from "./components/VideoPanel";
import { OperationPanel } from "./components/OperationPanel";
@@ -154,22 +155,11 @@ export default async function VideoDetailPage({
const registry = getRegistry();
const existingQueues = registry.activeQueueNames();
- const activeJobs = registry
- .list()
- .filter(
- (j) =>
- j.channelSlug === slug &&
- j.videoId === id &&
- (j.status === "running" || j.status === "queued"),
- )
- .map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "queued" | "running",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- }));
+ // Through the one builder, so this list has the same progress bars /jobs does
+ // (it used to drop `progress`, `tasks`, `drainable` and the reorder bounds).
+ const activeJobs = await liveJobRows(
+ (j) => j.channelSlug === slug && j.videoId === id,
+ );
const defaultQueueKey = config.platform
? platformQueueKey(config.platform)
: queueKeyForUrl(config.url);
diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts
@@ -1,13 +1,18 @@
import {
getRegistry,
- isDrainableKind,
type JobRecord,
} from "yt-dlp-transcript-common/jobs/registry";
+import { getScheduler, type QueueView } from "yt-dlp-transcript-common/jobs/scheduler";
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 { getSettings } from "yt-dlp-transcript-common/lib/settings";
import { isGateHeld } from "yt-dlp-transcript-common/lib/pauseGates";
@@ -15,7 +20,8 @@ import {
diskGate,
type DiskGateReason,
} from "yt-dlp-transcript-common/lib/diskSpace";
-import type { RunningJobsListItem } from "../components/RunningJobsList";
+import type { JobRowView } from "../jobRowView";
+import { fromEntry, fromRecord, orderLiveRows, 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.
@@ -80,8 +86,18 @@ export type ActiveLaneView = {
};
export type ActiveJobsPayload = {
- jobs: RunningJobsListItem[];
- channels: { slug: string; displayName: string }[];
+ // 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;
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
@@ -89,106 +105,20 @@ export type ActiveJobsPayload = {
lanes: ActiveLaneView[];
};
-// Estimate seconds remaining as: remaining tasks × average measured task
-// duration, converted to wall-clock via the effective concurrency observed so
-// far. Returns undefined until at least one task has completed (no average
-// yet) or when there's no work left. `now` is injected so the same snapshot is
-// used across all jobs in a single payload build.
-function computeEtaSeconds(
- job: JobRecord,
- remaining: number,
- now: number,
- remainingAudioSeconds?: number,
-): number | undefined {
- const count = job.completedTaskCount ?? 0;
- const totalMs = job.completedTaskMs ?? 0;
- if (count < 1 || remaining <= 0 || job.startedAt === undefined) {
- return undefined;
- }
- // Prefer an AUDIO-HOUR estimate where the work is proportional to length.
- // Averaging tasks assumes every unit costs about the same, which is true for
- // downloads and wildly false for digests: this corpus is ~77k videos and ~77k
- // audio-hours, and a channel of 9-hour VODs and a channel of 10-minute clips
- // have the same task count and a 50x difference in cost. A task average would
- // therefore quote an ETA that is wrong by more than an order of magnitude at
- // exactly the moment an operator most needs it — the start of an 80-day run.
- const doneAudio = job.completedTaskAudioSeconds ?? 0;
- if (doneAudio > 0 && remainingAudioSeconds && remainingAudioSeconds > 0) {
- const secondsPerAudioSecond = totalMs / 1000 / doneAudio;
- const elapsedMs = Math.max(1, now - job.startedAt);
- const concurrency = Math.max(1, totalMs / elapsedMs);
- return (remainingAudioSeconds * secondsPerAudioSecond) / concurrency;
- }
- const avgProcMs = totalMs / count; // measured average per task
- const elapsedMs = Math.max(1, now - job.startedAt);
- const concurrency = Math.max(1, totalMs / elapsedMs); // effective parallelism
- // Round up to whole parallel waves: the final wave runs its tasks
- // concurrently, so e.g. 2 tasks left across 4 workers still costs ~one full
- // task — not "half a task". Without the ceil the tail is underestimated.
- const waves = Math.ceil(remaining / concurrency);
- return (waves * avgProcMs) / 1000;
-}
-
-function computeJobProgressView(
- job: JobRecord,
- stat: ChannelStat | undefined,
+// 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,
-): RunningJobsListItem["progress"] {
- const snap = job.progress;
- if (!snap || !stat) return undefined;
- // `current` is RE-COUNTED from disk (readChannelStat), never reported by the
- // runner — which is why the digest metric needed its own on-disk counter
- // (digestCount) rather than a number the batch could have just told us.
- // A runner-reported `current` WINS. The disk re-count below cannot see a
- // regeneration — a regenerated digest is rewritten in place, so the file
- // count never moves and the bar sits at 0% for the whole job. Only the runner
- // knows it did the work. Downloads and transcripts report nothing and keep
- // the disk re-count, unchanged.
- const current =
- snap.current ??
- (snap.metric === "downloads"
- ? stat.downloadCount
- : snap.metric === "digests"
- ? (stat.digestCount ?? 0)
- : stat.transcriptCount);
- const range = Math.max(0, snap.target - snap.initial);
- const advance = Math.max(0, current - snap.initial);
- const pct =
- range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100));
- const remaining = Math.max(0, snap.target - current);
- return {
- metric: snap.metric,
- initial: snap.initial,
- current,
- target: snap.target,
- pct,
- etaSeconds: computeEtaSeconds(
- job,
- remaining,
- now,
- snap.remainingAudioSeconds,
- ),
- };
-}
-
-// Builds the Active Jobs screen payload from the in-memory registry plus
-// on-disk channel stats. Shared by the server-rendered page (initial state)
-// and the /api/jobs/active poll route (live updates).
-export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
- const registry = getRegistry();
- const activeRecords = registry
- .list()
- .filter((j) => j.status === "running" || j.status === "queued")
- .sort((a, b) => {
- if (a.status === b.status) return 0;
- return a.status === "running" ? -1 : 1;
- });
-
+): Promise<JobRowView[]> {
const channelSlugs = Array.from(
new Set(
- activeRecords
- .map((j) => j.channelSlug)
- .filter((s): s is string => Boolean(s)),
+ records.map((j) => j.channelSlug).filter((s): s is string => Boolean(s)),
),
);
@@ -207,63 +137,92 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
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. position 0 is
- // the running head; queued jobs occupy 1..max. A parallel ("") job has
- // position -1 and is never reorderable.
+ // 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).
+ // 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 j of activeRecords) {
- const pos = registry.positionInQueue(j.id);
- positions.set(j.id, pos);
- if (pos >= 0) {
- queueMaxPos.set(j.queueKey, Math.max(queueMaxPos.get(j.queueKey) ?? 0, pos));
- }
+ 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);
}
+ 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(
+ 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);
+}
+
+// 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,
+ };
+}
+
+// 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.
+export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
+ const registry = getRegistry();
+ const scheduler = getScheduler();
const now = Date.now();
- const jobs: RunningJobsListItem[] = activeRecords.map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "running" | "queued",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- progress:
- j.status === "running" && j.channelSlug
- ? computeJobProgressView(j, channelStats.get(j.channelSlug), now)
- : undefined,
- tasks: j.tasks?.map((t) => ({
- id: t.id,
- label: t.label,
- kind: t.kind,
- fraction: t.fraction,
- detail: t.detail,
- startedAt: t.startedAt,
- workerId: t.workerId,
- workerName: t.workerId ? workerNames.get(t.workerId) : undefined,
- phase: t.phase,
- probeStartedAt: t.probeStartedAt,
- probeEtaSeconds: t.probeEtaSeconds,
- })),
- draining: j.draining === true,
- drainable:
- j.status === "running" && isDrainableKind(j.kind) && j.draining !== true,
- background: j.background === true,
- // A queued job can move up/promote if it isn't the first queued (position >
- // 1, since the running head is at 0), and down if it isn't the last in its
- // queue. Running and parallel ("") jobs are never reorderable.
- canMoveUp: j.status === "queued" && (positions.get(j.id) ?? -1) > 1,
- canMoveDown:
- j.status === "queued" &&
- (positions.get(j.id) ?? -1) >= 1 &&
- (positions.get(j.id) ?? -1) < (queueMaxPos.get(j.queueKey) ?? 0),
- }));
+ // Snapshot the scheduler ONCE — each queues() call returns fresh arrays, so
+ // nothing below can mutate it under us.
+ const views = scheduler.queues();
- const channels = channelSlugs.map((slug) => ({
- slug,
- displayName: channelStats.get(slug)?.config.name ?? slug,
- }));
+ const all = registry.list();
+ const activeRecords = all
+ .filter((j) => j.status === "running" || j.status === "queued")
+ .sort((a, b) => {
+ if (a.status === b.status) return 0;
+ return a.status === "running" ? -1 : 1;
+ });
+
+ const jobs = orderLiveRows(await rowsForRecords(activeRecords, views, now));
+ // The ~30 s finished window, built from the same registry read.
+ const recentRecords = all.filter(
+ (j) =>
+ typeof j.endedAt === "number" && j.endedAt >= now - RECENT_MS,
+ );
+ 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();
// "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
@@ -279,7 +238,7 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
message: diskStatus.message,
};
- return { jobs, channels, disk, lanes: buildLanes(jobs) };
+ return { jobs, recent, builtAt: now, disk, lanes: buildLanes(jobs) };
}
// The four lanes, from state this process already holds — getAutoRunnerStatus
@@ -288,14 +247,15 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> {
//
// The wording comes from the SAME two helpers the console uses. A runner's
// reason is its own idleReason, mapped to words here rather than on the client
-// so /jobs/active, the widget and the dashboard cannot describe one state three
+// so /jobs, the widget and the dashboard cannot describe one state three
// ways.
-function buildLanes(jobs: RunningJobsListItem[]): ActiveLaneView[] {
+function buildLanes(jobs: JobRowView[]): ActiveLaneView[] {
const settings = getSettings();
const runningOn = (keys: ReadonlyArray<string>): number => {
const set = new Set(keys);
- return jobs.filter((j) => j.status === "running" && set.has(j.queueKey))
- .length;
+ return jobs.filter(
+ (j) => j.status === "running" && j.queueKey !== undefined && set.has(j.queueKey),
+ ).length;
};
const lanes: ActiveLaneView[] = [];
diff --git a/editor/app/jobs/components/ActiveJobsLive.tsx b/editor/app/jobs/components/ActiveJobsLive.tsx
@@ -7,7 +7,8 @@ import type {
ActiveLaneView,
} from "../active/buildActiveJobs";
import { jobKindLabel } from "../jobKindLabels";
-import { RunningJobsList, type RunningJobsListItem } from "./RunningJobsList";
+import { RunningJobsList } from "./RunningJobsList";
+import type { JobRowView } from "../jobRowView";
import { LANE_DOT, LANE_TEXT, LANE_WORD } from "../../components/lanes/laneState";
import { DrainJobButton } from "./DrainJobButton";
import { CancelJobButton } from "./CancelJobButton";
@@ -42,12 +43,12 @@ export function ActiveJobsLive({ initial }: { initial: ActiveJobsPayload }) {
};
}, []);
- const { jobs, channels, lanes } = payload;
+ const { jobs, lanes } = payload;
- const jobsBySlug = new Map<string, RunningJobsListItem[]>();
+ const jobsBySlug = new Map<string, JobRowView[]>();
// Channel-less jobs (e.g. the cross-channel auto-queue runners) are grouped by
// kind so each gets its own labeled section instead of a generic "Other".
- const jobsByKind = new Map<string, RunningJobsListItem[]>();
+ const jobsByKind = new Map<string, JobRowView[]>();
for (const job of jobs) {
if (job.channelSlug) {
const list = jobsBySlug.get(job.channelSlug) ?? [];
@@ -60,6 +61,14 @@ export function ActiveJobsLive({ initial }: { initial: ActiveJobsPayload }) {
}
}
+ // The payload's `channels` field went with the one row type (slice 8c commit
+ // 1) — this screen is the only thing that read it, and it is deleted in the
+ // next commit. Until then the slug stands in for the display name.
+ const channels = [...jobsBySlug.keys()].map((slug) => ({
+ slug,
+ displayName: slug,
+ }));
+
// A LANE IS NOT A JOB, so an empty work list is not an empty page. The strip
// still renders: "every lane is idle and here is why" is the answer this
// screen was previously unable to give at all.
@@ -131,7 +140,7 @@ function LaneStrip({
lanes: ActiveLaneView[];
// kind -> the channel-less jobs of that kind. A lane has at most one (its own
// long-lived runner or sweep job); its controls ride on the lane's line.
- jobs: Map<string, RunningJobsListItem[]>;
+ jobs: Map<string, JobRowView[]>;
}) {
if (lanes.length === 0) return null;
return (
@@ -153,7 +162,7 @@ function LaneRow({
job,
}: {
lane: ActiveLaneView;
- job: RunningJobsListItem | undefined;
+ job: JobRowView | undefined;
}) {
return (
<section
diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx
@@ -7,6 +7,7 @@ import type {
JobProgressMetric,
JobTaskKind,
} from "yt-dlp-transcript-common/jobs/registry";
+import type { JobRowTask, JobRowView } from "../jobRowView";
import { JobLogTail } from "../[id]/components/JobLogTail";
import { jobKindLabel } from "../jobKindLabels";
import { DrainJobButton } from "./DrainJobButton";
@@ -14,66 +15,15 @@ import { CancelJobButton } from "./CancelJobButton";
import { ForceReleaseJobButton } from "./ForceReleaseJobButton";
import { ReorderJobButtons } from "./ReorderJobButtons";
-export type RunningJobsTask = {
- id: string;
- label: string;
- // Imported for the same reason as `metric` below: a re-spelled literal here
- // would not fail the build when JobTaskKind grew a member.
- kind: JobTaskKind;
- fraction?: number;
- detail?: string;
- // Epoch ms when this sub-operation started, for the live "running for" timer.
- startedAt: number;
- // The worker running this transcription (transcribe tasks only).
- workerId?: string;
- workerName?: string;
- // Audio-integrity probe state (download tasks only). When phase is "probing"
- // yt-dlp is paused for an ffmpeg integrity check; the bar fills against the
- // estimated probe duration instead of the (frozen) download fraction.
- phase?: "downloading" | "probing";
- probeStartedAt?: number;
- probeEtaSeconds?: number;
-};
-
-export type RunningJobsListItem = {
- id: string;
- kind: string;
- status: "queued" | "running";
- queueKey: string;
- channelSlug?: string;
- videoId?: string;
- progress?: {
- // Imported, NOT re-spelled: a literal copy here silently drifted from
- // JobProgressMetric and would not fail the build when the union grew.
- metric: JobProgressMetric;
- initial: number;
- current: number;
- target: number;
- pct: number;
- // Estimated seconds remaining (remaining tasks × measured avg per task).
- // Undefined until the first task completes, or when no work remains.
- etaSeconds?: number;
- };
- tasks?: RunningJobsTask[];
- draining?: boolean;
- drainable?: boolean;
- // Background work (e.g. an auto-download unit) queues BEHIND a manual job on
- // the same platform queue. Shown as an "auto" badge, and — when queued — as a
- // hint that a foreground job (a clicked Sync) is being let through first.
- background?: boolean;
- // Reorder affordances for a QUEUED job: whether it can move toward (up/promote)
- // or away from (down) the head of its queue. Both false for running jobs and
- // for parallel ("") jobs that aren't serialized.
- canMoveUp?: boolean;
- canMoveDown?: boolean;
-};
-
type Props = {
- jobs: RunningJobsListItem[];
+ jobs: JobRowView[];
hideChannelSlug?: boolean;
hideVideoId?: boolean;
};
+// THE CARD RENDERER for a page's own jobs (channel, video, build, operation
+// pages). /jobs draws the same JobRowView as a table row; both read the one
+// type, and both get their rows from the one builder (liveJobRows).
export function RunningJobsList({ jobs, hideChannelSlug, hideVideoId }: Props) {
if (jobs.length === 0) return null;
return (
@@ -103,7 +53,7 @@ function JobRow({
hideChannelSlug,
hideVideoId,
}: {
- job: RunningJobsListItem;
+ job: JobRowView;
hideChannelSlug?: boolean;
hideVideoId?: boolean;
}) {
@@ -233,7 +183,7 @@ const METRIC_FILL_BY_TASK: Record<JobTaskKind, string> = {
backfill: "bg-warning",
};
-function TaskProgressBar({ task }: { task: RunningJobsTask }) {
+function TaskProgressBar({ task }: { task: JobRowTask }) {
// How long this task has been running. Null until mounted (see useNow);
// formatDuration returns "" for 0, so the just-started case shows "0:00".
const now = useNow();
@@ -336,7 +286,7 @@ const METRIC_FILL: Record<JobProgressMetric, string> = {
function JobProgressBar({
progress,
}: {
- progress: NonNullable<RunningJobsListItem["progress"]>;
+ progress: NonNullable<JobRowView["progress"]>;
}) {
const label = `${METRIC_LABELS[progress.metric]}: ${progress.current} / ${progress.target}`;
// Append an ETA once the batch has a measured average. formatDuration returns
diff --git a/editor/app/jobs/jobRowView.ts b/editor/app/jobs/jobRowView.ts
@@ -0,0 +1,94 @@
+import type {
+ JobStatus,
+ JobProgressMetric,
+ JobTaskKind,
+} from "yt-dlp-transcript-common/jobs/registry";
+
+// A ROW IS ONE JOB, whichever of three places knows about it: the registry (a
+// 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.
+// 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.
+
+export type JobRowStatus = JobStatus | "archived" | "evicted";
+
+export type StuckReason =
+ | "record-evicted"
+ | "terminal-but-holding-slot"
+ | "possibly-stalled";
+
+export type JobRowTask = {
+ id: string;
+ label: string;
+ // Imported for the same reason as `metric` below: a re-spelled literal here
+ // would not fail the build when JobTaskKind grew a member.
+ kind: JobTaskKind;
+ fraction?: number;
+ detail?: string;
+ // Epoch ms when this sub-operation started, for the live "running for" timer.
+ startedAt: number;
+ // The worker running this transcription (transcribe tasks only).
+ workerId?: string;
+ workerName?: string;
+ // Audio-integrity probe state (download tasks only). When phase is "probing"
+ // yt-dlp is paused for an ffmpeg integrity check; the bar fills against the
+ // estimated probe duration instead of the (frozen) download fraction.
+ phase?: "downloading" | "probing";
+ probeStartedAt?: number;
+ probeEtaSeconds?: number;
+};
+
+export type JobRowProgress = {
+ // Imported, NOT re-spelled: a literal copy here silently drifted from
+ // JobProgressMetric and would not fail the build when the union grew.
+ metric: JobProgressMetric;
+ initial: number;
+ current: number;
+ target: number;
+ pct: number;
+ // Estimated seconds remaining (remaining tasks × measured avg per task).
+ // Undefined until the first task completes, or when no work remains.
+ etaSeconds?: number;
+};
+
+export type JobRowView = {
+ id: string;
+ kind: string; // "" for an old log with no sidecar; the cell renders "—"
+ status: JobRowStatus;
+ queueKey?: string; // undefined = unknown (old log); "" = parallel ("(immediate)")
+ channelSlug?: string;
+ videoId?: string;
+ progress?: JobRowProgress;
+ tasks?: JobRowTask[];
+ draining?: boolean;
+ drainable?: boolean;
+ // Background work (e.g. an auto-download unit) queues BEHIND a manual job on
+ // the same platform queue. Shown as an "auto" badge, and — when queued — as a
+ // hint that a foreground job (a clicked Sync) is being let through first.
+ background?: boolean;
+ // Reorder affordances for a QUEUED job: whether it can move toward (up/promote)
+ // or away from (down) the head of its queue. Both false for running jobs and
+ // for parallel ("") jobs that aren't serialized.
+ canMoveUp?: boolean;
+ canMoveDown?: boolean;
+ // The history side (listJobs.ts).
+ queuedAt?: number;
+ startedAt?: number;
+ endedAt?: number;
+ exitCode?: number;
+ logSize?: number;
+ inRegistry?: boolean;
+ replayable?: boolean;
+ // The slot side (the scheduler's view of the same id).
+ position?: number; // 0 = running head, 1.. = queued; absent for parallel/history
+ heldMs?: number;
+ stuck?: { reason: StuckReason; heldMs: number };
+ pid?: number;
+ lastLogLine?: string; // stuck rows only — see readLastLogLine
+ // Which adapter built it. Never rendered; tests and the merge read it.
+ source: "registry" | "archive" | "slot";
+};
diff --git a/editor/app/jobs/jobRows.test.ts b/editor/app/jobs/jobRows.test.ts
@@ -0,0 +1,226 @@
+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 {
+ fromEntry,
+ fromRecord,
+ isLive,
+ mergeJobRows,
+ orderLiveRows,
+ type FromRecordContext,
+} from "./jobRows";
+
+// Run with: pnpm -C editor exec tsx --test "app/**/*.test.ts"
+//
+// The adapters, the order and the merge — no registry, no scheduler, no disk.
+// Records are built the way common/jobs/registry.test.ts builds them: a literal
+// JobRecord with a fresh id.
+
+function rec(over: Partial<JobRecord> = {}): JobRecord {
+ return {
+ id: newJobId(),
+ kind: "whisper-all",
+ queueKey: "q",
+ status: "running",
+ queuedAt: 1_000,
+ logPath: "",
+ ...over,
+ };
+}
+
+const ctx = (over: Partial<FromRecordContext> = {}): FromRecordContext => ({
+ now: 10_000,
+ workerNames: new Map(),
+ position: -1,
+ queueMax: 0,
+ ...over,
+});
+
+test("fromRecord carries every field the card renderer used to read", () => {
+ const j = rec({
+ channelSlug: "a",
+ videoId: "v1",
+ startedAt: 2_000,
+ background: true,
+ tasks: [
+ {
+ id: "v1",
+ label: "One",
+ kind: "transcribe",
+ startedAt: 3_000,
+ workerId: "w1",
+ },
+ ],
+ });
+ const row = fromRecord(j, ctx({ workerNames: new Map([["w1", "Only"]]) }));
+ assert.equal(row.id, j.id);
+ assert.equal(row.kind, "whisper-all");
+ assert.equal(row.status, "running");
+ assert.equal(row.queueKey, "q");
+ assert.equal(row.channelSlug, "a");
+ assert.equal(row.videoId, "v1");
+ assert.equal(row.background, true);
+ assert.equal(row.draining, false);
+ assert.equal(row.tasks?.[0].workerName, "Only");
+ assert.equal(row.inRegistry, true);
+ assert.equal(row.source, "registry");
+ assert.equal(row.startedAt, 2_000);
+ assert.equal(row.queuedAt, 1_000);
+});
+
+test("fromRecord: drainable needs running, a drainable kind and no drain", () => {
+ assert.equal(fromRecord(rec(), ctx()).drainable, true);
+ assert.equal(fromRecord(rec({ draining: true }), ctx()).drainable, false);
+ assert.equal(fromRecord(rec({ status: "queued" }), ctx()).drainable, false);
+ // check-availability has no jobKinds entry, so it is not drainable.
+ assert.equal(
+ fromRecord(rec({ kind: "check-availability" }), ctx()).drainable,
+ false,
+ );
+});
+
+test("fromRecord: reorder bounds come from position and queueMax", () => {
+ const queued = (position: number, queueMax: number) =>
+ fromRecord(rec({ status: "queued" }), ctx({ position, queueMax }));
+ // Position 1 is the first QUEUED job (0 is the running head): nowhere up.
+ assert.equal(queued(1, 3).canMoveUp, false);
+ assert.equal(queued(2, 3).canMoveUp, true);
+ assert.equal(queued(3, 3).canMoveDown, false);
+ assert.equal(queued(2, 3).canMoveDown, true);
+ // A parallel ("") job is in no queue at all.
+ const parallel = fromRecord(rec({ status: "queued", queueKey: "" }), ctx());
+ assert.equal(parallel.canMoveUp, false);
+ assert.equal(parallel.canMoveDown, false);
+ assert.equal(parallel.position, undefined);
+});
+
+test("fromRecord: no channel stat means no progress bar, and spec means replayable", () => {
+ const j = rec({
+ channelSlug: "a",
+ progress: { metric: "transcripts", initial: 0, target: 10 },
+ });
+ assert.equal(fromRecord(j, ctx()).progress, undefined);
+ const withStat = fromRecord(
+ j,
+ ctx({
+ stat: {
+ config: { name: "A" },
+ transcriptCount: 5,
+ downloadCount: 0,
+ } as unknown as NonNullable<FromRecordContext["stat"]>,
+ }),
+ );
+ assert.equal(withStat.progress?.current, 5);
+ assert.equal(withStat.progress?.pct, 50);
+ assert.equal(fromRecord(rec(), ctx()).replayable, false);
+ assert.equal(
+ fromRecord(rec({ spec: { kind: "whisper-all" } as never }), ctx())
+ .replayable,
+ true,
+ );
+});
+
+test("fromEntry keeps the history fields, archived included", () => {
+ const e: JobListEntry = {
+ id: newJobId(),
+ kind: "sync",
+ channelSlug: "a",
+ status: "archived",
+ queuedAt: 1,
+ endedAt: 9,
+ inRegistry: false,
+ replayable: true,
+ logPath: "/x.log",
+ logSize: 4096,
+ };
+ const row = fromEntry(e);
+ assert.equal(row.status, "archived");
+ assert.equal(row.logSize, 4096);
+ assert.equal(row.inRegistry, false);
+ assert.equal(row.replayable, true);
+ assert.equal(row.source, "archive");
+ // An old log with no sidecar has no kind at all; the cell renders "—".
+ assert.equal(fromEntry({ ...e, kind: undefined }).kind, "");
+});
+
+test("orderLiveRows: running, then queued IN QUEUE ORDER, then what just ended", () => {
+ const rows = orderLiveRows([
+ { id: "q2", kind: "", status: "queued", queueKey: "q", position: 2, source: "registry" },
+ { id: "run", kind: "", status: "running", queuedAt: 5, source: "registry" },
+ { id: "q1", kind: "", status: "queued", queueKey: "q", position: 1, source: "registry" },
+ { id: "done", kind: "", status: "done", endedAt: 99, source: "registry" },
+ ]);
+ assert.deepEqual(
+ rows.map((r) => r.id),
+ ["run", "q1", "q2", "done"],
+ );
+});
+
+test("orderLiveRows: two queues interleave by name, then by position", () => {
+ const rows = orderLiveRows([
+ { id: "b1", kind: "", status: "queued", queueKey: "qB", position: 1, source: "registry" },
+ { id: "a2", kind: "", status: "queued", queueKey: "qA", position: 2, source: "registry" },
+ { id: "a1", kind: "", status: "queued", queueKey: "qA", position: 1, source: "registry" },
+ ]);
+ assert.deepEqual(
+ rows.map((r) => r.id),
+ ["a1", "a2", "b1"],
+ );
+});
+
+test("orderLiveRows: two running jobs keep registry order (newest queued first)", () => {
+ const rows = orderLiveRows([
+ { id: "old", kind: "", status: "running", queuedAt: 1, source: "registry" },
+ { id: "new", kind: "", status: "running", queuedAt: 2, source: "registry" },
+ ]);
+ assert.deepEqual(
+ rows.map((r) => r.id),
+ ["new", "old"],
+ );
+});
+
+test("mergeJobRows: the live row wins by id, and the head stays first", () => {
+ const live = [
+ {
+ id: "A",
+ kind: "sync",
+ status: "running" as const,
+ progress: { metric: "transcripts" as const, initial: 0, current: 1, target: 2, pct: 50 },
+ source: "registry" as const,
+ },
+ { id: "B", kind: "sync", status: "done" as const, endedAt: 5, source: "registry" as const },
+ ];
+ const history = [
+ { id: "A", kind: "sync", status: "running" as const, logSize: 10, source: "archive" as const },
+ { id: "B", kind: "sync", status: "done" as const, logSize: 10, source: "archive" as const },
+ { id: "C", kind: "sync", status: "failed" as const, logSize: 10, source: "archive" as const },
+ ];
+ const rows = mergeJobRows(live, history);
+ assert.deepEqual(
+ rows.map((r) => r.id),
+ ["A", "B", "C"],
+ );
+ // One <tr> per job, and the one that is drawn is the one with progress.
+ assert.equal(rows[0].source, "registry");
+ assert.equal(rows[0].progress?.pct, 50);
+ // A recent row stays ahead of the history page that also lists it.
+ assert.equal(rows[1].source, "registry");
+});
+
+test("isLive: anything non-terminal, plus a stuck or evicted slot", () => {
+ assert.equal(isLive({ id: "1", kind: "", status: "running", source: "registry" }), true);
+ assert.equal(isLive({ id: "1", kind: "", status: "queued", source: "registry" }), true);
+ assert.equal(isLive({ id: "1", kind: "", status: "evicted", source: "slot" }), true);
+ assert.equal(isLive({ id: "1", kind: "", status: "done", source: "archive" }), false);
+ assert.equal(
+ isLive({
+ id: "1",
+ kind: "",
+ status: "done",
+ stuck: { reason: "terminal-but-holding-slot", heldMs: 1 },
+ source: "registry",
+ }),
+ true,
+ );
+});
diff --git a/editor/app/jobs/jobRows.ts b/editor/app/jobs/jobRows.ts
@@ -0,0 +1,236 @@
+import { isDrainableKind } from "yt-dlp-transcript-common/jobs/jobKinds";
+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";
+
+// THE ADAPTERS, THE ORDER AND THE MERGE — pure, and directive-free by
+// construction so both the server builder and the "use client" table can import
+// them. The precedent is components/lanes/laneState.ts: its only value imports
+// are things with no imports of their own, so nothing drags the registry, the
+// scheduler or `fs` into the client bundle. Here that is isDrainableKind
+// (common/jobs/jobKinds.ts, zero imports, already in the client graph via
+// jobKindLabels.ts) and jobIdTime (common/jobs/ulid.ts, zero imports).
+// Everything else is `import type`.
+
+// How long a finished job stays in the LIVE head. Long enough for the 5 s pulse
+// to re-render the tail with the same row twice over; short enough that the head
+// is still "now".
+export const RECENT_MS = 30_000;
+
+// Estimate seconds remaining as: remaining tasks × average measured task
+// duration, converted to wall-clock via the effective concurrency observed so
+// far. Returns undefined until at least one task has completed (no average
+// yet) or when there's no work left. `now` is injected so the same snapshot is
+// used across all jobs in a single payload build.
+function computeEtaSeconds(
+ job: JobRecord,
+ remaining: number,
+ now: number,
+ remainingAudioSeconds?: number,
+): number | undefined {
+ const count = job.completedTaskCount ?? 0;
+ const totalMs = job.completedTaskMs ?? 0;
+ if (count < 1 || remaining <= 0 || job.startedAt === undefined) {
+ return undefined;
+ }
+ // Prefer an AUDIO-HOUR estimate where the work is proportional to length.
+ // Averaging tasks assumes every unit costs about the same, which is true for
+ // downloads and wildly false for digests: this corpus is ~77k videos and ~77k
+ // audio-hours, and a channel of 9-hour VODs and a channel of 10-minute clips
+ // have the same task count and a 50x difference in cost. A task average would
+ // therefore quote an ETA that is wrong by more than an order of magnitude at
+ // exactly the moment an operator most needs it — the start of an 80-day run.
+ const doneAudio = job.completedTaskAudioSeconds ?? 0;
+ if (doneAudio > 0 && remainingAudioSeconds && remainingAudioSeconds > 0) {
+ const secondsPerAudioSecond = totalMs / 1000 / doneAudio;
+ const elapsedMs = Math.max(1, now - job.startedAt);
+ const concurrency = Math.max(1, totalMs / elapsedMs);
+ return (remainingAudioSeconds * secondsPerAudioSecond) / concurrency;
+ }
+ const avgProcMs = totalMs / count; // measured average per task
+ const elapsedMs = Math.max(1, now - job.startedAt);
+ const concurrency = Math.max(1, totalMs / elapsedMs); // effective parallelism
+ // Round up to whole parallel waves: the final wave runs its tasks
+ // concurrently, so e.g. 2 tasks left across 4 workers still costs ~one full
+ // task — not "half a task". Without the ceil the tail is underestimated.
+ const waves = Math.ceil(remaining / concurrency);
+ return (waves * avgProcMs) / 1000;
+}
+
+function computeJobProgressView(
+ job: JobRecord,
+ stat: ChannelStat | undefined,
+ now: number,
+): JobRowView["progress"] {
+ const snap = job.progress;
+ if (!snap || !stat) return undefined;
+ // `current` is RE-COUNTED from disk (readChannelStat), never reported by the
+ // runner — which is why the digest metric needed its own on-disk counter
+ // (digestCount) rather than a number the batch could have just told us.
+ // A runner-reported `current` WINS. The disk re-count below cannot see a
+ // regeneration — a regenerated digest is rewritten in place, so the file
+ // count never moves and the bar sits at 0% for the whole job. Only the runner
+ // knows it did the work. Downloads and transcripts report nothing and keep
+ // the disk re-count, unchanged.
+ const current =
+ snap.current ??
+ (snap.metric === "downloads"
+ ? stat.downloadCount
+ : snap.metric === "digests"
+ ? (stat.digestCount ?? 0)
+ : stat.transcriptCount);
+ const range = Math.max(0, snap.target - snap.initial);
+ const advance = Math.max(0, current - snap.initial);
+ const pct =
+ range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100));
+ const remaining = Math.max(0, snap.target - current);
+ return {
+ metric: snap.metric,
+ initial: snap.initial,
+ current,
+ target: snap.target,
+ pct,
+ etaSeconds: computeEtaSeconds(
+ job,
+ remaining,
+ now,
+ snap.remainingAudioSeconds,
+ ),
+ };
+}
+
+export type FromRecordContext = {
+ now: number;
+ stat?: ChannelStat;
+ workerNames: ReadonlyMap<string, string>;
+ // 0 = running head, 1.. = queued; -1 for a parallel ("") job the scheduler
+ // never tracked.
+ position: number;
+ // The last index in this record's queue, for the reorder bounds.
+ queueMax: number;
+};
+
+// The live adapter: a registry record as a row.
+export function fromRecord(j: JobRecord, ctx: FromRecordContext): JobRowView {
+ return {
+ id: j.id,
+ kind: j.kind,
+ status: j.status,
+ queueKey: j.queueKey,
+ channelSlug: j.channelSlug,
+ videoId: j.videoId,
+ progress:
+ j.status === "running" && j.channelSlug
+ ? computeJobProgressView(j, ctx.stat, ctx.now)
+ : undefined,
+ tasks: j.tasks?.map((t) => ({
+ id: t.id,
+ label: t.label,
+ kind: t.kind,
+ fraction: t.fraction,
+ detail: t.detail,
+ startedAt: t.startedAt,
+ workerId: t.workerId,
+ workerName: t.workerId ? ctx.workerNames.get(t.workerId) : undefined,
+ phase: t.phase,
+ probeStartedAt: t.probeStartedAt,
+ probeEtaSeconds: t.probeEtaSeconds,
+ })),
+ draining: j.draining === true,
+ drainable:
+ j.status === "running" && isDrainableKind(j.kind) && j.draining !== true,
+ background: j.background === true,
+ // A queued job can move up/promote if it isn't the first queued (position >
+ // 1, since the running head is at 0), and down if it isn't the last in its
+ // queue. Running and parallel ("") jobs are never reorderable.
+ canMoveUp: j.status === "queued" && ctx.position > 1,
+ canMoveDown:
+ j.status === "queued" &&
+ ctx.position >= 1 &&
+ ctx.position < ctx.queueMax,
+ queuedAt: j.queuedAt,
+ startedAt: j.startedAt,
+ endedAt: j.endedAt,
+ exitCode: j.exitCode,
+ inRegistry: true,
+ replayable: Boolean(j.spec),
+ position: ctx.position >= 0 ? ctx.position : undefined,
+ source: "registry",
+ };
+}
+
+// The history adapter: a .log + sidecar entry as a row. Field for field —
+// listJobs already resolved the registry record where there was one.
+export function fromEntry(e: JobListEntry): JobRowView {
+ return {
+ id: e.id,
+ kind: e.kind ?? "",
+ status: e.status,
+ queueKey: e.queueKey,
+ channelSlug: e.channelSlug,
+ videoId: e.videoId,
+ queuedAt: e.queuedAt,
+ startedAt: e.startedAt,
+ endedAt: e.endedAt,
+ exitCode: e.exitCode,
+ logSize: e.logSize,
+ inRegistry: e.inRegistry,
+ replayable: e.replayable,
+ source: "archive",
+ };
+}
+
+function liveRank(r: JobRowView): number {
+ if (r.stuck) return 0;
+ if (r.status === "running") return 1;
+ if (r.status === "queued") return 2;
+ return 3;
+}
+
+// The head's order: what is wrong first, then what is working, then what is
+// waiting IN THE ORDER IT WILL RUN, then what just finished.
+//
+// Queued rows sort by (queue name, position) rather than by recency — finding 8
+// of plans/editor-ia-slice-8c.md. That is a change: /jobs/active ordered queued
+// jobs by the registry's queuedAt, which put the job that will run next
+// anywhere in the list. "2nd in line" in the Queue cell only means anything if
+// the rows are in line.
+export function orderLiveRows(rows: JobRowView[]): JobRowView[] {
+ return [...rows].sort((a, b) => {
+ const ra = liveRank(a);
+ const rb = liveRank(b);
+ if (ra !== rb) return ra - rb;
+ if (ra === 0) return (b.heldMs ?? 0) - (a.heldMs ?? 0);
+ // Running: registry order, newest queued first — unchanged.
+ if (ra === 1) return (b.queuedAt ?? 0) - (a.queuedAt ?? 0);
+ if (ra === 2) {
+ const byQueue = (a.queueKey ?? "").localeCompare(b.queueKey ?? "");
+ if (byQueue !== 0) return byQueue;
+ return (a.position ?? 0) - (b.position ?? 0);
+ }
+ return (b.endedAt ?? 0) - (a.endedAt ?? 0);
+ });
+}
+
+// LIVE WINS BY ID: the same job is in the registry and on disk at once, and its
+// live row is the one with progress. The head keeps its order; the tail follows
+// in the order the directory listing gave it, minus anything already above.
+export function mergeJobRows(
+ live: JobRowView[],
+ history: JobRowView[],
+): JobRowView[] {
+ const seen = new Set(live.map((r) => r.id));
+ return [...live, ...history.filter((r) => !seen.has(r.id))];
+}
+
+// Is this row still moving? Drives the 1 s poll: a page with nothing live makes
+// no requests at all.
+export function isLive(r: JobRowView): boolean {
+ return (
+ r.status === "running" ||
+ r.status === "queued" ||
+ r.status === "evicted" ||
+ Boolean(r.stuck)
+ );
+}
diff --git a/editor/app/operations/[id]/page.tsx b/editor/app/operations/[id]/page.tsx
@@ -15,7 +15,7 @@ import {
type DigestAppDescriptor,
} from "yt-dlp-transcript-common/lib/digestApps";
import type { AutoQueueKind } from "yt-dlp-transcript-common/jobs/autoQueueState";
-import { getRegistry } from "yt-dlp-transcript-common/jobs/registry";
+import { liveJobRows } from "../../jobs/active/buildActiveJobs";
import { buildAutoQueueStatusPayload } from "../status";
import { buildSyncRow } from "../syncRow";
import { buildSchedulerStatusPayload } from "../../scheduler/status";
@@ -214,24 +214,12 @@ export default async function OperationPage({
const laneJobKinds =
runnerKind || !laneId ? [] : JOB_KINDS_BY_LANE[laneId];
+ // Through the one builder, so this list has the same progress bars /jobs does
+ // (it used to drop `progress`, `tasks`, `drainable` and the reorder bounds).
const activeJobs =
laneJobKinds.length === 0
? []
- : getRegistry()
- .list()
- .filter(
- (j) =>
- laneJobKinds.includes(j.kind) &&
- (j.status === "running" || j.status === "queued"),
- )
- .map((j) => ({
- id: j.id,
- kind: j.kind,
- status: j.status as "queued" | "running",
- queueKey: j.queueKey,
- channelSlug: j.channelSlug,
- videoId: j.videoId,
- }));
+ : await liveJobRows((j) => laneJobKinds.includes(j.kind));
// THE CHANNELS WITH THIS OPERATION'S WORK. The cast is safe only because an
// id the module does not know returns [] — diarization, both attribution
diff --git a/editor/app/operations/components/OperationDetail.tsx b/editor/app/operations/components/OperationDetail.tsx
@@ -14,10 +14,8 @@ import { railStates } from "./railStates";
import { useHydrated, useOperationsStatus } from "./useOperationsStatus";
import { type Channel } from "./dispatch";
import { StateBand } from "../../components/pipelines/StateBand";
-import {
- RunningJobsList,
- type RunningJobsListItem,
-} from "../../jobs/components/RunningJobsList";
+import { RunningJobsList } from "../../jobs/components/RunningJobsList";
+import type { JobRowView } from "../../jobs/jobRowView";
// ONE OPERATION, IN FULL.
//
@@ -50,7 +48,7 @@ export function OperationDetail({
// Jobs on this operation's queues, SSR-rendered. Empty renders nothing at
// all (RunningJobsList returns null), so a quiet page carries no second
// in-flight list beside the lane's own.
- activeJobs: RunningJobsListItem[];
+ activeJobs: JobRowView[];
// Set for the operations a RUNNER dispatches; null for everything else.
// Resolved on the server FROM THE DESCRIPTOR (`op.runner`) so this component
// never re-decides which is which — and so an operation the registry gains
@@ -167,7 +165,7 @@ function SweepOperationView({
id: string;
laneId: SweepLaneId;
data: AutoQueueStatusPayload;
- activeJobs: RunningJobsListItem[];
+ activeJobs: JobRowView[];
onRefresh: () => Promise<void>;
operationSettings?: ReactNode;
laneSettings?: ReactNode;
diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx
@@ -8,7 +8,7 @@ import type {
ActiveJobsPayload,
DiskStatusView,
} from "../../jobs/active/buildActiveJobs";
-import type { RunningJobsListItem } from "../../jobs/components/RunningJobsList";
+import type { JobRowView } from "../../jobs/jobRowView";
import type {
JobProgressMetric,
JobTaskKind,
@@ -994,7 +994,7 @@ function ActiveJobsStrip({
headingProgress,
eta,
}: {
- jobs: RunningJobsListItem[];
+ jobs: JobRowView[];
showTitle: boolean;
compact: boolean;
jobBar: boolean;
@@ -1037,7 +1037,7 @@ function JobRow({
headingProgress,
eta,
}: {
- job: RunningJobsListItem;
+ job: JobRowView;
compact: boolean;
jobBar: boolean;
headingProgress: boolean;
@@ -1118,7 +1118,7 @@ const METRIC_PREFIX: Record<JobProgressMetric, string> = {
// Shared by the job heading (headingProgress) and the job bar's caption so the
// label/ETA formatting stays in one place.
function jobProgressText(
- progress: NonNullable<RunningJobsListItem["progress"]>,
+ progress: NonNullable<JobRowView["progress"]>,
showEta: boolean,
): string {
const label = `${METRIC_PREFIX[progress.metric]}${progress.current}/${progress.target}`;
@@ -1136,7 +1136,7 @@ function JobProgressBar({
progress,
showEta,
}: {
- progress: NonNullable<RunningJobsListItem["progress"]>;
+ progress: NonNullable<JobRowView["progress"]>;
showEta: boolean;
}) {
const label = `${METRIC_PREFIX[progress.metric]}${progress.current}/${progress.target}`;
@@ -1173,7 +1173,7 @@ function JobProgressBar({
function TaskBar({
task,
}: {
- task: NonNullable<RunningJobsListItem["tasks"]>[number];
+ task: NonNullable<JobRowView["tasks"]>[number];
}) {
const now = useNow();
const probing = task.phase === "probing";