Archilyzer · Source

archilyzer

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

commit cf3bd181fc752aad1df56945e7f403d05159951a
parent 2695aa4949d217d0eb20d21e471b37a31d3ddf95
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Mon, 27 Apr 2026 23:54:35 -0400

feat(editor): named job queues with sequential pipelines

Replaces the binary per-channel mutex with a named-queue concept. Jobs
sharing a queueKey run sequentially in submission order; jobs in different
queues run in parallel. Each pipeline/whisper/build action accepts an
optional queueKey, defaulting to channel:<slug> (preserves today's
per-channel serialization) or "build" for build actions.

The registry exposes enqueue / finalize / cancel / listQueues / activeQueueNames
and tracks queues as FIFO arrays. JobStatus gains "queued"; status flows
queued → running → done|failed|cancelled. Cancelling a queued job
removes it from its queue without disturbing the running head; cancelling
a running job advances the next queued job.

UI: a QueuePicker beside each panel's actions lets the user select an
existing queue or create a new one. StreamActionLog and JobLogTail render
a "Queued in <queue>, position N" banner that disappears when the job
becomes head-of-queue. The jobs list adds a Queue column, a queued status
badge, and a Cancel button on queued rows.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

Diffstat:
Mcommon/components/StreamActionLog.tsx | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/listJobs.ts | 11+++++++++--
Mcommon/jobs/registry.ts | 159++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcommon/jobs/streamCommand.ts | 236+++++++++++++++++++++++++++++++++++++++++++------------------------------------
Aeditor/app/_components/QueuePicker.tsx | 104+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/api/jobs/[id]/log/route.ts | 13+++++++++++--
Aeditor/app/build/_components/BuildButtons.tsx | 63+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/build/buildAction.ts | 14++++++++++----
Meditor/app/build/page.tsx | 47+++++++----------------------------------------
Meditor/app/channels/[slug]/_components/PipelinePanel.tsx | 21+++++++++++++++++----
Meditor/app/channels/[slug]/_components/WhisperPanel.tsx | 17++++++++++++++---
Meditor/app/channels/[slug]/page.tsx | 10++++++++--
Meditor/app/channels/[slug]/pipelineActions.ts | 16++++++++++++----
Meditor/app/channels/[slug]/whisperActions.ts | 12++++++++----
Meditor/app/jobs/[id]/_components/JobLogTail.tsx | 22++++++++++++++++++++--
Meditor/app/jobs/[id]/page.tsx | 13+++++++++++--
Meditor/app/jobs/page.tsx | 24+++++++++++++++++++-----
Meditor/app/page.tsx | 2+-
Aeditor/cypress/e2e/queues.cy.ts | 115+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-a/config.json | 6++++++
Aeditor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-b/config.json | 6++++++
21 files changed, 748 insertions(+), 213 deletions(-)

