Archilyzer · Source

archilyzer

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

commit 1cf1a7f7410854c0cf8160f256b762f8ce95ba48
parent 3499cbd7ba8e03020d04f4f05a3377f8d8c69a85
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Tue,  6 Oct 2026 11:41:20 -0400

Merge r18/publish-lane (slice S3: the publish status view and planner, stage enqueueing with on-disk preconditions, the publish lane runner, settings.publish and site publish policy, publish-* job kinds, publish status/now)

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

Diffstat:
MSETTINGS.md | 35+++++++++++++++++++++++++++++++++++
MSITE.md | 17++++++++++++++++-
Mcommon/bin/archilyzer.ts | 17+++++++++++++++--
Acommon/bin/publishNow.test.ts | 75+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/bin/publishNow.ts | 127+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/bootQueuedJobs.test.ts | 40++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/bootQueuedJobs.ts | 12++++++++++--
Acommon/jobs/jobDetail.test.ts | 42++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/jobDetail.ts | 13++++++++++++-
Mcommon/jobs/jobKinds.test.ts | 72++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/jobs/jobKinds.ts | 178+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/autoQueueTypes.ts | 11+++++++++++
Mcommon/lib/fileSchemaDocs.ts | 7+++++--
Mcommon/lib/pauseGates.test.ts | 22+++++++++++++++++++++-
Mcommon/lib/pauseGates.ts | 16+++++++++++++---
Mcommon/lib/settingsDocs.ts | 8++++++++
Mcommon/lib/settingsSchema.test.ts | 79++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/settingsSchema.ts | 113+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/site.ts | 5+++++
Mcommon/lib/siteSchema.test.ts | 62+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcommon/lib/siteSchema.ts | 65++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Acommon/publish/publishLaneState.ts | 73+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishPlan.test.ts | 543+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/publish/publishPlan.ts | 876+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
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+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/views/publishStatus.ts | 41+++++++++++++++++++++++++++++++++++++++++
Meditor/CHANGELOG.md | 2++
Meditor/app/components/lanes/pauseControl.tsx | 15+++++++++++++++
Meditor/app/jobs/jobReplayRegistry.ts | 20++++++++++++++++++++
Meditor/app/operations/actions.ts | 4+++-
Meditor/instrumentation.ts | 13+++++++++++++
Mplans/release-18.md | 152+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Msettings.json.example | 11+++++++++++
38 files changed, 4392 insertions(+), 16 deletions(-)

