Archilyzer · Source

archilyzer

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

commit 212f9175ecd5bafc3473237009643033a22db53f
parent 1b9d3716ed34696c7345bdbc1b23b5077f45279d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Wed, 26 Aug 2026 00:40:00 -0400

operationJobs: one way to run one operation over one channel

There were four, and no two agreed. The sweep had a per-channel runner; the
arbiter called it for backfill and a different one for digest; the editor's two
stage cards each carried a hand-written copy of the same runManagedFunction
call; the /channels group buttons called those copies.

They had drifted the way copies do. The editor's backfill summary printed
`deferred` and `blocked`; the sweep's printed `skipped`; neither printed the
other's. So "what did that run actually do" had a different answer depending on
which button started it — an operator watching a sweep could not see that 40
videos deferred over the diarization length limit, and one watching a card could
not see that 3 were skipped mid-rewrite. One line now carries all of them.
missingInput stays its own number, never summed (backfillKinds.ts's header: it
is ~91x the reachable figure on this corpus).

THE SPLIT THAT MATTERS: these runners start a job and RETURN ITS STREAM. They
do not drain it.

Draining is the SWEEP's property — sequential by design, and awaiting is what
lets its loop re-plan against real results — and it is exactly what a stage card
must not do: the editor consumes the returned stream client-side for a live log,
so a runner that drained would have consumed it first and left every card's log
dead. The drain moved to the three call sites that want it (backfillSweep,
digestSweep, arbiter), where "sequential by design" is a true statement about
the caller rather than a property silently imposed on everyone. The arbiter now
drains explicitly: it dispatches one unit per lane per pass and re-plans every
pass, so returning before the job finished would hand the next pass the same
channel.

And `onDone`, because common/ CANNOT import next/cache. revalidatePath is the
editor's business; the hook fires it at job end without this module knowing Next
exists.

The digest runner absorbs digestChannelAction's extras — queueKey, order,
limitCount, force, and the remoteEnabled fail-fast. The fail-fast MOVED rather
than staying in the editor: the arbiter can ask for the metered lane too, and a
guard only one of two callers performs is not a guard. spec.params is the
superset both copies stored between them; undefined fields drop out of JSON, so
a sweep's job stores exactly what it stored before and both
jobReplayRegistry.ts handlers read unchanged.

runOperationChannelJob is the dispatcher the arbiter and the group buttons
wanted: an operation id and a slug, without also having to know which of two
runners that id belongs to. The arbiter's ternary on DIGEST_KIND_ID is gone, and
with it a second derivation of the digest lane from settings — `unit.lane` is
live since the laneFor change, so the lane choice reads it.

laneForOperation moved to its own module (controller/operationLane.ts): the
arbiter plans against it and operationJobs reserves a queue key with it, while
the arbiter calls operationJobs. Left in arbiter.ts that was an import cycle.
Verified by loading all four modules under tsx.

The /channels backfill station deliberately does NOT go through the dispatcher:
that station is the LANE, not one operation on it — KIND_FOR says so by mapping
it to the lane's job kind — so it keeps calling backfillChannelAction with no
kindIds. The digest station keeps its LOCAL-lane pin, now stated as an explicit
lane rather than a defaulted one.

digestBucketAction is untouched: it takes an explicit id set the runner has no
parameter for.

tsc clean both packages; common 826/826.

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

Diffstat:
Mcommon/controller/arbiter.ts | 75+++++++++++++++++++++++++++++++--------------------------------------------
Mcommon/controller/backfillSweep.ts | 83++++++++++++++-----------------------------------------------------------------
Mcommon/controller/digestSweep.ts | 89+++++++++++--------------------------------------------------------------------
Mcommon/controller/laneForOperation.test.ts | 5++++-
Acommon/controller/operationJobs.ts | 335+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/controller/operationLane.ts | 31+++++++++++++++++++++++++++++++
Meditor/app/channels/[slug]/backfillActions.ts | 75++++++++++++++++++++-------------------------------------------------------
Meditor/app/channels/[slug]/digestActions.ts | 108+++++++++++++++++++------------------------------------------------------------
Meditor/app/channels/groupActions.ts | 25++++++++++++++++++++-----
9 files changed, 493 insertions(+), 333 deletions(-)

