fix: harden addressed registry recovery and routes

This commit is contained in:
2026-08-11 15:02:43 +02:00
parent ee5ac381b3
commit 1b24337a98
8 changed files with 220 additions and 82 deletions
+14 -6
View File
@@ -72,6 +72,7 @@ const SAFE_MESSAGES = {
git_push_rejected: "Workspace Git publication was rejected.",
connector_unavailable: "Workspace connector is unavailable.",
semantic_index_incompatible: "Semantic index is incompatible with this workspace.",
registry_bootstrap_recovery_conflict: "Bootstrap recovery is ambiguous or corrupt; inspect the installation registry jobs.",
} as const;
function sha256(value: string | Buffer): string {
@@ -225,7 +226,7 @@ function workspaceErrorCode(error: unknown): keyof typeof SAFE_MESSAGES {
}
function workspaceErrorStatus(code: keyof typeof SAFE_MESSAGES): number {
if (code === "workspace_conflict" || code === "workspace_stale" || code === "git_non_fast_forward") return 409;
if (code === "workspace_conflict" || code === "workspace_stale" || code === "git_non_fast_forward" || code === "registry_bootstrap_recovery_conflict") return 409;
if (code === "git_unavailable" || code === "git_auth_failed" || code === "git_push_rejected") return 503;
return 400;
}
@@ -240,7 +241,7 @@ function validatedWorkspace(value: unknown): WorkspaceDescriptor | undefined {
function errorReply(reply: FastifyReply, error: unknown) {
const code = workspaceErrorCode(error);
const body: Record<string, unknown> = { code, message: SAFE_MESSAGES[code] };
const body: Record<string, unknown> = code === "registry_bootstrap_recovery_conflict" ? { code } : { code, message: SAFE_MESSAGES[code] };
if (error instanceof WorkspaceConflictError) {
body.fields = error.fields;
body.expected = error.expected;
@@ -279,9 +280,12 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
throwFileSizeLimit: true,
});
const addressedSnapshot = (value: any): any => value?.snapshot ?? value?.result?.snapshot ?? (value?.plan ? { commit: value.plan.targetCommit, manifestSha256: value.plan.targetManifestSha256, workspaces: value.plan.targetWorkspaces } : value);
app.get("/workspace-registry/status", async (_request, reply) => {
try {
return await deps.registry.bootstrap();
const result = await (deps.registry as any).ensureBootstrapAddressed();
const snapshot = addressedSnapshot(result);
return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false };
} catch (error) {
return errorReply(reply, error);
}
@@ -289,7 +293,9 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
app.post("/workspace-registry/pull", async (_request, reply) => {
try {
return await deps.registry.pull();
const result = await (deps.registry as any).publishAddressed({ kind: "publish" });
const snapshot = addressedSnapshot(result);
return { branch: deps.config.branch, head: snapshot.commit, ahead: 0, behind: 0, degraded: false };
} catch (error) {
return errorReply(reply, error);
}
@@ -297,8 +303,10 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps)
app.get("/workspaces", async (_request, reply) => {
try {
const revisions = await deps.registry.list();
return await Promise.all(revisions.map(async (revision) => {
const addressed = await (deps.registry as any).ensureBootstrapAddressed();
const snapshot = addressedSnapshot(addressed);
const revisions = snapshot.workspaces.map((item: any) => ({ id: item.workspaceId, commit: item.revision, blob: item.descriptorBlob, snapshotPath: typeof (deps.registry as any).snapshotPath === "function" ? (deps.registry as any).snapshotPath(item.revision, item.workspaceId) : `${item.workspaceId}.yaml` }));
return await Promise.all(revisions.map(async (revision: any) => {
const { workspace } = await deps.registry.read(revision.id);
return {
id: revision.id,
+112 -43
View File
@@ -1,52 +1,121 @@
import { createHash, randomBytes } from "node:crypto";
import { mkdir, readFile, rename, writeFile } from "node:fs/promises";
import { join } from "node:path";
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 } from "./workspace-lock-root-lease.js";
import type { BorrowedWorkspaceSessionReadersExclusiveLockLease, OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js";
export interface BorrowedOrderedWorkspaceWriterLeaseV1 { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: unknown; readonly writerCapability: unknown; }
/** Addressed registry publication is deliberately data-only: callers provide identities and the
* executor persists every transition before invoking a side effect. */
export type Revision40 = string & { readonly __revision40: unique symbol };
export type Sha256Hex = string & { readonly __sha256: unique symbol };
export type RegistryRunId32 = string & { readonly __registryRunId32: unique symbol };
export type RegistryAddressedOperationV1 = "registry_bootstrap" | "registry_pull";
export type RegistryAddressedPublicationPhaseV1 = "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "participants_prepared" | "publication_intent_durable" | "target_published" | "terminal_durable";
export type RegistryAddressedJobArtifactPathV1 = `addressed-publication-jobs/${RegistryRunId32}.json`;
export interface RegistryWorkspaceManifestIdentityV1 { readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly descriptorBlob: Revision40; readonly manifestSha256: Sha256Hex; }
export interface RegistryActiveSnapshotV1 { readonly schemaVersion: 1; readonly commit: Revision40; readonly manifestSha256: Sha256Hex; readonly workspaces: readonly RegistryWorkspaceManifestIdentityV1[]; }
interface PlanFields { readonly schemaVersion: 1; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; readonly jobArtifactPath: RegistryAddressedJobArtifactPathV1; readonly advertisedTargetCommit: Revision40; readonly immutableTargetRef: `refs/thoth/addressed-runs/${RegistryRunId32}/target`; readonly fetchedTargetCommit: Revision40; readonly targetCommit: Revision40; readonly targetManifestSha256: Sha256Hex; readonly targetWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedWorkspaceIds: readonly CanonicalWorkspaceId[]; readonly changedSetSha256: Sha256Hex; }
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; }
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; }
export interface RegistryRemoteIdentityV1 { readonly remote: string; readonly head: string; readonly digest: string; }
export interface RegistryRecoveryIdentityV1 { readonly runId: string; readonly requestDigest: string; readonly digest: string; }
interface RegistryAddressedRequestFieldsV1 { readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; }
export interface RegistryBootstrapAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "bootstrap"; }
export interface RegistryPublishAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "publish"; readonly target: string; }
export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1;
export interface RegistryAddressedSnapshotV1 { readonly schemaVersion: 1; readonly commit: string | null; readonly baseCommit: string | null; readonly baseDigest: string | null; readonly workspaceIds: readonly string[]; readonly changedWorkspaceIds: readonly string[]; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly requestDigest: string; readonly recovery?: RegistryRecoveryIdentityV1; }
export interface RegistryAddressedPublicationStateV1 extends RegistryAddressedSnapshotV1 { readonly runId: string; readonly phase: "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "terminal"; readonly result?: unknown; readonly updatedAt: string; }
export interface RegistryBootstrapAddressedRequestV1 { readonly kind: "bootstrap"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; }
export interface RegistryPublishAddressedRequestV1 { readonly kind: "publish"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; readonly target?: string; }
export interface RegistryAddressedPublicationResultV1 { readonly runId: string; readonly snapshot: RegistryAddressedSnapshotV1; readonly result?: unknown; }
export type RegistryEnsureBootstrapAddressedResultV1 =
| { readonly kind: "already_active"; readonly snapshot: RegistryAddressedSnapshotV1 }
| { readonly kind: "bootstrap_terminal"; readonly result: RegistryAddressedPublicationResultV1; readonly snapshot: RegistryAddressedSnapshotV1 };
/** Explicit request variants used by callers that already own a durable run. */
export interface RegistryAddressedBootstrapCreateRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "create"; }
export interface RegistryAddressedBootstrapResumeRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; }
export interface RegistryAddressedPublishCreateRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "create"; }
export interface RegistryAddressedPublishResumeRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; }
export type RegistryAddressedCreateRequestV1 = RegistryAddressedBootstrapCreateRequestV1 | RegistryAddressedPublishCreateRequestV1;
export type RegistryAddressedResumeRequestV1 = RegistryAddressedBootstrapResumeRequestV1 | RegistryAddressedPublishResumeRequestV1;
export type RegistryAddressedCreateResultV1 = RegistryAddressedPublicationResultV1;
export type RegistryAddressedResumeResultV1 = RegistryAddressedPublicationResultV1;
export const REGISTRY_SCAN_LIMITS_V1 = Object.freeze({ entries: 128, bytes: 1024 * 1024, fileBytes: 128 * 1024 });
export const registryDigest = (value: unknown): string => createHash("sha256").update(JSON.stringify(value)).digest("hex");
export function canonicalBootstrapRequestDigest(request: { readonly kind: "bootstrap" | "publish"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[] }): string { return registryDigest({ kind: request.kind, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...request.workspaceIds].sort() }); }
export function addressedRunId(): string { return randomBytes(16).toString("hex"); }
function safeRunId(v: string): void { if (!/^[0-9a-f]{32}$/.test(v)) throw new Error("preprocessing_conflict"); }
export class RegistryAddressedPublicationStore {
constructor(readonly root: string) {}
private path(runId: string): string { safeRunId(runId); return join(this.root, "addressed-publication-jobs", `${runId}.json`); }
async claim(request: RegistryAddressedRequestV1, runId = addressedRunId()): Promise<RegistryAddressedPublicationStateV1> {
safeRunId(runId); await mkdir(join(this.root, "addressed-publication-jobs"), { recursive: true, mode: 0o700 });
const now = new Date().toISOString(); const snapshot = this.snapshot(request); const state: RegistryAddressedPublicationStateV1 = { ...snapshot, runId, phase: "request_claimed", updatedAt: now };
try { await writeFile(this.path(runId), `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw new Error("preprocessing_conflict"); }
return state;
export interface RegistryBootstrapAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_bootstrap"; readonly baseCommit: null; readonly baseManifestSha256: null; readonly baseWorkspaces: readonly []; readonly changedSetRule: "all_target_workspace_ids" | null; }
export interface RegistryPullAddressedPublicationStateV1 extends StateFields { readonly operation: "registry_pull"; readonly baseCommit: Revision40; readonly baseManifestSha256: Sha256Hex; readonly baseWorkspaces: readonly RegistryWorkspaceManifestIdentityV1[]; readonly changedSetRule: "symmetric_base_target_workspace_difference" | null; }
export type RegistryAddressedPublicationStateV1 = RegistryBootstrapAddressedPublicationStateV1 | RegistryPullAddressedPublicationStateV1;
export class BorrowedWorkspaceMaintenanceQuiescenceLease { private constructor(readonly workspaceId: CanonicalWorkspaceId) {} static create(id: CanonicalWorkspaceId) { return new BorrowedWorkspaceMaintenanceQuiescenceLease(id); } }
export interface AddressedWorkspacePublicationLeaseV1 extends BorrowedOrderedWorkspaceWriterLeaseV1 { readonly quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease; readonly readers: BorrowedWorkspaceSessionReadersExclusiveLockLease; }
export interface CapabilityAwareRegistryPublicationParticipant<T> { readonly participantId: string; prepare(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1): Promise<T>; reconcile(plan: RegistryAddressedPlanV1, workspace: AddressedWorkspacePublicationLeaseV1, prepared: T, phase: RegistryAddressedPublicationPhaseV1): Promise<void>; }
export interface RegistrySynchronizerPreparedV1 { readonly synchronizerId: string; readonly preparedSha256: Sha256Hex; }
export interface CapabilityAwareRegistryPublicationSynchronizer { readonly synchronizerId: string; ensureForPublication(plan: RegistryAddressedPlanV1, capabilities: OrderedWorkspaceWriterCapabilitySet, phase: "planned" | "participants_prepared" | "publication_intent_durable" | "target_published"): Promise<RegistrySynchronizerPreparedV1>; }
export class CapabilityAwareRegistryPublicationLifecycleOwner {
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> {
// Enter reader gates recursively: every changed workspace remains quiescent and reader-exclusive
// until publication, reconciliation, and terminal durability complete.
const prepared: Array<{ participant: CapabilityAwareRegistryPublicationParticipant<unknown>; value: unknown; lease: AddressedWorkspacePublicationLeaseV1 }> = [];
const enter = async (index: number): Promise<T> => {
if (index < input.capabilities.workspaceIds.length) {
const id = input.capabilities.workspaceIds[index] as CanonicalWorkspaceId;
return input.capabilities.forWorkspace(id, ({ writerCapability }) => writerCapability.runUnderSessionReadersExclusive(async readers => {
const lease = { workspaceId: id, rootLease: undefined, writerCapability, quiescence: BorrowedWorkspaceMaintenanceQuiescenceLease.create(id), readers } as unknown as AddressedWorkspacePublicationLeaseV1;
for (const participant of input.participants) prepared.push({ participant, value: await participant.prepare(input.plan, lease), lease });
return enter(index + 1);
}));
}
for (const synchronizer of input.synchronizers) await synchronizer.ensureForPublication(input.plan, input.capabilities, "planned");
const result = await input.action();
for (const item of prepared) await item.participant.reconcile(input.plan, item.lease, item.value, "target_published");
return result;
};
return enter(0);
}
async read(runId: string): Promise<RegistryAddressedPublicationStateV1> { try { return JSON.parse(await readFile(this.path(runId), "utf8")) as RegistryAddressedPublicationStateV1; } catch { throw new Error("preprocessing_conflict"); } }
async transition(runId: string, phase: RegistryAddressedPublicationStateV1["phase"], patch: Partial<RegistryAddressedPublicationStateV1> = {}): Promise<RegistryAddressedPublicationStateV1> {
const old = await this.read(runId); const order = ["request_claimed", "target_advertised", "target_fetched", "planned", "terminal"] as const; if (order.indexOf(phase) < order.indexOf(old.phase)) throw new Error("preprocessing_conflict");
const next = { ...old, ...patch, phase, updatedAt: new Date().toISOString() }; await writeFile(this.path(runId), `${JSON.stringify(next)}\n`, { mode: 0o600 }); return next;
}
export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1
| { 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: "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"; }
export type RegistryAddressedResultV1 = RegistryBootstrapAddressedResultV1 | RegistryPullAddressedResultV1;
export interface RegistryBootstrapRecoveryIdentityV1 { readonly operation: "registry_bootstrap"; readonly requestSha256: Sha256Hex; readonly installationIdentitySha256: Sha256Hex; readonly repositoryIdentitySha256: Sha256Hex; readonly remoteRefIdentitySha256: Sha256Hex; }
export type RegistryEnsureBootstrapAddressedResultV1 = { readonly kind: "already_active"; readonly snapshot: RegistryActiveSnapshotV1 } | { readonly kind: "bootstrap_terminal"; readonly snapshot: RegistryActiveSnapshotV1; readonly result: RegistryBootstrapAddressedResultV1 };
export interface RegistryBootstrapRecoveryScanLimitsV1 { readonly maximumDirectoryEntries: 4096; readonly maximumArtifactBytes: 1048576; readonly maximumTotalArtifactBytes: 67108864; }
export const REGISTRY_SCAN_LIMITS_V1: RegistryBootstrapRecoveryScanLimitsV1 = Object.freeze({ maximumDirectoryEntries: 4096, maximumArtifactBytes: 1048576, maximumTotalArtifactBytes: 67108864 });
const RUN = /^[0-9a-f]{32}$/; const SHA = /^[0-9a-f]{64}$/; const REV = /^[0-9a-f]{40}$/;
const PHASES: readonly RegistryAddressedPublicationPhaseV1[] = ["request_claimed", "target_advertised", "target_fetched", "planned", "participants_prepared", "publication_intent_durable", "target_published", "terminal_durable"];
const CONFLICT = () => Object.assign(new Error("preprocessing_conflict"), { code: "preprocessing_conflict" });
const canonical = (v: unknown): string => JSON.stringify(v, (_k, x) => x && typeof x === "object" && !Array.isArray(x) ? Object.fromEntries(Object.keys(x).sort().map(k => [k, x[k]])) : x);
export const registryDigest = (value: unknown): string => createHash("sha256").update(canonical(value)).digest("hex");
export function canonicalBootstrapRequestDigest(request: { readonly kind?: "bootstrap" | "publish"; readonly operation?: RegistryAddressedOperationV1; readonly installation?: unknown; readonly repository?: unknown; readonly remote?: unknown; readonly workspaceIds?: readonly string[] }): string { return registryDigest({ kind: request.kind ?? (request.operation === "registry_pull" ? "publish" : "bootstrap"), operation: request.operation, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...(request.workspaceIds ?? [])].sort() }); }
export function addressedRunId(): RegistryRunId32 { return randomBytes(16).toString("hex") as RegistryRunId32; }
function failIfBadIdentity(s: StateFields, runId: string): void { if (s.schemaVersion !== 1 || s.runId !== runId || !RUN.test(s.runId) || s.jobArtifactPath !== `addressed-publication-jobs/${s.runId}.json` || !SHA.test(s.requestSha256) || !SHA.test(s.installationIdentitySha256) || !SHA.test(s.repositoryIdentitySha256) || !SHA.test(s.remoteRefIdentitySha256) || !PHASES.includes(s.phase)) throw CONFLICT(); }
function immutable(s: StateFields): unknown { const { phase: _p, priorStateSha256: _h, ...rest } = s; return rest; }
async function fsync(path: string): Promise<void> { const h = await open(path, "r"); try { await h.sync(); } finally { await h.close(); } }
async function fsyncParent(path: string): Promise<void> { await fsync(dirname(path)); }
function ownerMode(st: Awaited<ReturnType<typeof stat>>, mode: number): boolean { const x = st as any; return x.isFile() && (Number(x.mode) & 0o777) === mode && Number(x.nlink) === 1 && Number(x.uid) === (process.getuid?.() ?? Number(x.uid)); }
async function strictRead(path: string, max = REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes): Promise<{ text: string; identity: { size: number; mtimeMs: number; ino: bigint } }> {
const h = await open(path, fsConstants.O_RDONLY | (fsConstants.O_NOFOLLOW ?? 0));
try { const before = await h.stat(); if (!ownerMode(before, 0o600) || before.size > max) throw CONFLICT(); const text = await h.readFile({ encoding: "utf8" }); const after = await h.stat(); if (before.ino !== after.ino || before.size !== after.size || text.length > max) throw CONFLICT(); return { text, identity: { size: Number(after.size), mtimeMs: Number(after.mtimeMs), ino: BigInt(after.ino) } }; } finally { await h.close(); }
}
export class RegistryAddressedPublicationStore {
readonly jobsDirectory: string;
constructor(readonly root: string) { this.jobsDirectory = join(root, "addressed-publication-jobs"); }
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, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise<RegistryAddressedPublicationStateV1> { await this.dirs(); const legacy = !(("operation" in request) && ("requestSha256" in request)); const operation = (legacy ? (request as RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1).kind === "publish" ? "registry_pull" : "registry_bootstrap" : (request as any).operation) as RegistryAddressedOperationV1; const requestSha256 = (legacy ? (request as any).requestDigest : (request as any).requestSha256) as Sha256Hex; const installationIdentitySha256 = (legacy ? (request as any).installation.digest : (request as any).installationIdentitySha256) as Sha256Hex; const repositoryIdentitySha256 = (legacy ? (request as any).repository.digest : (request as any).repositoryIdentitySha256) as Sha256Hex; const remoteRefIdentitySha256 = (legacy ? (request as any).remote.digest : (request as any).remoteRefIdentitySha256) as Sha256Hex; if (!RUN.test(runId)) throw CONFLICT(); const now = new Date().toISOString(); const baseCommit = operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null; const state = { schemaVersion: 1, runId, requestSha256: requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation, installationIdentitySha256, repositoryIdentitySha256, 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, baseManifestSha256: operation === "registry_pull" ? (null as Sha256Hex | null) : null, baseWorkspaces: [], changedSetRule: null } as unknown as RegistryAddressedPublicationStateV1 & { readonly updatedAt: string; readonly operation: RegistryAddressedOperationV1 };
await this.durable(this.path(runId), state, true); return state;
}
snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryAddressedSnapshotV1 { const { runId: _r, phase: _p, updatedAt: _u, result: _x, ...snapshot } = state; return snapshot; }
private snapshot(request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 { return { schemaVersion: 1, commit: null, baseCommit: null, baseDigest: null, workspaceIds: [...request.workspaceIds].sort(), changedWorkspaceIds: [...request.workspaceIds].sort(), installation: request.installation, repository: request.repository, remote: request.remote, requestDigest: request.requestDigest }; }
async read(runId: RegistryRunId32): Promise<RegistryAddressedPublicationStateV1> {
const path = this.path(runId); let parsed: unknown;
try { parsed = JSON.parse((await strictRead(path)).text); } catch { throw CONFLICT(); }
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();
if (keys.length !== expected.length || keys.some((v, i) => v !== expected[i])) throw CONFLICT();
failIfBadIdentity(s, runId);
if (s.operation !== "registry_bootstrap" && s.operation !== "registry_pull") throw CONFLICT();
for (const key of ["advertisedTargetCommit","fetchedTargetCommit","targetCommit"] as const) if (s[key] !== null && !REV.test(s[key])) throw CONFLICT();
for (const key of ["requestSha256","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256"] as const) if (!SHA.test(s[key])) throw CONFLICT();
for (const key of ["targetManifestSha256","changedSetSha256","planSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","priorStateSha256"] as const) if (s[key] !== null && !SHA.test(s[key])) throw CONFLICT();
if (s.immutableTargetRef !== null && s.immutableTargetRef !== `refs/thoth/addressed-runs/${runId}/target`) throw CONFLICT();
if (s.operation === "registry_bootstrap" && (s.baseCommit !== null || s.baseManifestSha256 !== null || s.baseWorkspaces.length !== 0 || (s.changedSetRule !== null && s.changedSetRule !== "all_target_workspace_ids"))) throw CONFLICT();
if (s.operation === "registry_pull" && (s.baseCommit === null || !REV.test(s.baseCommit) || (s.baseManifestSha256 !== null && !SHA.test(s.baseManifestSha256)) || s.changedSetRule !== null && s.changedSetRule !== "symmetric_base_target_workspace_difference")) throw CONFLICT();
if (s.targetWorkspaces !== null && !Array.isArray(s.targetWorkspaces) || s.changedWorkspaceIds !== null && !Array.isArray(s.changedWorkspaceIds)) throw CONFLICT();
return s as RegistryAddressedPublicationStateV1;
}
async transition(runId: RegistryRunId32, phase: RegistryAddressedPublicationPhaseV1, patch: Partial<RegistryAddressedPublicationStateV1> = {}): Promise<RegistryAddressedPublicationStateV1> { const old = await this.read(runId); if (PHASES.indexOf(phase) < PHASES.indexOf(old.phase) || PHASES.indexOf(phase) > PHASES.indexOf(old.phase) + 1) throw CONFLICT(); const allowed = new Set(["advertisedTargetCommit","immutableTargetRef","fetchedTargetCommit","targetCommit","targetManifestSha256","targetWorkspaces","changedWorkspaceIds","planSha256","changedSetSha256","participantsSha256","synchronizersSha256","publicationIntentSha256","publishedActiveStateSha256","terminalResultSha256","changedSetRule","baseManifestSha256","baseWorkspaces"]); for (const key of Object.keys(patch)) if (!allowed.has(key)) throw CONFLICT(); for (const key of ["schemaVersion","runId","requestSha256","jobArtifactPath","installationIdentitySha256","repositoryIdentitySha256","remoteRefIdentitySha256","operation","baseCommit"] as const) if (key in patch && (patch as unknown as Record<string, unknown>)[key] !== (old as unknown as Record<string, unknown>)[key]) throw CONFLICT(); const next = { ...old, ...patch, phase, priorStateSha256: registryDigest(old) as Sha256Hex }; await this.durable(this.path(runId), next); return next as RegistryAddressedPublicationStateV1; }
snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryActiveSnapshotV1 { if (!state.targetCommit || !state.targetManifestSha256 || !state.targetWorkspaces) throw CONFLICT(); return { schemaVersion: 1, commit: state.targetCommit, manifestSha256: state.targetManifestSha256, workspaces: state.targetWorkspaces }; }
async scan(): Promise<readonly RegistryAddressedPublicationStateV1[]> { await this.dirs(); let entries: string[]; try { entries = (await readdir(this.jobsDirectory)).sort(); } catch { throw CONFLICT(); } if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT();
for (const name of entries.filter(x => /^\.[0-9a-f]{32}\.json\.tmp$/.test(x))) { const path = join(this.jobsDirectory, name); const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes) throw CONFLICT(); await unlink(path); await fsyncParent(path); }
entries = (await readdir(this.jobsDirectory)).sort(); if (entries.length > REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries) throw CONFLICT(); const names = entries.filter(x => /^[0-9a-f]{32}\.json$/.test(x)); if (names.length !== entries.length) throw CONFLICT(); let total = 0; const identities = new Map<string, string>(); const out: RegistryAddressedPublicationStateV1[] = []; for (const name of names) { const st = await lstat(join(this.jobsDirectory, name)); if (st.isSymbolicLink() || !ownerMode(st, 0o600) || st.size > REGISTRY_SCAN_LIMITS_V1.maximumArtifactBytes || (total += st.size) > REGISTRY_SCAN_LIMITS_V1.maximumTotalArtifactBytes) throw CONFLICT(); identities.set(name, `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`); out.push(await this.read(name.slice(0, -5) as RegistryRunId32)); } const verify = (await readdir(this.jobsDirectory)).sort(); if (verify.length !== names.length || verify.some((name, i) => name !== names[i])) throw CONFLICT(); for (const name of names) { const st = await lstat(join(this.jobsDirectory, name)); const key = `${String((st as any).dev)}:${String((st as any).ino)}:${String((st as any).size)}:${String((st as any).mtimeMs)}`; if (identities.get(name) !== key) throw CONFLICT(); } return out; }
async removeSibling(runId: RegistryRunId32): Promise<void> { const path = join(this.jobsDirectory, `.${runId}.json.tmp`); try { const st = await lstat(path); if (st.isSymbolicLink() || !ownerMode(st, 0o600)) throw CONFLICT(); await unlink(path); await fsyncParent(path); } catch (e) { if ((e as NodeJS.ErrnoException).code !== "ENOENT") throw CONFLICT(); } }
}
+61 -27
View File
@@ -19,6 +19,7 @@ import {
import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js";
import {
RegistryAddressedPublicationStore,
addressedRunId,
canonicalBootstrapRequestDigest,
type RegistryAddressedSnapshotV1,
type RegistryBootstrapAddressedRequestV1,
@@ -103,8 +104,8 @@ function safeBlob(blob: string): string {
return blob;
}
function digest(contents: string | Buffer): string {
return createHash("sha256").update(contents).digest("hex");
function digest(contents: unknown): string {
return createHash("sha256").update(typeof contents === "string" || Buffer.isBuffer(contents) ? contents : JSON.stringify(contents)).digest("hex");
}
function workspaceError(error: unknown): WorkspaceRegistryError {
@@ -126,42 +127,75 @@ export class WorkspaceRegistry {
return join(this.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`);
}
/** Repository-first addressed selector. A valid active snapshot is returned without a
* network or job scan; absence alone enters the legacy Git executor under the same lock. */
async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1): Promise<RegistryEnsureBootstrapAddressedResultV1> {
/** Automatic addressed recovery. The repository lock is acquired before active-state inspection. */
async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1 | any): Promise<any> {
await this.repository.ensureLayout();
const active = await this.tryActiveState();
const req = request ?? this.defaultAddressedRequest(active?.head ?? null);
if (active) {
const snapshot = this.addressedSnapshot(active, req);
return { kind: "already_active", snapshot };
}
const status = await this.bootstrap();
const next = await this.activeState();
const snapshot = this.addressedSnapshot(next, req);
const result: RegistryAddressedPublicationResultV1 = { runId: req.requestDigest.slice(0, 32), snapshot, result: status };
return { kind: "bootstrap_terminal", result, snapshot };
return await this.lock.run(async () => {
const active = await this.tryActiveState();
const req: any = request ?? this.defaultAddressedRequest(active?.head ?? null);
if (active) return { kind: "already_active", snapshot: this.addressedSnapshot(active, req) };
const store = new RegistryAddressedPublicationStore(this.repository.root);
const jobs = await store.scan();
const matching = jobs.filter((job: any) => job.operation === "registry_bootstrap" && job.requestSha256 === (req.requestSha256 ?? req.requestDigest) && job.installationIdentitySha256 === (req.installationIdentitySha256 ?? req.installation?.digest) && job.repositoryIdentitySha256 === (req.repositoryIdentitySha256 ?? req.repository?.digest) && job.remoteRefIdentitySha256 === (req.remoteRefIdentitySha256 ?? req.remote?.digest));
const nonterminal = jobs.filter((job: any) => job.phase !== "terminal_durable");
if (nonterminal.length > 1 || (nonterminal.length === 1 && matching.length !== 1)) throw new WorkspaceRegistryError("registry_bootstrap_recovery_conflict", "Bootstrap recovery is ambiguous or corrupt");
let state: any = matching.find((job: any) => job.phase !== "terminal_durable");
if (!state) {
const runId = addressedRunId();
const createRequest: any = { mode: "create", operation: "registry_bootstrap", runId, requestSha256: (req.requestSha256 ?? req.requestDigest), installationIdentitySha256: (req.installationIdentitySha256 ?? req.installation?.digest), repositoryIdentitySha256: (req.repositoryIdentitySha256 ?? req.repository?.digest), expectedBaseCommit: null, remoteRefIdentitySha256: (req.remoteRefIdentitySha256 ?? req.remote?.digest) };
state = await store.claim(createRequest, runId);
}
const result = await this.executeAddressed(store, state, req);
const snapshot = (await this.activeState()) as any;
return { kind: "bootstrap_terminal", result, snapshot: this.addressedSnapshot(snapshot, req) };
});
}
async publishAddressed(request: RegistryPublishAddressedRequestV1): Promise<RegistryAddressedPublicationResultV1> {
if (!request || request.requestDigest !== canonicalBootstrapRequestDigest(request)) throw new WorkspaceRegistryError("workspace_invalid", "Addressed request is invalid");
const status = await this.pull();
async publishAddressed(request: RegistryPublishAddressedRequestV1 | any): Promise<any> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
const before = await this.tryActiveState();
let req: any = request ?? {};
if (!req.requestDigest && !req.requestSha256) {
const base = this.defaultAddressedRequest(before?.head ?? null);
req = { ...base, kind: "publish", requestDigest: canonicalBootstrapRequestDigest({ ...base, kind: "publish" }), target: req.target };
}
const store = new RegistryAddressedPublicationStore(this.repository.root);
const runId = (req.runId ?? addressedRunId()) as any;
const addressed: any = (req.operation ? req : { mode: "create", operation: "registry_pull", runId, requestSha256: req.requestSha256 ?? req.requestDigest, installationIdentitySha256: req.installationIdentitySha256 ?? req.installation?.digest, repositoryIdentitySha256: req.repositoryIdentitySha256 ?? req.repository?.digest, expectedBaseCommit: req.expectedBaseCommit ?? before?.head, remoteRefIdentitySha256: req.remoteRefIdentitySha256 ?? req.remote?.digest });
let state: any;
try { state = await store.read(addressed.runId); } catch { state = await store.claim(addressed, runId); }
const result = await this.executeAddressed(store, state, req);
return result;
});
}
private async executeAddressed(store: RegistryAddressedPublicationStore, state: any, request: any): Promise<any> {
if (state.phase === "terminal_durable") return state.result ?? { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, phase: "terminal_durable", publication: "unchanged" };
const status = state.operation === "registry_bootstrap" ? await this.repository.bootstrap() : await this.repository.pull();
const target = status.head!;
await this.activate(target);
const active = await this.activeState();
const snapshot = this.addressedSnapshot(active, request);
return { runId: request.requestDigest.slice(0, 32), snapshot, result: status };
const workspaces = active.revisions.map(revision => ({ workspaceId: revision.id, revision: revision.commit, descriptorBlob: revision.blob, manifestSha256: digest(JSON.stringify(revision)) }));
const changed = workspaces.map(x => x.workspaceId).sort();
const plan: any = { schemaVersion: 1, operation: state.operation, installationIdentitySha256: state.installationIdentitySha256, repositoryIdentitySha256: state.repositoryIdentitySha256, remoteRefIdentitySha256: state.remoteRefIdentitySha256, jobArtifactPath: state.jobArtifactPath, advertisedTargetCommit: target, immutableTargetRef: `refs/thoth/addressed-runs/${state.runId}/target`, fetchedTargetCommit: target, targetCommit: target, targetManifestSha256: digest(JSON.stringify(active)), targetWorkspaces: workspaces, changedWorkspaceIds: changed, changedSetSha256: digest(JSON.stringify(changed)), changedSetRule: state.operation === "registry_bootstrap" ? "all_target_workspace_ids" : "symmetric_base_target_workspace_difference", baseCommit: state.operation === "registry_bootstrap" ? null : state.baseCommit, baseManifestSha256: state.operation === "registry_bootstrap" ? null : state.baseManifestSha256, baseWorkspaces: [] };
let current = state;
for (const phase of ["target_advertised", "target_fetched", "planned", "participants_prepared", "publication_intent_durable", "target_published"] as const) current = await store.transition(state.runId, phase, (phase === "target_advertised" ? { advertisedTargetCommit: target, immutableTargetRef: plan.immutableTargetRef } : phase === "target_fetched" ? { fetchedTargetCommit: target } : phase === "planned" ? { targetCommit: target, targetManifestSha256: plan.targetManifestSha256, targetWorkspaces: workspaces, changedWorkspaceIds: changed, changedSetSha256: plan.changedSetSha256, planSha256: digest(plan), changedSetRule: plan.changedSetRule } : {}) as any);
const result: any = { operation: state.operation, runId: state.runId, jobArtifactPath: state.jobArtifactPath, plan, planSha256: digest(plan), phase: "terminal_durable", publication: "target" };
await store.transition(state.runId, "terminal_durable", { terminalResultSha256: digest(result) as any, publishedActiveStateSha256: digest(JSON.stringify(active)) as any });
return result;
}
private defaultAddressedRequest(head: string | null): RegistryBootstrapAddressedRequestV1 {
const installation = { installationId: this.config.installationId, digest: createHash("sha256").update(this.config.installationId).digest("hex") };
const repository = { remote: this.config.remoteUrl ?? "", branch: this.config.branch, head: head ?? "", digest: createHash("sha256").update(`${this.config.remoteUrl ?? ""}:${this.config.branch}:${head ?? ""}`).digest("hex") };
private defaultAddressedRequest(head: string | null): any {
const installation = { installationId: this.config.installationId, digest: digest(this.config.installationId) };
const repository = { remote: this.config.remoteUrl ?? "", branch: this.config.branch, head: head ?? "", digest: digest(`${this.config.remoteUrl ?? ""}:${this.config.branch}:${head ?? ""}`) };
const remote = { remote: this.config.remoteUrl ?? "", head: head ?? "", digest: repository.digest };
const base = { kind: "bootstrap" as const, installation, repository, remote, workspaceIds: [] as string[] };
return { ...base, requestDigest: canonicalBootstrapRequestDigest(base) };
}
private addressedSnapshot(state: ActiveState, request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 {
private addressedSnapshot(state: ActiveState, request: any): RegistryAddressedSnapshotV1 {
const ids = state.revisions.map(revision => revision.id).sort();
return { schemaVersion: 1, commit: state.head, baseCommit: null, baseDigest: null, workspaceIds: ids, changedWorkspaceIds: ids, installation: request.installation, repository: { ...request.repository, head: state.head }, remote: { ...request.remote, head: state.head }, requestDigest: request.requestDigest };
return { schemaVersion: 1, commit: state.head as any, manifestSha256: digest(JSON.stringify(state)) as any, workspaces: ids.map(id => { const revision = state.revisions.find(x => x.id === id)!; return { workspaceId: id as any, revision: revision.commit as any, descriptorBlob: revision.blob as any, manifestSha256: digest(JSON.stringify(revision)) as any }; }), requestDigest: request.requestDigest ?? request.requestSha256 };
}
async bootstrap(): Promise<GitStatus> {
+1 -1
View File
@@ -14,6 +14,6 @@ export type WorkspaceErrorCode =
| "workspace_invalid" | "binding_missing" | "workspace_not_activatable"
| "workspace_stale" | "workspace_conflict" | "git_unavailable"
| "git_auth_failed" | "git_non_fast_forward" | "git_push_rejected"
| "connector_unavailable" | "semantic_index_incompatible";
| "connector_unavailable" | "semantic_index_incompatible" | "registry_bootstrap_recovery_conflict";
export type { WorkspaceV3 } from "./schema.js";
@@ -1 +1,11 @@
import { mkdir, writeFile } from "node:fs/promises"; const root=process.env.JOB_ROOT; if (!root) throw new Error("JOB_ROOT required"); await mkdir(root,{recursive:true,mode:0o700}); const id=process.env.RUN_ID??"a".repeat(32); await writeFile(`${root}/${id}.json`,JSON.stringify({runId:id,phase:"request_claimed"})+"\n",{flag:"wx",mode:0o600}); process.stdout.write(JSON.stringify({ok:true,runId:id}));
import { mkdir, open, writeFile, readFile, rm } from "node:fs/promises";
import { constants } from "node:fs";
const jobs = process.env.JOB_ROOT;
if (!jobs) throw new Error("JOB_ROOT required");
const id = process.env.RUN_ID ?? "a".repeat(32);
const barrier = process.env.BARRIER;
await mkdir(jobs, { recursive: true, mode: 0o700 });
if (barrier) { await writeFile(`${barrier}/${process.pid}.ready`, "ready", { flag: "wx", mode: 0o600 }); while (true) { try { await readFile(`${barrier}/release`); break; } catch { await new Promise(r => setTimeout(r, 5)); } } }
let winner = false;
try { const h = await open(`${jobs}/${id}.json`, constants.O_WRONLY | constants.O_CREAT | constants.O_EXCL | constants.O_NOFOLLOW, 0o600); await h.writeFile(JSON.stringify({ schemaVersion: 1, runId: id, phase: "request_claimed" }) + "\n"); await h.sync(); await h.close(); winner = true; } catch (e) { if (e?.code !== "EEXIST") throw e; }
process.stdout.write(JSON.stringify({ ok: true, runId: id, winner }) + "\n");
@@ -1 +1 @@
import {describe,it,expect} from "vitest"; import * as publication from "../src/workspaces/registry-publication.js"; describe("pull job exports",()=>it("exports addressed state",()=>expect(publication.REGISTRY_SCAN_LIMITS_V1.entries).toBeGreaterThan(0)));
import {describe,it,expect} from "vitest"; import * as publication from "../src/workspaces/registry-publication.js"; describe("pull job exports",()=>it("exports addressed state and exact scan bounds",()=>expect(publication.REGISTRY_SCAN_LIMITS_V1.maximumDirectoryEntries).toBe(4096)));
+4 -2
View File
@@ -116,7 +116,7 @@ const revision: WorkspaceRevision = {
snapshotPath: "/registry/snapshots/psd-clinical.yaml",
};
type RegistryFake = Pick<WorkspaceRegistry, "bootstrap" | "pull" | "list" | "read" | "publish">;
type RegistryFake = Pick<WorkspaceRegistry, "bootstrap" | "pull" | "list" | "read" | "publish"> & { ensureBootstrapAddressed: ReturnType<typeof vi.fn>; publishAddressed: ReturnType<typeof vi.fn> };
function registryFake(overrides: Partial<RegistryFake> = {}): RegistryFake {
return {
@@ -129,6 +129,8 @@ function registryFake(overrides: Partial<RegistryFake> = {}): RegistryFake {
list: vi.fn(async () => [revision]),
read: vi.fn(async () => ({ workspace, revision })),
publish: vi.fn(async () => revision),
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) }] } })),
...overrides,
};
}
@@ -218,7 +220,7 @@ test("returns a redacted registry status and pulls without Git credential detail
expect(status.statusCode).toBe(200);
expect(status.json()).toEqual({
branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: true, lastError: "git_auth_failed",
branch: "main", head: revision.commit, ahead: 0, behind: 0, degraded: false,
});
expect(pull.statusCode).toBe(200);
expect(JSON.stringify([status.json(), pull.json()])).not.toMatch(/token|password|ssh:\/\//i);
@@ -1 +1,16 @@
import {describe,it,expect} from "vitest"; describe("addressed publication process contract",()=>{it("has a bounded run id",()=>expect("a".repeat(32)).toHaveLength(32));});
import { describe, expect, it } from "vitest";
import { mkdtemp, mkdir, readdir, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { spawn } from "node:child_process";
const worker = join(process.cwd(), "test/fixtures/workspace-registry-addressed-worker.mjs");
async function run(env: Record<string,string>) { return await new Promise<string>((resolve, reject) => { const p = spawn(process.execPath, [worker], { env: { ...process.env, ...env }, stdio: ["ignore", "pipe", "pipe"] }); let out = ""; p.stdout.on("data", b => out += b); p.on("error", reject); p.on("exit", c => c === 0 ? resolve(out) : reject(new Error(`worker ${c}`))); }); }
describe("addressed publication process ownership", () => {
it("serializes two claimers and leaves one durable final artifact", async () => {
const root = await mkdtemp(join(process.env.TMPDIR ?? "/tmp", "thoth-addressed-process-")); const barrier = await mkdir(join(root, "barrier"), { recursive: true }).then(() => join(root, "barrier")); const jobs = join(root, "jobs");
const env = { JOB_ROOT: jobs, BARRIER: barrier, RUN_ID: "a".repeat(32) };
const a = run(env), b = run(env);
for (let i = 0; i < 100; i++) { if ((await readdir(barrier)).filter(x => x.endsWith(".ready")).length === 2) break; await new Promise(r => setTimeout(r, 5)); }
await writeFile(join(barrier, "release"), "go", { flag: "wx", mode: 0o600 }); const [one, two] = await Promise.all([a,b]); const results = [JSON.parse(one), JSON.parse(two)];
expect(results.filter(x => x.winner)).toHaveLength(1); expect(await readdir(jobs)).toEqual([`${"a".repeat(32)}.json`]);
}, 5000);
});