// 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; enqueue: (req: StageRequest, opts: { background: true }) => Promise; settings: () => PublishSettings; now: () => number; // Resolves after `ms`, or as soon as a signal fires. sleep: (ms: number, signals: AbortSignal[]) => Promise; }; function realSleep(ms: number, signals: AbortSignal[]): Promise { 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) => `${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 { const mem = publishLaneMemory(); const passStart = deps.now(); const runId = newPublishRunId("lane", passStart); const attempted = new Set(); const failedBuilds = new Set(); 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((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 { 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 { 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 { 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); }