// 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 and the charts config — the SETTINGS by the keys the index // reads (inputSig.ts indexSettingsSig), never the file's mtime, which // every pause click moves; 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//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 { indexSettingsSig } from "./inputSig"; 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 { 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): Promise { 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(); function lite(r: Pick): 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, live: readonly JobRecord[] = getRegistry().list(), ): Promise { const byId = new Map(); 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; }; /** The newest ingest end per channel (and overall) from the job metas. */ export function ingestFromMetas(metas: readonly MetaLite[]): IngestSignals { const byChannel: Record = {}; 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> { const out: Record = {}; 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(); 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(); const GIT_MEMO_MS = 30_000; async function memoized(key: string, now: number, read: () => Promise): Promise { 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 { 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 = { ...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(); 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.chartsConfigFile), ); // The settings are judged by the keys the index reads, not the file's mtime // (every pause click writes the file). const settingsSig = indexSettingsSig(settings); 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, settingsSig }, 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 { const now = opts.now ?? Date.now(); return buildPublishStatus(await readPublishInputs(paths, { ...opts, now }), now); }