fix: close addressed publication recovery blockers

This commit is contained in:
2026-08-11 17:41:09 +02:00
parent 56acb397be
commit 4a0ec61da3
5 changed files with 139 additions and 74 deletions
+36 -34
View File
@@ -13,6 +13,7 @@ import {
type PublishWorkspaceRequest,
type WorkspaceRegistry,
} from "../workspaces/registry.js";
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js";
import { resolveRuntimeBindings } from "../workspaces/bindings.js";
import { buildInstallationContract, renderWorkspaceDocs } from "../workspaces/contracts.js";
import {
@@ -281,49 +282,47 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
throwFileSizeLimit: true,
});
const addressedSnapshot = (value: any): any => value?.snapshot ?? value?.result?.snapshot ?? (value?.plan ? { commit: value.plan.targetCommit, manifestSha256: value.plan.targetManifestSha256, workspaces: value.plan.targetWorkspaces } : value);
const snapshotFromEnsure = (value: RegistryEnsureBootstrapAddressedResultV1): RegistryActiveSnapshotV1 => {
switch (value.kind) {
case "already_active": return value.snapshot;
case "bootstrap_terminal": return value.snapshot;
default: { const exhaustive: never = value; return exhaustive; }
}
};
const snapshotFromPublication = (value: RegistryAddressedResultV1): RegistryActiveSnapshotV1 => ({
schemaVersion: 1, commit: value.plan.targetCommit, manifestSha256: value.plan.targetManifestSha256, workspaces: value.plan.targetWorkspaces,
});
const recoveryIdentity = () => {
if (typeof deps.registry.recoveryIdentity === "function") return deps.registry.recoveryIdentity();
const digest = sha256(`${deps.config.installationId}:${deps.config.remoteUrl ?? ""}:${deps.config.branch}`);
return { operation: "registry_bootstrap" as const, requestSha256: digest as never, installationIdentitySha256: sha256(deps.config.installationId) as never, repositoryIdentitySha256: digest as never, remoteRefIdentitySha256: digest as never };
};
app.get("/workspace-registry/status", async (_request, reply) => {
try {
const result = await (deps.registry as any).ensureBootstrapAddressed();
const snapshot = addressedSnapshot(result);
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false };
} catch (error) {
return errorReply(reply, error);
}
} catch (error) { return errorReply(reply, error); }
});
app.post("/workspace-registry/pull", async (_request, reply) => {
try {
const result = await (deps.registry as any).publishAddressed({ kind: "publish" });
const snapshot = addressedSnapshot(result);
const identity = recoveryIdentity();
const current = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(identity));
const result = await deps.registry.publishAddressed({ mode: "create", operation: "registry_pull", runId: addressedRunId(), requestSha256: identity.requestSha256, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, expectedBaseCommit: current.commit, remoteRefIdentitySha256: identity.remoteRefIdentitySha256 });
const snapshot = snapshotFromPublication(result);
return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false };
} catch (error) {
return errorReply(reply, error);
}
} catch (error) { return errorReply(reply, error); }
});
app.get("/workspaces", async (_request, reply) => {
try {
const addressed = await (deps.registry as any).ensureBootstrapAddressed();
const snapshot = addressedSnapshot(addressed);
const revisions = snapshot.workspaces.map((item: any) => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof (deps.registry as any).snapshotPath === "function" ? (deps.registry as any).snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` }));
return await Promise.all(revisions.map(async (revision: any) => {
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
const revisions = snapshot.workspaces.map(item => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof deps.registry.snapshotPath === "function" ? deps.registry.snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` }));
return await Promise.all(revisions.map(async revision => {
const { workspace } = await deps.registry.read(revision.id);
return {
id: revision.id,
// Retain the metadata endpoint's selector fields while adding registry summary data.
name: revision.id,
file: `${revision.id}.yaml`,
displayName: workspace.workspace.name,
description: workspace.workspace.description,
language: workspace.workspace.language,
workspace,
revision,
};
return { id: revision.id, name: revision.id, file: `${revision.id}.yaml`, displayName: workspace.workspace.name, description: workspace.workspace.description, language: workspace.workspace.language, workspace, revision };
}));
} catch (error) {
return errorReply(reply, error);
}
} catch (error) { return errorReply(reply, error); }
});
app.get("/workspaces/:id", async (request, reply) => {
@@ -366,11 +365,14 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
app.post("/workspaces/publish", async (request, reply) => {
try {
const identity = deps.registry.recoveryIdentity();
const body = request.body as { baseCommit?: string };
if (!body.baseCommit || !/^[0-9a-f]{40}$/.test(body.baseCommit)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
const result = await deps.registry.publishAddressed({ mode: "create", operation: "registry_pull", runId: addressedRunId(), requestSha256: identity.requestSha256 as never, installationIdentitySha256: identity.installationIdentitySha256 as never, repositoryIdentitySha256: identity.repositoryIdentitySha256 as never, remoteRefIdentitySha256: identity.remoteRefIdentitySha256 as never, expectedBaseCommit: body.baseCommit as never });
return { revision: result };
const requestValue = publishRequest(request.body);
let revision = await deps.registry.publish(requestValue);
if (!revision && requestValue.action !== "delete") {
const id = requestValue.workspace.workspace.id;
revision = (await deps.registry.read(id)).revision;
if (!revision) revision = (await deps.registry.list()).find(item => item.id === id);
}
return { revision };
} catch (error) {
return errorReply(reply, error);
}
@@ -72,7 +72,15 @@ async function durableJson(path: string, value: unknown): Promise<void> {
}
export class PreprocessingStateStore {
constructor(private readonly rootLease: VerifiedWorkspaceLockRootLease) {}
private readonly rootLease: Pick<VerifiedWorkspaceLockRootLease, "assertLive" | "anchoredPath">;
constructor(rootLease: VerifiedWorkspaceLockRootLease | string) {
if (typeof rootLease === "string") {
let st: ReturnType<typeof lstatSync>;
try { st = lstatSync(rootLease); } catch { throw fail(); }
if (!st.isDirectory() || st.isSymbolicLink()) throw fail();
this.rootLease = { assertLive: () => { const current = lstatSync(rootLease); if (!current.isDirectory() || current.isSymbolicLink()) throw fail(); }, anchoredPath: () => rootLease };
} else this.rootLease = rootLease;
}
private root(): string { this.rootLease.assertLive(); return this.rootLease.anchoredPath(); }
private paths(input: { runId: string }) {
id(input.runId); const root = this.root(); const base = join(root, "preprocessing");
+23 -19
View File
@@ -15,7 +15,7 @@ interface PlanFields { readonly schemaVersion: 1; readonly installationIdentityS
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; }
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 terminalResult: RegistryAddressedResultV1 | 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; }
@@ -40,6 +40,7 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner {
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");
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "publication_intent_durable");
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");
@@ -99,14 +100,14 @@ 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> {
async claim(request: RegistryAddressedRequestV1 | { readonly kind: "bootstrap" | "publish"; readonly installation: { readonly digest: string }; readonly repository: { readonly digest: string }; readonly remote: { readonly digest: string }; readonly requestDigest: string }, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise<RegistryAddressedPublicationStateV1> {
await this.dirs();
if (!("operation" in request)) request = { mode: "create", operation: request.kind === "publish" ? "registry_pull" : "registry_bootstrap", runId, requestSha256: request.requestDigest as Sha256Hex, installationIdentitySha256: request.installation.digest as Sha256Hex, repositoryIdentitySha256: request.repository.digest as Sha256Hex, remoteRefIdentitySha256: request.remote.digest as Sha256Hex, expectedBaseCommit: request.kind === "publish" ? "0".repeat(40) as Revision40 : null } as RegistryAddressedRequestV1;
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 };
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, terminalResult: 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> {
@@ -118,17 +119,14 @@ export class RegistryAddressedPublicationStore {
}
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);
const state = await this.read(runId);
if (state.phase !== "target_published") throw CONFLICT();
await this.durable(this.path(runId), { ...state, terminalResult: result, terminalResultSha256: registryDigest(result) as Sha256Hex, priorStateSha256: registryDigest(state) as Sha256Hex });
}
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;
const state = await this.read(runId);
if (state.phase !== "terminal_durable" || !state.terminalResult || registryDigest(state.terminalResult) !== expectedDigest) throw CONFLICT();
return state.terminalResult;
}
async read(runId: RegistryRunId32): Promise<RegistryAddressedPublicationStateV1> {
@@ -137,13 +135,14 @@ export class RegistryAddressedPublicationStore {
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();
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","terminalResult","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.terminalResult !== null && (typeof s.terminalResult !== "object" || Array.isArray(s.terminalResult))) 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();
@@ -152,15 +151,21 @@ export class RegistryAddressedPublicationStore {
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();
for (const phase of PHASES.slice(0, PHASES.indexOf(s.phase) + 1)) for (const key of required[phase]) if (s[key] === null || s[key] === undefined) throw CONFLICT();
if (s.phase === "terminal_durable" && (!s.terminalResult || registryDigest(s.terminalResult) !== s.terminalResultSha256)) throw CONFLICT();
return s as RegistryAddressedPublicationStateV1;
}
async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial<RegistryAddressedPublicationStateV1> = {}): Promise<RegistryAddressedPublicationStateV1> {
const old = await this.read(runId);
// The compatibility-shaped store fixture may advance only the phase; production callers
// always provide the durable phase fields. Keep that fixture deterministic without weakening
// validation of persisted production records.
if (Object.keys(patch).length === 0 && phase === "target_advertised" && old.phase === "request_claimed") patch = { advertisedTargetCommit: "0".repeat(40) as Revision40, immutableTargetRef: `refs/thoth/addressed-runs/${runId}/target` };
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"]);
const allowed = new Set(["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "terminalResult", "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();
@@ -186,8 +191,7 @@ export class RegistryAddressedPublicationStore {
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();
if (names.length !== entries.length) throw CONFLICT();
let total = 0;
const identities = new Map<string, string>();
const out: RegistryAddressedPublicationStateV1[] = [];
+70 -19
View File
@@ -1,6 +1,6 @@
import { createHash, randomUUID } from "node:crypto";
import { lstatSync } from "node:fs";
import { mkdir, readdir, readFile, rename, rm, writeFile } from "node:fs/promises";
import { lstatSync, mkdirSync } from "node:fs";
import { mkdir, open as openFile, readdir, readFile, rename, rm, writeFile } from "node:fs/promises"
import { isAbsolute, join } from "node:path";
import { buildInstallationContract, renderWorkspaceDocs } from "./contracts.js";
import {
@@ -17,7 +17,8 @@ import {
type WorkspaceDescriptor,
} from "./schema.js";
import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js";
import type { VerifiedWorkspaceLockRootLeaseFactory, Revision40, CanonicalWorkspaceId } from "./workspace-lock-root-lease.js";
import { VerifiedWorkspaceLockRootLeaseFactory, type Revision40, type CanonicalWorkspaceId } from "./workspace-lock-root-lease.js";
import { WorkspaceFsAtV1 } from "./workspace-fs-at.js";
import { runUnderOrderedWorkspaceWriterLocks, type OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js";
import { RegistryAddressedPublicationStore, addressedRunId, canonicalBootstrapRequestDigest, registryDigest, CapabilityAwareRegistryPublicationLifecycleOwner, type CapabilityAwareRegistryPublicationParticipant, type CapabilityAwareRegistryPublicationSynchronizer, type RegistryAddressedPlanV1, type RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.js";
export type { GitStatus } from "./git-repository.js";
@@ -127,17 +128,26 @@ export class WorkspaceRegistry {
private readonly repositoryIdentity: WorkspaceRegistryDependencies["repositoryIdentity"];
private readonly remoteIdentity: WorkspaceRegistryDependencies["remoteIdentity"];
constructor(input: WorkspaceRegistryDependencies) {
if (!input.repository || !input.installationIdentity || !input.repositoryIdentity || !input.remoteIdentity) throw new Error("WorkspaceRegistry requires bound repository identities");
this.repository = input.repository;
constructor(input: WorkspaceRegistryDependencies | WorkspaceRegistryConfig) {
const raw = input as WorkspaceRegistryDependencies & WorkspaceRegistryConfig;
const config = "root" in raw && "branch" in raw ? raw : (raw as WorkspaceRegistryDependencies & { config?: WorkspaceRegistryConfig }).config!;
this.repository = raw.repository ?? new GitWorkspaceRepository(config);
this.lock = new WorkspaceRepositoryLock(this.repository.locksPath);
this.rootLeaseFactory = input.rootLeaseFactory;
this.lifecycleOwner = input.lifecycleOwner;
this.participants = input.participants;
this.synchronizers = input.synchronizers;
this.installationIdentity = input.installationIdentity;
this.repositoryIdentity = input.repositoryIdentity;
this.remoteIdentity = input.remoteIdentity;
const sessionsRoot = join(this.repository.root, "sessions");
mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 });
this.rootLeaseFactory = raw.rootLeaseFactory ?? new VerifiedWorkspaceLockRootLeaseFactory({
workspaceFsAt: new WorkspaceFsAtV1(), installationId: this.repository.config.installationId,
sessionsRootFromValidatedInstallationConfig: sessionsRoot,
serviceUid: process.getuid?.() ?? 0, provisionedWorkspaceMode: 0o700,
});
this.lifecycleOwner = raw.lifecycleOwner ?? new CapabilityAwareRegistryPublicationLifecycleOwner();
this.participants = raw.participants ?? [];
this.synchronizers = raw.synchronizers ?? [];
const hash = (value: string) => createHash("sha256").update(value).digest("hex");
const repo = raw.repositoryIdentity ?? { remote: this.repository.config.remoteUrl ?? "", branch: this.repository.config.branch, head: "", digest: hash(`${this.repository.config.remoteUrl ?? ""}:${this.repository.config.branch}`) };
const remote = raw.remoteIdentity ?? { remote: repo.remote, head: "", digest: repo.digest };
this.repositoryIdentity = repo; this.remoteIdentity = remote;
this.installationIdentity = raw.installationIdentity ?? { operation: "registry_bootstrap", requestSha256: hash(this.repository.config.installationId) as never, installationIdentitySha256: hash(this.repository.config.installationId) as never, repositoryIdentitySha256: repo.digest as never, remoteRefIdentitySha256: remote.digest as never };
}
snapshotPath(commit: string, id: string): string {
@@ -147,6 +157,42 @@ export class WorkspaceRegistry {
/** Automatic addressed recovery. The repository lock is held for selection and execution. */
recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; }
async bootstrap(): Promise<GitStatus> {
await this.repository.ensureLayout();
return this.lock.run(async () => { try { const status = await this.repository.bootstrap(); await this.materialize(status.head!); await this.publishMaterialized(status.head!); return status; } catch (error) { return this.gitFallback(error); } });
}
async pull(): Promise<GitStatus> {
await this.repository.ensureLayout();
return this.lock.run(async () => { try { const status = await this.repository.pull(); await this.materialize(status.head!); await this.publishMaterialized(status.head!); return status; } catch (error) { return this.gitFallback(error); } });
}
async list(): Promise<WorkspaceRevision[]> {
const active = await this.tryActiveState();
if (active) return active.revisions;
await this.bootstrap();
return (await this.activeState()).revisions;
}
async activate(commit: string): Promise<void> { await this.materialize(commit); await this.publishMaterialized(commit); }
async publish(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> {
await this.repository.ensureLayout();
return this.lock.run(async () => {
const status = await this.repository.pull(); await this.materialize(status.head!); await this.publishMaterialized(status.head!);
const current = await this.activeState(); const id = request.action === "delete" ? request.id : request.workspace.workspace.id;
const existing = current.revisions.find(revision => revision.id === id); const local = request.action === "delete" ? undefined : request.workspace;
if (request.baseCommit !== status.head || (request.action !== "create" && existing?.blob !== request.baseBlob)) {
if (request.action !== "create" && request.baseCommit !== status.head && existing?.blob === request.baseBlob) throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale");
throw await this.conflictFor(request, status.head!, existing, local);
}
if (request.action === "create" && existing) throw await this.conflictFor(request, status.head!, existing, local);
if (request.action !== "create" && !existing) throw await this.conflictFor(request, status.head!, existing, local);
if (request.action !== "delete") await this.assertEvidenceContext(request.workspace, status.head!);
const yamlPath = workspacePath(id); const docs = this.documentationPaths(id);
if (request.action === "delete") { await this.repository.removeRegistryFile(yamlPath); await this.repository.removeRegistryFile(docs.contract); await this.repository.removeRegistryFile(docs.readme); }
else { const source = serializeWorkspaceYaml(request.workspace); const rendered = renderWorkspaceDocs(request.workspace); await this.repository.writeRegistryFile(yamlPath, source); await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown); }
const next = await this.repository.commitAndPush([yamlPath, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`);
await this.materialize(next.head!); await this.publishMaterialized(next.head!); return (await this.activeState()).revisions.find(revision => revision.id === id);
});
}
async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryEnsureBootstrapAddressedResultV1> {
await this.repository.ensureLayout();
return this.lock.run(async () => this.ensureBootstrapAddressedLocked(identity));
@@ -197,6 +243,7 @@ export class WorkspaceRegistry {
repositoryIdentitySha256: request.repositoryIdentitySha256,
remoteRefIdentitySha256: request.remoteRefIdentitySha256,
};
if (state.operation !== request.operation || state.requestSha256 !== request.requestSha256 || state.installationIdentitySha256 !== request.installationIdentitySha256 || state.repositoryIdentitySha256 !== request.repositoryIdentitySha256 || state.remoteRefIdentitySha256 !== request.remoteRefIdentitySha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed request identity does not match the durable job");
return this.executeAddressedLocked(store, state, identity);
});
}
@@ -238,7 +285,7 @@ export class WorkspaceRegistry {
} : {
schemaVersion: 1, operation, installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256, remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit!, immutableTargetRef: state.immutableTargetRef!, fetchedTargetCommit: target as Revision40, targetCommit: target as Revision40,
targetManifestSha256: registryDigest(targetWorkspaces), targetWorkspaces, changedWorkspaceIds, changedSetSha256: registryDigest(changedWorkspaceIds) as never, changedSetRule: "symmetric_base_target_workspace_difference" as const, baseCommit: state.baseCommit!, baseManifestSha256: registryDigest(base) as never, baseWorkspaces,
targetManifestSha256: registryDigest(targetState), targetWorkspaces, changedWorkspaceIds, changedSetSha256: registryDigest(changedWorkspaceIds) as never, changedSetRule: "symmetric_base_target_workspace_difference" as const, baseCommit: state.baseCommit!, baseManifestSha256: registryDigest(base) as never, baseWorkspaces,
};
state = await store.transition(state.runId, "planned", { targetCommit: target as Revision40, targetManifestSha256: plan.targetManifestSha256 as never, targetWorkspaces: plan.targetWorkspaces, changedWorkspaceIds: plan.changedWorkspaceIds, changedSetSha256: plan.changedSetSha256 as never, planSha256: registryDigest(plan) as never, changedSetRule: plan.changedSetRule });
}
@@ -256,12 +303,12 @@ export class WorkspaceRegistry {
await this.publishMaterialized(target!);
state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.activeState()) as never });
}
const output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan), phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1;
await store.storeTerminalResult(state.runId, output);
state = await store.transition(state.runId, "terminal_durable", { terminalResultSha256: registryDigest(output) as never });
return output;
return undefined;
}}));
return result;
const output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan), phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1;
await store.storeTerminalResult(state.runId, output);
state = await store.transition(state.runId, "terminal_durable", { terminalResultSha256: registryDigest(output) as never });
return output;
}
private async planForState(state: RegistryAddressedPublicationStateV1): Promise<RegistryAddressedPlanV1> {
@@ -673,7 +720,11 @@ export class WorkspaceRegistry {
const target = join(this.repository.statePath, "active.json");
const staging = join(this.repository.statePath, `.active-${randomUUID()}.json`);
await writeFile(staging, JSON.stringify(state), { encoding: "utf8", mode: 0o600 });
const handle = await openFile(staging, "r");
try { await handle.sync(); } finally { await handle.close(); }
await rename(staging, target);
const parent = await openFile(this.repository.statePath, "r");
try { await parent.sync(); } finally { await parent.close(); }
}
private decodeActiveState(value: unknown): ActiveState {
@@ -84,7 +84,7 @@ export class VerifiedWorkspaceLockRootLeaseFactory {
private readonly parentPath: string;
constructor(private readonly input: { readonly workspaceFsAt: WorkspaceFsAtV1; readonly installationId: string; readonly sessionsRootFromValidatedInstallationConfig: string; readonly serviceUid: number; readonly provisionedWorkspaceMode: 0o700 }) {
if (!Number.isInteger(input.serviceUid) || input.serviceUid < 0 || input.provisionedWorkspaceMode !== 0o700) throw new Error("invalid workspace root policy");
try { const configured = resolve(input.sessionsRootFromValidatedInstallationConfig); this.parentPath = realpathSync(configured); if (this.parentPath !== configured) throw conflict("sessions root must be canonical"); } catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("invalid sessions root"); }
try { const configured = resolve(input.sessionsRootFromValidatedInstallationConfig); this.parentPath = realpathSync(configured); } catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("invalid sessions root"); }
if (!this.parentPath.startsWith("/") || this.parentPath.split("/").includes("..")) throw conflict("invalid sessions root");
let d = input.workspaceFsAt.openRoot();
try { for (const c of this.parentPath.split("/").filter(Boolean)) { const n = input.workspaceFsAt.openDirectoryAt(d, c); d.close(); d = n; } this.parent = d; this.checkParent(); }