Archilyzer · Source

archilyzer

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

commit 848ce0f545872bee66d30bf0065fe75288a60fbe
parent 25db381f5bf3533f0a68d29ad9272fc1668a31e5
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue, 11 Aug 2026 20:44:49 -0400

Group a site's channels, and give each group the pipeline

The editor has never rendered the groups a site files its channels
into — not the buckets, not their authored descriptions. This is the
data half of showing them, plus running one pipeline stage over a
whole group.

- activeSlugsForKinds() generalizes the sync-only running-job guard, so
  a group fan-out can dedupe against in-flight work of its own kind.
  The registry has no dedupe; without it a second click enqueues a
  second whisper-all for the same channel.
- queueForSlugs() is the one bulk fan-out loop, shared with "Sync all".
  The lastSyncAllAt stamp stays OUTSIDE it: a group sweep is not a full
  sweep and must not claim one in the monitor widget.
- channelGroupSections buckets like MCP's list_channels does and totals
  each station through ONE per-channel reader, so a group figure can
  never disagree with the channel page. digestWorkOf not noDigest,
  laneEntriesOf not Object.values, and never a sum across stations.
- groupActions takes IDS, never a slug list, and re-reads membership at
  click time.

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

Diffstat:
Mcommon/jobs/syncJobs.ts | 24+++++++++++++++++++++++-
Meditor/app/channels/actions.ts | 51++++++++++++++++++++-------------------------------
Aeditor/app/channels/groupActions.ts | 174+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/lib/channelGroupSections.test.ts | 378+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/lib/channelGroupSections.ts | 299+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Aeditor/app/channels/lib/queueForSlugs.ts | 43+++++++++++++++++++++++++++++++++++++++++++
6 files changed, 937 insertions(+), 32 deletions(-)

