fix addressed registry recovery blockers
This commit is contained in:
@@ -8,7 +8,7 @@ import { z } from "zod";
|
||||
import type { WorkspaceRegistryConfig } from "../workspaces/types.js";
|
||||
import { publishAddressedByAuthor, type WorkspaceAuthorGitService } from "../workspaces/author-git-service.js";
|
||||
import { WorkspaceRegistryError } from "../workspaces/git-repository.js";
|
||||
import { addressedRunId } from "../workspaces/registry-publication.js";
|
||||
import { addressedRunId, registryDigest } from "../workspaces/registry-publication.js";
|
||||
import {
|
||||
WorkspaceConflictError,
|
||||
type PublishWorkspaceRequest,
|
||||
@@ -396,7 +396,7 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
const id = requestValue.action === "delete" ? requestValue.id : requestValue.workspace.workspace.id;
|
||||
const addressed = {
|
||||
mode: "create" as const, operation: "registry_pull" as const, runId: addressedRunId(),
|
||||
requestSha256: sha256(JSON.stringify(requestValue)) as never,
|
||||
requestSha256: registryDigest(requestValue) as never,
|
||||
installationIdentitySha256: identity.installationIdentitySha256,
|
||||
repositoryIdentitySha256: identity.repositoryIdentitySha256,
|
||||
remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { WorkspaceConflictError, claimAddressedPublication, type PublishWorkspaceRequest, type WorkspaceRegistry } from "./registry.js";
|
||||
import { WorkspaceConflictError, addressedRunRef, advertiseAddressedPublication, claimAddressedPublication, type PublishWorkspaceRequest, type WorkspaceRegistry } from "./registry.js";
|
||||
import type { RegistryAddressedRequestV1, RegistryAddressedResultV1 } from "./registry-publication.js";
|
||||
import { GitWorkspaceRepository, WorkspaceRegistryError, WorkspaceRepositoryLock } from "./git-repository.js";
|
||||
import { parseWorkspaceYaml, renderWorkspaceDocs, serializeWorkspaceYaml, validateOperationalWorkspace, type CanonicalWorkspace, type WorkspaceDescriptor } from "./schema.js";
|
||||
@@ -31,53 +31,76 @@ export class WorkspaceAuthorGitService {
|
||||
}
|
||||
|
||||
async publish(request: PublishWorkspaceRequest): Promise<WorkspaceAuthorGitResult> {
|
||||
const prepared = await this.prepare(request);
|
||||
await this.pushPrepared(prepared);
|
||||
return prepared;
|
||||
}
|
||||
|
||||
/** Stage and commit the author mutation locally, without remote network after pull. */
|
||||
async prepare(request: PublishWorkspaceRequest): Promise<WorkspaceAuthorGitResult> {
|
||||
await this.repository.ensureLayout();
|
||||
return this.lock.run(() => this.#prepareLocked(request));
|
||||
}
|
||||
|
||||
/** Recover a commit created just before a process kill, without pulling or rewriting it. */
|
||||
async recoverPrepared(request: PublishWorkspaceRequest): Promise<WorkspaceAuthorGitResult | undefined> {
|
||||
await this.repository.ensureLayout();
|
||||
return this.lock.run(async () => {
|
||||
const status = await this.repository.pull();
|
||||
let status;
|
||||
try { status = await this.repository.status(); } catch { return undefined; }
|
||||
const head = status.head;
|
||||
if (!head || await this.repository.parentOf(head) !== request.baseCommit) return undefined;
|
||||
const id = request.action === "delete" ? request.id : request.workspace.workspace.id;
|
||||
const path = `workspaces/${id}.yaml`;
|
||||
const current = await this.currentDescriptor(status.head!, path);
|
||||
const currentBlob = current ? await this.repository.blob(path) : undefined;
|
||||
|
||||
if (request.baseCommit !== status.head) {
|
||||
if (request.action !== "create" && currentBlob === request.baseBlob) {
|
||||
throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale");
|
||||
try {
|
||||
if (request.action === "delete") {
|
||||
const paths = await this.repository.workspacePathsAt(head);
|
||||
if (paths.includes(path)) return undefined;
|
||||
return { id, commit: head, blob: request.baseBlob };
|
||||
}
|
||||
const base = await this.currentDescriptor(request.baseCommit, path);
|
||||
throw this.conflict(request, status.head!, current, currentBlob, base);
|
||||
}
|
||||
if (request.action === "create" && current) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
if (request.action !== "create" && !current) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
if (request.action !== "create") {
|
||||
if (currentBlob !== request.baseBlob) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
}
|
||||
if (request.action !== "delete") await this.assertEvidence(request.workspace, status.head!);
|
||||
|
||||
const docs = {
|
||||
contract: `workspace-docs/${id}/contract.env.example`,
|
||||
readme: `workspace-docs/${id}/README.md`,
|
||||
};
|
||||
if (request.action === "delete") {
|
||||
await this.repository.removeRegistryFile(path);
|
||||
await this.repository.removeRegistryFile(docs.contract);
|
||||
await this.repository.removeRegistryFile(docs.readme);
|
||||
} else {
|
||||
await this.repository.writeRegistryFile(path, serializeWorkspaceYaml(request.workspace));
|
||||
const rendered = renderWorkspaceDocs(request.workspace);
|
||||
await this.repository.writeRegistryFile(docs.contract, rendered.envExample);
|
||||
await this.repository.writeRegistryFile(docs.readme, rendered.markdown);
|
||||
}
|
||||
const published = await this.repository.commitAndPush(
|
||||
[path, docs.contract, docs.readme],
|
||||
request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`,
|
||||
);
|
||||
const commit = published.head;
|
||||
if (!commit) throw new WorkspaceRegistryError("git_unavailable", "Workspace Git service is unavailable");
|
||||
if (request.action === "delete") return { id, commit, blob: request.baseBlob };
|
||||
return { id, commit, blob: await this.repository.blob(path) };
|
||||
const candidate = parseWorkspaceYaml(await this.repository.readWorkspaceAt(head, path));
|
||||
if (serializeWorkspaceYaml(candidate) !== serializeWorkspaceYaml(request.workspace)) return undefined;
|
||||
return { id, commit: head, blob: await this.repository.blobAt(head, path) };
|
||||
} catch { return undefined; }
|
||||
});
|
||||
}
|
||||
|
||||
/** Push exactly the commit previously returned by prepare (safe to repeat). */
|
||||
async pushPrepared(prepared: Pick<WorkspaceAuthorGitResult, "commit">): Promise<void> {
|
||||
await this.repository.ensureLayout();
|
||||
await this.lock.run(() => this.repository.pushExact(prepared.commit));
|
||||
}
|
||||
|
||||
async #prepareLocked(request: PublishWorkspaceRequest): Promise<WorkspaceAuthorGitResult> {
|
||||
const status = await this.repository.pull();
|
||||
const id = request.action === "delete" ? request.id : request.workspace.workspace.id;
|
||||
const path = `workspaces/${id}.yaml`;
|
||||
const current = await this.currentDescriptor(status.head!, path);
|
||||
const currentBlob = current ? await this.repository.blob(path) : undefined;
|
||||
if (request.baseCommit !== status.head) {
|
||||
if (request.action !== "create" && currentBlob === request.baseBlob) throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale");
|
||||
const base = await this.currentDescriptor(request.baseCommit, path);
|
||||
throw this.conflict(request, status.head!, current, currentBlob, base);
|
||||
}
|
||||
if (request.action === "create" && current) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
if (request.action !== "create" && !current) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
if (request.action !== "create" && currentBlob !== request.baseBlob) throw this.conflict(request, status.head!, current, currentBlob, current);
|
||||
if (request.action !== "delete") await this.assertEvidence(request.workspace, status.head!);
|
||||
const docs = { contract: `workspace-docs/${id}/contract.env.example`, readme: `workspace-docs/${id}/README.md` };
|
||||
if (request.action === "delete") {
|
||||
await this.repository.removeRegistryFile(path); await this.repository.removeRegistryFile(docs.contract); await this.repository.removeRegistryFile(docs.readme);
|
||||
} else {
|
||||
await this.repository.writeRegistryFile(path, serializeWorkspaceYaml(request.workspace));
|
||||
const rendered = renderWorkspaceDocs(request.workspace);
|
||||
await this.repository.writeRegistryFile(docs.contract, rendered.envExample); await this.repository.writeRegistryFile(docs.readme, rendered.markdown);
|
||||
}
|
||||
const committed = await this.repository.commitOnly([path, docs.contract, docs.readme], request.action === "delete" ? `Delete workspace ${id}` : `Publish workspace ${id}`);
|
||||
const commit = committed.head;
|
||||
if (!commit) throw new WorkspaceRegistryError("git_unavailable", "Workspace Git service is unavailable");
|
||||
if (request.action === "delete") return { id, commit, blob: request.baseBlob };
|
||||
return { id, commit, blob: await this.repository.blob(path) };
|
||||
}
|
||||
|
||||
private async currentDescriptor(commit: string, path: string): Promise<WorkspaceDescriptor | undefined> {
|
||||
try {
|
||||
return parseWorkspaceYaml(await this.repository.readWorkspaceAt(commit, path));
|
||||
@@ -117,7 +140,27 @@ export async function publishAddressedByAuthor(
|
||||
request: PublishWorkspaceRequest,
|
||||
addressed: Extract<RegistryAddressedRequestV1, { readonly mode: "create" }>,
|
||||
): Promise<RegistryAddressedResultV1> {
|
||||
await claimAddressedPublication(registry, addressed);
|
||||
await author.publish(request);
|
||||
return registry.publishAddressed({ ...addressed, mode: "resume" });
|
||||
const claimed = await claimAddressedPublication(registry, addressed);
|
||||
const resume = { ...addressed, mode: "resume" as const, runId: claimed.runId };
|
||||
if (claimed.phase === "terminal_durable") return registry.publishAddressed(resume);
|
||||
|
||||
let target = claimed.advertisedTargetCommit;
|
||||
if (claimed.phase === "request_claimed") {
|
||||
// Handle a kill between update-ref and the state-file transition: the immutable
|
||||
// local ref is sufficient evidence to finish the same deterministic run.
|
||||
const existingRef = await addressedRunRef(registry, claimed.runId);
|
||||
if (existingRef) {
|
||||
target = existingRef as typeof target;
|
||||
await advertiseAddressedPublication(registry, claimed.runId, existingRef);
|
||||
} else {
|
||||
const prepared = await author.recoverPrepared(request) ?? await author.prepare(request);
|
||||
target = prepared.commit as typeof target;
|
||||
// The ref and durable advertisement precede the push, so either crash boundary
|
||||
// resumes by pushing this exact commit rather than authoring a new one.
|
||||
await advertiseAddressedPublication(registry, claimed.runId, prepared.commit);
|
||||
}
|
||||
}
|
||||
if (!target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed author run has no pinned target");
|
||||
await author.pushPrepared({ commit: target });
|
||||
return registry.publishAddressed(resume);
|
||||
}
|
||||
|
||||
@@ -132,6 +132,11 @@ export class GitWorkspaceRepository {
|
||||
private validateRevision(revision: string): void {
|
||||
if (!/^[0-9a-f]{40}$/.test(revision)) throw new WorkspaceRegistryError("workspace_invalid", "Workspace revision is invalid");
|
||||
}
|
||||
async parentOf(revision: string): Promise<string | undefined> {
|
||||
this.validateRevision(revision);
|
||||
return (await this.gitOptional(["rev-parse", `${revision}^`]))?.trim();
|
||||
}
|
||||
|
||||
async workspacePathsAt(revision: string): Promise<string[]> {
|
||||
this.validateRevision(revision);
|
||||
const output = await this.git(["ls-tree", "-r", "--name-only", revision, "--", "workspaces"]);
|
||||
@@ -221,6 +226,24 @@ export class GitWorkspaceRepository {
|
||||
await rm(join(this.repoPath, path), { force: true });
|
||||
}
|
||||
|
||||
/** Push an already-created commit by its exact object ID. Repeating this is idempotent. */
|
||||
async pushExact(revision: string): Promise<void> {
|
||||
this.validateRevision(revision);
|
||||
try {
|
||||
await this.git(["push", "origin", `${revision}:refs/heads/${this.config.branch}`]);
|
||||
} catch (error) {
|
||||
throw this.sanitizeGitError(error);
|
||||
}
|
||||
}
|
||||
|
||||
/** Create a local publication commit without contacting the remote. */
|
||||
async commitOnly(paths: readonly string[], message: string): Promise<GitStatus> {
|
||||
if (paths.length === 0 || paths.some((path) => !this.isRegistryArtifactPath(path))) throw new WorkspaceRegistryError("workspace_invalid", "Workspace repository path is invalid");
|
||||
await this.git(["add", "--", ...paths]);
|
||||
await this.git(["commit", "-m", message], this.publicationIdentity());
|
||||
return await this.status();
|
||||
}
|
||||
|
||||
/** Commit and push a fixed set of validated artifact paths without exposing Git output. */
|
||||
async commitAndPush(paths: readonly string[], message: string): Promise<GitStatus> {
|
||||
if (paths.length === 0 || paths.some((path) => !this.isRegistryArtifactPath(path))) {
|
||||
|
||||
@@ -30,37 +30,29 @@ export interface AddressedWorkspacePublicationLeaseV1 extends BorrowedOrderedWor
|
||||
export interface CapabilityAwareRegistryPublicationParticipant<T> { readonly participantId: string; prepare(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1): Promise<T>; reconcile(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1, prepared: T, phase: RegistryAddressedPublicationPhaseV1): Promise<void>; }
|
||||
export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: string; readonly preparedSha256: Sha256Hex; }
|
||||
export interface CapabilityAwareRegistryPublicationSynchronizer { readonly synchronizerId: string; ensureForPublication(plan: RegistryAddressedPlanV1, capabilities: OrderedWorkspaceWriterCapabilitySet, phase: "planned" | "participants_prepared" | "publication_intent_durable" | "target_published"): Promise<RegistrySynchronizerPreparedV1>; }
|
||||
const preparedPublicationRuns = new WeakMap<CapabilityAwareRegistryPublicationLifecycleOwner, Array<{ readonly participant: CapabilityAwareRegistryPublicationParticipant<unknown>; readonly value: unknown; readonly lease: AddressedWorkspacePublicationLeaseV1 }>>();
|
||||
|
||||
export class CapabilityAwareRegistryPublicationLifecycleOwner {
|
||||
constructor() {}
|
||||
async run<T>(input: { readonly plan: RegistryAddressedPlanV1; readonly capabilities: OrderedWorkspaceWriterCapabilitySet; readonly participants: readonly CapabilityAwareRegistryPublicationParticipant<unknown>[]; readonly synchronizers: readonly CapabilityAwareRegistryPublicationSynchronizer[]; readonly action: () => Promise<T>; }): Promise<T> {
|
||||
// Reader gates are recursively nested so every changed workspace remains quiescent and
|
||||
// reader-exclusive until the publication callback, reconciliation, and terminal durability
|
||||
// have all settled. No participant or synchronizer is allowed to acquire a capability.
|
||||
// Reader gates are recursively nested and held until action returns. Production action
|
||||
// callbacks own all post-publication work (including terminal durability); this owner has
|
||||
// no post-action work which could accidentally outlive the quiescence gates.
|
||||
const prepared: Array<{ readonly participant: CapabilityAwareRegistryPublicationParticipant<unknown>; readonly value: unknown; readonly lease: AddressedWorkspacePublicationLeaseV1 }> = [];
|
||||
const enter = async (index: number): Promise<T> => {
|
||||
if (index >= input.capabilities.workspaceIds.length) {
|
||||
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "planned");
|
||||
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "participants_prepared");
|
||||
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");
|
||||
return result;
|
||||
preparedPublicationRuns.set(this, prepared);
|
||||
try { return await input.action(); }
|
||||
finally { preparedPublicationRuns.delete(this); }
|
||||
}
|
||||
const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId;
|
||||
return input.capabilities.forWorkspace(id, ({ rootLease, writerCapability }) =>
|
||||
writerCapability.runUnderSessionReadersExclusive(async readers => {
|
||||
const lease: AddressedWorkspacePublicationLeaseV1 = {
|
||||
workspaceId: id,
|
||||
rootLease,
|
||||
writerCapability,
|
||||
quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id),
|
||||
readers,
|
||||
};
|
||||
for (const participant of input.participants) {
|
||||
const value = await participant.prepare(input.plan, lease);
|
||||
prepared.push({ participant, value, lease });
|
||||
}
|
||||
const lease: AddressedWorkspacePublicationLeaseV1 = { workspaceId: id, rootLease, writerCapability, quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id), readers };
|
||||
for (const participant of input.participants) prepared.push({ participant, value: await participant.prepare(input.plan, lease), lease });
|
||||
return enter(index + 1);
|
||||
}),
|
||||
);
|
||||
@@ -69,6 +61,13 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner {
|
||||
}
|
||||
}
|
||||
|
||||
/** Run participant reconciliation while the lifecycle owner's nested gates are held. */
|
||||
export async function reconcilePreparedPublication(owner: CapabilityAwareRegistryPublicationLifecycleOwner, plan: RegistryAddressedPlanV1): Promise<void> {
|
||||
const prepared = preparedPublicationRuns.get(owner);
|
||||
if (!prepared) throw CONFLICT();
|
||||
for (const item of prepared) await item.participant.reconcile(plan, item.lease, item.value, "target_published");
|
||||
}
|
||||
|
||||
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 }
|
||||
|
||||
@@ -20,7 +20,7 @@ import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.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 RegistryPullAddressedPlanV1, type RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.js";
|
||||
import { RegistryAddressedPublicationStore, addressedRunId, canonicalBootstrapRequestDigest, registryDigest, CapabilityAwareRegistryPublicationLifecycleOwner, reconcilePreparedPublication, type CapabilityAwareRegistryPublicationParticipant, type CapabilityAwareRegistryPublicationSynchronizer, type RegistryAddressedPlanV1, type RegistryPullAddressedPlanV1, 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";
|
||||
|
||||
export interface WorkspaceRevision {
|
||||
@@ -221,7 +221,7 @@ export class WorkspaceRegistrySnapshotReader {
|
||||
readPinned(id: string, commit: string) { return this.ops.readPinned(id, commit); }
|
||||
}
|
||||
const snapshotReaders = new WeakMap<WorkspaceRegistry, WorkspaceRegistrySnapshotReader>();
|
||||
const addressedClaimers = new WeakMap<WorkspaceRegistry, (request: RegistryAddressedRequestV1) => Promise<void>>();
|
||||
const addressedClaimers = new WeakMap<WorkspaceRegistry, (request: RegistryAddressedRequestV1) => Promise<RegistryAddressedPublicationStateV1>>();
|
||||
|
||||
export function workspaceRegistrySnapshotReader(registry: WorkspaceRegistry): WorkspaceRegistrySnapshotReader {
|
||||
const reader = snapshotReaders.get(registry);
|
||||
@@ -230,12 +230,30 @@ export function workspaceRegistrySnapshotReader(registry: WorkspaceRegistry): Wo
|
||||
}
|
||||
|
||||
/** Package-private author/publication coordinator seam: claim before author Git network. */
|
||||
export async function claimAddressedPublication(registry: WorkspaceRegistry, request: RegistryAddressedRequestV1): Promise<void> {
|
||||
export async function claimAddressedPublication(registry: WorkspaceRegistry, request: RegistryAddressedRequestV1): Promise<RegistryAddressedPublicationStateV1> {
|
||||
const claim = addressedClaimers.get(registry);
|
||||
if (!claim) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized");
|
||||
return claim(request);
|
||||
}
|
||||
|
||||
/** Authoring barrier helpers intentionally stay outside the frozen registry class API. */
|
||||
export async function addressedRunRef(registry: WorkspaceRegistry, runId: string): Promise<string | undefined> {
|
||||
return registryContext(registry).repository.runRef(runId);
|
||||
}
|
||||
export async function advertiseAddressedPublication(registry: WorkspaceRegistry, runId: RegistryAddressedRequestV1["runId"], target: string): Promise<RegistryAddressedPublicationStateV1> {
|
||||
await registryContext(registry).repository.ensureLayout();
|
||||
return registryContext(registry).lock.run(async () => {
|
||||
const store = new RegistryAddressedPublicationStore(registryContext(registry).repository.root);
|
||||
const state = await store.read(runId);
|
||||
if (state.phase !== "request_claimed") {
|
||||
if (state.advertisedTargetCommit !== target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed target conflicts with durable author run");
|
||||
return state;
|
||||
}
|
||||
await registryContext(registry).repository.ensureRunRef(runId, target);
|
||||
return store.transition(runId, "target_advertised", { advertisedTargetCommit: target as Revision40, immutableTargetRef: `refs/thoth/addressed-runs/${runId}/target` });
|
||||
});
|
||||
}
|
||||
|
||||
export class WorkspaceRegistry {
|
||||
private readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory;
|
||||
private readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner;
|
||||
@@ -296,14 +314,19 @@ export class WorkspaceRegistry {
|
||||
return { kind: "bootstrap_terminal", result: result as Extract<RegistryAddressedResultV1, { operation: "registry_bootstrap" }>, snapshot: this.#addressedSnapshot(published, identity.requestSha256) };
|
||||
}
|
||||
|
||||
async #claimAddressed(request: RegistryAddressedRequestV1): Promise<void> {
|
||||
async #claimAddressed(request: RegistryAddressedRequestV1): Promise<RegistryAddressedPublicationStateV1> {
|
||||
if (request.mode !== "create") throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Only create requests may claim a run");
|
||||
await registryContext(this).repository.ensureLayout();
|
||||
await registryContext(this).lock.run(async () => {
|
||||
return registryContext(this).lock.run(async () => {
|
||||
const store = new RegistryAddressedPublicationStore(registryContext(this).repository.root);
|
||||
// A retry may receive a fresh proposed run ID, but the request digest is the
|
||||
// idempotency key. Reuse exactly one matching author run, including terminal replays.
|
||||
const existing = (await store.scan()).filter(job => job.operation === request.operation && job.requestSha256 === request.requestSha256 && job.installationIdentitySha256 === request.installationIdentitySha256 && job.repositoryIdentitySha256 === request.repositoryIdentitySha256 && job.remoteRefIdentitySha256 === request.remoteRefIdentitySha256);
|
||||
if (existing.length > 1) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed request has multiple durable runs");
|
||||
if (existing.length === 1) return existing[0]!;
|
||||
let base: ActiveState | undefined;
|
||||
if (request.operation === "registry_pull") base = await this.#snapshotState(request.expectedBaseCommit);
|
||||
await store.claim(request, base ? { commit: base.head as Revision40, manifestSha256: registryDigest(base) as never, workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)) } : null);
|
||||
return store.claim(request, base ? { commit: base.head as Revision40, manifestSha256: registryDigest(base) as never, workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)) } : null);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -340,8 +363,8 @@ export class WorkspaceRegistry {
|
||||
async #executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryAddressedResultV1> {
|
||||
let state = initial;
|
||||
if (state.phase === "terminal_durable") {
|
||||
const pinned = await registryContext(this).repository.runRef(state.runId);
|
||||
if (pinned !== state.advertisedTargetCommit) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed target ref drifted");
|
||||
// Terminal state is an authenticated durable result. Replay it without consulting
|
||||
// mutable Git refs or the network; ref drift is only relevant to reconciliation.
|
||||
return store.readTerminalResult(state.runId, state.terminalResultSha256!);
|
||||
}
|
||||
const operation = state.operation;
|
||||
@@ -409,25 +432,29 @@ export class WorkspaceRegistry {
|
||||
}
|
||||
if (state.phase === "publication_intent_durable") {
|
||||
const activeBefore = await this.#tryActiveState();
|
||||
if (activeBefore?.head === target) {
|
||||
publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target";
|
||||
} else {
|
||||
await this.#publishSnapshotPointer(target!);
|
||||
}
|
||||
const activeDigest = activeBefore ? registryDigest(activeBefore) : undefined;
|
||||
const baseAllowed = operation === "registry_bootstrap"
|
||||
? activeBefore === undefined
|
||||
: (activeBefore !== undefined && activeBefore.head === state.baseCommit && activeDigest === state.baseManifestSha256);
|
||||
const targetAllowed = activeBefore !== undefined && activeBefore.head === target && activeDigest === state.targetManifestSha256;
|
||||
if (!baseAllowed && !targetAllowed) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry identity conflicts with addressed publication");
|
||||
if (targetAllowed) publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target";
|
||||
else await this.#publishSnapshotPointer(target!);
|
||||
state = await store.transition(state.runId, "target_published", { publishedActiveStateSha256: registryDigest(await this.#activeState()) as never });
|
||||
}
|
||||
if (state.phase === "target_published") {
|
||||
// A restart at this barrier must verify the exact active target before success.
|
||||
const active = await this.#tryActiveState();
|
||||
if (!active || active.head !== target || registryDigest(active) !== state.targetManifestSha256) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Active registry target drifted");
|
||||
if (resumedAtTargetPublished) publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target";
|
||||
for (const synchronizer of registryContext(this).synchronizers) await synchronizer.ensureForPublication(plan, capabilities, "target_published");
|
||||
await reconcilePreparedPublication(registryContext(this).lifecycleOwner, plan);
|
||||
output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable" as const, publication } as RegistryAddressedResultV1;
|
||||
await store.storeTerminalResult(state.runId, output);
|
||||
state = await store.read(state.runId);
|
||||
}
|
||||
return undefined;
|
||||
}});
|
||||
// The lifecycle callback has now reconciled target state. Persist terminal while
|
||||
// ordered writer/root capabilities remain owned, immediately before settlement.
|
||||
if (state.phase === "target_published") {
|
||||
if (resumedAtTargetPublished) {
|
||||
publication = operation === "registry_pull" && state.baseCommit === target ? "unchanged" : "reconciled_target";
|
||||
}
|
||||
output = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: registryDigest(plan) as never, phase: "terminal_durable" as const, publication } as RegistryAddressedResultV1;
|
||||
await store.storeTerminalResult(state.runId, output);
|
||||
state = await store.read(state.runId);
|
||||
}
|
||||
});
|
||||
if (!output) output = await store.readTerminalResult(state.runId, state.terminalResultSha256!);
|
||||
return output;
|
||||
|
||||
Reference in New Issue
Block a user