Archilyzer · Source

archilyzer

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

commit b975a7ba447882eef0c465da20cdc865e97fbfc0
parent 62b4cfc1753d1302b4518ba6deedfb497bc15b1b
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Wed, 29 Jul 2026 19:51:07 -0400

Price the sweep in chunks, and stop yielding to CPU workers

Three cost estimates disagreed 27 / 90 / 151 s per audio-hour and one of
them was in the shipped plan. They were never in conflict: a CHUNK is one
model call, chunk density varies 4x across the corpus, and re-pricing all
three per chunk reproduces them to within 2%. The corpus is 191,116
chunks (computed free from statsByPath.cueCount, present on all 73,367
transcribed videos), so the honest range is ~25 days idle to ~55
contended, not 81. digest-plan now reports 53.5.

Also fixes a bug I shipped yesterday. transcriptionActivity() yielded the
GPU for any busy `kind: "local"` worker, but this box runs one GPU worker
beside two pinned to `device: "cpu"` — so at parallelTranscriptions 2 the
digest lane stopped dead for transcription competing for zero shaders.
The queued-job signal had to be narrowed too or the device check would
have been dead code: it only covers acquisition GAPS now, when no local
worker is busy at all. New setting digest.yieldToCpuWorkers, default off;
absence falls to off because for an old settings file that is the fix,
not a change of intent.

And a bug found on the way: the settings form rebuilt the digest block
field-by-field behind an `as` cast, so saving any unrelated setting
DISARMED AN ARMED CORPUS SWEEP. Spread the current block and drop the
cast, so the next added field is a type error instead.

ollama's own load/prefill/decode timings were being discarded; they are
recorded and logged now, wall beside engine. Nothing reaches
DigestProvenance, so no digest is invalidated. Verified live: gemma2:9b
OFFLOADS 1,029 MB at 8192 ctx on this 8 GB card while qwen2.5:7b is
fully resident — its 134 re-priced days are partly this card, not the
model.

386 common tests (was 374); tsc clean in common, editor and export.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

Diffstat:
Mcommon/bin/digest-plan.ts | 91+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------
Mcommon/controller/digestPlan.test.ts | 175++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mcommon/controller/digestPlan.ts | 193++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------
Mcommon/controller/digestSweep.ts | 10+++++++---
Mcommon/controller/digestTarget.ts | 30+++++++++++++++++++++++-------
Acommon/controller/digestYield.test.ts | 155+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/controller/digestYield.ts | 124+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Mcommon/jobs/workerPool.ts | 10++++++++++
Mcommon/lib/digestApps.ts | 68++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcommon/lib/settings.ts | 37++++++++++++++++++++++++++++++-------
Mcommon/lib/transcriptWindow.ts | 30++++++++++++++++++++++++++++++
Meditor/app/settings/actions.ts | 126+++++++++++++++++++++++++++++++++++++++++++++++++------------------------------
Meditor/app/settings/components/SettingsForm.tsx | 23+++++++++++++++++++++++
13 files changed, 933 insertions(+), 139 deletions(-)

