// Pure parsers that turn a line of yt-dlp / whisper shell output into a // per-operation progress update. No I/O — kept dependency-free so they can be // 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; // 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)); } function formatClock(totalSeconds: number): string { const s = Math.max(0, Math.floor(totalSeconds)); const hh = Math.floor(s / 3600); const mm = Math.floor((s % 3600) / 60); const ss = s % 60; const pad = (n: number) => String(n).padStart(2, "0"); return hh > 0 ? `${hh}:${pad(mm)}:${pad(ss)}` : `${mm}:${pad(ss)}`; } // Bytes/sec → yt-dlp-style human rate, e.g. 4533988 → "4.33MiB/s". function formatBytesPerSec(bytesPerSec: number): string { const units = ["B", "KiB", "MiB", "GiB", "TiB"]; let v = bytesPerSec; let i = 0; while (v >= 1024 && i < units.length - 1) { v /= 1024; i += 1; } const digits = v >= 100 || i === 0 ? 0 : 2; return `${v.toFixed(digits)}${units[i]}/s`; } // Marker + template for yt-dlp's --progress-template. We emit a structured, // numeric line (rather than scrape the human "[download] …%" line) because the // fields are far easier and more robust to parse — see parseStructuredDownloadProgress. // The marker is deliberately distinct from downloadOneManaged's DLOM_ARCHIVE so // the archive-line parser never confuses the two. export const DOWNLOAD_PROGRESS_MARKER = "DLOM_PROGRESS"; export const DOWNLOAD_PROGRESS_TEMPLATE = `download:${DOWNLOAD_PROGRESS_MARKER}` + " status=%(progress.status)s" + " downloaded=%(progress.downloaded_bytes)s" + " total=%(progress.total_bytes,progress.total_bytes_estimate)s" + " frag=%(progress.fragment_index)s/%(progress.fragment_count)s" + " speed=%(progress.speed)s" + " eta=%(progress.eta)s"; // yt-dlp renders unavailable template fields as the literal "NA". Treat that // (and empty / non-finite) as "unknown". function num(value: string | undefined): number | undefined { if (value === undefined || value === "" || value === "NA") return undefined; const n = Number.parseFloat(value); return Number.isFinite(n) ? n : undefined; } // Structured progress line emitted by DOWNLOAD_PROGRESS_TEMPLATE, e.g. // DLOM_PROGRESS status=downloading downloaded=31744 total=117526 frag=NA/NA speed=7519285.37 eta=0 // 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 => line.match(new RegExp(`\\b${key}=(\\S+)`))?.[1]; const status = field("status"); const downloaded = num(field("downloaded")); const total = num(field("total")); const speed = num(field("speed")); const eta = num(field("eta")); const fragStr = field("frag"); // e.g. "3/6" or "NA/NA" const [fragIndex, fragCount] = (fragStr ?? "").split("/").map(num); let fraction: number | undefined; if (status === "finished") { fraction = 1; } else if (total !== undefined && total > 0 && downloaded !== undefined) { fraction = clamp01(downloaded / total); } else if ( fragIndex !== undefined && fragCount !== undefined && fragCount > 0 ) { fraction = clamp01(fragIndex / fragCount); } const detailParts: string[] = []; if (speed !== undefined) detailParts.push(formatBytesPerSec(speed)); if (eta !== undefined) { const shownEta = adjustEtaSeconds ? adjustEtaSeconds(eta) : eta; detailParts.push(`ETA ${formatClock(shownEta)}`); } if (fragIndex !== undefined && fragCount !== undefined) { detailParts.push(`frag ${fragIndex}/${fragCount}`); } if (fraction === undefined && detailParts.length === 0) return null; const out: ProgressUpdate = {}; if (fraction !== undefined) out.fraction = fraction; if (detailParts.length > 0) out.detail = detailParts.join(" · "); return out; } // yt-dlp download progress, e.g. // [download] 0.1% of ~ 1.04GiB at 2.31MiB/s ETA Unknown (frag 3/2599) // [download] 100% of 123.95KiB in 00:00:01 // Returns null for non-progress lines (Destination:, Downloading item, …). export function parseDownloadProgress(line: string): ProgressUpdate | null { // Managed downloads emit a structured --progress-template line; prefer it. const structured = parseStructuredDownloadProgress(line); if (structured) return structured; // Legacy fallback: scrape yt-dlp's human "[download] …%" line. Still used by // non-managed paths (single-video / subs-only) and the e2e fake. if (!line.includes("[download]")) return null; const pctMatch = line.match(/\[download\]\s+([\d.]+)%/); const fragMatch = line.match(/\(frag\s+(\d+)\/(\d+)\)/); const rateMatch = line.match(/at\s+([\d.]+\s*[KMG]?i?B\/s)/); const etaMatch = line.match(/ETA\s+([\d:]+)/); let fraction: number | undefined; if (pctMatch) fraction = clamp01(Number.parseFloat(pctMatch[1]) / 100); const detailParts: string[] = []; if (rateMatch) detailParts.push(rateMatch[1].replace(/\s+/g, "")); if (etaMatch) detailParts.push(`ETA ${etaMatch[1]}`); if (fragMatch) { detailParts.push(`frag ${fragMatch[1]}/${fragMatch[2]}`); // When the total size is unknown there's no usable percent; fall back to // the fragment ratio so the bar still advances. if (fraction === undefined) { const done = Number.parseInt(fragMatch[1], 10); const total = Number.parseInt(fragMatch[2], 10); if (total > 0) fraction = clamp01(done / total); } } if (fraction === undefined && detailParts.length === 0) return null; const out: ProgressUpdate = {}; if (fraction !== undefined) out.fraction = fraction; if (detailParts.length > 0) out.detail = detailParts.join(" · "); 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; // Subtitle-only downloads (youtube handling runs yt-dlp with --skip-download) // fetch one or more tiny .vtt files. yt-dlp can't report a byte total for // them, so DLOM_PROGRESS carries total=NA while downloading and the per-byte // fraction never moves. Instead, count the announced subtitle tracks and // advance the bar one step per completed track. `subsTotal > 0` means we're in // this mode; it's cleared the moment a real media download starts (a // "downloading" line with a concrete byte total — e.g. a youtube no-subs // fallback that switches to an audio download). let subsTotal = 0; let subsDone = 0; 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; } // yt-dlp announces the subtitle tracks it's about to fetch, e.g. // [info] : Downloading subtitles: en-orig, en // Use the count as the denominator for the step-progress bar below. const subsAnnounce = line.match(/Downloading subtitles:\s*(.+?)\s*$/); if (subsAnnounce) { const tracks = subsAnnounce[1] .split(",") .map((t) => t.trim()) .filter(Boolean); if (tracks.length > 0) { subsTotal = tracks.length; subsDone = 0; } return null; } if (line.includes(DOWNLOAD_PROGRESS_MARKER)) { const field = (key: string): string | undefined => line.match(new RegExp(`\\b${key}=(\\S+)`))?.[1]; const status = field("status"); const total = num(field("total")); // A "downloading" line with a real byte total means actual media is // streaming (not a subtitle file) — leave subtitle-step mode. if (status === "downloading" && total !== undefined && total > 0) { subsTotal = 0; } if (subsTotal > 0) { if (status === "finished" && subsDone < subsTotal) subsDone += 1; const structured = parseStructuredDownloadProgress( line, adjustEtaSeconds, ); const out: ProgressUpdate = { fraction: clamp01(subsDone / subsTotal), }; const parts: string[] = []; if (structured?.detail) parts.push(structured.detail); parts.push(`subs ${subsDone}/${subsTotal}`); out.detail = parts.join(" · "); return out; } } 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. export function createTranscribeProgressParser(): { feed: (line: string) => ProgressUpdate | null; } { let totalSec: number | undefined; return { feed(line: string): ProgressUpdate | null { const headerMatch = line.match(/\(\d+\s+samples,\s+([\d.]+)\s+sec\)/); if (headerMatch) { totalSec = Number.parseFloat(headerMatch[1]); return null; } const tsMatch = line.match( /^\s*\[(\d{2}):(\d{2}):(\d{2})\.\d+\s*-->\s*(\d{2}):(\d{2}):(\d{2})\.\d+\]/, ); if (!tsMatch) return null; const endSec = Number.parseInt(tsMatch[4], 10) * 3600 + Number.parseInt(tsMatch[5], 10) * 60 + Number.parseInt(tsMatch[6], 10); if (totalSec === undefined || totalSec <= 0) { // No header yet — show position but not a fraction. return { detail: formatClock(endSec) }; } return { fraction: clamp01(endSec / totalSec), detail: `${formatClock(endSec)} / ${formatClock(totalSec)}`, }; }, }; } // chough output is bar-based, not segment-timestamp based. It prints a header // audio: 28595.7s • chunks: 60s • format: json // then repeatedly overwrites a progress bar of full (█) and light (░) blocks // with a trailing ETA, e.g. // ████░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░░ ETA 38m 57s // Progress = filled blocks / total blocks; the ETA text is surfaced as detail. export function createChoughProgressParser(): { feed: (line: string) => ProgressUpdate | null; } { return { feed(line: string): ProgressUpdate | null { // Consume (and ignore) the audio header line. if (/\baudio:\s*[\d.]+s\b/.test(line)) return null; const filled = (line.match(/█/g) ?? []).length; const empty = (line.match(/░/g) ?? []).length; const total = filled + empty; if (total === 0) return null; const out: ProgressUpdate = { fraction: clamp01(filled / total) }; // ETA text up to an ANSI clear ("[K") or end of line. const etaMatch = line.match(/ETA\s+([0-9hdms\s]+?)\s*(?:\x1b?\[K|$)/); const eta = etaMatch?.[1]?.trim(); if (eta) out.detail = `ETA ${eta}`; return out; }, }; } // parakeet-stitch.mjs (the overlapping-segment wrapper) prints one stderr line // per window, e.g. // parakeet-stitch: segment 3/12 @360s — transcribing (avg 41.2s/seg, ETA 6:10) // Progress = completed segment index / total segments. Once a segment has // finished the line carries a per-video ETA (avg measured time per segment × // remaining segments), which we surface verbatim in the detail. // scripts/diarize-sherpa.py, windowed mode. It already printed everything a bar // needs and nothing was reading it: // // diarize-sherpa: 6.20 h -> 9 window(s) of 45 min (15s overlap), windowed mode // diarize-sherpa: window 3/9 (90-135 min): 412 turns, 6 local speaker(s) in 210.3s // // WHY THIS IS WORTH A PARSER AT ALL. A backfill task got the DOWNLOAD parser, // which matches nothing in this output, so a diarization showed an indeterminate // bar — for up to 46 minutes on the 6-hour files, while the job-level bar sat at // "0/9 videos" because no video had finished yet. Both readouts were technically // honest and together they said nothing: the operator could not tell a working // diarization from a wedged one, which is the exact case windowing was built for. // // The per-window line is printed AFTER the window is processed and carries its // elapsed time, so `i` is a COMPLETED count and the times give a measured ETA // rather than a guess. Windows are near-identical in cost (same duration, same // model), which is what makes a flat mean good enough here — the parakeet parser // takes its ETA from the producer for the same reason the producer knows better // there and does not here. // // Short files take the non-windowed path and print no per-window line at all. // They stay indeterminate on purpose: there is no progress to report, and the // whole run is over in the time the windowed path spends on one window. export function createDiarizeProgressParser(): { feed: (line: string) => ProgressUpdate | null; } { let elapsedSum = 0; let elapsedCount = 0; return { feed(line: string): ProgressUpdate | null { if (!line.includes("diarize-sherpa:")) return null; // The plan line, before any window has run. Worth consuming so the bar is // determinate from the start rather than jumping from a spinner to 1/9. const planned = line.match(/->\s*(\d+)\s+window\(s\)/); if (planned) { const n = Number.parseInt(planned[1], 10); if (Number.isFinite(n) && n > 0) { return { fraction: 0, detail: `0/${n} windows` }; } return null; } const m = line.match(/window\s+(\d+)\s*\/\s*(\d+)/i); if (!m) return null; const i = Number.parseInt(m[1], 10); const n = Number.parseInt(m[2], 10); if (!Number.isFinite(n) || n <= 0 || !Number.isFinite(i)) return null; // "skipping" is the zero-samples window: it reports an index but no time, // and counting its (absent) duration would drag the mean down. const secs = line.match(/in\s+([\d.]+)s\b/); if (secs) { const v = Number.parseFloat(secs[1]); if (Number.isFinite(v) && v > 0) { elapsedSum += v; elapsedCount++; } } let detail = `window ${i}/${n}`; if (elapsedCount > 0 && i < n) { const eta = Math.round((elapsedSum / elapsedCount) * (n - i)); if (eta > 0) detail += ` · ETA ${formatClock(eta)}`; } return { fraction: clamp01(i / n), detail }; }, }; } export function createParakeetProgressParser(): { feed: (line: string) => ProgressUpdate | null; } { return { feed(line: string): ProgressUpdate | null { const m = line.match(/segment\s+(\d+)\s*\/\s*(\d+)/i); if (!m) return null; const i = Number.parseInt(m[1], 10); const n = Number.parseInt(m[2], 10); if (!Number.isFinite(n) || n <= 0) return null; let detail = `segment ${i}/${n}`; const eta = line.match(/ETA\s+([\d:]+)/)?.[1]; if (eta) detail += ` · ETA ${eta}`; return { fraction: clamp01(i / n), detail }; }, }; } // --------------------------------------------------------------------------- // rsync --info=progress2, for the relocate job // --------------------------------------------------------------------------- // ONE LINE OF `rsync --info=progress2`, parsed. // // rsync rewrites this line in place with a carriage return, so a chunk read off // the child's stdout holds several of them; the caller splits on /[\r\n]/ and // feeds each part. The shape (measured, not assumed — rsync 3.x, no // --human-readable): // // 30,000,000 85% 1.03GB/s 0:00:00 (xfr#1, to-chk=1/3) // // Four fields: bytes transferred so far (grouped with commas), rsync's own // percentage, a rate, and an ETA as h:mm:ss. The trailing `(xfr#…)` appears only // once a file completes and is deliberately not parsed — `to-chk` counts FILES, // and a relocate is priced in bytes. // // WHY THE PERCENTAGE IS NOT THE FRACTION WE REPORT. rsync's own percentage is // against the total it has scanned SO FAR: with incremental recursion (the // default) an early line reads "85%" of a tree it has only half enumerated, and // the bar then walks backwards. The relocate controller has already measured the // whole tree for its space check, so it divides by that instead and the fraction // only ever moves forward. This parser returns both and lets the caller choose. export type RsyncProgress = { // Bytes rsync says it has transferred so far. bytes: number; // rsync's own percentage, 0..100. See above for why it is not the fraction. percent: number; // Verbatim, e.g. "1.03GB/s". Not re-formatted: it is already the unit an // operator watching a copy reads in. rate: string; // Seconds, from rsync's h:mm:ss ETA. etaSeconds: number; }; const RSYNC_PROGRESS_RE = /^\s*([\d,._ ]*\d)\s+(\d{1,3})%\s+(\S+)\s+(\d+):([0-5]?\d):([0-5]?\d)/; export function parseRsyncProgress(line: string): RsyncProgress | null { const m = line.match(RSYNC_PROGRESS_RE); if (!m) return null; // Grouping separators vary with the locale rsync was built against; strip // everything that is not a digit rather than assuming a comma. const bytes = Number.parseInt(m[1].replace(/\D/g, ""), 10); const percent = Number.parseInt(m[2], 10); if (!Number.isFinite(bytes) || !Number.isFinite(percent)) return null; const etaSeconds = Number.parseInt(m[4], 10) * 3600 + Number.parseInt(m[5], 10) * 60 + Number.parseInt(m[6], 10); return { bytes, percent, rate: m[3], etaSeconds }; } // Is this line ONLY a progress redraw? Used by the relocate controller to keep // several thousand carriage-return redraws out of a job log that an operator // reads afterwards — the parsed decile lines say the same thing in eleven lines. export function isRsyncProgressLine(line: string): boolean { return parseRsyncProgress(line) !== null; } // The one wording of a copy's progress, shared by the job task's detail string // and the decile line in the log — so the /jobs row and the log cannot word the // same instant differently. // // "12.3 GB of 45.6 GB · 27 % · 110.50MB/s · ETA 5:32" // // `totalBytes` is the measured source tree, not rsync's running total. export function formatRsyncProgressDetail( p: RsyncProgress, totalBytes: number, formatBytes: (n: number) => string, ): string { const pct = totalBytes > 0 ? Math.round(clamp01(p.bytes / totalBytes) * 100) : p.percent; const bits = [ `${formatBytes(p.bytes)} of ${totalBytes > 0 ? formatBytes(totalBytes) : "?"}`, `${pct} %`, p.rate, ]; if (p.etaSeconds > 0) bits.push(`ETA ${formatClock(p.etaSeconds)}`); return bits.join(" · "); } // 0..1 against the MEASURED tree, with rsync's own percentage as the fallback // for a caller that never measured one. export function rsyncProgressFraction( p: RsyncProgress, totalBytes: number, ): number { return totalBytes > 0 ? clamp01(p.bytes / totalBytes) : clamp01(p.percent / 100); }