// ONE WAY TO RUN ONE OPERATION OVER ONE CHANNEL. // // Before this file there were four, and no two agreed. Each sweep had a // per-channel runner; the arbiter called one 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 SEQUENTIAL caller wants — awaiting is what lets a 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 call sites that want it, 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_OPERATION_ID } from "../lib/operations"; import { countOperationWork, runOperationBatch, type DigestOrder, } from "./operationBatch"; 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. A stage card's per-row "run only this" // passes one; the lane runner passes one; the group button passes none, // because that station IS the lane rather than one operation on it. kindIds?: string[]; // Only these videos (intersected with disk by the batch). A stage card passes // none; the video page passes one. In `spec.params` so a replay stays scoped. ids?: 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. A background dispatcher sets // it; a stage card does not. background?: boolean; onStarted?: (jobId: string) => void; onDone?: () => void; }; export async function runBackfillChannelJob( opts: BackfillChannelJobOptions, ): Promise { const { paths, channelSlug, kindIds, ids } = opts; const result = await runManagedFunction({ kind: BACKFILL_CHANNEL_KIND, queueKey: resolveQueueKey(BACKFILL_QUEUE, opts.queueKey), paths, channelSlug, background: opts.background, // ALL THREE fields, because jobReplayRegistry.ts replays this job from // them: the kind scope, the queue override and the id scope all have to // survive a replay — a replayed per-video run that widened to the whole // channel would be a different job wearing the same record. Undefined // fields drop out of the JSON, so a run with none stores none. spec: { kind: BACKFILL_CHANNEL_KIND, slug: channelSlug, params: { kindIds, queueKey: opts.queueKey, ids }, }, fn: async (onLog, signal, setProgress, ctx) => { const before = await countOperationWork("backfill", paths, channelSlug, { operationIds: kindIds, ids, }); const batch = await runOperationBatch({ lane: "backfill", channelSlug, paths, operationIds: kindIds, ids, 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/operations.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.reacquireHandedOff > 0 ? `, ${batch.reacquireHandedOff} handed to auto-transcribe` : "") : "") + (batch.diskFloorHit ? " (stopped re-acquiring at the disk floor)" : "") + (batch.stoppedForRelocation ? " (stopped: a relocation is in flight)" : "") + ".", ); 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; // Only these videos (intersected with disk by the batch), as on the backfill // job. In `spec.params` so a replay stays scoped. Absent = the channel. ids?: string[]; 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 { 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 lane runner 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, ids: opts.ids, }, }, 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 countOperationWork("digest", paths, channelSlug, { digestLane: lane, ids: opts.ids, }) ).reachable; // 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 runOperationBatch({ lane: "digest", channelSlug, paths, digestLane: lane, ids: opts.ids, digestOrder: 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)" : "") + (batch.stoppedForRelocation ? " (stopped: a relocation is in flight)" : "") + ".", ); // 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; // Only these videos; absent = the whole channel. ids?: 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 /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. It was open-coded // twice before — the retired arbiter with a ternary on DIGEST_OPERATION_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 { const { paths, channelSlug, operation } = opts; if (operation === DIGEST_OPERATION_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, ids: opts.ids, queueKey: opts.queueKey, background: opts.background, onStarted: opts.onStarted, onDone: opts.onDone, }); } return runBackfillChannelJob({ paths, channelSlug, kindIds: [operation], ids: opts.ids, // 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, }); }