diff --git a/common/components/StreamActionLog.tsx b/common/components/StreamActionLog.tsx @@ -11,6 +11,12 @@ type Props = { cancelAction?: (id: string) => Promise<{ ok: boolean }>; }; +type StatusPoll = { + status: string; + queueKey?: string; + queuePosition?: number; +}; + export function StreamActionLog({ trigger, buttonLabel, @@ -23,6 +29,7 @@ export function StreamActionLog({ const [error, setError] = useState<string | null>(null); const [jobId, setJobId] = useState<string | null>(null); const [cancelling, setCancelling] = useState(false); + const [poll, setPoll] = useState<StatusPoll | null>(null); const preRef = useRef<HTMLPreElement | null>(null); const stickToBottomRef = useRef(true); @@ -46,10 +53,42 @@ export function StreamActionLog({ if (running) stickToBottomRef.current = true; }, [running]); + // While a job is queued or running, poll the job log endpoint just to + // surface queue/position info (the actual log content streams over the + // ReadableStream). When the job becomes terminal, stop polling. + useEffect(() => { + if (!jobId || !running) return; + let cancelled = false; + let timer: ReturnType<typeof setTimeout> | null = null; + async function tick() { + try { + const res = await fetch( + `/api/jobs/${encodeURIComponent(jobId!)}/log?from=999999999`, + { cache: "no-store" }, + ); + if (!res.ok) return; + const data = (await res.json()) as StatusPoll; + if (cancelled) return; + setPoll(data); + if (data.status === "queued" || data.status === "running") { + timer = setTimeout(tick, 1000); + } + } catch { + if (!cancelled) timer = setTimeout(tick, 2000); + } + } + tick(); + return () => { + cancelled = true; + if (timer) clearTimeout(timer); + }; + }, [jobId, running]); + async function handleClick() { setError(null); setLog(""); setJobId(null); + setPoll(null); stickToBottomRef.current = true; setRunning(true); try { @@ -116,6 +155,17 @@ export function StreamActionLog({ {error} </div> )} + {running && poll?.status === "queued" && poll.queueKey && ( + <div + data-testid={testId ? `${testId}-queue-banner` : undefined} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-zinc-50 dark:bg-zinc-900 px-3 py-2 text-sm text-zinc-700 dark:text-zinc-300" + > + Queued in <code className="font-mono">{poll.queueKey}</code> + {typeof poll.queuePosition === "number" && poll.queuePosition > 0 + ? `, position ${poll.queuePosition}` + : ""} + </div> + )} {(log || running) && ( <pre ref={preRef} diff --git a/common/jobs/listJobs.ts b/common/jobs/listJobs.ts @@ -7,8 +7,10 @@ export type JobListEntry = { id: string; kind?: string; channelSlug?: string; + queueKey?: string; status: JobStatus | "archived"; - startedAt: number; + queuedAt: number; + startedAt?: number; endedAt?: number; exitCode?: number; inRegistry: boolean; @@ -53,7 +55,9 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { id, kind: live_.kind, channelSlug: live_.channelSlug, + queueKey: live_.queueKey, status: live_.status, + queuedAt: live_.queuedAt, startedAt: live_.startedAt, endedAt: live_.endedAt, exitCode: live_.exitCode, @@ -65,6 +69,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { out.push({ id, status: "archived", + queuedAt: mtime, startedAt: mtime, endedAt: mtime, inRegistry: false, @@ -82,7 +87,9 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { id: r.id, kind: r.kind, channelSlug: r.channelSlug, + queueKey: r.queueKey, status: r.status, + queuedAt: r.queuedAt, startedAt: r.startedAt, endedAt: r.endedAt, exitCode: r.exitCode, @@ -93,7 +100,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { } } - return out.sort((a, b) => b.startedAt - a.startedAt); + return out.sort((a, b) => b.queuedAt - a.queuedAt); } export async function readLogChunk( diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -1,14 +1,20 @@ import type { ChildProcess } from "node:child_process"; -export type JobStatus = "running" | "done" | "failed" | "cancelled"; +export type JobStatus = + | "queued" + | "running" + | "done" + | "failed" + | "cancelled"; export type JobRecord = { id: string; kind: string; - mutexKey?: string; + queueKey: string; channelSlug?: string; status: JobStatus; - startedAt: number; + queuedAt: number; + startedAt?: number; endedAt?: number; exitCode?: number; logPath: string; @@ -16,30 +22,28 @@ export type JobRecord = { abortController?: AbortController; }; -class JobRegistry { - private jobs = new Map<string, JobRecord>(); - private mutexes = new Set<string>(); - - acquireMutex(key: string): boolean { - if (this.mutexes.has(key)) return false; - this.mutexes.add(key); - return true; - } +export type QueueSnapshot = { + name: string; + running?: JobRecord; + queued: JobRecord[]; +}; - releaseMutex(key: string): void { - this.mutexes.delete(key); - } +type StartFn = () => void; +type CancelFn = () => void; - hasMutex(key: string): boolean { - return this.mutexes.has(key); - } +class JobRegistry { + private jobs = new Map<string, JobRecord>(); + private queues = new Map<string, JobRecord[]>(); + private starts = new WeakMap<JobRecord, StartFn>(); + private cancels = new WeakMap<JobRecord, CancelFn>(); register(record: JobRecord): void { this.jobs.set(record.id, record); - // Bound the registry. Drop oldest finished records past 100. + // Bound the registry. Drop oldest finished records past 100, but never + // evict records that are queued or running. if (this.jobs.size > 100) { const finished = Array.from(this.jobs.values()) - .filter((j) => j.status !== "running") + .filter((j) => j.status !== "running" && j.status !== "queued") .sort((a, b) => (a.endedAt ?? 0) - (b.endedAt ?? 0)); for (const drop of finished.slice(0, this.jobs.size - 100)) { this.jobs.delete(drop.id); @@ -53,22 +57,117 @@ class JobRegistry { list(): JobRecord[] { return Array.from(this.jobs.values()).sort( - (a, b) => b.startedAt - a.startedAt, + (a, b) => b.queuedAt - a.queuedAt, ); } + // Append a job to its queue. If it's the only entry, start it immediately + // and mark "running". Otherwise leave it as "queued" and remember the + // start/cancel callbacks for when it becomes head-of-queue or is cancelled. + enqueue( + record: JobRecord, + callbacks: { start: StartFn; onCancel: CancelFn }, + ): { willRunNow: boolean; position: number } { + let q = this.queues.get(record.queueKey); + if (!q) { + q = []; + this.queues.set(record.queueKey, q); + } + q.push(record); + this.starts.set(record, callbacks.start); + this.cancels.set(record, callbacks.onCancel); + if (q.length === 1) { + record.status = "running"; + record.startedAt = Date.now(); + callbacks.start(); + return { willRunNow: true, position: 0 }; + } + return { willRunNow: false, position: q.length - 1 }; + } + + // Idempotent: marks the job terminal (if not already), splices it out of + // its queue, and starts the next queued job in that queue. + finalize( + id: string, + status: "done" | "failed" | "cancelled", + exitCode?: number, + ): void { + const job = this.jobs.get(id); + if (!job) return; + if (job.status === "queued" || job.status === "running") { + job.status = status; + job.endedAt = Date.now(); + if (typeof exitCode === "number") job.exitCode = exitCode; + } + const q = this.queues.get(job.queueKey); + if (!q) return; + const idx = q.indexOf(job); + if (idx >= 0) q.splice(idx, 1); + if (q.length === 0) { + this.queues.delete(job.queueKey); + return; + } + const next = q[0]; + if (next.status === "queued") { + next.status = "running"; + next.startedAt = Date.now(); + const startFn = this.starts.get(next); + if (startFn) startFn(); + } + } + cancel(id: string): boolean { const job = this.jobs.get(id); - if (!job || job.status !== "running") return false; - job.status = "cancelled"; - job.abortController?.abort(); - if (job.child) { - job.child.kill("SIGTERM"); - setTimeout(() => { - if (job.child && !job.child.killed) job.child.kill("SIGKILL"); - }, 5_000).unref?.(); + if (!job) return false; + if (job.status === "queued") { + const onCancel = this.cancels.get(job); + if (onCancel) onCancel(); + this.finalize(id, "cancelled"); + return true; + } + if (job.status === "running") { + // Set status synchronously so the streamCommand exit handler sees + // "cancelled" and skips its own status update. SIGTERM the child; + // streamCommand's .finally calls finalize(cancelled) which is + // idempotent and will start the next queued job. + job.status = "cancelled"; + job.abortController?.abort(); + if (job.child) { + job.child.kill("SIGTERM"); + setTimeout(() => { + if (job.child && !job.child.killed) job.child.kill("SIGKILL"); + }, 5_000).unref?.(); + } + return true; + } + return false; + } + + // Snapshot of queues for UI. Sorted by queue name. + listQueues(): QueueSnapshot[] { + const out: QueueSnapshot[] = []; + for (const [name, q] of this.queues.entries()) { + const running = q.find((j) => j.status === "running"); + const queued = q.filter((j) => j.status === "queued"); + out.push({ name, running, queued }); } - return true; + return out.sort((a, b) => a.name.localeCompare(b.name)); + } + + // Names of queues that currently have any non-terminal jobs. Used by the + // UI to populate the QueuePicker dropdown. + activeQueueNames(): string[] { + return Array.from(this.queues.keys()).sort(); + } + + // Position of the job in its queue (0 = currently running). Returns -1 if + // the job is no longer in any queue (terminal). + positionInQueue(id: string): number { + const job = this.jobs.get(id); + if (!job) return -1; + const q = this.queues.get(job.queueKey); + if (!q) return -1; + return q.indexOf(job); } } diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -11,9 +11,8 @@ export type StreamActionResult = type CommonOpts = { kind: string; - mutexKey: string; + queueKey: string; paths: Paths; - busyMessage?: string; channelSlug?: string; }; @@ -37,7 +36,7 @@ async function ensureJobsDir(paths: Paths): Promise<void> { function makeJob( kind: string, - mutexKey: string, + queueKey: string, paths: Paths, channelSlug?: string, ): { id: string; logPath: string; record: JobRecord } { @@ -46,10 +45,10 @@ function makeJob( const record: JobRecord = { id, kind, - mutexKey, + queueKey, channelSlug, - status: "running", - startedAt: Date.now(), + status: "queued", + queuedAt: Date.now(), logPath, }; getRegistry().register(record); @@ -60,143 +59,164 @@ export async function runManagedCommand( opts: RunManagedCommandOpts, ): Promise<StreamActionResult> { const registry = getRegistry(); - if (!registry.acquireMutex(opts.mutexKey)) { - return { - ok: false, - error: opts.busyMessage ?? "A build is already currently running!", - }; - } await ensureJobsDir(opts.paths); const { id, logPath, record } = makeJob( opts.kind, - opts.mutexKey, + opts.queueKey, opts.paths, opts.channelSlug, ); - const child = execa(opts.command, opts.args, { - cwd: opts.cwd, - env: opts.env, - all: true, - buffer: false, - reject: false, - }); - record.child = child; + let controller: ReadableStreamDefaultController<string> | null = null; + let fileStream: WriteStream | null = null; + let cancelledBeforeStart = false; - const fileStream = createWriteStream(logPath); const stream = new ReadableStream<string>({ - start(controller) { - teeChildToControllerAndFile( - child, - fileStream, - controller, - registry, - record, - ); + start(c) { + controller = c; }, cancel() { - // Client disconnected. Job continues to write to its log file so a - // future viewer can re-attach (Phase 8). + // Client disconnected. Job continues to write to its log file. }, }); - return { ok: true, jobId: id, stream }; -} - -function teeChildToControllerAndFile( - child: ReturnType<typeof execa>, - fileStream: WriteStream, - controller: ReadableStreamDefaultController<string>, - registry: ReturnType<typeof getRegistry>, - record: JobRecord, -): void { - child.all?.on("data", (chunk: Buffer) => { - fileStream.write(chunk); - try { - controller.enqueue(chunk.toString("utf8")); - } catch { - // Controller closed; client gone. + const start = () => { + if (cancelledBeforeStart) { + try { + controller?.close(); + } catch {} + return; } - }); - child.all?.on("error", () => {}); - - child - .then((result) => { - if (record.status === "cancelled") return; - record.exitCode = result.exitCode ?? undefined; - record.status = result.exitCode === 0 ? "done" : "failed"; - }) - .catch((err) => { - if (record.status !== "cancelled") { - record.status = "failed"; - try { - controller.enqueue(`\n[error] ${(err as Error).message}\n`); - } catch {} - } - }) - .finally(() => { - record.endedAt = Date.now(); - registry.releaseMutex(record.mutexKey!); - fileStream.end(); + fileStream = createWriteStream(logPath); + const child = execa(opts.command, opts.args, { + cwd: opts.cwd, + env: opts.env, + all: true, + buffer: false, + reject: false, + }); + record.child = child; + + child.all?.on("data", (chunk: Buffer) => { + fileStream!.write(chunk); try { - controller.close(); + controller?.enqueue(chunk.toString("utf8")); } catch {} }); + child.all?.on("error", () => {}); + + child + .then((result) => { + if (record.status === "cancelled") { + registry.finalize(id, "cancelled"); + return; + } + const status = result.exitCode === 0 ? "done" : "failed"; + registry.finalize(id, status, result.exitCode ?? undefined); + }) + .catch((err) => { + if (record.status === "cancelled") { + registry.finalize(id, "cancelled"); + return; + } + try { + controller?.enqueue(`\n[error] ${(err as Error).message}\n`); + } catch {} + registry.finalize(id, "failed"); + }) + .finally(() => { + fileStream?.end(); + try { + controller?.close(); + } catch {} + }); + }; + + const onCancel = () => { + cancelledBeforeStart = true; + try { + controller?.close(); + } catch {} + }; + + registry.enqueue(record, { start, onCancel }); + return { ok: true, jobId: id, stream }; } export async function runManagedFunction( opts: RunManagedFunctionOpts, ): Promise<StreamActionResult> { const registry = getRegistry(); - if (!registry.acquireMutex(opts.mutexKey)) { - return { - ok: false, - error: opts.busyMessage ?? "A build is already currently running!", - }; - } await ensureJobsDir(opts.paths); const { id, logPath, record } = makeJob( opts.kind, - opts.mutexKey, + opts.queueKey, opts.paths, opts.channelSlug, ); - const abort = new AbortController(); - record.abortController = abort; - const fileStream = createWriteStream(logPath); + let controller: ReadableStreamDefaultController<string> | null = null; + let fileStream: WriteStream | null = null; + let cancelledBeforeStart = false; + const stream = new ReadableStream<string>({ - start(controller) { - const onLog = (line: string) => { - const text = line.endsWith("\n") ? line : `${line}\n`; - fileStream.write(text); - try { - controller.enqueue(text); - } catch {} - }; - - opts - .fn(onLog, abort.signal) - .then(() => { - if (record.status !== "cancelled") record.status = "done"; - }) - .catch((err) => { - if (record.status !== "cancelled") { - record.status = "failed"; - onLog(`[error] ${(err as Error).message}`); - } - }) - .finally(() => { - record.endedAt = Date.now(); - registry.releaseMutex(record.mutexKey!); - fileStream.end(); - try { - controller.close(); - } catch {} - }); + start(c) { + controller = c; }, cancel() {}, }); + const start = () => { + if (cancelledBeforeStart) { + try { + controller?.close(); + } catch {} + return; + } + fileStream = createWriteStream(logPath); + const abort = new AbortController(); + record.abortController = abort; + + const onLog = (line: string) => { + const text = line.endsWith("\n") ? line : `${line}\n`; + fileStream!.write(text); + try { + controller?.enqueue(text); + } catch {} + }; + + opts + .fn(onLog, abort.signal) + .then(() => { + if (record.status === "cancelled") { + registry.finalize(id, "cancelled"); + return; + } + registry.finalize(id, "done"); + }) + .catch((err) => { + if (record.status === "cancelled") { + registry.finalize(id, "cancelled"); + return; + } + onLog(`[error] ${(err as Error).message}`); + registry.finalize(id, "failed"); + }) + .finally(() => { + fileStream?.end(); + try { + controller?.close(); + } catch {} + }); + }; + + const onCancel = () => { + cancelledBeforeStart = true; + try { + controller?.close(); + } catch {} + }; + + registry.enqueue(record, { start, onCancel }); return { ok: true, jobId: id, stream }; } diff --git a/editor/app/_components/QueuePicker.tsx b/editor/app/_components/QueuePicker.tsx @@ -0,0 +1,104 @@ +"use client"; + +import { useState } from "react"; + +type Props = { + value: string; + onChange: (next: string) => void; + defaultQueueKey: string; + existingQueues: string[]; + testId?: string; +}; + +const NEW_QUEUE_OPTION = "__new__"; + +export function QueuePicker({ + value, + onChange, + defaultQueueKey, + existingQueues, + testId, +}: Props) { + const [creating, setCreating] = useState(false); + const [draft, setDraft] = useState(""); + + const options = Array.from( + new Set([defaultQueueKey, ...existingQueues, value].filter(Boolean)), + ).sort(); + + if (creating) { + return ( + <div className="flex items-center gap-2"> + <label className="text-xs uppercase tracking-wide text-zinc-500"> + Queue + </label> + <input + type="text" + value={draft} + onChange={(e) => setDraft(e.target.value)} + autoFocus + placeholder="new-queue-name" + data-testid={testId ? `${testId}-input` : undefined} + className="text-sm font-mono px-2 py-1 rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900" + /> + <button + type="button" + onClick={() => { + const trimmed = draft.trim(); + if (trimmed) { + onChange(trimmed); + setCreating(false); + setDraft(""); + } + }} + data-testid={testId ? `${testId}-save` : undefined} + className="text-xs px-2 py-1 rounded bg-zinc-900 dark:bg-zinc-100 text-zinc-100 dark:text-zinc-900" + > + Use + </button> + <button + type="button" + onClick={() => { + setCreating(false); + setDraft(""); + }} + className="text-xs text-zinc-500 underline" + > + cancel + </button> + </div> + ); + } + + return ( + <div className="flex items-center gap-2"> + <label + htmlFor={testId ? `${testId}-select` : undefined} + className="text-xs uppercase tracking-wide text-zinc-500" + > + Queue + </label> + <select + id={testId ? `${testId}-select` : undefined} + value={value} + onChange={(e) => { + const next = e.target.value; + if (next === NEW_QUEUE_OPTION) { + setCreating(true); + return; + } + onChange(next); + }} + data-testid={testId} + className="text-sm font-mono px-2 py-1 rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900" + > + {options.map((opt) => ( + <option key={opt} value={opt}> + {opt} + </option> + ))} + <option value={NEW_QUEUE_OPTION}>+ new queue…</option> + </select> + </div> + ); +} diff --git a/editor/app/api/jobs/[id]/log/route.ts b/editor/app/api/jobs/[id]/log/route.ts @@ -22,7 +22,16 @@ export async function GET( const logPath = path.join(paths.jobsDir, `${id}.log`); const { content, nextOffset } = await readLogChunk(logPath, from); - const record = getRegistry().get(id); + const registry = getRegistry(); + const record = registry.get(id); const status = record ? record.status : "archived"; - return NextResponse.json({ content, nextOffset, status }); + const queueKey = record?.queueKey; + const queuePosition = record ? registry.positionInQueue(id) : -1; + return NextResponse.json({ + content, + nextOffset, + status, + queueKey, + queuePosition, + }); } diff --git a/editor/app/build/_components/BuildButtons.tsx b/editor/app/build/_components/BuildButtons.tsx @@ -0,0 +1,63 @@ +"use client"; + +import { useState } from "react"; +import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; +import { QueuePicker } from "../../_components/QueuePicker"; +import { cancelJobAction } from "../../jobs/actions"; +import { buildExportAction, buildIndexAction } from "../buildAction"; + +type Props = { + existingQueues: string[]; +}; + +export function BuildButtons({ existingQueues }: Props) { + const defaultQueueKey = "build"; + const [queueKey, setQueueKey] = useState(defaultQueueKey); + + return ( + <div className="flex flex-col gap-8"> + <QueuePicker + value={queueKey} + onChange={setQueueKey} + defaultQueueKey={defaultQueueKey} + existingQueues={existingQueues} + testId="build-queue" + /> + <section className="flex flex-col gap-3"> + <div> + <h2 className="text-lg font-semibold">Build index</h2> + <p className="text-sm text-zinc-500"> + Re-scans <code>transcripts/channels/</code> and rewrites paginated + JSON in <code>export/public/</code>. Cheap when nothing changed + (mtime short-circuit). + </p> + </div> + <StreamActionLog + trigger={() => buildIndexAction(queueKey)} + cancelAction={cancelJobAction} + buttonLabel="Build index" + runningLabel="Building index…" + testId="build-index" + /> + </section> + + <section className="flex flex-col gap-3 border-t border-zinc-200 dark:border-zinc-800 pt-6"> + <div> + <h2 className="text-lg font-semibold">Build static export</h2> + <p className="text-sm text-zinc-500"> + Spawns <code>pnpm run build</code> in <code>export/</code> — runs + the index build, then <code>next build</code> to produce the + static site at <code>export/out/</code>. + </p> + </div> + <StreamActionLog + trigger={() => buildExportAction(queueKey)} + cancelAction={cancelJobAction} + buttonLabel="Build static export" + runningLabel="Building static export…" + testId="build-export" + /> + </section> + </div> + ); +} diff --git a/editor/app/build/buildAction.ts b/editor/app/build/buildAction.ts @@ -9,11 +9,15 @@ import { type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; -export async function buildIndexAction(): Promise<StreamActionResult> { +const DEFAULT_BUILD_QUEUE = "build"; + +export async function buildIndexAction( + queueKey?: string, +): Promise<StreamActionResult> { const paths = getPaths(); const result = await runManagedFunction({ kind: "build-index", - mutexKey: "build", + queueKey: queueKey?.trim() || DEFAULT_BUILD_QUEUE, paths, fn: async (onLog) => { await buildIndex({ paths, onLog }); @@ -23,11 +27,13 @@ export async function buildIndexAction(): Promise<StreamActionResult> { return result; } -export async function buildExportAction(): Promise<StreamActionResult> { +export async function buildExportAction( + queueKey?: string, +): Promise<StreamActionResult> { const paths = getPaths(); return runManagedCommand({ kind: "build-export", - mutexKey: "build", + queueKey: queueKey?.trim() || DEFAULT_BUILD_QUEUE, paths, cwd: paths.exportDir, command: "pnpm", diff --git a/editor/app/build/page.tsx b/editor/app/build/page.tsx @@ -1,47 +1,14 @@ -import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; -import { cancelJobAction } from "../jobs/actions"; -import { buildExportAction, buildIndexAction } from "./buildAction"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { BuildButtons } from "./_components/BuildButtons"; + +export const dynamic = "force-dynamic"; export default function BuildPage() { + const existingQueues = getRegistry().activeQueueNames(); return ( - <div className="flex flex-col gap-8"> + <div className="flex flex-col gap-4"> <h1 className="text-2xl font-semibold">Build</h1> - - <section className="flex flex-col gap-3"> - <div> - <h2 className="text-lg font-semibold">Build index</h2> - <p className="text-sm text-zinc-500"> - Re-scans <code>transcripts/channels/</code> and rewrites paginated - JSON in <code>export/public/</code>. Cheap when nothing changed - (mtime short-circuit). - </p> - </div> - <StreamActionLog - trigger={buildIndexAction} - cancelAction={cancelJobAction} - buttonLabel="Build index" - runningLabel="Building index…" - testId="build-index" - /> - </section> - - <section className="flex flex-col gap-3 border-t border-zinc-200 dark:border-zinc-800 pt-6"> - <div> - <h2 className="text-lg font-semibold">Build static export</h2> - <p className="text-sm text-zinc-500"> - Spawns <code>pnpm run build</code> in <code>export/</code> — runs - the index build, then <code>next build</code> to produce the - static site at <code>export/out/</code>. - </p> - </div> - <StreamActionLog - trigger={buildExportAction} - cancelAction={cancelJobAction} - buttonLabel="Build static export" - runningLabel="Building static export…" - testId="build-export" - /> - </section> + <BuildButtons existingQueues={existingQueues} /> </div> ); } diff --git a/editor/app/channels/[slug]/_components/PipelinePanel.tsx b/editor/app/channels/[slug]/_components/PipelinePanel.tsx @@ -1,6 +1,8 @@ "use client"; +import { useState } from "react"; import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; +import { QueuePicker } from "../../../_components/QueuePicker"; import { cancelJobAction } from "../../../jobs/actions"; import { downloadAction, @@ -11,9 +13,13 @@ import { type Props = { slug: string; hasUrl: boolean; + existingQueues: string[]; }; -export function PipelinePanel({ slug, hasUrl }: Props) { +export function PipelinePanel({ slug, hasUrl, existingQueues }: Props) { + const defaultQueueKey = `channel:${slug}`; + const [queueKey, setQueueKey] = useState(defaultQueueKey); + if (!hasUrl) { return ( <p className="text-sm text-zinc-500"> @@ -24,13 +30,20 @@ export function PipelinePanel({ slug, hasUrl }: Props) { } return ( <div className="flex flex-col gap-6"> + <QueuePicker + value={queueKey} + onChange={setQueueKey} + defaultQueueKey={defaultQueueKey} + existingQueues={existingQueues} + testId="pipeline-queue" + /> <div className="flex flex-col gap-2"> <Heading title="Store playlist" desc="Fetch the channel's full URL list (--flat-playlist --skip-download --print url) and save it for later runs." /> <StreamActionLog - trigger={() => storePlaylistAction(slug)} + trigger={() => storePlaylistAction(slug, queueKey)} cancelAction={cancelJobAction} buttonLabel="Store playlist" runningLabel="Storing playlist…" @@ -43,7 +56,7 @@ export function PipelinePanel({ slug, hasUrl }: Props) { desc="Read the saved playlist, prefilter against the archive, then run yt-dlp on whatever's left. Use for first-time imports and resumes." /> <StreamActionLog - trigger={() => downloadAction(slug)} + trigger={() => downloadAction(slug, queueKey)} cancelAction={cancelJobAction} buttonLabel="Download from playlist" runningLabel="Downloading…" @@ -56,7 +69,7 @@ export function PipelinePanel({ slug, hasUrl }: Props) { desc="Quick incremental fetch using --lazy-playlist + --break-on-existing. Run regularly to grab the newest videos." /> <StreamActionLog - trigger={() => syncAction(slug)} + trigger={() => syncAction(slug, queueKey)} cancelAction={cancelJobAction} buttonLabel="Sync" runningLabel="Syncing…" diff --git a/editor/app/channels/[slug]/_components/WhisperPanel.tsx b/editor/app/channels/[slug]/_components/WhisperPanel.tsx @@ -2,6 +2,7 @@ import { useState } from "react"; import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; +import { QueuePicker } from "../../../_components/QueuePicker"; import { cancelJobAction } from "../../../jobs/actions"; import { retryFailuresAction, @@ -12,18 +13,28 @@ import { type Props = { slug: string; + existingQueues: string[]; }; -export function WhisperPanel({ slug }: Props) { +export function WhisperPanel({ slug, existingQueues }: Props) { + const defaultQueueKey = `channel:${slug}`; + const [queueKey, setQueueKey] = useState(defaultQueueKey); return ( <div className="flex flex-col gap-6"> + <QueuePicker + value={queueKey} + onChange={setQueueKey} + defaultQueueKey={defaultQueueKey} + existingQueues={existingQueues} + testId="whisper-queue" + /> <div className="flex flex-col gap-2"> <Heading title="Transcribe missing" desc="Run whisper-cli over every video that has audio but no transcript.json. Failures append to channels/<slug>/failed-transcriptions." /> <StreamActionLog - trigger={() => transcribeMissingAction(slug)} + trigger={() => transcribeMissingAction(slug, queueKey)} cancelAction={cancelJobAction} buttonLabel="Transcribe missing" runningLabel="Transcribing…" @@ -36,7 +47,7 @@ export function WhisperPanel({ slug }: Props) { desc="Re-run whisper for any video listed in failed-transcriptions." /> <StreamActionLog - trigger={() => retryFailuresAction(slug)} + trigger={() => retryFailuresAction(slug, queueKey)} cancelAction={cancelJobAction} buttonLabel="Retry failures" runningLabel="Retrying…" diff --git a/editor/app/channels/[slug]/page.tsx b/editor/app/channels/[slug]/page.tsx @@ -2,6 +2,7 @@ import Link from "next/link"; import { notFound } from "next/navigation"; import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; import { ChannelFormClient } from "../_components/ChannelFormClient"; import { DeleteChannelForm } from "../_components/DeleteChannelForm"; import { PipelinePanel } from "./_components/PipelinePanel"; @@ -25,6 +26,7 @@ export default async function ChannelDetailPage({ const update = updateChannelAction.bind(null, slug); const del = deleteChannelAction.bind(null, slug); + const existingQueues = getRegistry().activeQueueNames(); return ( <div className="flex flex-col gap-6"> @@ -53,13 +55,17 @@ export default async function ChannelDetailPage({ <section className="flex flex-col gap-3 border-t border-zinc-200 dark:border-zinc-800 pt-6"> <h2 className="text-lg font-semibold">Pipeline</h2> - <PipelinePanel slug={slug} hasUrl={!!config.url} /> + <PipelinePanel + slug={slug} + hasUrl={!!config.url} + existingQueues={existingQueues} + /> </section> {config.handling === "transcribe" && ( <section className="flex flex-col gap-3 border-t border-zinc-200 dark:border-zinc-800 pt-6"> <h2 className="text-lg font-semibold">Transcription</h2> - <WhisperPanel slug={slug} /> + <WhisperPanel slug={slug} existingQueues={existingQueues} /> </section> )} diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -9,10 +9,15 @@ import { type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +function defaultQueueKey(slug: string): string { + return `channel:${slug}`; +} + async function runPipelineAction( slug: string, mode: "store-playlist" | "download-from-playlist" | "sync", kind: string, + queueKey?: string, ): Promise<StreamActionResult> { const paths = getPaths(); const channelConfig = await readChannelConfig(paths, slug); @@ -24,10 +29,9 @@ async function runPipelineAction( } return runManagedFunction({ kind, - mutexKey: `channel:${slug}`, + queueKey: queueKey?.trim() || defaultQueueKey(slug), paths, channelSlug: slug, - busyMessage: `Channel "${slug}" already has an operation running`, fn: async (onLog, signal) => { await runYtdlp({ channelSlug: slug, @@ -45,22 +49,26 @@ async function runPipelineAction( export async function storePlaylistAction( slug: string, + queueKey?: string, ): Promise<StreamActionResult> { - return runPipelineAction(slug, "store-playlist", "store-playlist"); + return runPipelineAction(slug, "store-playlist", "store-playlist", queueKey); } export async function downloadAction( slug: string, + queueKey?: string, ): Promise<StreamActionResult> { return runPipelineAction( slug, "download-from-playlist", "download-from-playlist", + queueKey, ); } export async function syncAction( slug: string, + queueKey?: string, ): Promise<StreamActionResult> { - return runPipelineAction(slug, "sync", "sync"); + return runPipelineAction(slug, "sync", "sync", queueKey); } diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts @@ -9,16 +9,20 @@ import { type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +function defaultQueueKey(slug: string): string { + return `channel:${slug}`; +} + export async function transcribeMissingAction( slug: string, + queueKey?: string, ): Promise<StreamActionResult> { const paths = getPaths(); return runManagedFunction({ kind: "whisper-all", - mutexKey: `channel:${slug}`, + queueKey: queueKey?.trim() || defaultQueueKey(slug), paths, channelSlug: slug, - busyMessage: `Channel "${slug}" already has an operation running`, fn: async (onLog, signal) => { const result = await runWhisperBatch({ channelSlug: slug, @@ -37,14 +41,14 @@ export async function transcribeMissingAction( export async function retryFailuresAction( slug: string, + queueKey?: string, ): Promise<StreamActionResult> { const paths = getPaths(); return runManagedFunction({ kind: "whisper-retry", - mutexKey: `channel:${slug}`, + queueKey: queueKey?.trim() || defaultQueueKey(slug), paths, channelSlug: slug, - busyMessage: `Channel "${slug}" already has an operation running`, fn: async (onLog, signal) => { const result = await runWhisperBatch({ channelSlug: slug, diff --git a/editor/app/jobs/[id]/_components/JobLogTail.tsx b/editor/app/jobs/[id]/_components/JobLogTail.tsx @@ -19,6 +19,8 @@ type LogResponse = { content: string; nextOffset: number; status: string; + queueKey?: string; + queuePosition?: number; }; export function JobLogTail({ jobId, initiallyRunning, testId }: Props) { @@ -26,6 +28,8 @@ export function JobLogTail({ jobId, initiallyRunning, testId }: Props) { const [status, setStatus] = useState<string>( initiallyRunning ? "running" : "loaded", ); + const [queueKey, setQueueKey] = useState<string | undefined>(); + const [queuePosition, setQueuePosition] = useState<number>(-1); const [cancelling, setCancelling] = useState(false); const offsetRef = useRef(0); const preRef = useRef<HTMLPreElement | null>(null); @@ -60,7 +64,11 @@ export function JobLogTail({ jobId, initiallyRunning, testId }: Props) { if (data.content) setLog((prev) => prev + data.content); offsetRef.current = data.nextOffset; setStatus(data.status); - if (data.status === "running") { + setQueueKey(data.queueKey); + setQueuePosition( + typeof data.queuePosition === "number" ? data.queuePosition : -1, + ); + if (data.status === "running" || data.status === "queued") { timer = setTimeout(tick, 600); } } catch { @@ -84,6 +92,7 @@ export function JobLogTail({ jobId, initiallyRunning, testId }: Props) { } const isRunning = status === "running"; + const isQueued = status === "queued"; return ( <div className="flex flex-col gap-2" data-testid={testId}> @@ -94,7 +103,16 @@ export function JobLogTail({ jobId, initiallyRunning, testId }: Props) { {status} </span> </div> - {isRunning && ( + {isQueued && queueKey && queuePosition > 0 && ( + <div + className="text-xs text-zinc-500" + data-testid={testId ? `${testId}-queue-banner` : undefined} + > + Queued in <code className="font-mono">{queueKey}</code>, position{" "} + {queuePosition} + </div> + )} + {(isRunning || isQueued) && ( <button type="button" onClick={handleCancel} diff --git a/editor/app/jobs/[id]/page.tsx b/editor/app/jobs/[id]/page.tsx @@ -42,7 +42,14 @@ export default async function JobDetailPage({ )} <dl className="text-sm grid grid-cols-2 sm:grid-cols-4 gap-3"> <Cell label="Status" value={job.status} /> - <Cell label="Started" value={new Date(job.startedAt).toLocaleString()} /> + <Cell + label="Started" + value={ + job.startedAt + ? new Date(job.startedAt).toLocaleString() + : new Date(job.queuedAt).toLocaleString() + } + /> <Cell label="Ended" value={job.endedAt ? new Date(job.endedAt).toLocaleString() : "—"} @@ -57,7 +64,9 @@ export default async function JobDetailPage({ </p> <JobLogTail jobId={id} - initiallyRunning={job.status === "running"} + initiallyRunning={ + job.status === "running" || job.status === "queued" + } testId="job-log" /> </div> diff --git a/editor/app/jobs/page.tsx b/editor/app/jobs/page.tsx @@ -17,6 +17,8 @@ function fmtDuration(ms: number): string { function statusColor(status: string): string { switch (status) { + case "queued": + return "bg-zinc-200 text-zinc-700 dark:bg-zinc-700 dark:text-zinc-200"; case "running": return "bg-blue-100 text-blue-800 dark:bg-blue-900 dark:text-blue-200"; case "done": @@ -32,7 +34,9 @@ function statusColor(status: string): string { export default async function JobsPage() { const jobs = await listAllJobs(getPaths()); - const hasNonTerminal = jobs.some((j) => j.status === "running"); + const hasNonTerminal = jobs.some( + (j) => j.status === "running" || j.status === "queued", + ); return ( <div className="flex flex-col gap-4"> <JobsAutoRefresh hasNonTerminal={hasNonTerminal} /> @@ -57,6 +61,7 @@ export default async function JobsPage() { <th className="text-left font-medium px-3 py-2">ID</th> <th className="text-left font-medium px-3 py-2">Kind</th> <th className="text-left font-medium px-3 py-2">Channel</th> + <th className="text-left font-medium px-3 py-2">Queue</th> <th className="text-left font-medium px-3 py-2">Status</th> <th className="text-left font-medium px-3 py-2">Started</th> <th className="text-left font-medium px-3 py-2">Duration</th> @@ -66,9 +71,12 @@ export default async function JobsPage() { </thead> <tbody> {jobs.map((j) => { + const startedAt = j.startedAt ?? j.queuedAt; const dur = j.endedAt - ? j.endedAt - j.startedAt - : Date.now() - j.startedAt; + ? j.endedAt - startedAt + : j.startedAt + ? Date.now() - j.startedAt + : 0; return ( <tr key={j.id} @@ -99,6 +107,12 @@ export default async function JobsPage() { "—" )} </td> + <td + className="px-3 py-2 font-mono text-xs" + data-testid={`job-row-queue-${j.id}`} + > + {j.queueKey ?? "—"} + </td> <td className="px-3 py-2"> <span className={`text-xs uppercase tracking-wide px-2 py-0.5 rounded ${statusColor(j.status)}`} @@ -107,7 +121,7 @@ export default async function JobsPage() { </span> </td> <td className="px-3 py-2 text-xs text-zinc-500"> - {new Date(j.startedAt).toLocaleString()} + {new Date(startedAt).toLocaleString()} </td> <td className="px-3 py-2 text-xs text-zinc-500"> {fmtDuration(dur)} @@ -116,7 +130,7 @@ export default async function JobsPage() { {j.logSize.toLocaleString()} B </td> <td className="px-3 py-2 text-right"> - {j.status === "running" && ( + {(j.status === "running" || j.status === "queued") && ( <CancelJobButton jobId={j.id} testId={`cancel-${j.id}`} diff --git a/editor/app/page.tsx b/editor/app/page.tsx @@ -27,7 +27,7 @@ export default async function Dashboard() { const lastBuild = await getLastBuildTime(paths); const runningJobs = getRegistry() .list() - .filter((j) => j.status === "running"); + .filter((j) => j.status === "running" || j.status === "queued"); return ( <div className="flex flex-col gap-6"> diff --git a/editor/cypress/e2e/queues.cy.ts b/editor/cypress/e2e/queues.cy.ts @@ -0,0 +1,115 @@ +describe("Named job queues", () => { + it("default queue is channel:<slug>", () => { + cy.resetData("slow-pipeline-channel"); + cy.visit("/channels/slow-channel"); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-log", { timeout: 15_000 }).should( + "contain.text", + "cypress-slow", + ); + cy.visit("/jobs"); + cy.get('[data-testid^="job-row-queue-"]') + .first() + .should("have.text", "channel:slow-channel"); + }); + + it("queues a second job in the same queue, runs sequentially", () => { + cy.resetData("two-slow-channels"); + + cy.visit("/channels/slow-a"); + cy.findByTestId("pipeline-queue").select("+ new queue…"); + cy.findByTestId("pipeline-queue-input").clear().type("qShared"); + cy.findByTestId("pipeline-queue-save").click(); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-log", { timeout: 15_000 }).should( + "contain.text", + "cypress-slow", + ); + + cy.visit("/channels/slow-b"); + cy.findByTestId("pipeline-queue").select("qShared"); + cy.findByTestId("pipeline-sync").find("button").click(); + // Second job should display the queued banner. + cy.findByTestId("pipeline-sync-queue-banner", { timeout: 10_000 }).should( + "contain.text", + "qShared", + ); + + cy.visit("/jobs"); + // Two rows, both in queue qShared. Most recent first. + cy.get('[data-testid^="job-row-queue-"]') + .first() + .should("have.text", "qShared"); + cy.get('[data-testid^="job-row-queue-"]') + .eq(1) + .should("have.text", "qShared"); + + // Second row (slow-a, queued first) is currently running. First row + // (slow-b, queued second) is queued. + cy.contains("tr", "slow-a").contains(/^running$/i); + cy.contains("tr", "slow-b").contains(/^queued$/i); + + // Cancel the running job; the queued one should advance to running. + cy.contains("tr", "slow-a") + .find('[data-testid^="cancel-"]') + .click(); + cy.contains("tr", "slow-b").contains(/^running$/i, { timeout: 10_000 }); + }); + + it("runs jobs in different queues in parallel", () => { + cy.resetData("two-slow-channels"); + + cy.visit("/channels/slow-a"); + cy.findByTestId("pipeline-queue").select("+ new queue…"); + cy.findByTestId("pipeline-queue-input").clear().type("qA"); + cy.findByTestId("pipeline-queue-save").click(); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-log", { timeout: 15_000 }).should( + "contain.text", + "cypress-slow", + ); + + cy.visit("/channels/slow-b"); + cy.findByTestId("pipeline-queue").select("+ new queue…"); + cy.findByTestId("pipeline-queue-input").clear().type("qB"); + cy.findByTestId("pipeline-queue-save").click(); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-log", { timeout: 15_000 }).should( + "contain.text", + "cypress-slow", + ); + + cy.visit("/jobs"); + cy.contains("tr", "slow-a").contains(/^running$/i); + cy.contains("tr", "slow-b").contains(/^running$/i); + }); + + it("cancels a queued job without disturbing the one running ahead of it", () => { + cy.resetData("two-slow-channels"); + + cy.visit("/channels/slow-a"); + cy.findByTestId("pipeline-queue").select("+ new queue…"); + cy.findByTestId("pipeline-queue-input").clear().type("qShared"); + cy.findByTestId("pipeline-queue-save").click(); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-log", { timeout: 15_000 }).should( + "contain.text", + "cypress-slow", + ); + + cy.visit("/channels/slow-b"); + cy.findByTestId("pipeline-queue").select("qShared"); + cy.findByTestId("pipeline-sync").find("button").click(); + cy.findByTestId("pipeline-sync-queue-banner", { timeout: 10_000 }); + + cy.visit("/jobs"); + cy.contains("tr", "slow-b").contains(/^queued$/i); + + cy.contains("tr", "slow-b") + .find('[data-testid^="cancel-"]') + .click(); + cy.contains("tr", "slow-b").contains(/^cancelled$/i, { timeout: 10_000 }); + // slow-a (in front) keeps running. + cy.contains("tr", "slow-a").contains(/^running$/i); + }); +}); diff --git a/editor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-a/config.json b/editor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-a/config.json @@ -0,0 +1,6 @@ +{ + "handling": "youtube", + "name": "Slow A", + "url": "https://www.youtube.com/@slow-a/videos", + "ytdlpExtraArgs": ["--cy-slow"] +} diff --git a/editor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-b/config.json b/editor/cypress/fixtures/test-transcripts/two-slow-channels/channels/slow-b/config.json @@ -0,0 +1,6 @@ +{ + "handling": "youtube", + "name": "Slow B", + "url": "https://www.youtube.com/@slow-b/videos", + "ytdlpExtraArgs": ["--cy-slow"] +}