diff --git a/backend/native/workspace-fs-at/workspace_fs_at.cc b/backend/native/workspace-fs-at/workspace_fs_at.cc index b8de192e..58e8118f 100644 --- a/backend/native/workspace-fs-at/workspace_fs_at.cc +++ b/backend/native/workspace-fs-at/workspace_fs_at.cc @@ -1,8 +1,12 @@ #include #include #include +#ifdef __linux__ +#include +#endif #include #include +#include #include #include #include @@ -85,6 +89,45 @@ static napi_value openatFn(napi_env env,napi_callback_info info) { if(kind=="regular_lock" && (!S_ISREG(st.st_mode)||(st.st_mode&0777)!=0600||st.st_uid!=(uid_t)::geteuid()||st.st_nlink!=1)) {::close(fd);return fail(env,"openat",EPERM);} return result(env,fd,kind=="directory",st); } +static napi_value openFileAtFn(napi_env env,napi_callback_info info) { + size_t n=1; napi_value a[1]; if(!getInput(env,info,a,&n)||n!=1)return fail(env,"openat",EINVAL); + napi_value pv,nv,av; if(napi_get_named_property(env,a[0],"parent",&pv)!=napi_ok||napi_get_named_property(env,a[0],"name",&nv)!=napi_ok||napi_get_named_property(env,a[0],"access",&av)!=napi_ok)return fail(env,"openat",EINVAL); + Handle* ph=nullptr; std::string name,access; if(!requireHandle(env,pv,&ph,"openat")||!ph->directory||!str(env,nv,name)||!validComponent(name)||!str(env,av,access)||(access!="read"&&access!="create"))return fail(env,"openat",EINVAL); + int flags=O_RDONLY|O_CLOEXEC|O_NOFOLLOW; if(access=="create") flags=O_WRONLY|O_CREAT|O_EXCL|O_CLOEXEC|O_NOFOLLOW; + int fd=openRetry(ph->fd,name.c_str(),flags,0600); if(fd<0)return fail(env,"openat",errno); + struct stat st; if(statRetry(fd,&st)<0){int e=errno;::close(fd);return fail(env,"fstat",e);} + if(!S_ISREG(st.st_mode)||(st.st_mode&0777)!=0600||st.st_uid!=(uid_t)::geteuid()||st.st_nlink!=1){::close(fd);return fail(env,"openat",EPERM);} + return result(env,fd,false,st); +} +static napi_value readFileFn(napi_env env,napi_callback_info info) { + size_t n=2; napi_value a[2]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=2)return fail(env,"read",EINVAL); + Handle*h=nullptr; int64_t max=0; if(!requireHandle(env,a[0],&h,"read")||h->directory||napi_get_value_int64(env,a[1],&max)!=napi_ok||max<0||max>67108864)return fail(env,"read",EINVAL); + struct stat st; if(statRetry(h->fd,&st)<0)return fail(env,"fstat",errno); if(!S_ISREG(st.st_mode)||(uint64_t)st.st_size>(uint64_t)max)return fail(env,"read",EFBIG); + void* data=nullptr; napi_value out; if(napi_create_buffer(env,(size_t)st.st_size,&data,&out)!=napi_ok)return fail(env,"read",ENOMEM); + size_t got=0; while(got<(size_t)st.st_size){ssize_t rc=::read(h->fd,(char*)data+got,(size_t)st.st_size-got);if(rc<0&&errno==EINTR)continue;if(rc<=0)return fail(env,"read",rc==0?EIO:errno);got+=(size_t)rc;} return out; +} +static napi_value writeFileFn(napi_env env,napi_callback_info info) { + size_t n=2; napi_value a[2]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=2)return fail(env,"write",EINVAL); + Handle*h=nullptr; void* data=nullptr; size_t len=0; if(!requireHandle(env,a[0],&h,"write")||h->directory||napi_get_buffer_info(env,a[1],&data,&len)!=napi_ok||len>67108864)return fail(env,"write",EINVAL); + size_t off=0; while(offfd,(char*)data+off,len-off);if(rc<0&&errno==EINTR)continue;if(rc<=0)return fail(env,"write",rc==0?EIO:errno);off+=(size_t)rc;} return nullptr; +} +static napi_value fsyncFileFn(napi_env env,napi_callback_info info) { + size_t n=1; napi_value a[1]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=1)return fail(env,"fsync",EINVAL); Handle*h=nullptr; if(!requireHandle(env,a[0],&h,"fsync")||h->directory)return fail(env,"fsync",EBADF); int rc; do{rc=::fsync(h->fd);}while(rc<0&&errno==EINTR); if(rc<0)return fail(env,"fsync",errno); return nullptr; +} +static napi_value renameAtFn(napi_env env,napi_callback_info info) { + size_t n=4; napi_value a[4]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=4)return fail(env,"renameat",EINVAL); Handle*h=nullptr; std::string from,to; bool replace=false; bool b=false; if(!requireHandle(env,a[0],&h,"renameat")||!h->directory||!str(env,a[1],from)||!str(env,a[2],to)||!validComponent(from)||!validComponent(to)||napi_get_value_bool(env,a[3],&b)!=napi_ok)return fail(env,"renameat",EINVAL); replace=b; +#ifdef __linux__ + if(!replace){ int rc=(int)syscall(SYS_renameat2,h->fd,from.c_str(),h->fd,to.c_str(),1 /* RENAME_NOREPLACE */); if(rc<0)return fail(env,"renameat",errno); return nullptr; } +#endif + if(!replace){ struct stat st; if(::fstatat(h->fd,to.c_str(),&st,AT_SYMLINK_NOFOLLOW)==0)return fail(env,"renameat",EEXIST); if(errno!=ENOENT)return fail(env,"renameat",errno); } + int rc; do{rc=::renameat(h->fd,from.c_str(),h->fd,to.c_str());}while(rc<0&&errno==EINTR); if(rc<0)return fail(env,"renameat",errno); return nullptr; +} +static napi_value listDirectoryFn(napi_env env,napi_callback_info info) { + size_t n=1; napi_value a[1]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=1)return fail(env,"readdir",EINVAL); Handle*h=nullptr; if(!requireHandle(env,a[0],&h,"readdir")||!h->directory)return fail(env,"readdir",EBADF); int dupfd=::dup(h->fd); if(dupfd<0)return fail(env,"readdir",errno); DIR*d=::fdopendir(dupfd); if(!d){int e=errno;::close(dupfd);return fail(env,"readdir",e);} napi_value out; if(napi_create_array(env,&out)!=napi_ok){::closedir(d);return fail(env,"readdir",ENOMEM);} uint32_t index=0; errno=0; while(dirent* ent=::readdir(d)){ if(strcmp(ent->d_name,".")==0||strcmp(ent->d_name,"..")==0)continue; napi_value s; napi_create_string_utf8(env,ent->d_name,NAPI_AUTO_LENGTH,&s); napi_set_element(env,out,index++,s); if(index>4096){::closedir(d);return fail(env,"readdir",EOVERFLOW);} } int e=errno;::closedir(d);if(e)return fail(env,"readdir",e);return out; +} +static napi_value unlinkAtFn(napi_env env,napi_callback_info info) { + size_t n=2; napi_value a[2]; if(napi_get_cb_info(env,info,&n,a,nullptr,nullptr)!=napi_ok||n!=2)return fail(env,"unlinkat",EINVAL); Handle*h=nullptr; std::string name; if(!requireHandle(env,a[0],&h,"unlinkat")||!h->directory||!str(env,a[1],name)||!validComponent(name))return fail(env,"unlinkat",EINVAL); int rc; do{rc=::unlinkat(h->fd,name.c_str(),0);}while(rc<0&&errno==EINTR); if(rc<0)return fail(env,"unlinkat",errno); return nullptr; +} static napi_value mkdiratFn(napi_env env,napi_callback_info info){size_t n=3;napi_value a[3];napi_get_cb_info(env,info,&n,a,nullptr,nullptr);Handle*h;std::string s;int32_t m;if(n!=3||!requireHandle(env,a[0],&h,"workspace-fs-at")||!h->directory||!str(env,a[1],s)||napi_get_value_int32(env,a[2],&m)!=napi_ok||m!=0700||!validComponent(s))return fail(env,"mkdirat",EINVAL);int rc;do{rc=::mkdirat(h->fd,s.c_str(),0700);}while(rc<0&&errno==EINTR);if(rc<0)return fail(env,"mkdirat",errno);return nullptr;} static napi_value fstatatFn(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 s;if(n!=2||!requireHandle(env,a[0],&h,"workspace-fs-at")||!h->directory||!str(env,a[1],s)||!validComponent(s))return fail(env,"fstatat",EINVAL);struct stat st;int rc;do{rc=::fstatat(h->fd,s.c_str(),&st,AT_SYMLINK_NOFOLLOW);}while(rc<0&&errno==EINTR);if(rc<0)return fail(env,"fstatat",errno);return statObj(env,st);} static napi_value fsyncFn(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||!requireHandle(env,a[0],&h,"workspace-fs-at")||!h->directory)return fail(env,"fsync",EBADF);int rc;do{rc=::fsync(h->fd);}while(rc<0&&errno==EINTR); @@ -95,6 +138,6 @@ static napi_value fsyncFn(napi_env env,napi_callback_info info){size_t n=1;napi_ 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=nullptr; if(n!=1) return fail(env,"close",EINVAL,"ERR_WORKSPACE_FS_AT_ARGUMENT"); auto hr=getHandle(env,a[0],&h); if(hr!=HANDLE_OK) return fail(env,"close",EINVAL,hr==HANDLE_CLOSED?"ERR_WORKSPACE_FS_AT_HANDLE_CLOSED":"ERR_WORKSPACE_FS_AT_ARGUMENT"); if(h->borrows) return fail(env,"close",EBUSY); auto *st=stateFor(env); { std::lock_guard g(st->mutex); h->closed=true; st->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||!requireHandle(env,a[0],&h,"workspace-fs-at")||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 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||!requireHandle(env,a[0],&w,"dup2")||!requireHandle(env,a[1],&r,"dup2")||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},{"openFileAt",0,openFileAtFn,0,0,0,napi_default,0},{"readFile",0,readFileFn,0,0,0,napi_default,0},{"writeFile",0,writeFileFn,0,0,0,napi_default,0},{"fsyncFile",0,fsyncFileFn,0,0,0,napi_default,0},{"renameAt",0,renameAtFn,0,0,0,napi_default,0},{"unlinkAt",0,unlinkAtFn,0,0,0,napi_default,0},{"listDirectory",0,listDirectoryFn,0,0,0,napi_default,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,12,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;} } NAPI_MODULE(NODE_GYP_MODULE_NAME,init) diff --git a/backend/src/native/workspace-fs-at-binding.d.ts b/backend/src/native/workspace-fs-at-binding.d.ts index 1d82718f..00bee4e1 100644 --- a/backend/src/native/workspace-fs-at-binding.d.ts +++ b/backend/src/native/workspace-fs-at-binding.d.ts @@ -6,7 +6,14 @@ export interface NativeWorkspaceFsAtStatV1 { readonly device: bigint; readonly i 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; + openat(input: { readonly parent: NativeWorkspaceFsAtHandleV1|null; readonly name: "/"|NativeWorkspaceFsAtComponentV1; readonly kind: "directory"|"regular_lock"|"regular_file"; readonly createMode: 0|0o600 }): NativeWorkspaceFsAtOpenResultV1; + openFileAt(input: { readonly parent: NativeWorkspaceFsAtHandleV1; readonly name: NativeWorkspaceFsAtComponentV1; readonly access: "read"|"create" }): NativeWorkspaceFsAtOpenResultV1; + readFile(handle: NativeWorkspaceFsAtHandleV1, maxBytes: number): Uint8Array; + writeFile(handle: NativeWorkspaceFsAtHandleV1, bytes: Uint8Array): void; + fsyncFile(handle: NativeWorkspaceFsAtHandleV1): void; + renameAt(parent: NativeWorkspaceFsAtHandleV1, from: NativeWorkspaceFsAtComponentV1, to: NativeWorkspaceFsAtComponentV1, replace: boolean): void; + unlinkAt(parent: NativeWorkspaceFsAtHandleV1, name: NativeWorkspaceFsAtComponentV1): void; + listDirectory(handle: NativeWorkspaceFsAtHandleV1): readonly string[]; mkdirat(parent: NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1,mode:0o700):void; fstatat(parent:NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1):NativeWorkspaceFsAtStatV1; fsyncDirectory(handle:NativeWorkspaceFsAtHandleV1):void; diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index c908b060..0a396ddd 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -1,12 +1,9 @@ import { createHash, randomBytes } from "node:crypto"; -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, + workspaceFsAtFriend, } from "./workspace-fs-at.js"; import type { RuntimeConfigLease } from "./runtime-config-lease.js"; import { @@ -27,7 +24,6 @@ export interface PreprocessingRunStateV1 { } 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; } @@ -49,99 +45,57 @@ function id(v: string): void { if (typeof v !== "string" || !RUN_ID.test(v)) thr function checkSha(v: string): void { if (typeof v !== "string" || !SHA256.test(v)) throw fail("invalid digest"); } 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(); - return got.length === keys.length && got.every((key, i) => key === [...keys].sort()[i]); + const got = Object.keys(value as object).sort(); const expected = [...keys].sort(); + return got.length === expected.length && got.every((key, i) => key === expected[i]); } -function regular0600(path: string): void { - const st = lstatSync(path); - if (!st.isFile() || (st.mode & 0o777) !== 0o600 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink !== 1) throw fail(); +function validDir(directory: OwnedWorkspaceFsAtDirectory): void { + const st = directory.stat(); if ((st.mode & 0o170000) !== 0o040000 || (st.mode & 0o7777) !== 0o700 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink < 2n) throw fail(); } -async function fsyncParent(path: string): Promise { - const p = await openFile(dirname(path), "r"); - try { await p.sync(); } finally { await p.close(); } -} -async function durableJson(path: string, value: unknown): Promise { - const tmp = `${path}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`; - try { - await writeFile(tmp, `${JSON.stringify(value)}\n`, { mode: 0o600, flag: "wx" }); - regular0600(tmp); - const handle = await openFile(tmp, "r"); - try { await handle.sync(); } finally { await handle.close(); } - await rename(tmp, path); - await fsyncParent(path); - regular0600(path); - } catch (error) { await rm(tmp, { force: true }).catch(() => undefined); throw error; } +function validFile(file: OwnedWorkspaceFsAtRegularFile, max = MAX_FILE_BYTES): void { + const st = file.stat(); if ((st.mode & 0o170000) !== 0o100000 || (st.mode & 0o7777) !== 0o600 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink !== 1n || st.device === 0n || st.inode === 0n || max < 0) throw fail(); } +/** All state access is relative to the transferred retained root descriptor. */ export class PreprocessingStateStore { - private readonly rootLease: { assertLive(): void; anchoredPath(): string }; - constructor(rootLease: VerifiedWorkspaceLockRootLease | string) { - if (typeof rootLease === "string") { - let st: ReturnType; - try { st = lstatSync(rootLease); } catch { throw fail(); } - if (!st.isDirectory() || st.isSymbolicLink()) throw fail(); - this.rootLease = { assertLive: () => { const current = lstatSync(rootLease); if (!current.isDirectory() || current.isSymbolicLink()) throw fail(); }, anchoredPath: () => rootLease }; - } else this.rootLease = { assertLive: () => workspaceRootInternals.assertLive(rootLease), anchoredPath: () => { throw fail(); } }; + constructor(private readonly rootLease: VerifiedWorkspaceLockRootLease) {} + private dirs(action: (fs: WorkspaceFsAtV1, base: OwnedWorkspaceFsAtDirectory, jobs: OwnedWorkspaceFsAtDirectory, candidates: OwnedWorkspaceFsAtDirectory, reviews: OwnedWorkspaceFsAtDirectory) => T): T { + return workspaceRootInternals.withRoot(this.rootLease, (fs, root) => { + validDir(root); let base: OwnedWorkspaceFsAtDirectory | undefined; let jobs: OwnedWorkspaceFsAtDirectory | undefined; let candidates: OwnedWorkspaceFsAtDirectory | undefined; let reviews: OwnedWorkspaceFsAtDirectory | undefined; + const mkdir = (parent: OwnedWorkspaceFsAtDirectory, name: string): OwnedWorkspaceFsAtDirectory => { try { return workspaceFsAtFriend.openDirectory(parent, name); } catch { try { workspaceFsAtFriend.mkdir(parent, name); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; } return workspaceFsAtFriend.openDirectory(parent, name); } }; + try { base = mkdir(root, "preprocessing"); jobs = mkdir(base, "jobs"); candidates = mkdir(base, "fk-candidates"); reviews = mkdir(base, "fk-reviews"); [base, jobs, candidates, reviews].forEach(validDir); for (const d of [reviews, candidates, jobs, base, root]) workspaceFsAtFriend.fsyncDirectory(d); return action(fs, base, jobs, candidates, reviews); } + catch (error) { if (error instanceof Error && error.name === "PreprocessingConflictError") throw error; throw fail(); } + finally { for (const d of [reviews, candidates, jobs, base]) { try { d?.close(); } catch {} } } + }); } - 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`) }; + private file(directory: OwnedWorkspaceFsAtDirectory, name: string, access: "read"|"create"): OwnedWorkspaceFsAtRegularFile { + try { const f = workspaceFsAtFriend.openFile(directory, name, access); validFile(f); return f; } catch { throw fail(); } } - private async ensureDirs(p: ReturnType): Promise { - for (const path of [p.base, p.jobs, dirname(p.candidate), dirname(p.review)]) { - try { await mkdir(path, { mode: 0o700 }); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw fail(); } - const st = lstatSync(path); if (!st.isDirectory() || (st.mode & 0o777) !== 0o700 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink < 2) throw fail(); - } + private read(directory: OwnedWorkspaceFsAtDirectory, name: string): Uint8Array { + const f = this.file(directory, name, "read"); try { const bytes = workspaceFsAtFriend.readFile(f, MAX_FILE_BYTES); if (bytes.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds"); return bytes; } catch { throw fail(); } finally { try { f.close(); } catch {} } } - private assertFile(path: string, max = MAX_FILE_BYTES): void { regular0600(path); if (lstatSync(path).size > max) throw fail("preprocessing_bounds"); } - private async assertBounds(jobs: string): Promise { - let entries: string[]; try { entries = await readdir(jobs); } catch { throw fail(); } - if (entries.length > MAX_ENTRIES) throw fail("preprocessing_bounds"); let total = 0; - for (const name of entries) { if (!RUN_ID.test(name.replace(/\.json$/, "")) || !name.endsWith(".json")) throw fail("preprocessing_bounds"); const path = join(jobs, name); const st = lstatSync(path); if (!st.isFile() || st.nlink !== 1 || (st.mode & 0o777) !== 0o600) throw fail("preprocessing_bounds"); if (st.size > MAX_FILE_BYTES || (total += st.size) > MAX_AGGREGATE_BYTES) throw fail("preprocessing_bounds"); } + private writeExclusive(directory: OwnedWorkspaceFsAtDirectory, name: string, bytes: Uint8Array): void { + if (!(bytes instanceof Uint8Array) || bytes.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds"); const f = this.file(directory, name, "create"); + try { workspaceFsAtFriend.writeFile(f, bytes); workspaceFsAtFriend.fsyncFile(f); } catch { throw fail(); } finally { try { f.close(); } catch {} } + workspaceFsAtFriend.fsyncDirectory(directory); } + private writeReplace(directory: OwnedWorkspaceFsAtDirectory, name: string, bytes: Uint8Array): void { + const tmp = `.${name}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`; try { this.writeExclusive(directory, tmp, bytes); workspaceFsAtFriend.rename(directory, tmp, name, true); workspaceFsAtFriend.fsyncDirectory(directory); } catch { try { workspaceFsAtFriend.unlink(directory, tmp); } catch {} throw fail(); } + } + private rawState(jobs: OwnedWorkspaceFsAtDirectory, runId: string): PreprocessingRunStateV1 { id(runId); const names = workspaceFsAtFriend.listDirectory(jobs); if (names.length > MAX_ENTRIES) throw fail("preprocessing_bounds"); let total = 0; for (const name of names) { if (!name.endsWith(".json") || !RUN_ID.test(name.slice(0, -5))) throw fail("preprocessing_bounds"); const bytes = this.read(jobs, name); total += bytes.byteLength; if (bytes.byteLength > MAX_FILE_BYTES || total > MAX_AGGREGATE_BYTES) throw fail("preprocessing_bounds"); } let parsed: unknown; try { parsed = JSON.parse(new TextDecoder().decode(this.read(jobs, `${runId}.json`))); } catch { throw fail(); } if (!this.validState(parsed)) throw fail(); return parsed; } + private rawReview(reviews: OwnedWorkspaceFsAtDirectory, runId: string): FkReviewRecordV1 | undefined { try { const parsed = JSON.parse(new TextDecoder().decode(this.read(reviews, `${runId}.json`))); if (!strictObject(parsed, ["candidate", "reviewSha256", "runId", "recordedAt", ...(parsed && "annotationSha256" in parsed ? ["annotationSha256"] : [])])) throw fail(); return parsed as unknown as FkReviewRecordV1; } catch { return undefined; } } async create(input: CreateRunInput): Promise { - const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); const p = this.paths({ runId }); - if (!WORKSPACE.test(input.workspaceId) || !REVISION.test(input.revision) || !OPERATIONS.has(input.operation)) throw fail(); - await this.ensureDirs(p); const now = new Date().toISOString(); - const state: PreprocessingRunStateV1 = { schemaVersion: 1, runId, workspaceId: input.workspaceId, revision: input.revision, operation: input.operation, phase: "created", artifacts: [], createdAt: now, updatedAt: now }; - try { await durableJson(p.path, state); } catch { throw fail(); } - return state; - } - async loadForResume(input: ResumeRunInput): Promise { - const p = this.paths(input); let value: unknown; - try { await this.assertBounds(p.jobs); this.assertFile(p.path); value = JSON.parse(await readFile(p.path, "utf8")); } catch { throw fail(); } - if (!this.validState(value) || value.runId !== input.runId || value.workspaceId !== input.workspaceId || value.revision !== input.revision || value.operation !== input.operation) throw fail("preprocessing_resume_mismatch"); - return value; + const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); if (!WORKSPACE.test(input.workspaceId) || !REVISION.test(input.revision) || !OPERATIONS.has(input.operation)) throw fail(); + return this.dirs((_fs, _base, jobs) => { const now = new Date().toISOString(); const state: PreprocessingRunStateV1 = { schemaVersion: 1, runId, workspaceId: input.workspaceId, revision: input.revision, operation: input.operation, phase: "created", artifacts: [], createdAt: now, updatedAt: now }; const bytes = new TextEncoder().encode(`${JSON.stringify(state)} +`); try { this.writeExclusive(jobs, `${runId}.json`, bytes); return state; } catch (error) { try { const prior = this.rawState(jobs, runId); if (prior.workspaceId === input.workspaceId && prior.revision === input.revision && prior.operation === input.operation && prior.runId === runId) return prior; } catch {} throw fail(); } }); } + async loadForResume(input: ResumeRunInput): Promise { id(input.runId); if (!WORKSPACE.test(input.workspaceId) || !REVISION.test(input.revision) || !OPERATIONS.has(input.operation)) throw fail("preprocessing_resume_mismatch"); return this.dirs((_fs, _base, jobs) => { const state = this.rawState(jobs, input.runId); if (state.workspaceId !== input.workspaceId || state.revision !== input.revision || state.operation !== input.operation) throw fail("preprocessing_resume_mismatch"); return state; }); } async load(input: ResumeRunInput): Promise { return this.loadForResume(input); } - async transition(runId: string, transition: RunTransition): Promise { - const p = this.paths({ runId }); const current = await this.readRaw(p.path); - if (!this.validState(current) || current.runId !== runId || !strictObject(transition, ["phase"]) || typeof transition.phase !== "string") throw fail(); - const old = PHASES.indexOf(current.phase as never); const next = PHASES.indexOf(transition.phase as never); - if (next < 0 || old < 0 || next < old || current.phase === "terminal") throw fail(); - const updated: PreprocessingRunStateV1 = { ...current, phase: transition.phase, updatedAt: new Date().toISOString() }; - await durableJson(p.path, updated); return updated; - } - async writeFkCandidate(runId: string, yaml: Uint8Array): Promise { - const p = this.paths({ runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds"); - await this.ensureDirs(p); const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity; - try { await writeFile(p.candidate, yaml, { flag: "wx", mode: 0o600 }); this.assertFile(p.candidate); const h = await openFile(p.candidate, "r"); try { await h.sync(); } finally { await h.close(); } await fsyncParent(p.candidate); } catch { throw fail(); } - return artifact; - } - async recordFkReview(runId: string, review: FkReviewInput): Promise { - const p = this.paths({ runId }); - if (!strictObject(review, ["candidate", "reviewSha256", ...(review && "annotationSha256" in review ? ["annotationSha256"] : [])]) || !review?.candidate || !strictObject(review.candidate, ["kind", "digest", "bytes"]) || review.candidate.kind !== "fk-candidate" || !Number.isSafeInteger(review.candidate.bytes) || review.candidate.bytes < 0 || review.candidate.bytes > MAX_FILE_BYTES) throw fail(); - checkSha(review.reviewSha256); if (review.annotationSha256 !== undefined) checkSha(review.annotationSha256); - const out: FkReviewRecordV1 = { ...review, runId, recordedAt: new Date().toISOString() }; await this.ensureDirs(p); await durableJson(p.review, out); return out; - } - private async readRaw(path: string): Promise { try { this.assertFile(path); return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } } - private validState(value: unknown): value is PreprocessingRunStateV1 { - if (!strictObject(value, ["schemaVersion", "runId", "workspaceId", "revision", "operation", "phase", "artifacts", "createdAt", "updatedAt"])) return false; - const x = value as Record; if (x.schemaVersion !== 1 || typeof x.runId !== "string" || !RUN_ID.test(x.runId as string) || typeof x.workspaceId !== "string" || !WORKSPACE.test(x.workspaceId as string) || typeof x.revision !== "string" || !REVISION.test(x.revision as string) || typeof x.operation !== "string" || !OPERATIONS.has(x.operation as string) || typeof x.phase !== "string" || !PHASES.includes(x.phase as never) || !Array.isArray(x.artifacts) || typeof x.createdAt !== "string" || typeof x.updatedAt !== "string") return false; - if (x.artifacts.length > MAX_ENTRIES) return false; - return (x.artifacts as unknown[]).every(a => strictObject(a, ["kind", "digest", "bytes"]) && typeof (a as Record).kind === "string" && typeof (a as Record).digest === "string" && SHA256.test((a as Record).digest as string) && Number.isSafeInteger((a as Record).bytes) && ((a as Record).bytes as number) >= 0 && ((a as Record).bytes as number) <= MAX_FILE_BYTES); - } + async transition(runId: string, transition: RunTransition): Promise { id(runId); if (!strictObject(transition, ["phase"]) || typeof transition.phase !== "string") throw fail(); return this.dirs((_fs, _base, jobs) => { const current = this.rawState(jobs, runId); const old = PHASES.indexOf(current.phase as never); const next = PHASES.indexOf(transition.phase as never); if (old < 0 || next < 0 || next < old || current.phase === "terminal") throw fail(); const updated: PreprocessingRunStateV1 = { ...current, phase: transition.phase, updatedAt: new Date().toISOString() }; this.writeReplace(jobs, `${runId}.json`, new TextEncoder().encode(`${JSON.stringify(updated)} +`)); return updated; }); } + async writeFkCandidate(runId: string, yaml: Uint8Array): Promise { id(runId); if (!(yaml instanceof Uint8Array) || yaml.byteLength > MAX_FILE_BYTES) throw fail("preprocessing_bounds"); return this.dirs((_fs, _base, jobs, candidates) => { const state = this.rawState(jobs, runId); const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity; try { this.writeExclusive(candidates, `${runId}.yaml`, yaml); } catch { throw fail(); } if (state.runId !== runId) throw fail(); return artifact; }); } + async recordFkReview(runId: string, review: FkReviewInput): Promise { id(runId); if (!strictObject(review, ["candidate", "reviewSha256", ...(review && "annotationSha256" in review ? ["annotationSha256"] : [])]) || !review?.candidate || !strictObject(review.candidate, ["kind", "digest", "bytes"]) || review.candidate.kind !== "fk-candidate" || !Number.isSafeInteger(review.candidate.bytes) || review.candidate.bytes < 0 || review.candidate.bytes > MAX_FILE_BYTES) throw fail(); checkSha(review.reviewSha256); if (review.annotationSha256 !== undefined) checkSha(review.annotationSha256); return this.dirs((_fs, _base, jobs, candidates, reviews) => { this.rawState(jobs, runId); let candidate: Uint8Array; try { candidate = this.read(candidates, `${runId}.yaml`); } catch { throw fail(); } if (review.candidate.bytes !== candidate.byteLength || review.candidate.digest !== digest(candidate)) throw fail("preprocessing_review_mismatch"); const prior = this.rawReview(reviews, runId); if (prior) { const same = prior.runId === runId && JSON.stringify(prior.candidate) === JSON.stringify(review.candidate) && prior.reviewSha256 === review.reviewSha256 && prior.annotationSha256 === review.annotationSha256; if (!same) throw fail("preprocessing_review_mismatch"); return prior; } const out: FkReviewRecordV1 = { ...review, runId, recordedAt: new Date().toISOString() }; try { this.writeExclusive(reviews, `${runId}.json`, new TextEncoder().encode(`${JSON.stringify(out)} +`)); } catch { throw fail(); } return out; }); } + private validState(value: unknown): value is PreprocessingRunStateV1 { if (!strictObject(value, ["schemaVersion", "runId", "workspaceId", "revision", "operation", "phase", "artifacts", "createdAt", "updatedAt"])) return false; const x = value as Record; if (x.schemaVersion !== 1 || typeof x.runId !== "string" || !RUN_ID.test(x.runId) || typeof x.workspaceId !== "string" || !WORKSPACE.test(x.workspaceId) || typeof x.revision !== "string" || !REVISION.test(x.revision) || typeof x.operation !== "string" || !OPERATIONS.has(x.operation) || typeof x.phase !== "string" || !PHASES.includes(x.phase as never) || !Array.isArray(x.artifacts) || typeof x.createdAt !== "string" || typeof x.updatedAt !== "string" || x.artifacts.length > MAX_ENTRIES) return false; return (x.artifacts as unknown[]).every(a => strictObject(a, ["kind", "digest", "bytes"]) && typeof (a as Record).kind === "string" && typeof (a as Record).digest === "string" && SHA256.test((a as Record).digest as string) && Number.isSafeInteger((a as Record).bytes) && (a as Record).bytes as number >= 0 && (a as Record).bytes as number <= MAX_FILE_BYTES); } } 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; } @@ -150,60 +104,62 @@ export interface EvidenceLockedChildRequest { readonly kind: "evidence_preproces export type WorkspaceLockedChildRequest = DwhLockedChildRequest | SchemaLockedChildRequest | EvidenceLockedChildRequest; export interface WorkspaceLockedChildResult { readonly exitCode: number; readonly stdout: Uint8Array; readonly stderr: Uint8Array; } +const borrowedState = new WeakMap(); export class BorrowedWorkspaceSessionReadersExclusiveLockLease { - private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} + private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) { borrowedState.set(this, { live: true }); } static [INTERNAL_STATE](id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } - assertLive(): void { if (!this.live) throw fail(); } - invalidate(): void { this.live = false; } } +function invalidateBorrowed(value: BorrowedWorkspaceSessionReadersExclusiveLockLease): void { const state = borrowedState.get(value); if (!state) throw fail(); state.live = false; } + +interface WriterState { live: boolean; settled: boolean; readerExclusive: boolean; spawnActive: boolean; poisoned: boolean; root: VerifiedWorkspaceLockRootLease; writer: WorkspaceRootLock; } +const writerState = new WeakMap(); 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(); workspaceRootInternals.assertLive(this.root); this.writer.assertPath(); } - invalidateForSettlement(): void { this.settled = true; } + readonly workspaceId: CanonicalWorkspaceId; + readonly rootIdentity: WorkspaceLockRootIdentityV1; + private constructor(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { + this.workspaceId = id; this.rootIdentity = identity; writerState.set(this, { live: true, settled: false, readerExclusive: false, spawnActive: false, poisoned: false, 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 workspaceRootInternals.acquireReadersExclusive(this.root); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); } - 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; } - this.readerExclusive = false; - if (cleanupError) { this.poisoned = true; this.live = false; throw fail(); } - } + const state = assertCapability(this); if (state.readerExclusive || state.spawnActive) throw fail(); state.readerExclusive = true; let lock: WorkspaceRootLock; + try { lock = await workspaceRootInternals.acquireReadersExclusive(state.root); } catch { state.readerExclusive = false; state.poisoned = true; throw fail(); } + const borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease[INTERNAL_STATE](this.workspaceId, this.rootIdentity); let callbackError: unknown; let result: T | undefined; + try { result = await action(borrowed); } catch (error) { callbackError = error; } + invalidateBorrowed(borrowed); let cleanupError: unknown; try { lock.close(); } catch (error) { cleanupError = error; } + state.readerExclusive = false; if (cleanupError) { state.poisoned = true; state.live = false; } + if (callbackError) throw callbackError; if (cleanupError) throw fail(); return result as T; } async spawnChild(request: WorkspaceLockedChildRequest): Promise { - this.assertLive(); if (!this.readerExclusive || this.spawnActive || !request || request.workspaceId !== this.workspaceId) throw fail(); + const state = assertCapability(this); if (!state.readerExclusive || state.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; - if (cfg.workspaceId !== this.workspaceId || cfg.workspaceRevision !== request.revision || typeof cfg.path !== "string") throw fail(); - this.spawnActive = true; + const cfg = request.runtimeConfig; if (cfg.workspaceId !== this.workspaceId || cfg.workspaceRevision !== request.revision || typeof cfg.path !== "string") throw fail(); state.spawnActive = true; try { - 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 workspaceRootInternals.spawn(this.root, 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; } + let argv: string[]; + switch (request.kind) { case "dwh_preprocess": argv = ["-m", "tht.cli", "preprocess", "dwh", "--steps", request.stage, "--json", "-c", cfg.path]; break; case "schema_preprocess": argv = ["-m", "tht.cli", "schema", request.stage === "fk_suggest" ? "suggest-fks" : request.stage === "fk_check" ? "check" : "index", "--json", "-c", cfg.path]; break; case "evidence_preprocess": argv = ["-m", "tht.cli", "preprocess", "evidence", "--json", "-c", cfg.path]; break; default: throw fail(); } + return await workspaceRootInternals.withRoot(state.root, (_fs, retainedRoot) => state.writer.spawn(retainedRoot, 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 { state.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 assertCapability(cap: WorkspaceWriterLockCapability): WriterState { const state = writerState.get(cap); if (!state || !state.live || state.settled || state.poisoned) throw fail(); workspaceRootInternals.assertLive(state.root); state.writer.assertPath(); return state; } +function settleCapability(cap: WorkspaceWriterLockCapability): void { const state = writerState.get(cap); if (state) state.settled = true; } +async function closeCapability(cap: WorkspaceWriterLockCapability): Promise { const state = writerState.get(cap); if (!state || !state.live) return; if (state.readerExclusive || state.spawnActive) { state.poisoned = true; throw fail(); } state.live = false; let error: unknown; try { state.writer.close(); } catch (e) { error = e; } try { await state.root.close(); } catch (e) { error ??= e; } if (error) throw fail(); } 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; } +const orderedState = new WeakMap }>(); export class OrderedWorkspaceWriterCapabilitySet { - private live = true; private constructor(private readonly caps: Map) {} + private constructor(caps: Map) { orderedState.set(this, { live: true, 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: makeBorrowedRootLease(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); } + get workspaceIds(): readonly CanonicalWorkspaceId[] { const state = orderedState.get(this); if (!state?.live) throw fail(); return [...state.caps.keys()]; } + async forWorkspace(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise): Promise { const state = orderedState.get(this); if (!state?.live) throw fail(); const cap = state.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: makeBorrowedRootLease(cap.rootIdentity, () => { if (!orderedState.get(this)?.live) throw fail(); }), writerCapability: cap }); } async forEachWorkspace(action: (lease: OrderedWorkspaceCapability) => Promise): Promise { return Promise.all(this.workspaceIds.map(id => this.forWorkspace(id, action))); } } +function settleSet(set: OrderedWorkspaceWriterCapabilitySet): void { const state = orderedState.get(set); if (!state) return; state.live = false; for (const cap of state.caps.values()) settleCapability(cap); } export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise): Promise { const sorted = [...rootLeases].sort((a, b) => a.identity.workspaceId.localeCompare(b.identity.workspaceId)); if (new Set(sorted.map(x => x.identity.workspaceId)).size !== sorted.length) throw fail(); const caps: WorkspaceWriterLockCapability[] = []; let set: OrderedWorkspaceWriterCapabilitySet | undefined; let result: T | undefined; let callbackError: unknown; @@ -223,14 +179,14 @@ export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonl } if (!callbackError) { set = OrderedWorkspaceWriterCapabilitySet[INTERNAL_STATE](new Map(caps.map(c => [c.workspaceId, c]))); - for (const capability of caps) capability.assertWriterPath(); - try { result = await action(set); for (const capability of caps) capability.assertWriterPath(); } + for (const capability of caps) assertCapability(capability); + try { result = await action(set); for (const capability of caps) assertCapability(capability); } catch (error) { callbackError = error; } } } catch (error) { callbackError ??= error; } - set?.invalidate(); + if (set) settleSet(set); for (const cap of [...caps].reverse()) { - try { await cap.close(); } catch (error) { callbackError ??= error; } + try { await closeCapability(cap); } catch (error) { callbackError ??= error; } } if (callbackError) throw callbackError; return result as T; } diff --git a/backend/src/workspaces/workspace-fs-at.ts b/backend/src/workspaces/workspace-fs-at.ts index c6cff631..e59fc7a9 100644 --- a/backend/src/workspaces/workspace-fs-at.ts +++ b/backend/src/workspaces/workspace-fs-at.ts @@ -82,9 +82,33 @@ export class WorkspaceFsAtV1 { export type WorkspaceRootFriend = { assertPath(root: OwnedWorkspaceFsAtDirectory, lock: OwnedWorkspaceFsAtRegularFile, name: LockFileName, identity: WorkspaceLockRootIdentityV1Like): void; spawn(writer: OwnedWorkspaceFsAtRegularFile, root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv): Promise<{ exitCode: number; stdout: Uint8Array; stderr: Uint8Array }>; + openDirectory(parent: OwnedWorkspaceFsAtDirectory, name: string): OwnedWorkspaceFsAtDirectory; + mkdir(parent: OwnedWorkspaceFsAtDirectory, name: string): void; + openFile(parent: OwnedWorkspaceFsAtDirectory, name: string, access: "read"|"create"): OwnedWorkspaceFsAtRegularFile; + readFile(file: OwnedWorkspaceFsAtRegularFile, maxBytes: number): Uint8Array; + writeFile(file: OwnedWorkspaceFsAtRegularFile, bytes: Uint8Array): void; + fsyncFile(file: OwnedWorkspaceFsAtRegularFile): void; + rename(parent: OwnedWorkspaceFsAtDirectory, from: string, to: string, replace: boolean): void; + unlink(parent: OwnedWorkspaceFsAtDirectory, name: string): void; + listDirectory(directory: OwnedWorkspaceFsAtDirectory): readonly string[]; + fsyncDirectory(directory: OwnedWorkspaceFsAtDirectory): void; }; export type WorkspaceLockRootIdentityV1Like = { readonly device: bigint; readonly inode: bigint }; export const workspaceFsAtFriend: WorkspaceRootFriend = { + openDirectory(parent, name) { + return wrapDirectory(binding.openat({ parent: rawDirectory(parent), name: component(name), kind: "directory", createMode: 0 })); + }, + mkdir(parent, name) { binding.mkdirat(rawDirectory(parent), component(name), 0o700); }, + openFile(parent, name, access) { + return wrapLock(binding.openFileAt({ parent: rawDirectory(parent), name: component(name), access })); + }, + readFile(file, maxBytes) { if (!rawHandles.has(file)) throw new Error("workspace descriptor is closed"); return binding.readFile(rawHandles.get(file)!, maxBytes); }, + writeFile(file, bytes) { if (!rawHandles.has(file)) throw new Error("workspace descriptor is closed"); binding.writeFile(rawHandles.get(file)! , bytes); }, + fsyncFile(file) { if (!rawHandles.has(file)) throw new Error("workspace descriptor is closed"); binding.fsyncFile(rawHandles.get(file)!); }, + rename(parent, from, to, replace) { binding.renameAt(rawDirectory(parent), component(from), component(to), replace); }, + fsyncDirectory(directory) { binding.fsyncDirectory(rawDirectory(directory)); }, + unlink(parent, name) { binding.unlinkAt(rawDirectory(parent), component(name)); }, + listDirectory(directory) { return binding.listDirectory(rawDirectory(directory)); }, assertPath(root, lock, name, identity) { const current = fsStatAtNoFollow(root, name); const opened = lock.stat(); if (current.device !== opened.device || current.inode !== opened.inode || current.mode !== opened.mode || current.uid !== opened.uid || current.nlink !== opened.nlink) throw new Error("preprocessing_conflict: lock pathname identity changed"); diff --git a/backend/src/workspaces/workspace-lock-root-lease.ts b/backend/src/workspaces/workspace-lock-root-lease.ts index 347f538c..0b0a340a 100644 --- a/backend/src/workspaces/workspace-lock-root-lease.ts +++ b/backend/src/workspaces/workspace-lock-root-lease.ts @@ -1,4 +1,4 @@ -import { resolve } from "node:path"; +import { lstatSync, realpathSync } from "node:fs"; import { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, OwnedWorkspaceFsAtRegularFile, workspaceFsAtFriend, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js"; declare const canonicalWorkspaceIdBrand: unique symbol; declare const revision40Brand: unique symbol; @@ -17,15 +17,15 @@ function exactParent(stat: WorkspaceFsAtStatV1, uid: number): boolean { return e function sameIdentity(a: {device: bigint|number; inode: bigint|number}, b: {device: bigint; inode: bigint}): boolean { return BigInt(a.device) === b.device && BigInt(a.inode) === b.inode; } type InternalLock = { 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 }>; }; +const writerAdmissions = new Set(); class RootLock implements InternalLock { private live = true; - constructor(private readonly root: OwnedWorkspaceFsAtDirectory, private readonly lock: OwnedWorkspaceFsAtRegularFile, private readonly name: "writer.lock"|"session-readers.lock", private readonly identity: WorkspaceLockRootIdentityV1, private readonly fs: WorkspaceFsAtV1, private readonly activeKey?: string) {} + constructor(private readonly root: OwnedWorkspaceFsAtDirectory, private readonly lock: OwnedWorkspaceFsAtRegularFile, private readonly name: "writer.lock"|"session-readers.lock", private readonly identity: WorkspaceLockRootIdentityV1, private readonly fs: WorkspaceFsAtV1, private readonly admissionKey?: string) {} assertPath(): void { if (!this.live) throw conflict(); workspaceFsAtFriend.assertPath(this.root, this.lock, this.name, this.identity); } flock(kind: WorkspaceFlockKindV1, wait: WorkspaceFlockWaitV1): void { this.assertPath(); this.fs.flockOwnedLock(this.lock, kind, wait); this.assertPath(); } spawn(root: OwnedWorkspaceFsAtDirectory, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return workspaceFsAtFriend.spawn(this.lock, root, executable, args, environment); } - close(): void { if (!this.live) return; this.live = false; try { this.lock.close(); } finally { if (this.activeKey) activeWriterLocks.delete(this.activeKey); } } + close(): void { if (!this.live) return; this.live = false; try { this.lock.close(); } finally { if (this.admissionKey) writerAdmissions.delete(this.admissionKey); } } } -const activeWriterLocks = new Set(); type RootState = { fs: WorkspaceFsAtV1; parent: OwnedWorkspaceFsAtDirectory; root: OwnedWorkspaceFsAtDirectory; serviceUid: number; owner: symbol; identity: WorkspaceLockRootIdentityV1; live: boolean }; const rootState = new WeakMap(); @@ -59,13 +59,12 @@ export class VerifiedWorkspaceLockRootLease { try { const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); lock = new RootLock(this.root, owned, "session-readers.lock", this.identity, this.fs); lock.flock("shared", "nonblocking"); assertRootLive(this); rootState.get(this)!.live = false; return Reflect.construct(WorkspaceSessionReadersLockLease, [lock, this.root, this.identity]) as WorkspaceSessionReadersLockLease; } catch (e) { try { lock?.close(); } catch { this.live = false; rootState.get(this)!.live = false; try { this.root.close(); } catch {} } throw conflict(); } } - transfer(): VerifiedWorkspaceLockRootLease { assertRootLive(this); if (this.borrowed) throw conflict("workspace root is borrowed"); this.live = false; rootState.get(this)!.live = false; return legacyLease(Reflect.construct(VerifiedWorkspaceLockRootLease, [this.fs, this.parent, this.root, this.identity, this.serviceUid, this.owner]) as VerifiedWorkspaceLockRootLease); } + transfer(): VerifiedWorkspaceLockRootLease { assertRootLive(this); if (this.borrowed) throw conflict("workspace root is borrowed"); this.live = false; rootState.get(this)!.live = false; return Reflect.construct(VerifiedWorkspaceLockRootLease, [this.fs, this.parent, this.root, this.identity, this.serviceUid, this.owner]) as VerifiedWorkspaceLockRootLease; } async close(): Promise { if (!this.live) return; while (this.borrowed) await new Promise(r => setTimeout(r, 1)); this.live = false; rootState.get(this)!.live = false; try { this.root.close(); } catch { throw conflict(); } } } -function legacyLease(raw: VerifiedWorkspaceLockRootLease): VerifiedWorkspaceLockRootLease { const proxy = new Proxy(raw, { get(target, property, receiver) { if (property === "assertLive") return () => workspaceRootInternals.assertLive(proxy); if (property === "acquireWriterLock") return () => workspaceRootInternals.acquireWriter(proxy); if (property === "acquireSessionReadersExclusive") return () => workspaceRootInternals.acquireReadersExclusive(proxy); if (property === "spawnChild") return (lock: InternalLock, executable: string, args: readonly string[], env?: NodeJS.ProcessEnv) => workspaceRootInternals.spawn(proxy, lock, executable, args, env); if (property === "fsync") return () => undefined; if (property === "anchoredPath") return () => { throw conflict(); }; return Reflect.get(target, property, receiver); } }); const state = rootState.get(raw)!; rootState.set(proxy, state); return proxy; } -function makeRoot(fs: WorkspaceFsAtV1, parent: OwnedWorkspaceFsAtDirectory, root: OwnedWorkspaceFsAtDirectory, id: WorkspaceLockRootIdentityV1, uid: number, owner: symbol): VerifiedWorkspaceLockRootLease { return legacyLease(Reflect.construct(VerifiedWorkspaceLockRootLease, [fs, parent, root, id, uid, owner]) as VerifiedWorkspaceLockRootLease); } +function makeRoot(fs: WorkspaceFsAtV1, parent: OwnedWorkspaceFsAtDirectory, root: OwnedWorkspaceFsAtDirectory, id: WorkspaceLockRootIdentityV1, uid: number, owner: symbol): VerifiedWorkspaceLockRootLease { return Reflect.construct(VerifiedWorkspaceLockRootLease, [fs, parent, root, id, uid, owner]) as VerifiedWorkspaceLockRootLease; } export class VerifiedWorkspaceLockRootLeaseFactory { private readonly owner = Symbol("workspace-root-factory"); private readonly parent: OwnedWorkspaceFsAtDirectory; private provisionTail: Promise = Promise.resolve(); @@ -73,8 +72,9 @@ export class VerifiedWorkspaceLockRootLeaseFactory { if (!Number.isInteger(input.serviceUid) || input.serviceUid < 0 || input.provisionedWorkspaceMode !== 0o700) throw new Error("invalid workspace root policy"); const configured = input.sessionsRootFromValidatedInstallationConfig; if (!configured.startsWith("/") || configured.split("/").some(c => c === "" ? false : c === "." || c === "..")) throw conflict("invalid sessions root"); + let canonical: string; try { const leaf = lstatSync(configured); if (leaf.isSymbolicLink()) throw conflict("invalid sessions root"); canonical = realpathSync(configured); } catch (e) { if ((e as Error).name === "PreprocessingConflictError") throw e; throw conflict("invalid sessions root"); } let d = input.workspaceFsAt.openRoot(); - try { for (const c of configured.split("/").filter(Boolean)) { const n = input.workspaceFsAt.openDirectoryAt(d, c); d.close(); d = n; } this.parent = d; this.checkParent(); } + try { for (const c of canonical.split("/").filter(Boolean)) { const n = input.workspaceFsAt.openDirectoryAt(d, c); d.close(); d = n; } this.parent = d; this.checkParent(); } catch (e) { try { d.close(); } catch {} throw conflict("invalid sessions root"); } } private checkParent(): void { try { if (!exactParent(this.parent.stat(), this.input.serviceUid)) throw conflict("installation root identity changed"); } catch (e) { if ((e as Error).name === "PreprocessingConflictError") throw e; throw conflict("installation root identity changed"); } } @@ -88,7 +88,7 @@ export class VerifiedWorkspaceLockRootLeaseFactory { let release!: () => void; const prior = this.provisionTail; this.provisionTail = new Promise(r => { release = r; }); await prior; try { this.checkParent(); try { return await this.open(input); } catch (e) { if ((e as {code?: string}).code !== "ENOENT") throw e; } this.checkParent(); try { this.input.workspaceFsAt.mkdirAt(this.parent, input.workspaceId, 0o700); } catch (e) { if ((e as {code?: string}).code !== "EEXIST") throw conflict(); } - this.checkParent(); const winner = await this.open(input); try { this.input.workspaceFsAt.fsyncDirectory((winner as unknown as {root: OwnedWorkspaceFsAtDirectory}).root); } catch { await winner.close().catch(() => undefined); throw conflict(); } this.input.workspaceFsAt.fsyncDirectory(this.parent); return winner; + this.checkParent(); const winner = await this.open(input); try { this.input.workspaceFsAt.fsyncDirectory(rootState.get(winner)!.root); } catch { await winner.close().catch(() => undefined); throw conflict(); } this.input.workspaceFsAt.fsyncDirectory(this.parent); return winner; } finally { release(); } } } @@ -97,7 +97,9 @@ export class VerifiedWorkspaceLockRootLeaseFactory { // identity, borrow, shared reader acquisition, transfer and close. export const workspaceRootInternals = { assertLive(root: VerifiedWorkspaceLockRootLease): void { assertRootLive(root) }, - async acquireWriter(root: VerifiedWorkspaceLockRootLease): Promise { workspaceRootInternals.assertLive(root); const x = root as unknown as { fs: WorkspaceFsAtV1; root: OwnedWorkspaceFsAtDirectory; serviceUid: number; identity: WorkspaceLockRootIdentityV1 }; const key = `${x.identity.device}:${x.identity.inode}`; if (activeWriterLocks.has(key)) return Promise.reject(conflict()); const owned = x.fs.openOrCreateLockAt(x.root, "writer.lock", 0o600); const lock = new RootLock(x.root, owned, "writer.lock", x.identity, x.fs, key); try { lock.flock("exclusive", "nonblocking"); activeWriterLocks.add(key); return Promise.resolve(lock); } catch { try { lock.close(); } catch {} return Promise.reject(conflict()); } }, + withRoot(root: VerifiedWorkspaceLockRootLease, action: (fs: WorkspaceFsAtV1, directory: OwnedWorkspaceFsAtDirectory) => T): T { + assertRootLive(root); const state = rootState.get(root); if (!state) throw conflict(); return action(state.fs, state.root); + }, + async acquireWriter(root: VerifiedWorkspaceLockRootLease): Promise { workspaceRootInternals.assertLive(root); const x = root as unknown as { fs: WorkspaceFsAtV1; root: OwnedWorkspaceFsAtDirectory; identity: WorkspaceLockRootIdentityV1 }; const key = `${x.identity.device}:${x.identity.inode}`; if (writerAdmissions.has(key)) throw conflict(); const owned = x.fs.openOrCreateLockAt(x.root, "writer.lock", 0o600); const lock = new RootLock(x.root, owned, "writer.lock", x.identity, x.fs, key); try { lock.flock("exclusive", "nonblocking"); writerAdmissions.add(key); return Promise.resolve(lock); } catch { try { lock.close(); } catch {} return Promise.reject(conflict()); } }, async acquireReadersExclusive(root: VerifiedWorkspaceLockRootLease): Promise { workspaceRootInternals.assertLive(root); const x = root as unknown as { fs: WorkspaceFsAtV1; root: OwnedWorkspaceFsAtDirectory; identity: WorkspaceLockRootIdentityV1 }; const owned = x.fs.openOrCreateLockAt(x.root, "session-readers.lock", 0o600); const lock = new RootLock(x.root, owned, "session-readers.lock", x.identity, x.fs); try { lock.flock("exclusive", "nonblocking"); return Promise.resolve(lock); } catch { try { lock.close(); } catch {} return Promise.reject(conflict()); } }, - spawn(root: VerifiedWorkspaceLockRootLease, lock: InternalLock, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { workspaceRootInternals.assertLive(root); const x = root as unknown as { root: OwnedWorkspaceFsAtDirectory }; return lock.spawn(x.root, executable, args, environment); }, }; diff --git a/backend/test/workspace-lock-root-lease.test.ts b/backend/test/workspace-lock-root-lease.test.ts index af97f2c3..3397c97b 100644 --- a/backend/test/workspace-lock-root-lease.test.ts +++ b/backend/test/workspace-lock-root-lease.test.ts @@ -17,7 +17,7 @@ describe("retained canonical workspace root", () => { it("fails closed when the canonical pathname is replaced after retention", async () => { const parent = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(parent); const f = factory(parent); const lease = await f.acquireOrProvision(f.canonicalInput("abc-workspace")); renameSync(join(parent, "abc-workspace"), join(parent, "old")); mkdirSync(join(parent, "abc-workspace"), { mode: 0o700 }); - await expect(lease.acquireWriterLock()).rejects.toThrow(/preprocessing/); await lease.close(); + const { runUnderWorkspaceWriterLock } = await import("../src/workspaces/preprocessing-state.js"); await expect(runUnderWorkspaceWriterLock(lease, async () => undefined)).rejects.toThrow(/preprocessing/); await lease.close(); }); it("does not accept a symlink or wrong ownership/mode root", async () => { const parent = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(parent); const other = mkdtempSync(join(process.cwd(), "thoth-other-")); roots.push(other); symlinkSync(other, join(parent, "abc-workspace")); diff --git a/backend/test/workspace-preprocessing-state.test.ts b/backend/test/workspace-preprocessing-state.test.ts index 36c4f404..c63580ea 100644 --- a/backend/test/workspace-preprocessing-state.test.ts +++ b/backend/test/workspace-preprocessing-state.test.ts @@ -1,35 +1,16 @@ import { describe, expect, it, afterEach } from "vitest"; -import { mkdtemp, readFile, chmod, writeFile, symlink, lstat } from "node:fs/promises"; +import { mkdtempSync, renameSync, mkdirSync, rmSync, readFileSync, chmodSync, writeFileSync } from "node:fs"; import { join } from "node:path"; -import { tmpdir } from "node:os"; +import { WorkspaceFsAtV1 } from "../src/workspaces/workspace-fs-at.js"; import { PreprocessingStateStore } from "../src/workspaces/preprocessing-state.js"; - +import { VerifiedWorkspaceLockRootLeaseFactory } from "../src/workspaces/workspace-lock-root-lease.js"; const roots: string[] = []; -afterEach(async () => { for (const root of roots.splice(0)) await import("node:fs/promises").then(fs => fs.rm(root, { recursive: true, force: true })); }); -async function makeStore() { const root = await mkdtemp(join(tmpdir(), "thoth-state-")); roots.push(root); return { root, store: new PreprocessingStateStore(root) }; } -const input = { workspaceId: "abc-workspace" as never, revision: "a".repeat(40), operation: "schema" }; - +afterEach(async () => { for (const root of roots.splice(0)) rmSync(root, { recursive: true, force: true }); }); +async function makeStore() { const parent = mkdtempSync(join(process.cwd(), "thoth-state-")); roots.push(parent); const factory = new VerifiedWorkspaceLockRootLeaseFactory({ workspaceFsAt: new WorkspaceFsAtV1(), installationId: "test", sessionsRootFromValidatedInstallationConfig: parent, serviceUid: process.getuid!(), provisionedWorkspaceMode: 0o700 }); const lease = await factory.acquireOrProvision(factory.canonicalInput("abc-workspace")); return { root: join(parent, "abc-workspace"), store: new PreprocessingStateStore(lease), lease }; } +const input = { workspaceId: "abc-workspace" as never, revision: "a".repeat(40) as never, operation: "schema" }; describe("durable preprocessing state", () => { - it("creates and replays exact state with immutable identity", async () => { - const { root, store } = await makeStore(); const state = await store.create(input); - const replay = await store.loadForResume({ ...input, runId: state.runId }); - expect(replay).toEqual(state); await expect(store.transition(state.runId, { phase: "introspected", workspaceId: "evil" } as never)).rejects.toThrow("preprocessing_conflict"); - }); - it("rejects tampered, traversal, mode and version state before returning it", async () => { - const { root, store } = await makeStore(); const state = await store.create(input); const path = join(root, "preprocessing", "jobs", `${state.runId}.json`); - const original = JSON.parse(await readFile(path, "utf8")); await writeFile(path, JSON.stringify({ ...original, schemaVersion: 9 })); - await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow(); - await writeFile(path, JSON.stringify(original)); await chmod(path, 0o644); await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow("preprocessing_conflict"); - await expect(store.load({ ...input, runId: "../" + state.runId })).rejects.toThrow(); - }); - it("writes candidate artifacts atomically with bounded bytes and immutable links", async () => { - const { root, store } = await makeStore(); const state = await store.create(input); const bytes = new TextEncoder().encode("tables: []\n"); - const artifact = await store.writeFkCandidate(state.runId, bytes); expect(artifact.bytes).toBe(bytes.byteLength); - const candidate = lstat(join(root, "preprocessing", "fk-candidates", `${state.runId}.yaml`)); expect((await candidate).nlink).toBe(1); - await expect(store.writeFkCandidate(state.runId, new Uint8Array(1 << 20))).rejects.toThrow("preprocessing_conflict"); - }); - it("rejects a symlinked state root", async () => { - const real = await mkdtemp(join(tmpdir(), "thoth-state-real-")); const link = join(tmpdir(), `thoth-state-link-${Date.now()}`); roots.push(real, link); await symlink(real, link); - expect(() => new PreprocessingStateStore(link)).toThrow("preprocessing_conflict"); - }); + it("creates and replays exact state with immutable identity", async () => { const { store, lease } = await makeStore(); const state = await store.create(input); expect(await store.create({ ...input, runId: state.runId })).toEqual(state); expect(await store.loadForResume({ ...input, runId: state.runId })).toEqual(state); await expect(store.transition(state.runId, { phase: "introspected", workspaceId: "evil" } as never)).rejects.toThrow("preprocessing_conflict"); await lease.close(); }); + it("rejects tampered, traversal, mode and version state before returning it", async () => { const { root, store, lease } = await makeStore(); const state = await store.create(input); const path = join(root, "preprocessing", "jobs", `${state.runId}.json`); const original = JSON.parse(readFileSync(path, "utf8")); writeFileSync(path, JSON.stringify({ ...original, schemaVersion: 9 })); await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow(); writeFileSync(path, JSON.stringify(original)); chmodSync(path, 0o644); await expect(store.load({ ...input, runId: state.runId })).rejects.toThrow("preprocessing_conflict"); await expect(store.load({ ...input, runId: "../" + state.runId })).rejects.toThrow(); await lease.close(); }); + it("writes candidate artifacts atomically and validates review identity", async () => { const { store, lease } = await makeStore(); const state = await store.create(input); const bytes = new TextEncoder().encode("tables: []\n"); const artifact = await store.writeFkCandidate(state.runId, bytes); expect(artifact.bytes).toBe(bytes.byteLength); await expect(store.recordFkReview(state.runId, { candidate: { ...artifact, digest: "not-a-digest" }, reviewSha256: "a".repeat(64) })).rejects.toThrow(); await lease.close(); }); + it("retained root replacement fails closed without writing replacement state", async () => { const { root, store, lease } = await makeStore(); renameSync(root, `${root}.old`); mkdirSync(root, { mode: 0o700 }); await expect(store.create(input)).rejects.toThrow("preprocessing_conflict"); expect(() => readFileSync(join(root, "preprocessing", "jobs"))).toThrow(); await lease.close().catch(() => undefined); }); }); diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index c0dacff2..ee60e1c3 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -60,8 +60,21 @@ def annotations_path(cfg: Config) -> Path: def refresh_catalog(cfg, *, dwh=None, output_path: Path | None = None): - _require_writer_capability() - """Run the existing catalog algorithm and persist its canonical output.""" + """Run the catalog algorithm only under the exact backend-bound root capability.""" + import os + import stat + cap = __import__("tht.workspace_writer_lock", fromlist=["require_workspace_writer_capability"]).require_workspace_writer_capability() + cfg_workspace = getattr(cfg, "_workspace_id", None) + cfg_revision = getattr(cfg, "_workspace_revision", None) + runtime = getattr(cfg, "runtime_identity", None) + if cfg_workspace != cap.workspace_id or cfg_revision != cap.revision or runtime is None or runtime.workspace_id != cap.workspace_id or runtime.workspace_revision != cap.revision: + raise RuntimeError("preprocessing_conflict") + try: + st = os.stat(cfg.paths.sessions, follow_symlinks=False) + except OSError as exc: + raise RuntimeError("preprocessing_conflict") from exc + if not stat.S_ISDIR(st.st_mode) or (st.st_dev, st.st_ino) != (cap.device, cap.inode) or st.st_uid != os.getuid() or (st.st_mode & 0o777) != 0o700: + raise RuntimeError("preprocessing_conflict") target = dwh if dwh is not None else build_dwh(cfg) physical = target.introspect() _add_examples(target, physical, cfg.examples) diff --git a/harness/tht/workspace_writer_lock.py b/harness/tht/workspace_writer_lock.py index 810ffc24..f719fbd4 100644 --- a/harness/tht/workspace_writer_lock.py +++ b/harness/tht/workspace_writer_lock.py @@ -11,8 +11,6 @@ import fcntl import os import re import stat -import struct -import sys from dataclasses import dataclass @@ -78,24 +76,11 @@ 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() - # 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. + # First prove that an independently opened description cannot acquire the + # lock. Then probe the inherited description itself. flock is an + # open-file-description lock: the second call succeeds only on the same + # description held by the backend and does not release it. try: - # 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 not in (errno.EACCES, errno.EAGAIN, errno.EWOULDBLOCK): @@ -106,6 +91,10 @@ def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: di except OSError: pass raise WorkspaceWriterConflict() + try: + fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError as exc: + raise WorkspaceWriterConflict() from exc finally: try: os.close(lock_fd)