// ONE CONSOLE FOR A PUBLISH RUN (release 18 S4). A run is several jobs on the // `publish` queue — the index update, a build, its deploy — each with its own // log; the /sites consoles (JobLane, StreamActionLog) follow ONE stream. This // joins the run's streams in order, each under a `=== ===` // line, and settles when the last one does: the run's verdict is its first // job that did not end `done`, else `done`. // // A CANCEL IS THE RUN'S. The console's Cancel is `cancelPublishRunAction` // (publishActions.ts), which cancels every live job carrying the run's id; and // here, a job of the run that ends `cancelled` cancels the run's jobs after it, // queued or running — else the console would wait on a build still queued // behind something else, for a run its operator called off. A part that was // ALREADY QUEUED by someone else (another tab's run, the lane's pass) is never // cancelled from here, never the console's job id, and not waited on: it is // another run's. // // Nothing here chains one job to another — the order is on disk (a build // carries `indexAfter`, a deploy `builtAfter`; publishStages.ts). A deploy // whose build failed ends `failed` itself, with the sentence in its log. // // Not a "use server" module: a server action file may only export async // functions, and these are plain helpers the actions call. import { getRegistry } from "yt-dlp-transcript-common/jobs/registry"; import type { JobDoneResult, StreamActionResult } from "yt-dlp-transcript-common/jobs/streamCommand"; export type RunPart = { // "Build site jeralyzer" — the stage label and its target. label: string; jobId: string; // Absent for a stage that was already queued (its own console has it). stream?: ReadableStream; done?: Promise; // Another run's job this run found queued (enqueuePublishRun `existing`). existing?: boolean; }; /** The parts as one stream + one verdict (the first job's id is the run's). */ export function followRun( parts: RunPart[], preface: string[] = [], ): Extract { if (parts.length === 0) throw new Error("followRun: a run with no jobs"); // The console's job: the first this run enqueued (its Cancel cancels the // run by that job's run id), else — every part another run's — the first. const first = parts.find((p) => !p.existing) ?? parts[0]; let resolveDone!: (r: JobDoneResult) => void; const done = new Promise((r) => { resolveDone = r; }); const stream = new ReadableStream({ async start(c) { let closed = false; const push = (text: string) => { if (closed) return; try { c.enqueue(text); } catch { closed = true; } }; for (const line of preface) push(line.endsWith("\n") ? line : `${line}\n`); let verdict: JobDoneResult | null = null; for (const [i, part] of parts.entries()) { if (parts.length > 1) push(`${i === 0 ? "" : "\n"}=== ${part.label} (job ${part.jobId}) ===\n`); if (!part.stream) { push(`[publish] ${part.label} was already queued — its log is on /jobs/${part.jobId}\n`); continue; } const reader = part.stream.getReader(); try { while (true) { const { value, done: end } = await reader.read(); if (end) break; if (value) push(value); } } catch { /* the job's stream tore down; its log is on disk */ } const term = part.done ? await part.done : null; if (term && term.status !== "done" && verdict === null) verdict = term; if (term?.status === "cancelled") { for (const later of parts.slice(i + 1)) if (!later.existing) getRegistry().cancel(later.jobId); } } resolveDone(verdict ?? { status: "done", jobId: first.jobId }); if (!closed) { try { c.close(); } catch { /* already closed */ } } }, cancel() { // The consumer left (navigation): the jobs run on, their logs on disk. for (const p of parts) void p.stream?.cancel().catch(() => {}); }, }); return { ok: true, jobId: first.jobId, stream, done }; }