commit b1188708cc2c59dd2c22ec1ab4ba7ca18d599557
parent 2cbba25df88d5d3e0e67883bfe54d44a84d25df0
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 30 May 2026 15:24:42 -0400
fix shards
Diffstat:
7 files changed, 216 insertions(+), 18 deletions(-)
diff --git a/common/components/StreamActionLog.tsx b/common/components/StreamActionLog.tsx
@@ -1,5 +1,6 @@
"use client";
+import { useRouter } from "next/navigation";
import { useCallback, useEffect, useLayoutEffect, useRef, useState } from "react";
import type { StreamActionResult } from "../jobs/streamCommand";
@@ -42,6 +43,7 @@ export function StreamActionLog({
disabled = false,
}: Props) {
const accessibleName = label ?? buttonLabel;
+ const router = useRouter();
const [running, setRunning] = useState(false);
const [log, setLog] = useState("");
const [error, setError] = useState<string | null>(null);
@@ -109,12 +111,14 @@ export function StreamActionLog({
setPoll(null);
stickToBottomRef.current = true;
setRunning(true);
+ let started = false;
try {
const result = await trigger();
if (!result.ok) {
setError(result.error);
return;
}
+ started = true;
setJobId(result.jobId);
const reader = result.stream.getReader();
while (true) {
@@ -126,6 +130,13 @@ export function StreamActionLog({
setError((e as Error).message);
} finally {
setRunning(false);
+ // A run that actually started (success, mid-stream error, or cancel)
+ // may have mutated server-derived data: a saved shard config, channel
+ // counts, bucket lists, failed-transcription lists. Re-fetch the server
+ // components so the page reflects it without a manual reload. The
+ // streamed log and any typed inputs are client state and survive
+ // router.refresh().
+ if (started) router.refresh();
}
}
diff --git a/common/controller/whisperBatch.ts b/common/controller/whisperBatch.ts
@@ -6,7 +6,9 @@ import type { AudioFormat } from "../lib/channelConfig";
import { VTT_FILENAME, WHISPER_FILENAME } from "../lib/videoStatus";
import { transcribeOneVideo } from "./transcribeOne";
import { resolveShardItems } from "./shard";
+import { countNotYetTranscribed } from "./channels";
import type { TaskTracker } from "../jobs/taskHooks";
+import type { JobProgress } from "../jobs/registry";
const { pathExists, readdir, appendFile, readFile, ensureFile } = fs;
@@ -20,6 +22,11 @@ export type WhisperBatchOptions = {
reverse?: boolean;
shardTotal?: number;
shardIndex?: number;
+ // When sharding is active, refine the job's progress total to this slice
+ // (progressBaseline + shard items still needing a transcript) so the Active
+ // Jobs bar measures the shard, not the whole channel.
+ setProgress?: (snap: JobProgress) => void;
+ progressBaseline?: number;
// When set, only consider these video IDs (intersected with what's on
// disk). Used by bucket-scoped actions like "Transcribe downloaded audio"
// where the caller has already determined the exact set.
@@ -51,6 +58,8 @@ export async function runWhisperBatch({
reverse = false,
shardTotal,
shardIndex,
+ setProgress,
+ progressBaseline,
ids,
onLog,
signal,
@@ -82,7 +91,7 @@ export async function runWhisperBatch({
return ids.filter((id) => onDisk.has(id));
})()
: allDirs;
- const { items: shardItems } = await resolveShardItems({
+ const shardResult = await resolveShardItems({
paths,
slug: channelSlug,
op: "transcribe-missing",
@@ -91,6 +100,27 @@ export async function runWhisperBatch({
shardIndex: shardIndex,
onLog: log,
});
+ const shardItems = shardResult.items;
+ // With a shard active, refine the job's progress total to this slice:
+ // baseline (channel-wide transcriptCount at run start) + the shard's items
+ // still needing a transcript. `current` is read channel-wide on the Active
+ // Jobs screen, so a baseline-relative, subset-sized range still reaches 100%
+ // when the shard finishes.
+ if (
+ shardResult.source !== "full" &&
+ setProgress &&
+ progressBaseline !== undefined
+ ) {
+ const remaining = await countNotYetTranscribed(paths, channelSlug, shardItems);
+ setProgress({
+ metric: "transcripts",
+ initial: progressBaseline,
+ target: progressBaseline + remaining,
+ });
+ log(
+ `Shard transcribe-missing: progress scoped to ${remaining} of ${shardItems.length} shard item(s) needing a transcript.`,
+ );
+ }
const videoDirs = [...shardItems];
if (reverse) videoDirs.reverse();
diff --git a/common/ytdlp/runYtdlp.ts b/common/ytdlp/runYtdlp.ts
@@ -22,6 +22,7 @@ import { backfillAvailabilityFromMetadata } from "../controller/backfillAvailabi
import { resolveShardItems } from "../controller/shard";
import { downloadOneManaged } from "./downloadOneManaged";
import type { TaskTracker } from "../jobs/taskHooks";
+import type { JobProgress } from "../jobs/registry";
export type YtdlpMode =
| "store-playlist"
@@ -56,6 +57,11 @@ export type RunYtdlpOpts = {
// shard. Saved/resumed via channels/<slug>/shard-download-missing.json.
shardTotal?: number;
shardIndex?: number;
+ // download-missing only: when sharding is active, the run refines the job's
+ // progress target to its own slice (progressBaseline + shard items left to
+ // fetch) so the Active Jobs bar measures the shard, not the whole channel.
+ setProgress?: (snap: JobProgress) => void;
+ progressBaseline?: number;
// download-one-audio only: full webpage URL of the single video to fetch.
singleVideoUrl?: string;
// download-one-audio only: override channelConfig.audioFormat for this run.
@@ -391,22 +397,49 @@ async function downloadPlaylistManaged(
tofetch.length = 0;
tofetch.push(...filteredTofetch);
- // Sharding is only meaningful for the missing-files mode (matches prior
- // behavior where `download-missing` accepted shardTotal / shardIndex).
- const items =
- prefilter === "destination-exists"
- ? (
- await resolveShardItems({
- paths: opts.paths,
- slug: opts.channelSlug,
- op: "download-missing",
- fullItems: tofetch,
- totalShards: opts.shardTotal,
- shardIndex: opts.shardIndex,
- onLog: (m) => opts.onLog(`${m}\n`),
- })
- ).items
- : tofetch;
+ // Sharding (download-missing): slice the *missing* set into a stable 1/N
+ // partition and persist it here — after the prefilter has determined what's
+ // missing, but BEFORE any download starts. Snapshotting the missing-set
+ // slice to disk is the whole point: the missing set shrinks as videos land,
+ // so on a resume run with the same (total, index) resolveShardItems returns
+ // the saved slice instead of re-slicing the now-smaller set, and the run
+ // picks up exactly the videos this shard still owes. Only the plain
+ // download-missing path is sharded — never the archive-prefilter
+ // (download-from-playlist) path, nor the id-allowlisted retry-bucket run.
+ const shardResult =
+ prefilter === "destination-exists" && !idAllowList
+ ? await resolveShardItems({
+ paths: opts.paths,
+ slug: opts.channelSlug,
+ op: "download-missing",
+ fullItems: tofetch,
+ totalShards: opts.shardTotal,
+ shardIndex: opts.shardIndex,
+ onLog: (m) => opts.onLog(`${m}\n`),
+ })
+ : null;
+ const items = shardResult ? shardResult.items : tofetch;
+
+ // With a shard active, refine the job's progress total to this slice:
+ // baseline (channel-wide downloadCount at run start) + the shard's items
+ // still needing download. `current` is read channel-wide on the Active Jobs
+ // screen, so a baseline-relative, subset-sized range still reaches 100% when
+ // the shard finishes.
+ if (
+ shardResult &&
+ shardResult.source !== "full" &&
+ opts.setProgress &&
+ opts.progressBaseline !== undefined
+ ) {
+ opts.setProgress({
+ metric: "downloads",
+ initial: opts.progressBaseline,
+ target: opts.progressBaseline + items.length,
+ });
+ opts.onLog(
+ `Shard download-missing: progress scoped to ${items.length} shard item(s) needing download.\n`,
+ );
+ }
if (items.length === 0) {
opts.onLog("Nothing to fetch.\n");
diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
+- **Shard configs now show up immediately, and a shard job's progress bar totals only its slice.** Previously, setting a shard (N/i) and running **Download videos** / **Transcribe missing** did persist the slice to `shard-*.json`, but the saved-shard indicator never updated until you manually reloaded the whole page — so it looked like nothing saved. The indicator on the Download/Transcribe/Diagnostics controls (and channel counts, buckets, and failed-lists generally) now refresh as soon as a run finishes — including after a cancel — without a reload. The download slice is still taken over the *missing* set and snapshotted to the file, so a resume run with the same N/i reuses that exact slice instead of re-slicing the now-smaller set. A sharded job's progress bar on `/jobs/active` now totals only its own slice rather than the whole channel.
- **Multi-site support (major change).** One editor instance can now power several public sites (e.g. "Jeralyzer" and "Rekietalyzer") over a single shared channel pool — a channel's downloads are stored once and reused by every site that includes it. A new **Sites** area (sidebar) manages each site's `sites/<id>/site.json`: its branding (title, header, description, tagline), social links, channel-group layout, and **which channels it exposes** (with a per-site group for each member, so the same channel can sit in different groups on different sites) plus its Cloudflare Pages project. Site-level branding/groups/social have moved **out** of global Settings, which now holds only operational config plus a new **Admin title** for the editor's own shell; the per-channel "Group" field is gone (grouping is configured per site). The **Charts** tab is per-site (pick the site at the top; each site has its own dashboard at `sites/<id>/chart-templates.json`). **Build index** / **Build stats dataset** now build the shared per-channel data once and a filtered bundle per site; the Deploy page's **Build static export** and **Deploy** gained a site selector and deploy each site to its own Cloudflare project. Upgrading an existing single-site install: the Sites page shows a one-click **Migrate** button (also a `migrate-to-sites` CLI) that lifts your current branding, channel groups, per-channel group assignments, and chart dashboard into a first site and reduces `settings.json` to operational config.
- **"Audio + Whisper (skip pipeline)" no longer mistakes a live-chat sidecar for audio.** yt-dlp writes the live-chat track as `audio.live_chat.json` (it ignores the subtitle output path), so an interrupted download could leave an `audio.live_chat.json.part` behind. One-click whisper saw the `audio.` prefix, decided audio was already on disk, skipped the download, and then failed with "no audio file found". Live-chat sidecars (`.part` or completed) are now excluded everywhere audio is detected, so whisper downloads the real audio and transcribes as expected.
- **Cleanup tasks on `/actionable`.** The global Actionable view now surfaces two housekeeping sections alongside the download/transcribe backlog: **channels with cleanable transcribed audio** (videos that have a whisper transcript but still keep their audio on disk) and **channels with extra audio formats** (videos with leftover audio files beside the channel's configured format), each with a per-channel count and an inline button. Because these delete files, the buttons pop a confirm dialog before queuing. The dashboard "Needs attention" card is unchanged — it still tracks only download/transcribe work.
diff --git a/editor/app/channels/[slug]/pipelineActions.ts b/editor/app/channels/[slug]/pipelineActions.ts
@@ -65,9 +65,14 @@ async function runPipelineAction(
paths,
channelSlug: slug,
fn: async (onLog, signal, setProgress, ctx) => {
+ // Baseline (channel-wide downloadCount at run start) handed to runYtdlp so
+ // a sharded download can refine the progress total to its own slice once
+ // the slice is known.
+ let progressBaseline: number | undefined;
if (mode !== "store-playlist") {
const stat = await readChannelStat(paths, slug);
if (stat) {
+ progressBaseline = stat.downloadCount;
let target: number;
if (mode === "retry-bucket") {
const bucketIds = options?.bucketIds ?? [];
@@ -105,6 +110,8 @@ async function runPipelineAction(
shardIndex: options?.shardIndex,
bucketIds: options?.bucketIds,
handlingOverride: options?.handlingOverride,
+ setProgress,
+ progressBaseline,
});
onLog("Regenerating channel report…");
try {
diff --git a/editor/app/channels/[slug]/whisperActions.ts b/editor/app/channels/[slug]/whisperActions.ts
@@ -80,6 +80,8 @@ export async function transcribeMissingAction(
strictAudioFormat: fmt !== undefined && strictAudioFormat === true,
shardTotal,
shardIndex,
+ setProgress,
+ progressBaseline: stat?.transcriptCount,
onLog,
signal,
drainSignal: ctx.drainSignal,
diff --git a/editor/e2e/shard.spec.ts b/editor/e2e/shard.spec.ts
@@ -1,5 +1,6 @@
+import { rm, writeFile } from "node:fs/promises";
import { test, expect } from "@playwright/test";
-import { pathExists, readJson, resetData } from "./helpers";
+import { pathExists, readJson, resetData, resolvePath } from "./helpers";
const AVAILABILITY_CHANNEL = "availability-test";
@@ -131,6 +132,13 @@ test.describe("Transcribe missing sharding", () => {
await expect(log).toContainText("Shard transcribe-missing: computed slice 1/2", {
timeout: 30_000,
});
+ // Part C: the job's progress total is scoped to the shard subset, not the
+ // whole channel (3 videos sorted, index 1 of 2 → 1 item still needs a
+ // transcript).
+ await expect(log).toContainText(
+ "Shard transcribe-missing: progress scoped to 1 of 1 shard item(s)",
+ { timeout: 30_000 },
+ );
await expect(log).toContainText("succeeded", { timeout: 30_000 });
const cfg = await readJson<ShardFile>(
@@ -141,6 +149,12 @@ test.describe("Transcribe missing sharding", () => {
// 3 videos (vidA, vidB, vidC) sorted; index 1 of 2 picks 1 of them.
expect(cfg.items).toHaveLength(1);
+ // Part B: the saved-shard pill reflects the just-saved config WITHOUT a
+ // manual page reload (StreamActionLog refreshes server data on completion).
+ await expect(
+ page.getByLabel("shard transcribe-missing saved").first(),
+ ).toBeVisible({ timeout: 10_000 });
+
// Only the sharded id should have a transcript.
let transcribed = 0;
for (const id of ["vidA", "vidB", "vidC"]) {
@@ -155,3 +169,103 @@ test.describe("Transcribe missing sharding", () => {
expect(transcribed).toBe(1);
});
});
+
+test.describe("Download missing sharding", () => {
+ const channelRoot = "test-transcripts/channels/test-transcribe";
+
+ // Make all three videos missing (clear their audio) and set a 3-URL playlist.
+ // The shard slices the *missing* set: vidA/vidB/vidC sorted, shard 2/0 owns
+ // indices 0 and 2 → vidA, vidC.
+ async function setUpMissing() {
+ for (const id of ["vidA", "vidB", "vidC"]) {
+ await rm(resolvePath(`${channelRoot}/data/${id}/audio.m4a`), {
+ force: true,
+ });
+ await rm(resolvePath(`${channelRoot}/data/${id}/transcript.json`), {
+ force: true,
+ });
+ }
+ await writeFile(
+ resolvePath(`${channelRoot}/playlist`),
+ [
+ "https://www.youtube.com/watch?v=vidA",
+ "https://www.youtube.com/watch?v=vidB",
+ "https://www.youtube.com/watch?v=vidC",
+ ].join("\n") + "\n",
+ );
+ }
+
+ test("shard 2/0 persists the missing-set slice before fetching and shows it without reload", async ({
+ page,
+ }) => {
+ await resetData("one-transcribe-channel-with-audio");
+ await setUpMissing();
+ await page.goto("/channels/test-transcribe");
+
+ await page.getByLabel("shard download-missing total").first().fill("2");
+ await page.getByLabel("shard download-missing index").first().fill("0");
+
+ await page.getByRole("button", { name: "Download videos" }).click();
+ const log = page.getByLabel("Download videos output");
+ // The slice is computed + persisted before any download starts.
+ await expect(log).toContainText(
+ "Shard download-missing: computed slice 0/2",
+ { timeout: 30_000 },
+ );
+ // Part C: the job's progress total is scoped to the shard subset.
+ await expect(log).toContainText(
+ "Shard download-missing: progress scoped to 2 shard item(s)",
+ { timeout: 30_000 },
+ );
+
+ // The shard config is on disk, snapshotting this shard's slice of the
+ // missing set (vidA, vidC) so a resume reuses it rather than re-slicing.
+ const cfg = await readJson<ShardFile>(
+ `${channelRoot}/shard-download-missing.json`,
+ );
+ expect(cfg.totalShards).toBe(2);
+ expect(cfg.shardIndex).toBe(0);
+ expect(cfg.items).toHaveLength(2);
+ expect(cfg.items.some((u) => u.includes("vidA"))).toBe(true);
+ expect(cfg.items.some((u) => u.includes("vidC"))).toBe(true);
+ expect(cfg.items.some((u) => u.includes("vidB"))).toBe(false);
+
+ // Part B: saved-shard pill appears without a manual reload.
+ await expect(
+ page.getByLabel("shard download-missing saved").first(),
+ ).toBeVisible({ timeout: 10_000 });
+ });
+
+ test("re-running with the same shard inputs reuses the saved slice", async ({
+ page,
+ }) => {
+ await resetData("one-transcribe-channel-with-audio");
+ await setUpMissing();
+ await page.goto("/channels/test-transcribe");
+
+ await page.getByLabel("shard download-missing total").first().fill("2");
+ await page.getByLabel("shard download-missing index").first().fill("0");
+ await page.getByRole("button", { name: "Download videos" }).click();
+ await expect(
+ page.getByLabel("Download videos output"),
+ ).toContainText("computed slice 0/2", { timeout: 30_000 });
+
+ // Wait for the first run to fully finish — the saved-shard pill appears via
+ // the post-run refresh — then start a fresh run. (Reload to get a clean
+ // run rather than racing a second click against the still-streaming first
+ // button; the no-reload guarantee is covered by the test above.) After
+ // reload the shard inputs are pre-filled from the saved config.
+ await expect(
+ page.getByLabel("shard download-missing saved").first(),
+ ).toBeVisible({ timeout: 30_000 });
+ await page.reload();
+
+ await page.getByRole("button", { name: "Download videos" }).click();
+ // resolveShardItems reports the saved slice (resume), proving a cancelled
+ // run picks up exactly the videos this shard still owes rather than
+ // re-slicing the now-smaller missing set.
+ await expect(
+ page.getByLabel("Download videos output"),
+ ).toContainText("using saved slice 0/2", { timeout: 30_000 });
+ });
+});