diff --git a/common/bin/digest-plan.ts b/common/bin/digest-plan.ts @@ -1,15 +1,24 @@ #!/usr/bin/env tsx -// Price the digest backfill, in audio-hours, before spending GPU-weeks on it. +// Price the digest backfill, in CHUNKS, before spending GPU-weeks on it. // // Reads the last duplicate-detection run and build:stats' cache; writes nothing. // Run it before and after a cheapness lever (a --near threshold change, a batch // of confirmed clusters) and diff `generate`: that number, not a video count, is // what the sweep's wall clock is made of. // +// A chunk is one model call and IS the unit of cost. Audio-hours are reported +// alongside because they are what a human thinks in, but chunk density varies 4× +// across the corpus, so a projection built on s/audio-hour is only valid for a +// slice with the corpus-average mix. Quoting one is what made three separate +// measurements look contradictory when per chunk they agreed to within 2%. +// // Flags: // --lane local|remote which engine identity to test freshness against -// --rate N seconds per audio-hour (default 90, the measured local -// rate from the 102-video validation run) +// --per-chunk N seconds per model call (default: the measured +// production rate, MEASURED_SECONDS_PER_CHUNK) +// --rate N LEGACY: seconds per audio-hour. Prices off audio-hours +// instead of chunks, so it can reproduce an old estimate; +// it is corpus-average only and wrong per channel. // --no-freshness skip the per-video sidecar read; reports the // from-scratch cost instead of the remaining one // --channels a,b restrict to these channel slugs @@ -19,9 +28,12 @@ import { getPaths } from "../lib/paths"; import { buildDigestSweepPlan, audioHours, + chunksPerAudioHour, sweepDays, + sweepDaysFromAudio, DIGEST_PLAN_ROLES, MEASURED_SECONDS_PER_AUDIO_HOUR, + MEASURED_SECONDS_PER_CHUNK, roleMustGenerate, } from "../controller/digestPlan"; import { parseFlags } from "./_parseFlags"; @@ -29,19 +41,34 @@ import { parseFlags } from "./_parseFlags"; const flags = parseFlags(process.argv.slice(2)); const lane = flags.lane === "remote" ? "remote" : "local"; -const rate = - flags.rate !== undefined - ? Number(flags.rate) - : MEASURED_SECONDS_PER_AUDIO_HOUR; +// --rate is honoured for reproducing an old audio-hour estimate; --per-chunk (or +// nothing) uses the real unit. +const legacyRate = flags.rate !== undefined ? Number(flags.rate) : null; +const perChunk = + flags["per-chunk"] !== undefined + ? Number(flags["per-chunk"]) + : MEASURED_SECONDS_PER_CHUNK; const top = flags.top !== undefined ? Number(flags.top) : 15; const asJson = flags.json === "true"; +// One projection function for every line below, so the header's stated basis and +// every number under it cannot disagree. +function days(chunks: number, audioSeconds: number): number { + return legacyRate !== null + ? sweepDaysFromAudio(audioSeconds, legacyRate) + : sweepDays(chunks, perChunk); +} + function hours(seconds: number): string { return audioHours(seconds).toLocaleString("en-US", { maximumFractionDigits: 0, }); } +function count(n: number): string { + return n.toLocaleString("en-US", { maximumFractionDigits: 0 }); +} + function pct(part: number, whole: number): string { return whole > 0 ? `${((part / whole) * 100).toFixed(1)}%` : "—"; } @@ -62,10 +89,16 @@ async function main(): Promise<void> { JSON.stringify( { ...plan, - secondsPerAudioHour: rate, + secondsPerChunk: perChunk, + secondsPerAudioHour: legacyRate ?? MEASURED_SECONDS_PER_AUDIO_HOUR, + pricedOn: legacyRate !== null ? "audio-hours" : "chunks", generateAudioHours: audioHours(plan.generateSeconds), sharedAudioHours: audioHours(plan.sharedSeconds), - sweepDays: sweepDays(plan.generateSeconds, rate), + generateChunksPerAudioHour: chunksPerAudioHour( + plan.generateChunks, + plan.generateSeconds, + ), + sweepDays: days(plan.generateChunks, plan.generateSeconds), }, null, 2, @@ -79,7 +112,10 @@ async function main(): Promise<void> { console.log(""); console.log( - `Digest sweep plan — ${lane} lane, ${rate}s per audio-hour` + + `Digest sweep plan — ${lane} lane, ` + + (legacyRate !== null + ? `${legacyRate}s per audio-hour (LEGACY basis)` + : `${perChunk}s per chunk at ${plan.maxCues} cues/chunk`) + (plan.freshnessChecked ? "" : ", FROM SCRATCH (freshness not checked)"), ); console.log( @@ -89,40 +125,55 @@ async function main(): Promise<void> { console.log("!! stats cache is stale — run build:stats."); console.log(""); - console.log("Remaining work by role (audio-hours):"); + console.log("Remaining work by role (chunks and audio-hours):"); for (const role of DIGEST_PLAN_ROLES) { const t = plan.remaining[role]; console.log( - ` ${role.padEnd(17)} ${hours(t.audioSeconds).padStart(8)} h ` + + ` ${role.padEnd(17)} ${count(t.chunks).padStart(9)} chunks ` + + `${hours(t.audioSeconds).padStart(8)} h ` + `${t.videos.toLocaleString().padStart(7)} videos ` + `${roleMustGenerate(role) ? "GENERATE" : "shared free"}`, ); } console.log( - ` ${"already fresh".padEnd(17)} ${hours(plan.fresh.audioSeconds).padStart(8)} h ` + + ` ${"already fresh".padEnd(17)} ${count(plan.fresh.chunks).padStart(9)} chunks ` + + `${hours(plan.fresh.audioSeconds).padStart(8)} h ` + `${plan.fresh.videos.toLocaleString().padStart(7)} videos skipped`, ); console.log( - ` ${"ineligible".padEnd(17)} ${"—".padStart(8)} ${plan.ineligible.toLocaleString().padStart(7)} videos no transcript`, + ` ${"ineligible".padEnd(17)} ${"—".padStart(9)} ${"—".padStart(8)} ` + + `${plan.ineligible.toLocaleString().padStart(7)} videos no transcript`, ); + if (plan.chunksEstimated > 0) + console.log( + ` !! ${plan.chunksEstimated.toLocaleString()} video(s) had no cueCount; their chunks were ESTIMATED from duration.`, + ); console.log(""); console.log( - `TO GENERATE : ${hours(plan.generateSeconds)} audio-hours → ${sweepDays(plan.generateSeconds, rate).toFixed(1)} sweep days`, + `TO GENERATE : ${count(plan.generateChunks)} chunks / ${hours(plan.generateSeconds)} audio-hours ` + + `(${chunksPerAudioHour(plan.generateChunks, plan.generateSeconds).toFixed(2)} chunks/audio-h) ` + + `→ ${days(plan.generateChunks, plan.generateSeconds).toFixed(1)} sweep days`, ); console.log( - `SAVED by sharing: ${hours(plan.sharedSeconds)} audio-hours (${pct(plan.sharedSeconds, eligible)} of eligible) → ${sweepDays(plan.sharedSeconds, rate).toFixed(1)} days avoided`, + `SAVED by sharing: ${count(plan.sharedChunks)} chunks / ${hours(plan.sharedSeconds)} audio-hours ` + + `(${pct(plan.sharedSeconds, eligible)} of eligible) → ` + + `${days(plan.sharedChunks, plan.sharedSeconds).toFixed(1)} days avoided`, ); console.log(""); const listed = top > 0 ? plan.channels.slice(0, top) : plan.channels; + // Chunk density per channel is the column that explains why the sweep's rate + // changes as the queue drains — it runs heaviest-first, and the heaviest + // channels are the CHEAPEST per audio-hour. console.log(`Heaviest channels (the sweep's work queue order):`); for (const c of listed) { - if (c.generateSeconds <= 0) continue; + if (c.generateChunks <= 0) continue; console.log( - ` ${c.channelSlug.padEnd(28)} ${hours(c.generateSeconds).padStart(7)} h generate ` + - `${hours(c.sharedSeconds).padStart(6)} h shared ` + - `${sweepDays(c.generateSeconds, rate).toFixed(1).padStart(5)} d`, + ` ${c.channelSlug.padEnd(28)} ${count(c.generateChunks).padStart(8)} chunks ` + + `${hours(c.generateSeconds).padStart(7)} h ` + + `${chunksPerAudioHour(c.generateChunks, c.generateSeconds).toFixed(2).padStart(5)} c/h ` + + `${days(c.generateChunks, c.generateSeconds).toFixed(1).padStart(5)} d`, ); } const rest = plan.channels.length - listed.length; diff --git a/common/controller/digestPlan.test.ts b/common/controller/digestPlan.test.ts @@ -3,12 +3,19 @@ import assert from "node:assert/strict"; import type { DigestClusterRole } from "./digestSharing"; import { audioHours, + chunksForVideo, + chunksPerAudioHour, classifyDigestRole, + CORPUS_CHUNKS_PER_AUDIO_HOUR, DIGEST_PLAN_ROLES, MEASURED_SECONDS_PER_AUDIO_HOUR, + MEASURED_SECONDS_PER_CHUNK, roleMustGenerate, sweepDays, + sweepDaysFromAudio, } from "./digestPlan"; +import { chunkCuesForContext, countCueChunks } from "../lib/transcriptWindow"; +import { DIGEST_OVERLAP_CUES, maxCuesForContext } from "../lib/digestPrompt"; // Run with: // pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/digestPlan.test.ts @@ -60,14 +67,110 @@ test("exactly one role is free; every other role costs GPU time", () => { assert.deepEqual(free, ["mirror-aligned"]); }); -test("audio-hours and the day projection agree with the measured rate", () => { +test("audio-hours and the legacy audio-hour projection still agree", () => { assert.equal(audioHours(3600), 1); // One audio-hour at 90 s/audio-hour is 90 seconds of wall clock. - assert.equal(sweepDays(3600, 90), 90 / 86400); - // The corpus figure, as a regression pin on the headline: 77,298 audio-hours - // at the measured rate is ~80 days, which is what the whole plan is about. - const days = sweepDays(77_298 * 3600, MEASURED_SECONDS_PER_AUDIO_HOUR); - assert.ok(days > 79 && days < 82, `expected ~80 sweep days, got ${days}`); + assert.equal(sweepDaysFromAudio(3600, 90), 90 / 86400); +}); + +// The unit fix, pinned. A chunk is one model call, and the projection is chunks × +// seconds-per-chunk — NOT audio-hours × seconds-per-audio-hour, which is what made +// 27, 90 and 151 s/audio-hour look like three irreconcilable measurements of the +// same box. +test("the projection is priced per chunk", () => { + assert.equal(sweepDays(86_400, 1), 1); + // The corpus headline, as a regression pin: 191,116 chunks at the measured + // production rate is ~55 days, and the 81 days the audio-hour model claimed was + // wrong on the validation run's OWN data. + const days = sweepDays(191_116, MEASURED_SECONDS_PER_CHUNK); + assert.ok(days > 54 && days < 56, `expected ~55 sweep days, got ${days}`); + // Idle-engine cost from bake-off round 2, re-priced: ~25 days, reproducing the + // 24.2 it originally claimed. Same corpus, same code, different contention. + const idle = sweepDays(191_116, 11.2); + assert.ok(idle > 24 && idle < 26, `expected ~25 idle days, got ${idle}`); + // gemma2's 5.4× — a cost decision, not a judgement call. + assert.ok(sweepDays(191_116, 60.6) / idle > 5); +}); + +// s/audio-hour is kept only as a derived convenience, and it must stay derived: a +// hand-edited second number here is exactly how the unit error happened. +test("the audio-hour rate is derived from the per-chunk cost", () => { + assert.equal( + MEASURED_SECONDS_PER_AUDIO_HOUR, + MEASURED_SECONDS_PER_CHUNK * CORPUS_CHUNKS_PER_AUDIO_HOUR, + ); + // And the two bases must agree on the corpus, since that is the mix the + // conversion factor was measured on. + const viaChunks = sweepDays(191_116, MEASURED_SECONDS_PER_CHUNK); + const viaAudio = sweepDaysFromAudio( + 77_298 * 3600, + MEASURED_SECONDS_PER_AUDIO_HOUR, + ); + assert.ok( + Math.abs(viaChunks - viaAudio) / viaChunks < 0.02, + `bases disagree: ${viaChunks} vs ${viaAudio}`, + ); +}); + +// The census band table, as arithmetic rather than prose: chunk density spans 4×, +// which is the entire reason a single s/audio-hour figure cannot be trusted. +test("chunk density varies fourfold across the corpus", () => { + const bands = [ + { label: "<15m", chunks: 41_960, audioSeconds: 6_203 * 3600 }, + { label: ">8h", chunks: 24_888, audioSeconds: 14_850 * 3600 }, + ]; + const short = chunksPerAudioHour(bands[0].chunks, bands[0].audioSeconds); + const long = chunksPerAudioHour(bands[1].chunks, bands[1].audioSeconds); + assert.ok(short > 6.5 && short < 7.1, `short band density ${short}`); + assert.ok(long > 1.5 && long < 1.9, `long band density ${long}`); + assert.ok(short / long > 3.5, `expected ~4x spread, got ${short / long}`); + // The result that inverts the intuition: the long band is 2.4× the audio of the + // short band but LESS work. + assert.ok(bands[1].audioSeconds > bands[0].audioSeconds * 2); + assert.ok(bands[1].chunks < bands[0].chunks); +}); + +// The plan's chunk count and the chunker's actual slicing must be the same +// function. Pinned against the REAL chunker, because a drift here silently +// mis-prices GPU-weeks. +test("the chunk-count helper matches the real chunker", () => { + const maxCues = maxCuesForContext(8192); + assert.equal(maxCues, 600); + const opts = { maxCues, overlapCues: DIGEST_OVERLAP_CUES }; + + // The plan's stated formula: step is 560 at the shipped config. + assert.equal(countCueChunks(0, opts), 0); + assert.equal(countCueChunks(600, opts), 1); + assert.equal(countCueChunks(601, opts), 2); + assert.equal(countCueChunks(1160, opts), 2); // 600 + 560 exactly + assert.equal(countCueChunks(1161, opts), 3); + + // …and it agrees with slicing real cues, at the boundaries and either side. + for (const n of [0, 1, 599, 600, 601, 1159, 1160, 1161, 5000, 12_345]) { + const cues = Array.from({ length: n }, (_, i) => ({ + start: i, + end: i + 1, + text: "x", + })); + assert.equal( + countCueChunks(n, opts), + chunkCuesForContext(cues, opts).length, + `chunk count disagrees with the chunker at ${n} cues`, + ); + } +}); + +test("a video with no recorded cueCount is estimated, and says so", () => { + const real = chunksForVideo({ cueCount: 601, duration: 3600 }, 600); + assert.deepEqual(real, { chunks: 2, estimated: false }); + + // Falls back to the census density rather than dropping the video from the + // bill — understating the sweep is the worse failure. + const guessed = chunksForVideo({ cueCount: null, duration: 4 * 3600 }, 600); + assert.equal(guessed.estimated, true); + assert.equal(guessed.chunks, Math.round(4 * CORPUS_CHUNKS_PER_AUDIO_HOUR)); + // Never zero: a transcribed video always costs at least one call. + assert.equal(chunksForVideo({ cueCount: null, duration: 30 }, 600).chunks, 1); }); // A saving counted in videos is the mistake this module exists to prevent, so @@ -81,3 +184,63 @@ test("audio-hours, not video count, is what a saving is measured in", () => { "one VOD channel outweighs thousands of short mirrors", ); }); + +// The census, as a live regression pin against the REAL corpus. +// +// The pure test above pins the formula; this one pins the number the whole +// re-pricing rests on. It re-derives the corpus chunk total from the same LMDB row +// the plan reads, using the same helper, so it fails if the plan and the chunker +// ever drift — or if a DEFAULT_DIGEST_* changes the shipped chunk size without +// anyone re-pricing the sweep. +// +// Skipped when there is no corpus on this machine (CI, a fresh checkout): a pin on +// data that isn't there would just be a broken test. +test("the corpus chunk census reproduces 191,116 at the shipped config", async (t) => { + const { existsSync } = await import("node:fs"); + const { getPaths } = await import("../lib/paths"); + const paths = getPaths(); + if (!existsSync(paths.lmdbPath)) { + t.skip("no LMDB corpus on this machine"); + return; + } + const { open } = await import("lmdb"); + const { STATS_SCHEMA_VERSION } = await import("../lib/stats"); + const root = open({ path: paths.lmdbPath, maxDbs: 12, compression: true }); + const meta = root.openDB<unknown, string>({ + name: "statsMeta", + encoding: "msgpack", + }); + if ((meta.get("schema") as number | undefined) !== STATS_SCHEMA_VERSION) { + t.skip("stats cache is at another schema version — run build:stats"); + return; + } + const statsByPath = root.openDB< + { metaMs: number; stat: import("../lib/stats").VideoStat }, + [string, string] + >({ name: "statsByPath", encoding: "msgpack" }); + + const maxCues = maxCuesForContext(8192); + let chunks = 0; + let videos = 0; + let audioSeconds = 0; + let missingCueCount = 0; + for (const { value } of statsByPath.getRange()) { + const stat = value.stat; + if (!stat.hasTranscript || !(stat.duration > 0)) continue; + videos++; + audioSeconds += stat.duration; + if (typeof stat.cueCount !== "number") missingCueCount++; + chunks += chunksForVideo(stat, maxCues).chunks; + } + + assert.equal(chunks, 191_116, "corpus chunk total moved"); + assert.equal(videos, 73_367, "transcribed video count moved"); + // The census cost nothing precisely because cueCount was already on every row. + assert.equal(missingCueCount, 0, "some rows lost their cueCount"); + // And the derived density constant must still describe this corpus. + const density = chunksPerAudioHour(chunks, audioSeconds); + assert.ok( + Math.abs(density - CORPUS_CHUNKS_PER_AUDIO_HOUR) < 0.01, + `census density ${density.toFixed(3)} != CORPUS_CHUNKS_PER_AUDIO_HOUR`, + ); +}); diff --git a/common/controller/digestPlan.ts b/common/controller/digestPlan.ts @@ -1,11 +1,18 @@ -// What the digest backfill actually costs, measured in AUDIO-HOURS. +// What the digest backfill actually costs, measured in CHUNKS — and reported in +// audio-hours too, because that is the unit a human thinks in. // // A saving counted in videos is close to meaningless here: the corpus is 77k // videos but 77k audio-HOURS, and the two are not proportional per channel — // mirrors skew long (Hasan VODs, Quartering re-uploads), shorts skew numerous. -// The sweep's wall clock is (audio-hours × seconds-per-audio-hour), so a plan -// that moves 4,000 videos and 200 audio-hours has moved nothing. Everything -// here is therefore denominated in seconds of audio. +// So audio-seconds are tracked per video and reported everywhere. +// +// But audio-hours are not what the GPU is billed in. ONE CHUNK IS ONE MODEL +// CALL, and chunk density varies FOURFOLD across the corpus — 6.8 chunks per +// audio-hour under 15 minutes against 1.7 over 8 hours. A single +// seconds-per-audio-hour figure therefore is not a stable unit, and that unit +// error is the whole reason three separate cost measurements looked +// irreconcilable (27 vs 90 vs 151 s/audio-hour) when re-pricing them per chunk +// reproduces all three to within ~2%. See MEASURED_SECONDS_PER_CHUNK. // // Two consumers, deliberately one implementation: // - `bin/digest-plan.ts`, to price the sweep before committing GPU-weeks to it @@ -16,8 +23,10 @@ // // Source of truth is the `statsByPath` LMDB sub-DB — build:stats' output, keyed // [channelSlug, videoDir], which is the batch runner's own enumeration unit and -// carries duration + hasTranscript in one scan. It can lag the corpus; the -// schema version is checked and a mismatch is reported rather than swallowed. +// carries duration + hasTranscript + cueCount in one scan. It can lag the corpus; +// the schema version is checked and a mismatch is reported rather than swallowed. +// `cueCount` being on that row is what makes the chunk census free: it needed no +// GPU time and no transcript reads to compute. import { open } from "lmdb"; import path from "node:path"; @@ -26,18 +35,65 @@ import type { VideoStat } from "../lib/stats"; import { STATS_SCHEMA_VERSION } from "../lib/stats"; import { isSectionFresh } from "../lib/digest"; import { loadDigest } from "../lib/digest-server"; +import { countCueChunks } from "../lib/transcriptWindow"; +import { DIGEST_OVERLAP_CUES } from "../lib/digestPrompt"; import { buildDigestClusterPlan, type DigestClusterPlan, type DigestClusterRole, } from "./digestSharing"; -import { resolveDigestTarget, type DigestLaneChoice } from "./digestTarget"; +import { + resolveDigestChunking, + resolveDigestTarget, + type DigestLaneChoice, +} from "./digestTarget"; -// The measured cost of the local lane, from the 102-video validation run on the -// real corpus (NOT the bake-off's 27 s, which was projected from long videos on -// an idle box — see plans/FACTS.md). Overridable, because it is the one number -// here that is an estimate rather than a measurement of this corpus. -export const MEASURED_SECONDS_PER_AUDIO_HOUR = 90; +// The measured cost of the local lane, PER CHUNK — one model call. +// +// Re-pricing every measurement taken so far in chunks reconciles all of them, +// which the audio-hour model does for none (× 191,116 corpus chunks): +// +// 11.2 s idle box, ENGINE time (bake-off round 2) → 24.8 days [24.2 ✓] +// 24.9 s contended box, WALL time (102-video run) → 55.1 days [81 ✗] +// 60.6 s idle box, engine, gemma2 (bake-off round 2) → 134 days [132 ✓] +// +// The old 27 → 90 s/audio-hour "contention penalty" factors exactly: 1.49× from +// the validation sample's length mix (3.63 vs 2.44 chunks/audio-hour) × 2.22× +// per-chunk cost = 3.30× against the observed 3.33×. So the honest range is ~25 +// days idle to ~55 contended — not 81. +// +// The DEFAULT is the contended wall figure, deliberately: it is the only one +// measured in production, on this corpus, with all the sidecar work a bake-off +// skips. It is pessimistic on an idle box by design. The 2.22× gap has three +// parts (idle engine, contended engine, and non-engine wall inflation from the +// yield gate's 3 s poll) which only a production run can separate. +export const MEASURED_SECONDS_PER_CHUNK = 24.9; + +// Corpus-wide chunk density, from the free census over `statsByPath.cueCount` — +// all 73,367 transcribed videos, none missing a cue count, at the shipped maxCues +// 600 / overlap 40. It cost zero GPU time and no transcript reads: +// +// band videos audio-h % audio chunks chunks/audio-h % chunks +// < 15 min 41,960 6,203 8.0% 41,960 6.76 22.0% +// 15–60 min 14,036 6,631 8.6% 22,446 3.39 11.7% +// 1–2 h 5,847 8,480 11.0% 22,353 2.64 11.7% +// 2–4 h 5,486 15,832 20.5% 33,106 2.09 17.3% +// 4–8 h 4,504 25,304 32.7% 46,363 1.83 24.3% +// > 8 h 1,534 14,850 19.2% 24,888 1.68 13.0% +// total 73,367 77,298 191,116 2.47 +// +// The counter-intuitive result worth keeping in view: 4 h+ videos are 52% of the +// AUDIO but only 37% of the WORK, while sub-hour videos are 16% of audio and 34% +// of chunks. Any reasoning that prices the sweep in audio-hours gets this +// backwards. +export const CORPUS_CHUNKS_PER_AUDIO_HOUR = 2.47; + +// DERIVED, not measured — kept only so `bin/digest-plan.ts --rate` and anything +// else thinking in audio-hours still has a number to start from. It is the +// corpus-average conversion and is wrong for any individual channel by up to 2.7× +// in either direction. +export const MEASURED_SECONDS_PER_AUDIO_HOUR = + MEASURED_SECONDS_PER_CHUNK * CORPUS_CHUNKS_PER_AUDIO_HOUR; // Why a video is or is not this sweep's work. The mirror split is the point: a // mirror is only free if its cues actually ALIGN with its canonical member's. @@ -66,15 +122,19 @@ export function roleMustGenerate(role: DigestPlanRole): boolean { export type DigestPlanTotals = { videos: number; audioSeconds: number; + // Model calls. THIS is what the sweep's wall clock is made of; audioSeconds is + // for human comprehension and for comparing against the corpus census. + chunks: number; }; function emptyTotals(): DigestPlanTotals { - return { videos: 0, audioSeconds: 0 }; + return { videos: 0, audioSeconds: 0, chunks: 0 }; } -function addTo(t: DigestPlanTotals, seconds: number): void { +function addTo(t: DigestPlanTotals, seconds: number, chunks: number): void { t.videos++; t.audioSeconds += seconds; + t.chunks += chunks; } export type DigestChannelPlan = { @@ -86,8 +146,15 @@ export type DigestChannelPlan = { // Indexed but not digestable: no transcript, or no usable duration. ineligible: number; // Convenience rollups over `remaining`. - generateSeconds: number; // what this channel costs the GPU - sharedSeconds: number; // what cluster sharing takes off the bill + generateSeconds: number; // audio this channel still has to be digested + sharedSeconds: number; // audio cluster sharing takes off the bill + generateChunks: number; // what this channel actually costs the GPU + sharedChunks: number; // model calls cluster sharing avoids + // Videos counted as remaining work whose cueCount is missing from the stats + // cache, i.e. whose chunk count had to be estimated from duration. Surfaced + // because it is the one way the chunk total can be wrong, and it must not be + // invisible. + chunksEstimated: number; }; export type DigestSweepPlan = { @@ -100,6 +167,12 @@ export type DigestSweepPlan = { ineligible: number; generateSeconds: number; sharedSeconds: number; + generateChunks: number; + sharedChunks: number; + chunksEstimated: number; + // Cues per chunk the census was computed at. Recorded because the chunk total + // is only meaningful against a context size: halving numCtx roughly doubles it. + maxCues: number; clusters: number; clusterMembersMapped: number; // True when build:stats' cache is not at the version this code expects, i.e. @@ -139,9 +212,36 @@ function mergeInto( for (const role of DIGEST_PLAN_ROLES) { into[role].videos += from[role].videos; into[role].audioSeconds += from[role].audioSeconds; + into[role].chunks += from[role].chunks; } } +// Chunks for one video, from the cue count build:stats already recorded. +// +// `cueCount` is null only for a row written before stats carried it. Rather than +// drop such a video from the bill (which would understate the sweep), fall back to +// the census density — and the caller counts how often that happened, so an +// estimate can never quietly become the headline. +export function chunksForVideo( + stat: Pick<VideoStat, "cueCount" | "duration">, + maxCues: number, +): { chunks: number; estimated: boolean } { + if (typeof stat.cueCount === "number" && stat.cueCount >= 0) { + return { + chunks: countCueChunks(stat.cueCount, { + maxCues, + overlapCues: DIGEST_OVERLAP_CUES, + }), + estimated: false, + }; + } + const hours = Math.max(0, stat.duration) / 3600; + return { + chunks: Math.max(1, Math.round(hours * CORPUS_CHUNKS_PER_AUDIO_HOUR)), + estimated: true, + }; +} + // The whole role decision, as one pure function so it can be asserted on // without an LMDB corpus behind it. // @@ -209,6 +309,10 @@ export async function buildDigestSweepPlan( rows.push({ videoDir, stat: value.stat }); } + // Chunk size for the lane being priced. Resolved ONCE: it depends only on the + // app's numCtx, so re-deriving it per channel would just be slower. + const { maxCues } = resolveDigestChunking({ lane: opts.lane }); + const channels: DigestChannelPlan[] = []; for (const [channelSlug, rows] of byChannel) { const entry: DigestChannelPlan = { @@ -218,6 +322,9 @@ export async function buildDigestSweepPlan( ineligible: 0, generateSeconds: 0, sharedSeconds: 0, + generateChunks: 0, + sharedChunks: 0, + chunksEstimated: 0, }; const resolved = checkFreshness @@ -237,13 +344,15 @@ export async function buildDigestSweepPlan( continue; } + const { chunks, estimated } = chunksForVideo(stat, maxCues); + if (resolved) { const record = await loadDigest(path.join(dataDir, videoDir)); const allFresh = resolved.sections.every((section) => isSectionFresh(record, section, resolved.target), ); if (allFresh) { - addTo(entry.fresh, stat.duration); + addTo(entry.fresh, stat.duration, chunks); continue; } } @@ -252,19 +361,30 @@ export async function buildDigestSweepPlan( clusterPlan.bySlug.get(stat.slug), alignedBySlug.get(stat.slug), ); - addTo(entry.remaining[planRole], stat.duration); + addTo(entry.remaining[planRole], stat.duration, chunks); + if (estimated) entry.chunksEstimated++; } for (const role of DIGEST_PLAN_ROLES) { - if (roleMustGenerate(role)) + if (roleMustGenerate(role)) { entry.generateSeconds += entry.remaining[role].audioSeconds; - else entry.sharedSeconds += entry.remaining[role].audioSeconds; + entry.generateChunks += entry.remaining[role].chunks; + } else { + entry.sharedSeconds += entry.remaining[role].audioSeconds; + entry.sharedChunks += entry.remaining[role].chunks; + } } channels.push(entry); } + // Ordered by CHUNKS, not audio-hours: chunks are what the queue actually spends. + // The two orders differ — long-VOD channels have the lowest chunk density — and + // ordering by audio would put the cheapest-per-hour work first while claiming to + // be heaviest-first. (It is still longest-first in practice, so early throughput + // will look better than the corpus average; see the census table above.) channels.sort( (a, b) => + b.generateChunks - a.generateChunks || b.generateSeconds - a.generateSeconds || a.channelSlug.localeCompare(b.channelSlug), ); @@ -276,6 +396,10 @@ export async function buildDigestSweepPlan( ineligible: 0, generateSeconds: 0, sharedSeconds: 0, + generateChunks: 0, + sharedChunks: 0, + chunksEstimated: 0, + maxCues, clusters: clusterPlan.clusters, clusterMembersMapped: clusterPlan.bySlug.size, statsSchemaStale, @@ -285,9 +409,13 @@ export async function buildDigestSweepPlan( mergeInto(totals.remaining, c.remaining); totals.fresh.videos += c.fresh.videos; totals.fresh.audioSeconds += c.fresh.audioSeconds; + totals.fresh.chunks += c.fresh.chunks; totals.ineligible += c.ineligible; totals.generateSeconds += c.generateSeconds; totals.sharedSeconds += c.sharedSeconds; + totals.generateChunks += c.generateChunks; + totals.sharedChunks += c.sharedChunks; + totals.chunksEstimated += c.chunksEstimated; } return totals; } @@ -313,10 +441,33 @@ export function audioHours(seconds: number): number { return seconds / 3600; } -// The projection the whole plan exists to produce. +// The projection the whole plan exists to produce, priced in the unit the GPU is +// actually billed in. export function sweepDays( + chunks: number, + secondsPerChunk: number = MEASURED_SECONDS_PER_CHUNK, +): number { + return (chunks * secondsPerChunk) / 86400; +} + +// The same projection from audio-hours, for the one caller that thinks in them +// (`bin/digest-plan.ts --rate`). Kept because a rate in s/audio-hour is what every +// past measurement was quoted in and a reader will want to reproduce them — but it +// is corpus-average only. Prefer sweepDays(). +export function sweepDaysFromAudio( audioSeconds: number, secondsPerAudioHour: number = MEASURED_SECONDS_PER_AUDIO_HOUR, ): number { return (audioHours(audioSeconds) * secondsPerAudioHour) / 86400; } + +// Chunks per audio-hour for an arbitrary slice of the plan — the density figure +// that makes a per-chunk cost comparable to the census. Report THIS beside a +// throughput number, never a bare s/audio-hour: mid-sweep the queue is +// longest-first and its density drifts from 1.7 toward 6.8 as it drains, so a +// running s/audio-hour will look like the projection is falling apart when it is +// only reaching shorter videos. +export function chunksPerAudioHour(chunks: number, audioSeconds: number): number { + const h = audioHours(audioSeconds); + return h > 0 ? chunks / h : 0; +} diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts @@ -44,6 +44,7 @@ import { countMissingDigests, runDigestBatch } from "./digestBatch"; import { audioHours, buildDigestSweepPlan, + chunksPerAudioHour, sweepDays, type DigestSweepPlan, } from "./digestPlan"; @@ -215,8 +216,10 @@ async function runSweepLoop( } onLog( `Pass ${pass}: ${work.length} channel(s), ` + + `${plan.generateChunks.toLocaleString()} chunks / ` + `${audioHours(plan.generateSeconds).toFixed(0)} audio-hours to generate ` + - `(~${sweepDays(plan.generateSeconds).toFixed(1)} days at the measured rate), ` + + `(~${sweepDays(plan.generateChunks).toFixed(1)} days at the measured rate, ` + + `${chunksPerAudioHour(plan.generateChunks, plan.generateSeconds).toFixed(2)} chunks/audio-h), ` + `${audioHours(plan.sharedSeconds).toFixed(0)} audio-hours covered by cluster sharing.`, ); @@ -228,8 +231,9 @@ async function runSweepLoop( return; } onLog( - `→ ${channel.channelSlug}: ${audioHours(channel.generateSeconds).toFixed(0)} audio-hours ` + - `(~${sweepDays(channel.generateSeconds).toFixed(1)} days).`, + `→ ${channel.channelSlug}: ${channel.generateChunks.toLocaleString()} chunks / ` + + `${audioHours(channel.generateSeconds).toFixed(0)} audio-hours ` + + `(~${sweepDays(channel.generateChunks).toFixed(1)} days).`, ); const outcome = await runDigestChannelJob({ paths, diff --git a/common/controller/digestTarget.ts b/common/controller/digestTarget.ts @@ -53,14 +53,18 @@ export type ResolvedDigestTarget = { context: DigestContext; }; -export async function resolveDigestTarget(opts: { - paths: Paths; - channelSlug: string; +// Which app a lane resolves to, and the chunk size that follows from its context. +// +// Split out of resolveDigestTarget because the PRICING pass needs the chunk size +// without needing a channel: the number of chunks a transcript costs depends only +// on `maxCues`, while a freshness target additionally needs the per-channel +// context note. Sharing this keeps the plan's chunk count and the chunker's +// actual slicing derived from one place — the alternative is a second copy of +// `maxCuesForContext(config.numCtx)` that silently stops matching the engine. +export function resolveDigestChunking(opts: { lane?: DigestLaneChoice; - // Overrides the lane's configured app. The bake-off harness drives this. appId?: string; - sections?: DigestSectionKind[]; -}): Promise<ResolvedDigestTarget> { +}): { app: DigestApp; config: DigestAppConfig; maxCues: number } { const digestSettings = getSettings().digest; const appId = opts.appId ?? @@ -69,8 +73,20 @@ export async function resolveDigestTarget(opts: { : digestSettings.localAppId); const app = getDigestApp(appId); const config: DigestAppConfig = digestSettings.apps[app.id] ?? {}; + return { app, config, maxCues: maxCuesForContext(config.numCtx) }; +} + +export async function resolveDigestTarget(opts: { + paths: Paths; + channelSlug: string; + lane?: DigestLaneChoice; + // Overrides the lane's configured app. The bake-off harness drives this. + appId?: string; + sections?: DigestSectionKind[]; +}): Promise<ResolvedDigestTarget> { + const digestSettings = getSettings().digest; + const { app, config, maxCues } = resolveDigestChunking(opts); const modelRequested = config.model?.trim() || app.defaultModel(); - const maxCues = maxCuesForContext(config.numCtx); const promptVariant = digestPromptVariant({ ...digestSettings, maxCues, diff --git a/common/controller/digestYield.test.ts b/common/controller/digestYield.test.ts @@ -0,0 +1,155 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + evaluateTranscriptionActivity, + workerContendsForGpu, +} from "./digestYield"; +import { defaultDigest, sanitizeDigest } from "../lib/settings"; + +// Run with: +// pnpm --filter yt-dlp-transcript-common exec tsx --test common/controller/digestYield.test.ts + +// This box, exactly: one worker on the engine's default device plus two pinned to +// CPU, all `kind: "local"` and all enabled. Reproducing it here is the point — the +// bug was invisible to a single-worker mental model. +const GPU = { kind: "local" as const, device: undefined }; +const CPU_A = { kind: "local" as const, device: "cpu" }; +const CPU_B = { kind: "local" as const, device: "cpu" }; +const CUDA = { kind: "local" as const, device: "cuda:0" }; +const REMOTE = { kind: "remote" as const, device: undefined }; + +function busy<T extends object>(w: T) { + return { ...w, busy: true }; +} +function idle<T extends object>(w: T) { + return { ...w, busy: false }; +} + +test("only an explicit cpu device is non-contending", () => { + // The fix: explicit "cpu" competes for zero GPU shaders. + assert.equal(workerContendsForGpu(CPU_A, false), false); + // Trimmed and case-folded: settings.json is hand-editable. + assert.equal(workerContendsForGpu({ ...CPU_A, device: " CPU " }, false), false); + // A device we did not set is the engine binary's OWN default, which may be the + // GPU. The unknown case must fail SAFE, i.e. still contend. + assert.equal(workerContendsForGpu(GPU, false), true); + assert.equal(workerContendsForGpu({ kind: "local", device: "" }, false), true); + assert.equal(workerContendsForGpu(CUDA, false), true); + // A remote worker runs on another machine and competes for nothing here. + assert.equal(workerContendsForGpu(REMOTE, false), false); + // Opting in restores the old blanket behaviour for every local worker. + assert.equal(workerContendsForGpu(CPU_A, true), true); + assert.equal(workerContendsForGpu(REMOTE, true), false); +}); + +// THE CLAIM OF THE FIX, both directions. Two CPU-pinned transcriptions running +// while the GPU worker sits idle must NOT stop the digest lane under the default, +// and flipping the setting must restore exactly today's behaviour. +test("busy CPU workers do not stall the digest lane by default", () => { + const workers = [idle(GPU), busy(CPU_A), busy(CPU_B)]; + // A real transcription batch IS running while those CPU workers work — that is + // what made the naive queued-job signal keep the bug alive. + const activity = evaluateTranscriptionActivity({ + workers, + transcriptionJobRunning: true, + yieldToCpuWorkers: false, + }); + assert.deepEqual(activity, { busy: false, reason: null }); + + const optedIn = evaluateTranscriptionActivity({ + workers, + transcriptionJobRunning: true, + yieldToCpuWorkers: true, + }); + assert.deepEqual(optedIn, { busy: true, reason: "local-worker" }); +}); + +test("a busy GPU worker still stops the digest lane", () => { + for (const yieldToCpuWorkers of [false, true]) { + assert.deepEqual( + evaluateTranscriptionActivity({ + workers: [busy(GPU), idle(CPU_A)], + transcriptionJobRunning: true, + yieldToCpuWorkers, + }), + { busy: true, reason: "local-worker" }, + `default-device worker must contend (yieldToCpuWorkers=${yieldToCpuWorkers})`, + ); + assert.deepEqual( + evaluateTranscriptionActivity({ + workers: [busy(CUDA), busy(CPU_A)], + transcriptionJobRunning: true, + yieldToCpuWorkers, + }), + { busy: true, reason: "local-worker" }, + `explicit cuda worker must contend (yieldToCpuWorkers=${yieldToCpuWorkers})`, + ); + } +}); + +// The gap signal is why this file has two signals at all: between worker +// acquisitions (audio extraction, model load) the pool shows nothing busy while +// the card may well be loading a model. Narrowing the CPU case must not cost it. +test("a running job with no busy worker still covers the acquisition gap", () => { + assert.deepEqual( + evaluateTranscriptionActivity({ + workers: [idle(GPU), idle(CPU_A), idle(CPU_B)], + transcriptionJobRunning: true, + yieldToCpuWorkers: false, + }), + { busy: true, reason: "queued-job" }, + ); + // …and an idle box yields to nothing. + assert.deepEqual( + evaluateTranscriptionActivity({ + workers: [idle(GPU)], + transcriptionJobRunning: false, + yieldToCpuWorkers: false, + }), + { busy: false, reason: null }, + ); +}); + +// A busy REMOTE worker is not local activity, so the gap signal is still the one +// that speaks — the remote's own box arbitrates its own card. +test("a busy remote worker does not suppress the gap signal", () => { + assert.deepEqual( + evaluateTranscriptionActivity({ + workers: [busy(REMOTE), idle(GPU)], + transcriptionJobRunning: true, + yieldToCpuWorkers: false, + }), + { busy: true, reason: "queued-job" }, + ); +}); + +// Absence must fall to OFF, and it must not disturb the master switch's own +// default. The two fields use opposite idioms on purpose. +test("yieldToCpuWorkers defaults off; yieldToTranscription defaults on", () => { + assert.equal(defaultDigest().yieldToCpuWorkers, false); + assert.equal(defaultDigest().yieldToTranscription, true); + + // A settings file written before either field existed. + const legacy = sanitizeDigest({}); + assert.equal(legacy.yieldToCpuWorkers, false); + assert.equal(legacy.yieldToTranscription, true); + + // Explicit values survive, in both directions, for both fields. + assert.equal(sanitizeDigest({ yieldToCpuWorkers: true }).yieldToCpuWorkers, true); + assert.equal( + sanitizeDigest({ yieldToCpuWorkers: false }).yieldToCpuWorkers, + false, + ); + assert.equal( + sanitizeDigest({ yieldToTranscription: false }).yieldToTranscription, + false, + ); + // Garbage is not truthy-coerced into opting in. + for (const v of ["true", 1, {}, [], null]) { + assert.equal( + sanitizeDigest({ yieldToCpuWorkers: v }).yieldToCpuWorkers, + false, + `${JSON.stringify(v)} must not enable yieldToCpuWorkers`, + ); + } +}); diff --git a/common/controller/digestYield.ts b/common/controller/digestYield.ts @@ -6,12 +6,16 @@ // network-bound), but its consequence for digest-local vs transcription is that // **ollama and whisper run concurrently on the same 8 GB card**. // -// That is not hypothetical: the bake-off projected 27 s per audio-hour on an -// idle box and the 102-video validation run measured **90 s** on a box also -// running auto-transcribe. Over a 77,000-audio-hour sweep the difference is -// roughly 80 days versus 24 — far and away the largest cost lever in the -// backfill, and about fifty times bigger than everything duplicate sharing can -// save (see bin/digest-plan.ts). +// That is not hypothetical, but the size of it is a HYPOTHESIS, not a +// measurement. The bake-off's 11.2 s per chunk was ENGINE time on an idle box; +// the 102-video validation run's 24.9 s per chunk was WALL time on a box also +// running auto-transcribe. Those are not the same quantity, so their 2.22× ratio +// is an upper bound on contention that also contains model-load, prefill and — +// crucially — this file's own 3 s poll idling. Re-priced per chunk the corpus is +// ~25 days idle to ~55 contended (see controller/digestPlan.ts). Splitting that +// gap needs a production run with the per-call load/prefill/decode breakdown that +// digestApps.ts now records; until then, do not quote a contention penalty as +// fact. // // The fix is deliberately NOT a scheduler. `digestBatch`'s limit() already // returns 0 to idle-wait on a pause, and runPool treats a zero limit as "hold, @@ -23,9 +27,10 @@ // starts just after a dispatch. Interrupting would waste the partial work; the // next dispatch sees the busy lane and holds. -import { getWorkerPool } from "../jobs/workerPool"; +import { getWorkerPool, type WorkerSummary } from "../jobs/workerPool"; import { getRegistry } from "../jobs/registry"; import { TRANSCRIPTION_QUEUE } from "../lib/queueKeys"; +import { getSettings } from "../lib/settings"; export type TranscriptionActivity = { busy: boolean; @@ -34,35 +39,104 @@ export type TranscriptionActivity = { reason: "local-worker" | "queued-job" | null; }; +// Does this busy worker actually compete for GPU shaders? +// +// The original check was `kind === "local"` alone, and that was a BUG. `local` +// distinguishes "on this box" from "delegated to another machine" — it says +// nothing about which processor the engine uses. This box runs three enabled local +// parakeet workers: one on the engine's default device and two pinned to +// `device: "cpu"`. At parallelTranscriptions 2 the CPU pair alone was enough to +// hold the digest lane at zero throughput for work using no shaders at all. +// +// The unknown case fails SAFE, in the direction that costs throughput rather than +// correctness: only an explicit "cpu" is treated as non-contending. No device set +// means parakeet-cli's own default, which may be the GPU. +export function workerContendsForGpu( + w: Pick<WorkerSummary, "kind" | "device">, + yieldToCpuWorkers: boolean, +): boolean { + if (w.kind !== "local") return false; + if (yieldToCpuWorkers) return true; + return w.device?.trim().toLowerCase() !== "cpu"; +} + // Is the local transcription lane using the GPU right now? // // Two signals, because one alone leaves a hole: -// - a busy LOCAL worker is whisper actually running on this box's card -// (`remote` workers delegate to another machine and compete for nothing -// here, so they are deliberately excluded); -// - a RUNNING job on TRANSCRIPTION_QUEUE covers the gaps between worker +// - a busy LOCAL worker on a GPU device is the transcription engine actually +// running on this box's card (`remote` workers delegate to another machine +// and compete for nothing here, so they are deliberately excluded, as are +// CPU-pinned workers unless `yieldToCpuWorkers` says otherwise); +// - a RUNNING job on TRANSCRIPTION_QUEUE covers the GAPS BETWEEN worker // acquisitions — audio extraction, model load, the moment between two // videos in a batch. Those gaps are exactly where a multi-minute digest // generation would otherwise slip in and hold the card. +// +// The second signal is deliberately consulted ONLY when no local worker is busy. +// It is a proxy for "a phase the pool cannot see", and if a worker IS busy then +// the pool has already told us the truth — including which DEVICE. Without that +// condition the device check above would be dead code on this box: a CPU-pinned +// transcription is a running job on the queue, so the digest lane would go on +// yielding to it through the gap signal and the fix would change nothing. +// +// Known, accepted narrow window: if one worker is busy on CPU while a second job +// is loading a model onto the GPU, this reports not-busy and one digest chunk may +// overlap it. A yield is not a lock — an in-flight generation was never +// interrupted either — and the alternative reinstates the bug. +// The whole decision, as a pure function of the two signals — so both directions +// of the CPU-worker fix can be asserted without a live pool, a live registry or a +// GPU. transcriptionActivity() below is the thin I/O wrapper that feeds it. +export function evaluateTranscriptionActivity(opts: { + workers: readonly Pick<WorkerSummary, "busy" | "kind" | "device">[]; + transcriptionJobRunning: boolean; + yieldToCpuWorkers: boolean; +}): TranscriptionActivity { + let anyLocalBusy = false; + for (const w of opts.workers) { + if (!w.busy || w.kind !== "local") continue; + anyLocalBusy = true; + if (workerContendsForGpu(w, opts.yieldToCpuWorkers)) { + return { busy: true, reason: "local-worker" }; + } + } + if (anyLocalBusy) return { busy: false, reason: null }; + if (opts.transcriptionJobRunning) return { busy: true, reason: "queued-job" }; + return { busy: false, reason: null }; +} + export function transcriptionActivity(): TranscriptionActivity { + // Read once: a setting flipping mid-scan would produce an answer that matches + // neither configuration. + let yieldToCpuWorkers = false; try { - for (const w of getWorkerPool().summary()) { - if (w.busy && w.kind === "local") { - return { busy: true, reason: "local-worker" }; - } - } + yieldToCpuWorkers = getSettings().digest.yieldToCpuWorkers; } catch { - // A pool that cannot be read must not wedge the digest lane: fail OPEN, i.e. - // keep generating. The cost of being wrong here is contention, not deadlock. + /* unreadable settings: the sanitizer's default is false, so keep it */ } + + // Each read is guarded separately and falls back to "nothing there", i.e. fail + // OPEN: a pool or registry that cannot be read must not wedge the digest lane. + // The cost of being wrong here is contention, not deadlock. + let workers: readonly WorkerSummary[] = []; try { - for (const job of getRegistry().list()) { - if (job.queueKey === TRANSCRIPTION_QUEUE && job.status === "running") { - return { busy: true, reason: "queued-job" }; - } - } + workers = getWorkerPool().summary(); } catch { - /* same: fail open */ + /* fail open */ } - return { busy: false, reason: null }; + let transcriptionJobRunning = false; + try { + transcriptionJobRunning = getRegistry() + .list() + .some( + (job) => + job.queueKey === TRANSCRIPTION_QUEUE && job.status === "running", + ); + } catch { + /* fail open */ + } + return evaluateTranscriptionActivity({ + workers, + transcriptionJobRunning, + yieldToCpuWorkers, + }); } diff --git a/common/jobs/workerPool.ts b/common/jobs/workerPool.ts @@ -55,6 +55,15 @@ export type WorkerSummary = { state: WorkerRuntimeState; degraded: boolean; enabled: boolean; // persisted intent (config.enabled) + // The engine's compute device, verbatim from the per-worker config ("cpu", + // "cuda:0", …). UNDEFINED MEANS UNKNOWN, not CPU: it is the engine binary's own + // default, which on parakeet-cli may well be the GPU. + // + // Exposed because digestYield() has to answer "is anything competing for GPU + // shaders right now?", and `kind: "local"` does not answer it — this box runs a + // GPU worker beside two explicitly CPU-pinned ones, and treating those as GPU + // contention stopped the digest lane dead for work that used no shaders at all. + device?: string; }; type PoolEntry = { @@ -480,6 +489,7 @@ export class WorkerPool { state: e.state, degraded: e.degraded, enabled: e.config.enabled, + device: e.config.config?.device, })) .sort((a, b) => a.priority - b.priority || a.id.localeCompare(b.id)); } diff --git a/common/lib/digestApps.ts b/common/lib/digestApps.ts @@ -56,11 +56,33 @@ export type DigestRunResult = { // Recorded in provenance, so a config of "qwen2.5" resolving to "qwen2.5:7b" // doesn't later look like a model change and trigger a needless regeneration. model: string; + // WALL time, bracketing the fetch — not engine time. The distinction is the + // whole reason the fields below exist: "27 s/audio-hour vs 90" was a comparison + // of an engine number against a wall number, and nothing recorded could tell + // them apart after the fact. durationMs: number; // Metered lanes only. costUsd?: number; inputTokens?: number; outputTokens?: number; + + // Engine-reported timing breakdown, when the engine reports one (ollama does, + // on every /api/chat response; the CLI lane does not). All optional and all + // ADDITIVE — nothing here reaches DigestProvenance, so recording them changes + // no freshness identity and invalidates no existing digest. + // + // Why it is worth the four lines: this splits a call into model-load, prefill + // and decode. Warmup and a cold model load then stop masquerading as + // contention, and a contended run shows WHERE it was hurt — prefill (shader + // competition) reads differently from decode (memory bandwidth). Comparing + // totalMs against durationMs also quantifies non-engine wall inflation, which + // on an idle box measured 1.004× and under contention is unbounded. + loadMs?: number; + promptEvalMs?: number; + evalMs?: number; + // The engine's own total. NOT the sum of the three above: ollama's + // total_duration also covers queueing and tokenization inside the server. + totalMs?: number; }; export type DigestApp = { @@ -200,6 +222,12 @@ const ollamaDirect: DigestApp = { message?: { content?: string }; prompt_eval_count?: number; eval_count?: number; + // Nanoseconds, all four. ollama has always returned these; they were simply + // being thrown away. + total_duration?: number; + load_duration?: number; + prompt_eval_duration?: number; + eval_duration?: number; }; const content = body.message?.content ?? ""; const data = extractJsonObject(content); @@ -208,21 +236,57 @@ const ollamaDirect: DigestApp = { `ollama returned unparseable content: ${content.slice(0, 300)}`, ); } + const durationMs = Date.now() - startedAt; + const loadMs = nsToMs(body.load_duration); + const promptEvalMs = nsToMs(body.prompt_eval_duration); + const evalMs = nsToMs(body.eval_duration); + const totalMs = nsToMs(body.total_duration); + // The breakdown goes in the LOG, not just the return value: a multi-week + // sweep's job logs are the only surviving record of how it actually behaved, + // and the last measurement's logs had already rotated away when the numbers + // needed reconciling. `wall` beside `engine` is what makes the yield gate's + // own idling visible. + const parts = [ + loadMs !== undefined ? `load ${secs(loadMs)}` : null, + promptEvalMs !== undefined ? `prefill ${secs(promptEvalMs)}` : null, + evalMs !== undefined ? `decode ${secs(evalMs)}` : null, + totalMs !== undefined ? `engine ${secs(totalMs)}` : null, + ].filter(Boolean); onLog?.( `ollama ${body.model ?? model}: ${body.prompt_eval_count ?? "?"} in / ${ body.eval_count ?? "?" - } out tokens in ${Math.round((Date.now() - startedAt) / 100) / 10}s (num_ctx ${numCtx})`, + } out tokens in ${secs(durationMs)} wall` + + (parts.length ? ` (${parts.join(", ")})` : "") + + ` (num_ctx ${numCtx})`, ); return { data, model: body.model ?? model, - durationMs: Date.now() - startedAt, + durationMs, inputTokens: body.prompt_eval_count, outputTokens: body.eval_count, + ...(loadMs !== undefined ? { loadMs } : {}), + ...(promptEvalMs !== undefined ? { promptEvalMs } : {}), + ...(evalMs !== undefined ? { evalMs } : {}), + ...(totalMs !== undefined ? { totalMs } : {}), }; }, }; +// ollama reports durations in NANOSECONDS. Absent or non-finite means the engine +// did not report that phase (an older build, or a cached-prompt call that skipped +// prefill entirely), and the field is then omitted rather than recorded as 0 — a +// zero here would look like a measurement, not a gap. +function nsToMs(ns: number | undefined): number | undefined { + return typeof ns === "number" && Number.isFinite(ns) && ns >= 0 + ? Math.round(ns / 1e6) + : undefined; +} + +function secs(ms: number): string { + return `${Math.round(ms / 100) / 10}s`; +} + // --------------------------------------------------------------------------- // claude-code — the metered overflow lane, OFF by default // --------------------------------------------------------------------------- diff --git a/common/lib/settings.ts b/common/lib/settings.ts @@ -216,13 +216,26 @@ export type DigestSettings = { // pattern), so a pause survives a restart with no boot hook — unlike // transcriptionsPaused, which needs editor/instrumentation.ts to re-apply it. digestsPaused: boolean; - // Yield the GPU to the transcription lane: while whisper is working, the + // Yield the GPU to the transcription lane: while transcription is working, the // digest batch's limit() returns 0 and the pool idle-waits. ON by default, // because `digest:local` is deliberately on a different queue from - // TRANSCRIPTION_QUEUE and so would otherwise run ollama and whisper on the - // same 8 GB card — measured at 90 s per audio-hour against the 27 s an idle - // box projected. See controller/digestYield.ts. + // TRANSCRIPTION_QUEUE and so would otherwise run ollama and the transcription + // engine on the same 8 GB card. See controller/digestYield.ts. yieldToTranscription: boolean; + // Whether a busy worker pinned to `device: "cpu"` counts as GPU contention. + // + // OFF by default, which is the FIX for a real bug: the yield originally tested + // only `kind === "local"`, so on a box with one GPU worker and two CPU-pinned + // ones (this box, at parallelTranscriptions 2) the digest lane stopped dead for + // transcription that competes for zero GPU shaders. + // + // Only an EXPLICIT "cpu" is treated as non-contending. A worker with no device + // set is using the engine binary's own default, which may be the GPU, so it + // still triggers the yield — the unknown case fails safe. + // + // Composes with `yieldToTranscription`: that is the master switch, this only + // narrows which workers it reacts to. + yieldToCpuWorkers: boolean; // The corpus-wide sweep is armed. Read at boot by the editor's instrumentation // hook, the same way the auto-transcribe/auto-download runners are, so a sweep // survives a server restart. It is persisted INTENT, not a cursor: the batch @@ -522,10 +535,13 @@ export function defaultDigest(): DigestSettings { // here would give the same number two homes and let them drift. apps: {}, digestsPaused: false, - // ON. Contention with whisper is the single largest cost in the backfill - // (90 s/audio-hour measured vs 27 projected on an idle box), so the safe - // default is to step aside; turning it off is the deliberate choice. + // ON. Real GPU contention with the transcription engine is a genuine cost + // (re-priced: 11.2 s/chunk idle against 24.9 s/chunk on a contended box), so + // the safe default is to step aside; turning it off is the deliberate choice. yieldToTranscription: true, + // OFF. A CPU-pinned worker is not GPU contention, and treating it as such + // stalled the digest lane for nothing. See DigestSettings.yieldToCpuWorkers. + yieldToCpuWorkers: false, // OFF. A corpus-wide sweep is GPU-weeks of work and is never armed by // default — an operator starts it. sweepEnabled: false, @@ -601,6 +617,13 @@ export function sanitizeDigest(value: unknown): DigestSettings { // settings file written before this field existed keeps the GPU-safe // behaviour instead of silently opting into contention. yieldToTranscription: r.yieldToTranscription !== false, + // The OPPOSITE idiom, and deliberately so: `=== true`, so absence falls to + // OFF. The field's absence means a settings file written before the CPU-worker + // bug was found, and for those files OFF is the FIXED behaviour, not a silent + // change of intent — nobody ever asked to stall the digest lane for a CPU + // transcription. `yieldToTranscription` still gates the whole thing, so the + // GPU-safe default is untouched. + yieldToCpuWorkers: r.yieldToCpuWorkers === true, sweepEnabled: r.sweepEnabled === true, sweepChannels: Array.isArray(r.sweepChannels) ? r.sweepChannels.filter( diff --git a/common/lib/transcriptWindow.ts b/common/lib/transcriptWindow.ts @@ -81,6 +81,36 @@ export function chunkCuesForContext( return out; } +// How many chunks chunkCuesForContext() WOULD produce for a transcript of +// `cueCount` cues — the same count, without needing the cues themselves. +// +// This exists because a chunk is the digest sweep's real unit of work (one chunk +// = one model call), and the pricing pass has only `statsByPath.cueCount` to work +// from — reading 73k transcripts to count slices would cost more than the thing +// it is pricing. Chunk density varies FOURFOLD across the corpus (6.8 +// chunks/audio-hour under 15 min against 1.7 over 8 h), which is why +// seconds-per-audio-hour is not a stable unit and three past cost estimates +// looked contradictory when they agreed to within 2% per chunk. +// +// Deliberately in this file, immediately below the chunker it mirrors: the two +// must not drift, and a test pins this against the real slicing. +export function countCueChunks( + cueCount: number, + opts: { maxCues?: number; overlapCues?: number } = {}, +): number { + // Clamp exactly as the chunker does, so an odd config can't make the two + // disagree about what "one chunk" means. + const maxCues = Math.max(1, Math.floor(opts.maxCues ?? 1200)); + const overlapCues = Math.max( + 0, + Math.min(Math.floor(opts.overlapCues ?? 0), maxCues - 1), + ); + const n = Math.max(0, Math.floor(cueCount)); + if (n === 0) return 0; + if (n <= maxCues) return 1; + return 1 + Math.ceil((n - maxCues) / (maxCues - overlapCues)); +} + // Render cues as snippet objects (clock + seconds + collapsed, capped text). // Empty cues are dropped. export function cuesToSnippets(cues: Cue[]): WindowSnippet[] { diff --git a/editor/app/settings/actions.ts b/editor/app/settings/actions.ts @@ -21,7 +21,10 @@ import { type SocialLink, } from "yt-dlp-transcript-common/lib/settings"; import { DEFAULT_TRANSCRIPTION_APP_ID } from "yt-dlp-transcript-common/lib/transcriptionApps"; -import { isDigestTimestampMode } from "yt-dlp-transcript-common/lib/digest"; +import { + isDigestSectionKind, + isDigestTimestampMode, +} from "yt-dlp-transcript-common/lib/digest"; import { DEFAULT_COOKIE_MODE, isCookieMode, @@ -41,7 +44,9 @@ export async function saveSettingsAction( ): Promise<SaveResult> { const adminTitle = String(formData.get("adminTitle") ?? "").trim(); const homepageUrl = String(formData.get("homepageUrl") ?? "").trim(); - const maxBytesRaw = String(formData.get("maxTranscriptPageBytes") ?? "").trim(); + const maxBytesRaw = String( + formData.get("maxTranscriptPageBytes") ?? "", + ).trim(); const cookiesFromBrowser = String( formData.get("cookiesFromBrowser") ?? "", ).trim(); @@ -112,10 +117,7 @@ export async function saveSettingsAction( error: "sleepBetweenDownloadsSeconds must be a number", }; } - if ( - sleepParsed < 0 || - sleepParsed > SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS - ) { + if (sleepParsed < 0 || sleepParsed > SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS) { return { ok: false, error: `sleepBetweenDownloadsSeconds must be between 0 and ${SLEEP_BETWEEN_DOWNLOADS_MAX_SECONDS}`, @@ -180,7 +182,10 @@ export async function saveSettingsAction( return { ok: false, error: "Social links payload is malformed" }; } const socialParsed = parseSocialLinks(socialInput); - if (Array.isArray(socialInput) && socialParsed.length !== socialInput.length) { + if ( + Array.isArray(socialInput) && + socialParsed.length !== socialInput.length + ) { return { ok: false, error: @@ -216,7 +221,14 @@ export async function saveSettingsAction( return { ok: false, error: "Digest app config payload is malformed" }; } } - const digestSectionsRaw = formData.getAll("digestSections").map(String); + // Filtered to KNOWN kinds here rather than leaning on writeSettings' sanitizer. + // sanitizeDigest would drop an unknown value anyway, but typing it honestly is + // what lets the digest block below be checked against DigestSettings instead of + // cast — and the cast is what hid the dropped-fields bug. + const digestSectionsRaw = formData + .getAll("digestSections") + .map(String) + .filter(isDigestSectionKind); // A hidden marker, because unchecked checkboxes are simply ABSENT from a // FormData: without it, a submit from any form that lacks the digest fields // would read remoteEnabled as false and silently reset the block. @@ -224,51 +236,69 @@ export async function saveSettingsAction( const digestTimestampModeRaw = String( formData.get("digestTimestampMode") ?? "", ).trim(); - const digestSettings = ( - digestFormPresent - ? { - remoteEnabled: formData.get("digestRemoteEnabled") === "on", - longTailSeconds: Number.parseInt( - String(formData.get("digestLongTailSeconds") ?? "").trim(), - 10, - ), - localAppId: - String(formData.get("digestLocalAppId") ?? "").trim() || - dD.localAppId, - remoteAppId: - String(formData.get("digestRemoteAppId") ?? "").trim() || - dD.remoteAppId, - apps: digestApps, - // Not edited by this form — the dashboard/channel controls own the pause. - digestsPaused: dD.digestsPaused, - spendCapUsd: Number.parseFloat( - String(formData.get("digestSpendCapUsd") ?? "").trim(), - ), - sections: digestSectionsRaw.length > 0 ? digestSectionsRaw : dD.sections, - // Prompt SHAPE — both freshness-affecting (see digestPromptVariant), - // so a silent reset here would invalidate every digest generated - // under a non-default shape. Each still FALLS BACK to the current - // value rather than to the default when the field is absent: a form - // that rebuilds this block but omits a field is exactly how the reset - // bug happens, and the fallback is what makes omission harmless. - timestampMode: isDigestTimestampMode(digestTimestampModeRaw) - ? digestTimestampModeRaw - : dD.timestampMode, - promptVariant: formData.has("digestPromptVariant") - ? String(formData.get("digestPromptVariant") ?? "").trim() - : dD.promptVariant, - } - : dD - ) as SiteSettings["digest"]; + const digestSettings: SiteSettings["digest"] = digestFormPresent + ? { + // Spread the CURRENT block first. Every field this form does not render + // must survive a save untouched, and before this spread they did not: + // the block was built field-by-field and cast, so `yieldToTranscription` + // and — much worse — `sweepEnabled`/`sweepChannels` were absent from the + // object, and sanitizeDigest re-derived them from nothing. Saving any + // unrelated setting DISARMED AN ARMED CORPUS SWEEP. The cast is gone + // too, so the next added field is a type error rather than a silent + // reset. + ...dD, + remoteEnabled: formData.get("digestRemoteEnabled") === "on", + longTailSeconds: Number.parseInt( + String(formData.get("digestLongTailSeconds") ?? "").trim(), + 10, + ), + localAppId: + String(formData.get("digestLocalAppId") ?? "").trim() || + dD.localAppId, + remoteAppId: + String(formData.get("digestRemoteAppId") ?? "").trim() || + dD.remoteAppId, + // Cast, not sanitized here: this is hand-parsed JSON from the form and + // writeSettings runs sanitizeDigestApps over it. Narrowed at exactly + // this field so every OTHER field in the block stays type-checked. + apps: digestApps as SiteSettings["digest"]["apps"], + // Not edited by this form — the dashboard/channel controls own the pause. + digestsPaused: dD.digestsPaused, + spendCapUsd: Number.parseFloat( + String(formData.get("digestSpendCapUsd") ?? "").trim(), + ), + sections: + digestSectionsRaw.length > 0 ? digestSectionsRaw : dD.sections, + // Prompt SHAPE — both freshness-affecting (see digestPromptVariant), + // so a silent reset here would invalidate every digest generated + // under a non-default shape. Each still FALLS BACK to the current + // value rather than to the default when the field is absent: a form + // that rebuilds this block but omits a field is exactly how the reset + // bug happens, and the fallback is what makes omission harmless. + timestampMode: isDigestTimestampMode(digestTimestampModeRaw) + ? digestTimestampModeRaw + : dD.timestampMode, + promptVariant: formData.has("digestPromptVariant") + ? String(formData.get("digestPromptVariant") ?? "").trim() + : dD.promptVariant, + // Rendered by the form, so read it from the form — but only when the + // digest fields are actually present, which the branch already assures. + yieldToCpuWorkers: formData.get("digestYieldToCpuWorkers") === "on", + } + : dD; const dB = defaultBuildPipeline(); const buildModeRaw = String(formData.get("buildMode") ?? "").trim(); const buildPipeline = { mode: isBuildMode(buildModeRaw) ? buildModeRaw : dB.mode, - maxParallelBuilds: - Number.parseInt(String(formData.get("maxParallelBuilds") ?? "").trim(), 10), - dockerImage: String(formData.get("dockerImage") ?? "").trim() || dB.dockerImage, - dockerfile: String(formData.get("dockerfile") ?? "").trim() || dB.dockerfile, + maxParallelBuilds: Number.parseInt( + String(formData.get("maxParallelBuilds") ?? "").trim(), + 10, + ), + dockerImage: + String(formData.get("dockerImage") ?? "").trim() || dB.dockerImage, + dockerfile: + String(formData.get("dockerfile") ?? "").trim() || dB.dockerfile, }; const next: SiteSettings = { diff --git a/editor/app/settings/components/SettingsForm.tsx b/editor/app/settings/components/SettingsForm.tsx @@ -506,6 +506,29 @@ export function SettingsForm({ initial, apps, digestApps }: Props) { defaultValue={initial.digest.promptVariant} hint="Free-text label for a non-default prompt shape, recorded in every section's provenance (trimmed, max 40 chars). Setting or changing it invalidates digests generated under a different label — which is exactly what makes a bake-off round re-run its sample instead of skipping it as fresh. Leave blank unless you are running one." /> + <label className="flex items-start gap-2 text-sm"> + <input + type="checkbox" + name="digestYieldToCpuWorkers" + defaultChecked={initial.digest.yieldToCpuWorkers} + className="mt-1" + /> + <span className="flex flex-col gap-1"> + <span className="font-medium"> + Also yield the GPU to CPU-only transcription workers + </span> + <span className="text-xs text-muted-foreground"> + Off by default. The digest lane steps aside while transcription + uses the card, but a worker pinned to{" "} + <code>device: cpu</code> competes for no GPU shaders at all — + treating it as contention held the digest lane at zero throughput + for nothing. A worker with <em>no</em> device set still counts, + because that is the engine binary&apos;s own default and it may be + the GPU. Turn this on to make every local worker block the digest + lane regardless of device. + </span> + </span> + </label> <details className="text-sm"> <summary className="cursor-pointer font-medium"> Per-engine configuration