diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 888f7d6f..93b54ffe 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -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); } diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index 6d8e1937..3df1bfdb 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -206,14 +206,30 @@ export async function runUnderOrderedWorkspaceWriterLocks(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(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise) { return runUnderOrderedWorkspaceWriterLocks([rootLease], set => set.forWorkspace(rootLease.identity.workspaceId, x => action(x.writerCapability))); } diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index f985b74e..530c5237 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -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"; } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 1804dd08..6a80de4f 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -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 { + private async bootstrap(): Promise { 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 { + private async pull(): Promise { 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 { + private async list(): Promise { const active = await this.tryActiveState(); if (active) return active.revisions; await this.bootstrap(); return (await this.activeState()).revisions; } - async activate(commit: string): Promise { await this.materialize(commit); await this.publishMaterialized(commit); } - async publish(request: PublishWorkspaceRequest): Promise { + private async activate(commit: string): Promise { await this.materialize(commit); await this.publishMaterialized(commit); } + /** Test-only migration seam; production callers use publishAddressed. */ + private async publish(request: PublishWorkspaceRequest): Promise { return this.publishWorkspace(request); } + private async publishWorkspace(request: PublishWorkspaceRequest): Promise { await this.repository.ensureLayout(); - return this.lock.run(async () => { + return this.lock.run(() => this.publishWorkspaceLocked(request)); + } + private async publishWorkspaceLocked(request: PublishWorkspaceRequest): Promise { 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 { @@ -226,6 +229,10 @@ export class WorkspaceRegistry { async publishAddressed(request: RegistryAddressedRequestV1): Promise { 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 & { 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 & { readonly mutation: PublishWorkspaceRequest }): Promise { + 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 { let state = initial; if (state.phase === "terminal_durable") return store.readTerminalResult(state.runId, state.terminalResultSha256!); diff --git a/backend/test/routes-workspaces.test.ts b/backend/test/routes-workspaces.test.ts index 86abf229..b21af38d 100644 --- a/backend/test/routes-workspaces.test.ts +++ b/backend/test/routes-workspaces.test.ts @@ -116,7 +116,7 @@ const revision: WorkspaceRevision = { snapshotPath: "/registry/snapshots/psd-clinical.yaml", }; -type RegistryFake = Pick & { ensureBootstrapAddressed: ReturnType; publishAddressed: ReturnType }; +type RegistryFake = Pick & { ensureBootstrapAddressed: ReturnType; publishAddressed: ReturnType; snapshotPath: ReturnType }; function registryFake(overrides: Partial = {}): RegistryFake { return { @@ -128,9 +128,10 @@ function registryFake(overrides: Partial = {}): 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(); });