fix addressed registry recovery invariants
This commit is contained in:
@@ -7,7 +7,7 @@ import { getPrincipal } from "../auth/auth.js";
|
||||
import type { PrincipalContext } from "../auth/principal.js";
|
||||
import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
||||
import type { ListModelsFn } from "./meta.js";
|
||||
import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js";
|
||||
import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js";
|
||||
|
||||
@@ -43,6 +43,7 @@ export function sessionRoutes(
|
||||
maintenanceBarrier: MaintenanceBarrier;
|
||||
},
|
||||
) {
|
||||
const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(d.workspaceRegistry); } catch { return d.workspaceRegistry as any; } };
|
||||
const lifecycleTails = new Map<string, Promise<void>>();
|
||||
const boundRuntimes = new Map<
|
||||
string,
|
||||
@@ -99,9 +100,9 @@ export function sessionRoutes(
|
||||
app.addHook("onResponse", async (req) => { admissionLeases.get(req)?.(); });
|
||||
|
||||
/** Include retained historical descriptors after the same addressed selector used by status/list. */
|
||||
const sessionRevisions = async () => {
|
||||
const sessionRevisions = async (): Promise<Awaited<ReturnType<ReturnType<typeof snapshotReader>["listRetainedSnapshots"]>>> => {
|
||||
await d.workspaceRegistry.ensureBootstrapAddressed(d.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(d.workspaceRegistry));
|
||||
return d.workspaceRegistry.listRetainedSnapshots();
|
||||
return snapshotReader().listRetainedSnapshots();
|
||||
};
|
||||
|
||||
const isNotFound = (error: unknown) =>
|
||||
@@ -141,7 +142,7 @@ export function sessionRoutes(
|
||||
throw error;
|
||||
}
|
||||
};
|
||||
let revisions: Awaited<ReturnType<typeof d.workspaceRegistry.listRetainedSnapshots>>;
|
||||
let revisions: Awaited<ReturnType<ReturnType<typeof snapshotReader>["listRetainedSnapshots"]>>;
|
||||
try {
|
||||
revisions = await sessionRevisions();
|
||||
} catch (registryError) {
|
||||
@@ -169,7 +170,7 @@ export function sessionRoutes(
|
||||
const saved = located.manifest as { workspace_id?: string; workspace_revision?: string };
|
||||
if (!saved.workspace_id || !saved.workspace_revision) return located;
|
||||
try {
|
||||
const pinned = await d.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision);
|
||||
const pinned = await snapshotReader().readPinned(saved.workspace_id, saved.workspace_revision);
|
||||
const workspace = validateOperationalWorkspace(pinned.workspace);
|
||||
return {
|
||||
...located,
|
||||
@@ -320,7 +321,7 @@ export function sessionRoutes(
|
||||
code: "workspace_revision_unavailable",
|
||||
});
|
||||
}
|
||||
let revisionLease: Awaited<ReturnType<WorkspaceRegistry["acquireSessionRevision"]>> | undefined;
|
||||
let revisionLease: Awaited<ReturnType<ReturnType<typeof snapshotReader>["acquireSessionRevision"]>> | undefined;
|
||||
let manifestPersisted = false;
|
||||
try {
|
||||
let workspaceConfigPath: string | undefined;
|
||||
@@ -331,11 +332,11 @@ export function sessionRoutes(
|
||||
if (requestedWorkspaceId) {
|
||||
try {
|
||||
const registry = d.workspaceRegistry as Partial<WorkspaceRegistry>;
|
||||
const resolved = typeof registry.acquireSessionRevision === "function"
|
||||
? await registry.acquireSessionRevision.call(d.workspaceRegistry, requestedWorkspaceId)
|
||||
: await d.workspaceRegistry.read(requestedWorkspaceId);
|
||||
const resolved = typeof (registry as any).acquireSessionRevision === "function"
|
||||
? await snapshotReader().acquireSessionRevision(requestedWorkspaceId)
|
||||
: await snapshotReader().read(requestedWorkspaceId);
|
||||
if ("markPersisted" in resolved && "abort" in resolved) {
|
||||
revisionLease = resolved as Awaited<ReturnType<WorkspaceRegistry["acquireSessionRevision"]>>;
|
||||
revisionLease = resolved as Awaited<ReturnType<ReturnType<typeof snapshotReader>["acquireSessionRevision"]>>;
|
||||
}
|
||||
if (!d.workspaceRuntimeSupport(resolved.workspace)) {
|
||||
return reply.code(409).send({
|
||||
@@ -473,7 +474,7 @@ export function sessionRoutes(
|
||||
const runner = runnerFor(scopedPrincipal);
|
||||
const revisions = await sessionRevisions();
|
||||
const lists = await Promise.all(revisions
|
||||
.map((revision) => runner.sessionList(revision.snapshotPath) as Promise<SessionRow[]>));
|
||||
.map((revision: { snapshotPath: string }) => runner.sessionList(revision.snapshotPath) as Promise<SessionRow[]>));
|
||||
const sessions = new Map<string, SessionRow>();
|
||||
for (const row of lists.flat()) {
|
||||
if (!sessions.has(row.id)) sessions.set(row.id, row);
|
||||
|
||||
@@ -3,13 +3,14 @@ import type { ThtRunner } from "../tht/tht-runner.js";
|
||||
import { getPrincipal } from "../auth/auth.js";
|
||||
import type { PrincipalContext } from "../auth/principal.js";
|
||||
import type { Settings } from "../settings/settings-store.js";
|
||||
import { workspaceRegistryRecoveryIdentity, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
|
||||
export function sqlRoutes(app: FastifyInstance, deps: {
|
||||
tht: ThtRunner; getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||
workspaceRegistry: WorkspaceRegistry;
|
||||
workspaceRegistryRecoveryIdentity?: () => ReturnType<typeof workspaceRegistryRecoveryIdentity>;
|
||||
}): void {
|
||||
const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(deps.workspaceRegistry); } catch { return deps.workspaceRegistry as any; } };
|
||||
const runnerFor = (principal: PrincipalContext): any => {
|
||||
const runner = deps.tht as any;
|
||||
return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner;
|
||||
@@ -21,14 +22,14 @@ export function sqlRoutes(app: FastifyInstance, deps: {
|
||||
const runner = runnerFor(principal);
|
||||
if (typeof runner.sessionShow !== "function") return { manifest: {}, workspace: legacyWorkspace };
|
||||
await deps.workspaceRegistry.ensureBootstrapAddressed(deps.workspaceRegistryRecoveryIdentity?.() ?? workspaceRegistryRecoveryIdentity(deps.workspaceRegistry));
|
||||
const revisions = await deps.workspaceRegistry.listRetainedSnapshots();
|
||||
const revisions = await snapshotReader().listRetainedSnapshots();
|
||||
for (const revision of revisions) {
|
||||
try {
|
||||
const manifest = await runner.sessionShow(id, revision.snapshotPath);
|
||||
if (!manifest) continue;
|
||||
const saved = manifest as { workspace_id?: string; workspace_revision?: string };
|
||||
if (saved.workspace_id && saved.workspace_revision) {
|
||||
const pinned = await deps.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision);
|
||||
const pinned = await snapshotReader().readPinned(saved.workspace_id, saved.workspace_revision);
|
||||
return {
|
||||
manifest,
|
||||
workspace: pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath,
|
||||
|
||||
@@ -6,7 +6,7 @@ import yauzl from "yauzl";
|
||||
import yazl from "yazl";
|
||||
import { z } from "zod";
|
||||
import type { WorkspaceRegistryConfig } from "../workspaces/types.js";
|
||||
import type { WorkspaceAuthorGitService } from "../workspaces/author-git-service.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 {
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
type PublishWorkspaceRequest,
|
||||
type WorkspaceRegistry,
|
||||
workspaceRegistryRecoveryIdentity,
|
||||
workspaceRegistrySnapshotReader,
|
||||
} from "../workspaces/registry.js";
|
||||
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js";
|
||||
import type { Revision40 } from "../workspaces/workspace-lock-root-lease.js";
|
||||
@@ -283,6 +284,13 @@ function publishRequest(value: unknown): PublishWorkspaceRequest {
|
||||
}
|
||||
|
||||
export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps): void {
|
||||
const snapshotReader = () => { try { return workspaceRegistrySnapshotReader(deps.registry); } catch { return deps.registry as any; } };
|
||||
const readPinned = async (id: string, commit: string) => {
|
||||
const reader = snapshotReader();
|
||||
if (typeof reader.readPinned === "function") return reader.readPinned(id, commit);
|
||||
const legacy = await reader.read(id);
|
||||
return { workspace: legacy.workspace, workspaceConfigPath: legacy.revision.snapshotPath };
|
||||
};
|
||||
app.register(multipart, {
|
||||
limits: { fileSize: deps.config.maxImportBytes, files: 1, fields: 0, parts: 1 },
|
||||
throwFileSizeLimit: true,
|
||||
@@ -325,7 +333,7 @@ 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: deps.snapshotPath(item.revision, item.workspaceId) }));
|
||||
return await Promise.all(revisions.map(async revision => {
|
||||
const pinned = await deps.registry.readPinned(revision.id, revision.commit);
|
||||
const pinned = await readPinned(revision.id, revision.commit);
|
||||
const workspace = pinned.workspace;
|
||||
const exactRevision = { ...revision, snapshotPath: pinned.workspaceConfigPath };
|
||||
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 };
|
||||
@@ -336,7 +344,11 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
app.get("/workspaces/:id", async (request, reply) => {
|
||||
try {
|
||||
const { id } = z.object({ id: workspaceId }).parse(request.params);
|
||||
return await deps.registry.read(id);
|
||||
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
|
||||
const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id);
|
||||
if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
|
||||
const pinned = await readPinned(id, item.revision);
|
||||
return { workspace: pinned.workspace, revision: { id, commit: item.revision, blob: item.descriptorBlob, snapshotPath: pinned.workspaceConfigPath } };
|
||||
} catch (error) {
|
||||
return errorReply(reply, error);
|
||||
}
|
||||
@@ -355,7 +367,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
app.post("/workspaces/:id/test", async (request, reply) => {
|
||||
try {
|
||||
const { id } = z.object({ id: workspaceId }).parse(request.params);
|
||||
const { workspace } = await deps.registry.read(id);
|
||||
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
|
||||
const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id);
|
||||
if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
|
||||
const { workspace } = await readPinned(id, item.revision);
|
||||
let operational: CanonicalWorkspace;
|
||||
try {
|
||||
operational = validateOperationalWorkspace(workspace);
|
||||
@@ -387,10 +402,16 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
|
||||
expectedBaseCommit: requestValue.baseCommit as Revision40,
|
||||
};
|
||||
// Authoring is intentionally split from activation: this service creates and pushes
|
||||
// the remote commit only. The addressed registry is the sole pointer owner.
|
||||
await deps.authorService.publish(requestValue);
|
||||
const result = await deps.registry.publishAddressed(addressed);
|
||||
// Claim the addressed identity before any author Git network. A crash after
|
||||
// push therefore leaves a durable run that can be resumed with this exact ID.
|
||||
let result: RegistryAddressedResultV1;
|
||||
if (Object.prototype.hasOwnProperty.call(deps.registry, "snapshotReader")) {
|
||||
result = await publishAddressedByAuthor(deps.registry, deps.authorService, requestValue, addressed);
|
||||
} else {
|
||||
// Test/dry-run registry doubles predate the package-private coordinator.
|
||||
await deps.authorService.publish(requestValue);
|
||||
result = await deps.registry.publishAddressed(addressed);
|
||||
}
|
||||
return { revision: revisionFromPublication(result, id) };
|
||||
} catch (error) {
|
||||
return errorReply(reply, error);
|
||||
@@ -400,7 +421,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
|
||||
app.get("/workspaces/:id/export", async (request, reply) => {
|
||||
try {
|
||||
const { id } = z.object({ id: workspaceId }).parse(request.params);
|
||||
const { workspace } = await deps.registry.read(id);
|
||||
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
|
||||
const item = snapshot.workspaces.find(candidate => candidate.workspaceId === id);
|
||||
if (!item) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
|
||||
const { workspace } = await readPinned(id, item.revision);
|
||||
const canonical = validateWorkspaceDescriptor(workspace);
|
||||
const bundle = await exportBundle(canonical);
|
||||
return reply
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { WorkspaceConflictError, type PublishWorkspaceRequest } from "./registry.js";
|
||||
import { WorkspaceConflictError, 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";
|
||||
|
||||
@@ -104,3 +105,19 @@ export class WorkspaceAuthorGitService {
|
||||
return new WorkspaceConflictError(fields, expected, actual, base, local, remote);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Repository-author/publication friend. The addressed claim is durable before the
|
||||
* author Git service performs pull/stage/push network, and activation always resumes
|
||||
* that exact run identity.
|
||||
*/
|
||||
export async function publishAddressedByAuthor(
|
||||
registry: WorkspaceRegistry,
|
||||
author: WorkspaceAuthorGitService,
|
||||
request: PublishWorkspaceRequest,
|
||||
addressed: Extract<RegistryAddressedRequestV1, { readonly mode: "create" }>,
|
||||
): Promise<RegistryAddressedResultV1> {
|
||||
await claimAddressedPublication(registry, addressed);
|
||||
await author.publish(request);
|
||||
return registry.publishAddressed({ ...addressed, mode: "resume" });
|
||||
}
|
||||
|
||||
@@ -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 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 terminalPublication: "target" | "reconciled_target" | "unchanged" | 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>; readonly afterPublication?: (result: T) => Promise<void>; }): 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>; }): 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,7 +45,6 @@ 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;
|
||||
@@ -107,18 +106,16 @@ 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()), base?: { readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] }): Promise<RegistryAddressedPublicationStateV1> {
|
||||
async claim(request: RegistryAddressedRequestV1, base: { readonly commit: Revision40; readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[] } | null): Promise<RegistryAddressedPublicationStateV1> {
|
||||
await this.dirs();
|
||||
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 (!RUN.test(request.runId) || request.mode !== "create") 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, 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;
|
||||
if (request.operation === "registry_pull" && base && request.expectedBaseCommit !== base.commit) throw CONFLICT();
|
||||
if (request.operation === "registry_bootstrap" && (base !== null || request.expectedBaseCommit !== null)) throw CONFLICT();
|
||||
if (request.operation === "registry_pull" && !base) throw CONFLICT();
|
||||
const state = { schemaVersion: 1, runId: request.runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${request.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, terminalPublication: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" ? request.expectedBaseCommit : null, baseManifestSha256: base?.manifestSha256 ?? null, baseWorkspaces: base?.workspaces ?? [], changedSetRule: null };
|
||||
await this.durable(this.path(request.runId), state, true);
|
||||
return state as unknown as RegistryAddressedPublicationStateV1;
|
||||
}
|
||||
|
||||
async storeTerminalResult(runId: RegistryRunId32, result: RegistryAddressedResultV1): Promise<void> {
|
||||
@@ -127,7 +124,7 @@ export class RegistryAddressedPublicationStore {
|
||||
// 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 });
|
||||
await this.transition(runId, "terminal_durable", { terminalResultSha256: expected, terminalPublication: result.publication });
|
||||
}
|
||||
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();
|
||||
@@ -150,7 +147,7 @@ export class RegistryAddressedPublicationStore {
|
||||
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;
|
||||
plan, planSha256: state.planSha256!, phase: "terminal_durable" as const, publication: state.terminalPublication! } as RegistryAddressedResultV1;
|
||||
if (registryDigest(result) !== expectedDigest) throw CONFLICT();
|
||||
return result;
|
||||
}
|
||||
@@ -161,12 +158,18 @@ 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","terminalPublication","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.terminalPublication !== null && s.terminalPublication !== "target" && s.terminalPublication !== "reconciled_target" && s.terminalPublication !== "unchanged") 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();
|
||||
// The advertisement is the sole target identity. Every later durable phase must
|
||||
// carry exactly that OID; accepting a different fetched or planned OID would make
|
||||
// resume capable of silently publishing an object other than the claimed target.
|
||||
if (s.advertisedTargetCommit !== null && s.fetchedTargetCommit !== null && s.fetchedTargetCommit !== s.advertisedTargetCommit) throw CONFLICT();
|
||||
if (s.advertisedTargetCommit !== null && s.targetCommit !== null && s.targetCommit !== s.advertisedTargetCommit) throw CONFLICT();
|
||||
for (const key of ["requestSha256","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256"] as const) if (!SHA.test(s[key])) throw CONFLICT();
|
||||
for (const key of ["targetManifestSha256","changedSetSha256","planSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256"] as const) if (s[key] !== null && !SHA.test(s[key])) throw CONFLICT();
|
||||
if (s.immutableTargetRef !== null && s.immutableTargetRef !== `refs/thoth/addressed-runs/${runId}/target`) throw CONFLICT();
|
||||
@@ -175,13 +178,13 @@ export class RegistryAddressedPublicationStore {
|
||||
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"],
|
||||
planned: ["targetCommit", "targetManifestSha256", "targetWorkspaces", "changedWorkspaceIds", "planSha256", "changedSetSha256", "changedSetRule"], participants_prepared: ["participantsSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"],
|
||||
};
|
||||
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();
|
||||
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"],
|
||||
participants_prepared: ["participantsSha256", "synchronizersSha256"], publication_intent_durable: ["publicationIntentSha256"], target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"],
|
||||
};
|
||||
const currentIndex = PHASES.indexOf(s.phase);
|
||||
let chainPrevious: string | null = null;
|
||||
@@ -195,26 +198,23 @@ export class RegistryAddressedPublicationStore {
|
||||
}
|
||||
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", "terminalPublication", "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();
|
||||
}
|
||||
if ("fetchedTargetCommit" in patch && old.advertisedTargetCommit !== null && patch.fetchedTargetCommit !== old.advertisedTargetCommit) throw CONFLICT();
|
||||
if ("targetCommit" in patch && old.advertisedTargetCommit !== null && patch.targetCommit !== old.advertisedTargetCommit) 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"],
|
||||
target_published: ["publishedActiveStateSha256"], terminal_durable: ["terminalResultSha256", "terminalPublication"],
|
||||
};
|
||||
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) {
|
||||
@@ -229,7 +229,7 @@ export class RegistryAddressedPublicationStore {
|
||||
participants_prepared: ["participantsSha256"],
|
||||
publication_intent_durable: ["publicationIntentSha256"],
|
||||
target_published: ["publishedActiveStateSha256"],
|
||||
terminal_durable: ["terminalResultSha256"],
|
||||
terminal_durable: ["terminalResultSha256", "terminalPublication"],
|
||||
};
|
||||
for (const key of required[phase]) if (next[key] === null || next[key] === undefined) throw CONFLICT();
|
||||
await this.durable(this.path(runId), next);
|
||||
|
||||
@@ -186,9 +186,48 @@ export function createWorkspaceRegistry(config: WorkspaceRegistryConfig, deps: W
|
||||
lifecycleOwner: input.lifecycleOwner, participants: input.participants, synchronizers: input.synchronizers,
|
||||
installationIdentity, repositoryIdentity, remoteIdentity,
|
||||
});
|
||||
// Compatibility bindings are instance-owned and intentionally absent from the
|
||||
// frozen WorkspaceRegistry prototype. New callers receive the separate reader.
|
||||
const reader = workspaceRegistrySnapshotReader(registry);
|
||||
Object.defineProperties(registry, {
|
||||
listRetainedSnapshots: { value: reader.listRetainedSnapshots.bind(reader) },
|
||||
read: { value: reader.read.bind(reader) },
|
||||
acquireSessionRevision: { value: reader.acquireSessionRevision.bind(reader) },
|
||||
readPinned: { value: reader.readPinned.bind(reader) },
|
||||
snapshotReader: { value: reader },
|
||||
});
|
||||
return registry;
|
||||
}
|
||||
|
||||
interface SnapshotReaderOps {
|
||||
listRetainedSnapshots(): Promise<WorkspaceRevision[]>;
|
||||
read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }>;
|
||||
acquireSessionRevision(id: string): Promise<SessionRevisionLease>;
|
||||
readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }>;
|
||||
}
|
||||
export class WorkspaceRegistrySnapshotReader {
|
||||
constructor(private readonly ops: SnapshotReaderOps) {}
|
||||
listRetainedSnapshots() { return this.ops.listRetainedSnapshots(); }
|
||||
read(id: string) { return this.ops.read(id); }
|
||||
acquireSessionRevision(id: string) { return this.ops.acquireSessionRevision(id); }
|
||||
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>>();
|
||||
|
||||
export function workspaceRegistrySnapshotReader(registry: WorkspaceRegistry): WorkspaceRegistrySnapshotReader {
|
||||
const reader = snapshotReaders.get(registry);
|
||||
if (!reader) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized");
|
||||
return reader;
|
||||
}
|
||||
|
||||
/** Package-private author/publication coordinator seam: claim before author Git network. */
|
||||
export async function claimAddressedPublication(registry: WorkspaceRegistry, request: RegistryAddressedRequestV1): Promise<void> {
|
||||
const claim = addressedClaimers.get(registry);
|
||||
if (!claim) throw new WorkspaceRegistryError("git_unavailable", "Workspace registry is not initialized");
|
||||
return claim(request);
|
||||
}
|
||||
|
||||
export class WorkspaceRegistry {
|
||||
private readonly rootLeaseFactory: VerifiedWorkspaceLockRootLeaseFactory;
|
||||
private readonly lifecycleOwner: CapabilityAwareRegistryPublicationLifecycleOwner;
|
||||
@@ -206,6 +245,14 @@ export class WorkspaceRegistry {
|
||||
this.lifecycleOwner = input.lifecycleOwner;
|
||||
this.participants = input.participants;
|
||||
this.synchronizers = input.synchronizers;
|
||||
const reader = new WorkspaceRegistrySnapshotReader({
|
||||
listRetainedSnapshots: () => this.#listRetainedSnapshots(),
|
||||
read: (id) => this.#read(id),
|
||||
acquireSessionRevision: (id) => this.#acquireSessionRevision(id),
|
||||
readPinned: (id, commit) => this.#readPinned(id, commit),
|
||||
});
|
||||
snapshotReaders.set(this, reader);
|
||||
addressedClaimers.set(this, (request) => this.#claimAddressed(request));
|
||||
}
|
||||
|
||||
async ensureBootstrapAddressed(identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryEnsureBootstrapAddressedResultV1> {
|
||||
@@ -234,28 +281,41 @@ export class WorkspaceRegistry {
|
||||
installationIdentitySha256: identity.installationIdentitySha256, repositoryIdentitySha256: identity.repositoryIdentitySha256,
|
||||
expectedBaseCommit: null, remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
|
||||
};
|
||||
state = await store.claim(request, request.runId);
|
||||
state = await store.claim(request, null);
|
||||
}
|
||||
const result = await this.#executeAddressedLocked(store, state, identity);
|
||||
const published = await this.#activeState();
|
||||
return { kind: "bootstrap_terminal", result: result as Extract<RegistryAddressedResultV1, { operation: "registry_bootstrap" }>, snapshot: this.#addressedSnapshot(published, identity.requestSha256) };
|
||||
}
|
||||
|
||||
async #claimAddressed(request: RegistryAddressedRequestV1): Promise<void> {
|
||||
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 () => {
|
||||
const store = new RegistryAddressedPublicationStore(registryContext(this).repository.root);
|
||||
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);
|
||||
});
|
||||
}
|
||||
|
||||
async publishAddressed(request: RegistryAddressedRequestV1): Promise<RegistryAddressedResultV1> {
|
||||
await registryContext(this).repository.ensureLayout();
|
||||
return registryContext(this).lock.run(async () => {
|
||||
const store = new RegistryAddressedPublicationStore(registryContext(this).repository.root);
|
||||
let state: RegistryAddressedPublicationStateV1;
|
||||
try { state = await store.read(request.runId); }
|
||||
catch {
|
||||
if (request.mode === "resume") {
|
||||
try { state = await store.read(request.runId); }
|
||||
catch { throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed resume requires an existing durable run"); }
|
||||
} else {
|
||||
let base: ActiveState | undefined;
|
||||
if (request.operation === "registry_pull" && request.mode === "create") {
|
||||
base = await this.#snapshotState(request.expectedBaseCommit);
|
||||
}
|
||||
state = await store.claim(request, request.runId, base ? {
|
||||
manifestSha256: registryDigest(base) as never,
|
||||
if (request.operation === "registry_pull") base = await this.#snapshotState(request.expectedBaseCommit);
|
||||
// claim is an exclusive O_CREAT|O_EXCL operation: an existing same-ID run is
|
||||
// always a create conflict, never an implicit replay.
|
||||
state = await store.claim(request, base ? {
|
||||
commit: base.head as Revision40, manifestSha256: registryDigest(base) as never,
|
||||
workspaces: base.revisions.map(revision => this.#manifestIdentity(revision)),
|
||||
} : undefined);
|
||||
} : null);
|
||||
}
|
||||
const identity: RegistryBootstrapRecoveryIdentityV1 = {
|
||||
operation: "registry_bootstrap", requestSha256: request.requestSha256,
|
||||
@@ -271,8 +331,14 @@ export class WorkspaceRegistry {
|
||||
|
||||
async #executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryAddressedResultV1> {
|
||||
let state = initial;
|
||||
if (state.phase === "terminal_durable") return store.readTerminalResult(state.runId, state.terminalResultSha256!);
|
||||
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");
|
||||
return store.readTerminalResult(state.runId, state.terminalResultSha256!);
|
||||
}
|
||||
const operation = state.operation;
|
||||
let publication: "target" | "reconciled_target" | "unchanged" = "target";
|
||||
const resumedAtTargetPublished = state.phase === "target_published";
|
||||
let target = state.advertisedTargetCommit ?? state.fetchedTargetCommit;
|
||||
if (state.phase === "request_claimed") {
|
||||
const status = operation === "registry_bootstrap" ? await registryContext(this).repository.bootstrap() : await registryContext(this).repository.pull();
|
||||
@@ -290,6 +356,12 @@ export class WorkspaceRegistry {
|
||||
}
|
||||
target = state.fetchedTargetCommit ?? target;
|
||||
if (!target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt");
|
||||
// Once advertised, the run-specific ref is immutable evidence. Every resume after
|
||||
// the fetch barrier revalidates it before reading or publishing any bytes.
|
||||
if (state.phase !== "target_advertised" && state.phase !== "request_claimed") {
|
||||
const pinned = await registryContext(this).repository.runRef(state.runId);
|
||||
if (pinned !== target) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Addressed target ref drifted");
|
||||
}
|
||||
if (state.phase === "target_fetched") {
|
||||
await this.#materialize(target);
|
||||
const targetState = await this.#snapshotState(target);
|
||||
@@ -311,30 +383,44 @@ export class WorkspaceRegistry {
|
||||
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 });
|
||||
}
|
||||
const plan = await this.#planForState(state);
|
||||
const leases = [];
|
||||
for (const id of plan.changedWorkspaceIds) leases.push(await registryContext(this).rootLeaseFactory.acquireOrProvision(registryContext(this).rootLeaseFactory.canonicalInput(id)));
|
||||
const leases: Awaited<ReturnType<VerifiedWorkspaceLockRootLeaseFactory["acquireOrProvision"]>>[] = [];
|
||||
try {
|
||||
for (const id of plan.changedWorkspaceIds) leases.push(await registryContext(this).rootLeaseFactory.acquireOrProvision(registryContext(this).rootLeaseFactory.canonicalInput(id)));
|
||||
} catch (error) {
|
||||
for (const lease of [...leases].reverse()) { try { await lease.close(); } catch {} }
|
||||
throw error;
|
||||
}
|
||||
let output: RegistryAddressedResultV1 | undefined;
|
||||
await runUnderOrderedWorkspaceWriterLocks(leases, capabilities => registryContext(this).lifecycleOwner.run({ plan, capabilities, participants: registryContext(this).participants, synchronizers: registryContext(this).synchronizers, action: async () => {
|
||||
if (state.phase === "planned") {
|
||||
state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(registryContext(this).participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(registryContext(this).synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never });
|
||||
}
|
||||
if (state.phase === "participants_prepared") {
|
||||
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.#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.
|
||||
await runUnderOrderedWorkspaceWriterLocks(leases, async capabilities => {
|
||||
await registryContext(this).lifecycleOwner.run({ plan, capabilities, participants: registryContext(this).participants, synchronizers: registryContext(this).synchronizers, action: async () => {
|
||||
if (state.phase === "planned") {
|
||||
state = await store.transition(state.runId, "participants_prepared", { participantsSha256: registryDigest(registryContext(this).participants.map(participant => participant.participantId)) as never, synchronizersSha256: registryDigest(registryContext(this).synchronizers.map(synchronizer => synchronizer.synchronizerId)) as never });
|
||||
}
|
||||
if (state.phase === "participants_prepared") {
|
||||
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") {
|
||||
const activeBefore = await this.#tryActiveState();
|
||||
if (activeBefore?.head === target) {
|
||||
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 });
|
||||
}
|
||||
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") {
|
||||
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;
|
||||
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;
|
||||
}
|
||||
@@ -352,7 +438,7 @@ export class WorkspaceRegistry {
|
||||
return { schemaVersion: 1, commit: state.head as RegistryActiveSnapshotV1["commit"], manifestSha256: registryDigest(state) as RegistryActiveSnapshotV1["manifestSha256"], workspaces: ids.map(id => { const revision = state.revisions.find(candidate => candidate.id === id)!; return this.#manifestIdentity(revision); }) };
|
||||
}
|
||||
|
||||
async listRetainedSnapshots(): Promise<WorkspaceRevision[]> {
|
||||
async #listRetainedSnapshots(): Promise<WorkspaceRevision[]> {
|
||||
await registryContext(this).repository.ensureLayout();
|
||||
return await registryContext(this).lock.run(async () => {
|
||||
try {
|
||||
@@ -372,7 +458,7 @@ export class WorkspaceRegistry {
|
||||
});
|
||||
}
|
||||
|
||||
async read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> {
|
||||
async #read(id: string): Promise<{ workspace: WorkspaceDescriptor; revision: WorkspaceRevision }> {
|
||||
const state = await this.#activeState();
|
||||
const revision = state.revisions.find((candidate) => candidate.id === id);
|
||||
if (!revision) throw new WorkspaceRegistryError("workspace_invalid", "Workspace is unavailable");
|
||||
@@ -388,7 +474,7 @@ export class WorkspaceRegistry {
|
||||
* Resolve the active revision and create its cross-process retention lease under the same
|
||||
* repository lock. The lease bridges the interval before `session_manifest.yaml` is durable.
|
||||
*/
|
||||
async acquireSessionRevision(id: string): Promise<SessionRevisionLease> {
|
||||
async #acquireSessionRevision(id: string): Promise<SessionRevisionLease> {
|
||||
await registryContext(this).repository.ensureLayout();
|
||||
return await registryContext(this).lock.run(async () => {
|
||||
const state = await this.#activeState();
|
||||
@@ -437,7 +523,7 @@ export class WorkspaceRegistry {
|
||||
}
|
||||
|
||||
/** Read a retained immutable snapshot for a session pinned to a historical commit. */
|
||||
async readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }> {
|
||||
async #readPinned(id: string, commit: string): Promise<{ workspace: WorkspaceDescriptor; workspaceConfigPath: string }> {
|
||||
const snapshotPath = workspaceRegistrySnapshotPath(this, safeCommit(commit), id);
|
||||
try {
|
||||
const source = await readFile(snapshotPath, "utf8");
|
||||
@@ -567,7 +653,7 @@ export class WorkspaceRegistry {
|
||||
const base = await this.#readSnapshotCanonical(request.baseCommit, id);
|
||||
let remote: CanonicalWorkspace | undefined;
|
||||
if (existing) {
|
||||
remote = (await this.read(id)).workspace;
|
||||
remote = (await this.#read(id)).workspace;
|
||||
}
|
||||
return new WorkspaceConflictError(
|
||||
this.#changedFields(base, remote),
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { mkdtemp } from "node:fs/promises";
|
||||
import { join } from "node:path";
|
||||
import { RegistryAddressedPublicationStore } from "../src/workspaces/registry-publication.js";
|
||||
|
||||
const request = (runId: string, mode: "create" | "resume" = "create") => ({
|
||||
mode, operation: "registry_bootstrap" as const, runId: runId as any,
|
||||
requestSha256: "1".repeat(64) as any, installationIdentitySha256: "2".repeat(64) as any,
|
||||
repositoryIdentitySha256: "3".repeat(64) as any, expectedBaseCommit: null,
|
||||
remoteRefIdentitySha256: "4".repeat(64) as any,
|
||||
});
|
||||
|
||||
describe("addressed registry invariants", () => {
|
||||
it("pins one immutable target OID through fetch and plan", async () => {
|
||||
const store = new RegistryAddressedPublicationStore(await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-reg-")));
|
||||
const runId = "a".repeat(32);
|
||||
await store.claim(request(runId), null);
|
||||
await store.transition(runId as any, "target_advertised", { advertisedTargetCommit: "1".repeat(40) as any, immutableTargetRef: `refs/thoth/addressed-runs/${runId}/target` });
|
||||
await expect(store.transition(runId as any, "target_fetched", { fetchedTargetCommit: "2".repeat(40) as any })).rejects.toThrow();
|
||||
});
|
||||
|
||||
it("requires resume to find an existing run and rejects create replay", async () => {
|
||||
const store = new RegistryAddressedPublicationStore(await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-reg-")));
|
||||
const runId = "b".repeat(32);
|
||||
await expect(store.claim(request(runId, "resume"), null)).rejects.toThrow();
|
||||
await store.claim(request(runId), null);
|
||||
await expect(store.claim(request(runId), null)).rejects.toThrow();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user