Archilyzer · Source

archilyzer

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

commit 35f1971861862aefd3ef9137659d7e473a10b2a9
parent 17b7d14391a6c10a018b9f71fb97059aef4f6153
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 19 Jun 2026 15:39:43 -0400

Compact bookmarks menu + management page, with batched job/widget/download work

Bookmarks panel on /jobs and /jobs/active is now a compact one-click run
menu (fire-and-forget, no inline shell/expanding log); rename, delete,
metadata, and the streaming "Run again" move to a new /jobs/bookmarks
management page reachable via a "Manage" link.

Also in this batch (see editor/CHANGELOG.md):
- Bookmark-a-job feature itself (JobSpec replay descriptor + bookmarks store)
- Smarter batch ETAs (parallel-wave rounding) and an audio-probe progress bar
- Indeterminate progress bars now read as an amber pulse
- Individually toggleable monitor-widget elements
- Audio-integrity probe reads the .part in place instead of copying it
- Corrupt-source videos no longer recorded as failed transcriptions

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

Diffstat:
Mcommon/components/StreamActionLog.tsx | 15++++++++++++---
Mcommon/controller/channelSnapshot.ts | 21++++++++++++++++++++-
Mcommon/controller/transcribeError.ts | 7+++++--
Mcommon/controller/transcribeOne.ts | 11+++++++++--
Mcommon/controller/whisperBatch.ts | 55++++++++++++++++++++++++++++++++++++++++++++++++++++---
Acommon/jobs/bookmarks.ts | 134+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/jobMeta.ts | 11++++++++++-
Acommon/jobs/jobSpec.ts | 60++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/listJobs.ts | 7+++++++
Mcommon/jobs/progressParsers.ts | 115+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/jobs/registry.ts | 30+++++++++++++++++++++++++++++-
Mcommon/jobs/streamCommand.ts | 20++++++++++++++++++--
Mcommon/jobs/taskHooks.ts | 9+++++++--
Mcommon/lib/paths.ts | 5+++++
Mcommon/ytdlp/audioCheckedDownload.ts | 210++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Meditor/CHANGELOG.md | 7+++++++
Meditor/app/channels/[slug]/components/RetryBucketControl.tsx | 6++++++
Meditor/app/channels/[slug]/components/VideoListPane.tsx | 16++++++++++------
Meditor/app/channels/[slug]/components/stages/DownloadStage.tsx | 2++
Meditor/app/channels/[slug]/components/stages/TranscribeStage.tsx | 1+
Meditor/app/channels/[slug]/lib/stageStatus.ts | 10++++++++++
Meditor/app/channels/[slug]/lib/videoRows.ts | 10+++++++++-
Meditor/app/channels/[slug]/lib/videoRowsServer.ts | 9++++++++-
Meditor/app/channels/[slug]/pipelineActions.ts | 57+++++++++++++++++++++++++++++++++++++++++++++++++++------
Meditor/app/channels/[slug]/whisperActions.ts | 37+++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/[id]/page.tsx | 6+++++-
Meditor/app/jobs/active/buildActiveJobs.ts | 10+++++++++-
Meditor/app/jobs/active/page.tsx | 8+++++++-
Aeditor/app/jobs/bookmarkActions.ts | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/bookmarks/page.tsx | 37+++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/BookmarkJobButton.tsx | 50++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/BookmarkRunButton.tsx | 94+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/BookmarksList.tsx | 204+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/jobs/components/BookmarksMenu.tsx | 45+++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/jobs/components/JobsTable.tsx | 10+++++++---
Meditor/app/jobs/components/RunningJobsList.tsx | 61+++++++++++++++++++++++++++++++++++++++++++++++++------------
Aeditor/app/jobs/loadBookmarks.ts | 22++++++++++++++++++++++
Meditor/app/jobs/page.tsx | 8+++++++-
Aeditor/app/jobs/runJobSpec.ts | 134+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/app/widget/builder/components/WidgetBuilder.tsx | 25+++++++++++++++++++++++++
Meditor/app/widget/components/MonitorWidget.tsx | 117+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Meditor/app/widget/lib/config.ts | 32++++++++++++++++++++++++++++++++
Meditor/app/workers/components/WorkersView.tsx | 4+++-
Meditor/e2e/audio-check-scenarios.spec.ts | 43+++++++++++++++++++++++++++++++++++++++++++
Aeditor/e2e/bookmarks.spec.ts | 269+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Meditor/e2e/whisper.spec.ts | 76+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Meditor/e2e/widget.spec.ts | 29+++++++++++++++++++++++++++++
47 files changed, 2119 insertions(+), 123 deletions(-)

