import path from "node:path"; import { readFile, unlink } from "node:fs/promises"; import { writeJsonAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; export const SHARD_OPS = [ "download-missing", "transcribe-missing", "availability", ] as const; export type ShardOp = (typeof SHARD_OPS)[number]; export type ShardConfig = { totalShards: number; shardIndex: number; items: string[]; createdAt: string; }; export function shardFile(paths: Paths, slug: string, op: ShardOp): string { return path.join(paths.channelsDir, slug, `shard-${op}.json`); } export async function loadShardConfig( paths: Paths, slug: string, op: ShardOp, ): Promise { try { const raw = await readFile(shardFile(paths, slug, op), "utf8"); const parsed = JSON.parse(raw) as Partial; if ( typeof parsed?.totalShards === "number" && parsed.totalShards >= 1 && typeof parsed?.shardIndex === "number" && parsed.shardIndex >= 0 && parsed.shardIndex < parsed.totalShards && Array.isArray(parsed.items) && parsed.items.every((x) => typeof x === "string") && typeof parsed.createdAt === "string" ) { return parsed as ShardConfig; } return null; } catch { return null; } } export async function saveShardConfig( paths: Paths, slug: string, op: ShardOp, cfg: ShardConfig, ): Promise { await writeJsonAtomic(shardFile(paths, slug, op), cfg); } export async function clearShardConfig( paths: Paths, slug: string, op: ShardOp, ): Promise { try { await unlink(shardFile(paths, slug, op)); return true; } catch { return false; } } // Sort lexicographically first so the slice is stable across machines that // happen to readdir in different orders. Same (items, total, index) on every // box → same slice. export function computeShardSlice( items: string[], totalShards: number, shardIndex: number, ): string[] { if (totalShards < 1 || shardIndex < 0 || shardIndex >= totalShards) { throw new Error( `Invalid shard config: total=${totalShards}, index=${shardIndex}`, ); } const sorted = [...items].sort(); return sorted.filter((_, i) => i % totalShards === shardIndex); } export type ResolveShardOpts = { paths: Paths; slug: string; op: ShardOp; fullItems: string[]; // When undefined: no inputs provided this run. If a saved slice exists, // operate on it; otherwise fall through to fullItems. totalShards?: number; shardIndex?: number; onLog?: (msg: string) => void; }; export type ResolveShardResult = { items: string[]; source: "saved" | "computed" | "full"; config: ShardConfig | null; }; // Decide what to operate on for a sharded run: // 1. If totalShards+shardIndex provided AND a saved slice exists with the // same (total, index): return saved (resume). // 2. Else if totalShards+shardIndex provided: compute, persist, return. // 3. Else if a saved slice exists: return saved (resume without re-supplying inputs). // 4. Else: return full list. export async function resolveShardItems({ paths, slug, op, fullItems, totalShards, shardIndex, onLog, }: ResolveShardOpts): Promise { const log = onLog ?? ((m: string) => console.log(m)); const saved = await loadShardConfig(paths, slug, op); const haveInputs = totalShards !== undefined && shardIndex !== undefined && totalShards >= 1 && shardIndex >= 0 && shardIndex < totalShards; if ( haveInputs && saved && saved.totalShards === totalShards && saved.shardIndex === shardIndex ) { log( `Shard ${op}: using saved slice ${saved.shardIndex + 1}/${saved.totalShards} (${saved.items.length} item(s), saved ${saved.createdAt}).`, ); return { items: saved.items, source: "saved", config: saved }; } if (haveInputs) { const slice = computeShardSlice(fullItems, totalShards!, shardIndex!); const cfg: ShardConfig = { totalShards: totalShards!, shardIndex: shardIndex!, items: slice, createdAt: new Date().toISOString(), }; await saveShardConfig(paths, slug, op, cfg); log( `Shard ${op}: computed slice ${shardIndex! + 1}/${totalShards} (${slice.length} of ${fullItems.length} item(s)); persisted.`, ); return { items: slice, source: "computed", config: cfg }; } if (saved) { log( `Shard ${op}: using saved slice ${saved.shardIndex + 1}/${saved.totalShards} (${saved.items.length} item(s), saved ${saved.createdAt}).`, ); return { items: saved.items, source: "saved", config: saved }; } log( `Shard ${op}: no shard inputs provided; running the full set (${fullItems.length} item(s)).`, ); return { items: fullItems, source: "full", config: null }; }