Archilyzer · Source

archilyzer

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

commit 8755dcb087111810edaa4f40ea245fa47d77841a
parent 4c74fff7c1fe6725c0c50c7843c9240eae406c7d
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue,  6 Oct 2026 11:27:20 -0400

publish: readPublishStatus, enqueueStage / enqueuePublishRun, and the publish lane's runner

- publishState.ts `readPublishInputs` / `readPublishStatus` (the one status
  every surface reads): the stamps and bundle problems (S1's readNeedsInput),
  the sites with their membership and policies, the plan's config mtimes,
  the ingest signal per channel (ingest-kind job metas ended `done`, read
  through the registry and the `.jobs` sidecars, each terminal meta once;
  plus each channel's report regeneration, the only trace of a lane unit),
  the live and the newest ended publish stages, the lane's memory, the
  running commit (git memoized)
- publishStages.ts `enqueueStage` (runManagedCommand over stageCommand on
  the `publish` queue, spec {kind, slug: target, params: the request};
  a duplicate — same kind, target and destination queued or running — is
  refused info with the existing id) and `enqueuePublishRun` (a plan at
  once, one run id, preconditions on the specs); `stageRequestFromSpec`
  through the stage row's own parser
- publishRunner.ts `startPublishRunner` (auto-publish, queueKey ""), the
  loop (a check every checkEveryMinutes; a pass when there is no stamp,
  the index is stale past refreshEveryMinutes, or a policy target is left
  stale and the last pass is that old) and `runPublishPass` (one stage at a
  time, background, re-planned before each; a hold, quiet hours or the lane
  switched off stop the dispatching, never the stage; a drain finishes the
  stage; a failed index update ends the pass, a failed build drops its
  deploy); stop / drain; publishLaneState.ts holds the lane's memory

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

Diffstat:
Acommon/publish/publishLaneState.ts | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishRunner.test.ts | 352+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishRunner.ts | 333+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishStages.test.ts | 195+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishStages.ts | 196+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishState.test.ts | 165+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishState.ts | 401+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
7 files changed, 1715 insertions(+), 0 deletions(-)

diff --git a/common/publish/publishLaneState.ts b/common/publish/publishLaneState.ts @@ -0,0 +1,73 @@ +// The publish lane's LIVE state (release 18): what the runner is doing, held +// in this process's memory so the status (publishState.ts) can say it without +// importing the runner (which reads the status — a cycle). +// +// On globalThis, like the auto-queue runners' singleton, so a dev-server +// module reload does not lose the running lane. The CLI never starts the +// runner: there `known` is false — the lane is the editor's. + +import type { PublishLaneLive } from "./publishPlan"; + +export type PublishPassRecord = { + runId: string; + startedAt: number; + endedAt: number | null; + stages: { kind: string; target: string; jobId: string | null; status: string }[]; + summary: string; +}; + +export type PublishLaneMemory = { + jobId: string | null; + running: boolean; + passRunning: boolean; + lastCheckAt: number | null; + nextCheckAt: number | null; + lastPassAt: number | null; + lastPass: PublishPassRecord | null; + lastDecision: string | null; + // The stage job the pass is waiting on, if any. + currentJobId: string | null; +}; + +const KEY = "__ytt_publish_lane__"; + +function fresh(): PublishLaneMemory { + return { + jobId: null, + running: false, + passRunning: false, + lastCheckAt: null, + nextCheckAt: null, + lastPassAt: null, + lastPass: null, + lastDecision: null, + currentJobId: null, + }; +} + +/** The lane's memory in this process (created on first ask). */ +export function publishLaneMemory(): PublishLaneMemory { + const g = globalThis as unknown as Record<string, PublishLaneMemory | undefined>; + return (g[KEY] ??= fresh()); +} + +/** Forget everything (tests). */ +export function resetPublishLaneMemory(): void { + (globalThis as unknown as Record<string, PublishLaneMemory | undefined>)[KEY] = fresh(); +} + +/** The lane as the status reads it. `known`: this process runs editor lanes. */ +export function publishLaneLive(known: boolean): PublishLaneLive { + const m = publishLaneMemory(); + return { + known, + running: m.running, + jobId: m.jobId, + passRunning: m.passRunning, + lastCheckAt: m.lastCheckAt, + nextCheckAt: m.nextCheckAt, + lastPassAt: m.lastPassAt, + lastPassSummary: m.lastPass?.summary ?? null, + lastDecision: m.lastDecision, + }; +} diff --git a/common/publish/publishRunner.test.ts b/common/publish/publishRunner.test.ts @@ -0,0 +1,352 @@ +import { test, beforeEach } from "node:test"; +import assert from "node:assert/strict"; +import { defaultPublish, type PublishSettings } from "../lib/settingsSchema"; +import type { JobDoneResult } from "../jobs/streamCommand"; +import { buildPublishStatus, type PublishInputs, type PublishSiteInput } from "./publishPlan"; +import { publishLaneMemory, resetPublishLaneMemory } from "./publishLaneState"; +import { publishLoop, runPublishPass, type PublishRunnerDeps } from "./publishRunner"; +import type { EnqueueStageResult } from "./publishStages"; +import type { StageRequest } from "./stages"; +import type { BuiltStamp, IndexStamp } from "./stamps"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/publishRunner.test.ts +// +// The publish lane's runner over a fake world: the "stages" are promises this +// test settles, and each one moves the world the way the real stage would +// (the index stamp, a built stamp). No process is spawned, nothing on disk. + +const MIN = 60_000; +const T0 = 2_000_000_000_000; + +function stamp(over: Partial<IndexStamp> = {}): IndexStamp { + return { + v: 1, + stampId: "s1", + generation: 1, + scannedAt: T0 - 120 * MIN, + builtAt: T0 - 119 * MIN, + templatesAt: T0 - 119 * MIN, + commit: null, + index: { shortCircuited: false, added: 0, changed: 0, removed: 0, heldChannels: [] }, + stats: { shortCircuited: false, notIndexedYet: 0, notIndexable: 0 }, + sites: { alpha: { siteFp: null, statsFp: null, inputSig: "a1" }, beta: { siteFp: null, statsFp: null, inputSig: "b1" } }, + hubSig: "h1", + ...over, + }; +} + +function built(target: string, sig: string, at: number): BuiltStamp { + return { + v: 1, + stampId: `${target}-${at}`, + target, + kind: "site", + indexStampId: "s1", + inputSig: sig, + builtAt: at, + commit: null, + branch: "main", + runner: "local", + audience: "public", + corpusGeneratedAt: null, + files: 1, + bytes: 1, + archivesStaged: 0, + }; +} + +function siteIn(id: string, policy: PublishSiteInput["policy"]): PublishSiteInput { + return { + siteId: id, + title: id, + private: false, + listed: true, + members: [`${id}-ch`], + built: built(id, `${id[0]}1`, T0 - 118 * MIN), + deployed: null, + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy, + cloudflareProject: id, + url: null, + }; +} + +type World = { + inputs: PublishInputs; + clock: number; + settings: PublishSettings; + enqueued: StageRequest[]; + // The stage in flight: settle it to let the pass go on. + pending: { req: StageRequest; settle: (s: JobDoneResult["status"]) => void } | null; + onEnqueue?: (req: StageRequest) => void; +}; + +function world(policyAlpha: PublishSiteInput["policy"] = "build"): World { + const settings = { ...defaultPublish(), enabled: true, refreshEveryMinutes: 60 }; + return { + clock: T0, + settings, + enqueued: [], + pending: null, + inputs: { + index: { stamp: stamp(), lastIngestDoneAt: T0 - 30 * MIN, configChangedAt: null }, + ingestByChannel: { "alpha-ch": T0 - 30 * MIN }, + commit: null, + sites: [siteIn("alpha", policyAlpha), siteIn("beta", "off")], + hub: { + built: null, + deployed: null, + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off", + cloudflareProject: null, + url: null, + }, + homepage: { + built: null, + deployed: null, + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off", + cloudflareProject: null, + url: null, + mainHead: null, + }, + settings, + jobs: [], + ended: [], + lane: { + known: true, + running: true, + jobId: "lane", + passRunning: false, + lastCheckAt: null, + nextCheckAt: null, + lastPassAt: null, + lastPassSummary: null, + lastDecision: null, + }, + }, + }; +} + +// Apply what a stage that ended `done` writes. +function applyDone(w: World, req: StageRequest): void { + if (req.kind === "update-index") { + w.inputs.index.stamp = stamp({ + stampId: "s2", + scannedAt: w.clock - 1000, + builtAt: w.clock, + sites: { + alpha: { siteFp: null, statsFp: null, inputSig: "a2" }, + beta: { siteFp: null, statsFp: null, inputSig: "b2" }, + }, + }); + } else if (req.kind === "build-site") { + const s = w.inputs.sites.find((x) => x.siteId === req.target)!; + s.built = built(s.siteId, `${s.siteId[0]}2`, w.clock); + } +} + +function deps(w: World): PublishRunnerDeps { + return { + readStatus: async () => buildPublishStatus({ ...w.inputs, settings: w.settings }, w.clock), + settings: () => w.settings, + now: () => w.clock, + sleep: async (ms) => { + w.clock += ms; + }, + enqueue: async (req): Promise<EnqueueStageResult> => { + w.enqueued.push(req); + w.onEnqueue?.(req); + let settle!: (s: JobDoneResult["status"]) => void; + const done = new Promise<JobDoneResult>((resolve) => { + settle = (status) => { + w.clock += MIN; + if (status === "done") applyDone(w, req); + w.pending = null; + resolve({ status, jobId: `J${w.enqueued.length}` }); + }; + }); + w.pending = { req, settle }; + return { ok: true, jobId: `J${w.enqueued.length}`, stream: new ReadableStream<string>(), done }; + }, + }; +} + +// Settle every stage the pass dispatches with `status` as it arrives. +function autoSettle(w: World, status: (req: StageRequest) => JobDoneResult["status"] = () => "done"): void { + w.onEnqueue = (req) => queueMicrotask(() => w.pending?.settle(status(req))); +} + +const signals = () => ({ signal: new AbortController().signal, drain: new AbortController().signal }); +const kinds = (w: World) => w.enqueued.map((r) => `${r.kind} ${r.target}`); + +beforeEach(() => resetPublishLaneMemory()); + +test("a pass: the index, then — re-planned against the new stamp — the policy's build; never forced", async () => { + const w = world("build"); + autoSettle(w); + const out = await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index", "build-site alpha"]); + assert.equal(out.ended, "nothing left to do"); + for (const r of w.enqueued) { + assert.equal(r.force, undefined); + assert.equal(r.runId, out.runId); + assert.match(r.runId, /^lane-/); + } + // The build was planned AFTER the index ran: no indexAfter needed. + assert.equal(w.enqueued[1].indexAfter, undefined); + const mem = publishLaneMemory(); + assert.equal(mem.passRunning, false); + assert.equal(mem.lastPassAt, w.clock); + assert.match(mem.lastPass?.summary ?? "", /2 stages/); +}); + +test("a hold between stages stops dispatching; the stage in flight is never killed", async () => { + const w = world("production"); + w.onEnqueue = (req) => + queueMicrotask(() => { + // The operator holds the lane while the index update runs. + if (req.kind === "update-index") w.settings = { ...w.settings, held: true }; + w.pending?.settle("done"); + }); + const out = await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index"]); + assert.match(out.ended, /the lane is held: no next stage dispatched/); + assert.equal(out.ran[0].status, "done", "the index stage ran to its end"); +}); + +test("quiet hours beginning between stages stop dispatching too", async () => { + const w = world("build"); + w.onEnqueue = () => + queueMicrotask(() => { + const h = new Date(w.clock + MIN).getHours(); + w.settings = { ...w.settings, quietHours: { start: h, end: (h + 2) % 24 } }; + w.pending?.settle("done"); + }); + const out = await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index"]); + assert.match(out.ended, /quiet hours began/); +}); + +test("a drain finishes the stage in flight and dispatches nothing after it", async () => { + const w = world("build"); + const drain = new AbortController(); + w.onEnqueue = () => + queueMicrotask(() => { + drain.abort(); + // The stage is still running when the drain lands; it ends after. + setTimeout(() => w.pending?.settle("done"), 5); + }); + const out = await runPublishPass(deps(w), { signal: new AbortController().signal, drain: drain.signal }, () => {}); + assert.deepEqual(kinds(w), ["update-index _index"]); + assert.equal(out.ran[0].status, "done", "the drain waited for the stage"); + assert.equal(out.ended, "the runner was drained"); +}); + +test("a stop does not wait for the stage: it keeps running as its own job", async () => { + const w = world("build"); + const stop = new AbortController(); + w.onEnqueue = () => queueMicrotask(() => stop.abort()); + const out = await runPublishPass(deps(w), { signal: stop.signal, drain: new AbortController().signal }, () => {}); + assert.equal(out.ran[0].status, "still running"); + assert.equal(out.ended, "the runner was stopped"); + assert.ok(w.pending, "the stage was not settled (nor killed) by the stop"); +}); + +test("a failed index update ends the pass; a failed build drops its deploy", async () => { + const w = world("production"); + autoSettle(w, (r) => (r.kind === "update-index" ? "failed" : "done")); + const a = await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index"]); + assert.match(a.ended, /the index update failed/); + + const v = world("production"); + autoSettle(v, (r) => (r.kind === "build-site" ? "failed" : "done")); + const b = await runPublishPass(deps(v), signals(), () => {}); + assert.deepEqual(kinds(v), ["update-index _index", "build-site alpha"]); + assert.equal(b.ended, "every stage left waits on a build that failed"); +}); + +test("a pass with the production policy deploys after the build", async () => { + const w = world("production"); + autoSettle(w); + await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index", "build-site alpha", "deploy-site alpha"]); + assert.equal(w.enqueued[2].preview, undefined); +}); + +test("a stage someone else queued makes the pass yield", async () => { + const w = world("build"); + w.onEnqueue = () => + queueMicrotask(() => { + w.inputs.jobs = [ + { id: "X", kind: "build-site", target: "beta", status: "queued", runId: "run-x", queuedAt: w.clock }, + ]; + w.pending?.settle("done"); + }); + const out = await runPublishPass(deps(w), signals(), () => {}); + assert.deepEqual(kinds(w), ["update-index _index"]); + assert.match(out.ended, /someone else queued/); +}); + +test("the loop: no stamp → a pass at once; then the refresh interval gates the next index update", async () => { + const w = world("build"); + w.inputs.index.stamp = null; + autoSettle(w); + const drain = new AbortController(); + let checks = 0; + const d = deps(w); + const lines: string[] = []; + await publishLoop( + { + ...d, + readStatus: async () => { + checks++; + // New data arrives after every check; stop after a few checks. + w.inputs.index.lastIngestDoneAt = w.clock; + w.inputs.ingestByChannel["alpha-ch"] = w.clock; + if (checks > 40) drain.abort(); + return d.readStatus(); + }, + }, + { signal: new AbortController().signal, drain: drain.signal }, + (l) => lines.push(l), + ); + const indexRuns = w.enqueued.filter((r) => r.kind === "update-index").length; + assert.ok(indexRuns >= 2, `the index was updated again after the interval (${indexRuns})`); + // checkEvery 10 min, refresh 60 min: never more often than every 60 min. + const indexTimes = lines.filter((l) => l.includes("update-index _index →")).length; + assert.equal(indexTimes, indexRuns); + assert.ok(w.clock - T0 >= (indexRuns - 1) * 60 * MIN, "updates were at least refreshEveryMinutes apart"); + assert.match(lines.at(-1) ?? "", /lane runner drained/); +}); + +test("the loop ends when the lane is switched off, and a held lane never starts a pass", async () => { + const w = world("build"); + w.settings = { ...w.settings, held: true }; + autoSettle(w); + let reads = 0; + const d = deps(w); + await publishLoop( + { + ...d, + readStatus: async () => { + if (++reads > 3) w.settings = { ...w.settings, enabled: false }; + return d.readStatus(); + }, + }, + signals(), + () => {}, + ); + assert.deepEqual(w.enqueued, [], "held: nothing dispatched"); + assert.equal(publishLaneMemory().lastDecision, "held"); +}); diff --git a/common/publish/publishRunner.ts b/common/publish/publishRunner.ts @@ -0,0 +1,333 @@ +// THE PUBLISH LANE'S RUNNER (release 18). +// +// The ingest lanes' pattern (controller/autoRunner.ts startAutoRunner): one +// long-lived `runManagedFunction` job, kind `auto-publish`, queueKey "" (it +// runs beside everything and is serialized against nothing). It: +// +// - wakes every `settings.publish.checkEveryMinutes`; +// - skips the check when the lane is held, in its quiet hours, or while any +// publish stage is queued or running (somebody's Publish now, a Build +// button, the CLI's lock aside); +// - runs a PASS when one is due (publishPlan.ts `publishPassDecision`): no +// index stamp yet; the index is stale and its last update is at least +// `refreshEveryMinutes` old; or a policy target is left stale and the +// last pass is at least that old; +// - a pass dispatches ONE stage at a time — `enqueueStage(…, {background: +// true})`, then awaits its job's end — re-reading the status and +// re-planning (`planPublishRun`, deploys by policy, NEVER forced) before +// each next stage, so the builds after an index update see the stamp it +// wrote. A stage the pass already ran is not run again in it; a failed +// index update ends the pass; a failed build drops its deploy. +// - the gate is re-checked between stages: a hold, quiet hours or the lane +// switched off stop the DISPATCHING — the stage in flight is never killed. +// A drain lets the stage in flight finish and ends the runner (`done`). +// Stop (a hard cancel of the runner job) ends the runner at once; a stage +// already running is its own job and finishes (Cancel it on /jobs). +// +// Started by `startPublishRunnerIfEnabled` from editor/instrumentation.ts, +// beside the auto-queue runners and below the idle-boot gate: an idle boot +// leaves it off. (Not from controller/autoRunner.ts startAutoRunnersIfEnabled: +// the dispatch layer may not import the publish layer.) + +import { getRegistry } from "../jobs/registry"; +import { runManagedFunction, type JobDoneResult } from "../jobs/streamCommand"; +import { getPaths, type Paths } from "../lib/paths"; +import { getSettings, type PublishSettings } from "../lib/settings"; +import { isInQuietHours } from "../jobs/syncScheduler"; +import { + planPublishRun, + publishPassDecision, + stepKey, + type PlanStep, + type PublishStatus, +} from "./publishPlan"; +import { publishLaneMemory, type PublishPassRecord } from "./publishLaneState"; +import { enqueueStage, newPublishRunId, type EnqueueStageResult } from "./publishStages"; +import { readPublishStatus } from "./publishState"; +import { ALL_TARGET } from "./stamps"; +import type { StageRequest } from "./stages"; + +export const AUTO_PUBLISH_KIND = "auto-publish"; + +// How often the sleeping runner looks up (an abort, a drain, the lane +// switched off, a new checkEveryMinutes). +const TICK_MS = 5_000; + +export type PublishRunnerDeps = { + readStatus: () => Promise<PublishStatus>; + enqueue: (req: StageRequest, opts: { background: true }) => Promise<EnqueueStageResult>; + settings: () => PublishSettings; + now: () => number; + // Resolves after `ms`, or as soon as a signal fires. + sleep: (ms: number, signals: AbortSignal[]) => Promise<void>; +}; + +function realSleep(ms: number, signals: AbortSignal[]): Promise<void> { + return new Promise((resolve) => { + if (signals.some((s) => s.aborted)) return resolve(); + const done = () => { + clearTimeout(timer); + for (const s of signals) s.removeEventListener("abort", done); + resolve(); + }; + const timer = setTimeout(done, ms); + for (const s of signals) s.addEventListener("abort", done, { once: true }); + }); +} + +export function defaultPublishRunnerDeps(paths: Paths = getPaths()): PublishRunnerDeps { + return { + readStatus: () => readPublishStatus(paths), + enqueue: (req, opts) => enqueueStage(paths, req, opts), + settings: () => getSettings().publish, + now: () => Date.now(), + sleep: realSleep, + }; +} + +/** Why the lane must not dispatch the next stage now, or null. */ +function gateShut(s: PublishSettings, now: number): string | null { + if (!s.enabled) return "the lane was switched off"; + if (s.held) return "the lane is held"; + if (s.quietHours && isInQuietHours(now, s.quietHours.start, s.quietHours.end)) return "quiet hours began"; + return null; +} + +const stepWords = (s: Pick<StageRequest, "kind" | "target" | "preview" | "to">) => + `${s.kind} ${s.target}${s.preview ? ` (preview ${s.preview})` : s.to === "local" ? " (local)" : ""}`; + +export type PassOutcome = { + runId: string; + ran: { kind: string; target: string; jobId: string | null; status: string }[]; + // Why the pass ended. + ended: string; +}; + +/** + * ONE PASS: dispatch the plan's stages one at a time until it is empty, the + * gate shuts, the runner is drained or stopped, or the index update fails. + */ +export async function runPublishPass( + deps: PublishRunnerDeps, + signals: { signal: AbortSignal; drain: AbortSignal }, + onLog: (line: string) => void, +): Promise<PassOutcome> { + const mem = publishLaneMemory(); + const passStart = deps.now(); + const runId = newPublishRunId("lane", passStart); + const attempted = new Set<string>(); + const failedBuilds = new Set<string>(); + let indexRan = false; + const record: PublishPassRecord = { runId, startedAt: passStart, endedAt: null, stages: [], summary: "running" }; + mem.passRunning = true; + mem.lastPass = record; + const outcome: PassOutcome = { runId, ran: [], ended: "" }; + onLog(`[publish] pass ${runId} started`); + try { + for (;;) { + if (signals.signal.aborted) { + outcome.ended = "the runner was stopped"; + break; + } + if (signals.drain.aborted) { + outcome.ended = "the runner was drained"; + break; + } + const shut = gateShut(deps.settings(), deps.now()); + if (shut) { + outcome.ended = `${shut}: no next stage dispatched`; + break; + } + const status = await deps.readStatus(); + if (status.busy) { + outcome.ended = "a publish stage someone else queued is waiting: the pass yields"; + break; + } + const plan = planPublishRun(status, { + deploys: "policy", + runStart: passStart, + exclude: attempted, + skipIndex: indexRan, + }); + const step: PlanStep | undefined = plan.steps.find( + (s) => !(s.kind.startsWith("deploy") && (failedBuilds.has(s.target) || failedBuilds.has(ALL_TARGET))), + ); + if (!step) { + outcome.ended = plan.steps.length > 0 ? "every stage left waits on a build that failed" : "nothing left to do"; + break; + } + const { reason, previewUrl: _previewUrl, ...rest } = step; + const req: StageRequest = { ...rest, runId }; + attempted.add(stepKey(step)); + const res = await deps.enqueue(req, { background: true }); + if (!res.ok) { + if (res.info) { + outcome.ended = `${stepWords(req)} is already queued: the pass yields`; + break; + } + onLog(`[publish] ${stepWords(req)}: not enqueued — ${res.error}`); + outcome.ran.push({ kind: req.kind, target: req.target, jobId: null, status: "refused" }); + continue; + } + void res.stream.cancel().catch(() => {}); + onLog(`[publish] ${stepWords(req)} → job ${res.jobId} (${reason})`); + mem.currentJobId = res.jobId; + const entry = { kind: req.kind, target: req.target, jobId: res.jobId, status: "running" }; + record.stages.push(entry); + outcome.ran.push(entry); + // A stop does not wait for the stage (its own job); a drain does. + const term: JobDoneResult | null = await Promise.race([ + res.done, + new Promise<null>((resolve) => { + if (signals.signal.aborted) return resolve(null); + signals.signal.addEventListener("abort", () => resolve(null), { once: true }); + }), + ]); + mem.currentJobId = null; + if (!term) { + entry.status = "still running"; + onLog(`[publish] the runner stopped; ${stepWords(req)} (job ${res.jobId}) keeps running as its own job`); + outcome.ended = "the runner was stopped"; + break; + } + entry.status = term.status; + onLog(`[publish] ${stepWords(req)}: ${term.status}`); + if (req.kind === "update-index") { + indexRan = true; + if (term.status !== "done") { + outcome.ended = `the index update ${term.status}: nothing is built from a stale index`; + break; + } + } else if (req.kind.startsWith("build") && term.status !== "done") { + failedBuilds.add(req.target); + } + if (term.status === "cancelled") { + outcome.ended = `${stepWords(req)} was cancelled: the pass stops`; + break; + } + } + } finally { + const endedAt = deps.now(); + const n = outcome.ran.length; + record.endedAt = endedAt; + record.summary = `${n} stage${n === 1 ? "" : "s"}${n ? ` (${outcome.ran.map((r) => `${r.kind} ${r.target}: ${r.status}`).join(", ")})` : ""} — ${outcome.ended || "ended"}`; + mem.passRunning = false; + mem.lastPassAt = endedAt; + onLog(`[publish] pass ${runId} ended: ${record.summary}`); + } + return outcome; +} + +/** + * The runner's loop: a check every `checkEveryMinutes`, a pass when one is + * due. Returns when the runner is stopped, drained, or the lane is switched + * off. + */ +export async function publishLoop( + deps: PublishRunnerDeps, + signals: { signal: AbortSignal; drain: AbortSignal }, + onLog: (line: string) => void, +): Promise<void> { + const mem = publishLaneMemory(); + onLog("[publish] lane runner started"); + while (!signals.signal.aborted && !signals.drain.aborted) { + const settings = deps.settings(); + if (!settings.enabled) { + onLog("[publish] the lane was switched off; the runner ends"); + break; + } + const now = deps.now(); + mem.lastCheckAt = now; + const status = await deps.readStatus(); + const decision = publishPassDecision({ + settings, + stamp: status.index.stamp, + indexFresh: status.index.fresh, + busy: status.busy, + now, + lastPassAt: mem.lastPassAt, + policyWork: status.plan.steps.some((s) => s.kind !== "update-index"), + }); + if (decision.reason !== mem.lastDecision) onLog(`[publish] check: ${decision.run ? "pass due — " : ""}${decision.reason}`); + mem.lastDecision = decision.reason; + if (decision.run) await runPublishPass(deps, signals, onLog); + if (signals.signal.aborted || signals.drain.aborted) break; + // Sleep to the next check, looking up every tick. + const until = deps.now() + deps.settings().checkEveryMinutes * 60_000; + mem.nextCheckAt = until; + while (!signals.signal.aborted && !signals.drain.aborted) { + const left = until - deps.now(); + if (left <= 0) break; + if (!deps.settings().enabled) break; + await deps.sleep(Math.min(TICK_MS, left), [signals.signal, signals.drain]); + } + } + mem.nextCheckAt = null; + onLog( + signals.drain.aborted + ? "[publish] lane runner drained" + : signals.signal.aborted + ? "[publish] lane runner stopped" + : "[publish] lane runner ended", + ); +} + +/** Why the runner will not start, or null when it will. */ +export function startPublishRunnerBlockedReason(settings: PublishSettings = getSettings().publish): string | null { + return settings.enabled + ? null + : "The publish lane is off. Turn it on (settings.publish.enabled), then start the runner."; +} + +/** + * Start the lane's runner job if the lane is on and none is running. + * Idempotent (it asks the live registry). Returns the job id, or null. + */ +export async function startPublishRunner( + paths: Paths = getPaths(), + deps: PublishRunnerDeps = defaultPublishRunnerDeps(paths), +): Promise<string | null> { + if (startPublishRunnerBlockedReason(deps.settings())) return null; + const mem = publishLaneMemory(); + if (mem.jobId && getRegistry().get(mem.jobId)?.status === "running") return mem.jobId; + const result = await runManagedFunction({ + kind: AUTO_PUBLISH_KIND, + queueKey: "", + paths, + fn: async (onLog, signal, _setProgress, ctx) => { + mem.jobId = ctx.jobId; + mem.running = true; + try { + await publishLoop(deps, { signal, drain: ctx.drainSignal }, onLog); + } finally { + mem.running = false; + mem.passRunning = false; + mem.currentJobId = null; + } + }, + }); + if (!result.ok) return null; + mem.jobId = result.jobId; + // Nobody reads the runner's stream; its log is on disk. + void result.stream.cancel().catch(() => {}); + return result.jobId; +} + +/** Start it at boot when the lane is on (editor/instrumentation.ts). */ +export async function startPublishRunnerIfEnabled(paths: Paths = getPaths()): Promise<string | null> { + return startPublishRunner(paths); +} + +/** Stop the runner (a hard cancel of its job). A stage in flight finishes as its own job. */ +export function stopPublishRunner(): boolean { + const id = publishLaneMemory().jobId; + if (!id) return false; + return getRegistry().cancel(id); +} + +/** Drain the runner: the stage in flight finishes, no next one starts, the job ends `done`. */ +export function drainPublishRunner(): boolean { + const id = publishLaneMemory().jobId; + if (!id) return false; + return getRegistry().requestDrain(id); +} diff --git a/common/publish/publishStages.test.ts b/common/publish/publishStages.test.ts @@ -0,0 +1,195 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { mkdtemp, readFile, rm } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { getRegistry, type JobRecord } from "../jobs/registry"; +import { getJobKind, jobKindLabel } from "../jobs/jobKinds"; +import { readJobMeta } from "../jobs/jobMeta"; +import { getPaths, type Paths } from "../lib/paths"; +import { + PUBLISH_QUEUE, + enqueuePublishRun, + enqueueStage, + findQueuedStage, + newPublishRunId, + stageIdentity, + stageRequestFromSpec, + stageSpec, +} from "./publishStages"; +import { STAGES, STAGE_KINDS, type StageRequest } from "./stages"; +import type { PublishPlan } from "./publishPlan"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/publishStages.test.ts +// +// Enqueueing publish stages (release 18): the spec is the request, a duplicate +// is refused with the existing id, and a run's preconditions ride on its specs. +// One test spawns the real stage child — with a request its own parser refuses +// (exit 2), so it touches nothing — to prove the job's shape end to end. + +test("the seven stages' job kinds are jobKinds.ts's, with the stages' labels", () => { + for (const kind of STAGE_KINDS) { + const jobKind = STAGES[kind].jobKind; + assert.ok(getJobKind(jobKind), `${jobKind} is registered`); + assert.equal(jobKindLabel(jobKind), STAGES[kind].label, `${jobKind} label`); + assert.equal(getJobKind(jobKind)?.replayable, true); + } +}); + +test("a stage's spec is its request, and reads back to the same request", () => { + const reqs: StageRequest[] = [ + { kind: "update-index", target: "_index", runId: "run-1" }, + { kind: "build-site", target: "jeralyzer", runId: "run-1", indexAfter: 123, force: true, skipArchives: true }, + { kind: "build-site", target: "_all", runId: "run-1", runner: "docker" }, + { kind: "deploy-site", target: "jeralyzer", runId: "run-1", preview: "r18", builtAfter: 456 }, + { kind: "deploy-site", target: "jeralyzer", runId: "run-1", to: "local" }, + { kind: "build-hub", target: "_hub", runId: "run-1", allowMissingMedia: true }, + { kind: "deploy-homepage", target: "_homepage", runId: "run-1", preview: "smoke" }, + ]; + for (const req of reqs) { + const spec = stageSpec(req); + assert.equal(spec.kind, `publish-${req.kind}`); + assert.equal(spec.slug, req.target); + assert.equal(spec.params?.runId, "run-1"); + assert.deepEqual(stageRequestFromSpec(JSON.parse(JSON.stringify(spec))), req); + } + // Not a stage, or a request the stage row would refuse. + assert.equal(stageRequestFromSpec({ kind: "sync", slug: "x" }), null); + assert.equal(stageRequestFromSpec({ kind: "publish-update-index", slug: "_index", params: {} }), null, "no run id"); + assert.equal(stageRequestFromSpec({ kind: "publish-update-index", slug: "jeralyzer", params: { runId: "r" } }), null); + assert.equal( + stageRequestFromSpec({ kind: "publish-build-site", slug: "a", params: { runId: "r", kind: "deploy-site" } }), + null, + "the params' kind must agree with the job kind", + ); + assert.match(newPublishRunId("lane", 1_700_000_000_000), /^lane-[0-9a-z]{9}-[0-9a-f]{8}$/); +}); + +function rec(over: Partial<JobRecord>): JobRecord { + return { + id: "J", + kind: "publish-build-site", + queueKey: PUBLISH_QUEUE, + status: "queued", + queuedAt: 1, + logPath: "/dev/null", + ...over, + }; +} + +test("a duplicate is the same kind + target (+ destination) queued or running on `publish`", () => { + const build: StageRequest = { kind: "build-site", target: "alpha", runId: "r2" }; + const records = [rec({ id: "B1", spec: stageSpec({ ...build, runId: "r1" }) })]; + assert.equal(findQueuedStage(records, build)?.id, "B1", "another run's queued build is the same stage"); + assert.equal(findQueuedStage([{ ...records[0], status: "running" }], build)?.id, "B1"); + for (const other of [ + { ...records[0], status: "done" as const }, + { ...records[0], queueKey: "build" }, + { ...records[0], spec: stageSpec({ ...build, target: "beta" }) }, + ]) { + assert.equal(findQueuedStage([other], build), undefined); + } + // A production deploy is not a preview deploy of the same site. + const prod: StageRequest = { kind: "deploy-site", target: "alpha", runId: "r" }; + const prev: StageRequest = { ...prod, preview: "preview" }; + const deploys = [rec({ id: "D1", kind: "publish-deploy-site", spec: stageSpec(prev) })]; + assert.equal(findQueuedStage(deploys, prod), undefined); + assert.equal(findQueuedStage(deploys, prev)?.id, "D1"); + assert.notEqual(stageIdentity(prod), stageIdentity(prev)); + const forced: StageRequest = { ...build, runId: "zzz", force: true }; + assert.equal(stageIdentity(build), stageIdentity(forced), "another run, forced: still the same stage"); +}); + +test("enqueueStage refuses a duplicate with info and the existing job's id, spawning nothing", async () => { + const registry = getRegistry(); + const existing = rec({ + id: "01EXISTINGPUBLISHSTAGE0001", + kind: "publish-update-index", + spec: stageSpec({ kind: "update-index", target: "_index", runId: "run-a" }), + }); + registry.register(existing); + try { + const res = await enqueueStage(getPaths(), { kind: "update-index", target: "_index", runId: "run-b" }); + assert.equal(res.ok, false); + assert.ok(!res.ok && res.info === true); + assert.ok(!res.ok && res.jobId === existing.id); + assert.match(!res.ok ? res.error : "", /Update the index _index is already queued/); + } finally { + existing.status = "cancelled"; + } +}); + +test("enqueuePublishRun enqueues the plan in order under one run id, carrying each step's preconditions", async () => { + const plan: PublishPlan = { + runStart: 1000, + steps: [ + { kind: "update-index", target: "_index", reason: "stale" }, + { kind: "build-site", target: "alpha", indexAfter: 1000, reason: "1 channel changed" }, + { kind: "deploy-site", target: "alpha", preview: "preview", builtAfter: 1000, reason: "its build", previewUrl: "https://preview.alpha.pages.dev" }, + ], + skipped: [{ kind: "build-site", target: "beta", reason: "stale, and its policy is off" }], + }; + const seen: StageRequest[] = []; + const res = await enqueuePublishRun(getPaths(), plan, { + runId: "run-test", + enqueue: async (_paths, req) => { + seen.push(req); + if (req.kind === "build-site") return { ok: false, info: true, jobId: "OLD", error: "already queued" }; + return { + ok: true, + jobId: `J-${req.kind}`, + stream: new ReadableStream<string>(), + done: new Promise(() => {}), + }; + }, + }); + assert.deepEqual( + seen.map((r) => [r.kind, r.target, r.runId, r.indexAfter, r.builtAfter, r.preview]), + [ + ["update-index", "_index", "run-test", undefined, undefined, undefined], + ["build-site", "alpha", "run-test", 1000, undefined, undefined], + ["deploy-site", "alpha", "run-test", undefined, 1000, "preview"], + ], + ); + for (const r of seen) assert.equal("reason" in r, false, "a request carries no plan words"); + assert.equal(res.runId, "run-test"); + assert.deepEqual(res.jobs, [ + { kind: "update-index", target: "_index", jobId: "J-update-index" }, + { kind: "build-site", target: "alpha", jobId: "OLD", existing: true }, + { kind: "deploy-site", target: "alpha", jobId: "J-deploy-site", previewUrl: "https://preview.alpha.pages.dev" }, + ]); + assert.deepEqual(res.skipped, plan.skipped); +}); + +test("a real stage job: runManagedCommand on the `publish` queue, the spec on its meta, the child's exit code kept", async () => { + const jobsDir = await mkdtemp(path.join(os.tmpdir(), "publish-stage-job-")); + try { + const paths: Paths = { ...getPaths(), jobsDir }; + // An empty run id: the stage row refuses it (exit 2) before any lock or + // file is touched. + const req: StageRequest = { kind: "update-index", target: "_index", runId: "" }; + const res = await enqueueStage(paths, req, { background: true }); + assert.ok(res.ok, "enqueued"); + if (!res.ok) return; + void res.stream.cancel(); + const live = getRegistry().get(res.jobId); + assert.equal(live?.queueKey, PUBLISH_QUEUE); + assert.equal(live?.kind, "publish-update-index"); + assert.equal(live?.background, true); + assert.equal(live?.channelSlug, undefined); + const term = await res.done; + assert.equal(term.status, "failed"); + assert.equal(getRegistry().get(res.jobId)?.exitCode, 2); + // The meta is written after the end; give the serial writer a moment. + let meta = await readJobMeta(paths, res.jobId); + for (let i = 0; i < 50 && meta?.status !== "failed"; i++) { + await new Promise((r) => setTimeout(r, 20)); + meta = await readJobMeta(paths, res.jobId); + } + assert.equal(meta?.status, "failed"); + assert.deepEqual(meta?.spec, stageSpec(req)); + assert.match(await readFile(path.join(jobsDir, `${res.jobId}.log`), "utf8"), /--run-id is required/); + } finally { + await rm(jobsDir, { recursive: true, force: true }); + } +}); diff --git a/common/publish/publishStages.ts b/common/publish/publishStages.ts @@ -0,0 +1,196 @@ +// ENQUEUEING PUBLISH STAGES (release 18): one stage = one job on the editor's +// `publish` queue, run as a child process. +// +// enqueueStage(paths, req) one stage — refused (info) when the same +// stage is already queued or running +// enqueuePublishRun(paths, plan) a whole plan at once (Publish now, the +// site page's Build & deploy) +// +// Each job is `runManagedCommand` over `stageCommand(paths, req)` (S1's +// stageRun.ts: `tsx bin/archilyzer.ts stage <kind> <target> …`, cwd common/), +// so Cancel kills the child — and the child's tree (next build, wrangler, +// docker). The queue is ONE FIFO (`publish`, concurrency 1): stages run one at +// a time, beside the publish lock the child takes. +// +// ORDER IS ON DISK, NOT IN MEMORY: a run's builds carry `indexAfter` and its +// deploys `builtAfter` (the plan sets them), and a child whose precondition is +// not met when it starts exits 3 — the job ends `failed` with the sentence in +// its log. That is intended: nothing here chains one job to another. +// +// The spec IS the request: `{kind: "publish-<stage>", slug: target, params: +// {runId, …req}}` (`parseJobSpec` requires `slug`; there is no `channelSlug`), +// which is what Retry (editor jobReplayRegistry.ts) re-enqueues from. + +import { getRegistry, type JobRecord } from "../jobs/registry"; +import type { JobSpec } from "../jobs/jobSpec"; +import { runManagedCommand, type StreamActionResult } from "../jobs/streamCommand"; +import { getPaths, type Paths } from "../lib/paths"; +import type { PlanSkip, PublishPlan } from "./publishPlan"; +import { stageCommand } from "./stageRun"; +import { + STAGES, + STAGE_FLAGS, + deployKindOf, + isStageKind, + parseStageArgs, + type StageRequest, +} from "./stages"; +import { newStampId } from "./stamps"; + +export const PUBLISH_QUEUE = "publish"; + +/** A fresh run id (one Publish now, one lane pass): `<prefix>-<time>-<rand>`. */ +export function newPublishRunId(prefix = "run", now = Date.now()): string { + return `${prefix}-${newStampId(now)}`; +} + +/** + * What makes two stage requests the same stage: kind, target and — for a + * deploy — its destination (production, local, or which preview). A + * production deploy queued behind a preview of the same site is another stage. + */ +export function stageIdentity(r: Pick<StageRequest, "kind" | "target" | "preview" | "to">): string { + const where = r.kind.startsWith("deploy") ? `|${deployKindOf(r)}|${r.preview ?? ""}` : ""; + return `${r.kind}|${r.target}${where}`; +} + +/** The request a stage job's spec carries, or null when it is not one. */ +export function stageRequestFromSpec(spec: JobSpec | null | undefined): StageRequest | null { + if (!spec || !spec.kind.startsWith("publish-")) return null; + const kind = spec.kind.slice("publish-".length); + if (!isStageKind(kind)) return null; + const p = spec.params ?? {}; + if (p.kind !== undefined && p.kind !== kind) return null; + // Through the stage row's own parser: one definition of a valid request. + const flags: Record<string, string | boolean | undefined> = {}; + for (const [flag, type] of Object.entries(STAGE_FLAGS)) { + const key = flag.replace(/-([a-z])/g, (_, c: string) => c.toUpperCase()); + const v = p[key]; + if (v === undefined || v === null) continue; + if (type === "boolean") { + if (v === true) flags[flag] = true; + } else if (typeof v === "string" || typeof v === "number") { + flags[flag] = String(v); + } + } + flags["run-id"] = typeof p.runId === "string" ? p.runId : undefined; + const parsed = parseStageArgs([kind, spec.slug], flags); + return "error" in parsed ? null : parsed; +} + +/** The spec a stage job records (its replay descriptor). */ +export function stageSpec(req: StageRequest): JobSpec { + const params: Record<string, unknown> = {}; + for (const [k, v] of Object.entries(req)) if (v !== undefined) params[k] = v; + return { kind: STAGES[req.kind].jobKind, slug: req.target, params }; +} + +/** A stage job already queued or running for the same stage, if any. */ +export function findQueuedStage( + records: readonly Pick<JobRecord, "id" | "kind" | "queueKey" | "status" | "spec">[], + req: StageRequest, +): Pick<JobRecord, "id" | "status"> | undefined { + const want = stageIdentity(req); + return records.find((r) => { + if (r.queueKey !== PUBLISH_QUEUE) return false; + if (r.status !== "queued" && r.status !== "running") return false; + if (r.kind !== STAGES[req.kind].jobKind) return false; + const other = stageRequestFromSpec(r.spec); + return other !== null && stageIdentity(other) === want; + }); +} + +export type EnqueueStageResult = + | Extract<StreamActionResult, { ok: true }> + | { ok: false; error: string; info?: boolean; jobId?: string }; + +/** + * Enqueue ONE stage on the `publish` queue. A duplicate (the same stage + * already queued or running) is refused with `info: true` and the existing + * job's id — nothing is enqueued twice. `background` puts it behind any + * foreground stage (the lane's are background; a click is foreground). + */ +export async function enqueueStage( + paths: Paths, + req: StageRequest, + opts: { background?: boolean; env?: NodeJS.ProcessEnv } = {}, +): Promise<EnqueueStageResult> { + const dup = findQueuedStage(getRegistry().list(), req); + if (dup) { + return { + ok: false, + info: true, + jobId: dup.id, + error: `${STAGES[req.kind].label} ${req.target} is already ${dup.status} (job ${dup.id})`, + }; + } + const cmd = stageCommand(paths, req, opts.env); + return runManagedCommand({ + kind: STAGES[req.kind].jobKind, + queueKey: PUBLISH_QUEUE, + paths, + spec: stageSpec(req), + ...(opts.background ? { background: true } : {}), + command: cmd.command, + args: cmd.args, + cwd: cmd.cwd, + env: cmd.env, + }); +} + +export type PublishRunJob = { + kind: StageRequest["kind"]; + target: string; + jobId: string; + previewUrl?: string; + // True when the same stage was already queued (its id is the existing job). + existing?: boolean; +}; + +export type PublishRunResult = { + runId: string; + jobs: PublishRunJob[]; + skipped: PlanSkip[]; + refused: { kind: StageRequest["kind"]; target: string; error: string }[]; +}; + +/** + * Enqueue a whole plan at once, in its order, under one run id. Nobody reads + * the jobs' streams here: each is released (its log is on disk), and the + * caller follows the ids (`/api/jobs/<id>/log`, `--wait`). + */ +export async function enqueuePublishRun( + paths: Paths = getPaths(), + plan: PublishPlan, + opts: { + runId?: string; + background?: boolean; + env?: NodeJS.ProcessEnv; + // The one-stage enqueue (tests inject it). + enqueue?: typeof enqueueStage; + } = {}, +): Promise<PublishRunResult> { + const enqueue = opts.enqueue ?? enqueueStage; + const runId = opts.runId ?? newPublishRunId(); + const result: PublishRunResult = { runId, jobs: [], skipped: [...plan.skipped], refused: [] }; + for (const step of plan.steps) { + const { reason: _reason, previewUrl, ...rest } = step; + const req: StageRequest = { ...rest, runId }; + const res = await enqueue(paths, req, { background: opts.background, env: opts.env }); + if (res.ok) { + void res.stream.cancel(); + result.jobs.push({ kind: req.kind, target: req.target, jobId: res.jobId, ...(previewUrl ? { previewUrl } : {}) }); + } else if (res.info && res.jobId) { + result.jobs.push({ + kind: req.kind, + target: req.target, + jobId: res.jobId, + existing: true, + ...(previewUrl ? { previewUrl } : {}), + }); + } else { + result.refused.push({ kind: req.kind, target: req.target, error: res.error }); + } + } + return result; +} diff --git a/common/publish/publishState.test.ts b/common/publish/publishState.test.ts @@ -0,0 +1,165 @@ +// The publish state READ (release 18): readPublishInputs over a scratch corpus +// — the sites and their policies, the config mtimes, the ingest job metas and +// the report regenerations, the live and the ended publish stages. +// +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/publishState.test.ts +// +// SETTINGS SEAM, as settingsSchema.test.ts: `getPaths()` memoizes at module +// scope, so every root is set before anything imports lib/paths. Never a real +// corpus. + +import { mkdirSync, mkdtempSync, utimesSync, writeFileSync } from "node:fs"; +import { rm } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { test, after } from "node:test"; +import assert from "node:assert/strict"; + +const ROOT = mkdtempSync(path.join(os.tmpdir(), "publish-state-")); +process.env.TRANSCRIPTS_DIR = path.join(ROOT, "transcripts"); +process.env.SITES_DIR = path.join(ROOT, "transcripts", "sites"); +process.env.EXPORT_INDEX_DIR = path.join(ROOT, "export-index"); +process.env.EXPORT_BUILDS_DIR = path.join(ROOT, "builds"); +process.env.SETTINGS_FILE = path.join(ROOT, "settings.json"); +process.env.CHARTS_CONFIG_FILE = path.join(ROOT, "charts.json"); +process.env.SEARCH_ALIASES_FILE = path.join(ROOT, "transcripts", "search-aliases.json"); +process.env.CURATED_TAGS_FILE = path.join(ROOT, "transcripts", "tags.json"); + +const { getPaths } = await import("../lib/paths"); +const { readPublishInputs, readPublishStatus, resetPublishStateCache } = await import("./publishState"); +const { stageSpec } = await import("./publishStages"); +type JobRecord = import("../jobs/registry").JobRecord; + +after(() => rm(ROOT, { recursive: true, force: true })); + +const paths = getPaths(); +const T = 1_900_000_000_000; +const sec = (ms: number) => ms / 1000; + +function writeJson(file: string, value: unknown, mtime?: number): void { + mkdirSync(path.dirname(file), { recursive: true }); + writeFileSync(file, JSON.stringify(value)); + if (mtime !== undefined) utimesSync(file, sec(mtime), sec(mtime)); +} + +function meta(id: string, m: Record<string, unknown>): void { + writeJson(path.join(paths.jobsDir, `${id}.meta.json`), { id, queueKey: "q", queuedAt: T - 100_000, ...m }); +} + +function setup(): void { + writeJson(paths.settingsFile, { publish: { enabled: true, hub: "build", previewBranch: "smoke" } }, T - 500_000); + writeJson( + path.join(paths.sitesDir, "alpha", "site.json"), + { + siteId: "alpha", + siteTitle: "Alpha", + channels: [{ slug: "a-one" }, { slug: "shared" }], + cloudflareProject: "alpha", + publish: { auto: "preview" }, + }, + T - 400_000, + ); + writeJson( + path.join(paths.sitesDir, "secret", "site.json"), + { siteId: "secret", siteTitle: "Secret", channels: [{ slug: "shared" }], audience: "private", publish: { auto: "production" } }, + T - 300_000, + ); + writeJson(path.join(paths.sitesDir, "alpha", "tags.json"), {}, T - 50_000); + // Ingest metas: a member's sync ended done; a running download (not yet an + // ingest); a non-ingest kind; a done ingest with no channel (the index only). + meta("01J00000000000000000000001", { kind: "sync", status: "done", channelSlug: "a-one", endedAt: T - 20_000 }); + meta("01J00000000000000000000002", { kind: "download-missing", status: "running", channelSlug: "shared" }); + meta("01J00000000000000000000003", { kind: "scan-media", status: "done", channelSlug: "shared", endedAt: T - 1_000 }); + meta("01J00000000000000000000004", { kind: "normalize-transcripts", status: "done", endedAt: T - 10_000 }); + // An ended publish stage: a deploy refused for want of its build. + meta("01J00000000000000000000005", { + kind: "publish-deploy-site", + queueKey: "publish", + status: "failed", + exitCode: 3, + endedAt: T - 5_000, + spec: stageSpec({ kind: "deploy-site", target: "alpha", runId: "run-x", preview: "smoke", builtAfter: T - 9_000 }), + }); + // A lane unit's only trace: the channel's report regenerated. + writeJson(path.join(paths.channelsDir, "shared", "snapshot.json"), {}, T - 15_000); + mkdirSync(path.join(paths.channelsDir, "a-one"), { recursive: true }); +} + +test("readPublishInputs: sites, policies, mtimes, ingest signals, ended and live stages", async () => { + setup(); + resetPublishStateCache(); + const live: JobRecord[] = [ + { + id: "01JLIVE0000000000000000001", + kind: "publish-update-index", + queueKey: "publish", + status: "running", + queuedAt: T - 2_000, + startedAt: T - 1_000, + logPath: "/dev/null", + spec: stageSpec({ kind: "update-index", target: "_index", runId: "run-live" }), + }, + // A finished one in the registry wins over its (absent) meta. + { + id: "01JLIVE0000000000000000002", + kind: "fetch-posts", + queueKey: "platform:x.com", + status: "done", + channelSlug: "shared", + queuedAt: T - 9_000, + endedAt: T - 3_000, + logPath: "/dev/null", + }, + ]; + const i = await readPublishInputs(paths, { now: T, live, laneKnown: false }); + + assert.deepEqual(i.sites.map((s) => [s.siteId, s.policy, s.private, s.members.join(",")]), [ + ["alpha", "preview", false, "a-one,shared"], + ["secret", "build", true, "shared"], + ]); + const alpha = i.sites[0]; + assert.equal(alpha.configChangedAt, T - 50_000, "its tags.json is the newest of its config"); + assert.equal(alpha.cloudflareProject, "alpha"); + assert.equal(i.sites[1].deployProblem !== null, true, "a private site is never deployed"); + assert.equal(i.settings.previewBranch, "smoke"); + assert.equal(i.hub.policy, "build"); + assert.equal(i.homepage.policy, "off"); + + // Ingest: the member's sync, the registry's fetch-posts, the report. + assert.equal(i.ingestByChannel["a-one"], T - 20_000); + assert.equal(i.ingestByChannel.shared, T - 3_000, "fetch-posts (registry) beats the snapshot (T-15s)"); + assert.equal(i.index.lastIngestDoneAt, T - 3_000); + // The index's config: the settings file is the newest of the index inputs + // here — and tags.json of a site is not one of them. + assert.equal(i.index.configChangedAt, T - 300_000); + + assert.deepEqual( + i.jobs.map((j) => [j.kind, j.target, j.status, j.runId]), + [["update-index", "_index", "running", "run-live"]], + ); + assert.deepEqual( + i.ended.map((e) => [e.kind, e.target, e.exitCode, e.builtAfter, e.preview]), + [["deploy-site", "alpha", 3, T - 9_000, "smoke"]], + ); + assert.equal(i.lane.known, false); + assert.equal(i.index.stamp, null); +}); + +test("readPublishStatus: the chips a fresh scratch corpus shows", async () => { + setup(); + resetPublishStateCache(); + const s = await readPublishStatus(paths, { now: T, live: [], laneKnown: false }); + assert.equal(s.index.chip.text, "no index yet"); + assert.deepEqual(s.plan.steps.map((x) => `${x.kind} ${x.target}`), [ + "update-index _index", + "build-site alpha", + "deploy-site alpha", + "build-site secret", + "build-hub _hub", + ]); + const alpha = s.sites.find((t) => t.target === "alpha")!; + assert.deepEqual(alpha.chips.deployed, { tone: "blocked", text: "waiting for its build" }); + assert.equal(alpha.previewUrl, "https://smoke.alpha.pages.dev"); + assert.equal(s.lane.enabled, true); + assert.equal(s.lane.due, true, "no stamp: a pass is due"); +}); diff --git a/common/publish/publishState.ts b/common/publish/publishState.ts @@ -0,0 +1,401 @@ +// THE PUBLISH STATE, READ (release 18): what `buildPublishStatus` folds. +// +// `readPublishInputs(paths)` reads, and only reads: +// - the stamps and the bundles' problems (`readNeedsInput`, stageBodies.ts — +// the same reader a stage child judges itself by); +// - the sites (membership = site.json `channels`) and their policies, the +// hub's and the homepage's (settings.publish); +// - the input files' mtimes the plan lists: for the index tags.json, +// search-aliases.json, duplicates*.json, every sites/*/site.json, +// homepage.json, the settings file and the charts config; for a site its +// site.json, its tags and aliases and the corpus-wide three; for the hub +// homepage.json and every site.json; for the homepage homepage.json; +// - the ingest signal, per channel: every job meta of an ingest kind +// (jobs/jobKinds.ts `isIngestKind`) that ENDED `done`, read through the +// registry and the `.jobs` sidecars (no new writer anywhere) — plus, for +// the lane units that make no job record at all (a transcription, a +// digest or a backfill dispatched by its runner), the channel's report +// regeneration (`channels/<slug>/snapshot.json`'s mtime): every runner unit +// requests one when it settles; +// - the live publish stages (queued / running on the `publish` queue) and +// the newest ENDED stage per kind + target; +// - the lane's live state (publishLaneState.ts) and the running commit. +// +// `readPublishStatus(paths)` is the one function every surface calls — the +// /sites Publish panel, `GET /api/ops/publish`, `archilyzer publish status` +// and the lane — and returns `buildPublishStatus(inputs, now)`. +// +// The signals decide WHEN to run; the stage children decide WHAT changed (the +// plan's ruling). Over-signalling costs a no-op stage; missing a signal costs +// a stale site — so every doubt is a signal. + +import { readdir, stat } from "node:fs/promises"; +import path from "node:path"; +import { mapConcurrent } from "../lib/concurrency"; +import { DUPLICATES_FILENAME, DUPLICATE_OVERRIDES_FILENAME } from "../lib/duplicates"; +import { getHomepageConfig } from "../lib/homepage"; +import { getPaths, type Paths } from "../lib/paths"; +import { PROJECT_URL } from "../lib/project"; +import { getSettings, normalizeHomepageUrl } from "../lib/settings"; +import { + isListedSite, + isPrivateSite, + listSites, + siteAliasesFile, + siteConfigFile, + siteTagsFile, + sitePublishPolicy, +} from "../lib/site"; +import { getRegistry, type JobRecord } from "../jobs/registry"; +import { readJobMeta } from "../jobs/jobMeta"; +import type { JobSpec } from "../jobs/jobSpec"; +import { isIngestKind, isPublishStageJobKind } from "../jobs/jobKinds"; +import { snapshotPath } from "../controller/channels"; +import { + buildPublishStatus, + type PublishInputs, + type PublishJob, + type PublishJobEnded, + type PublishSiteInput, + type PublishStatus, +} from "./publishPlan"; +import { publishLaneLive } from "./publishLaneState"; +import { checkoutInfo, mainHeadOf, readNeedsInput } from "./stageBodies"; +import { isStageKind, type StageKind } from "./stages"; + +// --------------------------------------------------------------------------- +// mtimes +// --------------------------------------------------------------------------- + +async function mtimeOf(file: string): Promise<number | null> { + try { + return (await stat(file)).mtimeMs; + } catch { + return null; + } +} + +function newest(...xs: (number | null | undefined)[]): number | null { + let out: number | null = null; + for (const x of xs) if (typeof x === "number" && (out === null || x > out)) out = x; + return out; +} + +/** duplicates.json, duplicates.overrides.json and any other duplicates*.json. */ +async function duplicatesFiles(paths: Pick<Paths, "transcriptsDir">): Promise<string[]> { + const out = new Set([ + path.join(paths.transcriptsDir, DUPLICATES_FILENAME), + path.join(paths.transcriptsDir, DUPLICATE_OVERRIDES_FILENAME), + ]); + try { + for (const name of await readdir(paths.transcriptsDir)) { + if (/^duplicates.*\.json$/.test(name)) out.add(path.join(paths.transcriptsDir, name)); + } + } catch { + /* no corpus dir: nothing changed */ + } + return [...out]; +} + +// --------------------------------------------------------------------------- +// Job metas (ingest + ended publish stages), read once per id +// --------------------------------------------------------------------------- + +type MetaLite = { + id: string; + kind: string; + queueKey: string; + status: string; + channelSlug?: string; + endedAt?: number; + exitCode?: number; + spec?: JobSpec; +}; + +const TERMINAL = new Set(["done", "failed", "cancelled"]); + +// A terminal meta never changes again (the boot pass only closes `queued` and +// `running` ones), so each is read once per process. +const terminalMetas = new Map<string, MetaLite>(); + +function lite(r: Pick<JobRecord, "id" | "kind" | "queueKey" | "status" | "channelSlug" | "endedAt" | "exitCode" | "spec">): MetaLite { + return { + id: r.id, + kind: r.kind, + queueKey: r.queueKey, + status: r.status, + ...(r.channelSlug ? { channelSlug: r.channelSlug } : {}), + ...(typeof r.endedAt === "number" ? { endedAt: r.endedAt } : {}), + ...(typeof r.exitCode === "number" ? { exitCode: r.exitCode } : {}), + ...(r.spec ? { spec: r.spec } : {}), + }; +} + +/** Every job this process and the `.jobs` sidecars know, the registry winning. */ +export async function readJobLites( + paths: Pick<Paths, "jobsDir">, + live: readonly JobRecord[] = getRegistry().list(), +): Promise<MetaLite[]> { + const byId = new Map<string, MetaLite>(); + let ids: string[] = []; + try { + ids = (await readdir(paths.jobsDir)) + .filter((n) => n.endsWith(".meta.json")) + .map((n) => n.slice(0, -".meta.json".length)); + } catch { + ids = []; + } + const onDisk = new Set(ids); + for (const id of terminalMetas.keys()) if (!onDisk.has(id)) terminalMetas.delete(id); + const liveIds = new Set(live.map((r) => r.id)); + const toRead = ids.filter((id) => !liveIds.has(id) && !terminalMetas.has(id)); + const read = await mapConcurrent(toRead, 32, (id) => readJobMeta(paths as Paths, id)); + for (const m of read) { + if (!m) continue; + const l = lite(m as unknown as JobRecord); + if (TERMINAL.has(l.status)) terminalMetas.set(l.id, l); + byId.set(l.id, l); + } + for (const id of ids) { + const hit = terminalMetas.get(id); + if (hit) byId.set(id, hit); + } + for (const r of live) byId.set(r.id, lite(r)); + return [...byId.values()]; +} + +/** Forget the meta memo (tests). */ +export function resetPublishStateCache(): void { + terminalMetas.clear(); + memo.clear(); +} + +export type IngestSignals = { + lastDoneAt: number | null; + byChannel: Record<string, number>; +}; + +/** The newest ingest end per channel (and overall) from the job metas. */ +export function ingestFromMetas(metas: readonly MetaLite[]): IngestSignals { + const byChannel: Record<string, number> = {}; + let lastDoneAt: number | null = null; + for (const m of metas) { + if (m.status !== "done" || typeof m.endedAt !== "number" || !isIngestKind(m.kind)) continue; + lastDoneAt = newest(lastDoneAt, m.endedAt); + if (m.channelSlug) byChannel[m.channelSlug] = Math.max(byChannel[m.channelSlug] ?? 0, m.endedAt); + } + return { lastDoneAt, byChannel }; +} + +/** Each channel's report regeneration (the lane units' only trace). */ +async function snapshotSignals(paths: Paths): Promise<Record<string, number>> { + const out: Record<string, number> = {}; + let slugs: string[] = []; + try { + slugs = (await readdir(paths.channelsDir, { withFileTypes: true })) + .filter((d) => d.isDirectory()) + .map((d) => d.name); + } catch { + return out; + } + await mapConcurrent(slugs, 32, async (slug) => { + const at = await mtimeOf(snapshotPath(paths, slug)); + if (at !== null) out[slug] = at; + }); + return out; +} + +function stageKindOfJob(kind: string): StageKind | null { + const k = kind.startsWith("publish-") ? kind.slice("publish-".length) : ""; + return isStageKind(k) ? k : null; +} + +const num = (v: unknown): number | undefined => (typeof v === "number" && Number.isFinite(v) ? v : undefined); +const str = (v: unknown): string | undefined => (typeof v === "string" && v ? v : undefined); +const dest = (v: unknown): "pages" | "local" | undefined => (v === "pages" || v === "local" ? v : undefined); + +/** The publish stages queued or running on the `publish` queue. */ +export function livePublishJobs(live: readonly JobRecord[]): PublishJob[] { + const out: PublishJob[] = []; + for (const r of live) { + if (r.status !== "queued" && r.status !== "running") continue; + if (!isPublishStageJobKind(r.kind)) continue; + const kind = stageKindOfJob(r.kind); + if (!kind) continue; + const p = r.spec?.params ?? {}; + out.push({ + id: r.id, + kind, + target: r.spec?.slug ?? str(p.target) ?? "?", + status: r.status, + runId: str(p.runId) ?? null, + queuedAt: r.queuedAt, + ...(r.startedAt ? { startedAt: r.startedAt } : {}), + ...(str(p.preview) ? { preview: str(p.preview) } : {}), + ...(dest(p.to) ? { to: dest(p.to) } : {}), + }); + } + return out.sort((a, b) => a.queuedAt - b.queuedAt); +} + +/** The newest ENDED publish stage per kind + target. */ +export function endedPublishJobs(metas: readonly MetaLite[]): PublishJobEnded[] { + const newestBy = new Map<string, PublishJobEnded>(); + for (const m of metas) { + if (!TERMINAL.has(m.status) || typeof m.endedAt !== "number") continue; + const kind = stageKindOfJob(m.kind); + if (!kind) continue; + const p = m.spec?.params ?? {}; + const target = m.spec?.slug ?? str(p.target); + if (!target) continue; + const e: PublishJobEnded = { + id: m.id, + kind, + target, + status: m.status as PublishJobEnded["status"], + exitCode: typeof m.exitCode === "number" ? m.exitCode : null, + endedAt: m.endedAt, + runId: str(p.runId) ?? null, + ...(num(p.indexAfter) !== undefined ? { indexAfter: num(p.indexAfter) } : {}), + ...(num(p.builtAfter) !== undefined ? { builtAfter: num(p.builtAfter) } : {}), + ...(str(p.preview) ? { preview: str(p.preview) } : {}), + ...(dest(p.to) ? { to: dest(p.to) } : {}), + }; + const key = `${kind}:${target}`; + const prev = newestBy.get(key); + if (!prev || prev.endedAt < e.endedAt) newestBy.set(key, e); + } + return [...newestBy.values()]; +} + +// --------------------------------------------------------------------------- +// git facts, memoized (a status poll must not fork git every second) +// --------------------------------------------------------------------------- + +const memo = new Map<string, { at: number; value: string | null }>(); +const GIT_MEMO_MS = 30_000; + +async function memoized(key: string, now: number, read: () => Promise<string | null>): Promise<string | null> { + const hit = memo.get(key); + if (hit && now - hit.at < GIT_MEMO_MS) return hit.value; + const value = await read().catch(() => null); + memo.set(key, { at: now, value }); + return value; +} + +// --------------------------------------------------------------------------- +// readPublishInputs / readPublishStatus +// --------------------------------------------------------------------------- + +export type ReadPublishOpts = { + now?: number; + // Whether this process holds the lane (the editor). The CLI passes false. + laneKnown?: boolean; + // The registry's records (tests inject them). + live?: readonly JobRecord[]; +}; + +export async function readPublishInputs( + paths: Paths = getPaths(), + opts: ReadPublishOpts = {}, +): Promise<PublishInputs> { + const now = opts.now ?? Date.now(); + const settings = getSettings(); + const sites = listSites(paths); + const needs = await readNeedsInput(paths); + const live = opts.live ?? getRegistry().list(); + const metas = await readJobLites(paths, live); + const ingest = ingestFromMetas(metas); + const snapshots = await snapshotSignals(paths); + const byChannel: Record<string, number> = { ...ingest.byChannel }; + for (const [slug, at] of Object.entries(snapshots)) byChannel[slug] = Math.max(byChannel[slug] ?? 0, at); + const lastIngestDoneAt = newest(ingest.lastDoneAt, ...Object.values(snapshots)); + + const dupes = await duplicatesFiles(paths); + const corpusWide = newest( + await mtimeOf(paths.globalTagsFile), + await mtimeOf(paths.globalAliasesFile), + ...(await Promise.all(dupes.map(mtimeOf))), + ); + const siteJsonAt = new Map<string, number | null>(); + for (const s of sites) siteJsonAt.set(s.siteId, await mtimeOf(siteConfigFile(paths, s.siteId))); + const homepageJsonAt = await mtimeOf(paths.homepageConfigFile); + const indexConfigAt = newest( + corpusWide, + ...siteJsonAt.values(), + homepageJsonAt, + await mtimeOf(paths.settingsFile), + await mtimeOf(paths.chartsConfigFile), + ); + + const siteInputs: PublishSiteInput[] = []; + for (const s of sites) { + const t = needs.sites[s.siteId]; + siteInputs.push({ + siteId: s.siteId, + title: s.siteTitle, + private: isPrivateSite(s), + listed: isListedSite(s), + members: s.channels.map((c) => c.slug), + built: t.built, + deployed: t.deployed, + bundleProblem: t.bundleProblem, + deployProblem: t.deployProblem ?? null, + pagesProblem: t.pagesProblem ?? null, + configChangedAt: newest( + siteJsonAt.get(s.siteId), + await mtimeOf(siteTagsFile(paths, s.siteId)), + await mtimeOf(siteAliasesFile(paths, s.siteId)), + corpusWide, + ), + policy: sitePublishPolicy(s), + cloudflareProject: s.cloudflareProject?.trim() || null, + url: s.siteUrl ?? null, + }); + } + + const mainHead = await memoized("mainHead", now, () => mainHeadOf(paths)); + const commit = await memoized("commit", now, async () => (await checkoutInfo(paths)).commit); + return { + index: { stamp: needs.index.stamp, lastIngestDoneAt, configChangedAt: indexConfigAt }, + ingestByChannel: byChannel, + commit, + sites: siteInputs, + hub: { + built: needs.hub.built, + deployed: needs.hub.deployed, + bundleProblem: needs.hub.bundleProblem, + deployProblem: null, + pagesProblem: needs.hub.pagesProblem ?? null, + configChangedAt: newest(homepageJsonAt, ...siteJsonAt.values()), + policy: settings.publish.hub, + cloudflareProject: getHomepageConfig(paths).cloudflareProject?.trim() || null, + url: normalizeHomepageUrl(settings.homepageUrl) || null, + }, + homepage: { + built: needs.homepage.built, + deployed: needs.homepage.deployed, + bundleProblem: needs.homepage.bundleProblem, + deployProblem: null, + pagesProblem: needs.homepage.pagesProblem ?? null, + configChangedAt: homepageJsonAt, + policy: settings.publish.homepage, + cloudflareProject: (await import("./build")).HOMEPAGE_PAGES_PROJECT, + url: PROJECT_URL, + mainHead, + }, + settings: settings.publish, + jobs: livePublishJobs(live), + ended: endedPublishJobs(metas), + lane: publishLaneLive(opts.laneKnown ?? true), + }; +} + +/** THE status every surface reads (the one function S4's route calls). */ +export async function readPublishStatus( + paths: Paths = getPaths(), + opts: ReadPublishOpts = {}, +): Promise<PublishStatus> { + const now = opts.now ?? Date.now(); + return buildPublishStatus(await readPublishInputs(paths, { ...opts, now }), now); +}