import type { JobTaskKind } from "./registry"; import type { JobRunContext } from "./streamCommand"; import { createDiarizeProgressParser, createDownloadProgressParser, } from "./progressParsers"; import { getTranscriptionApp } from "../lib/transcriptionApps"; // A handle for one in-flight sub-operation. `onLog` is a drop-in replacement // for the controller's existing `onLog`: it forwards every line to the shared // job log AND feeds it to the matching progress parser, updating the task's // per-operation progress bar. `end()` prunes the task (call it in a finally so // failures don't leak a task into record.tasks). export type TaskHandle = { onLog: (line: string) => void; // Push a parsed progress update directly (bypassing the text parser). Used by // remote workers, which stream already-parsed {fraction, detail} back rather // than raw engine output. No-op when the caller didn't opt into tracking. update: (patch: { fraction?: number; detail?: string }) => void; end: () => void; }; export type TaskTracker = { start: (init: { id: string; label: string; kind: JobTaskKind; // For transcribe tasks: the engine (app id) running this task, used to pick // the right progress parser. A remote worker streams pre-parsed progress, so // it passes no appId and the tracker uses a pass-through. appId?: string; // The worker that owns this task, surfaced on the Workers page. workerId?: string; // The video's duration. Recorded against the task's measured wall clock so // the job can report SECONDS PER AUDIO-HOUR — the unit a digest sweep is // estimated in, and one a task-count average cannot express. Omit where the // work is not proportional to length (a download is bytes, not minutes). audioSeconds?: number; }) => TaskHandle; }; type TaskCtx = Pick< JobRunContext, "addTask" | "updateTask" | "removeTask" | "recordTaskDone" >; // Build a tracker from a job's run context. `forwardLog` is the shared `onLog` // the controller already writes to. When `ctx` is undefined (a caller that // doesn't opt in) the tracker degrades to plain forwarding. export function makeTaskTracker( ctx: TaskCtx | undefined, forwardLog: (line: string) => void, ): TaskTracker { return { start({ id, label, kind, appId, workerId, audioSeconds }) { if (!ctx) { return { onLog: forwardLog, update: () => {}, end: () => {} }; } const startedAt = Date.now(); ctx.addTask({ id, label, kind, startedAt, workerId }); // Transcription progress output is app-specific (whisper's segment // timestamps vs chough's ETA bars), so pick the parser for THIS task's // worker engine. A transcribe task with no appId (e.g. a remote worker that // streams pre-parsed progress) falls through to the download parser, which // is harmless for non-matching lines. const transcribeParser = 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. // // Digest tasks get NO parser: a digest's per-chunk progress is discrete // (chunk 3 of 7), not a byte/second rate, and its log lines carry model // token counts that the download parser would happily misread as a // percentage. The task row shows an indeterminate bar instead, which is // honest. const downloadParser = transcribeParser || kind === "digest" ? null : createDownloadProgressParser(); // Backfill tasks are the diarization lane (and the attribution lanes, // which produce no progress output and simply never match). Chained // AHEAD of the download parser rather than replacing it: a backfill task // really can download, when allowRedownload re-acquires media for a // missing-input video, and that phase's progress is worth keeping. const diarizeParser = kind === "backfill" ? createDiarizeProgressParser() : null; let ended = false; const onLog = (line: string) => { forwardLog(line); // yt-dlp rewrites the progress line in place with carriage returns, so // split on both \r and \n to see each discrete update. for (const part of line.split(/[\r\n]+/)) { if (!part) continue; // Both parsers are null for a digest task (see above), so this must // tolerate having neither. The `!` this replaces was a lie that threw // on the FIRST log line of every digest job — the batch caught it as a // per-video failure, so a whole sweep would have reported "0 // generated, N failed" while looking like an engine problem. const update = transcribeParser ? transcribeParser.feed(part) : (diarizeParser?.feed(part) ?? downloadParser?.feed(part)); if (update) ctx.updateTask(id, update); } }; const update = (patch: { fraction?: number; detail?: string }) => { ctx.updateTask(id, patch); }; const end = () => { if (ended) return; ended = true; ctx.recordTaskDone(Date.now() - startedAt, audioSeconds); ctx.removeTask(id); }; return { onLog, update, end }; }, }; }