diff --git a/SETTINGS.md b/SETTINGS.md @@ -47,6 +47,7 @@ A copied example PINS every default it spells — including each lane's `autoQue | [`backfill`](#backfill) | object — see below | | [`attribution`](#attribution) | object — see below | | [`archiveOrg`](#archiveorg) | object — see below | +| [`publish`](#publish) | object — see below | ## `adminTitle` @@ -821,3 +822,37 @@ Default: "maxUploadKiBps": 0 } ``` + +## `publish` + +The publish LANE (release 18): a runner that, when the index is stale, updates it and then builds — and, where a site's own `publish.auto` (site.json) says so, deploys — what changed, one stage at a time on the `publish` queue. OFF by default; the manual stages (Publish now, Build, Deploy) work either way. See PublishSettings and common/publish/publishRunner.ts. + +#### `publish` + +| Key | Default | Description | +|---|---|---| +| `enabled` | `false` | Whether the publish lane's runner runs (`auto-publish` on /jobs). Off by default. Turning it on starts nothing by itself until the index is stale (or there is no index stamp yet) — see `refreshEveryMinutes`. | +| `held` | `false` | The lane's pause gate (lib/pauseGates.ts, lane `publish`). A hold stops the runner DISPATCHING: the stage in flight finishes, no next one starts. It never kills a stage. | +| `checkEveryMinutes` | `10` | How often (minutes) the runner wakes to ask whether a pass is due. Clamped to [1, 1440]; default 10. | +| `refreshEveryMinutes` | `360` | The least time (minutes) between two index updates the lane starts: a pass runs when the index is stale and the last update is at least this old, or when there is no index stamp. Clamped to [0, 43200]; default 360. 0 = whenever the index is stale. | +| `quietHours` | `null` | A local-clock window in which the lane starts no pass, `{ "start": 22, "end": 6 }` (hours [0,23], `[start, end)`, may wrap midnight), or null (default). A pass already running finishes its stage and then waits. | +| `runner` | `"local"` | Which build runner the lane's builds ask for: "local" (default; each site in turn, as a child of the editor) or "docker" (every stale site in containers — a host with a container engine only; in a container it is refused). See PUBLISH.md. | +| `previewBranch` | `"preview"` | The Pages preview branch a `preview` policy deploys to (`wrangler pages deploy --branch <previewBranch>`). Lowercase letters, digits and dashes, never `main`; an invalid name reads as the default, "preview". | +| `hub` | `"off"` | The hub's policy: "off" (default), "build", "preview" or "production". The hub is built when the index or the listed sites changed, and deployed as the policy says — to production only with a Pages project in homepage.json. | +| `homepage` | `"off"` | The homepage's policy, as `hub` (default "off"). Building the homepage also publishes the source mirror (the operator's scrub and denylist files must exist). | + +Default: + +```json +{ + "enabled": false, + "held": false, + "checkEveryMinutes": 10, + "refreshEveryMinutes": 360, + "quietHours": null, + "runner": "local", + "previewBranch": "preview", + "hub": "off", + "homepage": "off" +} +``` diff --git a/SITE.md b/SITE.md @@ -4,7 +4,7 @@ One public site: its branding, its channel grouping and which channels it exposes, persisted to `transcripts/sites/<id>/site.json` (the directory under `$SITES_DIR` when that is set). The schema is `common/lib/siteSchema.ts`. Global operational settings are `settings.json` — see [SETTINGS.md](SETTINGS.md). The PUBLIC `/site.json` a built site serves is a different file (`common/lib/siteDescriptor.ts`). -Every key is optional on read. A missing key reads as its default, an ill-typed one as its default (or is dropped, for the optional ones), and an unknown one is dropped on the next save. A save ALWAYS writes `siteId`, `siteTitle`, `siteDescription`, `headerTitle`, `homeTagline`, `groups`, `defaultGroupId` and `channels`; every other key is written only when it differs from its default (`socialLinks` whenever it is an array, even an empty one). A save is REFUSED when there is no channel group, when `defaultGroupId` names no group, or when a social link's SVG is not safe to inline. +Every key is optional on read. A missing key reads as its default, an ill-typed one as its default (or is dropped, for the optional ones), and an unknown one is dropped on the next save. A save ALWAYS writes `siteId`, `siteTitle`, `siteDescription`, `headerTitle`, `homeTagline`, `groups`, `defaultGroupId` and `channels`; every other key is written only when it differs from its default (`socialLinks` whenever it is an array, even an empty one). A save is REFUSED when there is no channel group, when `defaultGroupId` names no group, when a social link's SVG is not safe to inline, or when a public site's `publish.auto` deploys and it has no `cloudflareProject`. Regenerate this file with `pnpm --filter yt-dlp-transcript-common exec tsx bin/file-schemas-docs.ts`. @@ -34,6 +34,7 @@ Regenerate this file with `pnpm --filter yt-dlp-transcript-common exec tsx bin/f | [`transcriptDownloads`](#transcriptdownloads) | `true` | | [`archiveMaxBytes`](#archivemaxbytes) | absent | | [`hubUrl`](#huburl) | absent | +| [`publish`](#publish) | absent | ## `siteId` @@ -243,3 +244,17 @@ Default: absent Per-site override for the hub this site belongs under. Absent = the family default, `settings.json` `homepageUrl`. Published on the public `/site.json` and `/corpus.json` so a hub can tell member sites from arbitrary added origins; the header does not link to it (release 14). Default: absent + +## `publish` + +What the publish lane — and Publish now — may do with this site when it is stale (release 18): `{ "auto": "off" | "build" | "preview" | "production" }`. Absent = `off`: the lane leaves the site alone (a manual Build or Deploy still works, and "Build all stale" still builds it). `build` rebuilds its bundle; `preview` also deploys it to the Pages preview branch `settings.json` `publish.previewBranch` names; `production` deploys it to production. A private site is clamped to `build` (it is never deployed); `preview` and `production` need a `cloudflareProject` — a save without one is refused, and a file that says so anyway reads as `build`. Only a policy other than `off` is written. + +Default: absent + +#### `publish` + +Per entry — each entry spells its own values. + +| Key | Description | +|---|---| +| `auto` | The policy: `off` (default), `build`, `preview` or `production` — see `publish` above. | diff --git a/common/bin/archilyzer.ts b/common/bin/archilyzer.ts @@ -156,8 +156,21 @@ export const COMMANDS: Command[] = [ }); }, }, - // `publish status` and `publish now` read the publish state view — release 18 - // slice S3 adds their rows here. + // `publish status` and `publish now` read the publish state view (release 18 + // S3; bin/publishNow.ts). + { + path: ["publish", "status"], + usage: + "[--json] the publish status: the index, the lane, and per site / hub / homepage its policy and its built, deployed and live chips, then the plan Publish now would run", + flags: { json: "boolean" }, + run: async ({ flags }) => (await import("./publishNow")).publishStatusRow({ json: flags.json === true }), + }, + { + path: ["publish", "now"], + usage: + "run what Publish now runs — the index update when it is stale, then each policy target's build and deploy (site.json publish.auto; settings.json publish.hub / publish.homepage), never forced — one stage at a time IN THIS PROCESS under the publish lock, as the editor's queue would (the CLI has no queue); exit 0 when every stage ran or was a no-op", + run: async () => (await import("./publishNow")).publishNow(), + }, { path: ["stage"], usage: diff --git a/common/bin/publishNow.test.ts b/common/bin/publishNow.test.ts @@ -0,0 +1,75 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { COMMANDS } from "./archilyzer"; +import { argumentProblem, booleanFlags, resolveCommand, usage } from "./_cli"; +import { parseArgv } from "./_parseFlags"; +import { defaultPublish } from "../lib/settingsSchema"; +import { buildPublishStatus, type PublishInputs } from "../publish/publishPlan"; +import { statusLines } from "./publishNow"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test bin/publishNow.test.ts +// +// `publish status [--json]` and `publish now` (release 18 S3): their rows, and +// the status's lines over a hand-built status. `publish now` itself runs the +// stage bodies, which the stage tests (publish/stages*.test.ts) cover. + +test("publish status takes --json and nothing else; publish now takes nothing", () => { + const status = resolveCommand(COMMANDS, ["publish", "status"])?.command; + assert.deepEqual(status?.path, ["publish", "status"]); + assert.deepEqual(status?.flags, { json: "boolean" }); + const parsed = parseArgv(["publish", "status", "--json"], booleanFlags(COMMANDS)); + assert.deepEqual(parsed, { positionals: ["publish", "status"], flags: { json: true } }); + assert.equal(argumentProblem(status!, parsed.flags, []), null); + assert.match(argumentProblem(status!, {}, ["x"])!, /unexpected argument "x"/); + const now = resolveCommand(COMMANDS, ["publish", "now"])?.command; + assert.deepEqual(now?.path, ["publish", "now"]); + assert.match(argumentProblem(now!, { force: true }, [])!, /unknown flag --force/); + const u = usage(COMMANDS); + assert.match(u, /archilyzer publish status\s+\[--json\] the publish status/); + assert.match(u, /archilyzer publish now\s+run what Publish now runs .* IN THIS PROCESS under the publish lock, as the editor's queue would \(the CLI has no queue\)/); +}); + +test("the status reads as one row per target, then the plan", () => { + const empty = { + built: null, + deployed: null, + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off" as const, + cloudflareProject: null, + url: null, + }; + const i: PublishInputs = { + index: { stamp: null, lastIngestDoneAt: null, configChangedAt: null }, + ingestByChannel: {}, + commit: null, + sites: [{ ...empty, siteId: "alpha", title: "Alpha", private: false, listed: true, members: [], policy: "build" }], + hub: empty, + homepage: { ...empty, mainHead: null }, + settings: defaultPublish(), + jobs: [], + ended: [], + lane: { + known: false, + running: false, + jobId: null, + passRunning: false, + lastCheckAt: null, + nextCheckAt: null, + lastPassAt: null, + lastPassSummary: null, + lastDecision: null, + }, + }; + const lines = statusLines(buildPublishStatus(i, 1_900_000_000_000)); + assert.equal(lines[0], "index no index yet"); + assert.match(lines[1], /^lane {11}lane off — the lane is off$/); + assert.match(lines[3], /^alpha +\| \[build\] +\| built: update the index first \| deployed: not deployed \| live: not checked \| next: build$/); + assert.match(lines[4], /^_hub /); + assert.match(lines[5], /^_homepage /); + assert.equal(lines[7], "Publish now would run:"); + assert.match(lines[8], /^ update-index _index — no index stamp yet$/); + assert.match(lines[9], /^ build-site alpha — built after the index's first update$/); +}); diff --git a/common/bin/publishNow.ts b/common/bin/publishNow.ts @@ -0,0 +1,127 @@ +// `archilyzer publish status` and `archilyzer publish now` (release 18 S3). +// +// publish status [--json] the publish status every surface reads +// (publish/publishState.ts readPublishStatus): one +// row per target with the chips' words, the lane, +// and the plan Publish now would run +// publish now run that plan: update-index when stale, then each +// policy target's build and deploy — never forced +// +// THE CLI HAS NO QUEUE. Where the editor's Publish now enqueues the whole plan +// on its `publish` queue (enqueuePublishRun), `publish now` runs the plan's +// stages one after another IN THIS PROCESS, each under the publish lock — +// exactly what the queue would do: the same stage bodies, in the same order, +// with the same on-disk preconditions (a build carries `indexAfter`, a deploy +// `builtAfter`), so a stage whose precondition is not met refuses (exit 3) and +// the rest are still tried. The index update runs as the same child the +// editor spawns (its heap cap), as `publish index` does. + +import { setKillChildTrees, killChildTreesNow } from "../jobs/runChild"; +import { getPaths, type Paths } from "../lib/paths"; +import { planPublishRun, type PublishStatus, type PublishTargetStatus } from "../publish/publishPlan"; +import { readPublishStatus } from "../publish/publishState"; +import { STAGE_EXIT, runStage } from "../publish/stageRun"; +import type { StageRequest } from "../publish/stages"; +import { cliRunId, publishIndex } from "./publish"; + +type Out = { log: (s: string) => void; error: (s: string) => void }; + +function pad(s: string, n: number): string { + return s.length >= n ? s : s + " ".repeat(n - s.length); +} + +function targetRow(t: PublishTargetStatus): string { + const c = t.chips; + return [ + pad(t.target, 14), + pad(`[${t.policy}]`, 13), + `built: ${c.built.text}`, + `deployed: ${c.deployed.text}`, + `live: ${c.live.text}`, + `next: ${t.next}`, + ].join(" | "); +} + +/** The status as lines (exported for the test). */ +export function statusLines(s: PublishStatus): string[] { + const out: string[] = []; + out.push(`index ${s.index.chip.text}`); + out.push(`lane ${s.lane.chip.text} — ${s.lane.reason}`); + out.push(""); + for (const t of [...s.sites, s.hub, s.homepage]) out.push(targetRow(t)); + out.push(""); + if (s.plan.steps.length === 0) out.push("Publish now would run nothing."); + else { + out.push("Publish now would run:"); + for (const step of s.plan.steps) { + const where = step.preview ? ` --preview ${step.preview}` : step.to === "local" ? " --to local" : ""; + out.push(` ${step.kind} ${step.target}${where} — ${step.reason}`); + } + } + for (const skip of s.plan.skipped) out.push(` (skipped ${skip.kind} ${skip.target}: ${skip.reason})`); + return out; +} + +export async function publishStatusRow( + a: { json?: boolean; paths?: Paths }, + out: Out = console, +): Promise<number> { + const status = await readPublishStatus(a.paths ?? getPaths(), { laneKnown: false }); + if (a.json) out.log(JSON.stringify(status, null, 2)); + else for (const line of statusLines(status)) out.log(line); + return 0; +} + +/** + * Run Publish now's plan in this process, one stage at a time, as the + * editor's queue would. Exit: 0 when every stage ran or was a no-op, else + * the worst stage's code (130 at once on a cancel). + */ +export async function publishNow( + a: { paths?: Paths; signal?: AbortSignal } = {}, + out: Out = console, +): Promise<number> { + const paths = a.paths ?? getPaths(); + const status = await readPublishStatus(paths, { laneKnown: false }); + const plan = planPublishRun(status, { deploys: "policy" }); + const runId = cliRunId(); + for (const skip of plan.skipped) out.log(`[publish] skipped ${skip.kind} ${skip.target}: ${skip.reason}`); + if (plan.steps.length === 0) { + out.log("[publish] nothing to do: the index and every policy target are current"); + return 0; + } + out.log(`[publish] run ${runId}: ${plan.steps.map((s) => `${s.kind} ${s.target}`).join(", ")}`); + const ac = new AbortController(); + const onSignal = () => { + if (!ac.signal.aborted) return ac.abort(); + killChildTreesNow("SIGKILL"); + process.exit(STAGE_EXIT.cancelled); + }; + const signal = a.signal ?? ac.signal; + if (!a.signal) { + process.on("SIGINT", onSignal); + process.on("SIGTERM", onSignal); + } + setKillChildTrees(true); + let worst = 0; + try { + for (const step of plan.steps) { + if (signal.aborted) return STAGE_EXIT.cancelled; + const { reason: _reason, previewUrl: _previewUrl, ...rest } = step; + const req: StageRequest = { ...rest, runId }; + const code = + req.kind === "update-index" + ? await publishIndex({ paths, runId }) + : (await runStage(req, { paths, signal })).code; + if (code === STAGE_EXIT.cancelled) return code; + // A failure (1) outranks a refusal (3) or a usage error (2). + if (code !== 0 && worst !== STAGE_EXIT.failed) worst = code; + } + } finally { + if (!a.signal) { + process.off("SIGINT", onSignal); + process.off("SIGTERM", onSignal); + } + } + return worst; +} diff --git a/common/jobs/bootQueuedJobs.test.ts b/common/jobs/bootQueuedJobs.test.ts @@ -8,6 +8,7 @@ import type { JobMeta } from "./jobMeta"; import type { JobSpec } from "./jobSpec"; import { INTERRUPTED_REASON, + PUBLISH_RESTART_REASON, STORAGE_PASS_WAIT_MS, processIsAlive, settleAfterStoragePass, @@ -555,3 +556,42 @@ test("the queued pass leaves a queued meta whose writer is still alive", async ( await rm(f.root, { recursive: true, force: true }); } }); + +// RELEASE 18: a queued publish stage is never re-queued. The lane and Publish +// now re-derive stages from the stamps; a stage re-queued with its run's +// preconditions would wait on a run that is gone. +test("a queued publish stage is cancelled as `publish`, never re-queued, even with a fresh spec", async () => { + const stage: JobSpec = { + kind: "publish-build-site", + slug: "jeralyzer", + params: { kind: "build-site", target: "jeralyzer", runId: "r1", indexAfter: BOOT - HOUR }, + }; + const f = await fixture([ + { id: "P1", kind: "publish-build-site", queueKey: "publish", spec: stage }, + { id: "P2", kind: "publish-update-index", queueKey: "publish", spec: { kind: "publish-update-index", slug: "_index" } }, + { id: "F1", spec: SPEC }, + ]); + try { + const rq = recordingRequeue(); + const lines: string[] = []; + const res = await settleQueuedJobMetas({ + paths: f.paths, + requeue: rq.fn, + bootedAt: BOOT, + log: (l) => lines.push(l), + }); + assert.deepEqual(rq.calls, [SPEC]); + assert.deepEqual( + res.cancelled.map((c) => [c.id, c.category]).sort(), + [ + ["P1", "publish"], + ["P2", "publish"], + ], + ); + assert.equal((await f.read("P1")).cancelReason, PUBLISH_RESTART_REASON); + assert.match(PUBLISH_RESTART_REASON, /the publish lane re-derives stages from on-disk state/); + assert.match(lines.at(-1) ?? "", /cancelled 2 \(publish 2\)/); + } finally { + await rm(f.root, { recursive: true, force: true }); + } +}); diff --git a/common/jobs/bootQueuedJobs.ts b/common/jobs/bootQueuedJobs.ts @@ -4,6 +4,7 @@ import { writeFileAtomic } from "../lib/jsonFile-server"; import type { Paths } from "../lib/paths"; import { metaPath, readJobMeta, type JobMeta } from "./jobMeta"; import type { JobSpec } from "./jobSpec"; +import { isPublishStageJobKind } from "./jobKinds"; // THE BOOT PASS OVER STALE `queued` METAS (release 9, B4b). // @@ -23,7 +24,10 @@ import type { JobSpec } from "./jobSpec"; // 1. an idle boot (`requeue: null`) cancels everything — it must not resume // work; // 2. `sync` is never re-queued: the heartbeat re-derives due syncs itself, -// paced; +// paced; nor is a publish stage (release 18, `publish-*`): the publish +// lane and Publish now re-derive stages from the stamps on disk, and a +// stage re-queued with its run's preconditions would wait on a run that +// is gone; // 3. a meta queued more than REQUEUE_MAX_AGE_MS before this boot is stale; // 4. a meta with no replay spec cannot be re-queued; // 5. of what is left, only the NEWEST meta per kind + channel + params is @@ -48,7 +52,7 @@ export type RequeueFn = ( ) => Promise<{ ok: true; jobId: string } | { ok: false; error: string }>; export type CancelCategory = - "idle" | "sync" | "stale" | "no-spec" | "superseded" | "refused"; + "idle" | "sync" | "publish" | "stale" | "no-spec" | "superseded" | "refused"; export type BootQueuedResult = { requeued: { id: string; newId: string }[]; @@ -56,6 +60,8 @@ export type BootQueuedResult = { }; export const RESTART_REASON = "server restarted before it ran"; +export const PUBLISH_RESTART_REASON = + "server restarted; the publish lane re-derives stages from on-disk state"; export const REQUEUE_MAX_AGE_MS = 24 * 60 * 60 * 1000; // Key a spec by what it would DO: kind + channel + bucket + params, with the @@ -283,6 +289,8 @@ export async function settleQueuedJobMetas( "sync", "server restarted; the scheduler re-derives syncs", ); + } else if (isPublishStageJobKind(meta.kind)) { + await cancel(meta, "publish", PUBLISH_RESTART_REASON); } else if ( typeof meta.queuedAt !== "number" || opts.bootedAt - meta.queuedAt > maxAgeMs diff --git a/common/jobs/jobDetail.test.ts b/common/jobs/jobDetail.test.ts @@ -0,0 +1,42 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { jobSpecDetail } from "./jobDetail"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/jobDetail.test.ts + +test("a publish stage reads `run <last six of the run id> · <target>`", () => { + assert.equal( + jobSpecDetail("publish-build-site", { + kind: "publish-build-site", + slug: "jeralyzer", + params: { runId: "run-0000abc12-9f8e7d6c", kind: "build-site", target: "jeralyzer" }, + }), + "run 8e7d6c · jeralyzer", + ); + assert.equal( + jobSpecDetail("publish-update-index", { + kind: "publish-update-index", + slug: "_index", + params: { runId: "cli-abcdef" }, + }), + "run abcdef · _index", + ); +}); + +test("a publish stage with no run id (or no spec) adds nothing", () => { + assert.equal(jobSpecDetail("publish-deploy-hub", { kind: "publish-deploy-hub", slug: "_hub" }), undefined); + assert.equal(jobSpecDetail("publish-deploy-hub", undefined), undefined); +}); + +test("every other kind is unchanged: fetch-window keeps its phrase, the rest nothing", () => { + assert.equal( + jobSpecDetail("fetch-window", { + kind: "fetch-window", + slug: "c", + params: { requestedBy: "umtool", manifest: "m", clipId: "c1" }, + }), + "umtool · m/c1", + ); + assert.equal(jobSpecDetail("sync", { kind: "sync", slug: "c", params: { runId: "zzzzzzzz" } }), undefined); + assert.equal(jobSpecDetail(undefined, undefined), undefined); +}); diff --git a/common/jobs/jobDetail.ts b/common/jobs/jobDetail.ts @@ -7,16 +7,27 @@ import type { JobSpec } from "./jobSpec"; // for what. The answer is already on the job — the replay spec's params — so // this reads it rather than adding a second place to store it. // +// A PUBLISH STAGE (release 18, `publish-<stage>`) reads `run <last six of the +// run id> · <target>`: every stage one Publish now, one lane pass or one CLI +// command enqueued carries the same run id, so a run reads as a group on +// /jobs. The target is the spec's slug ("_index", a site id, "_hub", …). +// // Pure and total: a kind with nothing to add returns undefined and the row // renders exactly as it did. export function jobSpecDetail( kind: string | undefined, spec: JobSpec | null | undefined, ): string | undefined { - if (kind !== "fetch-window") return undefined; const p = spec?.params ?? {}; const s = (v: unknown): string | null => typeof v === "string" && v.trim() !== "" ? v.trim() : null; + if (typeof kind === "string" && kind.startsWith("publish-")) { + const runId = s(p.runId); + const target = s(spec?.slug); + if (!runId || !target) return undefined; + return `run ${runId.slice(-6)} · ${target}`; + } + if (kind !== "fetch-window") return undefined; const by = s(p.requestedBy); if (!by) return undefined; const what = [s(p.manifest), s(p.clipId)].filter(Boolean).join("/"); diff --git a/common/jobs/jobKinds.test.ts b/common/jobs/jobKinds.test.ts @@ -6,6 +6,9 @@ import { getJobKind, kindNeedsMedia, kindNeedsText, + isIngestKind, + isPublishStageJobKind, + jobKindIds, } from "./jobKinds"; // Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test jobs/jobKinds.test.ts @@ -211,3 +214,72 @@ test("the feed metadata backfill reads and writes the text tier only", () => { assert.equal(jobKindLabel("feed-metadata"), "Backfill feed metadata"); assert.equal(isDrainableKind("feed-metadata"), false); }); + +// RELEASE 18: the seven publish stages, the lane's runner, and the kinds the +// stages replace (label-only: archived metas still read). +const PUBLISH_LABELS: Record<string, string> = { + "publish-update-index": "Update the index", + "publish-build-site": "Build site", + "publish-deploy-site": "Deploy site", + "publish-build-hub": "Build hub", + "publish-deploy-hub": "Deploy hub", + "publish-build-homepage": "Build homepage", + "publish-deploy-homepage": "Deploy homepage", +}; + +test("release 18: the seven publish stages are replayable, undrainable, labelled", () => { + for (const [kind, label] of Object.entries(PUBLISH_LABELS)) { + const meta = getJobKind(kind); + assert.ok(meta, `${kind} is registered`); + assert.equal(meta.label, label); + assert.equal(meta.replayable, true, `${kind} is replayable`); + assert.equal(isDrainableKind(kind), false, `${kind} is a child process: cancel, not drain`); + assert.equal(kindNeedsMedia(kind), false); + assert.equal(isPublishStageJobKind(kind), true); + } + assert.equal(isPublishStageJobKind("publish-nothing"), false); + assert.equal(isPublishStageJobKind("build-hub"), false); + assert.equal(isPublishStageJobKind(undefined), false); + const runner = getJobKind("auto-publish"); + assert.ok(runner); + assert.equal(runner.label, "Auto-publish runner"); + assert.equal(runner.drainable, true); + assert.equal(runner.replayable, false); +}); + +test("release 18: the kinds no longer created keep their labels (seven of them none)", () => { + for (const kind of [ + "build-index", + "build-stats", + "build-export", + "build-deploy", + "build-all", + "build-deploy-all", + "deploy-export", + ]) { + assert.ok(getJobKind(kind), `${kind} is registered`); + // /jobs shows them by their raw kind, as it always has (jobs.spec reads it). + assert.equal(jobKindLabel(kind), kind); + assert.equal(getJobKind(kind)?.replayable, false); + } + assert.equal(jobKindLabel("build-hub"), "Build hub"); + assert.equal(jobKindLabel("build-deploy-homepage"), "Build & deploy homepage"); +}); + +test("release 18: every drainable kind but the lane runners is an ingest kind", () => { + const runners = new Set(["auto-transcribe", "auto-download", "auto-digest", "auto-backfill", "auto-publish"]); + for (const kind of jobKindIds()) { + if (!isDrainableKind(kind) || runners.has(kind)) continue; + assert.equal(isIngestKind(kind), true, `${kind} is drainable and per channel: an ingest kind`); + } + for (const kind of runners) assert.equal(isIngestKind(kind), false, kind); + // The lane's own unit, and the single writers, count too. + for (const kind of ["auto-download-unit", "import-one", "transcribe-one", "normalize-transcripts"]) { + assert.equal(isIngestKind(kind), true, kind); + } + // Nothing that publishes, moves or only reads. + for (const kind of ["publish-build-site", "build-export", "relocate-channel-media", "scan-media", "fetch-window"]) { + assert.equal(isIngestKind(kind), false, kind); + } + assert.equal(isIngestKind("no-such-kind"), false); +}); diff --git a/common/jobs/jobKinds.ts b/common/jobs/jobKinds.ts @@ -674,6 +674,121 @@ const JOB_KINDS: Record<string, JobKindMeta> = { replayable: false, queueKeyStrategy: "custom", }, + // THE PUBLISH STAGES (release 18): one job per stage on the `publish` queue, + // each a child process (`archilyzer stage <kind> <target>`, publish/ + // stageRun.ts) the editor spawns through runManagedCommand — so Cancel kills + // it, and there is nothing to drain. Replayable: the spec's params ARE the + // StageRequest (publish/publishStages.ts), and a replay re-asks the stage's + // preconditions on disk. The labels are the stages' own (publish/stages.ts + // STAGES[kind].label; jobKinds.test.ts holds them equal). No channelSlug and + // no `needsMedia`: a stage reads the index and the bundles, and the update- + // index child guards its own channels (the index build's hold). + "publish-update-index": { + kind: "publish-update-index", + label: "Update the index", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-build-site": { + kind: "publish-build-site", + label: "Build site", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-deploy-site": { + kind: "publish-deploy-site", + label: "Deploy site", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-build-hub": { + kind: "publish-build-hub", + label: "Build hub", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-deploy-hub": { + kind: "publish-deploy-hub", + label: "Deploy hub", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-build-homepage": { + kind: "publish-build-homepage", + label: "Build homepage", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + "publish-deploy-homepage": { + kind: "publish-deploy-homepage", + label: "Deploy homepage", + drainable: false, + replayable: true, + queueKeyStrategy: "custom", + }, + // THE PUBLISH LANE'S RUNNER (publish/publishRunner.ts): one long-lived job on + // queueKey "", like the four auto-queue runners. Drainable — a drain lets the + // stage in flight finish and dispatches no more — and never replayable: a + // loop is not a unit of work. + "auto-publish": { + kind: "auto-publish", + label: "Auto-publish runner", + drainable: true, + replayable: false, + queueKeyStrategy: "parallel", + }, + // THE KINDS RELEASE 18 NO LONGER CREATES. Their archived metas still read, + // so they stay known, labels exactly as they were: the six the hub and the + // homepage have carried since release 13 above, and these seven with NONE + // (/jobs has always shown them by their raw kind, and e2e reads that text). + "build-index": { + kind: "build-index", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "build-stats": { + kind: "build-stats", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "build-export": { + kind: "build-export", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "build-deploy": { + kind: "build-deploy", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "build-all": { + kind: "build-all", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "build-deploy-all": { + kind: "build-deploy-all", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, + "deploy-export": { + kind: "deploy-export", + drainable: false, + replayable: false, + queueKeyStrategy: "custom", + }, // A REPORT SITE'S EVIDENCE MEDIA (publish/reportMedia.ts): every clip its // published reports cite, cut from the media on disk, and every cited post // capture copied, into the site's report-media cache before its build. @@ -808,6 +923,69 @@ export function jobKindLabel(kind: string): string { return JOB_KINDS[kind]?.label ?? kind; } +// THE PUBLISH STAGES' JOB KINDS (release 18): `publish-<stage kind>`. +export function isPublishStageJobKind(kind: string | undefined): boolean { + return typeof kind === "string" && kind.startsWith("publish-") && JOB_KINDS[kind] !== undefined; +} + +// THE KINDS WHOSE `done` CAN CHANGE WHAT THE INDEX READS (release 18): what +// makes the index stale, and a site's channels "changed", to the publish +// status (publish/publishState.ts). The plan's words are "the drainable +// kinds", and every drainable kind that runs per channel is here +// (jobKinds.test.ts holds that) — plus the per-video and one-shot writers of +// the same text and posts that are not drainable: the auto-download lane's +// unit, a single import, transcription or download, the availability checks, +// the forum import, the feed backfill and the cues sweep. The lane RUNNERS are +// not: their meta is `running` for as long as the lane is, and their units +// either have their own job (auto-download-unit) or none at all +// (transcription, digest and backfill units — publishState.ts reads the +// channel's report regeneration for those). +const INGEST_KINDS: ReadonlySet<string> = new Set([ + "whisper-all", + "whisper-bucket-downloaded-no-transcript", + "whisper-bucket-auto-subs", + "purge-superseded-auto-subs", + "digest-channel-local", + "digest-channel-remote", + "digest-share-cluster", + "normalize-transcripts", + "redownload-incomplete-bucket", + "download-from-playlist", + "download-missing", + "download-missing-subs", + "import-one", + "import-archive-org", + "redownload-archive", + "retry-bucket", + "diarize-channel", + "backfill-channel", + "persist-videos", + "persist-kept", + "fetch-posts", + "import-forum-pages", + "capture-posts", + "check-post-availability", + "sync", + "metadata-scan", + "feed-metadata", + "auto-download-unit", + "whisper-video", + "transcribe-one", + "download-one-pipeline", + "check-availability", + "quick-availability-check", + "check-maybe-missing", +]); + +export function isIngestKind(kind: string | undefined): boolean { + return typeof kind === "string" && INGEST_KINDS.has(kind); +} + +// Every registered kind (for the tests that hold one table to another). +export function jobKindIds(): string[] { + return Object.keys(JOB_KINDS); +} + export function isDrainableKind(kind: string): boolean { return JOB_KINDS[kind]?.drainable ?? false; } diff --git a/common/lib/autoQueueTypes.ts b/common/lib/autoQueueTypes.ts @@ -208,3 +208,14 @@ export const LANES = [ ] as const; export type AutoQueueKind = (typeof LANES)[number]; + +// THE PIPELINE LANES (release 18): lanes that dispatch STAGES of a pipeline, +// not work picked per channel — today the one publish lane. A pipeline lane +// has a runner, a pause gate (lib/pauseGates.ts) and a console, but NO rule +// tree: it is deliberately NOT in LANES, because nine loops compile a channel +// tree per entry of LANES (channelPriority, channelWriters, the storage watch, +// operationBatch, the channel snapshot, the auto-queue schema …). Its +// settings are `settings.publish`, not an `autoQueue` policy. +export const PIPELINE_LANES = ["publish"] as const; + +export type PipelineLane = (typeof PIPELINE_LANES)[number]; diff --git a/common/lib/fileSchemaDocs.ts b/common/lib/fileSchemaDocs.ts @@ -27,6 +27,7 @@ import { RELATED_SITE_GROUP_FIELD_DOCS, SITE_CHANNEL_MEMBERSHIP_FIELD_DOCS, SITE_FIELD_DOCS, + SITE_PUBLISH_FIELD_DOCS, parseSite, type Site, } from "./siteSchema"; @@ -50,6 +51,7 @@ const SITE_NESTED: Partial<Record<keyof Site, KeyTable[]>> = { groups: [{ path: "groups[]", docs: CHANNEL_GROUP_FIELD_DOCS }], channels: [{ path: "channels[]", docs: SITE_CHANNEL_MEMBERSHIP_FIELD_DOCS }], relatedSites: [{ path: "relatedSites[]", docs: RELATED_SITE_GROUP_FIELD_DOCS }], + publish: [{ path: "publish", docs: SITE_PUBLISH_FIELD_DOCS }], }; export function renderSiteMarkdown(): string { @@ -78,8 +80,9 @@ export function renderSiteMarkdown(): string { "`homeTagline`, `groups`, `defaultGroupId` and `channels`; every other " + "key is written only when it differs from its default (`socialLinks` " + "whenever it is an array, even an empty one). A save is REFUSED when " + - "there is no channel group, when `defaultGroupId` names no group, or " + - "when a social link's SVG is not safe to inline.", + "there is no channel group, when `defaultGroupId` names no group, " + + "when a social link's SVG is not safe to inline, or when a public " + + "site's `publish.auto` deploys and it has no `cloudflareProject`.", ); out.push(""); out.push(REGENERATE); diff --git a/common/lib/pauseGates.test.ts b/common/lib/pauseGates.test.ts @@ -136,7 +136,7 @@ test("the backfill lane ships held, as its inverted field always made it", () => assert.equal(base.autoQueue.backfill.enabled, false, "and unarmed as well"); for (const lane of LANES) { if (lane === "backfill") continue; - assert.equal(base.autoQueue[lane].held, false, `${lane} ships free`); + assert.equal(isGateHeld(base, lane), false, `${lane} ships free`); } }); @@ -232,3 +232,23 @@ test("pauseLaneFor answers for every catalog id", () => { test("an id the catalog does not know has no lane", () => { assert.equal(pauseLaneFor("not-an-operation"), null); }); + +// RELEASE 18: the publish lane is a PIPELINE lane, not a LANES entry, and its +// gate is `settings.publish.held` — the one gate not on an `autoQueue` policy. +test("the publish gate round-trips through settings.publish.held, and touches no other gate", () => { + const base = defaultSiteSettings(); + assert.equal(isGateHeld(base, "publish"), false, "the publish lane ships free"); + const held = withGateHeld(base, "publish", true); + assert.equal(held.publish.held, true); + assert.equal(isGateHeld(held, "publish"), true); + // The rest of the block is kept (a spread, never a rebuilt literal). + assert.deepEqual({ ...held.publish, held: false }, base.publish); + assert.equal("publish" in held.autoQueue, false, "no autoQueue.publish is invented"); + for (const lane of LANES) { + assert.equal(isGateHeld(held, lane), isGateHeld(base, lane), `holding publish moved ${lane}`); + assert.equal(isGateHeld(withGateHeld(base, lane, true), "publish"), false, `holding ${lane} held publish`); + } + assert.equal(isGateHeld(withGateHeld(held, "publish", false), "publish"), false); + // A partial object (laneGuards.test casts one) answers free, never throws. + assert.equal(isGateHeld({} as SiteSettings, "publish"), false); +}); diff --git a/common/lib/pauseGates.ts b/common/lib/pauseGates.ts @@ -1,5 +1,5 @@ import type { SiteSettings } from "./settings"; -import type { AutoQueueKind } from "./autoQueueTypes"; +import type { AutoQueueKind, PipelineLane } from "./autoQueueTypes"; import { operationCatalog } from "./operations"; import { BACKFILL_QUEUE, @@ -46,7 +46,11 @@ import { // closed that gap, so this is now an alias and the two id spaces cannot drift. // The name survives because thirteen call sites read as "which lane's gate", // and because a gate is what this file is about. -export type PauseLane = AutoQueueKind; +// +// WIDENED BY ONE in release 18: the publish lane is a PIPELINE lane +// (autoQueueTypes.ts PIPELINE_LANES), not a LANES entry, and its gate is +// `settings.publish.held` — the one gate that is not on an `autoQueue` policy. +export type PauseLane = AutoQueueKind | PipelineLane; // WHICH LANE'S GATE HOLDS THIS OPERATION — the answer to "the operator is on // /operations/diarization and wants to hold it". @@ -65,7 +69,7 @@ export type PauseLane = AutoQueueKind; // OPERATION_BY_ID, which knows only the four registry entries. Sync is in that // walk and its answer is null: the scheduler's `enabled` is its own switch, not // a lane gate. -export function pauseLaneFor(operationId: string): PauseLane | null { +export function pauseLaneFor(operationId: string): AutoQueueKind | null { const op = operationCatalog().find((o) => o.id === operationId); if (!op) return null; if (op.runner) return op.runner; @@ -95,6 +99,9 @@ export function pauseLaneFor(operationId: string): PauseLane | null { // `autoQueue` is read defensively: laneGuards.test.ts casts a partial object to // SiteSettings, and this must answer for it the way it always has. export function isGateHeld(settings: SiteSettings, lane: PauseLane): boolean { + // The publish lane's gate is its own block's (release 18): there is no + // `autoQueue.publish`. + if (lane === "publish") return settings.publish?.held === true; return settings.autoQueue?.[lane]?.held === true; } @@ -115,6 +122,9 @@ export function withGateHeld( lane: PauseLane, held: boolean, ): SiteSettings { + if (lane === "publish") { + return { ...settings, publish: { ...settings.publish, held } }; + } return { ...settings, autoQueue: { diff --git a/common/lib/settingsDocs.ts b/common/lib/settingsDocs.ts @@ -18,6 +18,7 @@ import { BUILD_PIPELINE_SETTINGS_FIELD_DOCS, DIARIZATION_SETTINGS_FIELD_DOCS, PACING_SETTINGS_FIELD_DOCS, + PUBLISH_SETTINGS_FIELD_DOCS, DIGEST_SETTINGS_FIELD_DOCS, SAVED_VIDEO_BACKUP_SETTINGS_FIELD_DOCS, SOCIAL_LINK_FIELD_DOCS, @@ -244,6 +245,13 @@ export function blockTables(d: SiteSettings): Partial<Record<keyof SiteSettings, defaults: fromObject(d.archiveOrg), }, ], + publish: [ + { + path: "publish", + docs: PUBLISH_SETTINGS_FIELD_DOCS, + defaults: fromObject(d.publish), + }, + ], }; } diff --git a/common/lib/settingsSchema.test.ts b/common/lib/settingsSchema.test.ts @@ -37,6 +37,7 @@ import type { DiarizationSettings, DigestSettings, PacingSettingsBlock, + PublishSettings, ReportDebouncePreset, SavedVideoBackupSettings, SocialLink, @@ -100,13 +101,15 @@ type PreSchemaSiteSettings = { attribution: AttributionSettings; // How archive.org files are fetched: BitTorrent with seeding, else direct. archiveOrg: ArchiveOrgFetchSettings; + // The publish lane (release 18). + publish: PublishSettings; }; // Bracketed so the conditional does not distribute (see commit 8c43231). type Same<A, B> = [A] extends [B] ? ([B] extends [A] ? true : false) : false; const shapeUnchanged: Same<SiteSettings, PreSchemaSiteSettings> = true; -test("SiteSettings keeps its 35 fields, in file order", () => { +test("SiteSettings keeps its 36 fields, in file order", () => { assert.equal(shapeUnchanged, true); assert.deepEqual(Object.keys(siteSettingsSchema.shape), [ "adminTitle", @@ -144,6 +147,7 @@ test("SiteSettings keeps its 35 fields, in file order", () => { "backfill", "attribution", "archiveOrg", + "publish", ]); // A parsed object carries every key, in that order — writeSettings writes // exactly this, so the order is the on-disk order. @@ -489,3 +493,76 @@ test("the retired sweep fields migrate onto a lane ONLY when the lane is absent" assert.equal(withFile("{}").autoQueue.digest.enabled, false); assert.equal(withFile("{}").autoQueue.backfill.enabled, false); }); + +// RELEASE 18: the publish lane's block. +test("publish: defaults, and every field sanitized", () => { + const { sanitizePublish, defaultPublish } = S; + const d = defaults().publish; + assert.deepEqual(d, { + enabled: false, + held: false, + checkEveryMinutes: 10, + refreshEveryMinutes: 360, + quietHours: null, + runner: "local", + previewBranch: "preview", + hub: "off", + homepage: "off", + }); + assert.deepEqual(defaultPublish(), d); + for (const raw of [undefined, null, 3, "x", []]) assert.deepEqual(sanitizePublish(raw), d); + const good = sanitizePublish({ + enabled: true, + held: true, + checkEveryMinutes: 5, + refreshEveryMinutes: 0, + quietHours: { start: 22, end: 6 }, + runner: "docker", + previewBranch: "r18", + hub: "preview", + homepage: "production", + }); + assert.deepEqual(good, { + enabled: true, + held: true, + checkEveryMinutes: 5, + refreshEveryMinutes: 0, + quietHours: { start: 22, end: 6 }, + runner: "docker", + previewBranch: "r18", + hub: "preview", + homepage: "production", + }); + const bad = sanitizePublish({ + enabled: "yes", + held: 1, + checkEveryMinutes: 0, + refreshEveryMinutes: 10 ** 9, + quietHours: { start: 3, end: 3 }, + runner: "kubernetes", + previewBranch: "main", + hub: "always", + homepage: null, + }); + assert.deepEqual(bad, { + ...d, + checkEveryMinutes: 1, + refreshEveryMinutes: 43200, + }); + assert.equal(sanitizePublish({ checkEveryMinutes: 99999 }).checkEveryMinutes, 1440); + assert.equal(sanitizePublish({ quietHours: { start: 24, end: 6 } }).quietHours, null); + assert.equal(sanitizePublish({ quietHours: { start: 1 } }).quietHours, null); + assert.equal(sanitizePublish({ previewBranch: "Has Spaces" }).previewBranch, "preview"); + assert.equal(sanitizePublish({ previewBranch: " r18 " }).previewBranch, "r18"); +}); + +test("publish: an unknown key inside the block is dropped, and the block round-trips the schema", () => { + const parsed = siteSettingsSchema.parse({ + publish: { enabled: true, held: true, extra: 1, previewBranch: "smoke" }, + }); + assert.equal(parsed.publish.enabled, true); + assert.equal(parsed.publish.held, true); + assert.equal(parsed.publish.previewBranch, "smoke"); + assert.equal("extra" in parsed.publish, false); + assert.deepEqual(siteSettingsSchema.parse(JSON.parse(JSON.stringify(parsed))).publish, parsed.publish); +}); diff --git a/common/lib/settingsSchema.ts b/common/lib/settingsSchema.ts @@ -96,6 +96,7 @@ import { type DigestTimestampMode, } from "./digest"; import type { FieldDocs } from "./fieldDocs"; +import { previewBranchProblem } from "./pagesDeploy"; import { sanitizeSocial, type SocialSettings, @@ -1563,6 +1564,115 @@ export function clampPageBytes(value: unknown): number { } +// --- The publish lane (release 18) ------------------------------------------ + +// What the lane (and Publish now) may do to one target when it is stale: +// nothing, build it, build it and deploy it as a preview, or build it and +// deploy it to production. A site carries its own in site.json +// (`publish.auto`, lib/siteSchema.ts); the hub and the homepage carry theirs +// here. A manual Build or Deploy button never asks it. +export type PublishPolicy = "off" | "build" | "preview" | "production"; + +export const PUBLISH_POLICIES: readonly PublishPolicy[] = ["off", "build", "preview", "production"]; + +export function isPublishPolicy(v: unknown): v is PublishPolicy { + return v === "off" || v === "build" || v === "preview" || v === "production"; +} + +export type PublishQuietHours = { start: number; end: number }; + +// Each field is documented in PUBLISH_SETTINGS_FIELD_DOCS below. +export type PublishSettings = { + enabled: boolean; + held: boolean; + checkEveryMinutes: number; + refreshEveryMinutes: number; + quietHours: PublishQuietHours | null; + runner: "local" | "docker"; + previewBranch: string; + hub: PublishPolicy; + homepage: PublishPolicy; +}; + +export const PUBLISH_SETTINGS_FIELD_DOCS: FieldDocs<PublishSettings> = { + enabled: + "Whether the publish lane's runner runs (`auto-publish` on /jobs). Off by default. Turning it on starts nothing by itself until the index is stale (or there is no index stamp yet) — see `refreshEveryMinutes`.", + held: + "The lane's pause gate (lib/pauseGates.ts, lane `publish`). A hold stops the runner DISPATCHING: the stage in flight finishes, no next one starts. It never kills a stage.", + checkEveryMinutes: + "How often (minutes) the runner wakes to ask whether a pass is due. Clamped to [1, 1440]; default 10.", + refreshEveryMinutes: + "The least time (minutes) between two index updates the lane starts: a pass runs when the index is stale and the last update is at least this old, or when there is no index stamp. Clamped to [0, 43200]; default 360. 0 = whenever the index is stale.", + quietHours: + "A local-clock window in which the lane starts no pass, `{ \"start\": 22, \"end\": 6 }` (hours [0,23], `[start, end)`, may wrap midnight), or null (default). A pass already running finishes its stage and then waits.", + runner: + "Which build runner the lane's builds ask for: \"local\" (default; each site in turn, as a child of the editor) or \"docker\" (every stale site in containers — a host with a container engine only; in a container it is refused). See PUBLISH.md.", + previewBranch: + "The Pages preview branch a `preview` policy deploys to (`wrangler pages deploy --branch <previewBranch>`). Lowercase letters, digits and dashes, never `main`; an invalid name reads as the default, \"preview\".", + hub: + "The hub's policy: \"off\" (default), \"build\", \"preview\" or \"production\". The hub is built when the index or the listed sites changed, and deployed as the policy says — to production only with a Pages project in homepage.json.", + homepage: + "The homepage's policy, as `hub` (default \"off\"). Building the homepage also publishes the source mirror (the operator's scrub and denylist files must exist).", +}; + +export const PUBLISH_CHECK_EVERY_DEFAULT_MINUTES = 10; +export const PUBLISH_CHECK_EVERY_MAX_MINUTES = 1440; +export const PUBLISH_REFRESH_EVERY_DEFAULT_MINUTES = 360; +export const PUBLISH_REFRESH_EVERY_MAX_MINUTES = 43200; +export const PUBLISH_DEFAULT_PREVIEW_BRANCH = "preview"; + +export function defaultPublish(): PublishSettings { + return { + enabled: false, + held: false, + checkEveryMinutes: PUBLISH_CHECK_EVERY_DEFAULT_MINUTES, + refreshEveryMinutes: PUBLISH_REFRESH_EVERY_DEFAULT_MINUTES, + quietHours: null, + runner: "local", + previewBranch: PUBLISH_DEFAULT_PREVIEW_BRANCH, + hub: "off", + homepage: "off", + }; +} + +function sanitizeQuietHours(value: unknown): PublishQuietHours | null { + if (!value || typeof value !== "object" || Array.isArray(value)) return null; + const r = value as Record<string, unknown>; + const start = clampHourOrNull(r.start); + const end = clampHourOrNull(r.end); + // Both valid hours and a non-empty window, or no window at all. + if (start === null || end === null || start === end) return null; + return { start, end }; +} + +// Coerce a raw settings.publish into a clean PublishSettings: every field its +// default when missing or ill-typed, numbers clamped, the two switches true +// only when exactly true, a bad preview name the default. +export function sanitizePublish(value: unknown): PublishSettings { + const d = defaultPublish(); + if (!value || typeof value !== "object" || Array.isArray(value)) return d; + const r = value as Record<string, unknown>; + const previewBranch = + typeof r.previewBranch === "string" && previewBranchProblem(r.previewBranch) === null + ? r.previewBranch.trim() + : d.previewBranch; + return { + enabled: r.enabled === true, + held: r.held === true, + checkEveryMinutes: clampPositiveInt(r.checkEveryMinutes, d.checkEveryMinutes, PUBLISH_CHECK_EVERY_MAX_MINUTES), + refreshEveryMinutes: clampIntAllowZero( + r.refreshEveryMinutes, + d.refreshEveryMinutes, + PUBLISH_REFRESH_EVERY_MAX_MINUTES, + ), + quietHours: sanitizeQuietHours(r.quietHours), + runner: r.runner === "docker" ? "docker" : "local", + previewBranch, + hub: isPublishPolicy(r.hub) ? r.hub : d.hub, + homepage: isPublishPolicy(r.homepage) ? r.homepage : d.homepage, + }; +} + // Coerce a raw settings.transcriptionApps value into a clean keyed map of // AppInstanceConfig, dropping unknown/ill-typed fields. export function sanitizeTranscriptionApps( @@ -1705,6 +1815,9 @@ export const siteSettingsSchema = z.object({ archiveOrg: settingsField((v): ArchiveOrgFetchSettings => sanitizeArchiveOrg(v)).describe( "How archive.org files are fetched (controller/archiveOrgDownload.ts). Over BitTorrent with aria2c when the item's torrent carries the file — archive.org is the torrent's web seed, so the swarm takes load off archive.org — then seeded for a while; otherwise, or when the torrent stalls, a direct download from archive.org. Either way the file is verified against archive.org's sha1/md5. No yt-dlp.", ), + publish: settingsField((v): PublishSettings => sanitizePublish(v)).describe( + "The publish LANE (release 18): a runner that, when the index is stale, updates it and then builds — and, where a site's own `publish.auto` (site.json) says so, deploys — what changed, one stage at a time on the `publish` queue. OFF by default; the manual stages (Publish now, Build, Deploy) work either way. See PublishSettings and common/publish/publishRunner.ts.", + ), }); export type SiteSettings = z.infer<typeof siteSettingsSchema>; diff --git a/common/lib/site.ts b/common/lib/site.ts @@ -18,6 +18,7 @@ import { isValidSiteId, parseSite, parseSiteUrl, + sitePublishProblem, siteToDisk, type Site, } from "./siteSchema"; @@ -245,6 +246,10 @@ export async function writeSite( `Default group "${String(site.defaultGroupId)}" is not in the configured groups`, ); } + // A public site whose publish policy deploys needs a Pages project (release + // 18). A private one is clamped to "build" by siteToDisk, not refused. + const publishProblem = sitePublishProblem(site); + if (publishProblem) throw new Error(publishProblem); // undefined socialLinks = inherit the global default; only validate/persist a // key when the site explicitly overrides (an array, even empty). // A link whose SVG is unchanged from the file on disk is kept as it is; a new diff --git a/common/lib/siteSchema.test.ts b/common/lib/siteSchema.test.ts @@ -10,11 +10,14 @@ import { SITE_FIELD_DOCS, SITE_KEYS, channelsOnlyOnUnlistedSites, + clampSitePublishPolicy, isCitedSite, isListedSite, parseSite, parseSiteReports, siteFieldsSchema, + sitePublishPolicy, + sitePublishProblem, siteToDisk, type Site, } from "./siteSchema"; @@ -336,7 +339,9 @@ test("legacy publish: \"cited\" reads as search off and is rewritten as search: const legacy = parseSite("s", { publish: "cited" }); assert.equal(legacy.search, false); assert.equal(isCitedSite(legacy), true); - assert.equal("publish" in legacy, false); + // `publish` is a key again since release 18 (the lane's policy), and the + // legacy string is not one: it reads as absent and is never written back. + assert.equal(legacy.publish, undefined); const disk = siteToDisk(legacy); assert.equal(disk.search, false); assert.equal("publish" in disk, false); @@ -472,3 +477,58 @@ test("patchSite applies a patch to the site on disk now, keeps every other key, assert.equal(getSite("s", paths).siteId, "s"); assert.equal(fs.existsSync(siteConfigFile(paths, "other")), false); }); + +// RELEASE 18: the per-site publish policy. +test("publish.auto: off by default, the four values read, anything else is off", () => { + assert.equal(parseSite("s", {}).publish, undefined); + assert.equal(sitePublishPolicy(parseSite("s", {})), "off"); + const withProject = { cloudflareProject: "proj" }; + for (const auto of ["build", "preview", "production"] as const) { + assert.deepEqual(parseSite("s", { ...withProject, publish: { auto } }).publish, { auto }); + } + for (const bad of [{ auto: "off" }, { auto: "always" }, { auto: 3 }, "production", [], null, {}]) { + assert.equal(parseSite("s", { ...withProject, publish: bad }).publish, undefined, JSON.stringify(bad)); + } + // The legacy report-only switch is a string: it is not a policy. + const legacy = parseSite("s", { publish: "cited" }); + assert.equal(legacy.publish, undefined); + assert.equal(legacy.search, false); +}); + +test("publish.auto: a private site at most builds; no Pages project, no deploy policy", () => { + const priv = parseSite("s", { audience: "private", cloudflareProject: "p", publish: { auto: "production" } }); + assert.deepEqual(priv.publish, { auto: "build" }); + const noProject = parseSite("s", { publish: { auto: "preview" } }); + assert.deepEqual(noProject.publish, { auto: "build" }); + assert.deepEqual(parseSite("s", { publish: { auto: "build" } }).publish, { auto: "build" }); + assert.equal(clampSitePublishPolicy({ audience: "private" }, "preview"), "build"); + assert.equal(clampSitePublishPolicy({ cloudflareProject: " " }, "production"), "build"); + assert.equal(clampSitePublishPolicy({ cloudflareProject: "p" }, "production"), "production"); + assert.equal(clampSitePublishPolicy({}, "off"), "off"); +}); + +test("publish.auto: only a policy other than off is written, clamped", () => { + const base = parseSite("s", { cloudflareProject: "p" }); + assert.equal("publish" in siteToDisk(base), false); + assert.equal("publish" in siteToDisk({ ...base, publish: { auto: "off" } }), false); + assert.deepEqual(siteToDisk({ ...base, publish: { auto: "preview" } }).publish, { auto: "preview" }); + assert.deepEqual( + siteToDisk({ ...base, audience: "private", publish: { auto: "production" } }).publish, + { auto: "build" }, + ); +}); + +test("writeSite refuses a public deploy policy with no cloudflareProject, naming the key", async () => { + const paths = scratchPaths(await mkdtemp(path.join(os.tmpdir(), "site-"))); + const base = parseSite("s", {}); + await assert.rejects(writeSite({ ...base, publish: { auto: "production" } }, paths), /cloudflareProject/); + assert.equal(sitePublishProblem({ publish: { auto: "preview" } })?.includes("publish.auto"), true); + assert.equal(sitePublishProblem({ publish: { auto: "build" } }), null); + assert.equal(sitePublishProblem({ audience: "private", publish: { auto: "production" } }), null); + assert.equal(fs.existsSync(siteConfigFile(paths, "s")), false); + // A private site is clamped, not refused; a public one with a project saves. + await writeSite({ ...base, audience: "private", publish: { auto: "production" } }, paths); + assert.deepEqual(getSite("s", paths).publish, { auto: "build" }); + await writeSite({ ...base, cloudflareProject: "proj", publish: { auto: "production" } }, paths); + assert.deepEqual(getSite("s", paths).publish, { auto: "production" }); +}); diff --git a/common/lib/siteSchema.ts b/common/lib/siteSchema.ts @@ -38,7 +38,12 @@ import { } from "./channelGroups"; import { parseAccentSetting } from "./accent"; import { wordmarkLeadFor } from "./brand"; -import { parseSocialLinks, type SocialLink } from "./settingsSchema"; +import { + isPublishPolicy, + parseSocialLinks, + type PublishPolicy, + type SocialLink, +} from "./settingsSchema"; import { settingsField } from "./settingsFieldSchemas"; import type { FieldDocs } from "./fieldDocs"; import { REPORT_ID_RE, isReportId } from "./report/schema"; @@ -157,6 +162,7 @@ export type Site = { transcriptDownloads?: boolean; archiveMaxBytes?: number; hubUrl?: string; + publish?: SitePublish; }; export const SITE_FIELD_DOCS: FieldDocs<Site> = { @@ -204,8 +210,56 @@ export const SITE_FIELD_DOCS: FieldDocs<Site> = { "Per-site served-file size cap in bytes: any archive larger is dropped from what is served and flagged in the manifest, so a capped host (Cloudflare Pages: 25 MB) will not reject the deploy. 0 = no cap. Absent = the global default. Negative or non-numeric values are dropped.", hubUrl: "Per-site override for the hub this site belongs under. Absent = the family default, `settings.json` `homepageUrl`. Published on the public `/site.json` and `/corpus.json` so a hub can tell member sites from arbitrary added origins; the header does not link to it (release 14).", + publish: + "What the publish lane — and Publish now — may do with this site when it is stale (release 18): `{ \"auto\": \"off\" | \"build\" | \"preview\" | \"production\" }`. Absent = `off`: the lane leaves the site alone (a manual Build or Deploy still works, and \"Build all stale\" still builds it). `build` rebuilds its bundle; `preview` also deploys it to the Pages preview branch `settings.json` `publish.previewBranch` names; `production` deploys it to production. A private site is clamped to `build` (it is never deployed); `preview` and `production` need a `cloudflareProject` — a save without one is refused, and a file that says so anyway reads as `build`. Only a policy other than `off` is written.", +}; + +// Each field is documented in SITE_PUBLISH_FIELD_DOCS below. +export type SitePublish = { auto: PublishPolicy }; + +export const SITE_PUBLISH_FIELD_DOCS: FieldDocs<SitePublish> = { + auto: + "The policy: `off` (default), `build`, `preview` or `production` — see `publish` above.", }; +// The policy a site may have, given who it is for and where it can go: a +// private site at most builds, and a site with no Pages project cannot be +// deployed to one. Pure; parseSite and siteToDisk both apply it. +export function clampSitePublishPolicy( + site: Pick<Site, "audience" | "cloudflareProject">, + policy: PublishPolicy, +): PublishPolicy { + if (policy !== "preview" && policy !== "production") return policy; + if (isPrivateSite(site)) return "build"; + if (!site.cloudflareProject?.trim()) return "build"; + return policy; +} + +// THE ONE READER of a site's publish policy: absent is "off". +export function sitePublishPolicy(site: Pick<Site, "publish">): PublishPolicy { + return site.publish?.auto ?? "off"; +} + +// Why a save of this site's publish policy is refused, or null. A public site +// that asks the lane to deploy needs somewhere to deploy to; the sentence +// names the key. A private site is not refused: it is clamped to `build`. +export function sitePublishProblem( + site: Pick<Site, "audience" | "cloudflareProject" | "publish">, +): string | null { + const policy = site.publish?.auto; + if (policy !== "preview" && policy !== "production") return null; + if (isPrivateSite(site)) return null; + if (site.cloudflareProject?.trim()) return null; + return `publish.auto "${policy}" deploys the site, and it has no cloudflareProject — set the Pages project, or choose "build"`; +} + +function parseSitePublish(v: unknown): SitePublish | undefined { + if (!v || typeof v !== "object" || Array.isArray(v)) return undefined; + const auto = (v as Record<string, unknown>).auto; + if (!isPublishPolicy(auto) || auto === "off") return undefined; + return { auto }; +} + // siteId shares the group-id grammar: lowercase slug, used as a directory name. export const SITE_ID_RE = /^[a-z0-9][a-z0-9-]*$/; @@ -376,6 +430,9 @@ export const siteFieldsSchema = z.object({ ), archiveMaxBytes: settingsField(archiveMaxBytesOf).describe(d.archiveMaxBytes), hubUrl: settingsField(parseSiteUrl).describe(d.hubUrl), + // Clamped against `audience` and `cloudflareProject` in the object step + // below. The legacy `publish: "cited"` (a string) reads as absent here. + publish: settingsField(parseSitePublish).describe(d.publish), }); // The whole schema: the per-key object, then the three sibling-dependent @@ -404,6 +461,8 @@ export const siteSchema = z.preprocess(migrateLegacyPublish, siteFieldsSchema).t } return c; }), + // A private site at most builds; no Pages project, no deploy policy. + publish: s.publish ? { auto: clampSitePublishPolicy(s, s.publish.auto) } : undefined, }; const out: Partial<Record<keyof Site, unknown>> = {}; for (const key of SITE_KEYS) out[key] = resolved[key]; @@ -438,6 +497,8 @@ export function siteToDisk(site: Site): Site { const siteUrl = parseSiteUrl(site.siteUrl); const hubUrl = parseSiteUrl(site.hubUrl); const archiveMaxBytes = archiveMaxBytesOf(site.archiveMaxBytes); + const parsedPublish = parseSitePublish(site.publish); + const publishPolicy = parsedPublish ? clampSitePublishPolicy(site, parsedPublish.auto) : "off"; return { siteId: site.siteId, siteTitle: site.siteTitle, @@ -471,6 +532,8 @@ export function siteToDisk(site: Site): Site { ...(site.transcriptDownloads === false ? { transcriptDownloads: false } : {}), ...(archiveMaxBytes !== undefined ? { archiveMaxBytes } : {}), ...(hubUrl ? { hubUrl } : {}), + // Off is the default: only another policy is persisted (clamped). + ...(publishPolicy !== "off" ? { publish: { auto: publishPolicy } } : {}), }; } 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/publishPlan.test.ts b/common/publish/publishPlan.test.ts @@ -0,0 +1,543 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { readFile } from "node:fs/promises"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { defaultPublish, type PublishSettings } from "../lib/settingsSchema"; +import { + buildPublishStatus, + planPublishRun, + publishPassDecision, + stepKey, + type PublishInputs, + type PublishSiteInput, + type PublishStatus, +} from "./publishPlan"; +import { STAGES } from "./stages"; +import type { BuiltStamp, DeployedFile, IndexStamp } from "./stamps"; + +// Run with: pnpm --filter yt-dlp-transcript-common exec tsx --test publish/publishPlan.test.ts +// +// The publish status and plan (release 18) over hand-built inputs: every chip +// state the plan's tables name, the policies, the preconditions a run carries, +// and the lane's pass decision. Pure — no disk, no clock. + +const HERE = path.dirname(fileURLToPath(import.meta.url)); +const MIN = 60_000; +const NOW = 2_000_000_000_000; + +function stamp(over: Partial<IndexStamp> = {}): IndexStamp { + return { + v: 1, + stampId: "stamp-1", + generation: 7, + scannedAt: NOW - 60 * MIN, + builtAt: NOW - 55 * MIN, + templatesAt: NOW - 55 * MIN, + commit: "c1", + index: { shortCircuited: false, added: 0, changed: 0, removed: 0, heldChannels: [] }, + stats: { shortCircuited: false, notIndexedYet: 0, notIndexable: 0 }, + sites: { + alpha: { siteFp: null, statsFp: null, inputSig: "sig-alpha" }, + beta: { siteFp: null, statsFp: null, inputSig: "sig-beta" }, + }, + hubSig: "hub-sig", + ...over, + }; +} + +function built(target: string, over: Partial<BuiltStamp> = {}): BuiltStamp { + return { + v: 1, + stampId: `built-${target}`, + target, + kind: target === "_hub" ? "hub" : target === "_homepage" ? "homepage" : "site", + indexStampId: "stamp-1", + inputSig: target === "_hub" ? "hub-sig" : `sig-${target}`, + builtAt: NOW - 50 * MIN, + commit: "c1", + branch: "main", + runner: "local", + audience: "public", + corpusGeneratedAt: null, + files: 1, + bytes: 1, + archivesStaged: 0, + ...over, + }; +} + +function deployedTo(target: string, builtStampId: string, kind: "production" | "preview" = "production"): DeployedFile { + const rec = { + builtStampId, + builtAt: NOW - 50 * MIN, + kind, + url: `https://${target}.example`, + at: NOW - 40 * MIN, + liveCheck: null, + ...(kind === "preview" ? { branch: "preview", alias: `https://preview.${target}.pages.dev` } : {}), + }; + return kind === "production" + ? { v: 1, target, production: rec, previews: {} } + : { v: 1, target, previews: { preview: rec } }; +} + +function site(siteId: string, over: Partial<PublishSiteInput> = {}): PublishSiteInput { + return { + siteId, + title: siteId.toUpperCase(), + private: false, + listed: true, + members: [`${siteId}-ch`, "shared-ch"], + built: built(siteId), + deployed: deployedTo(siteId, `built-${siteId}`), + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off", + cloudflareProject: siteId, + url: `https://${siteId}.pages.dev`, + ...over, + }; +} + +function inputs(over: Partial<PublishInputs> = {}, settings: Partial<PublishSettings> = {}): PublishInputs { + return { + index: { stamp: stamp(), lastIngestDoneAt: null, configChangedAt: null }, + ingestByChannel: {}, + commit: "c1", + sites: [site("alpha"), site("beta")], + hub: { + built: built("_hub"), + deployed: deployedTo("_hub", "built-_hub"), + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off", + cloudflareProject: "hub", + url: null, + }, + homepage: { + built: built("_homepage"), + deployed: deployedTo("_homepage", "built-_homepage"), + bundleProblem: null, + deployProblem: null, + pagesProblem: null, + configChangedAt: null, + policy: "off", + cloudflareProject: "archilyzer", + url: null, + mainHead: null, + }, + settings: { ...defaultPublish(), ...settings }, + jobs: [], + ended: [], + lane: { + known: true, + running: false, + jobId: null, + passRunning: false, + lastCheckAt: null, + nextCheckAt: null, + lastPassAt: null, + lastPassSummary: null, + lastDecision: null, + }, + ...over, + }; +} + +const status = (i: PublishInputs): PublishStatus => buildPublishStatus(i, NOW); +const siteOf = (s: PublishStatus, id: string) => s.sites.find((t) => t.target === id)!; +const kinds = (s: PublishStatus["plan"]) => s.steps.map((x) => `${x.kind} ${x.target}`); + +// --- purity ------------------------------------------------------------------ + +test("the builder and the planner are pure: no fs, no clock, no singleton", async () => { + const src = await readFile(path.join(HERE, "publishPlan.ts"), "utf8"); + for (const banned of ["Date.now(", "getSettings(", "getPaths(", "getRegistry(", "getScheduler(", "readJobMeta("]) { + assert.equal(src.includes(banned), false, `publishPlan.ts calls ${banned}`); + } + const specs = [...src.matchAll(/from\s+"([^"]+)"/g)].map((m) => m[1]); + for (const spec of specs) { + assert.equal(spec.startsWith("node:"), false, `publishPlan.ts imports ${spec}`); + assert.ok(spec.startsWith("."), `publishPlan.ts imports a package (${spec})`); + } + // The same inputs give the same status. + assert.deepEqual(status(inputs()), status(inputs())); +}); + +// --- the index chip ------------------------------------------------------------ + +test("index: no stamp is stale and plans the update; every build is then after it", () => { + const s = status(inputs({ index: { stamp: null, lastIngestDoneAt: null, configChangedAt: null } }, {})); + assert.equal(s.index.fresh, false); + assert.equal(s.index.chip.text, "no index yet"); + assert.deepEqual(kinds(s.plan), ["update-index _index"]); + // With a site on "build", it is built after the index's first update. + const i = inputs({ index: { stamp: null, lastIngestDoneAt: null, configChangedAt: null } }); + i.sites[0].policy = "build"; + const p = status(i).plan; + assert.deepEqual(kinds(p), ["update-index _index", "build-site alpha"]); + assert.equal(p.steps[1].indexAfter, NOW); + // The chip a build shows with no stamp is the stage's own "blocked". + assert.equal(siteOf(status(i), "alpha").chips.built.text, "update the index first"); +}); + +test("index: fresh, stale by new data (after scannedAt), stale by a config file", () => { + assert.match(status(inputs()).index.chip.text, /^fresh · updated 55 min ago$/); + const byData = status(inputs({ index: { stamp: stamp(), lastIngestDoneAt: NOW - 30 * MIN, configChangedAt: null } })); + assert.equal(byData.index.fresh, false); + assert.equal(byData.index.chip.text, "stale: new data since the last index"); + // An ingest that ended BEFORE the scan began is in the index. + assert.equal( + status(inputs({ index: { stamp: stamp(), lastIngestDoneAt: NOW - 61 * MIN, configChangedAt: null } })).index.fresh, + true, + ); + const byConfig = status(inputs({ index: { stamp: stamp(), lastIngestDoneAt: null, configChangedAt: NOW - MIN } })); + assert.equal(byConfig.index.chip.text, "stale: a config file changed since the last index"); +}); + +test("index: a running or queued update reads busy", () => { + const job = { id: "J1", kind: "update-index" as const, target: "_index", runId: "r", queuedAt: NOW - MIN }; + assert.equal(status(inputs({ jobs: [{ ...job, status: "running" }] })).index.chip.text, "updating…"); + const q = status(inputs({ jobs: [{ ...job, status: "queued" }] })); + assert.equal(q.index.chip.text, "update queued"); + assert.equal(q.busy, true); +}); + +// --- the built chip ------------------------------------------------------------ + +test("built: fresh; never built; stale by channels, by data, by config; a bundle problem", () => { + assert.match(siteOf(status(inputs()), "alpha").chips.built.text, /^built 50 min ago$/); + assert.equal(siteOf(status(inputs()), "alpha").buildFreshness.state, "fresh"); + + const never = inputs(); + never.sites[0].built = null; + assert.equal(siteOf(status(never), "alpha").chips.built.text, "never built"); + + // Channels: an ingest of a MEMBER after the build. + const ch = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "elsewhere": NOW - MIN } }); + const a = siteOf(status(ch), "alpha"); + assert.deepEqual(a.buildStale, { reason: "channels", changedChannels: ["alpha-ch"] }); + assert.equal(a.chips.built.text, "stale: 1 channel changed (alpha-ch)"); + assert.equal(siteOf(status(ch), "beta").buildStale, undefined, "beta's members did not change"); + + // Data: the index says the site's inputs moved. + const data = inputs({ + index: { + stamp: stamp({ sites: { ...stamp().sites, alpha: { siteFp: null, statsFp: null, inputSig: "moved" } } }), + lastIngestDoneAt: null, + configChangedAt: null, + }, + }); + const d = siteOf(status(data), "alpha"); + assert.deepEqual(d.buildStale, { reason: "data", changedChannels: [] }); + assert.equal(d.chips.built.text, "stale: data changed"); + assert.equal(d.builtFromCurrentIndex, false); + assert.equal(siteOf(status(data), "beta").builtFromCurrentIndex, true); + + // Config: its site.json (or tags, aliases, duplicates) after the build. + const cfg = inputs(); + cfg.sites[0].configChangedAt = NOW - MIN; + assert.equal(siteOf(status(cfg), "alpha").chips.built.text, "stale: config changed"); + + const bundle = inputs(); + bundle.sites[0].bundleProblem = "no corpus.json"; + assert.equal(siteOf(status(bundle), "alpha").chips.built.text, "bundle: no corpus.json"); +}); + +test("built: channels changed is measured against builtCheckedAt — a no-op build clears it", () => { + const i = inputs({ ingestByChannel: { "alpha-ch": NOW - 30 * MIN } }); + assert.equal(siteOf(status(i), "alpha").buildStale?.reason, "channels"); + // A later no-op build found the bundle current at NOW - 20 min. + i.sites[0].built = built("alpha", { checkedAt: NOW - 20 * MIN }); + const a = siteOf(status(i), "alpha"); + assert.equal(a.buildStale, undefined); + assert.equal(a.buildFreshness.state, "fresh"); +}); + +test("built: code newer is said, never stale", () => { + const s = status(inputs({ commit: "c2" })); + const a = siteOf(s, "alpha"); + assert.equal(a.codeNewer, true); + assert.equal(a.buildFreshness.state, "fresh"); + assert.match(a.chips.built.text, /· code newer$/); + assert.equal(a.chips.built.tone, "ok"); + assert.equal(s.plan.steps.length, 0); +}); + +// --- deployed / live --------------------------------------------------------------- + +test("deployed: the policy's record; a newer build not deployed; never; private", () => { + const ok = siteOf(status(inputs()), "alpha"); + assert.equal(ok.deployedIsBuilt, true); + assert.match(ok.chips.deployed.text, /^production · 40 min ago$/); + + const newer = inputs(); + newer.sites[0].built = built("alpha", { stampId: "built-alpha-2" }); + const n = siteOf(status(newer), "alpha"); + assert.equal(n.deployedIsBuilt, false); + assert.equal(n.chips.deployed.text, "production: a newer build is not deployed"); + + const preview = inputs(); + preview.sites[0].policy = "preview"; + const p = siteOf(status(preview), "alpha"); + assert.equal(p.deployKind, "preview"); + assert.equal(p.chips.deployed.text, 'not on preview "preview"'); + assert.equal(p.previewUrl, "https://preview.alpha.pages.dev"); + + const priv = inputs(); + priv.sites[0] = site("alpha", { private: true, deployed: null, deployProblem: "private site", policy: "build" }); + assert.equal(siteOf(status(priv), "alpha").chips.deployed.text, "never deployed (private)"); + assert.equal(siteOf(status(inputs()), "alpha").chips.live.text, "not checked"); +}); + +test("deployed: a deploy refused for want of its build reads 'waiting for its build'", () => { + const i = inputs({ + ended: [ + { + id: "D1", + kind: "deploy-site", + target: "alpha", + status: "failed", + exitCode: 3, + endedAt: NOW - MIN, + runId: "r", + builtAfter: NOW - 2 * MIN, + }, + ], + }); + const a = siteOf(status(i), "alpha"); + assert.deepEqual(a.chips.deployed, { tone: "blocked", text: "waiting for its build" }); + i.ended[0] = { ...i.ended[0], builtAfter: undefined }; + assert.equal(siteOf(status(i), "alpha").chips.deployed.text, "last deploy refused: precondition not met"); +}); + +test("live: the record's live-check verdict", () => { + const i = inputs(); + const rec = i.sites[0].deployed!.production!; + const check = { + at: NOW, + url: "https://alpha.pages.dev/corpus.json", + plain: { status: 200 }, + busted: { status: 200 }, + expected: null, + }; + rec.liveCheck = { ...check, verdict: "ok" }; + assert.equal(siteOf(status(i), "alpha").chips.live.text, "live ok"); + rec.liveCheck = { ...check, verdict: "stale-edge" }; + assert.equal(siteOf(status(i), "alpha").chips.live.tone, "warn"); +}); + +// --- the plan: policies ------------------------------------------------------------ + +test("Publish now with every policy off is the index update alone", () => { + const i = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN }, index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null } }); + const s = status(i); + assert.deepEqual(kinds(s.plan), ["update-index _index"]); + assert.deepEqual( + s.plan.skipped.map((x) => `${x.kind} ${x.target}: ${x.reason}`), + ["build-site alpha: stale, and its policy is off"], + ); + assert.equal(siteOf(s, "alpha").next, "stale, and its policy is off"); +}); + +test("policies: build builds; preview and production also deploy, after the build, with builtAfter", () => { + const i = inputs({ + ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "beta-ch": NOW - 10 * MIN }, + index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null }, + }); + i.sites[0].policy = "preview"; + i.sites[1].policy = "production"; + const p = status(i).plan; + assert.deepEqual(kinds(p), [ + "update-index _index", + "build-site alpha", + "deploy-site alpha", + "build-site beta", + "deploy-site beta", + ]); + const [, ba, da, , db] = p.steps; + assert.equal(ba.indexAfter, NOW, "a build in a run with the index carries indexAfter"); + assert.equal(da.builtAfter, NOW, "a deploy whose build is in the run carries builtAfter"); + assert.equal(da.preview, "preview"); + assert.equal(da.previewUrl, "https://preview.alpha.pages.dev"); + assert.equal(db.preview, undefined, "production"); + assert.equal(db.builtAfter, NOW); + // The words the targets show. + assert.equal(siteOf(status(i), "alpha").next, 'build, then deploy to preview "preview"'); + + const buildOnly = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN } }); + buildOnly.sites[0].policy = "build"; + const b = status(buildOnly).plan; + assert.deepEqual(kinds(b), ["build-site alpha"]); + assert.equal(b.steps[0].indexAfter, undefined, "no index update in the run, no indexAfter"); +}); + +test("plan: a site with neither changed channels nor a signature mismatch is skipped", () => { + const i = inputs({ index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null } }); + for (const s of i.sites) s.policy = "production"; + assert.deepEqual(kinds(status(i).plan), ["update-index _index"]); +}); + +test("plan: a deploy-only step when the build is current but not deployed", () => { + const i = inputs(); + i.sites[0].policy = "production"; + i.sites[0].deployed = null; + const p = status(i).plan; + assert.deepEqual(kinds(p), ["deploy-site alpha"]); + assert.equal(p.steps[0].builtAfter, undefined); +}); + +test("private sites never deploy; a project-less site is skipped with the reason", () => { + const i = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "beta-ch": NOW - 10 * MIN } }); + // A private site's policy is clamped on read; even a stray production one + // is never planned as a deploy. + i.sites[0] = site("alpha", { private: true, policy: "production", deployProblem: "Site \"alpha\" is private" }); + i.sites[1] = site("beta", { policy: "preview", pagesProblem: "no Pages project", cloudflareProject: null }); + const p = status(i).plan; + assert.deepEqual(kinds(p), ["build-site alpha", "build-site beta"]); + assert.deepEqual( + p.skipped.map((x) => `${x.kind} ${x.target}`).sort(), + ["deploy-site alpha", "deploy-site beta"], + ); +}); + +test("plan: builds 'stale' builds every stale site whatever its policy, deploys only by policy", () => { + const i = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "beta-ch": NOW - 10 * MIN } }); + i.sites[1].policy = "production"; + const p = planPublishRun(status(i), { builds: "stale" }); + assert.deepEqual(kinds(p), ["build-site alpha", "build-site beta", "deploy-site beta"]); + assert.deepEqual(kinds(planPublishRun(status(i), { builds: "stale", deploys: "none" })), [ + "build-site alpha", + "build-site beta", + ]); +}); + +test("plan: the docker runner builds every stale site in one _all stage", () => { + const i = inputs({ ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "beta-ch": NOW - 10 * MIN } }, { runner: "docker" }); + i.sites[0].policy = "production"; + i.sites[1].policy = "build"; + const p = status(i).plan; + assert.deepEqual(kinds(p), ["build-site _all", "deploy-site alpha"]); + assert.equal(p.steps[0].runner, "docker"); + assert.equal(p.steps[1].builtAfter, NOW); +}); + +test("plan: the hub and the homepage follow settings.publish's policies", () => { + const i = inputs( + { + ingestByChannel: { "shared-ch": NOW - 10 * MIN }, + index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null }, + }, + { hub: "production", homepage: "build" }, + ); + const p = status(i).plan; + assert.deepEqual(kinds(p), [ + "update-index _index", + "build-hub _hub", + "deploy-hub _hub", + "build-homepage _homepage", + ]); + assert.equal(p.steps[1].indexAfter, NOW); + assert.equal(p.steps[2].builtAfter, NOW); + // Off (the default): nothing. + assert.deepEqual(kinds(status(inputs({ index: i.index, ingestByChannel: i.ingestByChannel })).plan), [ + "update-index _index", + ]); + // The homepage is stale when the index moved past its build. + const moved = inputs({ index: { stamp: stamp({ stampId: "stamp-2" }), lastIngestDoneAt: null, configChangedAt: null } }, { homepage: "build" }); + assert.deepEqual(kinds(status(moved).plan), ["build-homepage _homepage"]); +}); + +test("plan: never forced, and the lane's exclusions and skipIndex leave steps out", () => { + const i = inputs({ + ingestByChannel: { "alpha-ch": NOW - 10 * MIN }, + index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null }, + }); + i.sites[0].policy = "production"; + const s = status(i); + for (const step of s.plan.steps) assert.equal(step.force, undefined, `${step.kind} is not forced`); + const exclude = new Set([stepKey({ kind: "build-site", target: "alpha" })]); + assert.deepEqual(kinds(planPublishRun(s, { exclude, skipIndex: true })), ["deploy-site alpha"]); + assert.equal(stepKey({ kind: "deploy-site", target: "a", preview: "p" }), "deploy-site:a:preview:p"); + assert.equal(stepKey({ kind: "deploy-site", target: "a" }), "deploy-site:a:production"); + assert.equal(stepKey({ kind: "deploy-site", target: "a", to: "local" }), "deploy-site:a:local"); + // runStart is the caller's when given. + assert.equal(planPublishRun(s, { runStart: 5 }).steps[1].indexAfter, 5); +}); + +test("every plan step names a real stage and passes the stage's own argv parser", () => { + const i = inputs({ + ingestByChannel: { "alpha-ch": NOW - 10 * MIN, "shared-ch": NOW - 10 * MIN }, + index: { stamp: stamp(), lastIngestDoneAt: NOW - 10 * MIN, configChangedAt: null }, + }, { hub: "preview", homepage: "production" }); + i.sites[0].policy = "preview"; + for (const step of status(i).plan.steps) { + const stage = STAGES[step.kind]; + assert.ok(stage, step.kind); + const argv = stage.argv({ ...step, runId: "r1" }); + assert.equal(argv[0], "stage"); + assert.equal(argv[1], step.kind); + } +}); + +// --- the lane's pass decision ------------------------------------------------------- + +test("pass decision: off, held, quiet hours and busy never run", () => { + const base = { + settings: { ...defaultPublish(), enabled: true }, + stamp: null, + indexFresh: false, + busy: false, + now: NOW, + lastPassAt: null, + policyWork: false, + }; + assert.deepEqual(publishPassDecision(base), { run: true, reason: "no index stamp yet" }); + assert.equal(publishPassDecision({ ...base, settings: defaultPublish() }).reason, "the lane is off"); + assert.equal(publishPassDecision({ ...base, settings: { ...base.settings, held: true } }).reason, "held"); + assert.equal(publishPassDecision({ ...base, busy: true }).run, false); + const hour = new Date(NOW).getHours(); + const quiet = { start: hour, end: (hour + 1) % 24 }; + assert.equal(publishPassDecision({ ...base, settings: { ...base.settings, quietHours: quiet } }).reason, "quiet hours"); + const notQuiet = { start: (hour + 1) % 24, end: (hour + 2) % 24 }; + assert.equal(publishPassDecision({ ...base, settings: { ...base.settings, quietHours: notQuiet } }).run, true); +}); + +test("pass decision: a stale index waits out refreshEveryMinutes from the stamp's builtAt", () => { + const settings = { ...defaultPublish(), enabled: true, refreshEveryMinutes: 60 }; + const a = { settings, indexFresh: false, busy: false, now: NOW, lastPassAt: null, policyWork: false }; + assert.equal(publishPassDecision({ ...a, stamp: stamp({ builtAt: NOW - 61 * MIN }) }).run, true); + const early = publishPassDecision({ ...a, stamp: stamp({ builtAt: NOW - 30 * MIN }) }); + assert.equal(early.run, false); + assert.match(early.reason, /the index is stale; the next update is due 30 min from now/); + // 0 = whenever stale. + assert.equal( + publishPassDecision({ ...a, settings: { ...settings, refreshEveryMinutes: 0 }, stamp: stamp({ builtAt: NOW }) }).run, + true, + ); + // A fresh index with nothing to do never runs. + assert.equal(publishPassDecision({ ...a, indexFresh: true, stamp: stamp() }).run, false); + // Policy work left: once per interval since the last pass. + assert.equal(publishPassDecision({ ...a, indexFresh: true, stamp: stamp(), policyWork: true }).run, true); + assert.equal( + publishPassDecision({ ...a, indexFresh: true, stamp: stamp(), policyWork: true, lastPassAt: NOW - 10 * MIN }).run, + false, + ); +}); + +test("the status carries the lane's words and the plan Publish now would run", () => { + const s = status(inputs({}, { enabled: true })); + assert.equal(s.lane.enabled, true); + assert.equal(s.lane.due, false); + assert.match(s.lane.reason, /nothing to do/); + assert.equal(s.lane.chip.text, "lane on, runner not running"); + const off = status(inputs()); + assert.equal(off.lane.chip.text, "lane off"); + assert.equal(off.lane.blockedReason, "the lane is off"); + assert.deepEqual(off.plan.steps, []); +}); diff --git a/common/publish/publishPlan.ts b/common/publish/publishPlan.ts @@ -0,0 +1,876 @@ +// THE PUBLISH STATUS AND THE PUBLISH PLAN (release 18) — pure. +// +// One state builder (the plan's "Surfaces"): `publishState.ts` reads the +// stamps, the sites, the settings, the config mtimes, the ingest job metas and +// the live publish jobs into a `PublishInputs`; `buildPublishStatus(inputs, +// now)` folds them into the `PublishStatus` every surface reads — the /sites +// Publish panel, `archilyzer publish status`, `GET /api/ops/publish` and the +// lane; `planPublishRun(status, opts)` says what a run would enqueue — what +// Publish now enqueues all at once, and what the lane dispatches one stage at +// a time. +// +// PURE: no fs, no clock, no singleton. `now` is an argument. Every freshness +// question is the stage's own `needs()` (publish/stages.ts) over the +// `NeedsInput` built here — the status never re-implements a stage's rule; it +// only adds the two signals a stage child cannot see (`changedChannels` from +// the job metas, and the config mtimes) and the chips' words. +// +// WHY HERE AND NOT IN views/: the lane's runner (publish/publishRunner.ts) +// plans from it, and the publish layer may not import views/ (common/ +// architecture.test.ts). `views/publishStatus.ts` re-exports it as the view +// layer's name for it. SERVER-SIDE VALUES: it imports stages.ts, which imports +// the stamp readers; a client component imports its TYPES only. + +import type { PublishPolicy, PublishSettings } from "../lib/settingsSchema"; +import { previewAliasUrl } from "../lib/pagesDeploy"; +import { isInQuietHours } from "../jobs/syncScheduler"; +import { + STAGES, + builtCheckedAt, + deployKindOf, + type Freshness, + type NeedsInput, + type StageKind, + type StageRequest, + type TargetState, +} from "./stages"; +import { + ALL_TARGET, + HOMEPAGE_TARGET, + HUB_TARGET, + INDEX_TARGET, + deployRecordFor, + type BuiltStamp, + type DeployKind, + type DeployRecord, + type DeployedFile, + type IndexStamp, + type LiveCheck, +} from "./stamps"; + +// --------------------------------------------------------------------------- +// Inputs (what publishState.ts reads) +// --------------------------------------------------------------------------- + +/** A publish stage job, queued or running, on the `publish` queue. */ +export type PublishJob = { + id: string; + kind: StageKind; + target: string; + status: "queued" | "running"; + runId: string | null; + queuedAt: number; + startedAt?: number; + preview?: string; + to?: "pages" | "local"; +}; + +/** The newest publish stage job of a kind + target that ENDED. */ +export type PublishJobEnded = { + id: string; + kind: StageKind; + target: string; + status: "done" | "failed" | "cancelled"; + exitCode: number | null; + endedAt: number; + runId: string | null; + // The run's preconditions the job carried (a refusal for want of one reads + // "waiting for …"). + indexAfter?: number; + builtAfter?: number; + preview?: string; + to?: "pages" | "local"; +}; + +/** The lane runner's live state (publishLaneState.ts), as the editor holds it. */ +export type PublishLaneLive = { + // Whether this process can know (the CLI cannot: the lane is the editor's). + known: boolean; + running: boolean; + jobId: string | null; + passRunning: boolean; + lastCheckAt: number | null; + nextCheckAt: number | null; + lastPassAt: number | null; + lastPassSummary: string | null; + lastDecision: string | null; +}; + +export type PublishTargetInput = { + built: BuiltStamp | null; + deployed: DeployedFile | null; + bundleProblem: string | null; + // Never deployed anywhere (a private site), or null. + deployProblem: string | null; + // Cannot go to Cloudflare Pages (no project), or null. + pagesProblem: string | null; + // The newest mtime of a config file this target's build reads. + configChangedAt: number | null; + policy: PublishPolicy; + // The Pages project, for a preview's alias URL. + cloudflareProject: string | null; + // The public URL the live check reads. + url: string | null; +}; + +export type PublishSiteInput = PublishTargetInput & { + siteId: string; + title: string; + private: boolean; + listed: boolean; + // Member channel slugs (site.json `channels`). + members: string[]; +}; + +export type PublishInputs = { + index: { + stamp: IndexStamp | null; + // When the newest ingest ended (ms), or null. + lastIngestDoneAt: number | null; + // The newest mtime of an index input config file, or null. + configChangedAt: number | null; + }; + // Per channel: when its newest ingest ended (a job meta, or its report's + // regeneration — see publishState.ts). + ingestByChannel: Record<string, number>; + // The code this editor (or CLI) runs: a build of another commit reads + // "code newer" — never stale. + commit: string | null; + sites: PublishSiteInput[]; + hub: PublishTargetInput; + homepage: PublishTargetInput & { mainHead: string | null }; + settings: PublishSettings; + jobs: PublishJob[]; + ended: PublishJobEnded[]; + lane: PublishLaneLive; +}; + +// --------------------------------------------------------------------------- +// The status +// --------------------------------------------------------------------------- + +export type ChipTone = "ok" | "stale" | "blocked" | "busy" | "warn" | "off"; +export type Chip = { tone: ChipTone; text: string }; + +export type BuildStale = { + reason: "channels" | "data" | "config"; + changedChannels: string[]; +}; + +export type PublishTargetStatus = { + target: string; + kind: "site" | "hub" | "homepage"; + title: string; + indexFresh: boolean; + built: BuiltStamp | null; + // The bundle is what the CURRENT index would build (same signature). + builtFromCurrentIndex: boolean; + // The build stage's own answer (needs()). + buildFreshness: Freshness; + buildStale?: BuildStale; + // Built by another commit than the one running: says so, never stale. + codeNewer: boolean; + deployed: DeployedFile | null; + // Where the policy deploys (null: it does not), and that record. + deployKind: DeployKind | null; + deployRecord: DeployRecord | null; + // The record the chips read is of the bundle on disk. + deployedIsBuilt: boolean; + previewUrl?: string; + liveCheck?: LiveCheck; + policy: PublishPolicy; + // The Pages project (a preview's alias URL is built from it), or null. + cloudflareProject: string | null; + // The public URL the live check reads, or null. + url: string | null; + deployProblem?: string; + bundleProblem?: string; + // What happens next, in words ("build, then deploy to preview", …). + next: string; + running?: PublishJob; + queued: PublishJob[]; + // The newest ended build / deploy job (for a refusal's words). + lastBuild?: PublishJobEnded; + lastDeploy?: PublishJobEnded; + chips: { index: Chip; built: Chip; deployed: Chip; live: Chip }; +}; + +export type PublishLaneStatus = { + enabled: boolean; + held: boolean; + quietNow: boolean; + busy: boolean; + live: PublishLaneLive; + // Whether a pass would run now, and why (or why not). + due: boolean; + reason: string; + // Why the lane will not dispatch at all right now (off, held, quiet hours, + // a publish job queued or running), or null. + blockedReason: string | null; + chip: Chip; +}; + +export type PlanStep = Omit<StageRequest, "runId"> & { + reason: string; + previewUrl?: string; +}; + +export type PlanSkip = { kind: StageKind; target: string; reason: string }; + +export type PublishPlan = { + runStart: number; + steps: PlanStep[]; + skipped: PlanSkip[]; +}; + +export type PublishStatus = { + now: number; + commit: string | null; + index: { + stamp: IndexStamp | null; + freshness: Freshness; + fresh: boolean; + running?: PublishJob; + queued: PublishJob[]; + lastRun?: PublishJobEnded; + chip: Chip; + }; + sites: PublishTargetStatus[]; + hub: PublishTargetStatus; + homepage: PublishTargetStatus; + // Any publish stage queued or running. + busy: boolean; + jobs: PublishJob[]; + settings: PublishSettings; + lane: PublishLaneStatus; + // What Publish now would enqueue (deploys by policy). + plan: PublishPlan; + // What the stages' needs() read. + needs: NeedsInput; +}; + +const MIN = 60_000; + +function namesList(slugs: string[], max = 4): string { + const shown = slugs.slice(0, max).join(", "); + return slugs.length > max ? `${shown}, …` : shown; +} + +/** "3 min ago", "2 h ago", "4 d ago" — the chips' one clock word. */ +export function agoText(then: number, now: number): string { + const ms = Math.max(0, now - then); + if (ms < MIN) return "just now"; + if (ms < 60 * MIN) return `${Math.floor(ms / MIN)} min ago`; + if (ms < 48 * 60 * MIN) return `${Math.floor(ms / (60 * MIN))} h ago`; + return `${Math.floor(ms / (24 * 60 * MIN))} d ago`; +} + +function changedSince( + members: Iterable<string>, + byChannel: Record<string, number>, + since: number, +): string[] { + const out = new Set<string>(); + for (const slug of members) { + const at = byChannel[slug]; + if (at !== undefined && at > since) out.add(slug); + } + return [...out].sort(); +} + +function targetState(t: PublishTargetInput, changedChannels: string[]): TargetState { + return { + built: t.built, + deployed: t.deployed, + changedChannels, + configChangedAt: t.configChangedAt, + bundleProblem: t.bundleProblem, + deployProblem: t.deployProblem, + pagesProblem: t.pagesProblem, + }; +} + +/** The NeedsInput the stages' needs() read, with the two signals a child cannot see. */ +export function needsInputOf(i: PublishInputs): NeedsInput { + const sites: Record<string, TargetState> = {}; + for (const s of i.sites) { + const changed = s.built ? changedSince(s.members, i.ingestByChannel, builtCheckedAt(s.built)) : []; + sites[s.siteId] = targetState(s, changed); + } + const listedMembers = new Set(i.sites.filter((s) => s.listed).flatMap((s) => s.members)); + const hubChanged = i.hub.built + ? changedSince(listedMembers, i.ingestByChannel, builtCheckedAt(i.hub.built)) + : []; + const homeChanged = i.homepage.built + ? changedSince(Object.keys(i.ingestByChannel), i.ingestByChannel, builtCheckedAt(i.homepage.built)) + : []; + return { + index: { ...i.index }, + sites, + hub: targetState(i.hub, hubChanged), + homepage: { ...targetState(i.homepage, homeChanged), mainHead: i.homepage.mainHead }, + }; +} + +function req(kind: StageKind, target: string, extra: Partial<StageRequest> = {}): StageRequest { + return { kind, target, runId: "status", ...extra }; +} + +/** The build stage that builds `kind`'s target. */ +function buildKindOf(kind: PublishTargetStatus["kind"]): StageKind { + return kind === "site" ? "build-site" : kind === "hub" ? "build-hub" : "build-homepage"; +} + +function deployStageOf(kind: PublishTargetStatus["kind"]): StageKind { + return kind === "site" ? "deploy-site" : kind === "hub" ? "deploy-hub" : "deploy-homepage"; +} + +/** Where a policy deploys: production, the preview branch, or nowhere. */ +export function policyDeploy( + policy: PublishPolicy, + previewBranch: string, +): { to: DeployKind; preview?: string } | null { + if (policy === "production") return { to: "production" }; + if (policy === "preview") return { to: "preview", preview: previewBranch }; + return null; +} + +// The deploy record the chips read: the policy's destination when it deploys, +// else production, else local, else the newest preview. +function primaryRecord( + deployed: DeployedFile | null, + dest: { to: DeployKind; preview?: string } | null, +): DeployRecord | null { + if (!deployed) return null; + if (dest) return deployRecordFor(deployed, dest.to, dest.preview); + if (deployed.production) return deployed.production; + if (deployed.local) return deployed.local; + const previews = Object.values(deployed.previews).sort((a, b) => b.at - a.at); + return previews[0] ?? null; +} + +function previewUrlOf(t: PublishTargetInput, previewBranch: string): string | undefined { + const rec = t.deployed?.previews[previewBranch] ?? null; + if (rec?.alias) return rec.alias; + if (rec?.url) return rec.url; + if (t.policy === "preview" && t.cloudflareProject) return previewAliasUrl(t.cloudflareProject, previewBranch); + return undefined; +} + +function buildStaleOf( + t: PublishTargetInput, + state: TargetState, + sigMatches: boolean, +): BuildStale | undefined { + if (!t.built) return undefined; + if (state.changedChannels.length > 0) return { reason: "channels", changedChannels: state.changedChannels }; + if (t.configChangedAt !== null && t.configChangedAt > builtCheckedAt(t.built)) { + return { reason: "config", changedChannels: [] }; + } + if (!sigMatches) return { reason: "data", changedChannels: [] }; + return undefined; +} + +function liveChip(check: LiveCheck | undefined): Chip { + if (!check) return { tone: "off", text: "not checked" }; + switch (check.verdict) { + case "ok": + return { tone: "ok", text: "live ok" }; + case "stale-edge": + return { tone: "warn", text: "live: the edge serves an older build" }; + case "mismatch": + return { tone: "warn", text: "live: a different build answers" }; + case "unreachable": + return { tone: "warn", text: "live: unreachable" }; + default: + return { tone: "off", text: "live check skipped" }; + } +} + +// The words for a stage job that ended refused (exit 3) or failed. +function endedWords(e: PublishJobEnded, what: "build" | "deploy"): Chip | null { + if (e.status === "done") return null; + if (e.status === "cancelled") return { tone: "warn", text: `last ${what} cancelled` }; + if (e.exitCode === 3) { + if (what === "deploy" && e.builtAfter !== undefined) return { tone: "blocked", text: "waiting for its build" }; + if (what === "build" && e.indexAfter !== undefined) return { tone: "blocked", text: "waiting for the index" }; + return { tone: "blocked", text: `last ${what} refused: precondition not met` }; + } + return { tone: "warn", text: `last ${what} failed` }; +} + +function pickJobs(jobs: PublishJob[], kinds: StageKind[], target: string): { running?: PublishJob; queued: PublishJob[] } { + const mine = jobs.filter( + (j) => kinds.includes(j.kind) && (j.target === target || (j.target === ALL_TARGET && kinds.includes("build-site"))), + ); + return { + running: mine.find((j) => j.status === "running"), + queued: mine.filter((j) => j.status === "queued").sort((a, b) => a.queuedAt - b.queuedAt), + }; +} + +function newestEnded(ended: PublishJobEnded[], kind: StageKind, target: string): PublishJobEnded | undefined { + return ended + .filter((e) => e.kind === kind && (e.target === target || (kind === "build-site" && e.target === ALL_TARGET))) + .sort((a, b) => b.endedAt - a.endedAt)[0]; +} + +function targetStatus( + i: PublishInputs, + needs: NeedsInput, + now: number, + indexFresh: boolean, + indexChip: Chip, + kind: PublishTargetStatus["kind"], + target: string, + title: string, + t: PublishTargetInput, + state: TargetState, +): Omit<PublishTargetStatus, "next"> { + const stamp = needs.index.stamp; + const buildKind = buildKindOf(kind); + const buildFreshness = STAGES[buildKind].needs(needs, req(buildKind, target)); + let sigMatches = true; + let builtFromCurrentIndex = false; + if (t.built && stamp) { + if (kind === "site") sigMatches = stamp.sites[target] !== undefined && t.built.inputSig === stamp.sites[target].inputSig; + else if (kind === "hub") sigMatches = t.built.inputSig === stamp.hubSig; + else sigMatches = t.built.indexStampId === stamp.stampId; + builtFromCurrentIndex = sigMatches; + } + const buildStale = buildStaleOf(t, state, sigMatches); + const codeNewer = Boolean(t.built?.commit && i.commit && t.built.commit !== i.commit); + const dest = policyDeploy(t.policy, i.settings.previewBranch); + const deployRecord = primaryRecord(t.deployed, dest); + const deployedIsBuilt = Boolean(t.built && deployRecord && deployRecord.builtStampId === t.built.stampId); + const build = pickJobs(i.jobs, [buildKind], target); + const deploy = pickJobs(i.jobs, [deployStageOf(kind)], target); + const running = build.running ?? deploy.running; + const queued = [...build.queued, ...deploy.queued]; + const lastBuild = newestEnded(i.ended, buildKind, target); + const lastDeploy = newestEnded(i.ended, deployStageOf(kind), target); + + // --- built chip + let built: Chip; + if (build.running) built = { tone: "busy", text: "building…" }; + else if (build.queued.length > 0) built = { tone: "busy", text: "build queued" }; + else if (buildFreshness.state === "blocked") built = { tone: "blocked", text: buildFreshness.reason }; + else if (!t.built) built = { tone: "stale", text: "never built" }; + else if (buildStale?.reason === "channels") { + const n = buildStale.changedChannels.length; + built = { tone: "stale", text: `stale: ${n} channel${n === 1 ? "" : "s"} changed (${namesList(buildStale.changedChannels)})` }; + } else if (buildStale?.reason === "config") built = { tone: "stale", text: "stale: config changed" }; + else if (buildStale?.reason === "data") built = { tone: "stale", text: "stale: data changed" }; + else if (t.bundleProblem) built = { tone: "warn", text: `bundle: ${t.bundleProblem}` }; + else built = { tone: "ok", text: `built ${agoText(t.built.builtAt, now)}` }; + if (built.tone !== "busy" && lastBuild && (!t.built || lastBuild.endedAt > t.built.builtAt)) { + const w = endedWords(lastBuild, "build"); + if (w && built.tone !== "ok") built = { tone: w.tone, text: `${built.text} · ${w.text}` }; + } + if (codeNewer && t.built) built = { ...built, text: `${built.text} · code newer` }; + + // --- deployed chip + let deployed: Chip; + if (deploy.running) deployed = { tone: "busy", text: "deploying…" }; + else if (deploy.queued.length > 0) deployed = { tone: "busy", text: "deploy queued" }; + else if (t.deployProblem) deployed = { tone: "off", text: "never deployed (private)" }; + else if (!deployRecord) { + deployed = dest + ? { tone: "stale", text: dest.to === "preview" ? `not on preview "${dest.preview}"` : "never deployed" } + : { tone: "off", text: "not deployed" }; + } else { + const where = deployRecord.kind === "preview" ? `preview ${deployRecord.branch ?? ""}`.trim() : deployRecord.kind; + deployed = deployedIsBuilt + ? { tone: "ok", text: `${where} · ${agoText(deployRecord.at, now)}` } + : { tone: "stale", text: `${where}: a newer build is not deployed` }; + } + if (deployed.tone !== "busy" && lastDeploy && (!deployRecord || lastDeploy.endedAt > deployRecord.at)) { + const w = endedWords(lastDeploy, "deploy"); + if (w) deployed = w.tone === "blocked" ? w : { tone: w.tone, text: `${deployed.text} · ${w.text}` }; + } + + const liveCheck = deployRecord?.liveCheck ?? undefined; + return { + target, + kind, + title, + indexFresh, + built: t.built, + builtFromCurrentIndex, + buildFreshness, + ...(buildStale ? { buildStale } : {}), + codeNewer, + deployed: t.deployed, + deployKind: dest?.to ?? null, + deployRecord, + deployedIsBuilt, + ...(previewUrlOf(t, i.settings.previewBranch) ? { previewUrl: previewUrlOf(t, i.settings.previewBranch) } : {}), + ...(liveCheck ? { liveCheck } : {}), + policy: t.policy, + cloudflareProject: t.cloudflareProject, + url: t.url, + ...(t.deployProblem ? { deployProblem: t.deployProblem } : {}), + ...(t.bundleProblem ? { bundleProblem: t.bundleProblem } : {}), + ...(running ? { running } : {}), + queued, + ...(lastBuild ? { lastBuild } : {}), + ...(lastDeploy ? { lastDeploy } : {}), + chips: { index: indexChip, built, deployed, live: liveChip(liveCheck) }, + }; +} + +// --------------------------------------------------------------------------- +// The lane's decision +// --------------------------------------------------------------------------- + +export type PassDecision = { run: boolean; reason: string }; + +/** + * Whether the lane may dispatch at all now: off, held, quiet hours, or a + * publish stage someone else queued. Null when it may. + */ +export function laneBlockedReason( + settings: PublishSettings, + busy: boolean, + now: number, +): string | null { + if (!settings.enabled) return "the lane is off"; + if (settings.held) return "held"; + if (settings.quietHours && isInQuietHours(now, settings.quietHours.start, settings.quietHours.end)) { + return "quiet hours"; + } + if (busy) return "a publish stage is queued or running"; + return null; +} + +/** + * Is a pass due? When the lane may dispatch AND: there is no index stamp; or + * the index is stale and its last update is at least `refreshEveryMinutes` + * old; or the policies have work left (a target stale since the last update — + * a failed stage, a policy just switched on) and the last pass is at least + * that old. + */ +export function publishPassDecision(a: { + settings: PublishSettings; + stamp: IndexStamp | null; + indexFresh: boolean; + busy: boolean; + now: number; + lastPassAt: number | null; + policyWork: boolean; +}): PassDecision { + const blocked = laneBlockedReason(a.settings, a.busy, a.now); + if (blocked) return { run: false, reason: blocked }; + if (!a.stamp) return { run: true, reason: "no index stamp yet" }; + const refreshMs = a.settings.refreshEveryMinutes * MIN; + if (!a.indexFresh) { + const age = a.now - a.stamp.builtAt; + if (age >= refreshMs) return { run: true, reason: "the index is stale" }; + return { + run: false, + reason: `the index is stale; the next update is due ${agoText(a.now - (refreshMs - age), a.now).replace(" ago", "")} from now`, + }; + } + if (a.policyWork) { + if (a.lastPassAt === null || a.now - a.lastPassAt >= refreshMs) { + return { run: true, reason: "a target's policy has work left" }; + } + return { run: false, reason: "a target's policy has work left; waiting out the refresh interval" }; + } + return { run: false, reason: "nothing to do: the index and every policy target are current" }; +} + +// --------------------------------------------------------------------------- +// buildPublishStatus +// --------------------------------------------------------------------------- + +export function buildPublishStatus(i: PublishInputs, now: number): PublishStatus { + const needs = needsInputOf(i); + const indexFreshness = STAGES["update-index"].needs(needs, req("update-index", INDEX_TARGET)); + const indexFresh = indexFreshness.state === "fresh"; + const indexJobs = pickJobs(i.jobs, ["update-index"], INDEX_TARGET); + const lastIndex = newestEnded(i.ended, "update-index", INDEX_TARGET); + const stamp = i.index.stamp; + let indexChip: Chip; + if (indexJobs.running) indexChip = { tone: "busy", text: "updating…" }; + else if (indexJobs.queued.length > 0) indexChip = { tone: "busy", text: "update queued" }; + else if (!stamp) indexChip = { tone: "stale", text: "no index yet" }; + else if (indexFreshness.state !== "fresh") indexChip = { tone: "stale", text: `stale: ${indexFreshness.reason}` }; + else indexChip = { tone: "ok", text: `fresh · updated ${agoText(stamp.builtAt, now)}` }; + if (indexChip.tone !== "busy" && lastIndex && lastIndex.status !== "done" && (!stamp || lastIndex.endedAt > stamp.builtAt)) { + indexChip = { tone: "warn", text: `${indexChip.text} · last update ${lastIndex.status}` }; + } + + const sites = [...i.sites] + .sort((a, b) => a.siteId.localeCompare(b.siteId)) + .map((s) => + targetStatus(i, needs, now, indexFresh, indexChip, "site", s.siteId, s.title, s, needs.sites[s.siteId]), + ); + const hub = targetStatus(i, needs, now, indexFresh, indexChip, "hub", HUB_TARGET, "Hub", { ...i.hub, policy: i.settings.hub }, needs.hub); + const homepage = targetStatus( + i, + needs, + now, + indexFresh, + indexChip, + "homepage", + HOMEPAGE_TARGET, + "Homepage", + { ...i.homepage, policy: i.settings.homepage }, + needs.homepage, + ); + const busy = i.jobs.length > 0; + const partial = { + now, + commit: i.commit, + index: { + stamp, + freshness: indexFreshness, + fresh: indexFresh, + ...(indexJobs.running ? { running: indexJobs.running } : {}), + queued: indexJobs.queued, + ...(lastIndex ? { lastRun: lastIndex } : {}), + chip: indexChip, + }, + busy, + jobs: i.jobs, + settings: i.settings, + needs, + }; + // The plan reads the targets without their `next` (it writes it). + const draft = { + ...partial, + sites: sites.map((s) => ({ ...s, next: "" })), + hub: { ...hub, next: "" }, + homepage: { ...homepage, next: "" }, + lane: undefined as unknown as PublishLaneStatus, + plan: { runStart: now, steps: [], skipped: [] } as PublishPlan, + } satisfies PublishStatus; + const plan = planPublishRun(draft, { deploys: "policy" }); + const withNext = (t: PublishTargetStatus): PublishTargetStatus => ({ ...t, next: nextWords(t, plan) }); + const policyWork = plan.steps.some((s) => s.kind !== "update-index"); + const blockedReason = laneBlockedReason(i.settings, busy, now); + const decision = publishPassDecision({ + settings: i.settings, + stamp, + indexFresh, + busy, + now, + lastPassAt: i.lane.lastPassAt, + policyWork, + }); + const quietNow = Boolean( + i.settings.quietHours && isInQuietHours(now, i.settings.quietHours.start, i.settings.quietHours.end), + ); + const lane: PublishLaneStatus = { + enabled: i.settings.enabled, + held: i.settings.held, + quietNow, + busy, + live: i.lane, + due: decision.run, + reason: decision.reason, + blockedReason, + chip: laneChip(i, blockedReason, decision), + }; + return { + ...partial, + sites: draft.sites.map(withNext), + hub: withNext(draft.hub), + homepage: withNext(draft.homepage), + lane, + plan, + }; +} + +function laneChip(i: PublishInputs, blocked: string | null, d: PassDecision): Chip { + if (!i.settings.enabled) return { tone: "off", text: "lane off" }; + if (i.settings.held) return { tone: "warn", text: "lane held" }; + if (i.lane.passRunning) return { tone: "busy", text: "lane: pass running" }; + if (i.lane.known && !i.lane.running) return { tone: "warn", text: "lane on, runner not running" }; + if (blocked) return { tone: "off", text: `lane waiting: ${blocked}` }; + return d.run ? { tone: "busy", text: "lane: pass due" } : { tone: "ok", text: "lane idle" }; +} + +function nextWords(t: PublishTargetStatus, plan: PublishPlan): string { + if (t.running) return `${t.running.kind} running`; + if (t.queued.length > 0) return `${t.queued.map((j) => j.kind).join(", ")} queued`; + const mine = plan.steps.filter((s) => s.target === t.target || (s.target === ALL_TARGET && t.kind === "site")); + if (mine.length > 0) { + return mine + .map((s) => + s.kind.startsWith("build") + ? "build" + : s.to === "local" + ? "deploy locally" + : s.preview + ? `deploy to preview "${s.preview}"` + : "deploy to production", + ) + .join(", then "); + } + const skip = plan.skipped.find((s) => s.target === t.target); + if (skip) return skip.reason; + if (t.buildFreshness.state !== "fresh" && t.policy === "off") return "stale — the policy is off (Build rebuilds it)"; + return "up to date"; +} + +// --------------------------------------------------------------------------- +// planPublishRun +// --------------------------------------------------------------------------- + +export type PlanOptions = { + // "policy": deploy where the target's policy says (Publish now, the lane); + // "none": build only. + deploys?: "policy" | "none"; + // "policy": build targets whose policy is not off; "stale": every stale + // site, whatever its policy ("Build all stale"). The hub and the homepage + // follow their policies either way. + builds?: "policy" | "stale"; + // The run's start (the on-disk preconditions); defaults to status.now. + runStart?: number; + // Step keys (`stepKey`) to leave out: what the lane already ran this pass. + exclude?: ReadonlySet<string>; + // Leave the update-index step out (the lane, after its index stage). + skipIndex?: boolean; +}; + +/** A step's identity within a run: kind, target and destination. */ +export function stepKey(s: Pick<StageRequest, "kind" | "target" | "preview" | "to">): string { + const dest = s.kind.startsWith("deploy") ? `:${deployKindOf(s)}${s.preview ? `:${s.preview}` : ""}` : ""; + return `${s.kind}:${s.target}${dest}`; +} + +/** + * What a run enqueues, in order: update-index when the index is stale; then + * per site (by id) its build — when it has changed channels, its signature no + * longer matches the index, it was never built, its config changed or its + * bundle has a problem — and its deploy where the policy says; then the hub; + * then the homepage. NEVER forced: a fresh target is not in the plan, and a + * stage whose child finds it fresh is a no-op. + * + * Preconditions are on disk (stages.ts): a build carries `indexAfter = + * runStart` when update-index is in the run; a deploy `builtAfter = runStart` + * when its build is. + */ +export function planPublishRun(status: PublishStatus, opts: PlanOptions = {}): PublishPlan { + const runStart = opts.runStart ?? status.now; + const exclude = opts.exclude ?? new Set<string>(); + const needs = status.needs; + const settings = status.settings; + const steps: PlanStep[] = []; + const skipped: PlanSkip[] = []; + const add = (s: PlanStep) => { + if (!exclude.has(stepKey(s))) steps.push(s); + }; + + const indexStep = + !opts.skipIndex && + !status.index.fresh && + !exclude.has(stepKey({ kind: "update-index", target: INDEX_TARGET })); + if (indexStep) { + add({ + kind: "update-index", + target: INDEX_TARGET, + reason: status.index.freshness.state === "fresh" ? "the index is stale" : status.index.freshness.reason, + }); + } + const afterIndex = indexStep ? { indexAfter: runStart } : {}; + const noStamp = !status.index.stamp; + + // Whether a target's build belongs in the run, and why. + const buildWanted = ( + kind: StageKind, + t: PublishTargetStatus, + state: TargetState, + ): { want: boolean; reason: string } => { + const f = STAGES[kind].needs(needs, req(kind, t.target)); + if (f.state === "stale") return { want: true, reason: f.reason }; + if (f.state === "blocked") { + // No stamp yet: the index this run updates comes first. + if (indexStep && noStamp) return { want: true, reason: "built after the index's first update" }; + return { want: false, reason: f.reason }; + } + return { want: false, reason: "fresh" }; + }; + + // The deploy a policy asks for, after (or without) its build in this run. + const deployFor = (kind: StageKind, t: PublishTargetStatus, state: TargetState, building: boolean) => { + if (opts.deploys === "none") return; + const dest = policyDeploy(t.policy, settings.previewBranch); + if (!dest) return; + const where = dest.preview ? { preview: dest.preview } : {}; + const previewUrl = + dest.preview && t.cloudflareProject ? { previewUrl: previewAliasUrl(t.cloudflareProject, dest.preview) } : {}; + if (building) { + // Never deployable by configuration: said here, not left to the child. + const never = t.deployProblem ?? state.pagesProblem ?? null; + if (never) { + skipped.push({ kind, target: t.target, reason: never }); + return; + } + add({ kind, target: t.target, ...where, builtAfter: runStart, reason: "its build is in this run", ...previewUrl }); + return; + } + const f = STAGES[kind].needs(needs, req(kind, t.target, where)); + if (f.state === "stale") add({ kind, target: t.target, ...where, reason: f.reason, ...previewUrl }); + else if (f.state === "blocked") skipped.push({ kind, target: t.target, reason: f.reason }); + }; + + // --- sites: each one's build, then its deploy (by id) + const rows: { t: PublishTargetStatus; state: TargetState; build: PlanStep | null; deploys: boolean }[] = []; + for (const t of status.sites) { + const state = needs.sites[t.target]; + const policyOn = t.policy !== "off"; + const w = buildWanted("build-site", t, state); + const wantByPolicy = opts.builds === "stale" ? true : policyOn; + if (!wantByPolicy) { + if (w.want) skipped.push({ kind: "build-site", target: t.target, reason: "stale, and its policy is off" }); + continue; + } + if (!w.want && w.reason !== "fresh") skipped.push({ kind: "build-site", target: t.target, reason: w.reason }); + rows.push({ + t, + state, + build: w.want ? { kind: "build-site", target: t.target, ...afterIndex, reason: w.reason } : null, + deploys: policyOn, + }); + } + const builds = rows.filter((r) => r.build && !exclude.has(stepKey(r.build))); + if (settings.runner === "docker" && builds.length > 0) { + // The fan-out builds every stale site in containers, in ONE stage. + add({ + kind: "build-site", + target: ALL_TARGET, + runner: "docker", + ...afterIndex, + reason: `${builds.length} site${builds.length === 1 ? "" : "s"} to build (${namesList(builds.map((r) => r.t.target))}); the docker runner builds every stale site`, + }); + for (const r of rows) if (r.deploys) deployFor("deploy-site", r.t, r.state, r.build !== null); + } else { + for (const r of rows) { + if (r.build) add(r.build); + if (r.deploys) deployFor("deploy-site", r.t, r.state, r.build !== null); + } + } + + // --- the hub, then the homepage: their policies are settings.publish's + for (const [kind, deployKind, t, state] of [ + ["build-hub", "deploy-hub", status.hub, needs.hub], + ["build-homepage", "deploy-homepage", status.homepage, needs.homepage], + ] as const) { + if (t.policy === "off") continue; + let w = buildWanted(kind, t, state); + // Fresh by the current index, but the index this run updates moves the + // hub (its signature carries the stamp id) when a listed site's channels + // changed — and the homepage when any did. + if (!w.want && w.reason === "fresh" && indexStep && state.changedChannels.length > 0) { + const n = state.changedChannels.length; + w = { want: true, reason: `${n} channel${n === 1 ? "" : "s"} changed` }; + } + if (w.want) add({ kind, target: t.target, ...afterIndex, reason: w.reason }); + else if (w.reason !== "fresh") skipped.push({ kind, target: t.target, reason: w.reason }); + deployFor(deployKind, t, state, w.want); + } + return { runStart, steps, skipped }; +} 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); +} diff --git a/common/views/publishStatus.ts b/common/views/publishStatus.ts @@ -0,0 +1,41 @@ +// THE PUBLISH STATUS, as the view layer names it (release 18). +// +// The builder and the planner are pure and live in the publish layer +// (publish/publishPlan.ts) because the lane's runner plans from them, and the +// publish layer may not import views/. This module is the view layer's door to +// them: the /sites Publish panel, `GET /api/ops/publish` and the CLI's +// `publish status` import the payload TYPES from here, and the reader that +// fills it is publish/publishState.ts `readPublishStatus`. +// +// A client component imports TYPES only: the builder's module graph reaches +// the stamp readers. + +export { + agoText, + buildPublishStatus, + laneBlockedReason, + needsInputOf, + planPublishRun, + policyDeploy, + publishPassDecision, + stepKey, +} from "../publish/publishPlan"; +export type { + BuildStale, + Chip, + ChipTone, + PassDecision, + PlanOptions, + PlanSkip, + PlanStep, + PublishInputs, + PublishJob, + PublishJobEnded, + PublishLaneLive, + PublishLaneStatus, + PublishPlan, + PublishSiteInput, + PublishStatus, + PublishTargetInput, + PublishTargetStatus, +} from "../publish/publishPlan"; diff --git a/editor/CHANGELOG.md b/editor/CHANGELOG.md @@ -1,6 +1,8 @@ # Changelog ## [Unreleased] +- **The publish lane.** Publishing can run itself: turn on `publish.enabled` in settings and the lane checks every `checkEveryMinutes` (10) whether the index is stale; when it is — and its last update is at least `refreshEveryMinutes` (360) old — it updates it, then builds every site whose channels changed or whose data the new index moved, one stage at a time on the `publish` queue. What it may do with a site is the site's own: `site.json` `publish.auto` is `off` (the default: left alone), `build`, `preview` (built and deployed to the preview branch `publish.previewBranch`) or `production`; the hub and the homepage have `publish.hub` and `publish.homepage`. A private site is only ever built, and a site needs its Cloudflare Pages project before it may deploy. Hold the lane and the stage running finishes and no next one starts; quiet hours (`publish.quietHours`) do the same; Drain finishes the stage and ends the runner. The lane never forces a stage: a stage that finds its target current does nothing. On /jobs every stage of one run reads `run <id> · <target>`, and a stage still queued when the editor restarts is cancelled, never re-queued — the lane works out again what is stale from what is on disk. `archilyzer publish now` runs the same plan from the command line, one stage after another in its own process. +- **One index for every site.** The index is updated once and every site, the hub and the homepage are built from it; `archilyzer publish status` says, per site, whether its build is current — "stale: 3 channels changed (a, b, c)" as soon as a download, transcription or digest on one of its channels finishes, before any index runs; "stale: data changed" once the index has run and the site's data moved; "stale: config changed" after its site.json, tags or aliases changed — and whether what is deployed is that build, with a build made by older code marked "code newer" but not stale. - **Deploys are pinned and checked live.** wrangler is an exact dependency of the workspace (4.147.0), so a deploy runs the version installed with the code instead of whatever `pnpm dlx` fetched that day, and every deploy names its branch: production is `--branch main`, never taken from the checkout it ran in (where a "production" deploy from a feature branch used to land as a preview). The publish stages' deploy (release 18) refuses before wrangler runs when there is no Cloudflare credential at all — "set CLOUDFLARE_API_TOKEN in .env" — and says "REFUSED by Cloudflare — the API token was not accepted" when Cloudflare rejects one; it refuses a production deploy of a build made from a branch other than `main`. After each deploy it reads `corpus.json` at the site's address twice, as a visitor would and cache-busted, and records the verdict: ok, stale-edge (the deployment is right, Cloudflare's edge still serves an older copy), mismatch, or unreachable. A verdict short of ok is a warning in the log; the deploy itself succeeded. What each target last shipped, where, and how it read is kept in `deployed.json` beside its build. - **Withdrawn X posts ship tombstones.** While X posts are private, a public site's build no longer just leaves an X channel's posts out: at every path they were served from it ships an empty stand-in — the channel's posts manifest with no pages, and an empty page for each page the channel has — served uncached. The hub, which carries no posts, ships the same for every X channel a public site carries, with an empty posts manifest; a channel only on a private site, or on no site, is never named on the hub. Leaving a path out of a deploy does not take it off Cloudflare's edge, which kept serving a withdrawn copy for up to a week; a changed object at the same path replaces it. The hub's deploy reads each of those paths back. - **Publishing is stages, from the command line: `archilyzer publish`.** `publish index` updates the index — the LMDB index, the stats datasets and the chart templates, in one child process with an 8 GB heap — and writes an index stamp (`export/.export-index/stamp.json`) naming, for each site, a signature of everything that site's build reads. `publish build <id|all>` builds a site from that index (no data phase of its own) into its own bundle, `export/.export-builds/<id>/out`, and stamps it (`built.json`); a site whose bundle already matches the index is a no-op unless `--force`. `publish deploy <id|all> [--preview <branch>] [--to local]` ships that bundle — to Cloudflare Pages, or with `--to local` into the directory the docker `site` service serves — and records the deploy (`deployed.json`); deploying the same build again is a no-op unless `--force`. `all` passes over private sites and, to Pages, sites with no Pages project; any other site it cannot deploy is a failure, said after the rest are tried. `publish hub [--deploy]` and `publish homepage [--deploy]` do the same for the hub (`_hub/out`) and the homepage. A stage whose input is not there says so and exits 3: "update the index first", "no build of jeralyzer — archilyzer publish build jeralyzer". Production refuses a bundle built on a branch other than `main`, or with no branch recorded (a detached checkout; an image sets `ARCHILYZER_BRANCH`) — a preview of it is fine. Exit codes: 0 done or nothing to do, 1 failed, 2 usage, 3 precondition not met, 130 cancelled. diff --git a/editor/app/components/lanes/pauseControl.tsx b/editor/app/components/lanes/pauseControl.tsx @@ -76,6 +76,21 @@ const PAUSE_COPY: Record<PauseLane, { held: PauseCopy; free: PauseCopy }> = { "Resume digest generation. The running job picks up where it left off — it was holding, not stopped.", }, }, + // The publish lane (release 18): a hold stops it dispatching the next stage; + // the stage in flight finishes. + publish: { + free: { + label: "Hold the lane", + ariaLabel: "pause publishing", + title: + "Hold the publish lane: the stage in flight finishes and no next one starts. Manual Build, Deploy and Publish now still work.", + }, + held: { + label: "Resume the lane", + ariaLabel: "resume publishing", + title: "Resume the publish lane. Its next check decides whether a pass is due.", + }, + }, backfill: { free: { label: "Hold the lane", diff --git a/editor/app/jobs/jobReplayRegistry.ts b/editor/app/jobs/jobReplayRegistry.ts @@ -56,6 +56,10 @@ import { import { replayFetchWindowAction } from "../channels/[slug]/videos/[id]/videoActions"; import { reportsPrepareAction } from "../sites/lib/reportsPrepareAction"; import { reportsExportAction } from "../sites/lib/reportsExportAction"; +import { + enqueueStage, + stageRequestFromSpec, +} from "yt-dlp-transcript-common/publish/publishStages"; export type ReplayHandler = (spec: JobSpec) => Promise<StreamActionResult>; @@ -90,7 +94,23 @@ function params(spec: JobSpec): { return { p, queueKey: str(p.queueKey) }; } +// A PUBLISH STAGE (release 18): the spec IS the stage request, so a retry +// re-enqueues it as it was — same run id, same preconditions, which the child +// re-asks on disk. A spec the stage row would refuse is refused here. +const replayPublishStage: ReplayHandler = async (spec) => { + const req = stageRequestFromSpec(spec); + if (!req) return { ok: false, error: `Not a publish stage this editor can re-run (${spec.kind} ${spec.slug})` }; + return enqueueStage(getPaths(), req); +}; + export const JOB_REPLAY_HANDLERS: Record<string, ReplayHandler> = { + "publish-update-index": replayPublishStage, + "publish-build-site": replayPublishStage, + "publish-deploy-site": replayPublishStage, + "publish-build-hub": replayPublishStage, + "publish-deploy-hub": replayPublishStage, + "publish-build-homepage": replayPublishStage, + "publish-deploy-homepage": replayPublishStage, // A report site's evidence media. The spec's slug is the SITE (the job has // no channel); a re-run re-reads the site's reports and cuts only what its // cache does not already hold. diff --git a/editor/app/operations/actions.ts b/editor/app/operations/actions.ts @@ -197,7 +197,9 @@ async function setLaneHeld( try { const cur = getSettings(); if (isGateHeld(cur, lane) !== held) { - await saveSettings({ autoQueue: withGateHeld(cur, lane, held).autoQueue }); + const next = withGateHeld(cur, lane, held); + // The publish lane's gate is its own block (release 18), not a policy. + await saveSettings(lane === "publish" ? { publish: next.publish } : { autoQueue: next.autoQueue }); } } catch (e) { // REPORTED, not swallowed. The workers page's old best-effort persist diff --git a/editor/instrumentation.ts b/editor/instrumentation.ts @@ -223,6 +223,19 @@ export async function register() { /* a runner that fails to start must not block server readiness */ } + // THE PUBLISH LANE (release 18): its runner, when settings.publish.enabled. + // Beside the four above rather than inside startAutoRunnersIfEnabled — the + // dispatch layer may not import the publish layer — and, like them, below + // the idle gate: an idle boot leaves it off. + try { + const { startPublishRunnerIfEnabled } = await import( + "yt-dlp-transcript-common/publish/publishRunner" + ); + await startPublishRunnerIfEnabled(); + } catch { + /* a runner that fails to start must not block server readiness */ + } + // NO SECOND RESUME HOOK. There used to be two more here, one per sweep, and // each re-launched a corpus-wide pass from its own persisted flag. Both lanes // are runner lanes now, so `startAutoRunnersIfEnabled` above resumes all four diff --git a/plans/release-18.md b/plans/release-18.md @@ -1127,6 +1127,158 @@ the image built beside it; the file alone afterwards: 8/8 (S5 touches nothing it S5 is complete with this round. Left for the rollout, as recorded above: building `runtime-vulkan` and `runtime-cuda` once. +### Slice S3, as shipped — the publish status, the stage queue and the publish lane (2026-10-06) + +Branch `r18/publish-lane` off `r18/integration` `0ce00f76` (S1 and S5's image half merged; `1f07ca2f` — S2 and +S5 complete — merged in before the gates), worktree `~/Projects/r18-publish-lane` (editor 7201, test 7211, export +7210), one Opus implementer. Scratch files `s3-*` in the job's `tmp`. The plan is "Stages" (the two staleness +bullets), "Lane", "Surfaces: One state builder", the CLI's `publish status|now` and "Migration"; S1's "Seams for +the other slices" are what it codes against. No S1 file changed. + +**What it does.** +- **One status, one plan — `common/publish/publishPlan.ts` (pure).** `buildPublishStatus(inputs, now)` folds the + inputs into `PublishStatus`: per target (each site, `_hub`, `_homepage`) `indexFresh`, `built`, + `builtFromCurrentIndex`, the build stage's own `needs()` answer, `buildStale {reason: channels | data | config, + changedChannels}`, `codeNewer` (never stale), `deployed`, the policy's `deployKind` and record, + `deployedIsBuilt`, `previewUrl`, `liveCheck`, `policy`, `deployProblem`, `next`, `running`, `queued[]`, the + newest ended build/deploy job, and four chips (`index | built | deployed | live`, each `{tone, text}`); plus the + index (chip, freshness, its jobs), the lane (`due` + `reason`, `blockedReason`, chip) and `plan` — what Publish + now would run. Every freshness question is the stage's `needs()` over a `NeedsInput` built here; the status adds + only what a stage child cannot see: `changedChannels` (a member channel with an ingest ended after + `builtCheckedAt(built)` — so a no-op build clears it; the hub's are the listed sites' members, the homepage's + every channel) and the config mtimes. Chip words: "no index yet", "stale: new data since the last index", + "never built", "stale: N channels changed (a, b, …)", "stale: data changed", "stale: config changed", "built + 12 min ago · code newer", "production · 3 min ago", "production: a newer build is not deployed", "never + deployed (private)", and — from the newest ended stage, exit 3 — "waiting for its build" (a deploy that carried + `builtAfter`), "waiting for the index" (a build that carried `indexAfter`), else "last deploy refused: + precondition not met"; live: "live ok", "live: the edge serves an older build", …. +- **`planPublishRun(status, {deploys: policy | none, builds: policy | stale, runStart, exclude, skipIndex})`**: + update-index when the index is stale; per site (by id) its build when the stage's `needs()` says stale — changed + channels, a signature that no longer matches the index, never built, its config, its bundle — and its deploy + where `publish.auto` says (`preview` → `settings.publish.previewBranch`, with the alias URL); then the hub, then + the homepage, by `settings.publish.hub|homepage` (the index update in the run moves them when their channels + changed). Builds carry `indexAfter = runStart` when the index is in the run, deploys `builtAfter = runStart` + when their build is. A policy-off site is never built (`skipped`: "stale, and its policy is off") unless + `builds: "stale"` ("Build all stale"); a private or project-less target is never planned as a deploy (skipped + with the reason). Never forced. `settings.publish.runner: "docker"` plans one `build-site _all --runner docker`. + `views/publishStatus.ts` re-exports it: the view layer's name for it. +- **`common/publish/publishState.ts` `readPublishStatus(paths)`** — the one function the /sites panel, `GET + /api/ops/publish`, `archilyzer publish status` and the lane call. `readPublishInputs` reads the stamps and bundle + problems (S1's `readNeedsInput`), the sites (membership = `site.json` `channels`) and their policies, the config + mtimes of the plan's list (for the index tags.json, search-aliases.json, duplicates*.json, every site.json, + homepage.json, the settings file, 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, the live publish jobs, the newest ended stage per kind + target, the lane's memory, and the + running commit and `main`'s head (git, memoized 30 s). Job metas are read through the registry and the `.jobs` + sidecars, each terminal meta once per process. +- **`common/publish/publishStages.ts`.** `enqueueStage(paths, req, {background})`: one `runManagedCommand` over + S1's `stageCommand` on queue `publish`, spec `{kind: "publish-<stage>", slug: target, params: the request}`; a + duplicate — the same kind, target and destination queued or running — is refused `{ok: false, info: true, + jobId}`. `enqueuePublishRun(paths, plan, {runId})`: the plan in its order under one run id, each request + carrying its step's preconditions; returns `{runId, jobs: [{kind, target, jobId, previewUrl?, existing?}], + skipped, refused}`. `stageRequestFromSpec` reads a spec back through the stage row's own parser. +- **The lane — `common/publish/publishRunner.ts`.** `startPublishRunner` (kind `auto-publish`, queueKey `""`) + wakes every `checkEveryMinutes`; a pass is due when there is no stamp, when the index is stale and + `now − stamp.builtAt ≥ refreshEveryMinutes`, or when a policy target is left stale and the last pass is at least + that old; never while held, in quiet hours or while any publish stage is queued or running. A pass dispatches + ONE stage at a time (`background: true`) and awaits its job's end, re-reading the status and re-planning before + each next one (so the builds after its index update see the new stamp); a stage it ran is not run again in the + pass; 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, never a stage; a drain finishes the + stage in flight and ends the runner `done`; Stop ends the runner at once and leaves a running stage to finish as + its own job. `stopPublishRunner` / `drainPublishRunner` / `startPublishRunnerBlockedReason`; + `publishLaneState.ts` holds the lane's memory (last check, next check, last pass and its summary). The editor's + `instrumentation.ts` starts it below the idle-boot gate. +- **Settings and site.json.** `settings.publish {enabled: false, held: false, checkEveryMinutes: 10 [1–1440], + refreshEveryMinutes: 360 [0–43200; 0 = whenever stale], quietHours: {start, end} | null, runner: local, + previewBranch: "preview" (an invalid name reads as it), hub: off, homepage: off}`; `site.json` `publish: {auto: + off | build | preview | production}` — absent = off, only another policy written; a private site reads and is + written as `build`; a site with no `cloudflareProject` reads as `build`, and a save of `preview`/`production` + without one is refused ("publish.auto "production" deploys the site, and it has no cloudflareProject — …"). + SETTINGS.md, settings.json.example and SITE.md regenerated. +- **The pipeline lane.** `PIPELINE_LANES = ["publish"]` (`autoQueueTypes.ts`), not in `LANES`; `PauseLane = + AutoQueueKind | "publish"`; `isGateHeld` / `withGateHeld` read and write `settings.publish.held`. The editor's + lane pause action saves that block for `publish`, and the pause control has its words ("Hold the lane", + `pause publishing` / `resume publishing`). +- **Jobs.** The seven `publish-*` kinds (labels = the stages', replayable, not drainable) and `auto-publish` + (drainable); `build-index`, `build-stats`, `build-export`, `build-deploy`, `build-all`, `build-deploy-all`, + `deploy-export` known with NO label (/jobs shows their raw kind, as it always has; `jobs.spec` reads it) beside + the six hub/homepage kinds that keep theirs. `isIngestKind`. /jobs reads a stage as `run <last six of the run + id> · <target>`. A queued publish stage at boot is cancelled as `publish` — "server restarted; the publish lane + re-derives stages from on-disk state" — never re-queued. Retry re-enqueues a stage from its spec (same run id). +- **CLI.** `archilyzer publish status [--json]` (the index, the lane, a row per target with its policy and chips + and what is next, then the plan) and `archilyzer publish now` — the plan's stages one at a time IN THE CLI's + PROCESS under the publish lock, the index update as `publish index`'s child, as the editor's queue would run + them (the CLI has no queue); the rest are tried after a failure; exit 0, or the worst (1 over 3). + +**Deviations from the plan** (one sentence each): +1. The three modules are `common/publish/{publishState,publishStages,publishRunner}.ts`, not `controller/`, and + the pure builder is `publish/publishPlan.ts` re-exported by `views/publishStatus.ts`: `architecture.test.ts` + forbids controller → publish and publish → views, and the runner must plan. +2. The runner is started from `instrumentation.ts` beside `startAutoRunnersIfEnabled`, not inside it (dispatch may + not import publish); the idle boot leaves it off the same way. +3. "The drainable kinds" became an explicit `isIngestKind` set: every drainable per-channel kind (tested) plus the + non-drainable writers of the same text and posts (the auto-download unit, single imports, transcriptions and + downloads, the availability checks, the forum import, the feed backfill, the cues sweep); the lane runners' + metas stay `running` for days and are not in it. +4. A channel's `snapshot.json` mtime is also an ingest signal: the transcription, digest and backfill lanes' units + make no job record at all, and each requests the channel's report regeneration when it settles. +5. A duplicate stage is the same kind, target AND destination: a production deploy behind a preview of the same + site is not refused. +6. A pass is also due when a policy target is left stale (a failed stage, a policy just turned on) and the last + pass is `refreshEveryMinutes` old — else a fresh index would never retry it. +7. The lane re-plans before each stage instead of enqueueing a plan computed once. +8. Stop does not cancel the stage in flight (it is its own job, with its own Cancel); the plan did not say. +9. Files beyond the slice's list, each a line or a table entry: `common/lib/site.ts` (writeSite's refusal), + `common/lib/{settingsDocs,fileSchemaDocs}.ts` (the doc tables), `editor/app/operations/actions.ts` and + `editor/app/components/lanes/pauseControl.tsx` (the widened `PauseLane` needs its save and its words). + +**Exports added to S1's files:** none. + +**Found, not fixed.** +- The settings file is an index input (the plan's list), and the settings file is written by every pause click, + priority change and drive auto-pause: each makes the index "stale" — one short-circuited index update per + refresh interval, no rebuild (the signatures are unchanged). +- The snapshot signal is loose: a report regenerates after any channel job, failed ones included; the cost is a + short-circuited index and no-op builds. + +| commit | what | +|---|---| +| `8b55064a` | jobs: the seven `publish-*` kinds, `auto-publish`, the label-only kinds; `isIngestKind`; `run <id> · <target>`; boot category `publish` | +| `1df0e276` | settings: `settings.publish`, `site.json` `publish.auto`, `PIPELINE_LANES`, the publish gate; SETTINGS.md, SITE.md, example | +| `105c2663` | publish: `publishPlan.ts` — the status builder and the planner; `views/publishStatus.ts` | +| `3adafcb7` | merge `r18/integration` `1f07ca2f` (S2, S5) | +| `0081d1dd` | publish: `readPublishStatus`, `enqueueStage` / `enqueuePublishRun`, the runner, the lane's memory | +| `344ffee0` | editor: the runner at boot; Retry for publish stages | +| `0527be27` | cli: `publish status [--json]`, `publish now` | + +**Gates** (all from the worktree root): tsc clean at every commit; common **3339 passed** (57 new: publishPlan 23, +publishRunner 10, publishStages 6, publishState 2, publishNow 2, jobKinds 3, jobDetail 3, bootQueuedJobs 1, +settingsSchema 2, siteSchema 4, pauseGates 1); editor unit **142**; `test:scripts` **596 + 3 skipped**; mcp +**289**; export unit **116**; homepage unit **23**; `pnpm --filter editor exec next build` ok (51 s); `pnpm +--filter export exec next build` ok (27 s, over the committed fixture compose linked into the worktree's +`export/public` — the primary's holds a reports-only compose); `pnpm --filter homepage run build:nodata` ok +(15 s); umtool's capped build ok (19 s, link removed). e2e (editor suite, `s3-specs.txt`: lane-runner, auto-queue, jobs, jobs-filters, settings, sites-crud, site-scope, build, scheduler, operation-settings; the worktree's `export/public` seeded with the fixture compose, cleaned after): **88 passed, 0 failed, 6.2 min**. Numbers tool: none. CLI smoke over a scratch +corpus (`s3-smoke.sh`, one site on `build`): `publish status` → "no index yet", plan `update-index _index`, +`build-site smoke`; `publish now` → both ran (40 s, the build under `indexAfter`); `publish status` → "fresh", +"built just now", nothing to run; `publish now` again → "nothing to do"; the policy switched off → "stale: config +changed", "stale, and its policy is off". The scratch compose was removed from `export/public` after. + +**Left for the other slices.** S4: the /sites Publish panel, `/operations/publish`, `GET|POST /api/ops/publish` +read and call the names in "Seams" below; S4's e2e sets policies through `site.json` / `settings.publish`. +S5's doctor: the `source-repo` grade can key on `settings.publish.homepage` now. S6: PUBLISH.md (the lane, the +policies, `publish status|now`), FACTS. + +**Seams (the names S4 calls).** `readPublishStatus(paths?, {now?, laneKnown?})` → +`PublishStatus` (`publish/publishState.ts`; types from `views/publishStatus.ts`); `planPublishRun(status, +{deploys, builds, runStart?})` (Publish now = `{deploys: "policy"}`; Build all stale = `{builds: "stale", +deploys: "none"}`); `enqueuePublishRun(paths, plan, {runId?})` and `enqueueStage(paths, req, {background?})` +(`publish/publishStages.ts`, with `newPublishRunId`, `PUBLISH_QUEUE`, `stageRequestFromSpec`); the lane: +`startPublishRunner(paths?)`, `stopPublishRunner()`, `drainPublishRunner()`, +`startPublishRunnerBlockedReason(settings?)` (`publish/publishRunner.ts`), its hold `pauseLaneAction("publish")` / +`resumeLaneAction("publish")` (`editor/app/operations/actions.ts`) and `isGateHeld(settings, "publish")`; the +policies `sitePublishPolicy(site)`, `sitePublishProblem(site)`, `settings.publish`. + ## Rollout (Steps 1–7 above; "### As it went" is written as the rollout runs.) diff --git a/settings.json.example b/settings.json.example @@ -196,5 +196,16 @@ "maxPeers": 30, "maxDownloadKiBps": 0, "maxUploadKiBps": 0 + }, + "publish": { + "enabled": false, + "held": false, + "checkEveryMinutes": 10, + "refreshEveryMinutes": 360, + "quietHours": null, + "runner": "local", + "previewBranch": "preview", + "hub": "off", + "homepage": "off" } }