fix: enforce addressed registry lifecycle invariants
This commit is contained in:
+5
-3
@@ -47,6 +47,7 @@ export interface BuildAppDeps {
|
||||
|
||||
function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry {
|
||||
const registryConfig = config.workspaceRegistry;
|
||||
mkdirSync(registryConfig.root, { recursive: true, mode: 0o700 });
|
||||
const repository = new GitWorkspaceRepository(registryConfig);
|
||||
const sessionsRoot = join(registryConfig.root, "sessions");
|
||||
mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 });
|
||||
@@ -56,9 +57,10 @@ function createWorkspaceRegistry(config: AppConfig): WorkspaceRegistry {
|
||||
provisionedWorkspaceMode: 0o700,
|
||||
});
|
||||
const hash = (value: string) => createHash("sha256").update(value).digest("hex");
|
||||
const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(`${registryConfig.remoteUrl ?? ""}:${registryConfig.branch}`) };
|
||||
const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: repositoryIdentity.digest };
|
||||
const request = { kind: "bootstrap" as const, installation: { installationId: registryConfig.installationId, digest: hash(registryConfig.installationId) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] };
|
||||
const volumeIdentity = realpathSync(registryConfig.root);
|
||||
const repositoryIdentity = { remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch, head: "", digest: hash(JSON.stringify({ schemaVersion: 1, volume: volumeIdentity, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) };
|
||||
const remoteIdentity = { remote: registryConfig.remoteUrl ?? "", head: "", digest: hash(JSON.stringify({ schemaVersion: 1, remote: registryConfig.remoteUrl ?? "", branch: registryConfig.branch })) };
|
||||
const request = { kind: "bootstrap" as const, operation: "registry_bootstrap" as const, installation: { installationId: registryConfig.installationId, volume: volumeIdentity, digest: hash(JSON.stringify({ schemaVersion: 1, installationId: registryConfig.installationId, volume: volumeIdentity })) }, repository: repositoryIdentity, remote: remoteIdentity, workspaceIds: [] as string[] };
|
||||
return new WorkspaceRegistry({ rootLeaseFactory, lifecycleOwner: new CapabilityAwareRegistryPublicationLifecycleOwner(), participants: [], synchronizers: [], repository,
|
||||
installationIdentity: { operation: "registry_bootstrap", requestSha256: canonicalBootstrapRequestDigest(request) as never, installationIdentitySha256: request.installation.digest as never, repositoryIdentitySha256: repositoryIdentity.digest as never, remoteRefIdentitySha256: remoteIdentity.digest as never }, repositoryIdentity, remoteIdentity });
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ import {
|
||||
type PublishWorkspaceRequest,
|
||||
type WorkspaceRegistry,
|
||||
} from "../workspaces/registry.js";
|
||||
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1, RegistryWorkspaceMutationV1 } from "../workspaces/registry-publication.js";
|
||||
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js";
|
||||
import type { Revision40 } from "../workspaces/workspace-lock-root-lease.js";
|
||||
import { resolveRuntimeBindings } from "../workspaces/bindings.js";
|
||||
import { buildInstallationContract, renderWorkspaceDocs } from "../workspaces/contracts.js";
|
||||
@@ -320,8 +320,12 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
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, name: revision.id, file: `${revision.id}.yaml`, displayName: workspace.workspace.name, description: workspace.workspace.description, language: workspace.workspace.language, workspace, revision };
|
||||
const pinned = typeof deps.registry.readPinned === "function"
|
||||
? await deps.registry.readPinned(revision.id, revision.commit)
|
||||
: await deps.registry.read(revision.id);
|
||||
const workspace = pinned.workspace;
|
||||
const exactRevision = { ...revision, snapshotPath: "workspaceConfigPath" in pinned ? pinned.workspaceConfigPath : revision.snapshotPath };
|
||||
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: exactRevision };
|
||||
}));
|
||||
} catch (error) { return errorReply(reply, error); }
|
||||
});
|
||||
@@ -369,7 +373,6 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
const requestValue = publishRequest(request.body);
|
||||
const identity = recoveryIdentity();
|
||||
const id = requestValue.action === "delete" ? requestValue.id : requestValue.workspace.workspace.id;
|
||||
const mutation = requestValue as unknown as RegistryWorkspaceMutationV1;
|
||||
const addressed = {
|
||||
mode: "create" as const, operation: "registry_pull" as const, runId: addressedRunId(),
|
||||
requestSha256: sha256(JSON.stringify(requestValue)) as never,
|
||||
@@ -377,9 +380,8 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
repositoryIdentitySha256: identity.repositoryIdentitySha256,
|
||||
remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
|
||||
expectedBaseCommit: requestValue.baseCommit as Revision40,
|
||||
mutation,
|
||||
};
|
||||
const result = await deps.registry.publishAddressed(addressed);
|
||||
const result = await deps.registry.publishAddressed(addressed, { authoring: requestValue });
|
||||
return { revision: revisionFromPublication(result, id) };
|
||||
} catch (error) {
|
||||
return errorReply(reply, error);
|
||||
|
||||
@@ -16,7 +16,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 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; }
|
||||
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; }
|
||||
@@ -32,7 +32,7 @@ export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: strin
|
||||
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> {
|
||||
async run<T>(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant<unknown>[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise<T>; readonly afterPublication?: (result: T) => Promise<void>; }): 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.
|
||||
@@ -45,6 +45,7 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner {
|
||||
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");
|
||||
if (input.afterPublication) await input.afterPublication(result);
|
||||
return result;
|
||||
}
|
||||
const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId;
|
||||
@@ -69,15 +70,10 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner {
|
||||
}
|
||||
}
|
||||
|
||||
export type RegistryWorkspaceMutationV1 =
|
||||
| { readonly action: "create"; readonly workspace: CanonicalWorkspace; readonly baseCommit: Revision40 }
|
||||
| { readonly action: "update"; readonly workspace: CanonicalWorkspace; readonly baseCommit: Revision40; readonly baseBlob: Revision40 }
|
||||
| { readonly action: "delete"; readonly id: CanonicalWorkspaceId; readonly baseCommit: Revision40; readonly baseBlob: Revision40 };
|
||||
|
||||
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 mutation?: RegistryWorkspaceMutationV1 }
|
||||
| { 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"; }
|
||||
@@ -91,7 +87,10 @@ const PHASES: readonly RegistryAddressedPublicationPhaseV1[] = ["request_claimed
|
||||
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 canonicalBootstrapRequestDigest(request: { readonly kind?: "bootstrap" | "publish"; readonly operation?: RegistryAddressedOperationV1; readonly installation?: unknown; readonly repository?: unknown; readonly remote?: unknown; readonly workspaceIds?: readonly string[] }): string {
|
||||
const operation = request.operation ?? (request.kind === "publish" ? "registry_pull" : "registry_bootstrap");
|
||||
return registryDigest({ schemaVersion: 1, 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; }
|
||||
@@ -108,31 +107,52 @@ export class RegistryAddressedPublicationStore {
|
||||
private path(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.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 | { 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> {
|
||||
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()), base?: { readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] }): 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 (!("operation" in request)) {
|
||||
const legacyPublish = request.kind === "publish";
|
||||
request = { mode: "create", operation: legacyPublish ? "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: legacyPublish ? "0".repeat(40) as Revision40 : null } as RegistryAddressedRequestV1;
|
||||
if (legacyPublish && !base) base = { manifestSha256: "0".repeat(64) as Sha256Hex, workspaces: [] };
|
||||
}
|
||||
if (!RUN.test(runId)) throw CONFLICT();
|
||||
if (request.operation === "registry_pull" && !base) 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, terminalResult: 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, priorStateSha256: null, baseCommit: request.operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null, baseManifestSha256: base?.manifestSha256 ?? null, baseWorkspaces: base?.workspaces ?? [], 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> {
|
||||
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 });
|
||||
// The terminal result is a deterministic projection of the immutable plan. Persist only
|
||||
// its digest in the fixed state shape; recovery reconstructs and verifies the projection.
|
||||
const expected = registryDigest(result) as Sha256Hex;
|
||||
await this.transition(runId, "terminal_durable", { terminalResultSha256: expected });
|
||||
}
|
||||
private planFromState(state: RegistryAddressedPublicationStateV1): RegistryAddressedPlanV1 {
|
||||
if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces || !state.changedWorkspaceIds || !state.changedSetSha256 || !state.advertisedTargetCommit || !state.immutableTargetRef || !state.fetchedTargetCommit || !state.planSha256) throw CONFLICT();
|
||||
const common = { schemaVersion: 1 as const, installationIdentitySha256: state.installationIdentitySha256,
|
||||
repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256,
|
||||
jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: state.advertisedTargetCommit,
|
||||
immutableTargetRef: state.immutableTargetRef, fetchedTargetCommit: state.fetchedTargetCommit,
|
||||
targetCommit: state.targetCommit, targetManifestSha256: state.targetManifestSha256,
|
||||
targetWorkspaces: state.targetWorkspaces, changedWorkspaceIds: state.changedWorkspaceIds,
|
||||
changedSetSha256: state.changedSetSha256 };
|
||||
const plan = state.operation === "registry_bootstrap"
|
||||
? { ...common, operation: "registry_bootstrap" as const, changedSetRule: "all_target_workspace_ids" as const, baseCommit: null, baseManifestSha256: null, baseWorkspaces: [] as const }
|
||||
: { ...common, operation: "registry_pull" as const, changedSetRule: "symmetric_base_target_workspace_difference" as const,
|
||||
baseCommit: state.baseCommit!, baseManifestSha256: state.baseManifestSha256!, baseWorkspaces: state.baseWorkspaces };
|
||||
if (registryDigest(plan) !== state.planSha256) throw CONFLICT();
|
||||
return plan;
|
||||
}
|
||||
async readTerminalResult(runId: RegistryRunId32, expectedDigest: Sha256Hex): Promise<RegistryAddressedResultV1> {
|
||||
const state = await this.read(runId);
|
||||
if (state.phase !== "terminal_durable" || !state.terminalResult || registryDigest(state.terminalResult) !== expectedDigest) throw CONFLICT();
|
||||
return state.terminalResult;
|
||||
if (state.phase !== "terminal_durable" || state.terminalResultSha256 !== expectedDigest) throw CONFLICT();
|
||||
const plan = await this.planFromState(state);
|
||||
const result = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath,
|
||||
plan, planSha256: state.planSha256!, phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1;
|
||||
if (registryDigest(result) !== expectedDigest) throw CONFLICT();
|
||||
return result;
|
||||
}
|
||||
|
||||
async read(runId: RegistryRunId32): Promise<RegistryAddressedPublicationStateV1> {
|
||||
@@ -141,24 +161,36 @@ 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","terminalResult","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","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.priorStateSha256 !== null && !SHA.test(s.priorStateSha256)) throw CONFLICT();
|
||||
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();
|
||||
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();
|
||||
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 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();
|
||||
const introduced: Record<string, readonly string[]> = {
|
||||
request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"], target_fetched: ["fetchedTargetCommit"],
|
||||
planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"],
|
||||
participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"],
|
||||
};
|
||||
const currentIndex = PHASES.indexOf(s.phase);
|
||||
let chainPrevious: string | null = null;
|
||||
for (let index = 0; index <= currentIndex; index++) {
|
||||
const candidate: Record<string, unknown> = { ...s, phase: PHASES[index]!, priorStateSha256: chainPrevious };
|
||||
for (const later of PHASES.slice(index + 1)) for (const key of introduced[later]!) candidate[key] = null;
|
||||
if (index === currentIndex && s.priorStateSha256 !== chainPrevious) throw CONFLICT();
|
||||
chainPrevious = registryDigest(candidate);
|
||||
}
|
||||
return s as RegistryAddressedPublicationStateV1;
|
||||
}
|
||||
async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial<RegistryAddressedPublicationStateV1> = {}): Promise<RegistryAddressedPublicationStateV1> {
|
||||
@@ -171,8 +203,20 @@ export class RegistryAddressedPublicationStore {
|
||||
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", "terminalResult", "changedSetRule", "baseManifestSha256", "baseWorkspaces"]);
|
||||
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();
|
||||
// All durable facts are append-only. In particular the advertised OID/ref and base
|
||||
// identity may never be replaced during resume or by a competing caller.
|
||||
for (const key of ["advertisedTargetCommit", "immutableTargetRef", "fetchedTargetCommit", "targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "participantsSha256", "synchronizersSha256", "publicationIntentSha256", "publishedActiveStateSha256", "terminalResultSha256", "baseManifestSha256", "baseWorkspaces", "changedSetRule"] as const) {
|
||||
if (key in patch && old[key] !== null && old[key] !== undefined && JSON.stringify(patch[key]) !== JSON.stringify(old[key])) throw CONFLICT();
|
||||
}
|
||||
const phaseFields: Record<RegistryAddressedPublicationPhaseV1, readonly string[]> = {
|
||||
request_claimed: [], target_advertised: ["advertisedTargetCommit", "immutableTargetRef"],
|
||||
target_fetched: ["fetchedTargetCommit"], planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"],
|
||||
participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"],
|
||||
target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256"],
|
||||
};
|
||||
if (Object.keys(patch).some(key => !phaseFields[phase].includes(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();
|
||||
}
|
||||
|
||||
@@ -157,45 +157,6 @@ export class WorkspaceRegistry {
|
||||
/** Automatic addressed recovery. The repository lock is held for selection and execution. */
|
||||
recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; }
|
||||
|
||||
private 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); } });
|
||||
}
|
||||
private 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); } });
|
||||
}
|
||||
private async list(): Promise<WorkspaceRevision[]> {
|
||||
const active = await this.tryActiveState();
|
||||
if (active) return active.revisions;
|
||||
await this.bootstrap();
|
||||
return (await this.activeState()).revisions;
|
||||
}
|
||||
private async activate(commit: string): Promise<void> { await this.materialize(commit); await this.publishMaterialized(commit); }
|
||||
/** Test-only migration seam; production callers use publishAddressed. */
|
||||
private async publish(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> { return this.publishWorkspace(request); }
|
||||
private async publishWorkspace(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> {
|
||||
await this.repository.ensureLayout();
|
||||
return this.lock.run(() => this.publishWorkspaceLocked(request));
|
||||
}
|
||||
private async publishWorkspaceLocked(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> {
|
||||
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));
|
||||
@@ -226,23 +187,21 @@ export class WorkspaceRegistry {
|
||||
return { kind: "bootstrap_terminal", result: result as Extract<RegistryAddressedResultV1, { operation: "registry_bootstrap" }>, snapshot: this.addressedSnapshot(published, identity.requestSha256) };
|
||||
}
|
||||
|
||||
async publishAddressed(request: RegistryAddressedRequestV1): Promise<RegistryAddressedResultV1> {
|
||||
async publishAddressed(request: RegistryAddressedRequestV1, context?: { readonly authoring?: PublishWorkspaceRequest }): Promise<RegistryAddressedResultV1> {
|
||||
await this.repository.ensureLayout();
|
||||
return this.lock.run(async () => {
|
||||
// Authoring keeps the accepted HTTP CRUD payload, but crosses the same addressed
|
||||
// boundary as every other publication. The mutation is deliberately handled while
|
||||
// repository.lock is held; callers never receive the historical publish API.
|
||||
if ("mutation" in request && request.mutation) return this.publishWorkspaceAddressed(request as Extract<RegistryAddressedRequestV1, { readonly mode: "create"; readonly operation: "registry_pull" }> & { readonly mutation: PublishWorkspaceRequest });
|
||||
const store = new RegistryAddressedPublicationStore(this.repository.root);
|
||||
let state: RegistryAddressedPublicationStateV1;
|
||||
try { state = await store.read(request.runId); }
|
||||
catch {
|
||||
state = await store.claim(request, request.runId);
|
||||
let base: ActiveState | undefined;
|
||||
if (request.operation === "registry_pull" && request.mode === "create") {
|
||||
const base = await this.snapshotState(request.expectedBaseCommit);
|
||||
const baseWorkspaces = base.revisions.map(revision => this.manifestIdentity(revision));
|
||||
state = await store.setBase(state.runId, registryDigest(base) as never, baseWorkspaces);
|
||||
base = await this.snapshotState(request.expectedBaseCommit);
|
||||
}
|
||||
state = await store.claim(request, request.runId, base ? {
|
||||
manifestSha256: registryDigest(base) as never,
|
||||
workspaces: base.revisions.map(revision => this.manifestIdentity(revision)),
|
||||
} : undefined);
|
||||
}
|
||||
const identity: RegistryBootstrapRecoveryIdentityV1 = {
|
||||
operation: "registry_bootstrap", requestSha256: request.requestSha256,
|
||||
@@ -251,43 +210,30 @@ export class WorkspaceRegistry {
|
||||
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");
|
||||
if (request.operation === "registry_pull" && request.mode === "create" && state.baseCommit !== request.expectedBaseCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed base does not match the durable job");
|
||||
if (context?.authoring && state.phase !== "request_claimed") throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Author mutation cannot be replayed after claim");
|
||||
if (context?.authoring) await this.applyAuthorMutation(context.authoring);
|
||||
return this.executeAddressedLocked(store, state, identity);
|
||||
});
|
||||
}
|
||||
|
||||
private async publishWorkspaceAddressed(request: Extract<RegistryAddressedRequestV1, { readonly mode: "create"; readonly operation: "registry_pull" }> & { readonly mutation: PublishWorkspaceRequest }): Promise<RegistryAddressedResultV1> {
|
||||
const mutation = request.mutation;
|
||||
await this.publishWorkspaceLocked(mutation);
|
||||
const targetState = await this.activeState();
|
||||
const targetWorkspaces = targetState.revisions.map(revision => this.manifestIdentity(revision));
|
||||
const base = request.expectedBaseCommit ? await this.snapshotState(request.expectedBaseCommit) : undefined;
|
||||
const baseWorkspaces = base?.revisions.map(revision => this.manifestIdentity(revision)) ?? [];
|
||||
const baseIds = new Set(baseWorkspaces.map(item => item.workspaceId));
|
||||
const targetIds = new Set(targetWorkspaces.map(item => item.workspaceId));
|
||||
const changedWorkspaceIds = [...new Set([...baseIds, ...targetIds])].filter(id =>
|
||||
!baseIds.has(id) || !targetIds.has(id)
|
||||
|| registryDigest(baseWorkspaces.find(item => item.workspaceId === id)) !== registryDigest(targetWorkspaces.find(item => item.workspaceId === id)),
|
||||
).sort() as CanonicalWorkspaceId[];
|
||||
const plan: RegistryPullAddressedPlanV1 = {
|
||||
schemaVersion: 1, operation: "registry_pull",
|
||||
installationIdentitySha256: request.installationIdentitySha256,
|
||||
repositoryIdentitySha256: request.repositoryIdentitySha256,
|
||||
remoteRefIdentitySha256: request.remoteRefIdentitySha256,
|
||||
jobArtifactPath: `addressed-publication-jobs/${request.runId}.json`,
|
||||
advertisedTargetCommit: targetState.head as Revision40,
|
||||
immutableTargetRef: `refs/thoth/addressed-runs/${request.runId}/target`,
|
||||
fetchedTargetCommit: targetState.head as Revision40,
|
||||
targetCommit: targetState.head as Revision40,
|
||||
targetManifestSha256: registryDigest(targetState) as never,
|
||||
targetWorkspaces,
|
||||
changedWorkspaceIds,
|
||||
changedSetSha256: registryDigest(changedWorkspaceIds) as never,
|
||||
changedSetRule: "symmetric_base_target_workspace_difference",
|
||||
baseCommit: request.expectedBaseCommit ?? targetState.head as Revision40,
|
||||
baseManifestSha256: registryDigest(base) as never,
|
||||
baseWorkspaces,
|
||||
};
|
||||
return { operation: "registry_pull", runId: request.runId, jobArtifactPath: plan.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable", publication: "target" };
|
||||
private async applyAuthorMutation(request: PublishWorkspaceRequest): Promise<void> {
|
||||
const status = await this.repository.pull();
|
||||
const current = await this.tryActiveState();
|
||||
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 { await this.repository.writeRegistryFile(yamlPath, serializeWorkspaceYaml(request.workspace)); const rendered = renderWorkspaceDocs(request.workspace); await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown); }
|
||||
await this.repository.commitAndPush([yamlPath, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`);
|
||||
}
|
||||
|
||||
private async executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryAddressedResultV1> {
|
||||
@@ -334,7 +280,8 @@ export class WorkspaceRegistry {
|
||||
const plan = await this.planForState(state);
|
||||
const leases = [];
|
||||
for (const id of plan.changedWorkspaceIds) leases.push(await this.rootLeaseFactory.acquireOrProvision(this.rootLeaseFactory.canonicalInput(id)));
|
||||
const result = await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => {
|
||||
let output: RegistryAddressedResultV1 | undefined;
|
||||
await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => this.lifecycleOwner.run({ plan, capabilities, participants: this.participants, synchronizers: this.synchronizers, action: async () => {
|
||||
if (state.phase === "planned") {
|
||||
state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(this.participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(this.synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never });
|
||||
}
|
||||
@@ -342,14 +289,20 @@ export class WorkspaceRegistry {
|
||||
state = await store.transition(state.runId, "publication_intent_durable", { publicationIntentSha256: registryDigest({ runId: state.runId, planSha256: state.planSha256 }) as never });
|
||||
}
|
||||
if (state.phase === "publication_intent_durable") {
|
||||
await this.publishMaterialized(target!);
|
||||
await this.publishSnapshotPointer(target!);
|
||||
state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.activeState()) as never });
|
||||
}
|
||||
return undefined;
|
||||
}, afterPublication: async () => {
|
||||
// Terminal durability is deliberately the final operation while every writer, root,
|
||||
// quiescence, and reader-exclusive gate is still owned by this callback.
|
||||
if (state.phase === "target_published") {
|
||||
output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable" as const, publication: "target" as const } as RegistryAddressedResultV1;
|
||||
await store.storeTerminalResult(state.runId, output);
|
||||
state = await store.read(state.runId);
|
||||
}
|
||||
}}));
|
||||
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 });
|
||||
if (!output) output = await store.readTerminalResult(state.runId, state.terminalResultSha256!);
|
||||
return output;
|
||||
}
|
||||
|
||||
@@ -711,7 +664,7 @@ export class WorkspaceRegistry {
|
||||
|
||||
}
|
||||
|
||||
private async publishMaterialized(commit: string): Promise<void> {
|
||||
private async publishSnapshotPointer(commit: string): Promise<void> {
|
||||
const state = await this.snapshotState(commit);
|
||||
await this.writeActiveState({ head: state.head, revisions: state.revisions });
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user