diff --git a/common/jobs/syncJobs.ts b/common/jobs/syncJobs.ts @@ -16,5 +16,27 @@ export function activeSyncJobs(): JobRecord[] { } export function activeSyncSlugs(): Set<string> { - return new Set(activeSyncJobs().map((j) => j.channelSlug as string)); + return activeSlugsForKinds(["sync"]); +} + +// The same "is this channel already busy with this kind of work?" question for +// any job kind, which is what a bulk fan-out over a channel group needs: the job +// registry has NO dedupe of its own, so a second click would happily enqueue a +// second whisper-all for the same channel. Kind ids in use here: "sync", +// "download-missing", "whisper-all", "digest-channel-local", "backfill-channel". +// +// activeSyncSlugs() is expressed in terms of this rather than beside it, so the +// scheduler's dedupe and a group sweep's can never drift. +export function activeSlugsForKinds( + kinds: ReadonlyArray<string>, +): Set<string> { + const wanted = new Set(kinds); + const slugs = new Set<string>(); + for (const j of getRegistry().list()) { + if (!wanted.has(j.kind)) continue; + if (j.status !== "queued" && j.status !== "running") continue; + if (!j.channelSlug) continue; + slugs.add(j.channelSlug); + } + return slugs; } diff --git a/editor/app/channels/actions.ts b/editor/app/channels/actions.ts @@ -39,6 +39,7 @@ import { parseSiteMembershipsField, planSiteMembershipWrites, } from "./lib/siteMemberships"; +import { queueForSlugs, type QueueOutcome } from "./lib/queueForSlugs"; import { storePlaylistAction, syncAction } from "./[slug]/pipelineActions"; import { fetchPostsAction } from "./[slug]/socialActions"; import { prioritizeChannelDownloadAction } from "../auto-queue/actions"; @@ -358,10 +359,7 @@ export async function toggleChannelCleanupInclusionAction( return undefined; } -export type SyncAllResult = { - queued: string[]; - skipped: { slug: string; reason: string }[]; -}; +export type SyncAllResult = QueueOutcome; // Queue a sync for every eligible channel. Each sync decides for itself whether // it is due for a full sweep, so "Sync all" surfaces upstream deletions on @@ -372,40 +370,31 @@ export async function syncAllChannelsAction( ): Promise<SyncAllResult> { const paths = getPaths(); const channels = await listChannelConfigs(paths); + const bySlug = new Map(channels.map((c) => [c.slug, c.config])); const active = activeSyncSlugs(); - const queued: string[] = []; - const skipped: { slug: string; reason: string }[] = []; - for (const c of channels) { - if (!c.config.url) { - skipped.push({ slug: c.slug, reason: "no url" }); - continue; - } - if (c.config.excludeFromSync) { - skipped.push({ slug: c.slug, reason: "excluded from sync all" }); - continue; - } - if (active.has(c.slug)) { - skipped.push({ slug: c.slug, reason: "already running" }); - continue; - } - const result = await syncAction(c.slug, undefined, opts?.fullSweep); - if (!result.ok) { - skipped.push({ slug: c.slug, reason: result.error }); - continue; - } - queued.push(c.slug); - // Nobody will read the stream here — cancel it so the buffered chunks - // can be GC'd. The job keeps running and writes to its log file via - // runManagedFunction's onLog regardless. - void result.stream.cancel(); - } + const outcome = await queueForSlugs( + channels.map((c) => c.slug), + { + skip: (slug) => { + const config = bySlug.get(slug); + if (!config?.url) return "no url"; + if (config.excludeFromSync) return "excluded from sync all"; + if (active.has(slug)) return "already running"; + return null; + }, + run: (slug) => syncAction(slug, undefined, opts?.fullSweep), + }, + ); // Record the sweep's freshness marker for the monitor widget's last-sync // readout. Read-modify-write right before the write keeps the clobber window // vs. a concurrent scheduler tick minimal (single-user editor — acceptable). + // + // Deliberately OUTSIDE queueForSlugs: a per-GROUP sweep runs the same loop but + // is not a full sweep, and must not claim one in that readout. const state = await readSchedulerState(paths); state.lastSyncAllAt = Date.now(); await writeSchedulerState(paths, state); - return { queued, skipped }; + return outcome; } export async function deleteChannelAction( diff --git a/editor/app/channels/groupActions.ts b/editor/app/channels/groupActions.ts @@ -0,0 +1,174 @@ +"use server"; + +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"; +import { isValidGroupId } from "yt-dlp-transcript-common/lib/channelGroups"; +import { getPaths } from "yt-dlp-transcript-common/lib/paths"; +import { getSettings } from "yt-dlp-transcript-common/lib/settings"; +import { getSite, isValidSiteId } from "yt-dlp-transcript-common/lib/site"; +import { + laneOffFor, + slugsInGroup, + stationWorkFor, + type StationId, +} from "./lib/channelGroupSections"; +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"; + +// Run one pipeline stage over every channel in one of a site's groups. +// +// TAKES IDS, NEVER A CLIENT-SUPPLIED SLUG LIST. The membership is re-read from +// sites/<id>/site.json here, so a tampered call cannot queue work outside the +// group, and the fan-out stays correct if membership changed between the render +// and the click. +// +// FAN-OUT IS SAFE: common/lib/queueKeys.ts submits every non-empty queueKey with +// concurrency 1, so N transcribeMissingAction calls become N whisper-all records +// that run strictly one after another on the shared TRANSCRIPTION_QUEUE — no GPU +// thrash. Downloads serialize per platform, backfill shares BACKFILL_QUEUE. The +// `await` in the loop awaits ENQUEUE, not completion, so this returns promptly. +// Same shape as syncAllChannelsAction and the backfill sweep. + +export type GroupOpResult = { + group: { siteId: string; groupId: string }; + queued: string[]; + skipped: { slug: string; reason: string }[]; +}; + +// The job kind each station enqueues, for the "already running" dedupe. The job +// registry has no dedupe of its own. +const KIND_FOR: Record<StationId, string> = { + sync: "sync", + download: "download-missing", + transcribe: "whisper-all", + digest: "digest-channel-local", + backfill: "backfill-channel", +}; + +const RUN_FOR: Record< + StationId, + (slug: string) => Promise<StreamActionResult> +> = { + sync: (slug) => syncAction(slug, undefined), + // runPipelineAction already refuses with {ok:false, error} for a paused + // install, a per-platform 429 cooldown and low disk, so all three land in + // `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. + backfill: (slug) => backfillChannelAction(slug), +}; + +async function runGroupStation( + station: StationId, + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + const group = { siteId, groupId }; + if (!isValidSiteId(siteId)) { + return { group, queued: [], skipped: [{ slug: siteId, reason: "invalid site id" }] }; + } + if (!isValidGroupId(groupId)) { + return { group, queued: [], skipped: [{ slug: groupId, reason: "invalid group id" }] }; + } + const paths = getPaths(); + const settings = getSettings(); + // One read, outside the loop. Only throws on an invalid id, guarded above; a + // site with no config file parses to one with no channels, so an unknown id + // fans out over nothing rather than erroring. + const site = getSite(siteId, paths); + // Whole-group refusals answer once rather than repeating one identical line + // per channel. + if (laneOffFor(station, settings)) { + return { + group, + queued: [], + skipped: [{ slug: groupId, reason: "the lane is switched off in Settings" }], + }; + } + if (station === "download" && settings.downloadsPaused) { + return { + group, + queued: [], + skipped: [ + { + slug: groupId, + reason: "Downloads are paused — resume them from the dashboard.", + }, + ], + }; + } + + const members = new Set(slugsInGroup(site, groupId)); + const briefs = (await listChannelBriefs(paths)).filter((b) => + members.has(b.slug), + ); + const briefBySlug = new Map(briefs.map((b) => [b.slug, b])); + const active = activeSlugsForKinds([KIND_FOR[station]]); + + const outcome = await queueForSlugs( + briefs.map((b) => b.slug), + { + skip: (slug) => { + const brief = briefBySlug.get(slug); + if (!brief) return "channel not found"; + const work = stationWorkFor(station, brief, settings); + if (!work.eligible) return work.reason ?? "not applicable"; + if (active.has(slug)) return "already running"; + // work === null means the channel has never reported: run it rather + // than assume it is finished. Only a report that SAYS zero is a skip — + // which is what keeps a 4-channel group from becoming 18 no-op jobs. + // + // Sync is exempt: it carries no figure at all (syncAction decides for + // itself whether the channel is due), so its 0 means "no count", not + // "nothing to do", and gating on it would skip every channel. + if (station !== "sync" && work.work === 0) return "nothing to do"; + return null; + }, + run: RUN_FOR[station], + }, + ); + return { group, ...outcome }; +} + +export async function syncChannelGroupAction( + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + return runGroupStation("sync", siteId, groupId); +} + +export async function downloadChannelGroupAction( + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + return runGroupStation("download", siteId, groupId); +} + +export async function transcribeChannelGroupAction( + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + return runGroupStation("transcribe", siteId, groupId); +} + +export async function digestChannelGroupAction( + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + return runGroupStation("digest", siteId, groupId); +} + +export async function backfillChannelGroupAction( + siteId: string, + groupId: string, +): Promise<GroupOpResult> { + return runGroupStation("backfill", siteId, groupId); +} diff --git a/editor/app/channels/lib/channelGroupSections.test.ts b/editor/app/channels/lib/channelGroupSections.test.ts @@ -0,0 +1,378 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { + ChannelBrief, + ChannelStat, +} from "yt-dlp-transcript-common/controller/channels"; +import type { ChannelConfig } from "yt-dlp-transcript-common/lib/channelConfig"; +import type { ChannelSnapshot } from "yt-dlp-transcript-common/controller/channelSnapshot"; +import type { BackfillSnapshotEntry } from "yt-dlp-transcript-common/lib/backfillKinds"; +import type { SiteSettings } from "yt-dlp-transcript-common/lib/settings"; +import type { Site } from "yt-dlp-transcript-common/lib/site"; +import { normalizeBuckets } from "../[slug]/lib/stageStatus"; +import { + buildChannelGroupSections, + slugsInGroup, + type ChannelGroupSection, +} from "./channelGroupSections"; + +// Run from this directory: +// cd editor/app/channels/lib && ../../../../node_modules/.bin/tsx --test channelGroupSections.test.ts + +// Snapshot shape as WRITTEN TO DISK — same trick as channelFlow.test.ts. +function snapshotOf(patch: Partial<ChannelSnapshot> = {}): ChannelSnapshot { + return { + generatedAt: "2026-08-01T00:00:00.000Z", + totals: { videos: 100, transcribed: 40, downloaded: 60 }, + buckets: normalizeBuckets(undefined), + undownloadedIds: [], + ...patch, + }; +} + +function entry(patch: Partial<BackfillSnapshotEntry>): BackfillSnapshotEntry { + return { + ids: [], + missing: 0, + stale: 0, + missingInput: 0, + ...patch, + } as BackfillSnapshotEntry; +} + +function channel( + slug: string, + config: Partial<ChannelConfig> = {}, + snapshot: ChannelSnapshot | null = snapshotOf(), +): { stat: ChannelStat; brief: ChannelBrief } { + const merged: ChannelConfig = { + handling: "transcribe", + url: `https://example.test/${slug}`, + ...config, + }; + return { + stat: { + slug, + config: merged, + playlistCount: 0, + videoCount: 0, + transcriptCount: 0, + downloadCount: 0, + }, + brief: { slug, config: merged, snapshot }, + }; +} + +function siteOf(patch: Partial<Site> = {}): Site { + return { + siteId: "alpha", + siteTitle: "Alpha", + siteDescription: "", + headerTitle: "Alpha", + homeTagline: "", + groups: [], + defaultGroupId: "default", + channels: [], + ...patch, + }; +} + +// laneBackfillKinds / allBackfillKinds only read these two branches, so a +// partial cast exercises the real registry predicates. +function settingsOf(opts: { + diarization?: boolean; + attribution?: boolean; +}): SiteSettings { + return { + diarization: { + enabled: opts.diarization === true, + segModel: "seg.onnx", + embModel: "emb.onnx", + }, + attribution: { + enabled: opts.attribution === true, + textOnlyEnabled: opts.attribution === true, + diarizedEnabled: opts.attribution === true, + }, + } as unknown as SiteSettings; +} + +const LANE_ON = settingsOf({ diarization: true }); +const LANE_OFF = settingsOf({}); + +function build( + site: Site, + channels: { stat: ChannelStat; brief: ChannelBrief }[], + settings: SiteSettings = LANE_ON, +): ChannelGroupSection[] { + return buildChannelGroupSections( + site, + channels.map((c) => c.stat), + channels.map((c) => c.brief), + settings, + ); +} + +const GROUPS = [ + { id: "extended", name: "Extended Universe", selectedByDefault: false, order: 2 }, + { id: "archives", name: "Archives", selectedByDefault: true, order: 1 }, +]; + +test("groups render in sortGroups order, and empty groups are omitted", () => { + const a = channel("arch-a"); + const b = channel("ext-b"); + const sections = build( + siteOf({ + groups: [ + ...GROUPS, + { id: "guests", name: "Guest Appearances", selectedByDefault: true, order: 3 }, + ], + defaultGroupId: "archives", + channels: [ + { slug: "arch-a", groupId: "archives" }, + { slug: "ext-b", groupId: "extended" }, + ], + }), + [a, b], + ); + assert.deepEqual( + sections.map((s) => s.group.id), + ["archives", "extended"], + ); + assert.deepEqual(sections[0].channels.map((c) => c.slug), ["arch-a"]); +}); + +test("an unknown groupId folds onto the site's default group", () => { + const sections = build( + siteOf({ + groups: GROUPS, + defaultGroupId: "archives", + channels: [ + { slug: "stray", groupId: "no-such-group" }, + { slug: "ext-b", groupId: "extended" }, + ], + }), + [channel("stray"), channel("ext-b")], + ); + assert.deepEqual( + sections.map((s) => s.group.id), + ["archives", "extended"], + ); + assert.deepEqual(sections[0].channels.map((c) => c.slug), ["stray"]); +}); + +test("a site with no configured groups gets the one fallback section", () => { + const sections = build( + siteOf({ channels: [{ slug: "only" }] }), + [channel("only")], + ); + assert.equal(sections.length, 1); + assert.equal(sections[0].group.id, "default"); + assert.equal(sections[0].group.name, "All channels"); +}); + +test("within a group, membership order wins and absent order sorts last", () => { + const sections = build( + siteOf({ + channels: [ + { slug: "no-order" }, + { slug: "second", order: 2 }, + { slug: "first", order: 1 }, + ], + }), + [channel("no-order"), channel("second"), channel("first")], + ); + assert.deepEqual(sections[0].channels.map((c) => c.slug), [ + "first", + "second", + "no-order", + ]); +}); + +test("a channel with no snapshot lands in `unknown`, so the total is a floor", () => { + const sections = build( + siteOf({ channels: [{ slug: "reported" }, { slug: "silent" }] }), + [ + channel("reported", {}, snapshotOf({ undownloadedIds: ["v1", "v2"] })), + channel("silent", {}, null), + ], + ); + const { download } = sections[0]; + assert.deepEqual(download.eligible.sort(), ["reported", "silent"]); + assert.equal(download.total, 2); + assert.deepEqual(download.unknown, ["silent"]); +}); + +test("download excludes members-only/deleted/private ids", () => { + const sections = build( + siteOf({ channels: [{ slug: "c" }] }), + [ + channel( + "c", + {}, + snapshotOf({ + undownloadedIds: ["keep", "gone"], + excludedFromDownload: { + membersOnly: [], + deleted: ["gone"], + private: [], + }, + }), + ), + ], + ); + assert.equal(sections[0].download.total, 1); +}); + +test("a youtube-handling channel is not transcribe-eligible but is download-eligible", () => { + const sections = build( + siteOf({ channels: [{ slug: "yt" }, { slug: "whisper" }] }), + [ + channel("yt", { handling: "youtube" }, snapshotOf({ + undownloadedIds: ["a"], + buckets: normalizeBuckets({ downloadedNoTranscript: ["x", "y"] }), + })), + channel("whisper", { handling: "transcribe" }, snapshotOf({ + undownloadedIds: ["b"], + buckets: normalizeBuckets({ downloadedNoTranscript: ["z"] }), + })), + ], + ); + const { download, transcribe } = sections[0]; + assert.deepEqual(download.eligible.sort(), ["whisper", "yt"]); + assert.equal(download.total, 2); + assert.deepEqual(transcribe.eligible, ["whisper"]); + // The youtube channel's two awaiting-whisper videos are NOT in the figure — + // whisper will never run on them. + assert.equal(transcribe.total, 1); +}); + +test("a social channel is eligible for sync only", () => { + const sections = build( + siteOf({ channels: [{ slug: "posts" }] }), + [channel("posts", { sourceKind: "social", handling: "transcribe" })], + ); + const s = sections[0]; + assert.deepEqual(s.sync.eligible, ["posts"]); + assert.deepEqual(s.download.eligible, []); + assert.deepEqual(s.transcribe.eligible, []); + assert.deepEqual(s.digest.eligible, []); + assert.deepEqual(s.backfill.eligible, []); +}); + +test("sync skips a channel with no url or excluded from sync", () => { + const sections = build( + siteOf({ + channels: [{ slug: "ok" }, { slug: "nourl" }, { slug: "excluded" }], + }), + [ + channel("ok"), + channel("nourl", { url: undefined }), + channel("excluded", { excludeFromSync: true }), + ], + ); + assert.deepEqual(sections[0].sync.eligible, ["ok"]); + // Download needs a url too. + assert.deepEqual(sections[0].download.eligible.sort(), ["excluded", "ok"]); +}); + +test("backfill reports `off`, not 0, when no lane kind is enabled", () => { + const withWork = channel( + "c", + {}, + snapshotOf({ + backfill: { diarization: entry({ missing: 7 }) }, + }), + ); + const on = build(siteOf({ channels: [{ slug: "c" }] }), [withWork], LANE_ON); + assert.equal(on[0].backfill.laneOff, undefined); + assert.equal(on[0].backfill.total, 7); + + const off = build(siteOf({ channels: [{ slug: "c" }] }), [withWork], LANE_OFF); + // The snapshot still records 7 — an off lane is not a finished one, and the + // UI must be able to tell the difference. + assert.equal(off[0].backfill.laneOff, true); + assert.equal(off[0].backfill.total, 0); + assert.deepEqual(off[0].backfill.eligible, []); +}); + +test("digest and backfill totals are disjoint", () => { + const sections = build( + siteOf({ channels: [{ slug: "c" }] }), + [ + channel( + "c", + {}, + snapshotOf({ + backfill: { + // The per-kind map carries EVERY catalog operation, digest included. + digest: entry({ missing: 1200, stale: 4 }), + diarization: entry({ missing: 80, partial: 8 }), + }, + }), + ), + ], + ); + assert.equal(sections[0].digest.total, 1204); + assert.equal(sections[0].backfill.total, 88); +}); + +test("digest counts reachable work, never the noDigest bucket", () => { + const sections = build( + siteOf({ channels: [{ slug: "stale-report" }] }), + [ + channel( + "stale-report", + {}, + // A snapshot from before the registry: noDigest is populated, and there + // is no per-kind entry. digestWorkOf falls back to the bucket rather + // than reading the channel as fully digested. + snapshotOf({ buckets: normalizeBuckets({ noDigest: ["a", "b", "c"] }) }), + ), + ], + ); + assert.equal(sections[0].digest.total, 3); +}); + +test("digest skips a channel whose report says it has no transcripts", () => { + const sections = build( + siteOf({ channels: [{ slug: "none" }, { slug: "some" }, { slug: "silent" }] }), + [ + channel("none", {}, snapshotOf({ + totals: { videos: 10, transcribed: 0, downloaded: 10 }, + })), + channel("some", {}, snapshotOf({ + totals: { videos: 10, transcribed: 5, downloaded: 10 }, + backfill: { digest: entry({ missing: 5 }) }, + })), + // No report at all cannot claim "no transcripts" — it stays eligible. + channel("silent", {}, null), + ], + ); + const { digest } = sections[0]; + assert.deepEqual(digest.eligible.sort(), ["silent", "some"]); + assert.equal(digest.total, 5); + assert.deepEqual(digest.unknown, ["silent"]); +}); + +test("slugsInGroup buckets the same way the sections do", () => { + const site = siteOf({ + groups: GROUPS, + defaultGroupId: "archives", + channels: [ + { slug: "a", groupId: "archives" }, + { slug: "b", groupId: "extended" }, + { slug: "c", groupId: "no-such-group" }, + { slug: "d" }, + ], + }); + assert.deepEqual(slugsInGroup(site, "archives"), ["a", "c", "d"]); + assert.deepEqual(slugsInGroup(site, "extended"), ["b"]); + assert.deepEqual(slugsInGroup(site, "guests"), []); +}); + +test("slugsInGroup uses the fallback id when the site configures no groups", () => { + const site = siteOf({ channels: [{ slug: "a" }, { slug: "b" }] }); + assert.deepEqual(slugsInGroup(site, "default"), ["a", "b"]); + assert.deepEqual(slugsInGroup(site, "archives"), []); +}); diff --git a/editor/app/channels/lib/channelGroupSections.ts b/editor/app/channels/lib/channelGroupSections.ts @@ -0,0 +1,299 @@ +import type { ChannelBrief, ChannelStat } from "yt-dlp-transcript-common/controller/channels"; +import { isSocialChannel } from "yt-dlp-transcript-common/lib/channelConfig"; +import { + FALLBACK_GROUP, + resolveChannelGroupId, + sortGroups, + type ChannelGroup, +} from "yt-dlp-transcript-common/lib/channelGroups"; +import { + digestWorkOf, + excludedDownloadIdSet, +} from "yt-dlp-transcript-common/controller/channelSnapshot"; +import { + DIGEST_KIND_ID, + allBackfillKinds, + laneBackfillKinds, + laneEntriesOf, + reachableBackfillWork, +} from "yt-dlp-transcript-common/lib/backfillKinds"; +import type { SiteSettings } from "yt-dlp-transcript-common/lib/settings"; +import type { Site } from "yt-dlp-transcript-common/lib/site"; +import { normalizeBuckets } from "../[slug]/lib/stageStatus"; + +// Groups a site's channels into the sections /channels renders, and totals each +// section's pipeline work off the SAME snapshot readers the channel page's +// transit line uses — so a group figure can never disagree with the per-channel +// one. +// +// No "use server", no node:fs: imported by both the page (server) and the group +// actions, and unit-tested with plain tsx --test. The actions reuse +// stationWorkFor below, which is what keeps a button from acting on a different +// set than the number printed on it. + +export type StationId = "sync" | "download" | "transcribe" | "digest" | "backfill"; + +export const STATION_IDS: StationId[] = [ + "sync", + "download", + "transcribe", + "digest", + "backfill", +]; + +export type StationWork = { + // Slugs this operation applies to at all. A channel whose report says it has + // nothing to do IS in here — it is skipped at click time, where the read is + // fresh, and dropping it here would make the "+" floor unreadable. + eligible: string[]; + // The sum of ONE reader across `eligible`. + // + // NEVER summed with another station's total. Five separate figures on five + // separate controls, for the reason channelSnapshot.ts documents at length: a + // single "remaining" number puts every channel permanently at the top of + // every list. + total: number; + // Eligible channels with no snapshot at all. `total` is a FLOOR while this is + // non-empty, which the UI marks with a trailing "+". + unknown: string[]; + // digest/backfill only: the lane is switched off in settings. A snapshot's + // counts outlive the feature being switched off, so "off" is rendered instead + // of 0 (which reads as finished) or — (which reads as unknown). + laneOff?: boolean; +}; + +export type ChannelGroupSection = { + group: ChannelGroup; + // Membership order asc (absent sorts last), then slug. + channels: ChannelStat[]; + sync: StationWork; + download: StationWork; + transcribe: StationWork; + digest: StationWork; + backfill: StationWork; +}; + +// Whether a whole station is unrunnable because its lane is switched off. +// +// The digest kind declares `enabled: () => true` — it has no master switch, and +// its pause is honoured at DISPATCH precisely so a paused lane still reports +// what is outstanding. So digest is never "off" today; the question is asked of +// the registry rather than hardcoded so that if digest ever gains a real switch, +// this figure stops lying on its own. +export function laneOffFor( + station: StationId, + settings: SiteSettings, +): boolean { + if (station === "digest") { + return !allBackfillKinds(settings).some((k) => k.id === DIGEST_KIND_ID); + } + if (station === "backfill") return laneBackfillKinds(settings).length === 0; + return false; +} + +export type StationChannelWork = { + // Whether the operation applies to this channel at all. + eligible: boolean; + // How much work its report says there is, or null when it cannot say (no + // snapshot). Null is NOT zero: the figures above it are floors, and a click + // must not skip a channel that never reported. + work: number | null; + // Why the channel is not eligible, for the skip readout. + reason?: string; +}; + +// The single per-channel reader every group figure and every group button goes +// through. One derivation, so the label and the fan-out can never disagree. +export function stationWorkFor( + station: StationId, + brief: Pick<ChannelBrief, "config" | "snapshot">, + settings: SiteSettings, +): StationChannelWork { + const { config, snapshot } = brief; + if (laneOffFor(station, settings)) { + return { eligible: false, work: 0, reason: "the lane is switched off" }; + } + const social = isSocialChannel(config); + + if (station === "sync") { + // The same predicate "Sync every channel" applies. No figure: syncAction + // decides per channel whether it is due, so there is no count to promise. + if (!config.url) return { eligible: false, work: 0, reason: "no url" }; + if (config.excludeFromSync) { + return { eligible: false, work: 0, reason: "excluded from sync" }; + } + return { eligible: true, work: 0 }; + } + + if (social) { + return { eligible: false, work: 0, reason: "social account" }; + } + + if (station === "download") { + if (!config.url) return { eligible: false, work: 0, reason: "no url" }; + if (!snapshot) return { eligible: true, work: null }; + const excluded = excludedDownloadIdSet(snapshot); + return { + eligible: true, + work: (snapshot.undownloadedIds ?? []).filter((id) => !excluded.has(id)) + .length, + }; + } + + if (station === "transcribe") { + // A `youtube`-handling channel never runs whisper, so counting it would + // inflate the figure on a button that would skip it anyway. + if (config.handling !== "transcribe") { + return { eligible: false, work: 0, reason: "not set to transcribe" }; + } + if (!snapshot) return { eligible: true, work: null }; + const excluded = excludedDownloadIdSet(snapshot); + return { + eligible: true, + work: normalizeBuckets(snapshot.buckets).downloadedNoTranscript.filter( + (id) => !excluded.has(id), + ).length, + }; + } + + if (station === "digest") { + // Nothing to digest without transcripts — but a channel that has never + // reported cannot claim it has none, so it stays eligible and unknown. + if (!snapshot) return { eligible: true, work: null }; + if (snapshot.totals.transcribed <= 0) { + return { eligible: false, work: 0, reason: "no transcripts yet" }; + } + // digestWorkOf, NEVER buckets.noDigest: the bucket has no cues-staleness + // gate and no transcript gate, so a channel not re-reported since the + // operation registry landed reads as fully digested. + return { eligible: true, work: digestWorkOf(snapshot).reachable }; + } + + // Backfill. laneEntriesOf, NEVER Object.values: the per-kind map carries + // every catalog operation including digest, which has its own station right + // beside this one. The lane filter is what keeps the two figures disjoint. + if (!snapshot) return { eligible: true, work: null }; + return { + eligible: true, + work: laneEntriesOf(snapshot.backfill).reduce( + (n, e) => n + reachableBackfillWork(e), + 0, + ), + }; +} + +// The groups a site's channels actually fall into, in render order. Mirrors the +// bucketing MCP's list_channels does (handleListChannels), including its +// stray-bucket fallback, so the two never disagree about where a channel lives. +function bucketBySection( + site: Site, + stats: ReadonlyArray<ChannelStat>, +): { group: ChannelGroup; channels: ChannelStat[] }[] { + const hasGroups = site.groups.length > 0; + const order = new Map<string, number>(); + const groupOf = new Map<string, string>(); + for (const m of site.channels) { + if (typeof m.order === "number") order.set(m.slug, m.order); + groupOf.set(m.slug, resolveGroupIdFor(site, m.groupId, hasGroups)); + } + + const byGroup = new Map<string, ChannelStat[]>(); + for (const c of stats) { + // A stat with no membership row cannot happen for a site-scoped list, but + // fold it onto the default rather than dropping the channel. + const gid = + groupOf.get(c.slug) ?? resolveGroupIdFor(site, undefined, hasGroups); + const bucket = byGroup.get(gid) ?? []; + bucket.push(c); + byGroup.set(gid, bucket); + } + + const ordered = hasGroups ? sortGroups(site.groups) : [FALLBACK_GROUP]; + const knownIds = new Set(ordered.map((g) => g.id)); + const sections = ordered + .filter((g) => (byGroup.get(g.id)?.length ?? 0) > 0) + .map((group) => ({ group, channels: byGroup.get(group.id) ?? [] })); + // Defensive: an id resolveChannelGroupId could not fold onto a real group + // (e.g. a defaultGroupId naming a group that is no longer configured) gets a + // synthetic section rather than having its channels silently disappear. + for (const [gid, channels] of byGroup) { + if (knownIds.has(gid)) continue; + sections.push({ group: { ...FALLBACK_GROUP, id: gid }, channels }); + } + + for (const s of sections) { + s.channels.sort((a, b) => { + const ao = order.get(a.slug) ?? Number.POSITIVE_INFINITY; + const bo = order.get(b.slug) ?? Number.POSITIVE_INFINITY; + if (ao !== bo) return ao - bo; + return a.slug.localeCompare(b.slug); + }); + } + return sections; +} + +function resolveGroupIdFor( + site: Site, + groupId: string | undefined, + hasGroups: boolean, +): string { + return hasGroups + ? resolveChannelGroupId(groupId, site.groups, site.defaultGroupId) + : FALLBACK_GROUP.id; +} + +export function buildChannelGroupSections( + site: Site, + stats: ReadonlyArray<ChannelStat>, + briefs: ReadonlyArray<ChannelBrief>, + settings: SiteSettings, +): ChannelGroupSection[] { + const briefBySlug = new Map(briefs.map((b) => [b.slug, b])); + + return bucketBySection(site, stats).map(({ group, channels }) => { + const members = channels + .map((c) => briefBySlug.get(c.slug)) + .filter((b): b is ChannelBrief => !!b); + + const station = (id: StationId): StationWork => { + const work: StationWork = { eligible: [], total: 0, unknown: [] }; + if (laneOffFor(id, settings)) { + work.laneOff = true; + return work; + } + for (const b of members) { + const w = stationWorkFor(id, b, settings); + if (!w.eligible) continue; + work.eligible.push(b.slug); + if (w.work === null) work.unknown.push(b.slug); + else work.total += w.work; + } + return work; + }; + + return { + group, + channels, + sync: station("sync"), + download: station("download"), + transcribe: station("transcribe"), + digest: station("digest"), + backfill: station("backfill"), + }; + }); +} + +// The server-action half of the same bucketing: which slugs are in this group. +// Takes no snapshots, so a group button can never act on a different set than +// the rows it sits above — and re-reading it at click time keeps it correct if +// membership changed since the page rendered. +export function slugsInGroup(site: Site, groupId: string): string[] { + const hasGroups = site.groups.length > 0; + const out: string[] = []; + for (const m of site.channels) { + if (resolveGroupIdFor(site, m.groupId, hasGroups) === groupId) { + out.push(m.slug); + } + } + return out; +} diff --git a/editor/app/channels/lib/queueForSlugs.ts b/editor/app/channels/lib/queueForSlugs.ts @@ -0,0 +1,43 @@ +import type { StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; + +// The one bulk fan-out loop, shared by "Sync every channel" and the per-group +// stage buttons on /channels. +// +// Deliberately NOT in ../actions.ts: that file carries "use server", which means +// every non-type export in it must be an async server action — a shared helper +// taking callbacks cannot live there. + +export type QueueOutcome = { + queued: string[]; + skipped: { slug: string; reason: string }[]; +}; + +export async function queueForSlugs( + slugs: ReadonlyArray<string>, + opts: { + // Return a reason string to skip this slug, or null to run it. + skip?: (slug: string) => string | null | Promise<string | null>; + run: (slug: string) => Promise<StreamActionResult>; + }, +): Promise<QueueOutcome> { + const queued: string[] = []; + const skipped: { slug: string; reason: string }[] = []; + for (const slug of slugs) { + const reason = opts.skip ? await opts.skip(slug) : null; + if (reason) { + skipped.push({ slug, reason }); + continue; + } + const result = await opts.run(slug); + if (!result.ok) { + skipped.push({ slug, reason: result.error }); + continue; + } + queued.push(slug); + // Nobody will read the stream here — cancel it so the buffered chunks + // can be GC'd. The job keeps running and writes to its log file via + // runManagedFunction's onLog regardless. + void result.stream.cancel(); + } + return { queued, skipped }; +}