import type { ChannelConfig } from "../../lib/channelConfig"; import { digestWorkOf, excludedDownloadIdSet, type ChannelSnapshot, } from "../../controller/channelSnapshot"; import { digestCountOf } from "../../controller/channels"; import { backfillLaneEntriesOf, operationCatalog, operationLabel, operationsActionLabel, operationsGroupLabel, presentOperationWork, reachableOperationWork, type Operation, } from "../../lib/operations"; import type { OperationBand } from "./band"; import { buildChannelBands } from "./buildBands"; import { normalizeBuckets, type StageId, type StageStatus, type StageTone, } from "./stageStatus"; // THE CHANNEL LINE. // // A channel's lifecycle is not ten sibling cards — it is a conserved quantity // moving through stations. Every video enters at the playlist and either // advances or leaves the line. This module is the model behind that picture: // stations are stages, and the GAP between two stations carries the shortfall, // because the gap is the work. // // SNAPSHOT-ONLY, and deliberately pure. common/controller/noCorpusWalkInRender- // Paths.test.ts bans the corpus walk from render paths, and page.tsx documents // the multi-minute regression from generating a report inside a GET. The one // value this cannot derive from the snapshot — how many videos the playlist // file names — is passed IN by the caller (one small readFile), never read here. // // Two invariants this file exists to hold: // // 1. WORK ON THE LINE IS NEVER SUMMED WITH WORK OFF IT. `reachable` is what the // lane can do today (missing + stale + partial). `missingInput`, `deferred`, // `blocked`, `excludedFromDownload` and `untranscribable` have LEFT the line // and live in `sidings`. Four surfaces once summed Object.values(snapshot // .backfill) and put every channel permanently at the top of every list. // Different fields on different axes is what makes the mistake unspellable. // 2. UNKNOWN IS NOT ZERO. A third of the snapshots on disk predate `eligible`, // so `present`/`coverage` are `number | null` and a reader must render "—". // A 0 there reads as "nothing digested" on a fully digested channel. // // AND THE ONE THIS FILE USED TO BREAK. The `backfill` station set `through` and // `denominator` by SUMMING THREE OPERATIONS — the one thing invariant 1 forbids // everywhere else, hidden behind a station label that named the queue rather // than the work. On the live corpus it was adding audio passes (diarization: 4 // done of 11,338) to per-chunk model calls (attribution-text: 1 done of 11,338, // and its cost basis is the transcript CHUNK, not the video), and calling the result // "Backfill". A station now carries its group's OPERATIONS, each with its own // band and its own denominator, and the numeral above them belongs to exactly // one of them — see leadOf. export type FlowStationId = | "playlist" | "download" | "transcribe" | "digest" // Renamed with StageId — see stageStatus.ts. The station and the stage card it // links to must carry the same id or `?stage=${station.stage}` opens the wrong // panel. | "speakers"; // One pipeline drawn under a station. The band is the same instrument the // /channels strip and the /operations rail draw, at station scale — which is // what makes a figure here and a figure there impossible to disagree. export type StationOperation = { id: string; label: string; shortLabel: string; // What one video costs, in words — the cost basis. Printed wherever the operation is armed, so // an 11,337-video backlog of per-chunk model calls cannot read as a quiet row. costBasis: string; band: OperationBand; }; export type FlowStation = { id: FlowStationId; // DERIVED for a station that holds several operations, never hardcoded: a // lane holding a mix of groups falls back to "Derived data" rather than // advertising one member's name. See operationsGroupLabel. label: string; // Videos that have cleared this station. NULL when the snapshot cannot say — // see invariant 2 above. Only the digest and backfill stations can be null; // the rest come from `totals`, which every snapshot carries. through: number | null; // The eligible population. null = unknown. denominator: number | null; // through / denominator, 0..1. Null whenever either side is unknown — render // "—", never 0. coverage: number | null; running: boolean; tone: StageTone; // Which stage panel this station opens (?stage=). stage: StageId; // The pipelines that run at this station, in dependency order. Empty for // playlist, which is not an operation the registry dispatches or counts — it // keeps the plain coverage meter. operations: StationOperation[]; }; // A population that has LEFT the line: it is not work the lane can pick up, and // it must never be added to a gap's `reachable`. export type Siding = { label: string; count: number; stage: StageId; hint?: string; }; export type FlowGap = { from: FlowStationId; to: FlowStationId; // Work the lane can do today. NEVER summed with `sidings`. reachable: number; // "to download", "to transcribe", … label: string; // Where ?stage= sends you — the stage that OWNS this gap's work, i.e. the // destination station's stage. stage: StageId; // The destination station's off-line populations. Sidings hang below the gap // in the rendered line, on a different axis from `reachable`, so the layout // itself cannot sum them. sidings: Siding[]; }; export type ChannelFlow = { stations: FlowStation[]; gaps: FlowGap[]; // The `from` id of the largest reachable gap — the one the renderer promotes // typographically. Null when nothing is reachable anywhere. bottleneck: FlowStationId | null; // The single primary action offered on the page. See pickNext below for why // this is the FURTHEST UPSTREAM gap rather than the biggest one. next: { stage: StageId; label: string; count: number } | null; }; export type ComputeChannelFlowInput = { snapshot: ChannelSnapshot; // Reused, never recomputed: tone and running are stageStatus's job and a // second opinion about them is a second thing to keep in sync. stages: Record; config: ChannelConfig; failedVideoIds: string[]; // Enabled lane kinds. EMPTY MEANS THE LANE IS OFF, which is not the same as // finished — see the tone rule at the bottom of this file. laneOperations: Operation[]; // How many videos the channel's `playlist` file names, or null when there is // no playlist file. Read by the caller (countPlaylist) so this stays pure. playlistCount: number | null; }; function ratio(through: number | null, denominator: number | null): number | null { if (through == null || denominator == null || denominator <= 0) return null; return Math.min(1, through / denominator); } function siding( label: string, count: number, stage: StageId, hint?: string, ): Siding[] { return count > 0 ? [{ label, count, stage, ...(hint ? { hint } : {}) }] : []; } // THE PIPELINES DRAWN UNDER EACH STATION, off the registry. // // The band for a channel is the same fold over the same snapshot the corpus // rail uses, so a channel figure and a corpus figure cannot disagree about what // "downloaded" or "reachable" means. Everything else here — the label, the // column-width label, the cost basis — is read from the operation catalog // rather than restated, which is what stops this file drifting from the two // other surfaces that group the same operations. function stationOperations( snapshot: ChannelSnapshot, laneKindIds: ReadonlyArray, ): Map { const catalog = new Map(operationCatalog().map((o) => [o.id, o])); const bands = buildChannelBands(snapshot, laneKindIds); const out = new Map(); for (const band of bands) { const op = catalog.get(band.id); if (!op) continue; out.set(band.id, { id: band.id, label: op.label, shortLabel: op.shortLabel, costBasis: op.costBasis, band, }); } return out; } // The operation whose coverage the station's big numeral belongs to: the FIRST // in dependency order, which is the one every other member of the group either // consumes or runs beside. // // Explicitly NOT a sum, and not an average either. The three speaker operations // are three different populations with two different cost bases — one audio // pass per video against ~1 model call per transcript chunk — and any single // figure over all three is the mistake this station used to make. One member // owns the numeral; the rest state themselves, separately, in the foot. function leadOf( ops: ReadonlyArray, ): StationOperation | null { return ops[0] ?? null; } export function computeChannelFlow( input: ComputeChannelFlowInput, ): ChannelFlow { const { snapshot, stages, failedVideoIds, laneOperations, playlistCount, } = input; const buckets = normalizeBuckets(snapshot.buckets); const digestWarnings = snapshot.buckets?.digestWarnings ?? []; const totals = snapshot.totals ?? { videos: 0, transcribed: 0, downloaded: 0 }; const undownloadedIds = snapshot.undownloadedIds ?? []; const excluded = snapshot.excludedFromDownload; const excludedIds = excludedDownloadIdSet(snapshot); const actionableDownloadedNoTranscript = buckets.downloadedNoTranscript.filter( (id) => !excludedIds.has(id), ); // Read the digest operation through digestWorkOf — the operation registry's // entry. The old `noDigest` bucket had no cues-staleness gate and no // transcript gate, so it called deferred and blocked videos done. const digestWork = digestWorkOf(snapshot); // backfillLaneEntriesOf, never Object.values: the per-kind map now carries EVERY // catalog operation including digest (~75k videos on the live corpus), and // digest has its own station one step upstream. const laneEntries = backfillLaneEntriesOf(snapshot.backfill); // An empty kind list means the operator switched the feature off. That is not // an empty work list in the "finished" sense, and the tone rule below says so. const laneOff = laneOperations.length === 0; const laneKindIds = laneOperations.map((k) => k.id); const operationsById = stationOperations(snapshot, laneKindIds); const opsFor = (...ids: string[]): StationOperation[] => ids .map((id) => operationsById.get(id)) .filter((o): o is StationOperation => o != null); // The lane's own operations, in dependency order — the group the station is // named after, and the members its foot states one by one. const laneOps = laneOff ? [] : opsFor(...laneKindIds); const laneLead = leadOf(laneOps); const laneReachable = laneOff ? 0 : laneEntries.reduce((n, e) => n + reachableOperationWork(e), 0); const laneMissingInput = laneEntries.reduce((n, e) => n + e.missingInput, 0); // `?? 0` is load-bearing, not defensive: snapshots written before these fields // existed have neither, and .toLocaleString() on undefined throws in a render. const laneDeferred = laneEntries.reduce((n, e) => n + (e.deferred ?? 0), 0); const laneBlocked = laneEntries.reduce((n, e) => n + (e.blocked ?? 0), 0); // NO CROSS-OPERATION `present` OR `eligible` SUM LIVES HERE ANY MORE, and the // helper that made one is gone with it. It used to add diarization's coverage // to attribution-text's, which is an audio pass plus a per-chunk model call // over two different populations. The station reads its LEAD operation and // the foot states each member on its own — see leadOf. // // The four WORK counts above are still summed, and legitimately: a siding is // "how many videos have left the line for this reason", and that reason is // the same reason whichever operation reported it. // digestCountOf sums `digestEngines`, which 11 of the 65 live snapshots lack // entirely — it returns 0 for those, which would read as "nothing digested". // So it is only consulted when the map is actually present. const digestEnginesTotal = snapshot.digestEngines != null ? digestCountOf(snapshot) : null; // The listing, minus what the download filter settled. Null (no stored // playlist) stays null — "unknown", not zero. const listedAfterFilter = playlistCount == null ? null : Math.max(0, playlistCount - buckets.skippedByTitleFilter.length); const stationById: Record = { playlist: { id: "playlist", // THE DENOMINATOR EXCLUDES WHAT THE FILTER SETTLED. `totals.videos` no // longer counts settled stubs, so dividing it by the raw listing made a // filtered channel read 10 % covered forever — the 90 % being videos the // operator asked us not to fetch. The listing itself is unchanged and the // settled count is drawn as a siding under this station, so nothing is // hidden; it is just not counted as a shortfall. label: "Playlist", through: totals.videos, denominator: listedAfterFilter, coverage: ratio(totals.videos, listedAfterFilter), running: stages.playlist.running, tone: stages.playlist.tone, stage: "playlist", operations: [], }, download: { id: "download", label: "Download", through: totals.downloaded, denominator: totals.videos, coverage: ratio(totals.downloaded, totals.videos), running: stages.download.running, tone: stages.download.tone, stage: "download", operations: opsFor("download"), }, transcribe: { id: "transcribe", label: "Transcribe", through: totals.transcribed, denominator: totals.downloaded, coverage: ratio(totals.transcribed, totals.downloaded), running: stages.transcribe.running, tone: stages.transcribe.tone, stage: "transcribe", operations: opsFor("transcription"), }, digest: { id: "digest", label: "Digest", through: digestWork.present ?? digestEnginesTotal, denominator: digestWork.eligible, coverage: ratio(digestWork.present ?? digestEnginesTotal, digestWork.eligible), running: stages.digest.running, tone: stages.digest.tone, stage: "digest", operations: opsFor("digest"), }, speakers: { id: "speakers", // NAMES THE WORK, NOT THE QUEUE. "Backfill" is a scheduler key — three // operations happen to share it — and an operator cannot control, arm or // pause "a backfill". They can pause speaker work. The name is derived // from the group its members declare, so a lane that gains a kind from // another group degrades to "Derived data" instead of lying. label: operationsGroupLabel(laneKindIds), // THE NUMERAL BELONGS TO ONE OPERATION, not to a sum of three. See leadOf. // Read off the band rather than recomputed: the band IS presentBackfill- // Work over this snapshot, and a second derivation is a second thing that // can disagree with the strip on /channels. through: laneLead?.band.present ?? null, denominator: laneLead?.band.eligible ?? null, coverage: laneLead ? ratio(laneLead.band.present, laneLead.band.eligible) : null, running: stages.speakers.running, // A station whose lane is DISABLED is neutral — never "ok" and never // amber. An empty work list because a feature is off is not the same as // being finished, and colouring it green claims a thing nobody checked. tone: laneOff ? "neutral" : stages.speakers.tone, stage: "speakers", operations: laneOps, }, }; const order: FlowStationId[] = [ "playlist", "download", "transcribe", "digest", "speakers", ]; const stations = order.map((id) => stationById[id]); // Sidings belong to the DESTINATION station's stage: "4 need cookies" hangs // under "to download", "1,631 blocked" under "to digest". Same rule for every // gap, so nothing is homeless and nothing is counted twice. const sidingsOf: Record = { // The filtered-out siding hangs HERE, not under download: it explains a gap // between the listing and `totals.videos`, which is this station's gap. The // rule for every siding is that it belongs to the destination station of // the gap it explains. playlist: [ ...siding( "filtered out", buckets.skippedByTitleFilter.length, "diagnostics", "Declined by this channel's download filter and settled — they are in the listing and will not be downloaded. Change the filter to bring them back.", ), ], download: [ ...siding( "never fetched", snapshot.missingNeverFetched?.length ?? 0, "diagnostics", "Known to the roster, never downloaded, and gone from the current listing.", ), ...siding("members-only", excluded?.membersOnly?.length ?? 0, "diagnostics"), ...siding("deleted", excluded?.deleted?.length ?? 0, "diagnostics"), ...siding("private", excluded?.private?.length ?? 0, "diagnostics"), ...siding( "skipped by filter", buckets.skippedByFilter.length, "diagnostics", "Declined as currently live or upcoming; retried on a later sync.", ), // THE GAP IS REAL AND IT IS THIS STATION'S. A chat-only video is not // settled out of the listing (it is not in skippedByTitleFilter), so it // counts as expected here — and until its chat lands it has no directory, // so it is missing from totals.videos. That is exactly the gap a siding // exists to name. ...siding( "chat only, not fetched", buckets.chatOnlyPending.length, "diagnostics", "Livestreams the filter rejected on a channel set to keep the chat. The lane fetches the live chat last — after every real download — because it costs no media.", ), ...siding( "need cookies", buckets.needsCookies.length, "download", "Browser cookies could recover these.", ), ...siding("partial downloads", buckets.partialDownloads.length, "download"), ...siding( "corrupt source", buckets.corruptSource.length, "download", "Malformed source; needs re-downloading.", ), ...siding( "corrupt full source", buckets.corruptFullSource.length, "download", "Download completed but the audio stayed malformed. File kept for inspection; re-downloading is futile.", ), ], transcribe: [ ...siding("failed", failedVideoIds.length, "transcribe"), ...siding( "untranscribable", buckets.untranscribable.length, "diagnostics", "Marked untranscribable by hand — an intentional decision, not an anomaly.", ), ...siding( "incomplete transcript", buckets.incompleteTranscript.length, "transcribe", "The transcript covers a fraction of the runtime — the audio download truncated silently.", ), ...siding( "short audio", buckets.shortAudio.length, "transcribe", "The source served a truncated stream; the short file is kept so it is not re-downloaded into a loop.", ), ], digest: [ ...siding( "waiting on a transcript", digestWork.blocked, "digest", "Nothing the digest lane can do about these — the number falls on its own as transcription runs.", ), ...siding( "deferred", digestWork.deferred, "digest", deferredHintFor(laneOperations, "digest"), ), ...siding( "digest warnings", digestWarnings.length, "digest", "The digest pass recorded something a human should look at.", ), ], speakers: laneOff ? [] : [ ...siding( "needs media re-acquired", laneMissingInput, "speakers", "The source audio is gone; re-acquiring it is an opt-in re-download.", ), ...siding( "deferred", laneDeferred, "speakers", deferredHintFor(laneOperations), ), ...siding( "waiting on an earlier backfill", laneBlocked, "speakers", dependsOnHint(laneOperations), ), ], }; // Work owned by each DESTINATION station — the shortfall carried by the gap // that leads into it. const reachableInto: Record = { playlist: 0, download: undownloadedIds.length, transcribe: actionableDownloadedNoTranscript.length, digest: digestWork.reachable, speakers: laneReachable, }; const GAP_LABEL: Record = { playlist: "", download: "to download", transcribe: "to transcribe", digest: "to digest", // A VALUE tsc keys but does not spell, so it is hand-checked: // channel-line.spec.ts asserts this string. speakers: "to speakers", }; const gaps: FlowGap[] = []; for (let i = 0; i < order.length - 1; i++) { const from = order[i]; const to = order[i + 1]; gaps.push({ from, to, reachable: reachableInto[to], label: GAP_LABEL[to], stage: stationById[to].stage, sidings: sidingsOf[to], }); } const biggest = gaps.reduce( (best, g) => (g.reachable > (best?.reachable ?? 0) ? g : best), null, ); return { stations, gaps, bottleneck: biggest ? biggest.from : null, next: pickNext(gaps, laneKindIds), }; } // The ONE primary action. Deliberately the FURTHEST UPSTREAM gap with work // rather than the biggest one: the pipeline is a line, so 1,600 videos waiting // to be digested behind 113 that were never downloaded is not 1,600 jobs you can // start — clearing the upstream gap is what makes the downstream one shrink. // The bottleneck is still reported separately; it is the thing to LOOK at, not // necessarily the thing to press. function pickNext( gaps: FlowGap[], laneKindIds: ReadonlyArray, ): ChannelFlow["next"] { const ACTIONABLE: Partial> = { download: "Download missing", transcribe: "Transcribe pending", digest: "Digest channel", // A verb and its object, derived from the group — "Run speaker work". The // button used to read "Backfill", which is a queue key with no object and // nothing an operator recognises as a thing they wanted done. speakers: `Run ${operationsActionLabel(laneKindIds)}`, }; for (const gap of gaps) { if (gap.reachable <= 0) continue; const label = ACTIONABLE[gap.stage]; if (!label) continue; return { stage: gap.stage, label, count: gap.reachable }; } // Nothing reachable — including on a channel that has never been reported, // where the counts are all zero because nothing has looked yet. That case is // NOT offered here: NoReportYet already sits at the top of the page saying so // and carrying the button, and a second control with the same accessible name // is both a duplicate affordance and, as it turns out, a locator that matches // two elements. return null; } // What a deferred video of this kind is waiting for, from the registry rather // than hardcoded here. SpeakersStage used to say "too long to diarize", which // was true only while diarization was the sole kind that could defer. function deferredHintFor( kinds: Operation[], onlyId?: string, ): string | undefined { const hints = kinds .filter((k) => (onlyId ? k.id === onlyId : true)) .map((k) => k.deferredHint) .filter((h): h is string => Boolean(h)); const unique = [...new Set(hints)]; if (unique.length === 0) return undefined; return unique.map((h) => `Waiting because they are ${h}`).join(" · "); } // Which operations' output the lane's kinds are blocked on, by label — so a // blocked count says what it is waiting FOR rather than merely that it is stuck. function dependsOnHint(kinds: Operation[]): string | undefined { const labels = [ ...new Set(kinds.flatMap((k) => (k.dependsOn ?? []).map(operationLabel))), ]; if (labels.length === 0) return undefined; return `Waiting on ${labels.join(", ")}.`; }