fix: route workspace CRUD through addressed publication

This commit is contained in:
2026-08-11 18:30:06 +02:00
parent 4a0ec61da3
commit 29a2e507a9
5 changed files with 120 additions and 35 deletions
+20 -12
View File
@@ -13,7 +13,8 @@ import {
type PublishWorkspaceRequest,
type WorkspaceRegistry,
} from "../workspaces/registry.js";
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1 } from "../workspaces/registry-publication.js";
import type { RegistryActiveSnapshotV1, RegistryAddressedResultV1, RegistryEnsureBootstrapAddressedResultV1, RegistryWorkspaceMutationV1 } 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";
import {
@@ -292,11 +293,11 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
const snapshotFromPublication = (value: RegistryAddressedResultV1): RegistryActiveSnapshotV1 => ({
schemaVersion: 1, commit: value.plan.targetCommit, manifestSha256: value.plan.targetManifestSha256, workspaces: value.plan.targetWorkspaces,
});
const recoveryIdentity = () => {
if (typeof deps.registry.recoveryIdentity === "function") return deps.registry.recoveryIdentity();
const digest = sha256(`${deps.config.installationId}:${deps.config.remoteUrl ?? ""}:${deps.config.branch}`);
return { operation: "registry_bootstrap" as const, requestSha256: digest as never, installationIdentitySha256: sha256(deps.config.installationId) as never, repositoryIdentitySha256: digest as never, remoteRefIdentitySha256: digest as never };
const revisionFromPublication = (value: RegistryAddressedResultV1, id: string) => {
const item = value.plan.targetWorkspaces.find(candidate => candidate.workspaceId === id);
return item ? { id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: deps.registry.snapshotPath(item.revision, item.workspaceId) } : undefined;
};
const recoveryIdentity = () => deps.registry.recoveryIdentity();
app.get("/workspace-registry/status", async (_request, reply) => {
try {
const snapshot = snapshotFromEnsure(await deps.registry.ensureBootstrapAddressed(recoveryIdentity()));
@@ -366,13 +367,20 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
app.post("/workspaces/publish", async (request, reply) => {
try {
const requestValue = publishRequest(request.body);
let revision = await deps.registry.publish(requestValue);
if (!revision && requestValue.action !== "delete") {
const id = requestValue.workspace.workspace.id;
revision = (await deps.registry.read(id)).revision;
if (!revision) revision = (await deps.registry.list()).find(item => item.id === id);
}
return { revision };
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,
installationIdentitySha256: identity.installationIdentitySha256,
repositoryIdentitySha256: identity.repositoryIdentitySha256,
remoteRefIdentitySha256: identity.remoteRefIdentitySha256,
expectedBaseCommit: requestValue.baseCommit as Revision40,
mutation,
};
const result = await deps.registry.publishAddressed(addressed);
return { revision: revisionFromPublication(result, id) };
} catch (error) {
return errorReply(reply, error);
}
+22 -6
View File
@@ -206,14 +206,30 @@ export async function runUnderOrderedWorkspaceWriterLocks<T>(rootLeases: readonl
const sorted = [...rootLeases].sort((a, b) => a.identity.workspaceId.localeCompare(b.identity.workspaceId)); if (new Set(sorted.map(x => x.identity.workspaceId)).size !== sorted.length) throw fail();
const caps: WorkspaceWriterLockCapability[] = []; let set: OrderedWorkspaceWriterCapabilitySet | undefined; let result: T | undefined; let callbackError: unknown;
try {
for (const source of sorted) { const root = source.transfer(); let writer: WorkspaceRootLock | undefined; try { writer = await root.acquireWriterLock(); caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer)); } catch (error) { try { writer?.close(); } catch {} try { await root.close(); } catch {} throw error; } }
set = OrderedWorkspaceWriterCapabilitySet[INTERNAL_STATE](new Map(caps.map(c => [c.workspaceId, c])));
for (const capability of caps) capability.assertWriterPath();
try { result = await action(set); for (const capability of caps) capability.assertWriterPath(); }
catch (error) { callbackError = error; }
for (const source of sorted) {
const root = source.transfer();
let writer: WorkspaceRootLock | undefined;
try {
writer = await root.acquireWriterLock();
caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer));
} catch (error) {
try { writer?.close(); } catch {}
try { await root.close(); } catch {}
callbackError = error;
break;
}
}
if (!callbackError) {
set = OrderedWorkspaceWriterCapabilitySet[INTERNAL_STATE](new Map(caps.map(c => [c.workspaceId, c])));
for (const capability of caps) capability.assertWriterPath();
try { result = await action(set); for (const capability of caps) capability.assertWriterPath(); }
catch (error) { callbackError = error; }
}
} catch (error) { callbackError ??= error; }
set?.invalidate();
for (const cap of [...caps].reverse()) await cap.close().catch(() => undefined);
for (const cap of [...caps].reverse()) {
try { await cap.close(); } catch (error) { callbackError ??= error; }
}
if (callbackError) throw callbackError; return result as T;
}
export function runUnderWorkspaceWriterLock<T>(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise<T>) { return runUnderOrderedWorkspaceWriterLocks([rootLease], set => set.forWorkspace(rootLease.identity.workspaceId, x => action(x.writerCapability))); }
@@ -3,6 +3,7 @@ import { constants as fsConstants } from "node:fs";
import { lstat, mkdir, open, readFile, readdir, rename, rm, stat, unlink, writeFile } from "node:fs/promises";
import { dirname, join } from "node:path";
import type { CanonicalWorkspaceId, Revision40, Sha256Hex } from "./workspace-lock-root-lease.js";
import type { CanonicalWorkspace } from "./schema.js";
import type { BorrowedOrderedWorkspaceWriterLeaseV1, BorrowedWorkspaceSessionReadersExclusiveLockLease, OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js";
export type RegistryRunId32 = string & { readonly __registryRunId32: unique symbol };
@@ -68,10 +69,15 @@ 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 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: "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"; }
+50 -8
View File
@@ -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 RegistryAddressedRequestV1, type RegistryAddressedResultV1, type RegistryBootstrapRecoveryIdentityV1, type RegistryEnsureBootstrapAddressedResultV1, type RegistryActiveSnapshotV1, type RegistryWorkspaceManifestIdentityV1, type RegistryAddressedPublicationStateV1 } from "./registry-publication.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";
export type { GitStatus } from "./git-repository.js";
export interface WorkspaceRevision {
@@ -157,24 +157,28 @@ export class WorkspaceRegistry {
/** Automatic addressed recovery. The repository lock is held for selection and execution. */
recoveryIdentity(): RegistryBootstrapRecoveryIdentityV1 { return this.installationIdentity; }
async bootstrap(): Promise<GitStatus> {
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); } });
}
async pull(): Promise<GitStatus> {
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); } });
}
async list(): Promise<WorkspaceRevision[]> {
private async list(): Promise<WorkspaceRevision[]> {
const active = await this.tryActiveState();
if (active) return active.revisions;
await this.bootstrap();
return (await this.activeState()).revisions;
}
async activate(commit: string): Promise<void> { await this.materialize(commit); await this.publishMaterialized(commit); }
async publish(request: PublishWorkspaceRequest): Promise<WorkspaceRevision | undefined> {
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(async () => {
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;
@@ -190,7 +194,6 @@ export class WorkspaceRegistry {
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> {
@@ -226,6 +229,10 @@ export class WorkspaceRegistry {
async publishAddressed(request: RegistryAddressedRequestV1): 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); }
@@ -248,6 +255,41 @@ export class WorkspaceRegistry {
});
}
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 executeAddressedLocked(store: RegistryAddressedPublicationStore, initial: RegistryAddressedPublicationStateV1, identity: RegistryBootstrapRecoveryIdentityV1): Promise<RegistryAddressedResultV1> {
let state = initial;
if (state.phase === "terminal_durable") return store.readTerminalResult(state.runId, state.terminalResultSha256!);
+21 -8
View File
@@ -116,7 +116,7 @@ const revision: WorkspaceRevision = {
snapshotPath: "/registry/snapshots/psd-clinical.yaml",
};
type RegistryFake = Pick<WorkspaceRegistry, "bootstrap" | "pull" | "list" | "read" | "publish"> & { ensureBootstrapAddressed: ReturnType<typeof vi.fn>; publishAddressed: ReturnType<typeof vi.fn> };
type RegistryFake = Pick<WorkspaceRegistry, "read" | "recoveryIdentity"> & { ensureBootstrapAddressed: ReturnType<typeof vi.fn>; publishAddressed: ReturnType<typeof vi.fn>; snapshotPath: ReturnType<typeof vi.fn> };
function registryFake(overrides: Partial<RegistryFake> = {}): RegistryFake {
return {
@@ -128,9 +128,10 @@ function registryFake(overrides: Partial<RegistryFake> = {}): RegistryFake {
})),
list: vi.fn(async () => [revision]),
read: vi.fn(async () => ({ workspace, revision })),
publish: vi.fn(async () => revision),
recoveryIdentity: vi.fn(() => ({ operation: "registry_bootstrap", requestSha256: "c".repeat(64), installationIdentitySha256: "d".repeat(64), repositoryIdentitySha256: "e".repeat(64), remoteRefIdentitySha256: "f".repeat(64) })),
ensureBootstrapAddressed: vi.fn(async () => ({ kind: "already_active", snapshot: { schemaVersion: 1, commit: revision.commit, manifestSha256: "a".repeat(64), workspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })),
publishAddressed: vi.fn(async () => ({ plan: { targetCommit: revision.commit, targetManifestSha256: "a".repeat(64), targetWorkspaces: [{ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] } })),
snapshotPath: vi.fn(() => revision.snapshotPath),
...overrides,
};
}
@@ -278,7 +279,7 @@ test.each([
});
expect(response.body).not.toMatch(/migration_required|schema version/i);
}
expect(registry.publish).not.toHaveBeenCalled();
expect(registry.publishAddressed).not.toHaveBeenCalled();
});
test("runs diagnostics for a schema v3 workspace without external semantic bindings", async () => {
@@ -332,7 +333,7 @@ test("reports missing Evidence binding through the real test route without chang
variable,
})]));
expect(read).toHaveBeenCalledTimes(1);
expect(registry.publish).not.toHaveBeenCalled();
expect(registry.publishAddressed).not.toHaveBeenCalled();
expect(revision).toMatchObject({ commit: "a".repeat(40), blob: "b".repeat(40) });
} finally {
if (previous === undefined) delete process.env[variable];
@@ -352,7 +353,7 @@ test("returns a 409 field conflict instead of overwriting a changed workspace",
remote: { ...workspace, workspace: { ...workspace.workspace, description: "Remote description" } },
},
);
const registry = registryFake({ publish: vi.fn(async () => { throw conflict; }) });
const registry = registryFake({ publishAddressed: vi.fn(async () => { throw conflict; }) });
const app = appFor(registry);
const staleUpdate = {
action: "update",
@@ -378,7 +379,7 @@ test("returns a 409 field conflict instead of overwriting a changed workspace",
test("maps a stale registry commit to HTTP 409 without conflict payloads", async () => {
const registry = registryFake({
publish: vi.fn(async () => {
publishAddressed: vi.fn(async () => {
throw new WorkspaceRegistryError("workspace_stale", "Workspace revision is stale");
}),
});
@@ -395,6 +396,18 @@ test("maps a stale registry commit to HTTP 409 without conflict payloads", async
expect(res.json()).toEqual({ code: "workspace_stale", message: "Workspace revision is stale." });
});
test("publishes accepted CRUD through the addressed request and returns its terminal revision", async () => {
const publishAddressed = vi.fn(async (request: any) => ({
plan: { targetCommit: revision.commit, targetManifestSha256: "a".repeat(64), targetWorkspaces: [{ workspaceId: request.mutation.workspace.workspace.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: "b".repeat(64) }] },
}));
const registry = registryFake({ publishAddressed });
const app = appFor(registry);
const response = await app.inject({ method: "POST", url: "/workspaces/publish", payload: { action: "create", workspace, baseCommit: revision.commit } });
expect(response.statusCode).toBe(200);
expect(response.json()).toEqual({ revision: { id: workspace.workspace.id, commit: revision.commit, blob: revision.blob, snapshotPath: revision.snapshotPath } });
expect(publishAddressed).toHaveBeenCalledWith(expect.objectContaining({ mode: "create", operation: "registry_pull", mutation: expect.objectContaining({ action: "create", baseCommit: revision.commit }) }));
});
test("exports generated public artifacts without secret values", async () => {
const app = appFor(registryFake());
@@ -415,7 +428,7 @@ test("rejects a zip-slip import without publishing or writing a checkout file",
expect(res.statusCode).toBe(400);
expect(res.json()).toMatchObject({ code: "workspace_invalid" });
expect(registry.publish).not.toHaveBeenCalled();
expect(registry.publishAddressed).not.toHaveBeenCalled();
});
test("imports an exact generated bundle only as a browser draft", async () => {
@@ -426,7 +439,7 @@ test("imports an exact generated bundle only as a browser draft", async () => {
expect(res.statusCode).toBe(200);
expect(res.json()).toMatchObject({ draft: { workspace } });
expect(registry.publish).not.toHaveBeenCalled();
expect(registry.publishAddressed).not.toHaveBeenCalled();
});