diff --git a/common/controller/arbiter.ts b/common/controller/arbiter.ts @@ -34,6 +34,7 @@ import type { Paths } from "../lib/paths"; import { getPaths } from "../lib/paths"; import { getSettings } from "../lib/settings"; +import { DIGEST_REMOTE_QUEUE } from "../lib/queueKeys"; import { getRegistry } from "../jobs/registry"; import { runManagedFunction } from "../jobs/streamCommand"; import { drainStream } from "../jobs/drainStream"; @@ -46,13 +47,12 @@ import { } from "../jobs/autoQueuePolicy"; import { DIGEST_KIND_ID, - getBackfillKind, allBackfillKinds, operationLabel, type BackfillLane, } from "../lib/backfillKinds"; -import { runDigestChannelJob } from "./digestSweep"; -import { runBackfillChannelJob } from "./backfillSweep"; +import { laneForOperation } from "./operationLane"; +import { runOperationChannelJob } from "./operationJobs"; import { buildArbiterChannelWork, listArbiterChannelMeta, @@ -107,28 +107,6 @@ export type ArbiterUnit = { leafId: string; }; -// The lane an operation runs on, or null when the registry does not know it. -// -// ONE RESOLUTION, off the kind's own declaration. A kind whose lane depends on -// configuration says so with laneFor() — digest picks between the GPU-bound -// local queue and the metered network one, diarization between contending for -// the GPU and contending for cores — and a kind without one has a fixed lane. -// This used to special-case digest here, which meant the rule lived in two -// places and only digest's copy was live: diarization's laneFor was invisible -// to the arbiter. -// -// getBackfillKind, NOT operationCatalog(), and that is deliberate. Today -// getBackfillKind("download") is undefined, so download and transcription get -// no lane and planArbiterUnits skips them (see the `if (!lane) continue` below) -// — which is correct, because they are dispatched by their own runners. A -// catalog lookup would hand them a lane and the arbiter would start a backfill -// channel job for work no backfill kind can do. -export function laneForOperation(operation: string): BackfillLane | null { - const kind = getBackfillKind(operation); - if (!kind) return null; - return kind.laneFor?.(getSettings()) ?? kind.lane; -} - // Turn the policy trees into dispatchable units, in priority order. // // PURE, and separated from the loop for exactly that reason: this is where the @@ -237,30 +215,39 @@ async function runOneUnit( // metric a hand-clicked run produces, so an arbiter-driven pass is // inspectable with the tools that already exist rather than being an opaque // mega-job — and so this loop is a scheduler rather than a second runtime. - const outcome = - unit.operation === DIGEST_KIND_ID - ? await runDigestChannelJob({ - paths, - channelSlug: unit.channelSlug, - lane: getSettings().digest.remoteEnabled ? "remote" : "local", - // Behind anything an operator clicks by hand. - background: true, - }) - : await runBackfillChannelJob({ - paths, - channelSlug: unit.channelSlug, - kindIds: [unit.operation], - // RESERVE BY THE LANE'S OWN KEY, never invent one. This is what keeps - // digest and diarization overlapping instead of taking turns, which - // backfill.spec.ts pins. - queueKey: unit.lane.queueKey, - background: true, - }); + // + // ONE CALL, whichever operation it is. This used to be a ternary on + // DIGEST_KIND_ID with two hand-written option sets, which meant the arbiter + // had to be edited in step with the registry and re-derived the digest lane + // from settings a second time — while `unit.lane` beside it was already the + // answer. It is now: laneForOperation asks the kind's own laneFor, so + // `unit.lane.queueKey` IS the live lane, and the choice below reads it rather + // than re-asking getSettings. + const outcome = await runOperationChannelJob({ + paths, + channelSlug: unit.channelSlug, + operation: unit.operation, + // RESERVE BY THE LANE'S OWN KEY, never invent one. This is what keeps digest + // and diarization overlapping instead of taking turns, which backfill.spec + // pins. + queueKey: unit.lane.queueKey, + digest: { + lane: unit.lane.queueKey === DIGEST_REMOTE_QUEUE ? "remote" : "local", + }, + // Behind anything an operator clicks by hand. + background: true, + }); if (!outcome.ok) { // One channel failing to START is not the arbiter failing. onLog(`!! ${unit.channelSlug}: ${outcome.error ?? "failed to start"}`); return false; } + // WAIT FOR IT. The arbiter dispatches ONE unit per lane per pass and re-plans + // every pass, so a pass that returned before its job finished would re-plan + // against work still in flight and hand the next pass the same channel. The + // runners used to drain internally; they no longer do — see + // controller/operationJobs.ts for why that had to move to the callers. + await drainStream(outcome.stream); return true; } diff --git a/common/controller/backfillSweep.ts b/common/controller/backfillSweep.ts @@ -31,8 +31,6 @@ import { getSettings, writeSettings } from "../lib/settings"; import { getRegistry } from "../jobs/registry"; import { runManagedFunction } from "../jobs/streamCommand"; import { drainStream } from "../jobs/drainStream"; -import { makeTaskTracker } from "../jobs/taskHooks"; -import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; import { BACKFILL_QUEUE } from "../lib/queueKeys"; import { reachableBackfillWork, @@ -40,14 +38,17 @@ import { type BackfillSnapshotEntry, } from "../lib/backfillKinds"; import { listChannelStatsFromDisk, readChannelSnapshot } from "./channels"; -import { countBackfillWork, runBackfillBatch } from "./backfillBatch"; +import { countBackfillWork } from "./backfillBatch"; +import { runBackfillChannelJob } from "./operationJobs"; import type { AutoQueueOrder } from "../jobs/autoQueuePolicy"; import { buildRecencyKeys } from "./recencyIndex"; import { orderPlanByRecency } from "./planOrder"; import type { SweepPlanEntry } from "../lib/sweepPlan"; export const BACKFILL_SWEEP_KIND = "backfill-sweep"; -export const BACKFILL_CHANNEL_KIND = "backfill-channel"; +// The per-channel runner and its job kind live in controller/operationJobs.ts +// now — one runner for a hand-clicked run, a swept one and an arbiter-dispatched +// one. This module is the SWEEP: the plan, the loop, and the drain. // How long the loop waits before re-planning when a pass did no work. Generous: // nothing here is latency-sensitive and planning reads every channel's videos. @@ -259,70 +260,6 @@ export function countFromSnapshot( return { reachable, missingInput }; } -// Run ONE channel through the batch, as its own managed job — same kind, same -// queue and same progress metric a hand-clicked run produces, so a sweep is -// inspectable with the tools that already exist rather than being an opaque -// mega-job. -export async function runBackfillChannelJob(opts: { - paths: Paths; - channelSlug: string; - kindIds?: string[]; - queueKey?: string; - background?: boolean; - onStarted?: (jobId: string) => void; -}): Promise<{ ok: boolean; error?: string }> { - const { paths, channelSlug } = opts; - const result = await runManagedFunction({ - kind: BACKFILL_CHANNEL_KIND, - queueKey: opts.queueKey ?? BACKFILL_QUEUE, - paths, - channelSlug, - background: opts.background, - spec: { - kind: BACKFILL_CHANNEL_KIND, - slug: channelSlug, - params: { kindIds: opts.kindIds, queueKey: opts.queueKey }, - }, - fn: async (onLog, signal, setProgress, ctx) => { - const before = await countBackfillWork(paths, channelSlug, opts.kindIds); - const batch = await runBackfillBatch({ - channelSlug, - paths, - kindIds: opts.kindIds, - setProgress, - // The bar measures THIS run, from zero. Seeding it with what is already - // on disk cannot represent a regeneration, where no file count moves. - progressBaseline: 0, - progressTarget: before.reachable, - onLog, - signal, - drainSignal: ctx.drainSignal, - tracker: makeTaskTracker(ctx, onLog), - }); - onLog( - `Backfill ${channelSlug}: ${batch.succeeded} done, ${batch.fresh} already current, ` + - `${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` + - (batch.skipped > 0 - ? `; ${batch.skipped} skipped (their transcript is mid-rewrite)` - : "") + - (batch.reacquired > 0 - ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` - : "") + - (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + - ".", - ); - requestChannelSnapshot(paths, channelSlug); - }, - }); - if (!result.ok) return { ok: false, error: result.error }; - opts.onStarted?.(result.jobId); - // Wait for the channel to finish: the sweep is sequential by design, and the - // queue would serialize these anyway — awaiting makes that explicit and lets - // the loop re-plan against real results. - await drainStream(result.stream); - return { ok: true }; -} - async function runSweepLoop( paths: Paths, kindIds: string[] | undefined, @@ -415,12 +352,20 @@ async function runSweepLoop( live.channelJobId = jobId; }, }); - live.channelJobId = null; if (!outcome.ok) { + live.channelJobId = null; // One channel failing to START is not the sweep failing. onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`); continue; } + // DRAIN HERE, not in the runner. The sweep is sequential by design and the + // queue would serialize these anyway — awaiting makes that explicit and + // lets the loop re-plan against real results. It is the CALLER's property, + // which is why the runner returns the stream instead of consuming it: a + // stage card renders that same stream as a live log, and a runner that + // drained would leave every card's log dead. + await drainStream(outcome.stream); + live.channelJobId = null; didWork = true; } diff --git a/common/controller/digestSweep.ts b/common/controller/digestSweep.ts @@ -37,10 +37,6 @@ import { getSettings, writeSettings } from "../lib/settings"; import { getRegistry } from "../jobs/registry"; import { runManagedFunction } from "../jobs/streamCommand"; import { drainStream } from "../jobs/drainStream"; -import { makeTaskTracker } from "../jobs/taskHooks"; -import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; -import { DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE } from "../lib/queueKeys"; -import { countMissingDigests, runDigestBatch } from "./digestBatch"; import { audioHours, buildDigestSweepPlan, @@ -100,78 +96,10 @@ export type DigestSweepOptions = { channelSlugs?: string[]; }; -// Run ONE channel through the existing per-channel batch, as its own managed -// job — same kind, same queue and same progress metric a hand-clicked run -// produces, so the sweep is inspectable with the tools that already exist -// rather than being an opaque mega-job. -export async function runDigestChannelJob(opts: { - paths: Paths; - channelSlug: string; - lane: DigestLaneChoice; - background?: boolean; - // Called with the job id as soon as it exists, so a caller can drain this - // specific job rather than only the thing that launched it. - onStarted?: (jobId: string) => void; - onDone?: () => void; -}): Promise<{ ok: boolean; error?: string }> { - const { paths, channelSlug, lane } = opts; - const settings = getSettings(); - const remote = lane === "remote"; - const kind = remote ? "digest-channel-remote" : "digest-channel-local"; - - const result = await runManagedFunction({ - kind, - queueKey: remote ? DIGEST_REMOTE_QUEUE : DIGEST_LOCAL_QUEUE, - paths, - channelSlug, - background: opts.background, - spec: { kind, slug: channelSlug, params: { lane } }, - fn: async (onLog, signal, setProgress, ctx) => { - const missing = await countMissingDigests( - paths, - channelSlug, - undefined, - lane, - ); - const batch = await runDigestBatch({ - channelSlug, - paths, - lane, - ...(remote - ? { minDurationSeconds: settings.digest.longTailSeconds } - : {}), - setProgress, - // The bar measures THIS run, from zero. Seeding it with the count of - // digests already on disk was the old shape and it could not represent - // a regeneration, where the file count never moves. - progressBaseline: 0, - progressTarget: missing, - onLog, - signal, - drainSignal: ctx.drainSignal, - tracker: makeTaskTracker(ctx, onLog), - }); - onLog( - `Digest batch: ${batch.succeeded} generated, ${batch.fresh} already current, ` + - `${batch.shared} shared to mirrors, ${batch.misaligned} mirror(s) refused by the alignment gate, ` + - `${batch.skipped} skipped, ${batch.failed} failed; ` + - `${batch.engineCalls} model call(s), ${batch.warnings} warning(s)` + - (batch.costUsd > 0 ? `, $${batch.costUsd.toFixed(4)}` : "") + - (batch.spendCapped ? " (stopped at the spend cap)" : "") + - ".", - ); - requestChannelSnapshot(paths, channelSlug); - opts.onDone?.(); - }, - }); - if (!result.ok) return { ok: false, error: result.error }; - opts.onStarted?.(result.jobId); - // Wait for the channel to finish before returning: the sweep is sequential by - // design (one GPU), and the queue would serialize these anyway — awaiting - // makes that explicit and lets the loop re-plan against real results. - await drainStream(result.stream); - return { ok: true }; -} +// The per-channel runner lives in controller/operationJobs.ts now, beside the +// backfill one. It no longer drains — see the call site below for where that +// moved to, and why. +import { runDigestChannelJob } from "./operationJobs"; async function runSweepLoop( paths: Paths, @@ -264,12 +192,19 @@ async function runSweepLoop( live.channelJobId = jobId; }, }); - live.channelJobId = null; if (!outcome.ok) { + live.channelJobId = null; // One channel failing to START is not the sweep failing. onLog(`!! ${channel.channelSlug}: ${outcome.error ?? "failed to start"}`); continue; } + // DRAIN HERE, not in the runner. The sweep is sequential by design (one + // GPU) and the queue would serialize these anyway — awaiting makes that + // explicit and lets the loop re-plan against real results. It belongs to + // the CALLER: the editor's digest card renders this same stream as a live + // log, and a runner that drained would have consumed it first. + await drainStream(outcome.stream); + live.channelJobId = null; didWork = true; } diff --git a/common/controller/laneForOperation.test.ts b/common/controller/laneForOperation.test.ts @@ -7,6 +7,9 @@ // the env must be set before anything imports the module under test. // arbiter.test.ts is deliberately settings-free — it passes `laneOf` in — and // these cases are exactly the ones a stub would answer vacuously. +// +// Named for the function rather than for controller/operationLane.ts, which is +// the one-function module it now lives in. import { mkdtempSync, writeFileSync } from "node:fs"; import { rm } from "node:fs/promises"; @@ -25,7 +28,7 @@ process.env.TRANSCRIPTS_DIR = ROOT; const SETTINGS_FILE = path.join(ROOT, "settings.json"); process.env.SETTINGS_FILE = SETTINGS_FILE; -const { laneForOperation } = await import("./arbiter"); +const { laneForOperation } = await import("./operationLane"); const { DIGEST_LOCAL_QUEUE, DIGEST_REMOTE_QUEUE, BACKFILL_QUEUE } = await import("../lib/queueKeys"); diff --git a/common/controller/operationJobs.ts b/common/controller/operationJobs.ts @@ -0,0 +1,335 @@ +// ONE WAY TO RUN ONE OPERATION OVER ONE CHANNEL. +// +// Before this file there were four, and no two agreed. The sweep had a +// per-channel runner; the arbiter called it for backfill and a different one for +// digest; the editor's stage cards each had a hand-written copy of the same +// runManagedFunction call; the /channels group buttons called those copies. The +// copies had drifted in the way copies do — the editor's backfill summary line +// printed `deferred` and `blocked`, the sweep's printed `skipped`, and neither +// printed the other's — so "what did that run actually do" had a different +// answer depending on which button started it. +// +// THE SPLIT THAT MATTERS: these runners START a job and RETURN ITS STREAM. They +// do not drain it. +// +// Draining is what a SWEEP wants — it is sequential by design, and awaiting is +// what lets its loop re-plan against real results — and it is exactly what a +// stage card must not do: the editor consumes the returned stream client-side to +// draw a live log, and a runner that drained would have consumed it first and +// left every card's log dead. So the drain lives at the three call sites that +// want it (backfillSweep, digestSweep, arbiter), which is also where "sequential +// by design" is a true statement about the caller rather than a property +// silently imposed on everyone. +// +// AND `onDone`, because common/ CANNOT import next/cache. revalidatePath is the +// editor's business; the hook is how a server action gets it to fire at job end +// without this module knowing Next exists. The digest runner already had one. + +import type { Paths } from "../lib/paths"; +import { getSettings } from "../lib/settings"; +import { + runManagedFunction, + type StreamActionResult, +} from "../jobs/streamCommand"; +import { makeTaskTracker } from "../jobs/taskHooks"; +import { requestChannelSnapshot } from "../jobs/snapshotScheduler"; +import { + BACKFILL_QUEUE, + DIGEST_LOCAL_QUEUE, + DIGEST_REMOTE_QUEUE, + resolveQueueKey, +} from "../lib/queueKeys"; +import { DIGEST_KIND_ID } from "../lib/backfillKinds"; +import { countBackfillWork, runBackfillBatch } from "./backfillBatch"; +import { + countMissingDigests, + runDigestBatch, + type DigestOrder, +} from "./digestBatch"; +import type { DigestLaneChoice } from "./digestTarget"; +import { laneForOperation } from "./operationLane"; + +export const BACKFILL_CHANNEL_KIND = "backfill-channel"; + +function isDigestOrder(v: unknown): v is DigestOrder { + return v === "shortest-first" || v === "longest-first"; +} + +// ── The backfill lane, over one channel ──────────────────────────────────── + +export type BackfillChannelJobOptions = { + paths: Paths; + channelSlug: string; + // Absent = every enabled lane kind. The arbiter passes exactly one; a stage + // card's per-row "run only this" passes one; the group button passes none, + // because that station IS the lane rather than one operation on it. + kindIds?: string[]; + // An operator's queue override from the card's QueueControl. `undefined` is + // "no control wired up" and takes the lane's own key; "" means immediate. + queueKey?: string; + // Behind anything an operator asks for by hand. The sweeps and the arbiter + // set it; a stage card does not. + background?: boolean; + onStarted?: (jobId: string) => void; + onDone?: () => void; +}; + +export async function runBackfillChannelJob( + opts: BackfillChannelJobOptions, +): Promise<StreamActionResult> { + const { paths, channelSlug, kindIds } = opts; + const result = await runManagedFunction({ + kind: BACKFILL_CHANNEL_KIND, + queueKey: resolveQueueKey(BACKFILL_QUEUE, opts.queueKey), + paths, + channelSlug, + background: opts.background, + // BOTH fields, because jobReplayRegistry.ts reads `p.kindIds` to replay this + // job and the queue override has to survive a replay too. Undefined fields + // drop out of the JSON, so a run with neither stores neither. + spec: { + kind: BACKFILL_CHANNEL_KIND, + slug: channelSlug, + params: { kindIds, queueKey: opts.queueKey }, + }, + fn: async (onLog, signal, setProgress, ctx) => { + const before = await countBackfillWork(paths, channelSlug, kindIds); + const batch = await runBackfillBatch({ + channelSlug, + paths, + kindIds, + setProgress, + // The bar measures THIS run, from zero. Seeding it with what is already + // on disk cannot represent a regeneration, where no file count moves. + progressBaseline: 0, + progressTarget: before.reachable, + onLog, + signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), + }); + // ONE SUMMARY LINE, carrying every outcome the batch distinguishes. The + // two copies this replaces each printed a different subset — the card's + // had deferred and blocked, the sweep's had skipped — so an operator + // reading a sweep's log could not see that 40 videos deferred, and one + // reading a card's could not see that 3 were mid-rewrite. + // + // `missingInput` is printed SEPARATELY and never added to anything. See + // lib/backfillKinds.ts's header: on this corpus it is ~91x the reachable + // figure, and one summed "remaining" would be noise. + onLog( + `Backfill ${channelSlug}: ${batch.succeeded} done, ${batch.fresh} already current, ` + + `${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` + + (batch.skipped > 0 + ? `; ${batch.skipped} skipped (their transcript is mid-rewrite)` + : "") + + (batch.deferred > 0 + ? `; ${batch.deferred} deferred (over the diarization length limit)` + : "") + + (batch.blocked > 0 + ? `; ${batch.blocked} waiting on a prerequisite backfill` + : "") + + (batch.reacquired > 0 + ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` + : "") + + (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + + ".", + ); + requestChannelSnapshot(paths, channelSlug); + opts.onDone?.(); + }, + }); + // The job id as soon as it exists, so a caller can drain THIS job rather than + // only the thing that launched it. Not fired on failure: there is no job. + if (result.ok) opts.onStarted?.(result.jobId); + return result; +} + +// ── The digest lane, over one channel ────────────────────────────────────── + +export type DigestChannelJobOptions = { + paths: Paths; + channelSlug: string; + lane: DigestLaneChoice; + queueKey?: string; + order?: string; + limitCount?: number; + force?: boolean; + background?: boolean; + onStarted?: (jobId: string) => void; + onDone?: () => void; +}; + +// ONE JOB PER CHANNEL, deliberately: job logs keep only the newest 500 (30 days) +// and the in-memory registry keeps 100 records, so a job per video would evict +// the whole history of a sweep — including the running jobs' own logs. +export async function runDigestChannelJob( + opts: DigestChannelJobOptions, +): Promise<StreamActionResult> { + const { paths, channelSlug, lane } = opts; + const settings = getSettings(); + // FAIL FAST, before runManagedFunction, so an off metered lane produces an + // error the caller can show rather than a started job that dies. Lives here + // rather than in the editor action because the arbiter can ask for the remote + // lane too, and a check only one caller performs is a check. + if (lane === "remote" && !settings.digest.remoteEnabled) { + return { + ok: false, + error: + "The metered digest lane is off. Enable it in Settings → Digest before running it.", + }; + } + const remote = lane === "remote"; + const kind = remote ? "digest-channel-remote" : "digest-channel-local"; + + const result = await runManagedFunction({ + kind, + // Separate keys per lane so the two run CONCURRENTLY — the local lane is + // GPU-bound and the metered lane is network-bound, so serializing them + // would waste half the throughput of a multi-week sweep. + queueKey: resolveQueueKey( + remote ? DIGEST_REMOTE_QUEUE : DIGEST_LOCAL_QUEUE, + opts.queueKey, + ), + paths, + channelSlug, + background: opts.background, + // THE SUPERSET the two copies stored between them. jobReplayRegistry.ts + // reads p.lane/order/limitCount/force; undefined fields drop out of JSON, + // so a sweep's job stores `{lane}` exactly as it did before. + spec: { + kind, + slug: channelSlug, + params: { + queueKey: opts.queueKey, + lane, + order: opts.order, + limitCount: opts.limitCount, + force: opts.force, + }, + }, + fn: async (onLog, signal, setProgress, ctx) => { + // The LANE matters: the progress target must be computed against the + // engine that is actually about to run, not whichever one the local + // setting names. + const missing = await countMissingDigests( + paths, + channelSlug, + undefined, + lane, + ); + // The bar measures THIS RUN, from zero, and the batch reports its own + // `current`. It used to be seeded with the count of digests already on + // disk and left to the UI's disk re-count — which cannot move during a + // REGENERATION, because a regenerated digest is rewritten in place. The + // bar sat at 0% for whole jobs ({initial:39, current:39, target:72}). + const batch = await runDigestBatch({ + channelSlug, + paths, + lane, + order: isDigestOrder(opts.order) ? opts.order : undefined, + // The local lane takes everything up to the long-tail cutoff; the + // metered lane exists for the tail above it. Passing no window (the + // default) runs the whole channel on one lane, which is what a + // single-lane sweep wants. + ...(remote + ? { minDurationSeconds: settings.digest.longTailSeconds } + : {}), + limitCount: + typeof opts.limitCount === "number" && opts.limitCount > 0 + ? opts.limitCount + : undefined, + force: opts.force === true, + setProgress, + progressBaseline: 0, + progressTarget: missing, + onLog, + signal, + drainSignal: ctx.drainSignal, + tracker: makeTaskTracker(ctx, onLog), + }); + onLog( + `Digest batch: ${batch.succeeded} generated, ${batch.fresh} already current, ` + + `${batch.shared} shared to mirrors, ${batch.misaligned} mirror(s) refused by the alignment gate, ` + + `${batch.skipped} skipped, ${batch.failed} failed; ` + + `${batch.engineCalls} model call(s), ${batch.warnings} warning(s)` + + (batch.costUsd > 0 ? `, $${batch.costUsd.toFixed(4)}` : "") + + (batch.spendCapped ? " (stopped at the spend cap)" : "") + + ".", + ); + // ONCE, at job end — the digest kinds are in NO_REGEN_KINDS precisely so + // that per-video regens don't fire. See snapshotScheduler.ts. + requestChannelSnapshot(paths, channelSlug); + opts.onDone?.(); + }, + }); + if (result.ok) opts.onStarted?.(result.jobId); + return result; +} + +// ── The dispatcher ───────────────────────────────────────────────────────── + +export type OperationChannelJobOptions = { + paths: Paths; + channelSlug: string; + // A catalog operation id — "digest", "diarization", "attribution-text", … + operation: string; + queueKey?: string; + background?: boolean; + onStarted?: (jobId: string) => void; + onDone?: () => void; + // Digest's extras. Ignored for every other operation, and typed as its own + // object rather than flattened so that is obvious at the call site. + digest?: { + lane?: DigestLaneChoice; + order?: string; + limitCount?: number; + force?: boolean; + }; +}; + +// RUN ONE OPERATION OVER ONE CHANNEL, whichever operation it is. +// +// This is the seam the arbiter and the /channels group buttons wanted: a caller +// that has an operation id and a channel slug should not also have to know which +// of two runners that id belongs to, nor which queue key to reserve. Both were +// open-coded before — the arbiter with a ternary on DIGEST_KIND_ID, the group +// buttons with a per-station map — and both had to be edited in step with the +// registry. +export function runOperationChannelJob( + opts: OperationChannelJobOptions, +): Promise<StreamActionResult> { + const { paths, channelSlug, operation } = opts; + if (operation === DIGEST_KIND_ID) { + return runDigestChannelJob({ + paths, + channelSlug, + // The configured engine decides, so a caller that does not care gets the + // lane the operator actually has switched on. An explicit choice still + // wins — the group buttons pin "local" on purpose. + lane: + opts.digest?.lane ?? + (getSettings().digest.remoteEnabled ? "remote" : "local"), + order: opts.digest?.order, + limitCount: opts.digest?.limitCount, + force: opts.digest?.force, + queueKey: opts.queueKey, + background: opts.background, + onStarted: opts.onStarted, + onDone: opts.onDone, + }); + } + return runBackfillChannelJob({ + paths, + channelSlug, + kindIds: [operation], + // RESERVE BY THE LANE'S OWN KEY, never invent one. This is what keeps digest + // and diarization overlapping instead of taking turns, which backfill.spec + // pins. laneForOperation is live (it asks the kind's laneFor), so a caller + // passing no override gets the queue the operation would actually run on. + queueKey: opts.queueKey ?? laneForOperation(operation)?.queueKey, + background: opts.background, + onStarted: opts.onStarted, + onDone: opts.onDone, + }); +} diff --git a/common/controller/operationLane.ts b/common/controller/operationLane.ts @@ -0,0 +1,31 @@ +// WHICH LANE AN OPERATION RUNS ON. One function, one module. +// +// Its own file because two things need it and neither may import the other: the +// arbiter plans units against it, and controller/operationJobs.ts reserves a +// queue key with it — while the arbiter calls operationJobs to start the jobs it +// plans. Left in arbiter.ts, that was an import cycle. + +import { getSettings } from "../lib/settings"; +import { getBackfillKind, type BackfillLane } from "../lib/backfillKinds"; + +// The lane an operation runs on, or null when the registry does not know it. +// +// ONE RESOLUTION, off the kind's own declaration. A kind whose lane depends on +// configuration says so with laneFor() — digest picks between the GPU-bound +// local queue and the metered network one, diarization between contending for +// the GPU and contending for cores — and a kind without one has a fixed lane. +// This used to special-case digest here, which meant the rule lived in two +// places and only digest's copy was live: diarization's laneFor was invisible +// to the arbiter. +// +// getBackfillKind, NOT operationCatalog(), and that is deliberate. Today +// getBackfillKind("download") is undefined, so download and transcription get +// no lane and planArbiterUnits skips them — which is correct, because they are +// dispatched by their own runners. A +// catalog lookup would hand them a lane and the arbiter would start a backfill +// channel job for work no backfill kind can do. +export function laneForOperation(operation: string): BackfillLane | null { + const kind = getBackfillKind(operation); + if (!kind) return null; + return kind.laneFor?.(getSettings()) ?? kind.lane; +} diff --git a/editor/app/channels/[slug]/backfillActions.ts b/editor/app/channels/[slug]/backfillActions.ts @@ -3,69 +3,34 @@ import { revalidatePath } from "next/cache"; import { getPaths } from "yt-dlp-transcript-common/lib/paths"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; -import { BACKFILL_QUEUE, resolveQueueKey } from "yt-dlp-transcript-common/lib/queueKeys"; -import { runManagedFunction } from "yt-dlp-transcript-common/jobs/streamCommand"; -import { makeTaskTracker } from "yt-dlp-transcript-common/jobs/taskHooks"; -import { requestChannelSnapshot } from "yt-dlp-transcript-common/jobs/snapshotScheduler"; -import { - countBackfillWork, - runBackfillBatch, -} from "yt-dlp-transcript-common/controller/backfillBatch"; +import { runBackfillChannelJob } from "yt-dlp-transcript-common/controller/operationJobs"; // Run the backfill lane over ONE channel, from its stage card. // -// Deliberately the same kind, queue and progress metric the sweep's per-channel -// job uses (runBackfillChannelJob), so a hand-clicked run and a swept one are -// the same thing seen from /jobs — that is what makes a sweep inspectable with -// the tools that already exist. It is NOT that function, though, because a -// hand-clicked run must be foreground (the sweep's are `background: true`, -// deliberately behind anything an operator asks for by hand) and takes a queue -// override from the card's QueueControl. +// A THIN WRAPPER, and that is the whole of it now. This used to be a +// hand-written copy of the sweep's per-channel job — the same runManagedFunction +// call, the same batch, a summary line that had drifted (it printed `deferred` +// and `blocked` where the sweep's printed `skipped`, so the same run read +// differently depending on which button started it). The comment here claimed +// the duplication existed because a hand-clicked run must be foreground and take +// a queue override; both were already parameters of the function it was copying. +// +// What is left is what genuinely belongs to the editor: +// - NO `background`, so a hand-clicked run goes in front of the sweep's. +// - `onDone`, because common/ cannot import next/cache. The runner calls it at +// job end and revalidatePath happens here, where Next exists. +// - The stream comes back UNDRAINED, which the stage card's StreamActionLog +// consumes to draw a live log. A runner that drained would kill it. export async function backfillChannelAction( slug: string, queueKey?: string, kindIds?: string[], ): Promise<StreamActionResult> { - const paths = getPaths(); - const kind = "backfill-channel"; - return runManagedFunction({ - kind, - queueKey: resolveQueueKey(BACKFILL_QUEUE, queueKey), - paths, + return runBackfillChannelJob({ + paths: getPaths(), channelSlug: slug, - spec: { kind, slug, params: { queueKey, kindIds } }, - fn: async (onLog, signal, setProgress, ctx) => { - const before = await countBackfillWork(paths, slug, kindIds); - const batch = await runBackfillBatch({ - channelSlug: slug, - paths, - kindIds, - setProgress, - progressBaseline: 0, - progressTarget: before.reachable, - onLog, - signal, - drainSignal: ctx.drainSignal, - tracker: makeTaskTracker(ctx, onLog), - }); - // BOTH numbers, never their sum — see lib/backfillKinds.ts. - onLog( - `Backfill ${slug}: ${batch.succeeded} done, ${batch.fresh} already current, ` + - `${batch.failed} failed; ${batch.missingInput} still need their media re-acquired` + - (batch.deferred > 0 - ? `; ${batch.deferred} deferred (over the diarization length limit)` - : "") + - (batch.blocked > 0 - ? `; ${batch.blocked} waiting on a prerequisite backfill` - : "") + - (batch.reacquired > 0 - ? `; ${batch.reacquired} re-acquired, ${batch.reacquireCleaned} cleaned up` - : "") + - (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + - ".", - ); - requestChannelSnapshot(paths, slug); - revalidatePath(`/channels/${slug}`); - }, + kindIds, + queueKey, + onDone: () => revalidatePath(`/channels/${slug}`), }); } diff --git a/editor/app/channels/[slug]/digestActions.ts b/editor/app/channels/[slug]/digestActions.ts @@ -11,7 +11,6 @@ import { import { countMissingDigests, runDigestBatch, - type DigestOrder, } from "yt-dlp-transcript-common/controller/digestBatch"; import { runManagedFunction, @@ -32,17 +31,26 @@ import { type EffectiveDigest, } from "yt-dlp-transcript-common/lib/digest"; import path from "node:path"; +import { runDigestChannelJob } from "yt-dlp-transcript-common/controller/operationJobs"; +import type { DigestLaneChoice } from "yt-dlp-transcript-common/controller/digestTarget"; -function isDigestOrder(v: unknown): v is DigestOrder { - return v === "shortest-first" || v === "longest-first"; -} - -export type DigestLaneChoice = "local" | "remote"; +// RE-EXPORTED, not redeclared: the definition lives in +// controller/digestTarget.ts, and this alias is what the client components have +// always imported. Two declarations of the same two strings is how they would +// drift. +export type { DigestLaneChoice }; -// Run the digest sweep over one channel. ONE JOB PER CHANNEL, deliberately: -// job logs keep only the newest 500 (30 days) and the in-memory registry keeps -// 100 records, so a job per video would evict the whole history of a sweep — -// including the running jobs' own logs. +// Run the digest lane over one channel, from its stage card. +// +// A THIN WRAPPER over the shared runner, for the same reasons backfillActions.ts +// is: what it keeps is the editor's half — no `background`, so a hand-clicked run +// goes in front of the sweep's, and `onDone` for the revalidatePath that common/ +// cannot do (it may not import next/cache). The stream comes back UNDRAINED for +// the card's live log. +// +// The remoteEnabled fail-fast MOVED INTO THE RUNNER rather than staying here. +// The arbiter can ask for the metered lane too, and a guard only one of two +// callers performs is not a guard. export async function digestChannelAction( slug: string, lane: DigestLaneChoice = "local", @@ -51,79 +59,15 @@ export async function digestChannelAction( limitCount?: number, force?: boolean, ): Promise<StreamActionResult> { - const paths = getPaths(); - const settings = getSettings(); - if (lane === "remote" && !settings.digest.remoteEnabled) { - return { - ok: false, - error: - "The metered digest lane is off. Enable it in Settings → Digest before running it.", - }; - } - const remote = lane === "remote"; - const kind = remote ? "digest-channel-remote" : "digest-channel-local"; - return runManagedFunction({ - kind, - // Separate keys per lane so the two run CONCURRENTLY — the local lane is - // GPU-bound and the metered lane is network-bound, so serializing them would - // waste half the throughput of a multi-week sweep. - queueKey: resolveQueueKey( - remote ? DIGEST_REMOTE_QUEUE : DIGEST_LOCAL_QUEUE, - queueKey, - ), - paths, + return runDigestChannelJob({ + paths: getPaths(), channelSlug: slug, - spec: { - kind, - slug, - params: { queueKey, lane, order, limitCount, force }, - }, - fn: async (onLog, signal, setProgress, ctx) => { - // The LANE matters: the progress target must be computed against the - // engine that is actually about to run, not whichever one the local - // setting names. - const missing = await countMissingDigests(paths, slug, undefined, lane); - // The bar measures THIS RUN, from zero, and the batch reports its own - // `current`. It used to be seeded with the count of digests already on - // disk and left to the UI's disk re-count — which cannot move during a - // REGENERATION, because a regenerated digest is rewritten in place. The - // bar sat at 0% for whole jobs ({initial:39, current:39, target:72}). - const result = await runDigestBatch({ - channelSlug: slug, - paths, - lane, - order: isDigestOrder(order) ? order : undefined, - // The local lane takes everything up to the long-tail cutoff; the metered - // lane exists for the tail above it. Passing no window (the default) runs - // the whole channel on one lane, which is what a single-lane sweep wants. - ...(remote - ? { minDurationSeconds: settings.digest.longTailSeconds } - : {}), - limitCount: - typeof limitCount === "number" && limitCount > 0 ? limitCount : undefined, - force: force === true, - setProgress, - progressBaseline: 0, - progressTarget: missing, - onLog, - signal, - drainSignal: ctx.drainSignal, - tracker: makeTaskTracker(ctx, onLog), - }); - onLog( - `Digest batch: ${result.succeeded} generated, ${result.fresh} already current, ` + - `${result.shared} shared to mirrors, ${result.misaligned} mirror(s) refused by the alignment gate, ` + - `${result.skipped} skipped, ${result.failed} failed; ` + - `${result.engineCalls} model call(s), ${result.warnings} warning(s)` + - (result.costUsd > 0 ? `, $${result.costUsd.toFixed(4)}` : "") + - (result.spendCapped ? " (stopped at the spend cap)" : "") + - ".", - ); - // ONCE, at job end — the digest kinds are in NO_REGEN_KINDS precisely so - // that per-video regens don't fire. See snapshotScheduler.ts. - requestChannelSnapshot(paths, slug); - revalidatePath(`/channels/${slug}`); - }, + lane, + queueKey, + order, + limitCount, + force, + onDone: () => revalidatePath(`/channels/${slug}`), }); } diff --git a/editor/app/channels/groupActions.ts b/editor/app/channels/groupActions.ts @@ -1,5 +1,6 @@ "use server"; +import { revalidatePath } from "next/cache"; import { listChannelBriefs } from "yt-dlp-transcript-common/controller/channels"; import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; import { activeSlugsForKinds } from "yt-dlp-transcript-common/jobs/syncJobs"; @@ -16,8 +17,9 @@ import { import { queueForSlugs } from "./lib/queueForSlugs"; import { downloadMissingAction, syncAction } from "./[slug]/pipelineActions"; import { transcribeMissingAction } from "./[slug]/whisperActions"; -import { digestChannelAction } from "./[slug]/digestActions"; import { backfillChannelAction } from "./[slug]/backfillActions"; +import { runOperationChannelJob } from "yt-dlp-transcript-common/controller/operationJobs"; +import { DIGEST_KIND_ID } from "yt-dlp-transcript-common/lib/backfillKinds"; // Run one pipeline stage over every channel in one of a site's groups. // @@ -59,10 +61,23 @@ const RUN_FOR: Record< // `skipped` for free. download: (slug) => downloadMissingAction(slug), transcribe: (slug) => transcribeMissingAction(slug), - // The LOCAL lane, the one the sweep uses. The metered lane is behind a - // settings gate and a spend cap, so it is never what a group button starts. - digest: (slug) => digestChannelAction(slug, "local"), - // Every enabled lane kind. + // THE LOCAL LANE, PINNED. The metered lane is behind a settings gate and a + // spend cap, so it is never what a group button starts — which is why this + // names the lane explicitly instead of letting the dispatcher pick the + // configured one. Through runOperationChannelJob so a station that IS one + // operation goes through the one runner path. + digest: (slug) => + runOperationChannelJob({ + paths: getPaths(), + channelSlug: slug, + operation: DIGEST_KIND_ID, + digest: { lane: "local" }, + onDone: () => revalidatePath(`/channels/${slug}`), + }), + // EVERY ENABLED LANE KIND, and so NOT runOperationChannelJob: this station is + // the LANE, not one operation on it — KIND_FOR above says the same thing by + // mapping it to the lane's job kind. Passing no kindIds is what "all of them" + // means to the runner. backfill: (slug) => backfillChannelAction(slug), };