fix: close workspace ownership and API escapes

This commit is contained in:
2026-08-11 15:39:27 +02:00
parent 1b24337a98
commit ca37156a56
12 changed files with 202 additions and 77 deletions
@@ -1,5 +1,6 @@
#include <node_api.h>
#include <fcntl.h>
#include <sys/file.h>
#include <sys/stat.h>
#include <unistd.h>
#include <errno.h>
@@ -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)
+3 -2
View File
@@ -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",
+3 -1
View File
@@ -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";
}
+33 -28
View File
@@ -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<WorkspaceLockedChildResult>; };
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<string, unknown> {
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<void> {
}
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<T>(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise<T>): Promise<T> {
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<WorkspaceLockedChildResult> {
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<void> { 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<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) {}
static make(caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { return new OrderedWorkspaceWriterCapabilitySet(caps); }
static [INTERNAL_STATE](caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { 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<T>(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T> { 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<T>(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;
+9 -10
View File
@@ -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<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;
async claim(request: RegistryAddressedRequestV1, runId: RegistryRunId32 = ("runId" in request ? request.runId : addressedRunId())): Promise<RegistryAddressedPublicationStateV1> {
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<RegistryAddressedPublicationStateV1> {
const path = this.path(runId); let parsed: unknown;
+2 -4
View File
@@ -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<any> {
async ensureBootstrapAddressed(request?: any): Promise<any> {
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<any> {
async publishAddressed(request: any): Promise<any> {
await this.repository.ensureLayout();
return await this.lock.run(async () => {
const before = await this.tryActiveState();
+17 -7
View File
@@ -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<T>(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<T>(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,
@@ -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<string>();
/** 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<void> { 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<T>(action: (lease: { readonly identity: WorkspaceLockRootIdentityV1; assertLive(): void }) => Promise<T>): Promise<T> { 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<WorkspaceSessionReadersLockLease> {
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<WorkspaceRootLock> { 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<WorkspaceRootLock> { 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<WorkspaceRootLock> {
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<WorkspaceRootLock> { 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<void> { if (!this.live) return; while (this.borrowed) await new Promise<void>(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<VerifiedWorkspaceLockRootLease> { 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<VerifiedWorkspaceLockRootLease> { 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<VerifiedWorkspaceLockRootLease> { return this.open(input); }
async acquireOrProvision(input: CanonicalWorkspaceLockRootInput): Promise<VerifiedWorkspaceLockRootLease> {
if (input.owner !== this.owner) throw conflict();
+15 -1
View File
@@ -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/);
});
});
+5 -2
View File
@@ -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.
@@ -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)
+27 -5
View File
@@ -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
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)