Archilyzer · Source

archilyzer

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

commit 762642b24ae0f1efcbd96b6b4824412745bd062a
parent 2c13c9da898c003820e64b604594a206fd4ad286
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Fri,  9 Oct 2026 11:56:03 -0400

transcribe-file: durationMs is the engine's time (waitedMs the queue); a short file waits in the urgent tier

`pnpm ops transcribe`'s durationMs ran from the job's start, so it counted
the ffmpeg cut and the wait for a free worker — on a box whose GPU worker was
busy with the lane, minutes that read as a slow engine. It now runs from the
moment a worker takes the job (onWorker; the last attempt's, after a transport
retry) to the transcript, and the new `waitedMs` is the queue. The job's log
says both.

The pool already parks the lane's units (background) behind a manual
acquire, so a file transcription was ahead of the lane — but FIFO behind a
manual batch's next video. A cut of up to 15 minutes of audio
(URGENT_MAX_AUDIO_SEC, measured from the cut WAV's size before the acquire)
now asks the pool's existing "urgent" tier, ahead of every parked waiter; a
longer file keeps "foreground". The tier orders waiters only: nothing
running is interrupted. transcribeWithWorker takes a `tier` and passes it
to pool.acquire; no other caller changes.

Tests: transcribeFile 21 → 24 (the tier from the cut's length; waitedMs vs
durationMs with the GPU slot held; a file job parked LAST takes the freed
slot before a parked manual and a parked lane acquire — and fails with the
tier forced to "foreground").

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

Diffstat:
Mcommon/controller/transcribeFile.test.ts | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/transcribeFile.ts | 52+++++++++++++++++++++++++++++++++++++++++++++-------
Mcommon/controller/transcribeOne.ts | 5+++++
Meditor/CHANGELOG.md | 1+
Meditor/app/api/ops/transcribe/route.ts | 3++-
5 files changed, 144 insertions(+), 8 deletions(-)

diff --git a/common/controller/transcribeFile.test.ts b/common/controller/transcribeFile.test.ts @@ -10,6 +10,7 @@ import { symlink, writeFile, } from "node:fs/promises"; +import { existsSync, statSync } from "node:fs"; import os from "node:os"; import path from "node:path"; import type { Worker } from "../lib/workers"; @@ -130,7 +131,10 @@ const { enqueueTranscribeFile, offsetCues, parseTranscribeFileBody, + transcribeFileTier, transcribeWorkerFilter, + URGENT_MAX_AUDIO_SEC, + wavSeconds, windowOf, windowWavArgs, wordsFromTranscript, @@ -461,3 +465,90 @@ test("the guards answer before any job exists", async () => { assert.match((await runJob({ path: MEDIA, workerId: "off" })).error!, /"off" is disabled/); assert.equal(await readFile(SETTINGS_FILE, "utf8"), SETTINGS_TEXT); }); + +// --- the wait, and the tier ---------------------------------------------------- + +test("a cut's length decides its tier: up to 15 minutes of audio is urgent", () => { + assert.equal(wavSeconds(44), 0); + assert.equal(wavSeconds(44 + 32_000 * 90), 90); + assert.equal(URGENT_MAX_AUDIO_SEC, 900); + assert.equal(transcribeFileTier(0.1), "urgent"); + assert.equal(transcribeFileTier(900), "urgent"); + assert.equal(transcribeFileTier(901), "foreground"); +}); + +const { getWorkerPool } = await import("../jobs/workerPool"); +const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); +const onGpu = (w: Worker) => w.id === "gpu"; + +// Enqueue a job and resolve once it is PARKED on the pool (its log says it is +// waiting for a worker) — a fixed sleep raced the fake ffmpeg's start on a +// loaded machine. `finish` reads the result as runJob does. +async function parkedJob(body: Record<string, unknown>) { + const res = await enqueueTranscribeFile(body, { paths }); + if (!res.ok) throw new Error(`refused: ${res.error}`); + void res.stream.cancel(); + const logFile = path.join(paths.jobsDir, `${res.jobId}.log`); + const t0 = Date.now(); + while (!/Waiting for worker/.test(await readFile(logFile, "utf8").catch(() => ""))) { + if (Date.now() - t0 > 20_000) throw new Error("the job never parked"); + await sleep(20); + } + const finish = async () => { + const done = await res.done; + const log = await readFile(logFile, "utf8"); + const line = log.split("\n").find((l) => l.startsWith(TRANSCRIBE_RESULT_MARKER)); + return { + status: done.status, + log, + result: line ? JSON.parse(line.slice(TRANSCRIBE_RESULT_MARKER.length)) : null, + }; + }; + return { finish }; +} + +test("durationMs is the engine's time; the wait for a busy worker is waitedMs", async () => { + const pool = getWorkerPool(); + // Something else holds the GPU worker's one slot. + const held = await pool.acquire(undefined, { only: onGpu }); + const job = await parkedJob({ path: MEDIA, workerId: "gpu" }); + await sleep(500); + const released = Date.now(); + held.release(); + const run = await job.finish(); + assert.equal(run.status, "done", run.log); + const r = run.result; + assert.ok(r.waitedMs >= 450, `waitedMs ${r.waitedMs}`); + assert.ok(r.durationMs >= 0, `durationMs ${r.durationMs}`); + // The engine's clock starts when the worker is taken — after the release. + assert.ok(r.durationMs <= Date.now() - released, `durationMs ${r.durationMs} counts the wait`); + assert.match(run.log, /ahead of queued transcriptions/); +}); + +test("a short file goes ahead of parked transcriptions — the lane's and a manual batch's", async () => { + const pool = getWorkerPool(); + const held = await pool.acquire(undefined, { only: onGpu }); + const t0 = Date.now(); + await sleep(20); // so an engine run after this is visibly newer than t0 + // Did the file job's engine run before this slot was granted? + const engineRan = () => existsSync(ENGINE_ARGS) && statSync(ENGINE_ARGS).mtimeMs >= t0; + const order: string[] = []; + // Parked first: an auto-lane unit (background) and a manual batch's next + // video (foreground), both waiting for the GPU worker. + const lane = pool.acquire(undefined, { background: true, only: onGpu }).then((l) => { + order.push(`lane${engineRan() ? " after the file" : ""}`); + l.release(); + }); + const batch = pool.acquire(undefined, { only: onGpu }).then((l) => { + order.push(`batch${engineRan() ? " after the file" : ""}`); + l.release(); + }); + const job = await parkedJob({ path: MEDIA, workerId: "gpu" }); + held.release(); + const run = await job.finish(); + await Promise.all([lane, batch]); + assert.equal(run.status, "done", run.log); + // The file job, parked LAST, took the freed slot first; the two parked + // before it got it after, in their own order (manual before the lane). + assert.deepEqual(order, ["batch after the file", "lane after the file"]); +}); diff --git a/common/controller/transcribeFile.ts b/common/controller/transcribeFile.ts @@ -54,6 +54,7 @@ import { writeJsonAtomic } from "../lib/jsonFile-server"; import { getSettings } from "../lib/settings"; import { getWorkerPool, type WorkerFilter } from "../jobs/workerPool"; import { makeTaskTracker } from "../jobs/taskHooks"; +import type { SchedulerTier } from "../jobs/jobKinds"; import { runManagedFunction, type JobRunContext, @@ -80,6 +81,27 @@ export const TRANSCRIBE_FILE_BODY_KEYS = [ const AUDIO_NAME = "audio.wav"; // A WAV header with no samples after it: the window held no audio. const EMPTY_WAV_BYTES = 44; +// The cut is 16 kHz mono s16 (windowWavArgs): 32,000 bytes a second. +const WAV_BYTES_PER_SEC = 16_000 * 2; + +// A file transcription at most this long (seconds of audio) waits for a worker +// in the pool's "urgent" tier, ahead of every parked transcription — the +// lane's background units and a manual batch's next video alike. It is the +// quote check an agent is waiting on; one more long transcription in the +// queue is not. Longer files keep the manual ("foreground") tier: still ahead +// of the lane, behind batches queued before them. The tier only orders +// WAITERS: a transcription already running is never interrupted. +export const URGENT_MAX_AUDIO_SEC = 15 * 60; + +/** Seconds of audio in the cut WAV, from its size. */ +export function wavSeconds(bytes: number): number { + return Math.max(0, bytes - EMPTY_WAV_BYTES) / WAV_BYTES_PER_SEC; +} + +/** The worker-pool tier a file transcription of `audioSec` seconds asks for. */ +export function transcribeFileTier(audioSec: number): SchedulerTier { + return audioSec <= URGENT_MAX_AUDIO_SEC ? "urgent" : "foreground"; +} export type TranscribeFileRequest = { path: string; @@ -112,7 +134,11 @@ export type TranscribeFileResult = { worker: TranscribeWorkerInfo; transcriptFormat: TranscriptOutputFormat; transcribedAt: string; + // The engine's own time: from the moment a worker took the job to the + // transcript. The wait for a free worker is `waitedMs`, never in here. durationMs: number; + // How long the job waited for a free worker (the pool's queue). + waitedMs: number; cues: Cue[]; text: string; // Present only when the request asked for words: [] when the engine has none. @@ -441,7 +467,6 @@ export async function runTranscribeFile( opts: RunTranscribeFileOpts, ): Promise<TranscribeFileResult> { const { request: req, onLog, paths } = opts; - const started = Date.now(); const scratch = await mkdtemp(path.join(os.tmpdir(), "archilyzer-transcribe-")); try { const wav = path.join(scratch, AUDIO_NAME); @@ -465,12 +490,18 @@ export async function runTranscribeFile( } // Set by onWorker; a holder, so the closure's write is seen after the await. - const used: { worker?: Worker } = {}; + // `startedAt` is the moment a worker took it: the engine's clock starts + // there (the last attempt's, when a transport failure moved it). + const used: { worker?: Worker; startedAt?: number } = {}; + const audioSec = wavSeconds(wavStat.size); + const tier = transcribeFileTier(audioSec); onLog( - req.workerId - ? `Waiting for worker ${req.workerId}…` - : "Waiting for a free local worker…", + `${req.workerId ? `Waiting for worker ${req.workerId}` : "Waiting for a free local worker"}` + + (tier === "urgent" + ? ` (${Math.round(audioSec)}s of audio: ahead of queued transcriptions)…` + : "…"), ); + const asked = Date.now(); const label = `${path.basename(req.path)}${windowOf(req) ? " (window)" : ""}`; const outcome = await transcribeWithWorker({ paths, @@ -484,12 +515,15 @@ export async function runTranscribeFile( onLog, signal: opts.signal, only: transcribeWorkerFilter(req.workerId), + tier, onWorker: (w) => { used.worker = w; + used.startedAt = Date.now(); }, skipInlineDiarization: true, ...(req.words ? { words: true } : {}), }); + const finished = Date.now(); const worker = used.worker; if (outcome !== "transcribed" || !worker) { throw new Error( @@ -508,7 +542,8 @@ export async function runTranscribeFile( worker: describeTranscribeWorker(worker), transcriptFormat, transcribedAt: new Date().toISOString(), - durationMs: Date.now() - started, + durationMs: finished - (used.startedAt ?? asked), + waitedMs: (used.startedAt ?? asked) - asked, cues, text: cues.map((c) => c.text.trim()).filter(Boolean).join(" "), ...(req.words ? { words: wordsFromTranscript(raw, req.start ?? 0) } : {}), @@ -558,7 +593,10 @@ export async function enqueueTranscribeFile( onLog( `Transcribed with ${result.worker.appId} [${result.worker.id}]` + `${result.worker.model ? ` model ${result.worker.model}` : ""}: ` + - `${result.cues.length} cue(s) in ${(result.durationMs / 1000).toFixed(1)}s`, + `${result.cues.length} cue(s) in ${(result.durationMs / 1000).toFixed(1)}s` + + (result.waitedMs >= 1000 + ? ` (after ${(result.waitedMs / 1000).toFixed(1)}s waiting for the worker)` + : ""), ); if (req.out) { await writeJsonAtomic(req.out, result, { indent: 2 }); diff --git a/common/controller/transcribeOne.ts b/common/controller/transcribeOne.ts @@ -11,6 +11,7 @@ import { type WorkerFilter, } from "../jobs/workerPool"; import type { TaskTracker } from "../jobs/taskHooks"; +import type { SchedulerTier } from "../jobs/jobKinds"; import { normalizeTranscript } from "./normalizeTranscript"; import { diarizeOneVideo } from "./diarizeOne"; import { getSettings } from "../lib/settings"; @@ -376,6 +377,9 @@ export type TranscribeWithWorkerOptions = { // Auto-runner units pass true so they park BEHIND any manual (foreground) // acquire in the worker pool — a manual transcribe preempts queued auto work. background?: boolean; + // The pool tier outright, over `background`: a short one-off file + // transcription asks "urgent", ahead of every parked batch (transcribeFile). + tier?: SchedulerTier; // Narrows WHICH workers may take this video (the pool's `only`): a one-off // file transcription keeps to local workers, or to the one it was told to use. only?: WorkerFilter; @@ -411,6 +415,7 @@ export async function transcribeWithWorker( try { lease = await pool.acquire(acquireSignal, { background: opts.background, + ...(opts.tier ? { tier: opts.tier } : {}), ...(opts.only ? { only: opts.only } : {}), }); } catch (err) { diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] +- **A file transcription is timed from the engine, and a short one does not queue behind batches.** `pnpm ops transcribe`'s `durationMs` is now the engine's own time — from the moment a worker took the job — and the new `waitedMs` is how long it waited for a free worker; the job's log says both. A file of up to 15 minutes of audio (a window, usually) waits in the worker pool ahead of every parked transcription, a manual batch's next video included, not only ahead of the transcription lane; a longer file keeps its place behind batches queued before it. Nothing running is interrupted. Needs a restart of the editor. - **Heavy work takes turns, above a memory floor.** A publish stage's `next build` — a site's, the hub's, the homepage's — now waits for the machine's one heavy slot, which every e2e run and any `pnpm heavy -- <cmd>` (a video render) take too, and then until at least 6000 MB is available; the stage's log says whom it waits behind ("waiting for the heavy slot — held by …") or how much memory there is ("waiting for memory — 4210 MB available, the floor is 6000 MB"). Cancel still stops it. `HEAVY_MIN_FREE_MB` moves the floor (`0` turns it off) and `HEAVY=0` skips the gate. Needs a restart of the editor. - **A curated tag can exist on some sites only.** A tag's new **Sites** field on /tags (`sites` in `transcripts/tags.json`; `pnpm ops tags` takes it in a define) names the sites it exists on. Its rules then fire, and its pins apply, only to videos on those sites' channels, and every other site drops it from its records, its counts and its `/tags.json` — where **Hidden** only hid the chip. Empty is every site, as before. Setting it, or changing the channels of those sites, re-derives the corpus's tags once at the next index update. The Eva tags are what this is for: they belong on Anilyzer alone. - **The publish lane.** Publishing can run itself: turn it on at **/operations/publish** (the runner's Start, Drain and Stop, the hold, and the lane's settings; or `publish.enabled` in settings) and the lane checks every `checkEveryMinutes` (10) whether the index is stale; when it is — and its last update is at least `refreshEveryMinutes` (360) old — it updates it, then builds every site whose channels changed or whose data the new index moved, one stage at a time on the `publish` queue. What it may do with a site is the site's own — the **Publish policy** on the site's settings form, `site.json` `publish.auto` —: `off` (the default: left alone), `build`, `preview` (built and deployed to the preview branch `publish.previewBranch`) or `production`; the hub and the homepage have `publish.hub` and `publish.homepage`. A private site is only ever built, and a site needs its Cloudflare Pages project before it may deploy. Hold the lane and the stage running finishes and no next one starts; quiet hours (`publish.quietHours`) do the same; Drain finishes the stage and ends the runner. The lane never forces a stage: a stage that finds its target current does nothing. On /jobs every stage of one run reads `run <id> · <target>`, and a stage still queued when the editor restarts is cancelled, never re-queued — the lane works out again what is stale from what is on disk. `archilyzer publish now` runs the same plan from the command line, one stage after another in its own process. diff --git a/editor/app/api/ops/transcribe/route.ts b/editor/app/api/ops/transcribe/route.ts @@ -17,7 +17,8 @@ export const dynamic = "force-dynamic"; // // The job's log ends with the result as ONE line, `@@transcribe-result // {json}`: `{version, path, window, worker: {id, name, appId, model, device}, -// transcriptFormat, transcribedAt, durationMs, cues: [{start, end, text}], +// transcriptFormat, transcribedAt, durationMs (the engine's time), waitedMs (the +// wait for a free worker), cues: [{start, end, text}], // text}`, cue times on the SOURCE file's clock. `pnpm ops transcribe --wait` // prints that JSON on stdout. `out` writes it to that file as well. //