diff --git a/common/components/StreamActionLog.tsx b/common/components/StreamActionLog.tsx @@ -47,6 +47,9 @@ export function StreamActionLog({ const [running, setRunning] = useState(false); const [log, setLog] = useState(""); const [error, setError] = useState<string | null>(null); + // When the failure is really a neutral notice (info: true), render it as a + // status rather than a red error. + const [errorIsInfo, setErrorIsInfo] = useState(false); const [jobId, setJobId] = useState<string | null>(null); const [cancelling, setCancelling] = useState(false); const [poll, setPoll] = useState<StatusPoll | null>(null); @@ -106,6 +109,7 @@ export function StreamActionLog({ async function handleClick() { setError(null); + setErrorIsInfo(false); setLog(""); setJobId(null); setPoll(null); @@ -116,6 +120,7 @@ export function StreamActionLog({ const result = await trigger(); if (!result.ok) { setError(result.error); + setErrorIsInfo(result.info === true); return; } started = true; @@ -179,9 +184,13 @@ export function StreamActionLog({ </div> {error && ( <div - role="alert" - aria-label={`${accessibleName} error`} - className="rounded border border-red-300 bg-red-50 dark:border-red-800 dark:bg-red-950 px-3 py-2 text-sm text-red-700 dark:text-red-300" + role={errorIsInfo ? "status" : "alert"} + aria-label={`${accessibleName} ${errorIsInfo ? "notice" : "error"}`} + className={ + errorIsInfo + ? "rounded border border-zinc-300 bg-zinc-50 dark:border-zinc-700 dark:bg-zinc-900 px-3 py-2 text-sm text-zinc-700 dark:text-zinc-300" + : "rounded border border-red-300 bg-red-50 dark:border-red-800 dark:bg-red-950 px-3 py-2 text-sm text-red-700 dark:text-red-300" + } > {error} </div> diff --git a/common/controller/channelSnapshot.ts b/common/controller/channelSnapshot.ts @@ -62,6 +62,11 @@ export type ChannelSnapshot = { missingFromArchive: string[]; duplicateDirs: string[]; partialDownloads: string[]; + // Videos whose latest download-outcome.json is "failed-corrupt-source" (the + // audio-check pipeline gave up on a malformed source) and that still have no + // transcript. Surfaced as a download/source problem — NOT counted as a failed + // transcription. Optional: older snapshots lack it; readers default to []. + corruptSource: string[]; // Videos whose transcript rides on a non-canonical VTT name (e.g. only // transcript.en-US.vtt, or a foreign-only transcript.<lang>.vtt) instead of // the standard transcript.en.vtt — surfaced so the user can normalize/switch @@ -311,6 +316,7 @@ export async function generateChannelSnapshot( const untranscribable: string[] = []; const noMetadata: string[] = []; const partialDownloads: string[] = []; + const corruptSource: string[] = []; const nonStandardVtt: string[] = []; const skippedByFilter: string[] = []; let transcribedWithAudioBytes = 0; @@ -330,6 +336,15 @@ export async function generateChannelSnapshot( ) { skippedByFilter.push(id); } + // The audio-check pipeline gave up on a malformed source (left no real + // audio). Surface it as a download/source problem rather than letting the + // transcribe pass mislabel it as a failed transcription. + if ( + outcome?.status === "failed-corrupt-source" && + !isVideoTranscribed(files) + ) { + corruptSource.push(id); + } if ( targetAudioFile && files.audioFiles.length > 0 && @@ -447,6 +462,7 @@ export async function generateChannelSnapshot( } : undefined; + const corruptSourceSet = new Set(corruptSource); const snapshot: ChannelSnapshot = { generatedAt: new Date().toISOString(), totals: { @@ -462,10 +478,13 @@ export async function generateChannelSnapshot( transcribedWithAudio: transcribedWithAudio.sort(), untranscribable: untranscribable.sort(), noMetadata: noMetadata.sort(), - failedListed: failedListed.filter((id) => !excludedById.has(id)), + failedListed: failedListed.filter( + (id) => !excludedById.has(id) && !corruptSourceSet.has(id), + ), missingFromArchive: missingFromArchive.sort(), duplicateDirs: [], partialDownloads: partialDownloads.sort(), + corruptSource: corruptSource.sort(), nonStandardVtt: nonStandardVtt.sort(), skippedByFilter: skippedByFilter.sort(), }, diff --git a/common/controller/transcribeError.ts b/common/controller/transcribeError.ts @@ -1,14 +1,17 @@ // Why a transcription failed, so a batch can decide whether to retry it on a // different worker. "transport" = the worker/transport was at fault (remote // unreachable, 5xx) and the audio may still transcribe elsewhere; "transcription" -// = the audio itself failed and retrying on another worker is pointless. +// = the audio itself failed and retrying on another worker is pointless; +// "no-audio" = the video has no real audio file (a download/source problem, e.g. +// a corrupt-source / failed-extract download), so the batch skips it rather than +// recording a failure and a re-download will let it transcribe later. // // In its own module (no deps) so both transcribeOne.ts and the remote client // (remoteTranscribe.ts) can throw it without an import cycle. export class TranscribeError extends Error { constructor( message: string, - readonly failureClass: "transport" | "transcription", + readonly failureClass: "transport" | "transcription" | "no-audio", ) { super(message); this.name = "TranscribeError"; diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -79,12 +79,19 @@ export async function transcribeOneVideo( opts.strictAudio === true, ); if (!resolvedAudio) { + // A missing audio file is a download/source problem, not a transcription + // failure. Throw "no-audio" so the batch skips it (it transcribes once a + // good audio file exists) instead of permanently recording a failure. if (opts.strictAudio) { - throw new Error( + throw new TranscribeError( `no ${opts.audioFilename} found in ${opts.videoDir} (strict audio mode)`, + "no-audio", ); } - throw new Error(`no audio file found in ${opts.videoDir}`); + throw new TranscribeError( + `no audio file found in ${opts.videoDir}`, + "no-audio", + ); } if (resolvedAudio !== opts.audioFilename) { log(`Using ${resolvedAudio} instead of ${opts.audioFilename}`); diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts @@ -2,8 +2,14 @@ import path from "node:path"; import fs from "fs-extra"; import type { Paths } from "../lib/paths"; import type { AudioFormat } from "../lib/channelConfig"; -import { isVideoTranscribed, readVideoFiles } from "../lib/videoStatus"; +import { + isVideoDownloaded, + isVideoTranscribed, + readVideoFiles, +} from "../lib/videoStatus"; import { transcribeWithWorker } from "./transcribeOne"; +import { TranscribeError } from "./transcribeError"; +import { pruneFailedTranscriptions } from "./failedTranscriptions"; import { resolveShardItems } from "./shard"; import { countNotYetTranscribed } from "./channels"; import { getWorkerPool } from "../jobs/workerPool"; @@ -92,6 +98,28 @@ export async function runWhisperBatch({ const failedSet = new Set<string>(originalFailures); const allDirs = await readdir(dataDir); + // Auto-prune mislabeled failures: any id on the failed list whose dir has no + // real audio file was never download-ready — a corrupt-source / failed-extract + // download (not a transcription failure) that got recorded here. Drop those so + // they retry once re-downloaded and so the count reflects real failures. + const onDiskDirs = new Set(allDirs); + const bogusFailures: string[] = []; + for (const id of failedSet) { + if (!onDiskDirs.has(id)) continue; + if (!isVideoDownloaded(await readVideoFiles(path.join(dataDir, id)))) { + bogusFailures.push(id); + } + } + if (bogusFailures.length > 0) { + const { pruned } = await pruneFailedTranscriptions( + failureListFile, + new Set(bogusFailures), + ); + for (const id of bogusFailures) failedSet.delete(id); + log( + `Pruned ${pruned} undownloaded video(s) from failed-transcriptions (no audio on disk — will retry once re-downloaded).`, + ); + } // When the caller supplies an explicit id set (bucket-scoped run), filter // to that intersection so we never try to transcribe ids whose dirs aren't // on disk. Otherwise consider every video dir. @@ -109,9 +137,14 @@ export async function runWhisperBatch({ for (const id of candidateIds) { if (failedSet.has(id)) continue; const vp = path.join(dataDir, id); + const files = await readVideoFiles(vp); // Already transcribed if whisper ran (transcript.json) OR yt-dlp wrote an // English VTT (transcript.en.vtt or a regional/auto fallback like en-US). - if (isVideoTranscribed(await readVideoFiles(vp))) continue; + if (isVideoTranscribed(files)) continue; + // Not downloaded (no real audio file): a download candidate, not a + // transcription candidate. Attempting it would throw "no audio file found" + // and mislabel a download/source problem as a failed transcription. + if (!isVideoDownloaded(files)) continue; fullItems.push(id); } const shardResult = await resolveShardItems({ @@ -179,11 +212,19 @@ export async function runWhisperBatch({ // fallback like en-US). Without the VTT branch, videos on a youtube-handling // channel with only VTT auto-subs end up attempted -> fail with "no audio file // found" -> get added to failed-transcriptions on every run. - if (isVideoTranscribed(await readVideoFiles(videoPath))) { + const videoFiles = await readVideoFiles(videoPath); + if (isVideoTranscribed(videoFiles)) { log(`Transcription for ${videoDir} already exists`); skipped++; return; } + // No real audio file: a download/source problem, not a transcription + // candidate. Skip (don't fail) so a re-download lets it transcribe later. + if (!isVideoDownloaded(videoFiles)) { + log(`Skipping ${videoDir}: not downloaded (no audio file)`); + skipped++; + return; + } if (strict && !(await pathExists(path.join(videoPath, resolvedAudioFilename)))) { log(`Skipping ${videoDir}: no ${resolvedAudioFilename} (strict mode)`); skipped++; @@ -227,6 +268,14 @@ export async function runWhisperBatch({ skipped++; return; } + // A missing audio file is a download/source problem, not a transcription + // failure — skip it (it transcribes once a good audio file exists) rather + // than permanently recording it in failed-transcriptions. + if (err instanceof TranscribeError && err.failureClass === "no-audio") { + log(`Skipping ${videoDir}: no audio file (download/source problem)`); + skipped++; + return; + } log(`FAILED TO TRANSCRIBE ${videoDir}: ${String(err)}`); await appendFile(failureListFile, `${videoDir}\n`); failed++; diff --git a/common/jobs/bookmarks.ts b/common/jobs/bookmarks.ts @@ -0,0 +1,134 @@ +import fs from "node:fs"; +import path from "node:path"; +import type { Paths } from "../lib/paths"; +import { parseJobSpec, type JobSpec } from "./jobSpec"; + +// Persisted job bookmarks. A bookmark captures a JobSpec (see jobSpec.ts) plus a +// human label, so an operator can re-launch the same job kind on the same +// channel with one button. Re-running re-derives work from the channel's current +// state — the bookmark never stores a frozen video-id list. +// +// Modeled on common/jobs/workerDefaults.ts: a small JSON file outside the +// in-memory registry, written atomically (tmp + rename), with tolerant reads +// that coerce a missing/corrupt file to an empty list rather than crashing. + +export type JobBookmark = { + id: string; + name: string; + spec: JobSpec; + createdAt: number; + lastRunAt?: number; +}; + +type BookmarksFile = { v: 1; bookmarks: JobBookmark[] }; + +// Stable label derived from the spec, e.g. "whisper-all · HasanAbiVODs3" or +// "retry-bucket (partialDownloads) · cornbreadman". +export function defaultBookmarkName(spec: JobSpec): string { + const scope = spec.bucket ? `${spec.kind} (${spec.bucket})` : spec.kind; + return `${scope} · ${spec.slug}`; +} + +function newBookmarkId(): string { + return `bm-${Date.now().toString(36)}-${Math.random().toString(36).slice(2, 8)}`; +} + +// Two specs are "the same" job to avoid duplicate bookmarks when the operator +// bookmarks the same recurring job twice. params order is stable because each +// action builds them the same way. +function sameSpec(a: JobSpec, b: JobSpec): boolean { + return ( + a.kind === b.kind && + a.slug === b.slug && + a.bucket === b.bucket && + JSON.stringify(a.params ?? {}) === JSON.stringify(b.params ?? {}) + ); +} + +export async function readBookmarks(paths: Paths): Promise<JobBookmark[]> { + let raw: unknown; + try { + raw = JSON.parse(await fs.promises.readFile(paths.bookmarksFile, "utf8")); + } catch { + return []; + } + if (!raw || typeof raw !== "object") return []; + const arr = (raw as Record<string, unknown>).bookmarks; + if (!Array.isArray(arr)) return []; + const out: JobBookmark[] = []; + for (const item of arr) { + if (!item || typeof item !== "object") continue; + const b = item as Record<string, unknown>; + if (typeof b.id !== "string" || !b.id) continue; + const spec = parseJobSpec(b.spec); + if (!spec) continue; + out.push({ + id: b.id, + name: + typeof b.name === "string" && b.name + ? b.name + : defaultBookmarkName(spec), + spec, + createdAt: typeof b.createdAt === "number" ? b.createdAt : 0, + lastRunAt: typeof b.lastRunAt === "number" ? b.lastRunAt : undefined, + }); + } + return out; +} + +async function writeBookmarks( + paths: Paths, + bookmarks: JobBookmark[], +): Promise<void> { + const out: BookmarksFile = { v: 1, bookmarks }; + await fs.promises.mkdir(path.dirname(paths.bookmarksFile), { + recursive: true, + }); + const tmp = `${paths.bookmarksFile}.tmp-${process.pid}`; + await fs.promises.writeFile(tmp, JSON.stringify(out, null, 2) + "\n"); + await fs.promises.rename(tmp, paths.bookmarksFile); +} + +// Add a bookmark for `spec`, returning the new (or existing duplicate) entry. +// Idempotent: re-bookmarking an identical job returns the existing bookmark +// without creating a second one. +export async function addBookmark( + paths: Paths, + spec: JobSpec, +): Promise<JobBookmark> { + const existing = await readBookmarks(paths); + const dup = existing.find((b) => sameSpec(b.spec, spec)); + if (dup) return dup; + const bookmark: JobBookmark = { + id: newBookmarkId(), + name: defaultBookmarkName(spec), + spec, + createdAt: Date.now(), + }; + await writeBookmarks(paths, [bookmark, ...existing]); + return bookmark; +} + +export async function removeBookmark( + paths: Paths, + id: string, +): Promise<void> { + const existing = await readBookmarks(paths); + const next = existing.filter((b) => b.id !== id); + if (next.length !== existing.length) await writeBookmarks(paths, next); +} + +export async function updateBookmark( + paths: Paths, + id: string, + patch: Partial<Pick<JobBookmark, "name" | "lastRunAt">>, +): Promise<void> { + const existing = await readBookmarks(paths); + let changed = false; + const next = existing.map((b) => { + if (b.id !== id) return b; + changed = true; + return { ...b, ...patch }; + }); + if (changed) await writeBookmarks(paths, next); +} diff --git a/common/jobs/jobMeta.ts b/common/jobs/jobMeta.ts @@ -2,6 +2,7 @@ import path from "node:path"; import { readFile, writeFile } from "node:fs/promises"; import type { Paths } from "../lib/paths"; import type { JobRecord, JobStatus } from "./registry"; +import { parseJobSpec, type JobSpec } from "./jobSpec"; // A serializable sidecar (`<id>.meta.json`) written alongside each job's // `<id>.log`. The in-memory registry keeps only the 100 most-recent finished @@ -21,6 +22,9 @@ export type JobMeta = { startedAt?: number; endedAt?: number; exitCode?: number; + // Replay descriptor for bookmarkable jobs, preserved so an archived (evicted + // or post-restart) job can still be bookmarked from disk. See jobSpec.ts. + spec?: JobSpec; }; export function metaPath(paths: Paths, id: string): string { @@ -46,6 +50,7 @@ export async function writeJobMeta( startedAt: record.startedAt, endedAt: record.endedAt, exitCode: record.exitCode, + spec: record.spec, }; await writeFile(metaPath(paths, record.id), JSON.stringify(meta), "utf8"); } catch { @@ -70,7 +75,11 @@ export async function readJobMeta( if (!parsed || typeof parsed !== "object") return null; const m = parsed as Record<string, unknown>; if (typeof m.id !== "string" || typeof m.kind !== "string") return null; - return parsed as JobMeta; + const meta = parsed as JobMeta; + // Sanitize the replay descriptor: drop it if malformed so a stale sidecar + // can't surface an un-runnable bookmark. + meta.spec = m.spec === undefined ? undefined : (parseJobSpec(m.spec) ?? undefined); + return meta; } catch { return null; } diff --git a/common/jobs/jobSpec.ts b/common/jobs/jobSpec.ts @@ -0,0 +1,60 @@ +// A serializable description of a job, sufficient to re-launch ("bookmark") it. +// +// It is captured at job-creation time (stored on the JobRecord and in the +// <id>.meta.json sidecar) and copied into a bookmark when the user bookmarks a +// job. The dispatcher in editor/app/jobs/runJobSpec.ts maps a spec back to the +// server action that runs it. +// +// Re-running RE-DERIVES the work from the channel's CURRENT state rather than +// replaying a frozen list: bucket jobs store only the snapshot bucket category +// and the dispatcher reads the live bucket members at run time, so e.g. +// "retry partial downloads on cornbreadman" always acts on whatever is partial +// now. Flag-style parameters (queue, audio format, abort-on-error…) are captured +// verbatim in `params`. + +// The snapshot bucket categories whose current members a bookmarked job can be +// re-derived from. A subset of ChannelSnapshot["buckets"] — only the buckets the +// bucket-style actions (retry-bucket / whisper-bucket) actually operate on. +export type ReplayBucket = + | "partialDownloads" + | "noTranscript" + | "downloadedNoTranscript"; + +export type JobSpec = { + kind: string; + slug: string; + // Present only for bucket jobs: the snapshot bucket whose current members the + // job should operate on when re-run. + bucket?: ReplayBucket; + // Flag-style parameters captured verbatim (queueKey, reverse, audioFormat, + // strictAudioFormat, abortOnError, ignoreArchive, handlingOverride, + // targetFormat, shardTotal, shardIndex…). All JSON-serializable. + params?: Record<string, unknown>; +}; + +const REPLAY_BUCKETS: ReadonlySet<string> = new Set<ReplayBucket>([ + "partialDownloads", + "noTranscript", + "downloadedNoTranscript", +]); + +// Defensive parse for a spec read back from JSON (a sidecar or the bookmarks +// file). Returns null on anything malformed so a hand-edited or stale file can't +// crash a reader. Mirrors the tolerance of readJobMeta / readWorkerDefaults. +export function parseJobSpec(raw: unknown): JobSpec | null { + if (!raw || typeof raw !== "object") return null; + const r = raw as Record<string, unknown>; + if (typeof r.kind !== "string" || !r.kind) return null; + if (typeof r.slug !== "string" || !r.slug) return null; + const spec: JobSpec = { kind: r.kind, slug: r.slug }; + if (r.bucket !== undefined) { + if (typeof r.bucket !== "string" || !REPLAY_BUCKETS.has(r.bucket)) { + return null; + } + spec.bucket = r.bucket as ReplayBucket; + } + if (r.params && typeof r.params === "object" && !Array.isArray(r.params)) { + spec.params = r.params as Record<string, unknown>; + } + return spec; +} diff --git a/common/jobs/listJobs.ts b/common/jobs/listJobs.ts @@ -16,6 +16,9 @@ export type JobListEntry = { endedAt?: number; exitCode?: number; inRegistry: boolean; + // True when the job carries a replay descriptor (JobSpec) and can be + // bookmarked / re-run. See common/jobs/jobSpec.ts. + bookmarkable: boolean; logPath: string; logSize: number; }; @@ -71,6 +74,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { endedAt: live_.endedAt, exitCode: live_.exitCode, inRegistry: true, + bookmarkable: Boolean(live_.spec), logPath, logSize, }); @@ -94,6 +98,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { endedAt: meta.endedAt, exitCode: meta.exitCode, inRegistry: false, + bookmarkable: Boolean(meta.spec), logPath, logSize, }); @@ -105,6 +110,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { startedAt: mtime, endedAt: mtime, inRegistry: false, + bookmarkable: false, logPath, logSize, }); @@ -128,6 +134,7 @@ export async function listAllJobs(paths: Paths): Promise<JobListEntry[]> { endedAt: r.endedAt, exitCode: r.exitCode, inRegistry: true, + bookmarkable: Boolean(r.spec), logPath: r.logPath, logSize: 0, }); diff --git a/common/jobs/progressParsers.ts b/common/jobs/progressParsers.ts @@ -3,7 +3,29 @@ // unit-tested directly if a unit runner is ever added. Used by taskHooks.ts to // drive the per-task progress bars on the Active Jobs screen. -export type ProgressUpdate = { fraction?: number; detail?: string }; +export type ProgressUpdate = { + fraction?: number; + detail?: string; + // Audio-integrity phase of a download task. "probing" while yt-dlp is paused + // for an ffmpeg integrity probe; "downloading" otherwise. Only emitted by the + // managed download parser (createDownloadProgressParser). + phase?: "downloading" | "probing"; + // Estimated duration of the in-flight probe, in seconds — drives the + // determinate "scanning" bar. Undefined on the first probe (no prior sample). + probeEtaSeconds?: number; +}; + +// Structured marker the audio-checked download orchestrator emits around each +// integrity probe (mirrors DLOM_PROGRESS / DLOM_ARCHIVE). createDownloadProgressParser +// turns these into probe-phase updates and uses the measured probe durations to +// inflate yt-dlp's ETA (which is blind to the SIGSTOP pauses). See +// common/ytdlp/audioCheckedDownload.ts. +export const DLOM_PROBE_MARKER = "DLOM_PROBE"; + +// Smoothing factor for the per-probe duration trend (diffMs). Light enough that +// a single noisy sample (a probe skipped by the size gate, CPU contention) does +// not whipsaw the projected overhead. +const PROBE_DIFF_EWMA_ALPHA = 0.4; function clamp01(n: number): number { return Math.max(0, Math.min(1, n)); @@ -59,6 +81,9 @@ function num(value: string | undefined): number | undefined { // Returns null for any other line so the legacy parser can have a turn. export function parseStructuredDownloadProgress( line: string, + // Optional transform applied to yt-dlp's raw ETA (seconds) before it is + // formatted — used by the managed download parser to fold in probe overhead. + adjustEtaSeconds?: (etaSeconds: number) => number, ): ProgressUpdate | null { if (!line.includes(DOWNLOAD_PROGRESS_MARKER)) return null; const field = (key: string): string | undefined => @@ -87,7 +112,10 @@ export function parseStructuredDownloadProgress( const detailParts: string[] = []; if (speed !== undefined) detailParts.push(formatBytesPerSec(speed)); - if (eta !== undefined) detailParts.push(`ETA ${formatClock(eta)}`); + if (eta !== undefined) { + const shownEta = adjustEtaSeconds ? adjustEtaSeconds(eta) : eta; + detailParts.push(`ETA ${formatClock(shownEta)}`); + } if (fragIndex !== undefined && fragCount !== undefined) { detailParts.push(`frag ${fragIndex}/${fragCount}`); } @@ -140,6 +168,89 @@ export function parseDownloadProgress(line: string): ProgressUpdate | null { return out; } +// Stateful download parser for the managed (audio-checked) path. It wraps the +// pure parseStructuredDownloadProgress with two extras the pure parser can't do: +// +// 1. Probe phase. DLOM_PROBE markers (emitted by audioCheckedDownload while +// yt-dlp is SIGSTOPped for an integrity probe) flip the task to a "probing" +// phase carrying the estimated probe duration, then back to "downloading". +// +// 2. Probe-aware ETA. yt-dlp's ETA counts only active download time and is +// blind to the probe pauses. We extrapolate the per-probe duration forward +// as an arithmetic progression — each probe transcodes the whole, growing +// .part, so successive probes cost ~`diffMs` more — and add the projected +// overhead of all remaining probes on top of yt-dlp's ETA: +// N = floor(etaSeconds / intervalSeconds) // remaining probes +// overhead = N*lastProbeMs + diffMs*N*(N+1)/2 // ms +// diffMs is clamped ≥ 0 and EWMA-smoothed; a rollback (which shrinks the +// .part, so the next probe is cheaper) resets the baseline. +export function createDownloadProgressParser(): { + feed: (line: string) => ProgressUpdate | null; +} { + let lastProbeMs: number | undefined; + let diffMs = 0; // smoothed per-probe duration increase + let intervalSec: number | undefined; + + const adjustEtaSeconds = (etaSeconds: number): number => { + if (lastProbeMs === undefined || !intervalSec || intervalSec <= 0) { + return etaSeconds; + } + const n = Math.floor(etaSeconds / intervalSec); // remaining probes + if (n <= 0) return etaSeconds; + const overheadMs = n * lastProbeMs + (diffMs * (n * (n + 1))) / 2; + return etaSeconds + overheadMs / 1000; + }; + + return { + feed(line: string): ProgressUpdate | null { + if (line.includes(DLOM_PROBE_MARKER)) { + const field = (key: string): string | undefined => + line.match(new RegExp(`\\b${key}=(\\S+)`))?.[1]; + const status = field("status"); + if (status === "start") { + const estMs = num(field("estMs")); + return { + phase: "probing", + probeEtaSeconds: + estMs !== undefined && estMs > 0 ? estMs / 1000 : undefined, + }; + } + if (status === "done") { + // skip=1 means no real ffmpeg probe ran (size gate / copy timeout) — + // clear the phase but don't pollute the duration trend. + if (field("skip") !== "1") { + const d = num(field("durationMs")); + const i = num(field("intervalSec")); + const reset = field("reset") === "1"; + if (i !== undefined) intervalSec = i; + if (d !== undefined) { + if (reset || lastProbeMs === undefined) { + lastProbeMs = d; + diffMs = 0; + } else { + const rawDiff = Math.max(0, d - lastProbeMs); + diffMs = + diffMs === 0 + ? rawDiff + : PROBE_DIFF_EWMA_ALPHA * rawDiff + + (1 - PROBE_DIFF_EWMA_ALPHA) * diffMs; + lastProbeMs = d; + } + } + } + return { phase: "downloading" }; + } + return null; + } + const structured = parseStructuredDownloadProgress(line, adjustEtaSeconds); + if (structured) return structured; + // Legacy "[download] …%" lines (non-structured) carry no numeric ETA we + // can inflate, so just delegate. + return parseDownloadProgress(line); + }, + }; +} + // whisper-cli output is stateful: a header line declares the audio's total // length, and subsequent `[hh:mm:ss --> hh:mm:ss]` lines report the segment // just transcribed. Progress = latest end-timestamp / total length. diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -1,4 +1,5 @@ import type { ChildProcess } from "node:child_process"; +import type { JobSpec } from "./jobSpec"; export type JobStatus = | "queued" @@ -31,6 +32,15 @@ export type JobTask = { // For transcribe tasks: the worker running it (id into the worker pool), so // the Workers page can show which engine each task is on. workerId?: string; + // Audio-integrity phase of a download task. "probing" while yt-dlp is paused + // for an ffmpeg integrity probe (the determinate "scanning" bar), otherwise + // "downloading" (or undefined for tasks that never probe). + phase?: "downloading" | "probing"; + // Wall-clock start of the current probe, stamped when phase flips to + // "probing"; the UI fills the scanning bar against it. Cleared on resume. + probeStartedAt?: number; + // Estimated probe duration (seconds) for the scanning bar's denominator. + probeEtaSeconds?: number; }; export type JobRecord = { @@ -39,6 +49,9 @@ export type JobRecord = { queueKey: string; channelSlug?: string; videoId?: string; + // A serializable replay descriptor, set for bookmarkable job kinds. Its + // presence is what makes a job "bookmarkable" in the UI. See jobSpec.ts. + spec?: JobSpec; status: JobStatus; queuedAt: number; startedAt?: number; @@ -227,13 +240,28 @@ class JobRegistry { updateTask( jobId: string, taskId: string, - patch: { fraction?: number; detail?: string }, + patch: { + fraction?: number; + detail?: string; + phase?: "downloading" | "probing"; + probeEtaSeconds?: number; + }, ): void { const job = this.jobs.get(jobId); const task = job?.tasks?.find((t) => t.id === taskId); if (!task) return; if (patch.fraction !== undefined) task.fraction = patch.fraction; if (patch.detail !== undefined) task.detail = patch.detail; + if (patch.phase !== undefined && patch.phase !== task.phase) { + task.phase = patch.phase; + if (patch.phase === "probing") { + task.probeStartedAt = Date.now(); + task.probeEtaSeconds = patch.probeEtaSeconds; + } else { + task.probeStartedAt = undefined; + task.probeEtaSeconds = undefined; + } + } } removeTask(jobId: string, taskId: string): void { diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -16,6 +16,7 @@ import { shouldRequestSnapshot, } from "./snapshotScheduler"; import { writeJobMeta } from "./jobMeta"; +import type { JobSpec } from "./jobSpec"; // Mark a job's channel report dirty so the debounced scheduler regenerates the // snapshot — called both on each completed sub-operation and on the job's @@ -32,7 +33,10 @@ function requestSnapshotOnFinish(jobId: string, opts: CommonOpts): void { export type StreamActionResult = | { ok: true; jobId: string; stream: ReadableStream<string> } - | { ok: false; error: string }; + // `info: true` marks a non-error outcome that started no job (e.g. a + // re-derived bucket that's currently empty) so the UI can show it neutrally + // rather than as a red failure. + | { ok: false; error: string; info?: boolean }; type CommonOpts = { kind: string; @@ -40,6 +44,9 @@ type CommonOpts = { paths: Paths; channelSlug?: string; videoId?: string; + // When set, recorded on the job (and its meta sidecar) so the job can be + // bookmarked and re-launched later. See common/jobs/jobSpec.ts. + spec?: JobSpec; }; export type RunManagedCommandOpts = CommonOpts & { @@ -58,7 +65,12 @@ export type JobRunContext = { addTask: (task: JobTask) => void; updateTask: ( taskId: string, - patch: { fraction?: number; detail?: string }, + patch: { + fraction?: number; + detail?: string; + phase?: "downloading" | "probing"; + probeEtaSeconds?: number; + }, ) => void; removeTask: (taskId: string) => void; recordTaskDone: (durationMs: number) => void; @@ -83,6 +95,7 @@ function makeJob( paths: Paths, channelSlug?: string, videoId?: string, + spec?: JobSpec, ): { id: string; logPath: string; record: JobRecord } { const id = newJobId(); const logPath = path.join(paths.jobsDir, `${id}.log`); @@ -92,6 +105,7 @@ function makeJob( queueKey, channelSlug, videoId, + spec, status: "queued", queuedAt: Date.now(), logPath, @@ -117,6 +131,7 @@ export async function runManagedCommand( opts.paths, opts.channelSlug, opts.videoId, + opts.spec, ); const safe = makeSafeController<string>(); @@ -209,6 +224,7 @@ export async function runManagedFunction( opts.paths, opts.channelSlug, opts.videoId, + opts.spec, ); const safe = makeSafeController<string>(); diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts @@ -1,6 +1,6 @@ import type { JobTaskKind } from "./registry"; import type { JobRunContext } from "./streamCommand"; -import { parseDownloadProgress } from "./progressParsers"; +import { createDownloadProgressParser } from "./progressParsers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; // A handle for one in-flight sub-operation. `onLog` is a drop-in replacement @@ -59,6 +59,11 @@ export function makeTaskTracker( kind === "transcribe" && appId ? getTranscriptionApp(appId).makeProgressParser() : null; + // Download tasks use a stateful parser so audio-integrity probe markers + // (DLOM_PROBE) drive the probe phase and probe-aware ETA. Transcribe + // tasks without an appId (remote workers stream pre-parsed progress) fall + // through to it too — harmless for non-matching lines. + const downloadParser = transcribeParser ? null : createDownloadProgressParser(); let ended = false; const onLog = (line: string) => { forwardLog(line); @@ -68,7 +73,7 @@ export function makeTaskTracker( if (!part) continue; const update = transcribeParser ? transcribeParser.feed(part) - : parseDownloadProgress(part); + : downloadParser!.feed(part); if (update) ctx.updateTask(id, update); } }; diff --git a/common/lib/paths.ts b/common/lib/paths.ts @@ -24,6 +24,10 @@ export type Paths = { // disabled). Written by the Workers page "Set as default" button; applied by // the worker pool on first use. See common/jobs/workerDefaults.ts. workerDefaultsFile: string; + // Persisted job bookmarks: saved (kind, channel, params) specs an operator can + // re-launch with one button. Survives restarts, unlike the in-memory job + // registry. See common/jobs/bookmarks.ts. + bookmarksFile: string; lmdbPath: string; exportDir: string; // The dir Next.js serves at "/" (also holds checked-in static assets). For @@ -84,6 +88,7 @@ export function getPaths(): Paths { workerScratchDir: path.join(transcriptsDir, ".worker-scratch"), schedulerStateFile: path.join(transcriptsDir, ".scheduler", "state.json"), workerDefaultsFile: path.join(transcriptsDir, ".workers", "defaults.json"), + bookmarksFile: path.join(transcriptsDir, ".bookmarks", "bookmarks.json"), lmdbPath: path.join(transcriptsDir, "index.mdb"), exportDir, exportPublicDir, diff --git a/common/ytdlp/audioCheckedDownload.ts b/common/ytdlp/audioCheckedDownload.ts @@ -7,11 +7,15 @@ // in-progress `.part`, rolls back to a known-good snapshot on corruption, // and finally runs its own ffmpeg extract to the configured audio format. // -// Snapshot discipline: every periodic check captures `fstat(.part).size`, -// sends SIGSTOP, copies (reflink where supported), sends SIGCONT, then -// truncates the snapshot to the pre-SIGSTOP size to discard any kernel- -// buffered writes that landed during the snapshot. SIGCONT is always sent -// in a finally to avoid orphaning a suspended child. +// Snapshot discipline: every periodic check captures `fstat(.part).size` and +// sends SIGSTOP. In the default paused mode the child stays stopped across the +// integrity probe, so the probe reads the `.part` in place — no copy needed — +// and only a clean checkpoint copies `.part` → `.good` (reflink where +// supported), truncated to the pre-SIGSTOP size to discard any kernel-buffered +// writes that landed as SIGSTOP took effect. Legacy resume-during-probe mode +// instead copies a `.testing` snapshot up front (the child keeps writing during +// the probe) and promotes it to `.good` on success. SIGCONT is always sent in a +// finally to avoid orphaning a suspended child. import { constants as fsConstants } from "node:fs"; import { @@ -35,6 +39,7 @@ import type { Paths } from "../lib/paths"; import { transcodeAudio } from "../controller/transcode"; import type { StreamVerdict } from "./ffmpegStreamClassify"; import { probeAudioStream } from "./ffmpegStreamProbe"; +import { DLOM_PROBE_MARKER } from "../jobs/progressParsers"; const STDERR_TAIL_BYTES = 64 * 1024; const DEFAULT_SIZE_GATE_BYTES = 1024 * 1024; // 1 MiB minimum growth between checks @@ -75,6 +80,10 @@ export type CheckpointRecord = { bytes: number; verdict: StreamVerdict | null; action: CheckpointAction; + // Wall-clock duration of the ffmpeg integrity probe, in ms. Present only for + // checkpoints where a real probe ran (advance / rollback / restart / + // skip-copy-timeout); undefined for size-gate / missing-part skips. + durationMs?: number; }; export type AudioCheckOutcomeKind = @@ -358,9 +367,53 @@ function resumeChild( } } +// Reflink-copy `src` → `dst`, racing the copy against the copy timeout so we +// never hold SIGSTOP indefinitely, then truncate `dst` down to `preSize` to +// discard any bytes past the measured snapshot boundary (kernel-buffered writes +// that landed as SIGSTOP took effect). Returns "copy-timeout" if the copy +// outran the budget; the partial `dst` is then cleaned up best-effort. `dst` is +// always a copy of `src`, whose size only grows, so it is >= preSize and the +// truncate is always valid. +async function copyTruncated( + src: string, + dst: string, + preSize: number, + knobs: Knobs, +): Promise<"ok" | "copy-timeout"> { + const copyOp = (async () => { + await copyFile(src, dst, fsConstants.COPYFILE_FICLONE); + })(); + let timeoutId: ReturnType<typeof setTimeout> | null = null; + const timeout = new Promise<"timeout">((resolve) => { + timeoutId = setTimeout(() => resolve("timeout"), knobs.copyTimeoutMs); + }); + const winner = await Promise.race([ + copyOp.then(() => "copy" as const), + timeout, + ]); + if (timeoutId) clearTimeout(timeoutId); + if (winner === "timeout") { + // Best-effort clean up the partial copy. Don't await on copyOp completion; + // just leave it to settle in the background while we move on. + copyOp.catch(() => {}).then(() => rm(dst, { force: true }).catch(() => {})); + return "copy-timeout"; + } + const fh = await open(dst, "r+"); + try { + await fh.truncate(preSize); + } finally { + await fh.close(); + } + return "ok"; +} + async function snapshotPart( partFile: string, - testingFile: string, + // When null, the caller probes the held `.part` in place — no copy is taken. + // Only valid together with holdStopped (the child must stay SIGSTOPped across + // the probe, since the `.part` is read directly). When set, the snapshot is + // copied here (legacy resume-during-probe mode). + testingFile: string | null, pid: number | undefined, knobs: Knobs, onLog: (s: string) => void, @@ -398,45 +451,20 @@ async function snapshotPart( if (code === "ESRCH") return { ok: false, reason: "missing-part" }; throw err; } - // Re-stat post-stop in case more bytes flushed; we'll truncate the copy - // back to preSize after to discard those buffered writes. - let postSize = preSize; + // Re-stat post-stop purely to confirm the `.part` still exists; the child + // is stopped, so it shouldn't have vanished, but treat that as missing. try { - const st = await stat(partFile); - postSize = st.size; + await stat(partFile); } catch { - // Shouldn't happen — we're stopped — but treat as missing. return { ok: false, reason: "missing-part" }; } - // Copy with reflink where supported. Race against the copy timeout so we - // never hold SIGSTOP indefinitely. - const copyOp = (async () => { - await copyFile(partFile, testingFile, fsConstants.COPYFILE_FICLONE); - })(); - let timeoutId: ReturnType<typeof setTimeout> | null = null; - const timeout = new Promise<"timeout">((resolve) => { - timeoutId = setTimeout(() => resolve("timeout"), knobs.copyTimeoutMs); - }); - const winner = await Promise.race([ - copyOp.then(() => "copy" as const), - timeout, - ]); - if (timeoutId) clearTimeout(timeoutId); - if (winner === "timeout") { - // Best-effort clean up the partial snapshot. Don't await on copyOp - // completion; just leave it to settle in the background while we move - // on. - copyOp.catch(() => {}).then(() => rm(testingFile, { force: true }).catch(() => {})); - return { ok: false, reason: "copy-timeout" }; - } - // Truncate snapshot back to preSize. - if (postSize > preSize) { - const fh = await open(testingFile, "r+"); - try { - await fh.truncate(preSize); - } finally { - await fh.close(); - } + // Legacy resume-during-probe mode: take a reflink copy now, truncated back + // to preSize, because the child resumes and keeps writing during the probe. + // Default paused mode (testingFile null) skips this — the held `.part` is + // probed in place and only a clean checkpoint copies it into `.good`. + if (testingFile !== null) { + const copied = await copyTruncated(partFile, testingFile, preSize, knobs); + if (copied === "copy-timeout") return { ok: false, reason: "copy-timeout" }; } // Honour optional test-only debug pause so cancellation tests can wedge a // cancel into the SIGSTOP window deterministically. @@ -770,6 +798,12 @@ export async function runAudioCheckedYtdlp( let partFile: string | null = null; let lastCheckedBytes = 0; + // Most recent real ffmpeg probe duration (and the one before it), used to + // emit the estimated duration of the *next* probe in the DLOM_PROBE start + // marker so the UI's scanning bar has a denominator. The parser maintains + // its own smoothed trend for ETA inflation; this is just the display est. + let lastProbeMs = 0; + let prevProbeMs = 0; let watcherStop = false; // Set by checkpoint() as soon as a malformed verdict triggers a rollback // or restart. Persisted on a plain variable rather than a Promise so the @@ -797,11 +831,38 @@ export async function runAudioCheckedYtdlp( }); return; } - const testing = testingPath(partFile); + // A real probe pause begins here. Announce it so the UI flips the task to + // the "probing" phase, carrying the estimated duration (last probe plus + // its observed increase). finishProbe/skipProbe close the phase on every + // exit path below; durationMs from a real probe feeds the parser's trend. + const intervalSec = Math.round(knobs.intervalMs / 1000); + const probeStartedAt = Date.now(); + const estMs = + lastProbeMs > 0 ? lastProbeMs + Math.max(0, lastProbeMs - prevProbeMs) : 0; + opts.onLog(`${DLOM_PROBE_MARKER} status=start estMs=${Math.round(estMs)}\n`); + let probeDone = false; + const finishProbe = (reset: boolean): void => { + if (probeDone) return; + probeDone = true; + const durationMs = Date.now() - probeStartedAt; + prevProbeMs = lastProbeMs; + lastProbeMs = durationMs; + opts.onLog( + `${DLOM_PROBE_MARKER} status=done durationMs=${durationMs} intervalSec=${intervalSec} reset=${reset ? 1 : 0}\n`, + ); + }; + const skipProbe = (): void => { + if (probeDone) return; + probeDone = true; + opts.onLog(`${DLOM_PROBE_MARKER} status=done skip=1\n`); + }; // Default: keep the child SIGSTOPped across the probe below so it never // downloads bytes a malformed verdict would discard (429 minimisation). // resumeDuringProbe restores the legacy resume-immediately behavior. const holdStopped = !knobs.resumeDuringProbe; + // Paused mode probes the held `.part` in place (testing === null); legacy + // mode copies a `.testing` snapshot since the child keeps writing. + const testing = holdStopped ? null : testingPath(partFile); const snap = await snapshotPart( partFile, testing, @@ -821,6 +882,7 @@ export async function runAudioCheckedYtdlp( ? "skip-copy-timeout" : "skip-missing-part", }); + skipProbe(); return; } lastCheckedBytes = snap.bytes; @@ -829,38 +891,71 @@ export async function runAudioCheckedYtdlp( // intentionally kills it (SIGCONT + SIGTERM). `killed` gates that. let killed = false; try { - // Classify the snapshot. + // Classify the snapshot — or the held `.part` directly in paused mode. const probe = await probeAudioStream({ ffmpegBin: opts.paths.ffmpegBin, - file: testing, + file: testing ?? partFile, signal: opts.signal, onLog: opts.onLog, }); if (opts.signal.aborted) { - await rm(testing, { force: true }); + if (testing) await rm(testing, { force: true }); + skipProbe(); return; } if (probe.verdict === "clean" || probe.verdict === "partial") { - // Advance: rename testing → good, reset rollback counter. The + // Advance: record a known-good baseline, reset rollback counter. The // finally resumes the held child, so yt-dlp continues from where it // paused. const good = goodPath(partFile); - await rename(testing, good).catch(async () => { - // If rename failed (e.g. cross-device on weird setups), try - // copyFile+rm fallback. - await copyFile(testing, good).catch(() => {}); - await rm(testing, { force: true }); - }); + if (testing) { + // Legacy: promote the snapshot we already copied. + await rename(testing, good).catch(async () => { + // If rename failed (e.g. cross-device on weird setups), try + // copyFile+rm fallback. + await copyFile(testing, good).catch(() => {}); + await rm(testing, { force: true }); + }); + } else { + // Paused mode: the child is still SIGSTOPped, so the `.part` we just + // probed clean is stable — copy it into `.good` now (truncated to + // the probed size). This is the only copy in the paused path. + const copied = await copyTruncated( + partFile, + good, + snap.bytes, + knobs, + ); + if (copied === "copy-timeout") { + // Couldn't persist the baseline in time. The bytes are clean but + // we have no new snapshot, so don't advance — leave the rollback + // counter intact and try again next interval. + checkpoints.push({ + at: new Date().toISOString(), + bytes: snap.bytes, + verdict: probe.verdict, + action: "skip-copy-timeout", + durationMs: Date.now() - probeStartedAt, + }); + opts.onLog( + `Checkpoint clean at ${snap.bytes} bytes but timed out saving .good baseline; not advancing.\n`, + ); + finishProbe(false); + return; + } + } consecutiveRollbacks = 0; checkpoints.push({ at: new Date().toISOString(), bytes: snap.bytes, verdict: probe.verdict, action: "advance", + durationMs: Date.now() - probeStartedAt, }); opts.onLog( `Checkpoint OK at ${snap.bytes} bytes (verdict=${probe.verdict}).\n`, ); + finishProbe(false); return; } // Malformed. Decide rollback vs restart depending on whether we have @@ -873,12 +968,16 @@ export async function runAudioCheckedYtdlp( bytes: snap.bytes, verdict: probe.verdict, action: pendingDecision, + durationMs: Date.now() - probeStartedAt, }); opts.onLog( `Checkpoint MALFORMED at ${snap.bytes} bytes. ${haveGood ? "Rolling back to .good." : "No .good baseline; restarting from 0."}\n`, ); - // Drop the (failed) snapshot. - await rm(testing, { force: true }); + // reset=true: the rollback/restart shrinks the .part, so the next + // probe is cheaper — the parser must drop its duration trend. + finishProbe(true); + // Drop the (failed) snapshot, if legacy mode took one. + if (testing) await rm(testing, { force: true }); // Stop the watcher to avoid concurrent decisions. watcherStop = true; // Terminate yt-dlp, await exit, then mutate files. SIGCONT first so a @@ -903,6 +1002,9 @@ export async function runAudioCheckedYtdlp( await rm(partFile, { force: true }); } } finally { + // Safety net: if the probe threw before any explicit finish/skip, still + // close the probe phase so the UI bar doesn't stay stuck "probing". + skipProbe(); // Resume a child snapshotPart left stopped, unless we just killed it. // No-op in resumeDuringProbe mode (snapshotPart already resumed). if (holdStopped && !killed) resumeChild(child.pid, opts.onLog); diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,13 @@ # Changelog ## [Unreleased] +- **The Bookmarks panel is now a compact one-click run menu, with management moved to its own page.** The bookmarks block atop `/jobs` and `/jobs/active` used to render a heavy multi-row card per bookmark — inline rename, tags, created/last-run metadata, a confirm-gated delete, and a full streaming **Run again** button that popped a tall inline log box on every run — which buried the live job list below it. It's now a tight **wrap of quick-access buttons**, one per bookmark, labeled with the bookmark's name (the full kind/bucket/channel shows on hover). Clicking a button **launches the job fire-and-forget** — no inline shell, no expanding log — and the relaunched job simply appears in the live list below (it runs to completion server-side whether or not the page watches its stream); the button briefly reads *Launching… → Launched ✓*, an empty re-derived bucket still reads as a neutral *"Nothing to retry right now"* line, and a bookmark whose channel was deleted is disabled with a *channel missing* hint. The heavier controls — rename, delete, created/last-run metadata, and the full streaming **Run again** with its inline log — move to a new **`/jobs/bookmarks`** management page, reachable via the **Manage** link in the menu header. See `editor/app/jobs/components/BookmarksMenu.tsx`, `editor/app/jobs/components/BookmarkRunButton.tsx`, and `editor/app/jobs/bookmarks/page.tsx`. +- **ETAs are smarter in two edge cases, and audio-integrity probes now show their own progress bar.** Batch ETAs (downloads/transcripts on `/jobs/active` and the widget) now round the remaining work up to whole **parallel waves** instead of dividing straight through, so the tail no longer reads too low — e.g. *2 videos left across 4 workers* now estimates ~one full task, not "half a task." For **audio-checked downloads**, the per-task ETA used to come straight from yt-dlp, which only counts active download time and is blind to the periodic ffmpeg integrity **probe** pauses (each one transcodes the whole, growing `.part`, so later probes cost more). The download parser now measures each probe and **projects the remaining probe overhead** — extrapolating the per-probe duration as an arithmetic progression — onto yt-dlp's ETA, so a long audio-checked download no longer counts down faster than wall-clock. And while a probe runs (yt-dlp is paused, so the download bar would otherwise just **freeze**), the task now shows a distinct violet **"probing audio"** bar that fills against the estimated probe duration. See `editor/app/jobs/active/buildActiveJobs.ts`, `common/jobs/progressParsers.ts`, and `common/ytdlp/audioCheckedDownload.ts`. +- **Bookmark a job to re-run it with one click.** Any bookmarkable job now shows a **Bookmark** button — on the `/jobs` table rows, the job detail page, and the `/jobs/active` rows — and saved jobs appear in a **Bookmarks** panel atop both `/jobs` and `/jobs/active`, each with a **Run again** button (which streams the relaunched job's log) and **Delete**. So bookmarking a *whisper-all on HasanAbiVODs3* or a *retry partial downloads on cornbreadman* gives you a button that re-launches the same job on the same channel later. Re-running **re-derives the work from the channel's current state** rather than replaying a frozen list: bucket jobs (retry partial downloads, transcribe downloaded-no-transcript, missing-transcript-and-no-download) store only the snapshot bucket category, so "Run again" always acts on whatever's in that bucket *now* (and reports "Nothing to retry right now" when it's empty) — matching how *Transcribe missing* already re-scans. Every channel-scoped pipeline, whisper, transcode, and cleanup job kind is bookmarkable; ad-hoc checkbox selections and one-off URL imports are not (there's no stable set to re-derive). A bookmark captures the job's kind, channel, and its flags (queue, audio format, abort-on-error, etc.) as a small replay descriptor persisted on the job (and its `<id>.meta.json` sidecar, so even an archived job can be bookmarked) and saved to `transcripts/.bookmarks/bookmarks.json`. Each bookmark can be **renamed** (auto-named `kind · channel` by default — handy when two bookmarks differ only in their options), shows its **created / last-run** times, and **delete is confirm-gated**. An empty re-derived bucket reads as a **neutral notice** ("Nothing to retry right now") rather than a red error, and a bookmark whose channel was since deleted is **flagged and its Run-again disabled**. See `common/jobs/jobSpec.ts`, `common/jobs/bookmarks.ts`, and `editor/app/jobs/runJobSpec.ts`. +- **In-progress bars with no percentage now read clearly as "working, no number."** A task that hasn't reported a parseable percent yet (a download before yt-dlp's first percent, an engine before its header line) used to show a static 1/3-width bar in the same fill color as real progress — easy to misread as "stalled at ~33%." It's now a **full-width amber pulse**, visually distinct from the emerald/blue progress fill, so the indeterminate state is obvious. Applied consistently on the monitor widget, `/jobs/active`, and `/workers`. +- **The monitor widget's elements are now individually toggleable.** Each piece of the `/widget` view has its own GET param (and builder checkbox), so you can compose exactly the at-a-glance view you want and hand off the link. New flags: hide the per-job **batch/channel progress bar** (`jobbar=0`) to show only the individual per-task bars; optionally carry the batch progress (e.g. `5/10 · ~2m left`) as text in each job's **one-line heading** (`headtext=1`) so you keep the count/ETA after hiding the bar; hide all **time estimates** (`eta=0`, counts only); hide the **disk indicator** strip (`disk=0`); and collapse the Workers strip to **colored dots only** by hiding worker names (`wnames=0`). The combination `jobbar=0&headtext=1` yields the "only individual task bars under a text heading" layout. All are additive and back-compatible (existing `compact`/`titles`/`idle` links render identically) and omitted from shared URLs at their defaults. See `editor/app/widget/lib/config.ts`. +- **Audio-integrity check no longer copies the `.part` just to probe it.** In the default paused mode the download is held SIGSTOPped across each integrity probe, so the in-progress `.part` can't change underfoot — the probe now reads it in place instead of first reflink-copying it to a `.part.testing` snapshot. A copy is only taken when a checkpoint passes, to freeze the known-good `.part.good` baseline (and on malformed verdicts, no copy is taken at all). Legacy `resumeDuringProbe` mode is unchanged — it still snapshots up front since the child keeps writing during the probe. See `common/ytdlp/audioCheckedDownload.ts`. +- **Corrupt-source / undownloaded videos are no longer counted as failed transcriptions.** A video the downloader couldn't produce a real audio file for — most often a malformed source the audio-integrity check gave up on (`download-outcome.json` status `failed-corrupt-source`), or a download whose extract-to-`mp3` left no usable audio — used to get picked up by **Transcribe missing**, fail with "no audio file found", and be written to the channel's `failed-transcriptions` file *permanently* (so even a later clean re-download was skipped forever). These were download/source problems mislabeled as transcription failures, and they piled up. Now the transcribe pass **skips any video with no real audio file** instead of attempting it (the missing-audio case is classed `no-audio` and treated as a skip, not a failure), so it never lands in `failed-transcriptions`; it transcribes automatically once a good audio file exists. The pass also **auto-prunes** existing bogus entries — on each run, any id on the failed list whose dir currently has no audio is removed (it'll retry once re-downloaded). These videos are now surfaced *distinctly* rather than as failures: a new **corrupt-source** channel-report bucket feeds the **Download** stage summary ("N corrupt source (needs re-download)"), and in the video list they show a **corrupt_source** status (amber download dot, no red failure dot) and appear under the **No audio** filter. See `common/controller/whisperBatch.ts`, `common/controller/transcribeOne.ts`, and the new `corruptSource` snapshot bucket. - **Quick availability check: spot videos that vanished from a channel with one cheap yt-dlp call.** A channel's **Diagnostics → Availability checks** stage gained a **Quick check (flat playlist)** button that fetches the channel's *current* flat playlist (one `--flat-playlist --print url` pass, no per-video metadata) and diffs it against the videos already on disk. Any known video absent from the fresh listing is flagged **maybe-missing** — it was deleted, made private, or **unlisted** (unlisted videos legitimately drop out of channel listings, so a flag is a candidate, not a verdict). The set is persisted to `channels/<slug>/maybe-missing.json` and re-surfaced through the channel snapshot (intersected with on-disk dirs), so it survives report regens. This is the fast first pass that avoids running the full per-video `--dump-json` check across a whole channel just to notice the handful that disappeared. - **"Full-check unexpected": resolve maybe-missing videos without re-probing the whole channel.** Alongside the maybe-missing list, a **Full-check unexpected** button (with its own concurrency control) runs the full availability check on *only* the maybe-missing ids, skipping any already known permanently gone (`deleted` / `private` / `members_only`) and re-probing the rest (`unlisted` / `needs_auth` / `public` / `error` / never-checked, since their status may have changed). This efficiently resolves e.g. deleted-vs-unlisted for just the suspect videos. - **Per-video availability history.** Each video's `availability.json` now keeps a `history[]` timeline that appends an entry only when the observed availability *changes* — so a flip like `public → deleted` is preserved rather than overwritten. Entries are tagged by source (`check`, `backfill`, or `download` — the last captures a public video the uploader deleted, observed at download time, without disturbing the explicit-check fields). The video detail page renders this as an **Availability history** card (reverse-chronological, "No availability changes recorded" when empty). diff --git a/editor/app/channels/[slug]/components/RetryBucketControl.tsx b/editor/app/channels/[slug]/components/RetryBucketControl.tsx @@ -3,6 +3,7 @@ import { useState } from "react"; import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; import { HANDLING_VALUES } from "yt-dlp-transcript-common/lib/channelConfig"; +import type { ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; import { QueueControl } from "../../../components/QueueControl"; import { cancelJobAction } from "../../../jobs/actions"; import { retryBucketAction } from "../pipelineActions"; @@ -13,6 +14,9 @@ type Props = { actionLabel: string; defaultQueueKey: string; existingQueues: string[]; + // The snapshot bucket these ids represent. When set, the launched retry job + // is bookmarkable and re-runs against the current bucket members. + bucketKey?: ReplayBucket; }; export function RetryBucketControl({ @@ -21,6 +25,7 @@ export function RetryBucketControl({ actionLabel, defaultQueueKey, existingQueues, + bucketKey, }: Props) { const [queue, setQueue] = useState(defaultQueueKey); const [handlingOverride, setHandlingOverride] = useState(""); @@ -41,6 +46,7 @@ export function RetryBucketControl({ queue, abortOnError, handlingOverride || undefined, + bucketKey, ) } cancelAction={cancelJobAction} diff --git a/editor/app/channels/[slug]/components/VideoListPane.tsx b/editor/app/channels/[slug]/components/VideoListPane.tsx @@ -544,11 +544,12 @@ export function VideoListPane({ } function StatusGlyphs({ row }: { row: VideoRow }) { - const dlColor = row.partial - ? "bg-amber-500" - : row.downloaded - ? "bg-emerald-500" - : "bg-zinc-300 dark:bg-zinc-700"; + const dlColor = + row.partial || row.corruptSource + ? "bg-amber-500" + : row.downloaded + ? "bg-emerald-500" + : "bg-zinc-300 dark:bg-zinc-700"; const tcColor = row.transcoded === null ? "bg-transparent border border-zinc-300 dark:border-zinc-700" @@ -567,7 +568,10 @@ function StatusGlyphs({ row }: { row: VideoRow }) { aria-label={`status ${row.status}`} className="inline-flex items-center gap-1 flex-shrink-0" > - <span title="downloaded" className={`w-2 h-2 rounded-full ${dlColor}`} /> + <span + title={row.corruptSource ? "corrupt source (needs re-download)" : "downloaded"} + className={`w-2 h-2 rounded-full ${dlColor}`} + /> <span title="transcoded" className={`w-2 h-2 rounded-full ${tcColor}`} /> <span title="transcribed" className={`w-2 h-2 rounded-full ${trColor}`} /> {failureColor && ( diff --git a/editor/app/channels/[slug]/components/stages/DownloadStage.tsx b/editor/app/channels/[slug]/components/stages/DownloadStage.tsx @@ -365,6 +365,7 @@ function NoTranscriptList({ actionLabel="missing transcript and no download" defaultQueueKey={defaultQueueKey} existingQueues={existingQueues} + bucketKey="noTranscript" /> </div> ); @@ -409,6 +410,7 @@ function PartialDownloadsList({ actionLabel="resume partial downloads" defaultQueueKey={defaultQueueKey} existingQueues={existingQueues} + bucketKey="partialDownloads" /> </div> ); diff --git a/editor/app/channels/[slug]/components/stages/TranscribeStage.tsx b/editor/app/channels/[slug]/components/stages/TranscribeStage.tsx @@ -106,6 +106,7 @@ function BucketTranscribeSection({ queue, audioFormat === "any" ? undefined : audioFormat, audioFormat !== "any" && strictFormat, + "downloadedNoTranscript", ) } cancelAction={cancelJobAction} diff --git a/editor/app/channels/[slug]/lib/stageStatus.ts b/editor/app/channels/[slug]/lib/stageStatus.ts @@ -24,6 +24,7 @@ export function normalizeBuckets( missingFromArchive: raw?.missingFromArchive ?? [], duplicateDirs: raw?.duplicateDirs ?? [], partialDownloads: raw?.partialDownloads ?? [], + corruptSource: raw?.corruptSource ?? [], nonStandardVtt: raw?.nonStandardVtt ?? [], skippedByFilter: raw?.skippedByFilter ?? [], }; @@ -194,6 +195,15 @@ export function computeStageStatuses( ), ); } + if (buckets.corruptSource.length > 0) { + downloadParts.push( + pluralize( + buckets.corruptSource.length, + "corrupt source (needs re-download)", + "corrupt sources (need re-download)", + ), + ); + } const download: StageStatus = { id: "download", title: "Download", diff --git a/editor/app/channels/[slug]/lib/videoRows.ts b/editor/app/channels/[slug]/lib/videoRows.ts @@ -5,6 +5,7 @@ export type VideoRowStatus = | "untranscribable" | "partial_download" | "not_downloaded" + | "corrupt_source" | "failed"; export type VideoRow = { @@ -15,6 +16,10 @@ export type VideoRow = { transcoded: boolean | null; untranscribable: boolean; partial: boolean; + // The audio-check pipeline gave up on a malformed source (download-outcome + // "failed-corrupt-source"), leaving no real audio. A download/source problem — + // surfaced distinctly so it isn't conflated with a failed transcription. + corruptSource: boolean; failedTranscription: boolean; failedTranscoding: boolean; // Has at least one finalized audio file not in the channel's target format — @@ -55,7 +60,10 @@ export function filterRows(rows: VideoRow[], filter: VideoFilter): VideoRow[] { return rows.filter((r) => r.failedTranscription || r.failedTranscoding); case "no_audio": return rows.filter( - (r) => r.status === "no_audio" || r.status === "not_downloaded", + (r) => + r.status === "no_audio" || + r.status === "not_downloaded" || + r.status === "corrupt_source", ); case "downloaded_no_transcript": return rows.filter((r) => r.status === "downloaded_no_transcript"); diff --git a/editor/app/channels/[slug]/lib/videoRowsServer.ts b/editor/app/channels/[slug]/lib/videoRowsServer.ts @@ -40,6 +40,7 @@ export function computeVideoRows(input: ComputeRowsInput): VideoRow[] { const multipleAudioFormats = new Set(buckets.multipleAudioFormats); const untranscribable = new Set(buckets.untranscribable); const partial = new Set(buckets.partialDownloads); + const corruptSourceSet = new Set(buckets.corruptSource); const failedTranscription = new Set(input.failedTranscriptionIds); const failedTranscoding = new Set(input.failedTranscodingIds); @@ -57,7 +58,11 @@ export function computeVideoRows(input: ComputeRowsInput): VideoRow[] { const isUndownloaded = undownloaded.has(id) && !isOnDisk; const isUntranscribable = untranscribable.has(id); const isPartial = partial.has(id); - const isFailedT = failedTranscription.has(id); + const isCorruptSource = corruptSourceSet.has(id); + // A corrupt-source video is a download/source problem, not a failed + // transcription — don't let a stale failed-transcriptions entry (pruned on + // the next transcribe pass) mark it failed here. + const isFailedT = failedTranscription.has(id) && !isCorruptSource; const isFailedX = failedTranscoding.has(id); const inNoTranscript = noTranscript.has(id); const inDownloadedNoTranscript = downloadedNoTranscript.has(id); @@ -70,6 +75,7 @@ export function computeVideoRows(input: ComputeRowsInput): VideoRow[] { let status: VideoRowStatus; if (isFailedT || isFailedX) status = "failed"; + else if (isCorruptSource) status = "corrupt_source"; else if (isPartial) status = "partial_download"; else if (isUndownloaded) status = "not_downloaded"; else if (isUntranscribable) status = "untranscribable"; @@ -84,6 +90,7 @@ export function computeVideoRows(input: ComputeRowsInput): VideoRow[] { transcoded, untranscribable: isUntranscribable, partial: isPartial, + corruptSource: isCorruptSource, failedTranscription: isFailedT, failedTranscoding: isFailedX, wrongFormatAudio: diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts @@ -24,6 +24,7 @@ import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; import type { Paths } from "yt-dlp-transcript-common/lib/paths"; @@ -63,6 +64,8 @@ async function runPipelineAction( shardIndex?: number; bucketIds?: ReadonlyArray<string>; handlingOverride?: ChannelHandling; + // Replay descriptor, forwarded onto the job record so it can be bookmarked. + spec?: JobSpec; }, ): Promise<StreamActionResult> { const paths = getPaths(); @@ -84,6 +87,7 @@ async function runPipelineAction( queueKey: resolveQueueKey(downloadQueueKey(channelConfig), queueKey), paths, channelSlug: slug, + spec: options?.spec, fn: async (onLog, signal, setProgress, ctx) => { // Baseline (channel-wide downloadCount at run start) handed to runYtdlp so // a sharded download can refine the progress total to its own slice once @@ -148,7 +152,9 @@ export async function storePlaylistAction( slug: string, queueKey?: string, ): Promise<StreamActionResult> { - return runPipelineAction(slug, "store-playlist", "store-playlist", queueKey); + return runPipelineAction(slug, "store-playlist", "store-playlist", queueKey, { + spec: { kind: "store-playlist", slug, params: { queueKey } }, + }); } export async function downloadAction( @@ -161,7 +167,14 @@ export async function downloadAction( "download-from-playlist", "download-from-playlist", queueKey, - { abortOnError }, + { + abortOnError, + spec: { + kind: "download-from-playlist", + slug, + params: { queueKey, abortOnError }, + }, + }, ); } @@ -178,7 +191,17 @@ export async function downloadMissingAction( "download-missing", "download-missing", queueKey, - { ignoreArchive, abortOnError, shardTotal, shardIndex }, + { + ignoreArchive, + abortOnError, + shardTotal, + shardIndex, + spec: { + kind: "download-missing", + slug, + params: { queueKey, ignoreArchive, abortOnError, shardTotal, shardIndex }, + }, + }, ); } @@ -186,7 +209,9 @@ export async function syncAction( slug: string, queueKey?: string, ): Promise<StreamActionResult> { - return runPipelineAction(slug, "sync", "sync", queueKey); + return runPipelineAction(slug, "sync", "sync", queueKey, { + spec: { kind: "sync", slug, params: { queueKey } }, + }); } export async function downloadMissingSubsAction( @@ -199,7 +224,14 @@ export async function downloadMissingSubsAction( "download-missing-subs", "download-missing-subs", queueKey, - { abortOnError }, + { + abortOnError, + spec: { + kind: "download-missing-subs", + slug, + params: { queueKey, abortOnError }, + }, + }, ); } @@ -209,6 +241,11 @@ export async function retryBucketAction( queueKey?: string, abortOnError?: boolean, handlingOverride?: string, + // The snapshot bucket these ids came from. Passed by the named bucket controls + // (partial downloads / no-transcript) so the job can be bookmarked and re-run + // against the CURRENT bucket; omitted by ad-hoc checkbox selections, which are + // therefore not bookmarkable. + bucketKey?: ReplayBucket, ): Promise<StreamActionResult> { if (!Array.isArray(bucketIds) || bucketIds.length === 0) { return { ok: false, error: "No video IDs supplied for retry." }; @@ -223,12 +260,20 @@ export async function retryBucketAction( } handling = handlingOverride as ChannelHandling; } + const spec: JobSpec | undefined = bucketKey + ? { + kind: "retry-bucket", + slug, + bucket: bucketKey, + params: { queueKey, abortOnError, handlingOverride }, + } + : undefined; return runPipelineAction( slug, "retry-bucket", "retry-bucket", queueKey, - { bucketIds, handlingOverride: handling, abortOnError }, + { bucketIds, handlingOverride: handling, abortOnError, spec }, ); } diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts @@ -32,6 +32,7 @@ import { runManagedFunction, type StreamActionResult, } from "yt-dlp-transcript-common/jobs/streamCommand"; +import type { ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; function sanitizeAudioFormat(f: AudioFormat | undefined): AudioFormat | undefined { @@ -55,6 +56,18 @@ export async function transcribeMissingAction( queueKey: resolveQueueKey(TRANSCRIPTION_QUEUE, queueKey), paths, channelSlug: slug, + spec: { + kind: "whisper-all", + slug, + params: { + queueKey, + reverse, + audioFormat, + strictAudioFormat, + shardTotal, + shardIndex, + }, + }, fn: async (onLog, signal, setProgress, ctx) => { const stat = await readChannelStat(paths, slug); if (stat) { @@ -93,6 +106,10 @@ export async function transcribeBucketAction( queueKey?: string, audioFormat?: AudioFormat, strictAudioFormat?: boolean, + // The snapshot bucket these ids came from (downloadedNoTranscript). Passed by + // the named bucket control so the job is bookmarkable; omitted by ad-hoc + // selections. + bucketKey?: ReplayBucket, ): Promise<StreamActionResult> { const paths = getPaths(); const fmt = sanitizeAudioFormat(audioFormat); @@ -105,6 +122,14 @@ export async function transcribeBucketAction( queueKey: resolveQueueKey(TRANSCRIPTION_QUEUE, queueKey), paths, channelSlug: slug, + spec: bucketKey + ? { + kind: "whisper-bucket-downloaded-no-transcript", + slug, + bucket: bucketKey, + params: { queueKey, audioFormat, strictAudioFormat }, + } + : undefined, fn: async (onLog, signal, setProgress, ctx) => { const stat = await readChannelStat(paths, slug); if (stat) { @@ -144,6 +169,7 @@ export async function clearFailedTranscriptionsAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "clear-failed-transcriptions", slug, params: { queueKey } }, fn: async (onLog) => { const { cleared } = await clearFailedTranscriptions(paths, slug); onLog( @@ -168,6 +194,11 @@ export async function transcodeFailuresAction( queueKey: resolveQueueKey(TRANSCRIPTION_QUEUE, queueKey), paths, channelSlug: slug, + spec: { + kind: "transcode-failures", + slug, + params: { queueKey, targetFormat }, + }, fn: async (onLog, signal) => { const result = await runTranscodeFailures({ channelSlug: slug, @@ -194,6 +225,7 @@ export async function transcodeUntranscodedAction( queueKey: resolveQueueKey(TRANSCRIPTION_QUEUE, queueKey), paths, channelSlug: slug, + spec: { kind: "transcode-untranscoded", slug, params: { queueKey } }, fn: async (onLog, signal) => { const result = await runTranscodeUntranscoded({ channelSlug: slug, @@ -219,6 +251,7 @@ export async function removeFailedTranscodingsAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "remove-failed-transcodings", slug, params: { queueKey } }, fn: async (onLog, signal) => { const result = await removeFailedTranscodings({ channelSlug: slug, @@ -244,6 +277,7 @@ export async function clearFailedTranscodingsAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "clear-failed-transcodings", slug, params: { queueKey } }, fn: async (onLog) => { const ids = await loadFailedTranscodings(paths, slug); await clearFailedTranscodings(paths, slug); @@ -265,6 +299,7 @@ export async function cleanExtraAudioFormatsAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "clean-extra-audio-formats", slug, params: { queueKey } }, fn: async (onLog, signal) => { const result = await cleanExtraAudioFormats({ channelSlug: slug, @@ -293,6 +328,7 @@ export async function removeWrongFormatAudioAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "remove-wrong-format-audio", slug, params: { queueKey } }, fn: async (onLog, signal) => { const result = await removeWrongFormatAudio({ channelSlug: slug, @@ -318,6 +354,7 @@ export async function cleanAudioAction( queueKey: resolveQueueKey(channelQueueKey(slug), queueKey), paths, channelSlug: slug, + spec: { kind: "clean-audio-transcribed", slug, params: { queueKey } }, fn: async (onLog, signal) => { await cleanAudioFromTranscribed({ channelSlug: slug, diff --git a/editor/app/jobs/[id]/page.tsx b/editor/app/jobs/[id]/page.tsx @@ -5,6 +5,7 @@ import { notFound } from "next/navigation"; import { listAllJobs } from "yt-dlp-transcript-common/jobs/listJobs"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { JobLogTail } from "./components/JobLogTail"; +import { BookmarkJobButton } from "../components/BookmarkJobButton"; export const dynamic = "force-dynamic"; @@ -49,7 +50,10 @@ export default async function JobDetailPage({ <span>/</span> <span className="font-mono">{id}</span> </div> - <h1 className="text-2xl font-semibold">{job.kind ?? "job"}</h1> + <div className="flex items-center justify-between gap-3"> + <h1 className="text-2xl font-semibold">{job.kind ?? "job"}</h1> + {job.bookmarkable && <BookmarkJobButton jobId={id} />} + </div> {job.channelSlug && ( <p className="text-sm"> Channel:{" "} diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -48,7 +48,11 @@ function computeEtaSeconds( 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 - return (remaining / concurrency) * avgProcMs / 1000; + // 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( @@ -132,10 +136,14 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { 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, + bookmarkable: Boolean(j.spec), })); const channels = channelSlugs.map((slug) => ({ diff --git a/editor/app/jobs/active/page.tsx b/editor/app/jobs/active/page.tsx @@ -1,18 +1,24 @@ import type { Metadata } from "next"; import { buildActiveJobsPayload } from "./buildActiveJobs"; import { ActiveJobsLive } from "../components/ActiveJobsLive"; +import { BookmarksMenu } from "../components/BookmarksMenu"; +import { loadBookmarksView } from "../loadBookmarks"; export const dynamic = "force-dynamic"; export const metadata: Metadata = { title: "Active jobs" }; export default async function ActiveJobsPage() { - const initial = await buildActiveJobsPayload(); + const [initial, { bookmarks, missingSlugs }] = await Promise.all([ + buildActiveJobsPayload(), + loadBookmarksView(), + ]); return ( <div className="flex flex-col gap-4"> <div className="flex items-center justify-between"> <h1 className="text-2xl font-semibold">Active jobs</h1> </div> + <BookmarksMenu bookmarks={bookmarks} missingSlugs={missingSlugs} /> <ActiveJobsLive initial={initial} /> </div> ); diff --git a/editor/app/jobs/bookmarkActions.ts b/editor/app/jobs/bookmarkActions.ts @@ -0,0 +1,93 @@ +"use server"; + +import { revalidatePath } from "next/cache"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; +import { readJobMeta } from "yt-dlp-transcript-common/jobs/jobMeta"; +import { + addBookmark, + defaultBookmarkName, + readBookmarks, + removeBookmark, + updateBookmark, +} from "yt-dlp-transcript-common/jobs/bookmarks"; +import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import { runJobSpec } from "./runJobSpec"; + +export type BookmarkActionResult = { ok: true } | { ok: false; error: string }; + +function refreshJobsViews(): void { + revalidatePath("/jobs"); + revalidatePath("/jobs/active"); +} + +// Bookmark a job by copying its replay descriptor (from the live registry, or +// the on-disk sidecar for an archived job) into the bookmarks store. Jobs with +// no spec — ad-hoc selections, single imports — can't be bookmarked. +export async function bookmarkJobAction( + jobId: string, +): Promise<BookmarkActionResult> { + const paths = getPaths(); + let spec = getRegistry().get(jobId)?.spec; + if (!spec) { + const meta = await readJobMeta(paths, jobId); + spec = meta?.spec; + } + if (!spec) return { ok: false, error: "This job can't be bookmarked." }; + await addBookmark(paths, spec); + refreshJobsViews(); + return { ok: true }; +} + +export async function deleteBookmarkAction( + id: string, +): Promise<BookmarkActionResult> { + await removeBookmark(getPaths(), id); + refreshJobsViews(); + return { ok: true }; +} + +// Rename a bookmark. An empty/blank name resets it to the auto-generated default +// (kind · channel) so a bookmark always has a usable label. +export async function renameBookmarkAction( + id: string, + name: string, +): Promise<BookmarkActionResult> { + const paths = getPaths(); + const trimmed = name.trim().slice(0, 200); + if (trimmed) { + await updateBookmark(paths, id, { name: trimmed }); + } else { + const bookmark = (await readBookmarks(paths)).find((b) => b.id === id); + if (bookmark) { + await updateBookmark(paths, id, { + name: defaultBookmarkName(bookmark.spec), + }); + } + } + refreshJobsViews(); + return { ok: true }; +} + +// Re-launch a bookmarked job. Returns a StreamActionResult so the caller can +// stream the new job's log via StreamActionLog, exactly like the original +// stage-control buttons. +export async function runBookmarkAction( + id: string, +): Promise<StreamActionResult> { + const paths = getPaths(); + const bookmark = (await readBookmarks(paths)).find((b) => b.id === id); + if (!bookmark) return { ok: false, error: "Bookmark not found." }; + // Guard against a bookmark whose channel was deleted: re-running would scan a + // non-existent dir and do nothing useful. + if (!(await readChannelConfig(paths, bookmark.spec.slug))) { + return { + ok: false, + error: `Channel "${bookmark.spec.slug}" no longer exists.`, + }; + } + const result = await runJobSpec(bookmark.spec); + if (result.ok) await updateBookmark(paths, id, { lastRunAt: Date.now() }); + return result; +} diff --git a/editor/app/jobs/bookmarks/page.tsx b/editor/app/jobs/bookmarks/page.tsx @@ -0,0 +1,37 @@ +import type { Metadata } from "next"; +import Link from "next/link"; +import { BookmarksList } from "../components/BookmarksList"; +import { loadBookmarksView } from "../loadBookmarks"; + +export const dynamic = "force-dynamic"; + +export const metadata: Metadata = { title: "Bookmarks" }; + +// Full bookmark management: rename, delete, created / last-run metadata, and the +// streaming "Run again" button (with its inline log). The compact menu on /jobs +// and /jobs/active links here via "Manage". +export default async function BookmarksPage() { + const { bookmarks, missingSlugs } = await loadBookmarksView(); + return ( + <div className="flex flex-col gap-4"> + <div className="flex items-center gap-2 text-sm text-zinc-500"> + <Link + href="/jobs" + className="underline hover:text-zinc-900 dark:hover:text-zinc-100" + > + Jobs + </Link> + <span>/</span> + <span>Bookmarks</span> + </div> + <h1 className="text-2xl font-semibold">Bookmarks</h1> + {bookmarks.length === 0 ? ( + <p className="text-sm text-zinc-500 border border-dashed border-zinc-300 dark:border-zinc-700 rounded p-4"> + No bookmarks yet. + </p> + ) : ( + <BookmarksList bookmarks={bookmarks} missingSlugs={missingSlugs} /> + )} + </div> + ); +} diff --git a/editor/app/jobs/components/BookmarkJobButton.tsx b/editor/app/jobs/components/BookmarkJobButton.tsx @@ -0,0 +1,50 @@ +"use client"; + +import { useState, useTransition } from "react"; +import { useRouter } from "next/navigation"; +import { bookmarkJobAction } from "../bookmarkActions"; + +// Small inline button shown on a job (jobs table row, active-jobs row, job +// detail) when it carries a replay descriptor. Saves a bookmark, then shows a +// transient "Bookmarked ✓" state. Errors are surfaced as the button title. +export function BookmarkJobButton({ + jobId, + className, +}: { + jobId: string; + className?: string; +}) { + const router = useRouter(); + const [pending, startTransition] = useTransition(); + const [done, setDone] = useState(false); + const [error, setError] = useState<string | null>(null); + + function onClick() { + setError(null); + startTransition(async () => { + const r = await bookmarkJobAction(jobId); + if (r.ok) { + setDone(true); + router.refresh(); + } else { + setError(r.error); + } + }); + } + + return ( + <button + type="button" + onClick={onClick} + disabled={pending || done} + aria-label={`bookmark job ${jobId}`} + title={error ?? undefined} + className={ + className ?? + "px-2 py-1 rounded-md bg-zinc-200 dark:bg-zinc-800 text-zinc-900 dark:text-zinc-100 text-xs font-medium hover:opacity-90 disabled:opacity-50" + } + > + {done ? "Bookmarked ✓" : pending ? "Bookmarking…" : "Bookmark"} + </button> + ); +} diff --git a/editor/app/jobs/components/BookmarkRunButton.tsx b/editor/app/jobs/components/BookmarkRunButton.tsx @@ -0,0 +1,94 @@ +"use client"; + +import { useState, useTransition } from "react"; +import { useRouter } from "next/navigation"; +import type { JobBookmark } from "yt-dlp-transcript-common/jobs/bookmarks"; +import { runBookmarkAction } from "../bookmarkActions"; + +// One quick-access button in the compact bookmarks menu (BookmarksMenu). Unlike +// the full management UI on /jobs/bookmarks, this launches the job +// fire-and-forget: it never mounts StreamActionLog / an inline shell log. The +// re-launched job runs to completion server-side (streamCommand keeps the child +// and on-disk log alive even when no consumer reads the stream) and shows up in +// the live job list below. Errors and the neutral "nothing to do" (info) case +// are surfaced as a brief status line + the button title — no expanding log box. +export function BookmarkRunButton({ + bookmark, + channelMissing, +}: { + bookmark: JobBookmark; + channelMissing: boolean; +}) { + const router = useRouter(); + const [pending, startTransition] = useTransition(); + const [launched, setLaunched] = useState(false); + const [status, setStatus] = useState< + { kind: "error" | "info"; text: string } | null + >(null); + const { spec } = bookmark; + + const detail = [ + spec.kind, + spec.bucket ? `bucket ${spec.bucket}` : null, + spec.slug, + ] + .filter(Boolean) + .join(" · "); + + function onClick() { + if (channelMissing) return; + setStatus(null); + setLaunched(false); + startTransition(async () => { + const r = await runBookmarkAction(bookmark.id); + if (r.ok) { + // Release the stream reader; the job keeps running server-side. The + // refresh surfaces it in the live job list below. + void r.stream.cancel().catch(() => {}); + setLaunched(true); + router.refresh(); + } else { + setStatus({ kind: r.info ? "info" : "error", text: r.error }); + } + }); + } + + const label = channelMissing + ? bookmark.name + : launched + ? "Launched ✓" + : pending + ? "Launching…" + : bookmark.name; + + return ( + <div className="flex flex-col gap-1 max-w-full"> + <button + type="button" + onClick={onClick} + disabled={pending || channelMissing} + aria-label={`run bookmark ${bookmark.name}`} + title={status?.text ?? detail} + className="max-w-[16rem] truncate px-3 py-1.5 rounded-md bg-zinc-900 dark:bg-zinc-100 text-zinc-100 dark:text-zinc-900 text-xs font-medium hover:opacity-90 disabled:opacity-50 disabled:cursor-not-allowed" + > + {label} + </button> + {channelMissing ? ( + <span className="text-[11px] text-amber-700 dark:text-amber-300"> + channel missing + </span> + ) : status ? ( + <span + role={status.kind === "error" ? "alert" : "status"} + className={ + status.kind === "error" + ? "text-[11px] text-red-700 dark:text-red-300" + : "text-[11px] text-zinc-500" + } + > + {status.text} + </span> + ) : null} + </div> + ); +} diff --git a/editor/app/jobs/components/BookmarksList.tsx b/editor/app/jobs/components/BookmarksList.tsx @@ -0,0 +1,204 @@ +"use client"; + +import { useState, useTransition } from "react"; +import Link from "next/link"; +import { useRouter } from "next/navigation"; +import { StreamActionLog } from "yt-dlp-transcript-common/components/StreamActionLog"; +import type { JobBookmark } from "yt-dlp-transcript-common/jobs/bookmarks"; +import { cancelJobAction } from "../actions"; +import { + deleteBookmarkAction, + renameBookmarkAction, + runBookmarkAction, +} from "../bookmarkActions"; + +// The saved-jobs panel rendered on /jobs and /jobs/active. Each bookmark has a +// "Run again" button (streams the re-launched job's log), an inline rename, a +// confirm-gated delete, and created / last-run metadata. Bookmarks whose channel +// was deleted are flagged and can't be re-run. +export function BookmarksList({ + bookmarks, + missingSlugs = [], +}: { + bookmarks: JobBookmark[]; + missingSlugs?: string[]; +}) { + if (bookmarks.length === 0) return null; + const missing = new Set(missingSlugs); + return ( + <section + aria-label="Bookmarked jobs" + className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-zinc-50 dark:bg-zinc-900/40" + > + <h2 className="text-sm font-medium text-zinc-700 dark:text-zinc-300"> + Bookmarks + </h2> + <ul className="flex flex-col gap-3"> + {bookmarks.map((b) => ( + <BookmarkRow + key={b.id} + bookmark={b} + channelMissing={missing.has(b.spec.slug)} + /> + ))} + </ul> + </section> + ); +} + +function fmtDate(ms?: number): string { + if (!ms) return "—"; + return new Date(ms).toLocaleString(); +} + +function BookmarkRow({ + bookmark, + channelMissing, +}: { + bookmark: JobBookmark; + channelMissing: boolean; +}) { + const router = useRouter(); + const [pending, startTransition] = useTransition(); + const [confirmingDelete, setConfirmingDelete] = useState(false); + const [editing, setEditing] = useState(false); + const [draftName, setDraftName] = useState(bookmark.name); + const { spec } = bookmark; + + function onDelete() { + startTransition(async () => { + await deleteBookmarkAction(bookmark.id); + router.refresh(); + }); + } + + function onSaveName() { + setEditing(false); + startTransition(async () => { + await renameBookmarkAction(bookmark.id, draftName); + router.refresh(); + }); + } + + return ( + <li + aria-label={`bookmark ${bookmark.name}`} + className="flex flex-col gap-2 border-t border-zinc-200 dark:border-zinc-800 pt-3 first:border-t-0 first:pt-0" + > + <div className="flex flex-wrap items-center gap-2 text-sm"> + {editing ? ( + <> + <input + type="text" + value={draftName} + onChange={(e) => setDraftName(e.target.value)} + aria-label={`new name for ${bookmark.name}`} + className="rounded border border-zinc-300 dark:border-zinc-700 bg-white dark:bg-zinc-900 px-2 py-0.5 text-sm" + /> + <button + type="button" + onClick={onSaveName} + className="px-2 py-1 rounded-md bg-zinc-900 dark:bg-zinc-100 text-zinc-100 dark:text-zinc-900 text-xs font-medium hover:opacity-90" + > + Save + </button> + <button + type="button" + onClick={() => { + setEditing(false); + setDraftName(bookmark.name); + }} + className="px-2 py-1 rounded-md border border-zinc-300 dark:border-zinc-700 text-xs font-medium hover:bg-zinc-100 dark:hover:bg-zinc-800" + > + Cancel + </button> + </> + ) : ( + <> + <span className="font-medium">{bookmark.name}</span> + <button + type="button" + onClick={() => { + setDraftName(bookmark.name); + setEditing(true); + }} + aria-label={`rename bookmark ${bookmark.name}`} + className="px-2 py-0.5 rounded border border-zinc-300 dark:border-zinc-700 text-xs text-zinc-600 dark:text-zinc-400 hover:bg-zinc-100 dark:hover:bg-zinc-800" + > + Rename + </button> + </> + )} + <span className="font-mono text-xs text-zinc-500">{spec.kind}</span> + {spec.bucket && ( + <span className="font-mono text-xs text-zinc-500"> + bucket {spec.bucket} + </span> + )} + <Link + href={`/channels/${spec.slug}`} + className="font-mono text-xs underline hover:text-zinc-900 dark:hover:text-zinc-100" + > + {spec.slug} + </Link> + {channelMissing && ( + <span + aria-label={`channel missing for ${bookmark.name}`} + className="text-xs px-2 py-0.5 rounded bg-amber-100 text-amber-800 dark:bg-amber-900 dark:text-amber-200" + > + channel missing + </span> + )} + <div className="ml-auto flex items-center gap-2"> + {confirmingDelete ? ( + <> + <button + type="button" + onClick={onDelete} + disabled={pending} + aria-label={`confirm delete bookmark ${bookmark.name}`} + className="px-2 py-1 rounded-md bg-red-600 text-white text-xs font-medium hover:bg-red-700 disabled:opacity-50" + > + {pending ? "Deleting…" : "Confirm"} + </button> + <button + type="button" + onClick={() => setConfirmingDelete(false)} + className="px-2 py-1 rounded-md border border-zinc-300 dark:border-zinc-700 text-xs font-medium hover:bg-zinc-100 dark:hover:bg-zinc-800" + > + Cancel + </button> + </> + ) : ( + <button + type="button" + onClick={() => setConfirmingDelete(true)} + aria-label={`delete bookmark ${bookmark.name}`} + className="px-2 py-1 rounded-md border border-red-300 dark:border-red-800 text-xs font-medium text-red-700 dark:text-red-300 hover:bg-red-50 dark:hover:bg-red-950" + > + Delete + </button> + )} + </div> + </div> + <p className="text-xs text-zinc-500" suppressHydrationWarning> + Created {fmtDate(bookmark.createdAt)} · Last run{" "} + {bookmark.lastRunAt ? fmtDate(bookmark.lastRunAt) : "never"} + </p> + {channelMissing ? ( + <p className="text-xs text-amber-700 dark:text-amber-300"> + This channel no longer exists — delete the bookmark or recreate the + channel to run it again. + </p> + ) : ( + <StreamActionLog + trigger={() => runBookmarkAction(bookmark.id)} + cancelAction={cancelJobAction} + buttonLabel="Run again" + runningLabel="Running…" + label={`Run ${bookmark.name}`} + /> + )} + </li> + ); +} diff --git a/editor/app/jobs/components/BookmarksMenu.tsx b/editor/app/jobs/components/BookmarksMenu.tsx @@ -0,0 +1,45 @@ +import Link from "next/link"; +import type { JobBookmark } from "yt-dlp-transcript-common/jobs/bookmarks"; +import { BookmarkRunButton } from "./BookmarkRunButton"; + +// Compact bookmarks panel rendered on /jobs and /jobs/active: a quick-access +// row of run buttons (one per bookmark), each showing just enough to identify +// it. Launching is fire-and-forget — the inline shell log and the heavier +// rename/delete/metadata controls live on the linked /jobs/bookmarks page. +export function BookmarksMenu({ + bookmarks, + missingSlugs = [], +}: { + bookmarks: JobBookmark[]; + missingSlugs?: string[]; +}) { + if (bookmarks.length === 0) return null; + const missing = new Set(missingSlugs); + return ( + <section + aria-label="Bookmarked jobs" + className="flex flex-col gap-3 border border-zinc-200 dark:border-zinc-800 rounded-md p-3 bg-zinc-50 dark:bg-zinc-900/40" + > + <div className="flex items-center justify-between gap-2"> + <h2 className="text-sm font-medium text-zinc-700 dark:text-zinc-300"> + Bookmarks + </h2> + <Link + href="/jobs/bookmarks" + className="text-xs underline text-zinc-500 hover:text-zinc-900 dark:hover:text-zinc-100" + > + Manage + </Link> + </div> + <div className="flex flex-wrap items-start gap-2"> + {bookmarks.map((b) => ( + <BookmarkRunButton + key={b.id} + bookmark={b} + channelMissing={missing.has(b.spec.slug)} + /> + ))} + </div> + </section> + ); +} diff --git a/editor/app/jobs/components/JobsTable.tsx b/editor/app/jobs/components/JobsTable.tsx @@ -4,6 +4,7 @@ import { useEffect, useMemo, useState } from "react"; import Link from "next/link"; import type { JobListEntry } from "yt-dlp-transcript-common/jobs/listJobs"; import { CancelJobButton } from "./CancelJobButton"; +import { BookmarkJobButton } from "./BookmarkJobButton"; import { clearJobsFilters, defaultJobsFilters, @@ -274,9 +275,12 @@ export function JobsTable({ jobs }: { jobs: JobListEntry[] }) { {j.logSize.toLocaleString()} B </td> <td className="px-3 py-2 text-right"> - {(j.status === "running" || j.status === "queued") && ( - <CancelJobButton jobId={j.id} /> - )} + <div className="flex items-center justify-end gap-2"> + {j.bookmarkable && <BookmarkJobButton jobId={j.id} />} + {(j.status === "running" || j.status === "queued") && ( + <CancelJobButton jobId={j.id} /> + )} + </div> </td> </tr> ); diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx @@ -6,6 +6,7 @@ import { formatDuration } from "yt-dlp-transcript-common/lib/format"; import { JobLogTail } from "../[id]/components/JobLogTail"; import { DrainJobButton } from "./DrainJobButton"; import { CancelJobButton } from "./CancelJobButton"; +import { BookmarkJobButton } from "./BookmarkJobButton"; export type RunningJobsTask = { id: string; @@ -18,6 +19,12 @@ export type RunningJobsTask = { // 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 = { @@ -40,6 +47,8 @@ export type RunningJobsListItem = { tasks?: RunningJobsTask[]; draining?: boolean; drainable?: boolean; + // True when the job carries a replay descriptor and can be bookmarked. + bookmarkable?: boolean; }; type Props = { @@ -123,6 +132,7 @@ function JobRow({ </span> )} <div className="ml-auto flex items-center gap-2"> + {job.bookmarkable && <BookmarkJobButton jobId={job.id} />} {(job.drainable || job.draining) && ( <DrainJobButton jobId={job.id} draining={job.draining} /> )} @@ -170,16 +180,43 @@ function useNow(): number | null { } function TaskProgressBar({ task }: { task: RunningJobsTask }) { - const verb = task.kind === "download" ? "Downloading" : "Transcribing"; - const hasFraction = typeof task.fraction === "number"; - const pct = hasFraction ? Math.round((task.fraction as number) * 100) : 0; - const fillClass = - task.kind === "download" - ? "bg-emerald-300 dark:bg-emerald-700" - : "bg-emerald-500"; // 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(); + const probing = task.phase === "probing"; + const verb = probing + ? "Probing audio" + : task.kind === "download" + ? "Downloading" + : "Transcribing"; + // While probing, fill against the estimated probe duration (a distinct violet + // "scanning" bar) rather than the frozen download fraction. yt-dlp is paused, + // so the download fraction wouldn't advance anyway. Falls back to an + // indeterminate violet pulse on the first probe (no duration estimate yet). + const probeFraction = + probing && + now !== null && + typeof task.probeStartedAt === "number" && + typeof task.probeEtaSeconds === "number" && + task.probeEtaSeconds > 0 + ? Math.max( + 0, + Math.min(1, (now - task.probeStartedAt) / (task.probeEtaSeconds * 1000)), + ) + : undefined; + const hasFraction = probing + ? typeof probeFraction === "number" + : typeof task.fraction === "number"; + const fraction = probing ? probeFraction : task.fraction; + const pct = hasFraction ? Math.round((fraction as number) * 100) : 0; + const fillClass = probing + ? "bg-violet-400 dark:bg-violet-500" + : task.kind === "download" + ? "bg-emerald-300 dark:bg-emerald-700" + : "bg-emerald-500"; + const pulseClass = probing + ? "bg-violet-400 dark:bg-violet-500" + : "bg-amber-400 dark:bg-amber-500"; const elapsed = now === null ? null @@ -188,7 +225,7 @@ function TaskProgressBar({ task }: { task: RunningJobsTask }) { const meta = [ elapsed, hasFraction ? `${pct}%` : null, - task.detail, + probing ? "probing audio…" : task.detail, ] .filter(Boolean) .join(" · "); @@ -222,10 +259,10 @@ function TaskProgressBar({ task }: { task: RunningJobsTask }) { /> ) : ( // Indeterminate: no parseable progress yet (e.g. before whisper's - // header line or yt-dlp's first percent). - <div - className={`absolute inset-y-0 left-0 w-1/3 animate-pulse ${fillClass}`} - /> + // header line or yt-dlp's first percent, or the first audio probe). + // A full-width pulse reads as "working, no progress number" — distinct + // from the fill color so it can't be mistaken for a stalled ~33%. + <div className={`absolute inset-0 animate-pulse ${pulseClass}`} /> )} </div> </div> diff --git a/editor/app/jobs/loadBookmarks.ts b/editor/app/jobs/loadBookmarks.ts @@ -0,0 +1,22 @@ +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { readBookmarks, type JobBookmark } from "yt-dlp-transcript-common/jobs/bookmarks"; +import { readChannelConfig } from "yt-dlp-transcript-common/controller/channels"; + +// Load bookmarks for the /jobs and /jobs/active pages, plus the set of channel +// slugs that no longer exist — so the UI can flag bookmarks whose channel was +// deleted and disable their Run-again button. +export async function loadBookmarksView(): Promise<{ + bookmarks: JobBookmark[]; + missingSlugs: string[]; +}> { + const paths = getPaths(); + const bookmarks = await readBookmarks(paths); + const slugs = Array.from(new Set(bookmarks.map((b) => b.spec.slug))); + const checks = await Promise.all( + slugs.map( + async (s) => [s, Boolean(await readChannelConfig(paths, s))] as const, + ), + ); + const missingSlugs = checks.filter(([, ok]) => !ok).map(([s]) => s); + return { bookmarks, missingSlugs }; +} diff --git a/editor/app/jobs/page.tsx b/editor/app/jobs/page.tsx @@ -3,19 +3,25 @@ import { listAllJobs } from "yt-dlp-transcript-common/jobs/listJobs"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import { ClearArchivedButton } from "./components/ClearArchivedButton"; import { JobsTable } from "./components/JobsTable"; +import { BookmarksMenu } from "./components/BookmarksMenu"; +import { loadBookmarksView } from "./loadBookmarks"; export const dynamic = "force-dynamic"; export const metadata: Metadata = { title: "Jobs" }; export default async function JobsPage() { - const jobs = await listAllJobs(getPaths()); + const [jobs, { bookmarks, missingSlugs }] = await Promise.all([ + listAllJobs(getPaths()), + loadBookmarksView(), + ]); return ( <div className="flex flex-col gap-4"> <div className="flex items-center justify-between"> <h1 className="text-2xl font-semibold">Jobs</h1> <ClearArchivedButton /> </div> + <BookmarksMenu bookmarks={bookmarks} missingSlugs={missingSlugs} /> {jobs.length === 0 ? ( <p className="text-sm text-zinc-500 border border-dashed border-zinc-300 dark:border-zinc-700 rounded p-4"> No jobs have run yet. diff --git a/editor/app/jobs/runJobSpec.ts b/editor/app/jobs/runJobSpec.ts @@ -0,0 +1,134 @@ +// Server-only dispatcher: reached only through bookmarkActions ("use server"), +// so it is never bundled to the client. +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { readChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; +import type { JobSpec, ReplayBucket } from "yt-dlp-transcript-common/jobs/jobSpec"; +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; +import type { AudioFormat } from "yt-dlp-transcript-common/lib/channelConfig"; +import { + downloadAction, + downloadMissingAction, + downloadMissingSubsAction, + retryBucketAction, + storePlaylistAction, + syncAction, +} from "../channels/[slug]/pipelineActions"; +import { + cleanAudioAction, + cleanExtraAudioFormatsAction, + clearFailedTranscodingsAction, + clearFailedTranscriptionsAction, + removeFailedTranscodingsAction, + removeWrongFormatAudioAction, + transcodeFailuresAction, + transcodeUntranscodedAction, + transcribeBucketAction, + transcribeMissingAction, +} from "../channels/[slug]/whisperActions"; + +// The single place that maps a stored JobSpec back to the server action that +// runs it. Bucket jobs re-derive their work from the channel's CURRENT snapshot +// (not a frozen id list) via idsForBucket, so re-running "retry partial +// downloads" always acts on whatever is partial right now. + +async function idsForBucket( + slug: string, + bucket: ReplayBucket, +): Promise<string[]> { + const snap = await readChannelSnapshot(getPaths(), slug); + const ids = snap?.buckets?.[bucket]; + return Array.isArray(ids) ? ids : []; +} + +const str = (v: unknown): string | undefined => + typeof v === "string" ? v : undefined; +const bool = (v: unknown): boolean | undefined => + typeof v === "boolean" ? v : undefined; +const num = (v: unknown): number | undefined => + typeof v === "number" ? v : undefined; + +export async function runJobSpec(spec: JobSpec): Promise<StreamActionResult> { + const p = spec.params ?? {}; + const queueKey = str(p.queueKey); + switch (spec.kind) { + case "whisper-all": + return transcribeMissingAction( + spec.slug, + queueKey, + bool(p.reverse), + str(p.audioFormat) as AudioFormat | undefined, + bool(p.strictAudioFormat), + num(p.shardTotal), + num(p.shardIndex), + ); + case "whisper-bucket-downloaded-no-transcript": { + if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; + const ids = await idsForBucket(spec.slug, spec.bucket); + if (ids.length === 0) { + return { ok: false, error: "Nothing to transcribe right now.", info: true }; + } + return transcribeBucketAction( + spec.slug, + ids, + queueKey, + str(p.audioFormat) as AudioFormat | undefined, + bool(p.strictAudioFormat), + spec.bucket, + ); + } + case "retry-bucket": { + if (!spec.bucket) return { ok: false, error: "Bookmark is missing its bucket." }; + const ids = await idsForBucket(spec.slug, spec.bucket); + if (ids.length === 0) { + return { ok: false, error: "Nothing to retry right now.", info: true }; + } + return retryBucketAction( + spec.slug, + ids, + queueKey, + bool(p.abortOnError), + str(p.handlingOverride), + spec.bucket, + ); + } + case "download-from-playlist": + return downloadAction(spec.slug, queueKey, bool(p.abortOnError)); + case "download-missing": + return downloadMissingAction( + spec.slug, + queueKey, + bool(p.ignoreArchive), + bool(p.abortOnError), + num(p.shardTotal), + num(p.shardIndex), + ); + case "sync": + return syncAction(spec.slug, queueKey); + case "download-missing-subs": + return downloadMissingSubsAction(spec.slug, queueKey, bool(p.abortOnError)); + case "store-playlist": + return storePlaylistAction(spec.slug, queueKey); + case "clear-failed-transcriptions": + return clearFailedTranscriptionsAction(spec.slug, queueKey); + case "transcode-failures": + return transcodeFailuresAction( + spec.slug, + str(p.targetFormat) as AudioFormat, + queueKey, + ); + case "transcode-untranscoded": + return transcodeUntranscodedAction(spec.slug, queueKey); + case "remove-failed-transcodings": + return removeFailedTranscodingsAction(spec.slug, queueKey); + case "clear-failed-transcodings": + return clearFailedTranscodingsAction(spec.slug, queueKey); + case "clean-extra-audio-formats": + return cleanExtraAudioFormatsAction(spec.slug, queueKey); + case "remove-wrong-format-audio": + return removeWrongFormatAudioAction(spec.slug, queueKey); + case "clean-audio-transcribed": + return cleanAudioAction(spec.slug, queueKey); + default: + return { ok: false, error: `Cannot re-run job kind: ${spec.kind}` }; + } +} diff --git a/editor/app/widget/builder/components/WidgetBuilder.tsx b/editor/app/widget/builder/components/WidgetBuilder.tsx @@ -78,6 +78,31 @@ export function WidgetBuilder() { onChange={(v) => patch({ compact: v })} /> <Check + label="Show batch progress bar" + checked={config.jobBar} + onChange={(v) => patch({ jobBar: v })} + /> + <Check + label="Show progress in job heading" + checked={config.headingProgress} + onChange={(v) => patch({ headingProgress: v })} + /> + <Check + label="Show time estimates (ETA)" + checked={config.eta} + onChange={(v) => patch({ eta: v })} + /> + <Check + label="Show disk indicator" + checked={config.disk} + onChange={(v) => patch({ disk: v })} + /> + <Check + label="Show worker names" + checked={config.workerLabels} + onChange={(v) => patch({ workerLabels: v })} + /> + <Check label="Show section titles" checked={config.showTitles} onChange={(v) => patch({ showTitles: v })} diff --git a/editor/app/widget/components/MonitorWidget.tsx b/editor/app/widget/components/MonitorWidget.tsx @@ -103,12 +103,13 @@ export function MonitorWidget({ return ( <div className="flex flex-col gap-3 p-2 text-zinc-900 dark:text-zinc-100"> - {disk?.enabled && <DiskStrip disk={disk} />} + {config.disk && disk?.enabled && <DiskStrip disk={disk} />} {config.workers && ( <WorkersStrip workers={workers} paused={workersPayload?.paused ?? false} showTitle={config.showTitles} + showLabels={config.workerLabels} /> )} {config.jobs && ( @@ -116,6 +117,9 @@ export function MonitorWidget({ jobs={jobs} showTitle={config.showTitles} compact={config.compact} + jobBar={config.jobBar} + headingProgress={config.headingProgress} + eta={config.eta} /> )} </div> @@ -145,10 +149,12 @@ function WorkersStrip({ workers, paused, showTitle, + showLabels, }: { workers: WorkerView[]; paused: boolean; showTitle: boolean; + showLabels: boolean; }) { const busy = workers.filter((w) => w.busy).length; return ( @@ -177,12 +183,18 @@ function WorkersStrip({ key={w.id} title={`${w.name} — ${workerStateLabel(w)}`} aria-label={`worker ${w.name} ${workerStateLabel(w)}`} - className="flex items-center gap-1.5 rounded border border-zinc-200 dark:border-zinc-800 bg-white dark:bg-zinc-900 px-1.5 py-0.5 text-xs" + className={ + showLabels + ? "flex items-center gap-1.5 rounded border border-zinc-200 dark:border-zinc-800 bg-white dark:bg-zinc-900 px-1.5 py-0.5 text-xs" + : "flex items-center" + } > <span className={`inline-block h-2 w-2 shrink-0 rounded-full ${workerDotClass(w)}`} /> - <span className="truncate max-w-[10rem]">{w.name}</span> + {showLabels && ( + <span className="truncate max-w-[10rem]">{w.name}</span> + )} </li> ))} </ul> @@ -225,10 +237,16 @@ function ActiveJobsStrip({ jobs, showTitle, compact, + jobBar, + headingProgress, + eta, }: { jobs: RunningJobsListItem[]; showTitle: boolean; compact: boolean; + jobBar: boolean; + headingProgress: boolean; + eta: boolean; }) { const running = jobs.filter((j) => j.status === "running").length; const queued = jobs.filter((j) => j.status === "queued").length; @@ -250,7 +268,13 @@ function ActiveJobsStrip({ <ul className="flex flex-col gap-2"> {jobs.map((j) => ( <li key={j.id}> - <JobRow job={j} compact={compact} /> + <JobRow + job={j} + compact={compact} + jobBar={jobBar} + headingProgress={headingProgress} + eta={eta} + /> </li> ))} </ul> @@ -262,9 +286,15 @@ function ActiveJobsStrip({ function JobRow({ job, compact, + jobBar, + headingProgress, + eta, }: { job: RunningJobsListItem; compact: boolean; + jobBar: boolean; + headingProgress: boolean; + eta: boolean; }) { const statusClass = job.status === "running" @@ -284,8 +314,15 @@ function JobRow({ {job.channelSlug} </span> )} + {headingProgress && job.progress && ( + <span className="font-mono text-zinc-500 shrink-0"> + {jobProgressText(job.progress, eta)} + </span> + )} </div> - {job.progress && <JobProgressBar progress={job.progress} />} + {jobBar && job.progress && ( + <JobProgressBar progress={job.progress} showEta={eta} /> + )} {!compact && job.tasks && job.tasks.length > 0 && @@ -294,10 +331,33 @@ function JobRow({ ); } +// One-line textual summary of a job's batch progress, e.g. "↓ 5/10 · ~2m left". +// 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"]>, + showEta: boolean, +): string { + const label = + progress.metric === "downloads" + ? `↓ ${progress.current}/${progress.target}` + : `${progress.current}/${progress.target}`; + const remaining = progress.target - progress.current; + const etaText = + !showEta || remaining <= 0 + ? null + : typeof progress.etaSeconds === "number" + ? `~${formatDuration(Math.max(1, Math.round(progress.etaSeconds)))} left` + : "…"; + return etaText ? `${label} · ${etaText}` : label; +} + function JobProgressBar({ progress, + showEta, }: { progress: NonNullable<RunningJobsListItem["progress"]>; + showEta: boolean; }) { const label = progress.metric === "downloads" @@ -305,7 +365,7 @@ function JobProgressBar({ : `${progress.current}/${progress.target}`; const remaining = progress.target - progress.current; const etaText = - remaining <= 0 + !showEta || remaining <= 0 ? null : typeof progress.etaSeconds === "number" ? `~${formatDuration(Math.max(1, Math.round(progress.etaSeconds)))} left` @@ -339,19 +399,44 @@ function TaskBar({ task: NonNullable<RunningJobsListItem["tasks"]>[number]; }) { const now = useNow(); - const verb = task.kind === "download" ? "↓" : "✎"; - const hasFraction = typeof task.fraction === "number"; - const pct = hasFraction ? Math.round((task.fraction as number) * 100) : 0; - const fillClass = - task.kind === "download" + const probing = task.phase === "probing"; + const verb = probing ? "🔍" : task.kind === "download" ? "↓" : "✎"; + // While probing, fill against the estimated probe duration (violet "scanning" + // bar) rather than the frozen download fraction. See TaskProgressBar. + const probeFraction = + probing && + now !== null && + typeof task.probeStartedAt === "number" && + typeof task.probeEtaSeconds === "number" && + task.probeEtaSeconds > 0 + ? Math.max( + 0, + Math.min(1, (now - task.probeStartedAt) / (task.probeEtaSeconds * 1000)), + ) + : undefined; + const hasFraction = probing + ? typeof probeFraction === "number" + : typeof task.fraction === "number"; + const fraction = probing ? probeFraction : task.fraction; + const pct = hasFraction ? Math.round((fraction as number) * 100) : 0; + const fillClass = probing + ? "bg-violet-400 dark:bg-violet-500" + : task.kind === "download" ? "bg-emerald-300 dark:bg-emerald-700" : "bg-emerald-500"; + const pulseClass = probing + ? "bg-violet-400 dark:bg-violet-500" + : "bg-amber-400 dark:bg-amber-500"; const elapsed = now === null ? null : formatDuration(Math.max(0, Math.round((now - task.startedAt) / 1000))) || "0:00"; - const meta = [elapsed, hasFraction ? `${pct}%` : null, task.detail] + const meta = [ + elapsed, + hasFraction ? `${pct}%` : null, + probing ? "probing" : task.detail, + ] .filter(Boolean) .join(" · "); return ( @@ -375,9 +460,11 @@ function TaskBar({ style={{ width: `${pct}%` }} /> ) : ( - <div - className={`absolute inset-y-0 left-0 w-1/3 animate-pulse ${fillClass}`} - /> + // No fraction yet (just started, an engine that reports no %, or the + // first audio probe): a full-width pulse reads as "working, no + // progress number" — distinct from the emerald fill so it can't be + // mistaken for ~33%. + <div className={`absolute inset-0 animate-pulse ${pulseClass}`} /> )} </div> </div> diff --git a/editor/app/widget/lib/config.ts b/editor/app/widget/lib/config.ts @@ -17,6 +17,19 @@ export type WidgetConfig = { // When nothing is active, collapse to a tiny "Idle" line instead of the // full empty-state boxes. hideIdle: boolean; + // Show the batch/channel-level aggregate progress bar per job. Turn this off + // (with headingProgress on) to show only the individual per-task bars. + jobBar: boolean; + // Carry the batch progress (e.g. "5/10 · ~2m left") as text in each job's + // one-line heading — useful when the job bar is hidden. + headingProgress: boolean; + // Show ETA estimates wherever they appear (the job-bar caption and the + // heading progress text). Off = counts only, no "~Xm left". + eta: boolean; + // Show the free-disk indicator strip (when the disk gate is enabled). + disk: boolean; + // Show worker names in the Workers strip. Off = colored dots only (denser). + workerLabels: boolean; }; export const WIDGET_DEFAULTS: WidgetConfig = { @@ -27,6 +40,11 @@ export const WIDGET_DEFAULTS: WidgetConfig = { compact: false, showTitles: true, hideIdle: false, + jobBar: true, + headingProgress: false, + eta: true, + disk: true, + workerLabels: true, }; // Next's searchParams give each key as string | string[] | undefined. @@ -63,6 +81,11 @@ export function parseWidgetConfig(params: RawParams): WidgetConfig { compact: parseBool(params.compact, WIDGET_DEFAULTS.compact), showTitles: parseBool(params.titles, WIDGET_DEFAULTS.showTitles), hideIdle: first(params.idle) === "hide" ? true : WIDGET_DEFAULTS.hideIdle, + jobBar: parseBool(params.jobbar, WIDGET_DEFAULTS.jobBar), + headingProgress: parseBool(params.headtext, WIDGET_DEFAULTS.headingProgress), + eta: parseBool(params.eta, WIDGET_DEFAULTS.eta), + disk: parseBool(params.disk, WIDGET_DEFAULTS.disk), + workerLabels: parseBool(params.wnames, WIDGET_DEFAULTS.workerLabels), }; } @@ -81,5 +104,14 @@ export function buildWidgetQuery(config: WidgetConfig): string { if (config.showTitles !== WIDGET_DEFAULTS.showTitles) sp.set("titles", config.showTitles ? "1" : "0"); if (config.hideIdle !== WIDGET_DEFAULTS.hideIdle) sp.set("idle", "hide"); + if (config.jobBar !== WIDGET_DEFAULTS.jobBar) + sp.set("jobbar", config.jobBar ? "1" : "0"); + if (config.headingProgress !== WIDGET_DEFAULTS.headingProgress) + sp.set("headtext", config.headingProgress ? "1" : "0"); + if (config.eta !== WIDGET_DEFAULTS.eta) sp.set("eta", config.eta ? "1" : "0"); + if (config.disk !== WIDGET_DEFAULTS.disk) + sp.set("disk", config.disk ? "1" : "0"); + if (config.workerLabels !== WIDGET_DEFAULTS.workerLabels) + sp.set("wnames", config.workerLabels ? "1" : "0"); return sp.toString(); } diff --git a/editor/app/workers/components/WorkersView.tsx b/editor/app/workers/components/WorkersView.tsx @@ -382,7 +382,9 @@ function TaskRow({ task }: { task: WorkerTask }) { style={{ width: `${pct}%` }} /> ) : ( - <span className="absolute inset-y-0 left-0 w-1/3 animate-pulse bg-blue-500" /> + // Indeterminate (no parseable progress yet): a full-width amber pulse + // reads as "working, no number" rather than a stalled ~33%. + <span className="absolute inset-0 animate-pulse bg-amber-400 dark:bg-amber-500" /> )} </span> </div> diff --git a/editor/e2e/audio-check-scenarios.spec.ts b/editor/e2e/audio-check-scenarios.spec.ts @@ -424,6 +424,49 @@ test.describe("audio-checked download scenarios", () => { ); }); + test("integrity probe surfaces a 'probing' phase on /jobs/active and in the log", async ({ + page, + request, + }) => { + // Slow (~600ms) probe + fast chunks: while yt-dlp is paused for the probe + // the download task should report phase "probing" in the active-jobs + // payload (the data behind the violet scanning bar), and the orchestrator + // should emit DLOM_PROBE markers into the job log. + await resetData("audio-check-channel"); + await writeFakeConfig({ + mode: "happy", + slowProbe: true, + totalChunks: 24, + chunkDelayMs: 30, + }); + await triggerDownload(page); + + await expect + .poll( + async () => { + const res = await request.get(`${baseUrl}/api/jobs/active`); + if (!res.ok()) return false; + const body = await res.json(); + for (const job of body.jobs ?? []) { + for (const t of job.tasks ?? []) { + if (t.phase === "probing") return true; + } + } + return false; + }, + { timeout: 20_000, intervals: [50] }, + ) + .toBe(true); + + await waitForOutcome( + (o) => o.status === "ok-audio-checked", + "probing-phase outcome", + ); + + const log = page.getByLabel("Download videos output"); + await expect(log).toContainText("DLOM_PROBE status=start"); + }); + test("resumeDuringProbe keeps the download flowing through the probe (legacy)", async ({ page, }) => { diff --git a/editor/e2e/bookmarks.spec.ts b/editor/e2e/bookmarks.spec.ts @@ -0,0 +1,269 @@ +// Bookmark a running/finished job and re-launch it later. The key behavior is +// that re-running RE-DERIVES the work from the channel's current state rather +// than replaying a frozen video-id list: +// - whisper-all re-scans the disk for audio-without-transcript; +// - bucket jobs (retry partial downloads) read the CURRENT snapshot bucket. +// Both proofs change what the channel needs between the original run and the +// re-run, then assert the re-run acts on the new set. +// +// The bookmarks UI is split: a compact quick-access menu on /jobs and +// /jobs/active (one fire-and-forget run button per bookmark, NO inline log), +// and a full management page at /jobs/bookmarks (rename, delete, metadata, and +// the streaming "Run again" with its inline log). The behavioral re-derive +// proofs run against the management page's streaming log. + +import { rename, rm, writeFile } from "node:fs/promises"; +import { test, expect } from "@playwright/test"; +import { pathExists, resetData, resolvePath } from "./helpers"; + +const CHANNEL = "test-transcribe"; +const ROOT = `test-transcripts/channels/${CHANNEL}`; + +// Run "Transcribe missing" on the audio fixture and bookmark the resulting +// whisper-all job from the jobs table (newest bookmarkable row). Ends on /jobs, +// where the compact bookmarks menu is now shown. +async function transcribeAndBookmark(page: import("@playwright/test").Page) { + await page.goto(`/channels/${CHANNEL}`); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + await expect(page.getByLabel("Transcribe missing output")).toContainText( + "3 succeeded", + { timeout: 30_000 }, + ); + await page.goto("/jobs"); + await page + .getByRole("button", { name: /bookmark job/i }) + .first() + .click(); + await expect(page.getByRole("heading", { name: "Bookmarks" })).toBeVisible({ + timeout: 15_000, + }); +} + +test("compact menu: one-click run is fire-and-forget with no inline log, and links to Manage", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData("one-transcribe-channel-with-audio"); + await transcribeAndBookmark(page); + + const menu = page.getByRole("region", { name: "Bookmarked jobs" }); + // "Manage" links to the full management page. + await expect(menu.getByRole("link", { name: "Manage" })).toHaveAttribute( + "href", + "/jobs/bookmarks", + ); + // The menu is buttons only — no inline shell / expanding log. + await expect(menu.getByRole("log")).toHaveCount(0); + + const runBtn = menu.getByRole("button", { + name: "run bookmark whisper-all · test-transcribe", + }); + await expect(runBtn).toBeVisible(); + await runBtn.click(); + + // Fire-and-forget: the button confirms the launch and still no log appears. + await expect(runBtn).toHaveText("Launched ✓", { timeout: 15_000 }); + await expect(menu.getByRole("log")).toHaveCount(0); +}); + +test("management page: Run again re-derives the current missing set", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData("one-transcribe-channel-with-audio"); + await transcribeAndBookmark(page); + await page.goto("/jobs/bookmarks"); + + const row = page.getByLabel("bookmark whisper-all · test-transcribe", { + exact: true, + }); + await expect(row).toBeVisible(); + + // Change the channel's state: remove one transcript so exactly one video now + // needs transcribing again. A frozen-id replay would re-attempt all three; a + // re-derive does only the one that's currently missing. + await rm(resolvePath(`${ROOT}/data/vidA/transcript.json`)); + + await row.getByRole("button", { name: "Run again" }).click(); + await expect( + page.getByRole("log", { name: /Run whisper-all/ }), + ).toContainText("1 succeeded", { timeout: 30_000 }); + expect(await pathExists(`${ROOT}/data/vidA/transcript.json`)).toBe(true); +}); + +test("management page: rename a bookmark and see created / last-run metadata", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData("one-transcribe-channel-with-audio"); + await transcribeAndBookmark(page); + await page.goto("/jobs/bookmarks"); + + const row = page.getByLabel("bookmark whisper-all · test-transcribe", { + exact: true, + }); + await expect(row).toContainText("Created"); + await expect(row).toContainText("Last run never"); + + // Rename. + await row + .getByRole("button", { name: "rename bookmark whisper-all · test-transcribe" }) + .click(); + await row + .getByLabel("new name for whisper-all · test-transcribe") + .fill("Nightly transcribe"); + await row.getByRole("button", { name: "Save" }).click(); + + const renamed = page.getByLabel("bookmark Nightly transcribe", { + exact: true, + }); + await expect(renamed).toBeVisible({ timeout: 15_000 }); + + // Run again updates the last-run timestamp (away from "never"). + await renamed.getByRole("button", { name: "Run again" }).click(); + await expect( + page.getByRole("log", { name: /Run Nightly transcribe/ }), + ).toContainText("Whisper batch:", { timeout: 30_000 }); + await expect( + page.getByLabel("bookmark Nightly transcribe", { exact: true }), + ).not.toContainText("Last run never", { timeout: 15_000 }); +}); + +test("management page: Run again retries whatever is partial now, then delete it", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData("one-transcribe-channel-with-audio"); + // The retry flow reads the playlist to resolve each id's URL. Slug-style URLs + // keep extractVideoId() === the id without needing platform claim ids. + await writeFile( + resolvePath(`${ROOT}/playlist`), + "https://www.youtube.com/watch?v=vidA\nhttps://www.youtube.com/watch?v=vidB\n", + ); + + // Strand vidA mid-download so it surfaces in the partial-downloads bucket. + await rename( + resolvePath(`${ROOT}/data/vidA/audio.m4a`), + resolvePath(`${ROOT}/data/vidA/audio.m4a.part`), + ); + + // Original run: resume the partial download (vidA). + await page.goto(`/channels/${CHANNEL}`); + await page + .getByLabel("retry resume partial downloads bucket") + .getByRole("button", { name: /^Retry \(1\)$/ }) + .click(); + await expect( + page.getByLabel("Retry resume partial downloads output"), + ).toContainText("Managed download complete", { timeout: 30_000 }); + expect(await pathExists(`${ROOT}/data/vidA/audio.m4a`)).toBe(true); + + // Bookmark the retry-bucket job, then manage it on /jobs/bookmarks. + await page.goto("/jobs"); + await page + .getByRole("button", { name: /bookmark job/i }) + .first() + .click(); + await expect(page.getByRole("heading", { name: "Bookmarks" })).toBeVisible({ + timeout: 15_000, + }); + await page.goto("/jobs/bookmarks"); + const row = page.getByLabel( + "bookmark retry-bucket (partialDownloads) · test-transcribe", + { exact: true }, + ); + await expect(row).toBeVisible({ timeout: 15_000 }); + + // Now make a DIFFERENT video partial (vidB) and force a fresh snapshot, so the + // partial-downloads bucket currently holds vidB, not the original vidA. + await rename( + resolvePath(`${ROOT}/data/vidB/audio.m4a`), + resolvePath(`${ROOT}/data/vidB/audio.m4a.part`), + ); + await rm(resolvePath(`${ROOT}/snapshot.json`), { force: true }); + await page.goto(`/channels/${CHANNEL}`); // regenerates snapshot when absent + await expect( + page.getByRole("heading", { name: /Partial downloads \(1\)/ }), + ).toBeVisible(); + + // Run again: re-derives the current bucket → downloads vidB (NOT the original + // vidA). This is the proof that re-running is not a frozen-id replay. + await page.goto("/jobs/bookmarks"); + await row.getByRole("button", { name: "Run again" }).click(); + await expect( + page.getByRole("log", { name: /Run retry-bucket/ }), + ).toContainText("Managed download complete", { timeout: 30_000 }); + expect(await pathExists(`${ROOT}/data/vidB/audio.m4a`)).toBe(true); + + // Delete the bookmark (confirm-gated). + await row + .getByRole("button", { + name: "delete bookmark retry-bucket (partialDownloads) · test-transcribe", + exact: true, + }) + .click(); + await row + .getByRole("button", { + name: "confirm delete bookmark retry-bucket (partialDownloads) · test-transcribe", + exact: true, + }) + .click(); + await expect(row).toHaveCount(0, { timeout: 15_000 }); +}); + +test("management page: an empty re-derived bucket reads as a neutral notice, not an error", async ({ + page, +}) => { + test.setTimeout(60_000); + await resetData("one-transcribe-channel-with-audio"); + await writeFile( + resolvePath(`${ROOT}/playlist`), + "https://www.youtube.com/watch?v=vidA\n", + ); + await rename( + resolvePath(`${ROOT}/data/vidA/audio.m4a`), + resolvePath(`${ROOT}/data/vidA/audio.m4a.part`), + ); + + await page.goto(`/channels/${CHANNEL}`); + await page + .getByLabel("retry resume partial downloads bucket") + .getByRole("button", { name: /^Retry \(1\)$/ }) + .click(); + await expect( + page.getByLabel("Retry resume partial downloads output"), + ).toContainText("Managed download complete", { timeout: 30_000 }); + + await page.goto("/jobs"); + await page + .getByRole("button", { name: /bookmark job/i }) + .first() + .click(); + await expect(page.getByRole("heading", { name: "Bookmarks" })).toBeVisible({ + timeout: 15_000, + }); + await page.goto("/jobs/bookmarks"); + const row = page.getByLabel( + "bookmark retry-bucket (partialDownloads) · test-transcribe", + { exact: true }, + ); + await expect(row).toBeVisible({ timeout: 15_000 }); + + // vidA is fully downloaded now, so the partial-downloads bucket is empty. + // Refresh the snapshot to reflect that. + await rm(resolvePath(`${ROOT}/snapshot.json`), { force: true }); + await page.goto(`/channels/${CHANNEL}`); + await expect( + page.getByRole("heading", { name: /Partial downloads/ }), + ).toHaveCount(0); + + await page.goto("/jobs/bookmarks"); + await row.getByRole("button", { name: "Run again" }).click(); + // Neutral status, not a red alert. + await expect( + page.getByRole("status", { name: /Run retry-bucket.*notice/ }), + ).toContainText("Nothing to retry right now", { timeout: 30_000 }); + await expect( + page.getByRole("alert", { name: /Run retry-bucket/ }), + ).toHaveCount(0); +}); diff --git a/editor/e2e/whisper.spec.ts b/editor/e2e/whisper.spec.ts @@ -1,7 +1,35 @@ -import { readFile, writeFile } from "node:fs/promises"; +import { mkdir, readFile, writeFile } from "node:fs/promises"; import { test, expect } from "@playwright/test"; import { pathExists, resetData, resolvePath } from "./helpers"; +// A video the downloader couldn't produce real audio for: the dir exists but the +// only audio is an in-progress `.part`, and the audio-check pipeline recorded a +// `failed-corrupt-source` outcome. Helper writes that state plus a stale +// failed-transcriptions entry (the bug being fixed wrongly left these behind). +async function writeCorruptSourceVideo(id: string, alreadyFailed: boolean) { + const dir = `test-transcripts/channels/test-transcribe/data/${id}`; + await mkdir(resolvePath(dir), { recursive: true }); + await writeFile(resolvePath(`${dir}/audio.m4a.part`), "partial-bytes"); + await writeFile( + resolvePath(`${dir}/download-outcome.json`), + JSON.stringify({ + videoId: id, + status: "failed-corrupt-source", + startedAt: "2026-06-01T00:00:00.000Z", + finishedAt: "2026-06-01T00:01:00.000Z", + attempts: [], + }), + ); + if (alreadyFailed) { + await writeFile( + resolvePath( + "test-transcripts/channels/test-transcribe/failed-transcriptions", + ), + `${id}\n`, + ); + } +} + test("transcribes every audio file with no transcript", async ({ page }) => { await resetData("one-transcribe-channel-with-audio"); await page.goto("/channels/test-transcribe"); @@ -47,6 +75,52 @@ test("Transcribe missing excludes VTT-only videos instead of failing them", asyn ).toBe(false); }); +test("Transcribe missing skips corrupt-source videos and auto-prunes them from the failed list", async ({ + page, +}) => { + await resetData("one-transcribe-channel-with-audio"); + // vidBad has no real audio (only a .part) and is stale-listed as failed. + await writeCorruptSourceVideo("vidBad", true); + await page.goto("/channels/test-transcribe"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + // The three genuinely-downloaded videos transcribe; vidBad is never attempted. + await expect(page.getByLabel("Transcribe missing output")).toContainText( + "3 succeeded", + { timeout: 30_000 }, + ); + // vidBad was a download/source problem, not a transcription failure: it must + // be pruned from failed-transcriptions (so a re-download retries it), and it + // must not have been transcribed or re-recorded as failed. + const raw = await readFile( + resolvePath( + "test-transcripts/channels/test-transcribe/failed-transcriptions", + ), + "utf8", + ); + expect(raw).not.toContain("vidBad"); + expect( + await pathExists( + "test-transcripts/channels/test-transcribe/data/vidBad/transcript.json", + ), + ).toBe(false); +}); + +test("a corrupt-source video surfaces as corrupt_source, not a failed transcription", async ({ + page, +}) => { + await resetData("one-transcribe-channel-with-audio"); + await writeCorruptSourceVideo("vidBad", false); + await page.goto("/channels/test-transcribe"); + // The video row is labeled corrupt_source (download/source problem) rather + // than failed, and the Download stage summary calls it out for re-download. + await expect(page.getByLabel("status corrupt_source")).toBeVisible({ + timeout: 15_000, + }); + await expect( + page.getByText(/corrupt source \(needs re-download\)/), + ).toBeVisible(); +}); + test("verify reports nothing missing once transcribed", async ({ page }) => { await resetData("one-transcribe-channel-with-audio"); await page.goto("/channels/test-transcribe"); diff --git a/editor/e2e/widget.spec.ts b/editor/e2e/widget.spec.ts @@ -70,3 +70,32 @@ test("the builder keeps the shell and composes a link", async ({ page }) => { await expect(popup).toHaveURL(/\/widget\?workers=0$/); await popup.close(); }); + +test("display toggles serialize into the link and preview", async ({ page }) => { + await page.goto("/widget/builder"); + + const url = page.getByLabel("widget URL"); + const preview = page.locator('iframe[title="widget preview"]'); + await expect(url).toHaveValue(/\/widget$/); + + // Each new display flag is omitted at its default and appears once toggled, + // in both the copyable link and the live preview src. + await page.getByRole("checkbox", { name: "Show batch progress bar" }).uncheck(); + await expect(url).toHaveValue(/[?&]jobbar=0(&|$)/); + await expect(preview).toHaveAttribute("src", /[?&]jobbar=0(&|$)/); + + await page + .getByRole("checkbox", { name: "Show progress in job heading" }) + .check(); + await expect(url).toHaveValue(/[?&]headtext=1(&|$)/); + + await page.getByRole("checkbox", { name: "Show time estimates (ETA)" }).uncheck(); + await expect(url).toHaveValue(/[?&]eta=0(&|$)/); + + await page.getByRole("checkbox", { name: "Show disk indicator" }).uncheck(); + await expect(url).toHaveValue(/[?&]disk=0(&|$)/); + + await page.getByRole("checkbox", { name: "Show worker names" }).uncheck(); + await expect(url).toHaveValue(/[?&]wnames=0(&|$)/); + await expect(preview).toHaveAttribute("src", /[?&]wnames=0(&|$)/); +});