Archilyzer · Source

archilyzer

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

commit 2898c1f6816cb53b882bdac851e106af7420dc4e
parent 68fc72000590eee36090f4206d0696f261a29f70
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri, 12 Jun 2026 21:46:04 -0400

ETA

Diffstat:
Mcommon/jobs/registry.ts | 16++++++++++++++++
Mcommon/jobs/streamCommand.ts | 2++
Mcommon/jobs/taskHooks.ts | 9+++++++--
Meditor/CHANGELOG.md | 1+
Meditor/app/jobs/active/buildActiveJobs.ts | 32+++++++++++++++++++++++++++++---
Meditor/app/jobs/components/RunningJobsList.tsx | 18+++++++++++++++++-
Meditor/e2e/jobs-batch-tasks-drain.spec.ts | 32++++++++++++++++++++++++++++++++
7 files changed, 104 insertions(+), 6 deletions(-)

diff --git a/common/jobs/registry.ts b/common/jobs/registry.ts @@ -51,6 +51,12 @@ export type JobRecord = { draining?: boolean; progress?: JobProgress; tasks?: JobTask[]; + // Running totals over completed sub-operations, used to estimate time + // remaining (remaining tasks × average measured task duration). Accumulated + // as tasks finish — unlike `tasks`, these are NOT cleared, so the average + // reflects the whole batch. Both undefined until the first task completes. + completedTaskCount?: number; + completedTaskMs?: number; }; export type QueueSnapshot = { @@ -234,6 +240,16 @@ class JobRegistry { if (idx >= 0) job.tasks.splice(idx, 1); } + // Fold one finished sub-operation's wall-clock duration into the job's + // running totals (see JobRecord.completedTaskCount/Ms). No-op if the job is + // gone or the duration is nonsensical. + recordTaskDuration(jobId: string, durationMs: number): void { + const job = this.jobs.get(jobId); + if (!job || durationMs < 0) return; + job.completedTaskCount = (job.completedTaskCount ?? 0) + 1; + job.completedTaskMs = (job.completedTaskMs ?? 0) + durationMs; + } + // Snapshot of queues for UI. Sorted by queue name. listQueues(): QueueSnapshot[] { const out: QueueSnapshot[] = []; diff --git a/common/jobs/streamCommand.ts b/common/jobs/streamCommand.ts @@ -43,6 +43,7 @@ export type JobRunContext = { patch: { fraction?: number; detail?: string }, ) => void; removeTask: (taskId: string) => void; + recordTaskDone: (durationMs: number) => void; }; export type RunManagedFunctionOpts = CommonOpts & { @@ -232,6 +233,7 @@ export async function runManagedFunction( addTask: (task) => registry.addTask(id, task), updateTask: (taskId, patch) => registry.updateTask(id, taskId, patch), removeTask: (taskId) => registry.removeTask(id, taskId), + recordTaskDone: (ms) => registry.recordTaskDuration(id, ms), }; opts diff --git a/common/jobs/taskHooks.ts b/common/jobs/taskHooks.ts @@ -18,7 +18,10 @@ export type TaskTracker = { start: (init: { id: string; label: string; kind: JobTaskKind }) => TaskHandle; }; -type TaskCtx = Pick<JobRunContext, "addTask" | "updateTask" | "removeTask">; +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 @@ -32,7 +35,8 @@ export function makeTaskTracker( if (!ctx) { return { onLog: forwardLog, end: () => {} }; } - ctx.addTask({ id, label, kind, startedAt: Date.now() }); + const startedAt = Date.now(); + ctx.addTask({ id, label, kind, startedAt }); // Transcription progress output is app-specific (whisper's segment // timestamps vs chough's ETA bars), so pick the active app's parser. const transcribeParser = @@ -55,6 +59,7 @@ export function makeTaskTracker( const end = () => { if (ended) return; ended = true; + ctx.recordTaskDone(Date.now() - startedAt); ctx.removeTask(id); }; return { onLog, end }; diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **Batch jobs estimate the time remaining.** A running batch's overall progress bar on `/jobs/active` now shows an estimate of how long is left (e.g. `~4:30 left`), alongside the existing `Transcripts: 12 / 50` count. The estimate is **remaining tasks × the measured average time per task**: each download/transcription's real wall-clock duration is folded into per-job running totals as it finishes, and that average is converted to a wall-clock figure using the parallelism observed so far (so a 4-way parallel transcribe batch isn't estimated as if it ran one at a time). It reads as **estimating…** until the first sub-operation completes (no average yet), and disappears once no work remains. Jobs without per-task tracking (e.g. storing a playlist) show no estimate. - **Audio-checked downloads recover correctly when an interrupted attempt left a malformed `.part` and the video id isn't derivable from its URL.** For sources where the canonical id only appears after metadata (e.g. Odysee), the integrity-checked downloader's pre-check would correctly roll back a corrupt leftover `.part` to its `.good` snapshot (or discard it), but the post-launch file discovery then ignored that same directory as a "pre-existing" one — so the freshly re-downloaded audio was never found and the download was recorded as failed. The pre-check now reports the directory it acted on, and discovery scopes to it, so the resumed download finalizes as `ok-audio-checked`. Unrelated stale `.part`s from other videos' interrupted attempts are still ignored (clean pre-existing parts aren't reported), so the cross-video protection is unchanged. - **First-class support for multiple transcription apps (chough + whisper.cpp).** Transcription is no longer hardcoded to whisper.cpp. **Settings → Transcription** now has an **App** dropdown (whisper.cpp / chough) with per-app fields, replacing the old flat Binary / Model / Args command. Each app owns how it builds its command line, what file it writes, how its output JSON is parsed, and how its progress output is read — so adding another tool is a small code module (`common/lib/transcriptionApps.ts`). **chough** is supported in both **local** and **remote** modes (set a **Remote URL** to transcribe via a `chough --server`, passing `CHOUGH_URL`; leave it blank for local), with optional **Chunk size** (`-c`) and **Model** (`CHOUGH_MODEL`) fields; its binary defaults to the `CHOUGH_BIN` env var. whisper.cpp keeps its Binary / Model / custom-args template. Because chough writes its output to the exact `-o` path (no `.json` appended, unlike whisper-cli's `-of`), the runner now renames the app-declared output file — fixing transcripts that previously failed to materialize under chough. Transcript parsing is **format-aware and back-compatible**: existing whisper.cpp `transcript.json` files and new chough files coexist, each parsed correctly by content sniff (chough's seconds-based `chunk_data` vs whisper's millisecond `transcription` offsets), with the detected format recorded per video in `transcript.cues.json` (`transcriptFormat`) and a fallback to whisper.cpp for anything unrecognized — so re-indexing a mixed corpus (including across shard machines running different tools) just works. An existing `settings.json` migrates automatically: legacy `transcribeBin`/`transcribeArgs`/`transcribeModel` map onto the matching app (a binary named `chough` adopts the chough app; everything else becomes whisper.cpp, preserving a customized args template). - **Global default for parallel transcriptions.** A new **Settings → Parallel transcriptions** field sets how many videos a "Transcribe missing" / bucket run transcribes at once when its per-run Concurrency input is left blank (default **2**, clamped 1–16). This replaces the old `PARALLEL_TRANSCRIBE_LIMIT` env-var default of 4 as the source of the default: the value is now persisted in `settings.json`, shown as the placeholder in each channel's Concurrency input, and used as the server-side fallback. The per-run Concurrency input still overrides it for a single run. diff --git a/editor/app/jobs/active/buildActiveJobs.ts b/editor/app/jobs/active/buildActiveJobs.ts @@ -1,7 +1,7 @@ import { getRegistry, isDrainableKind, - type JobProgress, + type JobRecord, } from "yt-dlp-transcript-common/jobs/registry"; import { readChannelStat, @@ -15,10 +15,33 @@ export type ActiveJobsPayload = { channels: { slug: string; displayName: string }[]; }; +// Estimate seconds remaining as: remaining tasks × average measured task +// duration, converted to wall-clock via the effective concurrency observed so +// far. Returns undefined until at least one task has completed (no average +// yet) or when there's no work left. `now` is injected so the same snapshot is +// used across all jobs in a single payload build. +function computeEtaSeconds( + job: JobRecord, + remaining: number, + now: number, +): number | undefined { + const count = job.completedTaskCount ?? 0; + const totalMs = job.completedTaskMs ?? 0; + if (count < 1 || remaining <= 0 || job.startedAt === undefined) { + return undefined; + } + 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; +} + function computeJobProgressView( - snap: JobProgress | undefined, + job: JobRecord, stat: ChannelStat | undefined, + now: number, ): RunningJobsListItem["progress"] { + const snap = job.progress; if (!snap || !stat) return undefined; const current = snap.metric === "downloads" ? stat.downloadCount : stat.transcriptCount; @@ -26,12 +49,14 @@ function computeJobProgressView( const advance = Math.max(0, current - snap.initial); const pct = range === 0 ? 100 : Math.max(0, Math.min(100, (advance / range) * 100)); + const remaining = Math.max(0, snap.target - current); return { metric: snap.metric, initial: snap.initial, current, target: snap.target, pct, + etaSeconds: computeEtaSeconds(job, remaining, now), }; } @@ -66,6 +91,7 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { if (stat) channelStats.set(channelSlugs[i], stat); } + const now = Date.now(); const jobs: RunningJobsListItem[] = activeRecords.map((j) => ({ id: j.id, kind: j.kind, @@ -75,7 +101,7 @@ export async function buildActiveJobsPayload(): Promise<ActiveJobsPayload> { videoId: j.videoId, progress: j.status === "running" && j.channelSlug - ? computeJobProgressView(j.progress, channelStats.get(j.channelSlug)) + ? computeJobProgressView(j, channelStats.get(j.channelSlug), now) : undefined, tasks: j.tasks?.map((t) => ({ id: t.id, diff --git a/editor/app/jobs/components/RunningJobsList.tsx b/editor/app/jobs/components/RunningJobsList.tsx @@ -2,6 +2,7 @@ import Link from "next/link"; import { useState } from "react"; +import { formatDuration } from "yt-dlp-transcript-common/lib/format"; import { JobLogTail } from "../[id]/components/JobLogTail"; import { DrainJobButton } from "./DrainJobButton"; import { CancelJobButton } from "./CancelJobButton"; @@ -27,6 +28,9 @@ export type RunningJobsListItem = { current: number; target: number; pct: number; + // Estimated seconds remaining (remaining tasks × measured avg per task). + // Undefined until the first task completes, or when no work remains. + etaSeconds?: number; }; tasks?: RunningJobsTask[]; draining?: boolean; @@ -201,6 +205,15 @@ function JobProgressBar({ progress.metric === "downloads" ? `Downloads: ${progress.current} / ${progress.target}` : `Transcripts: ${progress.current} / ${progress.target}`; + // Append an ETA once the batch has a measured average. formatDuration returns + // "" for 0/falsy, so guard against printing a bare "·". + const remaining = progress.target - progress.current; + const etaText = + remaining <= 0 + ? null + : typeof progress.etaSeconds === "number" + ? `~${formatDuration(Math.max(1, Math.round(progress.etaSeconds)))} left` + : "estimating…"; const fillClass = progress.metric === "downloads" ? "bg-emerald-300 dark:bg-emerald-700" @@ -220,7 +233,10 @@ function JobProgressBar({ style={{ width: `${progress.pct}%` }} /> </div> - <div className="text-xs text-zinc-500">{label}</div> + <div className="flex items-baseline justify-between gap-2 text-xs text-zinc-500"> + <span>{label}</span> + {etaText ? <span className="shrink-0">{etaText}</span> : null} + </div> </div> ); } diff --git a/editor/e2e/jobs-batch-tasks-drain.spec.ts b/editor/e2e/jobs-batch-tasks-drain.spec.ts @@ -77,6 +77,38 @@ test("active jobs shows a per-operation progress bar that advances", async ({ .toBeGreaterThan(0); }); +test("active jobs shows an estimated time remaining once a task completes", async ({ + page, +}) => { + test.setTimeout(90_000); + await resetData(null); + // More videos than the default concurrency of 4, so some tasks finish while + // others are still in flight — i.e. completed > 0 and remaining > 0, the + // window in which an ETA can be shown. + const ids = ["slowope1", "slowope2", "slowope3", "slowope4", "slowope5", "slowope6"]; + await makeTranscribeChannel("eta-one", "Eta One", ids); + await invalidateCache(); + + await page.goto("/channels/eta-one"); + await page.getByRole("button", { name: "Transcribe missing" }).click(); + await page.goto("/jobs/active"); + + const section = page.locator("section[aria-label='Active jobs for Eta One']"); + // The batch-level progress bar is present (proves setProgress ran). + await expect(section.getByText(/Transcripts: \d+ \/ 6/)).toBeVisible({ + timeout: 15_000, + }); + + // Before any task completes the average is unknown → "estimating…"; once one + // finishes a concrete "~m:ss left" estimate appears. Poll for the concrete + // form (it implies we passed through, or skipped straight past, estimating). + await expect + .poll(async () => (await section.innerText()).replace(/\s+/g, " "), { + timeout: 60_000, + }) + .toMatch(/~\d+:\d{2} left/); +}); + test("draining a batch finishes in-flight work, skips the rest, and releases the queue", async ({ page, }) => {