Files
ThothII/backend/src/workspaces/registry-publication.ts
T

217 lines
28 KiB
TypeScript

import { createHash, randomBytes } from "node:crypto";
import { constants as fsConstants } from "node:fs";
import { lstat, mkdir, open, readFile, readdir, rename, rm, stat, unlink, writeFile } from "node:fs/promises";
import { dirname, join } from "node:path";
import type { CanonicalWorkspaceId, Revision40, Sha256Hex } from "./workspace-lock-root-lease.js";
import type { BorrowedOrderedWorkspaceWriterLeaseV1, BorrowedWorkspaceSessionReadersExclusiveLockLease, OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js";
export type RegistryRunId32 = string & { readonly __registryRunId32: unique symbol };
export type RegistryAddressedOperationV1 = "registry_bootstrap" | "registry_pull";
export type RegistryAddressedPublicationPhaseV1 = "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "participants_prepared" | "publication_intent_durable" | "target_published" | "terminal_durable";
export type RegistryAddressedJobArtifactPathV1 = `addressed-publication-jobs/${RegistryRunId32}.json`;
export interface RegistryWorkspaceManifestIdentityV1 { readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly descriptorBlob: Revision40; readonly manifestSha256: Sha256Hex; }
export interface RegistryActiveSnapshotV1 { readonly schemaVersion: 1; readonly commit: Revision40; readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[]; }
interface PlanFields { readonly schemaVersion: 1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly advertisedTargetCommit: Revision40; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target`; readonly fetchedTargetCommit: Revision40; readonly targetCommit: Revision40; readonly targetManifestSha256: Sha256Hex; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[]; readonly changedSetSha256: Sha256Hex; }
export interface RegistryBootstrapAddressedPlanV1 extends PlanFields { readonly operation: "registry_bootstrap"; readonly changedSetRule: "all_target_workspace_ids"; readonly baseCommit: null; readonly baseManifestSha256: null; readonly baseWorkspaces: readonly []; }
export interface RegistryPullAddressedPlanV1 extends PlanFields { readonly operation: "registry_pull"; readonly changedSetRule: "symmetric_base_target_workspace_difference"; readonly baseCommit: Revision40; readonly baseManifestSha256: Sha256Hex; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; }
export type RegistryAddressedPlanV1 = RegistryBootstrapAddressedPlanV1 | RegistryPullAddressedPlanV1;
interface StateFields { readonly schemaVersion: 1; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly phase: RegistryAddressedPublicationPhaseV1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly advertisedTargetCommit: Revision40 | null; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target` | null; readonly fetchedTargetCommit: Revision40 | null; readonly targetCommit: Revision40 | null; readonly targetManifestSha256: Sha256Hex | null; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[] | null; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[] | null; readonly planSha256: Sha256Hex | null; readonly changedSetSha256: Sha256Hex | null; readonly participantsSha256: Sha256Hex | null; readonly synchronizersSha256: Sha256Hex | null; readonly publicationIntentSha256: Sha256Hex | null; readonly publishedActiveStateSha256: Sha256Hex | null; readonly terminalResultSha256: Sha256Hex | null; readonly priorStateSha256: Sha256Hex | null; readonly baseCommit: Revision40 | null; readonly baseManifestSha256: Sha256Hex | null; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "all_target_workspace_ids" | "symmetric_base_target_workspace_difference" | null; }
export interface RegistryAddressedSnapshotV1 extends RegistryActiveSnapshotV1 { readonly requestDigest: string; }
export interface RegistryInstallationIdentityV1 { readonly installationId: string; readonly digest: string; }
export interface RegistryRepositoryIdentityV1 { readonly remote: string; readonly branch: string; readonly head: string; readonly digest: string; }
export interface RegistryRemoteIdentityV1 { readonly remote: string; readonly head: string; readonly digest: string; }
export interface RegistryAddressedPublicationResultV1 { readonly runId: string; readonly snapshot: RegistryAddressedSnapshotV1; readonly result?: unknown; }
export interface RegistryBootstrapAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_bootstrap"; readonly baseCommit: null; readonly baseManifestSha256: null; readonly baseWorkspaces: readonly []; readonly changedSetRule: "all_target_workspace_ids" | null; }
export interface RegistryPullAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_pull"; readonly baseCommit: Revision40; readonly baseManifestSha256: Sha256Hex; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "symmetric_base_target_workspace_difference" | null; }
export type RegistryAddressedPublicationStateV1 = RegistryBootstrapAddressedPublicationStateV1 | RegistryPullAddressedPublicationStateV1;
export class BorrowedWorkspaceMaintenanceQuiescenceLease { private constructor(readonly workspaceId: CanonicalWorkspaceId) {} static create(id: CanonicalWorkspaceId) { return new BorrowedWorkspaceMaintenanceQuiescenceLease(id); } }
export interface AddressedWorkspacePublicationLeaseV1 extends BorrowedOrderedWorkspaceWriterLeaseV1 { readonly quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease; readonly readers: BorrowedWorkspaceSessionReadersExclusiveLockLease; }
export interface CapabilityAwareRegistryPublicationParticipant<T> { readonly participantId: string; prepare(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1): Promise<T>; reconcile(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1, prepared: T, phase: RegistryAddressedPublicationPhaseV1): Promise<void>; }
export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: string; readonly preparedSha256: Sha256Hex; }
export interface CapabilityAwareRegistryPublicationSynchronizer { readonly synchronizerId: string; ensureForPublication(plan: RegistryAddressedPlanV1, capabilities: OrderedWorkspaceWriterCapabilitySet, phase: "planned" | "participants_prepared" | "publication_intent_durable" | "target_published"): Promise<RegistrySynchronizerPreparedV1>; }
export class CapabilityAwareRegistryPublicationLifecycleOwner {
constructor() {}
async run<T>(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant<unknown>[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise<T>; }): Promise<T> {
// Reader gates are recursively nested so every changed workspace remains quiescent and
// reader-exclusive until the publication callback, reconciliation, and terminal durability
// have all settled. No participant or synchronizer is allowed to acquire a capability.
const prepared: Array<{ readonly participant: CapabilityAwareRegistryPublicationParticipant<unknown>; readonly value: unknown; readonly lease: AddressedWorkspacePublicationLeaseV1 }> = [];
const enter = async (index: number): Promise<T> => {
if (index >= input.capabilities.workspaceIds.length) {
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "planned");
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "participants_prepared");
const result = await input.action();
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "target_published");
for (const item of prepared) await item.participant.reconcile(input.plan, item.lease, item.value, "target_published");
return result;
}
const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId;
return input.capabilities.forWorkspace(id, ({ rootLease, writerCapability }) =>
writerCapability.runUnderSessionReadersExclusive(async readers => {
const lease: AddressedWorkspacePublicationLeaseV1 = {
workspaceId: id,
rootLease,
writerCapability,
quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id),
readers,
};
for (const participant of input.participants) {
const value = await participant.prepare(input.plan, lease);
prepared.push({ participant, value, lease });
}
return enter(index + 1);
}),
);
};
return enter(0);
}
}
export type RegistryAddressedRequestV1 =
| { readonly mode: "create"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: null; readonly remoteRefIdentitySha256: Sha256Hex }
| { readonly mode: "resume"; readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex }
| { readonly mode: "create"; readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly expectedBaseCommit: Revision40; readonly remoteRefIdentitySha256: Sha256Hex }
| { readonly mode: "resume"; readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex };
export interface RegistryBootstrapAddressedResultV1 { readonly operation: "registry_bootstrap"; readonly runId: RegistryRunId32; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly plan: RegistryBootstrapAddressedPlanV1; readonly planSha256: Sha256Hex; readonly phase: "terminal_durable"; readonly publication: "target" | "reconciled_target" | "unchanged"; }
export interface RegistryPullAddressedResultV1 { readonly operation: "registry_pull"; readonly runId: RegistryRunId32; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly plan: RegistryPullAddressedPlanV1; readonly planSha256: Sha256Hex; readonly phase: "terminal_durable"; readonly publication: "target" | "reconciled_target" | "unchanged"; }
export type RegistryAddressedResultV1 = RegistryBootstrapAddressedResultV1 | RegistryPullAddressedResultV1;
export interface RegistryBootstrapRecoveryIdentityV1 { readonly operation: "registry_bootstrap"; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; }
export type RegistryEnsureBootstrapAddressedResultV1 = { readonly kind: "already_active"; readonly snapshot: RegistryActiveSnapshotV1 } | { readonly kind: "bootstrap_terminal"; readonly snapshot: RegistryActiveSnapshotV1; readonly result: RegistryBootstrapAddressedResultV1 };
export interface RegistryBootstrapRecoveryScanLimitsV1 { readonly maximumDirectoryEntries: 4096; readonly maximumArtifactBytes: 1048576; readonly maximumTotalArtifactBytes: 67108864; }
export const REGISTRY_SCAN_LIMITS_V1: RegistryBootstrapRecoveryScanLimitsV1 = Object.freeze({ maximumDirectoryEntries: 4096, maximumArtifactBytes: 1048576, maximumTotalArtifactBytes: 67108864 });
const RUN = /^[0-9a-f]{32}$/; const SHA = /^[0-9a-f]{64}$/; const REV = /^[0-9a-f]{40}$/;
const PHASES: readonly RegistryAddressedPublicationPhaseV1[] = ["request_claimed", "target_advertised", "target_fetched", "planned", "participants_prepared", "publication_intent_durable", "target_published", "terminal_durable"];
const CONFLICT = () => Object.assign(new Error("preprocessing_conflict"), { code: "preprocessing_conflict" });
const canonical = (v: unknown): string => JSON.stringify(v, (_k, x) => x && typeof x === "object" && !Array.isArray(x) ? Object.fromEntries(Object.keys(x).sort().map(k => [k, x[k]])) : x);
export const registryDigest = (value: unknown): string => createHash("sha256").update(canonical(value)).digest("hex");
export function canonicalBootstrapRequestDigest(request: { readonly kind?: "bootstrap" | "publish"; readonly operation?: RegistryAddressedOperationV1; readonly installation?: unknown; readonly repository?: unknown; readonly remote?: unknown; readonly workspaceIds?: readonly string[] }): string { return registryDigest({ kind: request.kind ?? (request.operation === "registry_pull" ? "publish" : "bootstrap"), operation: request.operation, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...(request.workspaceIds ?? [])].sort() }); }
export function addressedRunId(): RegistryRunId32 { return randomBytes(16).toString("hex") as RegistryRunId32; }
function failIfBadIdentity(s: StateFields, runId: string): void { if (s.schemaVersion !== 1 || s.runId !== runId || !RUN.test(s.runId) || s.jobArtifactPath !== `addressed-publication-jobs/${s.runId}.json` || !SHA.test(s.requestSha256) || !SHA.test(s.installationIdentitySha256) || !SHA.test(s.repositoryIdentitySha256) || !SHA.test(s.remoteRefIdentitySha256) || !PHASES.includes(s.phase)) throw CONFLICT(); }
function immutable(s: StateFields): unknown { const { phase: _p, priorStateSha256: _h, ...rest } = s; return rest; }
async function fsync(path: string): Promise<void> { const h = await open(path, "r"); try { await h.sync(); } finally { await h.close(); } }
async function fsyncParent(path: string): Promise<void> { await fsync(dirname(path)); }
function ownerMode(st: Awaited<ReturnType<typeof stat>>, mode: number): boolean { const x = st as any; return x.isFile() && (Number(x.mode) & 0o777) === mode && Number(x.nlink) === 1 && Number(x.uid) === (process.getuid?.() ?? Number(x.uid)); }
async function strictRead(path: string, max = REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes): Promise<{ text: string; identity: { size: number; mtimeMs: number; ino: bigint } }> {
const h = await open(path, fsConstants.O_RDONLY | (fsConstants.O_NOFOLLOW ?? 0));
try { const before = await h.stat(); if (!ownerMode(before, 0o600) || before.size > max) throw CONFLICT(); const text = await h.readFile({ encoding: "utf8" }); const after = await h.stat(); if (before.ino !== after.ino || before.size !== after.size || text.length > max) throw CONFLICT(); return { text, identity: { size: Number(after.size), mtimeMs: Number(after.mtimeMs), ino: BigInt(after.ino) } }; } finally { await h.close(); }
}
export class RegistryAddressedPublicationStore {
readonly jobsDirectory: string;
constructor(readonly root: string) { this.jobsDirectory = join(root, "addressed-publication-jobs"); }
private path(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.json`); }
private resultPath(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.result.json`); }
private async dirs(): Promise<void> { await mkdir(this.jobsDirectory, { recursive: true, mode: 0o700 }); const st = await lstat(this.jobsDirectory); if (st.isSymbolicLink() || !st.isDirectory() || (Number((st as any).mode) & 0o777) !== 0o700 || Number((st as any).uid) !== (process.getuid?.() ?? Number((st as any).uid))) throw CONFLICT(); }
private async durable(path: string, value: unknown, exclusive = false): Promise<void> { const name = path.split("/").pop()!; const tmp = join(this.jobsDirectory, `.${name}.tmp`); if (exclusive) { try { await stat(path); throw CONFLICT(); } catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw CONFLICT(); } } const bytes = `${canonical(value)}\n`; try { const h = await open(tmp, fsConstants.O_WRONLY | fsConstants.O_CREAT | fsConstants.O_EXCL | (fsConstants.O_NOFOLLOW ?? 0), 0o600); try { await h.writeFile(bytes); await h.sync(); } finally { await h.close(); } const st = await stat(tmp); if (!ownerMode(st, 0o600)) throw CONFLICT(); if (!exclusive) { const current = await lstat(path); if (current.isSymbolicLink() || !ownerMode(current, 0o600)) throw CONFLICT(); } await rename(tmp, path); await fsyncParent(path); } catch (e) { await rm(tmp, { force: true }).catch(() => undefined); if (exclusive && (e as NodeJS.ErrnoException)?.code === "EEXIST") throw CONFLICT(); throw CONFLICT(); } }
async claim(request: RegistryAddressedRequestV1, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise<RegistryAddressedPublicationStateV1> {
await this.dirs();
if (!RUN.test(runId)) throw CONFLICT();
const now = new Date().toISOString();
const state = { schemaVersion: 1, runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation: request.operation, installationIdentitySha256: request.installationIdentitySha256, repositoryIdentitySha256: request.repositoryIdentitySha256, remoteRefIdentitySha256: request.remoteRefIdentitySha256, advertisedTargetCommit: null, immutableTargetRef: null, fetchedTargetCommit: null, targetCommit: null, targetManifestSha256: null, targetWorkspaces: null, changedWorkspaceIds: null, planSha256: null, changedSetSha256: null, participantsSha256: null, synchronizersSha256: null, publicationIntentSha256: null, publishedActiveStateSha256: null, terminalResultSha256: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null, baseManifestSha256: null, baseWorkspaces: [], changedSetRule: null };
await this.durable(this.path(runId), state, true); return state as unknown as RegistryAddressedPublicationStateV1;
}
async setBase(runId: RegistryRunId32, baseManifestSha256: Sha256Hex, baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]): Promise<RegistryAddressedPublicationStateV1> {
const old = await this.read(runId);
if (old.phase !== "request_claimed" || old.operation !== "registry_pull") throw CONFLICT();
const next = { ...old, baseManifestSha256, baseWorkspaces, priorStateSha256: registryDigest(old) as Sha256Hex };
await this.durable(this.path(runId), next);
return next;
}
async storeTerminalResult(runId: RegistryRunId32, result: RegistryAddressedResultV1): Promise<void> {
await this.dirs();
const path = this.resultPath(runId);
await this.durable(path, { schemaVersion: 1, runId, result }, true);
}
async readTerminalResult(runId: RegistryRunId32, expectedDigest: Sha256Hex): Promise<RegistryAddressedResultV1> {
let parsed: unknown;
try { parsed = JSON.parse((await strictRead(this.resultPath(runId))).text); } catch { throw CONFLICT(); }
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw CONFLICT();
const value = parsed as { schemaVersion?: unknown; runId?: unknown; result?: unknown };
if (value.schemaVersion !== 1 || value.runId !== runId || !value.result || registryDigest(value.result) !== expectedDigest) throw CONFLICT();
return value.result as RegistryAddressedResultV1;
}
async read(runId: RegistryRunId32): Promise<RegistryAddressedPublicationStateV1> {
const path = this.path(runId); let parsed: unknown;
try { parsed = JSON.parse((await strictRead(path)).text); } catch { throw CONFLICT(); }
if (!parsed || typeof parsed !== "object" || Array.isArray(parsed)) throw CONFLICT();
const s = parsed as StateFields & { operation?: unknown };
const keys = Object.keys(parsed).sort();
const expected = ["schemaVersion","runId","requestSha256","jobArtifactPath","phase","operation","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256","advertisedTargetCommit","immutableTargetRef","fetchedTargetCommit","targetCommit","targetManifestSha256","targetWorkspaces","changedWorkspaceIds","planSha256","changedSetSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256","baseCommit","baseManifestSha256","baseWorkspaces","changedSetRule"].sort();
if (keys.length !== expected.length || keys.some((v, i) => v !== expected[i])) throw CONFLICT();
failIfBadIdentity(s, runId);
if (s.operation !== "registry_bootstrap" && s.operation !== "registry_pull") throw CONFLICT();
for (const key of ["advertisedTargetCommit","fetchedTargetCommit","targetCommit"] as const) if (s[key] !== null && !REV.test(s[key])) throw CONFLICT();
for (const key of ["requestSha256","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256"] as const) if (!SHA.test(s[key])) throw CONFLICT();
for (const key of ["targetManifestSha256","changedSetSha256","planSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256"] as const) if (s[key] !== null && !SHA.test(s[key])) throw CONFLICT();
if (s.immutableTargetRef !== null && s.immutableTargetRef !== `refs/thoth/addressed-runs/${runId}/target`) throw CONFLICT();
if (s.operation === "registry_bootstrap" && (s.baseCommit !== null || s.baseManifestSha256 !== null || s.baseWorkspaces.length !== 0 || (s.changedSetRule !== null && s.changedSetRule !== "all_target_workspace_ids"))) throw CONFLICT();
if (s.operation === "registry_pull" && (s.baseCommit === null || !REV.test(s.baseCommit) || (s.baseManifestSha256 !== null && !SHA.test(s.baseManifestSha256)) || s.changedSetRule !== null && s.changedSetRule !== "symmetric_base_target_workspace_difference")) throw CONFLICT();
if (s.targetWorkspaces !== null && !Array.isArray(s.targetWorkspaces) || s.changedWorkspaceIds !== null && !Array.isArray(s.changedWorkspaceIds)) throw CONFLICT();
const required: Record<RegistryAddressedPublicationPhaseV1, readonly (keyof StateFields)[]> = {
request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], target_fetched: ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit"],
planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"],
};
for (const key of required[s.phase]) if (s[key] === null || s[key] === undefined) throw CONFLICT();
return s as RegistryAddressedPublicationStateV1;
}
async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial<RegistryAddressedPublicationStateV1> = {}): Promise<RegistryAddressedPublicationStateV1> {
const old = await this.read(runId);
const from = PHASES.indexOf(old.phase);
const to = PHASES.indexOf(phase);
if (to !== from + 1) throw CONFLICT();
const allowed = new Set(["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "changedSetRule", "baseManifestSha256", "baseWorkspaces"]);
for (const key of Object.keys(patch)) if (!allowed.has(key)) throw CONFLICT();
for (const key of ["schemaVersion", "runId", "requestSha256", "jobArtifactPath", "installationIdentitySha256", "repositoryIdentitySha256", "remoteRefIdentitySha256", "operation", "baseCommit"] as const) {
if (key in patch && patch[key] !== old[key]) throw CONFLICT();
}
const next = { ...old, ...patch, phase, priorStateSha256: registryDigest(old) as Sha256Hex } as RegistryAddressedPublicationStateV1;
const required: Record<RegistryAddressedPublicationPhaseV1, readonly (keyof StateFields)[]> = {
request_claimed: [],
target_advertised: ["advertisedTargetCommit", "immutableTargetRef"],
target_fetched: ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit"],
planned: ["fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "changedSetSha256", "planSha256", "changedSetRule"],
participants_prepared: ["participantsSha256"],
publication_intent_durable: ["publicationIntentSha256"],
target_published: ["publishedActiveStateSha256"],
terminal_durable: ["terminalResultSha256"],
};
for (const key of required[phase]) if (next[key] === null || next[key] === undefined) throw CONFLICT();
await this.durable(this.path(runId), next);
return next;
}
snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryActiveSnapshotV1 { if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces) throw CONFLICT(); return { schemaVersion: 1, commit: state.targetCommit, manifestSha256: state.targetManifestSha256, workspaces: state.targetWorkspaces }; }
async scan(): Promise<readonly RegistryAddressedPublicationStateV1[]> { await this.dirs(); let entries: string[]; try { entries = (await readdir(this.jobsDirectory)).sort(); } catch { throw CONFLICT(); } if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT();
for (const name of entries.filter(x => /^\.[0-9a-f]{32}\.json\.tmp$/.test(x))) { const path = join(this.jobsDirectory, name); const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes) throw CONFLICT(); await unlink(path); await fsyncParent(path); }
entries = (await readdir(this.jobsDirectory)).sort();
if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT();
const names = entries.filter(x => /^[0-9a-f]{32}\.json$/.test(x));
const resultNames = entries.filter(x => /^[0-9a-f]{32}\.result\.json$/.test(x));
if (names.length + resultNames.length !== entries.length) throw CONFLICT();
let total = 0;
const identities = new Map<string, string>();
const out: RegistryAddressedPublicationStateV1[] = [];
for (const name of entries) {
const st = await lstat(join(this.jobsDirectory, name));
if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes || (total += st.size) > REGISTRY_SCAN_LIMITS_V1.maximumTotalArtifactBytes) throw CONFLICT();
identities.set(name, `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`);
}
for (const name of names) {
const state = await this.read(name.slice(0, -5) as RegistryRunId32);
if (state.phase === "terminal_durable") {
await this.readTerminalResult(state.runId, state.terminalResultSha256 as Sha256Hex);
}
out.push(state);
}
const verify = (await readdir(this.jobsDirectory)).sort();
if (verify.length !== entries.length || verify.some((name, i) => name !== entries[i])) throw CONFLICT();
for (const name of entries) {
const st = await lstat(join(this.jobsDirectory, name));
const key = `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`;
if (identities.get(name) !== key) throw CONFLICT();
}
return out;
}
async removeSibling(runId: RegistryRunId32): Promise<void> { const path = join(this.jobsDirectory, `.${runId}.json.tmp`); try { const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600)) throw CONFLICT(); await unlink(path); await fsyncParent(path); } catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw CONFLICT(); } }
}