From ca37156a5681089c0787850408ec72f398d9188b Mon Sep 17 00:00:00 2001 From: mptyl Date: Tue, 11 Aug 2026 15:35:42 +0200 Subject: [PATCH] fix: close workspace ownership and API escapes --- .../native/workspace-fs-at/workspace_fs_at.cc | 28 ++++++++- backend/package.json | 5 +- .../src/native/workspace-fs-at-binding.d.ts | 4 +- backend/src/workspaces/preprocessing-state.ts | 61 ++++++++++--------- .../src/workspaces/registry-publication.ts | 19 +++--- backend/src/workspaces/registry.ts | 6 +- backend/src/workspaces/workspace-fs-at.ts | 24 +++++--- .../workspaces/workspace-lock-root-lease.ts | 55 ++++++++++++----- .../test/workspace-lock-root-lease.test.ts | 16 ++++- docker/core.Dockerfile | 7 ++- harness/tests/test_workspace_writer_lock.py | 20 ++++++ harness/tht/workspace_writer_lock.py | 34 +++++++++-- 12 files changed, 202 insertions(+), 77 deletions(-) diff --git a/backend/native/workspace-fs-at/workspace_fs_at.cc b/backend/native/workspace-fs-at/workspace_fs_at.cc index d0b644a2..379f9a80 100644 --- a/backend/native/workspace-fs-at/workspace_fs_at.cc +++ b/backend/native/workspace-fs-at/workspace_fs_at.cc @@ -1,5 +1,6 @@ #include #include +#include #include #include #include @@ -108,7 +109,32 @@ static napi_value fsyncFn(napi_env env,napi_callback_info info){size_t n=1;napi_ if(rc<0)return fail(env,"fsync",errno);return nullptr;} static napi_value closeFn(napi_env env,napi_callback_info info){size_t n=1;napi_value a[1];napi_get_cb_info(env,info,&n,a,nullptr,nullptr);Handle*h;if(n!=1||!getHandle(env,a[0],&h))return fail(env,"close",EBADF);if(h->borrows)return fail(env,"close",EBUSY);h->closed=true;handles.erase(h);int rc=::close(h->fd);if(rc<0){int e=errno;if(e==EINTR)return fail(env,"close",e,"ERR_WORKSPACE_FS_AT_CLOSE_UNCERTAIN");return fail(env,"close",e);}return nullptr;} static napi_value withFdFn(napi_env env,napi_callback_info info){size_t n=2;napi_value a[2];napi_get_cb_info(env,info,&n,a,nullptr,nullptr);Handle*h; napi_valuetype t;if(n!=2||!getHandle(env,a[0],&h)||napi_typeof(env,a[1],&t)!=napi_ok||t!=napi_function)return fail(env,"borrow",EINVAL);h->borrows++;napi_value argv; napi_create_int32(env,h->fd,&argv); napi_value out; napi_value global; napi_get_global(env, &global); napi_status rc=napi_call_function(env, global,a[1],1,&argv,&out);h->borrows--;if(rc!=napi_ok)return nullptr;return out;} +static napi_value fdNumberForSynchronousBorrowFn(napi_env env,napi_callback_info info){ size_t n=1; napi_value a[1]; napi_get_cb_info(env,info,&n,a,nullptr,nullptr); Handle*h; if(n!=1||!getHandle(env,a[0],&h)) return fail(env,"fcntl",EBADF); napi_value out; napi_create_int32(env,h->fd,&out); return out; } +static napi_value fcntlFn(napi_env env,napi_callback_info info){ + size_t n=2; napi_value a[2]; napi_get_cb_info(env,info,&n,a,nullptr,nullptr); Handle*h; std::string op; + if(n!=2||!getHandle(env,a[0],&h)||!str(env,a[1],op)||h->directory) return fail(env,"fcntl",EINVAL); + if(op=="unlock") { +#ifdef F_OFD_SETLK + struct flock lk{}; lk.l_type=F_UNLCK; lk.l_whence=SEEK_SET; if(::fcntl(h->fd,F_OFD_SETLK,&lk)<0) return fail(env,"fcntl",errno); +#else + if(::flock(h->fd,LOCK_UN)<0) return fail(env,"fcntl",errno); +#endif + napi_value out; napi_create_string_utf8(env,"available",NAPI_AUTO_LENGTH,&out); return out; + } + if(op=="probe-exclusive-nonblocking" || op=="hold-exclusive") { +#ifdef F_OFD_SETLK + struct flock lk{}; lk.l_type=F_WRLCK; lk.l_whence=SEEK_SET; int rc=::fcntl(h->fd,F_OFD_SETLK,&lk); + if(rc==0) { if(op=="probe-exclusive-nonblocking") { lk.l_type=F_UNLCK; ::fcntl(h->fd,F_OFD_SETLK,&lk); } napi_value out; napi_create_string_utf8(env,op=="hold-exclusive"?"held":"available",NAPI_AUTO_LENGTH,&out); return out; } +#else + int rc=::flock(h->fd,LOCK_EX|LOCK_NB); + if(rc==0) { if(op=="probe-exclusive-nonblocking") ::flock(h->fd,LOCK_UN); napi_value out; napi_create_string_utf8(env,op=="hold-exclusive"?"held":"available",NAPI_AUTO_LENGTH,&out); return out; } +#endif + if(errno==EWOULDBLOCK||errno==EAGAIN||errno==EACCES) { napi_value out; napi_create_string_utf8(env,"held",NAPI_AUTO_LENGTH,&out); return out; } + return fail(env,"fcntl",errno); + } + return fail(env,"fcntl",EINVAL); +} static napi_value duplicateForChildStdioFn(napi_env env,napi_callback_info info){size_t n=4;napi_value a[4];napi_get_cb_info(env,info,&n,a,nullptr,nullptr);Handle*w,*r;int32_t wf,rf;if(n!=4||!getHandle(env,a[0],&w)||!getHandle(env,a[1],&r)||w->directory||!r->directory||napi_get_value_int32(env,a[2],&wf)!=napi_ok||napi_get_value_int32(env,a[3],&rf)!=napi_ok||wf<0||rf<0)return fail(env,"dup2",EINVAL);if(w->borrows||r->borrows)return fail(env,"dup2",EBUSY);int rc;do{rc=::dup2(w->fd,wf);}while(rc<0&&errno==EINTR);if(rc<0)return fail(env,"dup2",errno);do{rc=::dup2(r->fd,rf);}while(rc<0&&errno==EINTR);if(rc<0)return fail(env,"dup2",errno);return nullptr;} -static napi_value init(napi_env env,napi_value exports){napi_property_descriptor pub[]={{"openat",0,openatFn,0,0,0,napi_enumerable,0},{"mkdirat",0,mkdiratFn,0,0,0,napi_enumerable,0},{"fstatat",0,fstatatFn,0,0,0,napi_enumerable,0},{"fsyncDirectory",0,fsyncFn,0,0,0,napi_enumerable,0},{"close",0,closeFn,0,0,0,napi_enumerable,0}};napi_define_properties(env,exports,5,pub);napi_property_descriptor priv[]={{"withFd",0,withFdFn,0,0,0,napi_default,0},{"duplicateForChildStdio",0,duplicateForChildStdioFn,0,0,0,napi_default,0}};napi_define_properties(env,exports,2,priv);return exports;} +static napi_value init(napi_env env,napi_value exports){napi_property_descriptor pub[]={{"openat",0,openatFn,0,0,0,napi_enumerable,0},{"mkdirat",0,mkdiratFn,0,0,0,napi_enumerable,0},{"fstatat",0,fstatatFn,0,0,0,napi_enumerable,0},{"fsyncDirectory",0,fsyncFn,0,0,0,napi_enumerable,0},{"close",0,closeFn,0,0,0,napi_enumerable,0}};napi_define_properties(env,exports,5,pub);napi_property_descriptor priv[]={{"withFd",0,withFdFn,0,0,0,napi_default,0},{"duplicateForChildStdio",0,duplicateForChildStdioFn,0,0,0,napi_default,0},{"fdNumberForSynchronousBorrow",0,fdNumberForSynchronousBorrowFn,0,0,0,napi_default,0},{"fcntl",0,fcntlFn,0,0,0,napi_default,0}};napi_define_properties(env,exports,4,priv);return exports;} } NAPI_MODULE(NODE_GYP_MODULE_NAME,init) diff --git a/backend/package.json b/backend/package.json index 83766fdb..986d2811 100644 --- a/backend/package.json +++ b/backend/package.json @@ -5,12 +5,13 @@ "scripts": { "dev": "tsx watch src/server.ts", "prebuild": "node scripts/clean-dist.mjs", - "build": "npm run build:ts", + "build": "npm run build:native && npm run build:ts", "test": "vitest run", "start": "node dist/server.js", "build:native": "node scripts/build-workspace-fs-at.mjs", "postinstall": "npm run build:native", - "build:ts": "node scripts/clean-dist.mjs && tsc -p tsconfig.json" + "build:ts": "node scripts/clean-dist.mjs && tsc -p tsconfig.json", + "test:native": "npm run build:native && vitest run test/workspace-fs-at-native.test.ts" }, "dependencies": { "@fastify/cors": "^11.2.0", diff --git a/backend/src/native/workspace-fs-at-binding.d.ts b/backend/src/native/workspace-fs-at-binding.d.ts index 4da7cc20..e832d894 100644 --- a/backend/src/native/workspace-fs-at-binding.d.ts +++ b/backend/src/native/workspace-fs-at-binding.d.ts @@ -3,7 +3,7 @@ declare const nativeWorkspaceFsAtComponentBrand: unique symbol; export interface NativeWorkspaceFsAtHandleV1 { readonly [nativeWorkspaceFsAtHandleBrand]: true; } export type NativeWorkspaceFsAtComponentV1 = string & { readonly [nativeWorkspaceFsAtComponentBrand]: true }; export interface NativeWorkspaceFsAtStatV1 { readonly device: bigint; readonly inode: bigint; readonly mode: number; readonly uid: number; readonly gid: number; readonly nlink: bigint; } -export interface NativeWorkspaceFsAtErrorV1 extends Error { readonly code:string; readonly errno:number; readonly syscall:"openat"|"mkdirat"|"fstat"|"fstatat"|"fsync"|"borrow"|"dup2"|"close"; } +export interface NativeWorkspaceFsAtErrorV1 extends Error { readonly code:string; readonly errno:number; readonly syscall:"openat"|"mkdirat"|"fstat"|"fstatat"|"fsync"|"fcntl"|"close"; } export interface NativeWorkspaceFsAtOpenResultV1 { readonly handle: NativeWorkspaceFsAtHandleV1; readonly openedStat: NativeWorkspaceFsAtStatV1; } export interface WorkspaceFsAtBindingV1 { openat(input: { readonly parent: NativeWorkspaceFsAtHandleV1|null; readonly name: "/"|NativeWorkspaceFsAtComponentV1; readonly kind: "directory"|"regular_lock"; readonly createMode: 0|0o600 }): NativeWorkspaceFsAtOpenResultV1; @@ -11,4 +11,6 @@ export interface WorkspaceFsAtBindingV1 { fstatat(parent:NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1):NativeWorkspaceFsAtStatV1; fsyncDirectory(handle:NativeWorkspaceFsAtHandleV1):void; close(handle:NativeWorkspaceFsAtHandleV1):void; + fdNumberForSynchronousBorrow(handle: NativeWorkspaceFsAtHandleV1): number; + fcntl(handle: NativeWorkspaceFsAtHandleV1, operation: "probe-exclusive-nonblocking" | "hold-exclusive" | "unlock"): "held" | "available"; } diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index b1c3c1c8..665160b5 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -1,34 +1,37 @@ import { createHash, randomBytes } from "node:crypto"; -import { lstatSync, realpathSync } from "node:fs"; +import { lstatSync } from "node:fs"; import { mkdir, open as openFile, readFile, readdir, rename, rm, writeFile } from "node:fs/promises"; import { join, dirname, isAbsolute, relative } from "node:path"; import { spawn } from "node:child_process"; import { WorkspaceFsAtV1, type OwnedWorkspaceFsAtRegularFile, + type OwnedWorkspaceFsAtDirectory, } from "./workspace-fs-at.js"; +import type { RuntimeConfigLease } from "./runtime-config-lease.js"; import { BorrowedVerifiedWorkspaceLockRootLease, VerifiedWorkspaceLockRootLease, type CanonicalWorkspaceId, type WorkspaceLockRootIdentityV1, - WorkspaceRootLock, + Revision40, } from "./workspace-lock-root-lease.js"; export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; } export interface PreprocessingRunStateV1 { readonly schemaVersion: 1; readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; - readonly revision: string; readonly operation: string; readonly phase: string; + readonly revision: Revision40; readonly operation: string; readonly phase: string; readonly artifacts: readonly ArtifactIdentity[]; readonly createdAt: string; readonly updatedAt: string; } -export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; } -export interface ResumeRunInput { readonly runId: string; readonly workspaceId: string; readonly revision: string; readonly operation: string; } +export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly operation: string; } +export interface ResumeRunInput { readonly runId: string; readonly workspaceId: string; readonly revision: Revision40; readonly operation: string; } /** A transition is deliberately closed: identity and artifacts are never caller-writable. */ export interface RunTransition { readonly phase: string; } export interface FkReviewInput { readonly candidate: ArtifactIdentity; readonly reviewSha256: string; readonly annotationSha256?: string; } export interface FkReviewRecordV1 extends FkReviewInput { readonly runId: string; readonly recordedAt: string; } const digest = (x: Uint8Array | string) => createHash("sha256").update(x).digest("hex"); +const INTERNAL_STATE = Symbol("preprocessing-state-internal"); const RUN_ID = /^[0-9a-f]{32}$/; const SHA256 = /^(?:sha256:)?[0-9a-f]{64}$/; const REVISION = /^[0-9a-f]{40}$/; @@ -39,15 +42,9 @@ const MAX_FILE_BYTES = 1 << 20; const MAX_AGGREGATE_BYTES = 64 << 20; const MAX_ENTRIES = 4096; function fail(msg = "preprocessing_conflict"): Error { const e = new Error(msg); e.name = "PreprocessingConflictError"; return e; } +type WorkspaceRootLock = { assertPath(): void; close(): void; flock(kind: "shared" | "exclusive", wait: "blocking" | "nonblocking"): void; spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv): Promise; }; function id(v: string): void { if (typeof v !== "string" || !RUN_ID.test(v)) throw fail("invalid run id"); } function checkSha(v: string): void { if (typeof v !== "string" || !SHA256.test(v)) throw fail("invalid digest"); } -function checkedRoot(root: string | undefined): string { - if (!root || typeof root !== "string" || !isAbsolute(root) || root.includes("\0")) throw fail(); - const resolved = realpathSync(root); - const original = lstatSync(root); const st = lstatSync(resolved); - if (!original.isDirectory() || original.dev !== st.dev || original.ino !== st.ino || !st.isDirectory() || (st.mode & 0o777) !== 0o700 || st.nlink < 2) throw fail(); - return resolved; -} function strictObject(value: unknown, keys: readonly string[]): value is Record { if (!value || typeof value !== "object" || Array.isArray(value)) return false; const got = Object.keys(value as object).sort(); @@ -75,9 +72,8 @@ async function durableJson(path: string, value: unknown): Promise { } export class PreprocessingStateStore { - private readonly configuredRoot: string; - constructor(configuredRoot: string) { this.configuredRoot = checkedRoot(configuredRoot); } - private root(): string { return this.configuredRoot; } + constructor(private readonly rootLease: VerifiedWorkspaceLockRootLease) {} + private root(): string { this.rootLease.assertLive(); return this.rootLease.anchoredPath(); } private paths(input: { runId: string }) { id(input.runId); const root = this.root(); const base = join(root, "preprocessing"); return { base, jobs: join(base, "jobs"), path: join(base, "jobs", `${input.runId}.json`), candidate: join(base, "fk-candidates", `${input.runId}.yaml`), review: join(base, "fk-reviews", `${input.runId}.json`) }; @@ -138,15 +134,15 @@ export class PreprocessingStateStore { } } -export interface DwhLockedChildRequest { readonly kind: "dwh_preprocess"; readonly stage: "introspect" | "lsh"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; } -export interface SchemaLockedChildRequest { readonly kind: "schema_preprocess"; readonly stage: "fk_suggest" | "fk_check" | "schema_index"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; readonly reviewedArtifact: ArtifactIdentity | null; } -export interface EvidenceLockedChildRequest { readonly kind: "evidence_preprocess"; readonly stage: "http_publish"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; } +export interface DwhLockedChildRequest { readonly kind: "dwh_preprocess"; readonly stage: "introspect" | "lsh"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; } +export interface SchemaLockedChildRequest { readonly kind: "schema_preprocess"; readonly stage: "fk_suggest" | "fk_check" | "schema_index"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; readonly reviewedArtifact: ArtifactIdentity | null; } +export interface EvidenceLockedChildRequest { readonly kind: "evidence_preprocess"; readonly stage: "http_publish"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: Revision40; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: RuntimeConfigLease; readonly childRunId: string; } export type WorkspaceLockedChildRequest = DwhLockedChildRequest | SchemaLockedChildRequest | EvidenceLockedChildRequest; export interface WorkspaceLockedChildResult { readonly exitCode: number; readonly stdout: Uint8Array; readonly stderr: Uint8Array; } export class BorrowedWorkspaceSessionReadersExclusiveLockLease { private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} - static from(id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } + static [INTERNAL_STATE](id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } assertLive(): void { if (!this.live) throw fail(); } invalidate(): void { this.live = false; } } @@ -154,13 +150,14 @@ export class WorkspaceWriterLockCapability { private live = true; private settled = false; private readerExclusive = false; private spawnActive = false; private poisoned = false; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1, private readonly root: VerifiedWorkspaceLockRootLease, private readonly writer: WorkspaceRootLock) {} private assertLive(): void { if (!this.live || this.settled || this.poisoned) throw fail(); } + assertWriterPath(): void { this.assertLive(); this.root.assertLive(); this.writer.assertPath(); } invalidateForSettlement(): void { this.settled = true; } - static create(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); } + static [INTERNAL_STATE](id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); } async runUnderSessionReadersExclusive(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise): Promise { this.assertLive(); if (this.readerExclusive || this.spawnActive) throw fail(); this.readerExclusive = true; let lock: WorkspaceRootLock; try { lock = await this.root.acquireSessionReadersExclusive(); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); } - const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.from(this.workspaceId, this.rootIdentity); + const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease[INTERNAL_STATE](this.workspaceId, this.rootIdentity); try { return await action(borrowed); } catch (error) { throw error; } finally { borrowed.invalidate(); let cleanupError: unknown; try { lock.close(); } catch (error) { cleanupError = error; } @@ -171,22 +168,27 @@ export class WorkspaceWriterLockCapability { async spawnChild(request: WorkspaceLockedChildRequest): Promise { this.assertLive(); if (!this.readerExclusive || this.spawnActive || !request || request.workspaceId !== this.workspaceId) throw fail(); if (!RUN_ID.test(request.childRunId) || request.rootIdentity.device !== this.rootIdentity.device || request.rootIdentity.inode !== this.rootIdentity.inode || request.rootIdentity.workspaceId !== this.workspaceId) throw fail(); - const cfg = request.runtimeConfig as { workspaceId?: string; revision?: string; configPath?: string; path?: string } | null; - if (!cfg || cfg.workspaceId !== this.workspaceId || cfg.revision !== request.revision || typeof (cfg.configPath ?? cfg.path) !== "string") throw fail(); + const cfg = request.runtimeConfig; + if (cfg.workspaceId !== this.workspaceId || cfg.workspaceRevision !== request.revision || typeof cfg.path !== "string") throw fail(); this.spawnActive = true; try { - const configPath = (cfg.configPath ?? cfg.path)!; let argv: string[]; + const configPath = cfg.path; let argv: string[]; switch (request.kind) { case "dwh_preprocess": argv = ["-m", "tht.cli", "preprocess", "dwh", "--steps", request.stage, "--json", "-c", configPath]; break; case "schema_preprocess": argv = ["-m", "tht.cli", "schema", request.stage === "fk_suggest" ? "suggest-fks" : request.stage === "fk_check" ? "check" : "index", "--json", "-c", configPath]; break; case "evidence_preprocess": argv = ["-m", "tht.cli", "preprocess", "evidence", "--json", "-c", configPath]; break; default: throw fail(); } return await this.root.spawnChild(this.writer, process.env.THT_PYTHON ?? "python3", argv, { ...process.env, THOTH_WORKSPACE_ID: this.workspaceId, THOTH_WORKSPACE_REVISION: request.revision, THOTH_WORKSPACE_DEVICE: String(this.rootIdentity.device), THOTH_WORKSPACE_INODE: String(this.rootIdentity.inode) }); } finally { this.spawnActive = false; } } async close(): Promise { if (!this.live) return; if (this.readerExclusive || this.spawnActive) throw fail(); this.live = false; let error: unknown; try { this.writer.close(); } catch (e) { error = e; } try { await this.root.close(); } catch (e) { error ??= e; } if (error) throw fail(); } } -function makeWriterCapability(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return WorkspaceWriterLockCapability.create(id, identity, root, writer); } +function makeWriterCapability(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return WorkspaceWriterLockCapability[INTERNAL_STATE](id, identity, root, writer); } +export interface BorrowedOrderedWorkspaceWriterLeaseV1 { + readonly workspaceId: CanonicalWorkspaceId; + readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; + readonly writerCapability: WorkspaceWriterLockCapability; +} export interface OrderedWorkspaceCapability { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; } export class OrderedWorkspaceWriterCapabilitySet { private live = true; private constructor(private readonly caps: Map) {} - static make(caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } + static [INTERNAL_STATE](caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } invalidate(): void { this.live = false; for (const cap of this.caps.values()) cap.invalidateForSettlement(); } get workspaceIds(): readonly CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; } async forWorkspace(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise): Promise { if (!this.live) throw fail(); const cap = this.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: BorrowedVerifiedWorkspaceLockRootLease.make(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); } @@ -197,8 +199,11 @@ export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonl 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.make(new Map(caps.map(c => [c.workspaceId, c]))); result = await action(set); - } catch (error) { callbackError = 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; } + } catch (error) { callbackError ??= error; } set?.invalidate(); for (const cap of [...caps].reverse()) await cap.close().catch(() => undefined); if (callbackError) throw callbackError; return result as T; diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts index facac0d3..bf377aca 100644 --- a/backend/src/workspaces/registry-publication.ts +++ b/backend/src/workspaces/registry-publication.ts @@ -2,12 +2,9 @@ import { createHash, randomBytes } from "node:crypto"; 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; } +import type { CanonicalWorkspaceId, Revision40, Sha256Hex } from "./workspace-lock-root-lease.js"; +import type { BorrowedOrderedWorkspaceWriterLeaseV1, BorrowedWorkspaceSessionReadersExclusiveLockLease, OrderedWorkspaceWriterCapabilitySet } from "./preprocessing-state.js"; -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"; @@ -23,8 +20,6 @@ export interface RegistryAddressedSnapshotV1 extends RegistryActiveSnapshotV1 { 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 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 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; } @@ -56,7 +51,7 @@ export class CapabilityAwareRegistryPublicationLifecycleOwner { return enter(0); } } -export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1 +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 } @@ -90,8 +85,12 @@ export class RegistryAddressedPublicationStore { private path(runId: string): string { if (!RUN.test(runId)) throw CONFLICT(); return join(this.jobsDirectory, `${runId}.json`); } private async dirs(): Promise { 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 { 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 { 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; + async claim(request: RegistryAddressedRequestV1, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise { + await this.dirs(); + if (!RUN.test(runId)) throw CONFLICT(); + const now = new Date().toISOString(); + const state = { schemaVersion: 1, runId, requestSha256: request.requestSha256, jobArtifactPath: `addressed-publication-jobs/${runId}.json`, phase: "request_claimed" as const, operation: request.operation, installationIdentitySha256: request.installationIdentitySha256, repositoryIdentitySha256: request.repositoryIdentitySha256, remoteRefIdentitySha256: request.remoteRefIdentitySha256, advertisedTargetCommit: null, immutableTargetRef: null, fetchedTargetCommit: null, targetCommit: null, targetManifestSha256: null, targetWorkspaces: null, changedWorkspaceIds: null, planSha256: null, changedSetSha256: null, participantsSha256: null, synchronizersSha256: null, publicationIntentSha256: null, publishedActiveStateSha256: null, terminalResultSha256: null, priorStateSha256: null, baseCommit: request.operation === "registry_pull" && "expectedBaseCommit" in request ? request.expectedBaseCommit : null, baseManifestSha256: null, baseWorkspaces: [], changedSetRule: null }; + await this.durable(this.path(runId), state, true); return state as unknown as RegistryAddressedPublicationStateV1; } async read(runId: RegistryRunId32): Promise { const path = this.path(runId); let parsed: unknown; diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 2506f992..9dbac9f3 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -22,9 +22,7 @@ import { addressedRunId, canonicalBootstrapRequestDigest, type RegistryAddressedSnapshotV1, - type RegistryBootstrapAddressedRequestV1, type RegistryEnsureBootstrapAddressedResultV1, - type RegistryPublishAddressedRequestV1, type RegistryAddressedRequestV1, type RegistryAddressedPublicationResultV1, } from "./registry-publication.js"; @@ -128,7 +126,7 @@ export class WorkspaceRegistry { } /** Automatic addressed recovery. The repository lock is acquired before active-state inspection. */ - async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1 | any): Promise { + async ensureBootstrapAddressed(request?: any): Promise { await this.repository.ensureLayout(); return await this.lock.run(async () => { const active = await this.tryActiveState(); @@ -151,7 +149,7 @@ export class WorkspaceRegistry { }); } - async publishAddressed(request: RegistryPublishAddressedRequestV1 | any): Promise { + async publishAddressed(request: any): Promise { await this.repository.ensureLayout(); return await this.lock.run(async () => { const before = await this.tryActiveState(); diff --git a/backend/src/workspaces/workspace-fs-at.ts b/backend/src/workspaces/workspace-fs-at.ts index 19814a02..8cd3f2e2 100644 --- a/backend/src/workspaces/workspace-fs-at.ts +++ b/backend/src/workspaces/workspace-fs-at.ts @@ -4,7 +4,7 @@ import { spawn } from "node:child_process"; import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1, NativeWorkspaceFsAtComponentV1 } from "../native/workspace-fs-at-binding.js"; const require = createRequire(import.meta.url); const binding = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node") as WorkspaceFsAtBindingV1 & { - withFd(handle: NativeWorkspaceFsAtHandleV1, action: (fd: number) => void): void; + withFd(handle: NativeWorkspaceFsAtHandleV1, action: (fd: number) => T): T; duplicateForChildStdio(writer: NativeWorkspaceFsAtHandleV1, root: NativeWorkspaceFsAtHandleV1, writerFd: number, rootFd: number): void; }; export interface WorkspaceFsAtStatV1 extends NativeWorkspaceFsAtStatV1 {} @@ -36,21 +36,22 @@ abstract class Owned { } } export class OwnedWorkspaceFsAtDirectory extends Owned { - constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); } + private constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); } + static [INTERNAL](raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1): OwnedWorkspaceFsAtDirectory { return new OwnedWorkspaceFsAtDirectory(raw, stat, INTERNAL); } } export class OwnedWorkspaceFsAtRegularFile extends Owned { - constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); } - + private constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); } + static [INTERNAL](raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1): OwnedWorkspaceFsAtRegularFile { return new OwnedWorkspaceFsAtRegularFile(raw, stat, INTERNAL); } } function wrapDirectory(result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtDirectory { - try { if ((result.openedStat.mode & 0o170000) !== 0o040000) throw new Error("not a directory"); return new OwnedWorkspaceFsAtDirectory(result.handle, result.openedStat, INTERNAL); } + try { if ((result.openedStat.mode & 0o170000) !== 0o040000) throw new Error("not a directory"); return OwnedWorkspaceFsAtDirectory[INTERNAL](result.handle, result.openedStat); } catch (error) { try { binding.close(result.handle); } catch { /* preserve conversion error */ } throw error; } } function wrapLock(result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtRegularFile { try { const st = result.openedStat; if ((st.mode & 0o170000) !== 0o100000 || (st.mode & 0o777) !== 0o600 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink !== 1n) throw new Error("invalid lock identity"); - return new OwnedWorkspaceFsAtRegularFile(result.handle, st, INTERNAL); + return OwnedWorkspaceFsAtRegularFile[INTERNAL](result.handle, st); } catch (error) { try { binding.close(result.handle); } catch { /* preserve conversion error */ } throw error; } } function rawDirectory(value: OwnedWorkspaceFsAtDirectory): NativeWorkspaceFsAtHandleV1 { if (!rawHandles.has(value)) throw new Error("workspace descriptor is closed"); return rawHandles.get(value)!; } @@ -59,6 +60,12 @@ function withLockFd(value: OwnedWorkspaceFsAtRegularFile, action: (fd: number const count = borrowing.get(value) ?? 0; borrowing.set(value, count + 1); try { return binding.withFd(rawHandles.get(value)!, action as (fd: number) => void) as T; } finally { borrowing.set(value, count); } } +function anchoredDirectoryPath(value: OwnedWorkspaceFsAtDirectory): string { + if (!rawHandles.has(value)) throw new Error("workspace descriptor is closed"); + const count = borrowing.get(value) ?? 0; borrowing.set(value, count + 1); + try { return binding.withFd(rawHandles.get(value)!, fd => `${process.platform === "darwin" ? "/dev/fd" : "/proc/self/fd"}/${fd}`); } + finally { borrowing.set(value, count); } +} export class WorkspaceFsAtV1 { openRoot(): OwnedWorkspaceFsAtDirectory { return wrapDirectory(binding.openat({ parent: null, name: "/", kind: "directory", createMode: 0 })); } openDirectoryAt(parent: OwnedWorkspaceFsAtDirectory, name: string): OwnedWorkspaceFsAtDirectory { return wrapDirectory(binding.openat({ parent: rawDirectory(parent), name: component(name), kind: "directory", createMode: 0 })); } @@ -66,16 +73,19 @@ export class WorkspaceFsAtV1 { mkdirAt(parent: OwnedWorkspaceFsAtDirectory, name: string, mode: 0o700): void { binding.mkdirat(rawDirectory(parent), component(name), mode); } statAtNoFollow(parent: OwnedWorkspaceFsAtDirectory, name: string): WorkspaceFsAtStatV1 { return binding.fstatat(rawDirectory(parent), component(name)); } fsyncDirectory(directory: OwnedWorkspaceFsAtDirectory): void { binding.fsyncDirectory(rawDirectory(directory)); } + /** Returns a proc-fd path tied to the retained directory open description. */ + anchoredDirectoryPath(directory: OwnedWorkspaceFsAtDirectory): string { return anchoredDirectoryPath(directory); } flockOwnedLock(owned: OwnedWorkspaceFsAtRegularFile, kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { const fsExt: { flockSync(fd: number, operation: string): void } = require("fs-ext"); const operation = kind === "shared" ? (wait === "blocking" ? "sh" : "shnb") : (wait === "blocking" ? "ex" : "exnb"); withLockFd(owned, fd => fsExt.flockSync(fd, operation)); } + fcntlOwnedLock(owned: OwnedWorkspaceFsAtRegularFile, operation: "probe-exclusive-nonblocking" | "hold-exclusive" | "unlock"): "held" | "available" { if (!rawHandles.has(owned)) throw new Error("workspace descriptor is closed"); return binding.fcntl(rawHandles.get(owned)!, operation); } } /** Internal child boundary. It deliberately returns a ChildProcess result, never an FD. */ -export function spawnChildWithWorkspaceCapabilities( +export function spawnChildFromOwnedCapability( writer: OwnedWorkspaceFsAtRegularFile, root: OwnedWorkspaceFsAtDirectory, executable: string, diff --git a/backend/src/workspaces/workspace-lock-root-lease.ts b/backend/src/workspaces/workspace-lock-root-lease.ts index fb9dcac5..815af01c 100644 --- a/backend/src/workspaces/workspace-lock-root-lease.ts +++ b/backend/src/workspaces/workspace-lock-root-lease.ts @@ -1,20 +1,35 @@ import { lstatSync, realpathSync } from "node:fs"; import { resolve } from "node:path"; -import { spawnChildWithWorkspaceCapabilities, WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, type OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js"; -export type CanonicalWorkspaceId = string & { readonly __workspaceId: unique symbol }; +import { spawnChildFromOwnedCapability, WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, type OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js"; +declare const canonicalWorkspaceIdBrand: unique symbol; +declare const revision40Brand: unique symbol; +declare const sha256HexBrand: unique symbol; +export type CanonicalWorkspaceId = string & { readonly [canonicalWorkspaceIdBrand]: true }; +export type Revision40 = string & { readonly [revision40Brand]: true }; +export type Sha256Hex = string & { readonly [sha256HexBrand]: true }; export interface WorkspaceLockRootIdentityV1 { readonly schemaVersion: 1; readonly workspaceId: CanonicalWorkspaceId; readonly device: bigint; readonly inode: bigint; } export class CanonicalWorkspaceLockRootInput { private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly owner: symbol) {} } function conflict(message = "preprocessing_conflict"): Error { const e = new Error(message); e.name = "PreprocessingConflictError"; return e; } function exactRoot(stat: WorkspaceFsAtStatV1, uid: number): boolean { return (stat.mode & 0o170000) === 0o040000 && (stat.mode & 0o777) === 0o700 && stat.uid === uid && stat.nlink >= 2n; } function sameIdentity(a: {device: bigint|number; inode: bigint|number} | {dev: bigint|number; ino: bigint|number}, b: {device: bigint; inode: bigint}): boolean { const device = BigInt("device" in a ? a.device : a.dev); const inode = BigInt("inode" in a ? a.inode : a.ino); return device === b.device && inode === b.inode; } +const ROOT_INTERNAL = Symbol("workspace-root-internal"); +const activeWriterLocks = new Set(); + /** Internal opaque lock returned only after the root has been checked. */ -export class WorkspaceRootLock { +type InternalWorkspaceRootLock = { assertPath(): void; close(): void; flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void; spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv): Promise<{ exitCode: number; stdout: Uint8Array; stderr: Uint8Array }>; }; +class WorkspaceRootLock implements InternalWorkspaceRootLock { private live = true; - constructor(private readonly fs: WorkspaceFsAtV1, private readonly lock: OwnedWorkspaceFsAtRegularFile) {} - async spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { return spawnChildWithWorkspaceCapabilities(this.lock, root, executable, args, environment); } - flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { if (!this.live) throw conflict(); this.fs.flockOwnedLock(this.lock, kind, wait); } - close(): void { if (!this.live) return; this.live = false; this.lock.close(); } + constructor(private readonly fs: WorkspaceFsAtV1, private readonly root: OwnedWorkspaceFsAtDirectory, private readonly lock: OwnedWorkspaceFsAtRegularFile, private readonly name: "writer.lock" | "session-readers.lock", private readonly activeKey?: string) {} + assertPath(): void { + if (!this.live) throw conflict(); + const current = this.fs.statAtNoFollow(this.root, this.name); + const opened = this.lock.stat(); + if (current.device !== opened.device || current.inode !== opened.inode || current.mode !== opened.mode || current.uid !== opened.uid || current.nlink !== opened.nlink) throw conflict("preprocessing_conflict: lock pathname identity changed"); + } + async spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return spawnChildFromOwnedCapability(this.lock, root, executable, args, environment); } + flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { this.assertPath(); this.fs.flockOwnedLock(this.lock, kind, wait); if (kind === "exclusive" && wait === "nonblocking") { const probe = this.fs.fcntlOwnedLock(this.lock, "hold-exclusive"); if (probe !== "held") throw conflict(); } this.assertPath(); } + close(): void { if (!this.live) return; this.live = false; try { this.lock.close(); } finally { if (this.activeKey) activeWriterLocks.delete(this.activeKey); } } } export class BorrowedVerifiedWorkspaceLockRootLease { private constructor(readonly identity: WorkspaceLockRootIdentityV1, private readonly check: () => void) {} @@ -23,10 +38,10 @@ export class BorrowedVerifiedWorkspaceLockRootLease { } export class WorkspaceSessionReadersLockLease { private live = true; - private constructor(private readonly lock: WorkspaceRootLock, private readonly root: OwnedWorkspaceFsAtDirectory, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} + private constructor(private readonly lock: InternalWorkspaceRootLock, private readonly root: OwnedWorkspaceFsAtDirectory, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} transfer(): WorkspaceSessionReadersLockLease { if (!this.live) throw conflict(); this.live = false; return new WorkspaceSessionReadersLockLease(this.lock, this.root, this.rootIdentity); } async close(): Promise { if (!this.live) return; this.live = false; let failure: unknown; try { this.lock.close(); } catch (e) { failure = e; } try { this.root.close(); } catch (e) { failure ??= e; } if (failure) throw conflict(); } - static from(lock: WorkspaceRootLock, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1): WorkspaceSessionReadersLockLease { return new WorkspaceSessionReadersLockLease(lock, root, identity); } + static [ROOT_INTERNAL](lock: InternalWorkspaceRootLock, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1): WorkspaceSessionReadersLockLease { return new WorkspaceSessionReadersLockLease(lock, root, identity); } } export class VerifiedWorkspaceLockRootLease { private live = true; @@ -39,19 +54,29 @@ export class VerifiedWorkspaceLockRootLease { const retained = this.root.stat(); if (!sameIdentity(retained, this.identity) || !exactRoot(retained, this.serviceUid)) throw conflict("preprocessing_conflict: workspace root identity changed"); } assertLive(): void { this.assertPath(); } + /** Internal capability-scoped path; it resolves through the retained root FD. */ + anchoredPath(): string { this.assertPath(); return this.fs.anchoredDirectoryPath(this.root); } async borrow(action: (lease: { readonly identity: WorkspaceLockRootIdentityV1; assertLive(): void }) => Promise): Promise { this.assertPath(); this.borrowed++; try { return await action({ identity: this.identity, assertLive: () => this.assertPath() }); } finally { this.borrowed--; } } transfer(): VerifiedWorkspaceLockRootLease { this.assertPath(); if (this.borrowed) throw conflict("workspace root is borrowed"); this.live = false; return new VerifiedWorkspaceLockRootLease(this.owner, this.fs, this.root, this.identity, this.rootPath, this.serviceUid); } async acquireSessionReadersShared(): Promise { this.assertPath(); let lock: WorkspaceRootLock | undefined; - try { const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); lock = new WorkspaceRootLock(this.fs, owned); lock.flock("shared", "nonblocking"); this.assertPath(); this.live = false; return WorkspaceSessionReadersLockLease.from(lock, this.root, this.identity); } + try { const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); lock = new WorkspaceRootLock(this.fs, this.root, owned, "session-readers.lock"); lock.flock("shared", "nonblocking"); this.assertPath(); this.live = false; return WorkspaceSessionReadersLockLease[ROOT_INTERNAL](lock, this.root, this.identity); } catch (error) { try { lock?.close(); } catch { this.live = false; try { this.root.close(); } catch {} } throw conflict(); } } - async acquireWriterLock(): Promise { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "writer.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch (error) { try { lock.close(); } catch {} throw conflict(); } } - async acquireSessionReadersExclusive(): Promise { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, owned); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch { try { lock.close(); } catch {} throw conflict(); } } - async spawnChild(lock: WorkspaceRootLock, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return lock.spawn(this.root, executable, args, environment); } + async acquireWriterLock(): Promise { + this.assertPath(); + const key = `${this.identity.device}:${this.identity.inode}`; + if (activeWriterLocks.has(key)) throw conflict(); + const owned = this.fs.openOrCreateLockAt(this.root, "writer.lock", 0o600); + const lock = new WorkspaceRootLock(this.fs, this.root, owned, "writer.lock", key); + try { lock.assertPath(); lock.flock("exclusive", "nonblocking"); lock.assertPath(); activeWriterLocks.add(key); return lock; } + catch (error) { try { lock.close(); } catch {} throw conflict(); } + } + async acquireSessionReadersExclusive(): Promise { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, this.root, owned, "session-readers.lock"); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch { try { lock.close(); } catch {} throw conflict(); } } + async spawnChild(lock: InternalWorkspaceRootLock, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return lock.spawn(this.root, executable, args, environment); } fsync(): void { this.assertPath(); this.fs.fsyncDirectory(this.root); this.assertPath(); } async close(): Promise { if (!this.live) return; while (this.borrowed) await new Promise(resolve => setTimeout(resolve, 1)); this.live = false; try { this.root.close(); } catch { throw conflict(); } } - static from(owner: symbol, fs: WorkspaceFsAtV1, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, rootPath: string, uid: number): VerifiedWorkspaceLockRootLease { return new VerifiedWorkspaceLockRootLease(owner, fs, root, identity, rootPath, uid); } + static [ROOT_INTERNAL](owner: symbol, fs: WorkspaceFsAtV1, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, rootPath: string, uid: number): VerifiedWorkspaceLockRootLease { return new VerifiedWorkspaceLockRootLease(owner, fs, root, identity, rootPath, uid); } } export class VerifiedWorkspaceLockRootLeaseFactory { private readonly owner = Symbol("workspace-root-factory"); @@ -67,7 +92,7 @@ export class VerifiedWorkspaceLockRootLeaseFactory { } private checkParent(): void { try { const s = lstatSync(this.parentPath); const p = this.parent.stat(); if (!s.isDirectory() || !sameIdentity(s, p) || s.uid !== this.input.serviceUid || (s.mode & 0o777) !== 0o700) throw conflict("installation root identity changed"); } catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("installation root identity changed"); } } canonicalInput(workspaceId: string): CanonicalWorkspaceLockRootInput { if (!/^[a-z][a-z0-9-]{2,62}$/.test(workspaceId)) throw conflict(); return Object.assign(Object.create(CanonicalWorkspaceLockRootInput.prototype), { workspaceId: workspaceId as CanonicalWorkspaceId, owner: this.owner }) as CanonicalWorkspaceLockRootInput; } - private async open(input: CanonicalWorkspaceLockRootInput): Promise { if (input.owner !== this.owner) throw conflict(); this.checkParent(); let root: OwnedWorkspaceFsAtDirectory; try { root = this.input.workspaceFsAt.openDirectoryAt(this.parent, input.workspaceId); this.checkParent(); } catch { throw conflict(); } if (!exactRoot(root.stat(), this.input.serviceUid)) { root.close(); throw conflict(); } const identity = { schemaVersion: 1 as const, workspaceId: input.workspaceId, device: root.stat().device, inode: root.stat().inode }; const lease = VerifiedWorkspaceLockRootLease.from(this.owner, this.input.workspaceFsAt, root, identity, `${this.parentPath}/${input.workspaceId}`, this.input.serviceUid); try { lease.assertLive(); } catch { await lease.close().catch(() => undefined); throw conflict(); } return lease; } + private async open(input: CanonicalWorkspaceLockRootInput): Promise { if (input.owner !== this.owner) throw conflict(); this.checkParent(); let root: OwnedWorkspaceFsAtDirectory; try { root = this.input.workspaceFsAt.openDirectoryAt(this.parent, input.workspaceId); this.checkParent(); } catch { throw conflict(); } if (!exactRoot(root.stat(), this.input.serviceUid)) { root.close(); throw conflict(); } const identity = { schemaVersion: 1 as const, workspaceId: input.workspaceId, device: root.stat().device, inode: root.stat().inode }; const lease = VerifiedWorkspaceLockRootLease[ROOT_INTERNAL](this.owner, this.input.workspaceFsAt, root, identity, `${this.parentPath}/${input.workspaceId}`, this.input.serviceUid); try { lease.assertLive(); } catch { await lease.close().catch(() => undefined); throw conflict(); } return lease; } acquire(input: CanonicalWorkspaceLockRootInput): Promise { return this.open(input); } async acquireOrProvision(input: CanonicalWorkspaceLockRootInput): Promise { if (input.owner !== this.owner) throw conflict(); diff --git a/backend/test/workspace-lock-root-lease.test.ts b/backend/test/workspace-lock-root-lease.test.ts index 0e0d98fa..3e5b014b 100644 --- a/backend/test/workspace-lock-root-lease.test.ts +++ b/backend/test/workspace-lock-root-lease.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it, afterEach } from "vitest"; -import { mkdtempSync, renameSync, mkdirSync, rmSync, statSync, chmodSync, symlinkSync } from "node:fs"; +import { mkdtempSync, renameSync, mkdirSync, rmSync, statSync, chmodSync, symlinkSync, writeFileSync } from "node:fs"; import { join } from "node:path"; import { WorkspaceFsAtV1 } from "../src/workspaces/workspace-fs-at.js"; import { VerifiedWorkspaceLockRootLeaseFactory } from "../src/workspaces/workspace-lock-root-lease.js"; @@ -27,4 +27,18 @@ describe("retained canonical workspace root", () => { const parent = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(parent); const f = factory(parent); const a = await f.acquireOrProvision(f.canonicalInput("aaa-workspace")); const b = await f.acquireOrProvision(f.canonicalInput("bbb-workspace")); const { runUnderOrderedWorkspaceWriterLocks } = await import("../src/workspaces/preprocessing-state.js"); await runUnderOrderedWorkspaceWriterLocks([a, b], async set => { expect(set.workspaceIds).toEqual(["aaa-workspace", "bbb-workspace"]); }); await expect(b.close()).resolves.toBeUndefined(); }); + it("rejects writer pathname replacement without admitting a concurrent writer", async () => { + const parent = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(parent); + const f = factory(parent); const first = await f.acquireOrProvision(f.canonicalInput("abc-workspace")); + const second = await f.acquireOrProvision(f.canonicalInput("abc-workspace")); + const { runUnderWorkspaceWriterLock } = await import("../src/workspaces/preprocessing-state.js"); + const outcome = runUnderWorkspaceWriterLock(first, async () => { + renameSync(join(parent, "abc-workspace", "writer.lock"), join(parent, "abc-workspace", "writer.lock.old")); + writeFileSync(join(parent, "abc-workspace", "writer.lock"), "", { mode: 0o600 }); + await expect(runUnderWorkspaceWriterLock(second, async () => "entered")).rejects.toThrow(/preprocessing/); + return "done"; + }); + await expect(outcome).rejects.toThrow(/preprocessing/); + }); + }); diff --git a/docker/core.Dockerfile b/docker/core.Dockerfile index 79661db4..52a6b335 100644 --- a/docker/core.Dockerfile +++ b/docker/core.Dockerfile @@ -21,17 +21,20 @@ WORKDIR /src/backend COPY backend/package*.json ./ COPY backend/scripts ./scripts COPY backend/native ./native +COPY backend/src ./src +COPY backend/test ./test RUN npm ci RUN npm run build:native +RUN npm run test:native # ---- Stage 2: backend TypeScript -> dist ---- FROM node:22-bookworm@sha256:7725a5c2c83eed1d36258c66efae14b1ceccd021db9ed1d9559d3335ed3d68ed AS backend-build WORKDIR /src/backend COPY backend/package*.json ./ -RUN npm ci --ignore-scripts +COPY --from=backend-native-build /src/backend/node_modules ./node_modules COPY backend/src ./src COPY backend/scripts ./scripts COPY backend/tsconfig*.json ./ -RUN npm run build +RUN npm run build:ts RUN npm prune --omit=dev COPY --from=backend-native-build /src/backend/native/workspace-fs-at/build/Release/workspace_fs_at.node ./native/workspace-fs-at/build/Release/workspace_fs_at.node # fs-ext is the pinned runtime flock binding; carry its Node-22 build from the compiler stage. diff --git a/harness/tests/test_workspace_writer_lock.py b/harness/tests/test_workspace_writer_lock.py index 8dbd39f0..3676a999 100644 --- a/harness/tests/test_workspace_writer_lock.py +++ b/harness/tests/test_workspace_writer_lock.py @@ -1,6 +1,9 @@ from __future__ import annotations import os +import fcntl +import struct +import sys import pytest @@ -14,6 +17,9 @@ def test_verifier_checks_fd_identity(tmp_path): root=tmp_path / "root"; root.mkdir(mode=0o700); writer=root / "writer.lock"; writer.touch(mode=0o600) r=os.open(root, os.O_RDONLY); w=os.open(writer, os.O_RDWR) try: + fcntl.flock(w, fcntl.LOCK_EX | fcntl.LOCK_NB) + if hasattr(fcntl, "F_OFD_SETLK") and sys.platform.startswith("linux"): + fcntl.fcntl(w, fcntl.F_OFD_SETLK, struct.pack("hhqqi", fcntl.F_WRLCK, os.SEEK_SET, 0, 0, 0)) st=os.fstat(r); env={"THOTH_WORKSPACE_ID":"abc-workspace", "THOTH_WORKSPACE_REVISION":"a"*40, "THOTH_WORKSPACE_DEVICE":str(st.st_dev), "THOTH_WORKSPACE_INODE":str(st.st_ino)} cap=verify_workspace_writer_fds(writer_fd=w, root_fd=r, env=env) assert cap.inode == st.st_ino @@ -33,3 +39,17 @@ def test_verifier_rejects_unrelated_lock(tmp_path): verify_workspace_writer_fds(writer_fd=forged_fd, root_fd=root_fd, env=env) finally: os.close(forged_fd); os.close(root_fd) + + +def test_verifier_rejects_independently_opened_unlocked_writer_fd(tmp_path): + root = tmp_path / "root"; root.mkdir(mode=0o700) + writer = root / "writer.lock"; writer.touch(mode=0o600) + root_fd = os.open(root, os.O_RDONLY) + writer_fd = os.open(writer, os.O_RDWR) + try: + st = os.fstat(root_fd) + env = {"THOTH_WORKSPACE_ID": "abc-workspace", "THOTH_WORKSPACE_REVISION": "a" * 40, "THOTH_WORKSPACE_DEVICE": str(st.st_dev), "THOTH_WORKSPACE_INODE": str(st.st_ino)} + with pytest.raises(WorkspaceWriterConflict): + verify_workspace_writer_fds(writer_fd=writer_fd, root_fd=root_fd, env=env) + finally: + os.close(writer_fd); os.close(root_fd) diff --git a/harness/tht/workspace_writer_lock.py b/harness/tht/workspace_writer_lock.py index d4b9bb61..810ffc24 100644 --- a/harness/tht/workspace_writer_lock.py +++ b/harness/tht/workspace_writer_lock.py @@ -11,6 +11,8 @@ import fcntl import os import re import stat +import struct +import sys from dataclasses import dataclass @@ -76,14 +78,34 @@ def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: di lock = _fstat(lock_fd) if (lock.st_dev, lock.st_ino) != (writer.st_dev, writer.st_ino) or not stat.S_ISREG(lock.st_mode) or lock.st_uid != uid or (lock.st_mode & 0o777) != 0o600 or lock.st_nlink != 1: raise WorkspaceWriterConflict() - # A duplicate of the locked open description is re-lockable. An - # independently-opened description receives EWOULDBLOCK. + # Probe using a *different* open file description. An inherited FD 3 + # is valid only when its lock is already held: the independent probe + # must therefore receive EWOULDBLOCK. We deliberately never flock(3) + # here: doing so would turn an unheld, independently opened descriptor + # into an apparently valid capability. try: - fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + # Linux OFD locks identify the open file description rather than + # the process. The inherited FD 3 must already own this lock; + # an independent unlocked description can acquire the probe and + # is rejected. Darwin has no OFD constants, so retain flock's + # equivalent open-description probe there. + ofd_setlk = getattr(fcntl, "F_OFD_SETLK", None) if sys.platform.startswith("linux") else None + if ofd_setlk is not None: + lock_record = struct.pack("hhqqi", fcntl.F_WRLCK, os.SEEK_SET, 0, 0, 0) + fcntl.fcntl(lock_fd, ofd_setlk, lock_record) + unlock_record = struct.pack("hhqqi", fcntl.F_UNLCK, os.SEEK_SET, 0, 0, 0) + fcntl.fcntl(lock_fd, ofd_setlk, unlock_record) + raise WorkspaceWriterConflict() + fcntl.flock(lock_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) except OSError as exc: - if exc.errno in (errno.EACCES, errno.EAGAIN, errno.EWOULDBLOCK): - raise WorkspaceWriterConflict() from None - raise WorkspaceWriterConflict() from exc + if exc.errno not in (errno.EACCES, errno.EAGAIN, errno.EWOULDBLOCK): + raise WorkspaceWriterConflict() from exc + else: + try: + fcntl.flock(lock_fd, fcntl.LOCK_UN) + except OSError: + pass + raise WorkspaceWriterConflict() finally: try: os.close(lock_fd)