diff --git a/backend/native/workspace-fs-at/workspace_fs_at.cc b/backend/native/workspace-fs-at/workspace_fs_at.cc index 379f9a80..b8de192e 100644 --- a/backend/native/workspace-fs-at/workspace_fs_at.cc +++ b/backend/native/workspace-fs-at/workspace_fs_at.cc @@ -8,6 +8,8 @@ #include #include #include +#include +#include #include #ifdef __APPLE__ #include @@ -19,21 +21,14 @@ namespace { static const uint64_t MAGIC=0x5448545746534154ULL; static const char OWNER_TOKEN = 0; -struct Handle { uint64_t magic; const void* owner; int fd; bool directory; bool closed; unsigned borrows; }; -static std::unordered_set handles; -static void finalize(napi_env, void* data, void*) { - auto *h=static_cast(data); - if (!h) return; - if (handles.erase(h) && !h->closed) { h->closed=true; int rc = ::close(h->fd); (void)rc; } - delete h; -} -static bool getHandle(napi_env env,napi_value v,Handle** out) { - void *p=nullptr; - if (napi_get_value_external(env,v,&p)!=napi_ok || !p) return false; - auto *h=static_cast(p); - if (handles.find(h)==handles.end() || h->magic!=MAGIC || h->owner!=&OWNER_TOKEN || h->closed) return false; - *out=h; return true; -} +struct Handle { uint64_t magic; const void* owner; napi_env env; int fd; bool directory; bool closed; unsigned borrows; }; +struct EnvState { std::unordered_set handles; std::unordered_set all; std::mutex mutex; }; +static std::mutex statesMutex; +static std::unordered_map states; +static EnvState* stateFor(napi_env env) { std::lock_guard g(statesMutex); auto it=states.find(env); if(it!=states.end()) return it->second; auto *s=new EnvState; states.emplace(env,s); return s; } +static void finalize(napi_env env, void* data, void*) { auto *h=static_cast(data); if (!h) return; auto *st=stateFor(env); { std::lock_guard g(st->mutex); st->handles.erase(h); st->all.erase(h); if (!h->closed) { h->closed=true; (void)::close(h->fd); } } delete h; } +enum HandleResult { HANDLE_OK, HANDLE_CLOSED, HANDLE_ARGUMENT }; +static HandleResult getHandle(napi_env env,napi_value v,Handle** out) { void *p=nullptr; if (napi_get_value_external(env,v,&p)!=napi_ok || !p) return HANDLE_ARGUMENT; auto *h=static_cast(p); auto *st=stateFor(env); std::lock_guard g(st->mutex); if (st->all.find(h)==st->all.end() || h->magic!=MAGIC || h->owner!=&OWNER_TOKEN || h->env!=env) return HANDLE_ARGUMENT; if (h->closed || st->handles.find(h)==st->handles.end()) return HANDLE_CLOSED; *out=h; return HANDLE_OK; } static const char* errnoName(int e) { switch(e) { case EINVAL:return "EINVAL"; case EBADF:return "EBADF"; case EEXIST:return "EEXIST"; @@ -59,6 +54,7 @@ static napi_value error(napi_env env,const char* syscall,int e,const char* force return err; } static napi_value fail(napi_env env,const char*s,int e,const char* forced=nullptr){ napi_value x=error(env,s,e,forced); napi_throw(env,x); return nullptr; } +static bool requireHandle(napi_env env,napi_value v,Handle** out,const char* syscall) { auto r=getHandle(env,v,out); if(r==HANDLE_OK) return true; napi_throw(env,error(env,syscall,EINVAL,r==HANDLE_CLOSED?"ERR_WORKSPACE_FS_AT_HANDLE_CLOSED":"ERR_WORKSPACE_FS_AT_ARGUMENT")); return false; } static bool str(napi_env env,napi_value v,std::string& out) { size_t n=0; if(napi_get_value_string_utf8(env,v,nullptr,0,&n)!=napi_ok) return false; out.resize(n); size_t got=0; if(napi_get_value_string_utf8(env,v,out.data(),n+1,&got)!=napi_ok) return false; out.resize(got); return true; @@ -68,27 +64,16 @@ static bool validComponent(const std::string& s) { } static int openRetry(int dir,const char *name,int flags,mode_t mode) { int fd; do {fd=dir<0 ? ::open(name,flags,mode) : ::openat(dir,name,flags,mode);} while(fd<0&&errno==EINTR); return fd; } static int statRetry(int fd, struct stat *st) { int rc; do {rc=::fstat(fd,st);} while(rc<0&&errno==EINTR); return rc; } -static napi_value statObj(napi_env env,const struct stat& st) { - napi_value o,n; napi_create_object(env,&o); - napi_create_bigint_uint64(env,(uint64_t)st.st_dev,&n); napi_set_named_property(env,o,"device",n); - napi_create_bigint_uint64(env,(uint64_t)st.st_ino,&n); napi_set_named_property(env,o,"inode",n); - napi_create_uint32(env,(uint32_t)st.st_mode,&n); napi_set_named_property(env,o,"mode",n); - napi_create_uint32(env,(uint32_t)st.st_uid,&n); napi_set_named_property(env,o,"uid",n); - napi_create_uint32(env,(uint32_t)st.st_gid,&n); napi_set_named_property(env,o,"gid",n); - napi_create_bigint_uint64(env,(uint64_t)st.st_nlink,&n); napi_set_named_property(env,o,"nlink",n); return o; -} -static napi_value result(napi_env env,int fd,bool dir,const struct stat&st) { - auto*h=new Handle{MAGIC,&OWNER_TOKEN,fd,dir,false,0}; handles.insert(h); - napi_value e,o,s; napi_create_external(env,h,finalize,nullptr,&e); napi_create_object(env,&o); - napi_set_named_property(env,o,"handle",e); s=statObj(env,st); napi_set_named_property(env,o,"openedStat",s); return o; -} +static bool setStat(napi_env env,napi_value o,const char *name,uint64_t value,bool big) { napi_value n; napi_status rc=big?napi_create_bigint_uint64(env,value,&n):napi_create_uint32(env,(uint32_t)value,&n); if(rc!=napi_ok) return false; return napi_set_named_property(env,o,name,n)==napi_ok; } +static napi_value statObj(napi_env env,const struct stat& st) { napi_value o; if(napi_create_object(env,&o)!=napi_ok || !setStat(env,o,"device",(uint64_t)st.st_dev,true) || !setStat(env,o,"inode",(uint64_t)st.st_ino,true) || !setStat(env,o,"mode",(uint64_t)st.st_mode,false) || !setStat(env,o,"uid",(uint64_t)st.st_uid,false) || !setStat(env,o,"gid",(uint64_t)st.st_gid,false) || !setStat(env,o,"nlink",(uint64_t)st.st_nlink,true)) { napi_throw_error(env,"ERR_WORKSPACE_FS_AT_ARGUMENT","native result conversion failed"); return nullptr; } return o; } +static napi_value result(napi_env env,int fd,bool dir,const struct stat&st) { auto*h=new Handle{MAGIC,&OWNER_TOKEN,env,fd,dir,false,0}; auto *stt=stateFor(env); { std::lock_guard g(stt->mutex); stt->handles.insert(h); stt->all.insert(h); } napi_value e,o,so; if(napi_create_external(env,h,finalize,nullptr,&e)!=napi_ok || napi_create_object(env,&o)!=napi_ok || napi_set_named_property(env,o,"handle",e)!=napi_ok || !(so=statObj(env,st)) || napi_set_named_property(env,o,"openedStat",so)!=napi_ok) { { std::lock_guard g(stt->mutex); h->closed=true; stt->handles.erase(h); } (void)::close(fd); napi_throw_error(env,"ERR_WORKSPACE_FS_AT_ARGUMENT","native result conversion failed"); return nullptr; } return o; } static bool getInput(napi_env env,napi_callback_info info,napi_value *a,size_t *argc) { return napi_get_cb_info(env,info,argc,a,nullptr,nullptr)==napi_ok; } static napi_value openatFn(napi_env env,napi_callback_info info) { size_t argc=1; napi_value a[1]; if(!getInput(env,info,a,&argc)||argc!=1)return fail(env,"openat",EINVAL); napi_value pv,nv,kv,mv; 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],"kind",&kv)!=napi_ok||napi_get_named_property(env,a[0],"createMode",&mv)!=napi_ok)return fail(env,"openat",EINVAL); std::string name,kind; if(!str(env,nv,name)||!str(env,kv,kind))return fail(env,"openat",EINVAL); int32_t mode;if(napi_get_value_int32(env,mv,&mode)!=napi_ok)return fail(env,"openat",EINVAL); bool root=name=="/"; if(!root&&!validComponent(name))return fail(env,"openat",EINVAL); Handle*ph=nullptr; - if(root){if(kind!="directory"||mode!=0)return fail(env,"openat",EINVAL);} else if(!getHandle(env,pv,&ph)||!ph->directory)return fail(env,"openat",EBADF); + if(root){if(kind!="directory"||mode!=0)return fail(env,"openat",EINVAL);} else if(!requireHandle(env,pv,&ph,"openat")||!ph->directory)return fail(env,"openat",EBADF); int fd=-1; if(root) fd=openRetry(-1,"/",O_RDONLY|O_DIRECTORY|O_CLOEXEC|O_NOFOLLOW,0); else if(kind=="directory"&&mode==0) { fd=openRetry(ph->fd,name.c_str(),O_RDONLY|O_DIRECTORY|O_CLOEXEC|O_NOFOLLOW,0); } @@ -100,41 +85,16 @@ 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 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||!getHandle(env,a[0],&h)||!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||!getHandle(env,a[0],&h)||!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||!getHandle(env,a[0],&h)||!h->directory)return fail(env,"fsync",EBADF);int rc;do{rc=::fsync(h->fd);}while(rc<0&&errno==EINTR); +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); #ifdef __APPLE__ if(rc<0&&(errno==EINVAL||errno==ENOTSUP)){int frc;do{frc=::fcntl(h->fd,F_FULLFSYNC);}while(frc<0&&errno==EINTR);if(frc==0)return nullptr;rc=frc;} #endif 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},{"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;} +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;} } NAPI_MODULE(NODE_GYP_MODULE_NAME,init) diff --git a/backend/scripts/build-workspace-fs-at.mjs b/backend/scripts/build-workspace-fs-at.mjs index f5efd59a..bbab91b5 100644 --- a/backend/scripts/build-workspace-fs-at.mjs +++ b/backend/scripts/build-workspace-fs-at.mjs @@ -7,7 +7,16 @@ if (major !== 22 || Number(process.versions.napi ?? 0) < 8) throw new Error(`wor if (process.platform !== "linux" && process.platform !== "darwin") throw new Error(`workspace-fs-at unsupported platform: ${process.platform}`); const root = resolve(dirname(fileURLToPath(import.meta.url)), "../native/workspace-fs-at"); const npm = process.platform === "win32" ? "npm.cmd" : "npm"; -const result = spawnSync(npm, ["exec", "--offline", "--", "node-gyp@11.2.0", "rebuild", "--offline"], { cwd: root, stdio: "inherit" }); +const headerCandidates = [ + process.env.npm_config_nodedir, + "/usr/local/include/node", + "/usr/include/node", + resolve(process.execPath, "../../include/node"), +].filter((p) => p && existsSync(resolve(p, "node_api.h"))); +if (headerCandidates.length === 0) throw new Error("workspace-fs-at requires installed Node 22 headers; network downloads are forbidden"); +const headerRoot = headerCandidates[0].endsWith("/include/node") ? resolve(headerCandidates[0], "../..") : headerCandidates[0]; +const env = { ...process.env, npm_config_nodedir: headerRoot, npm_config_offline: "true", npm_config_node_gyp: undefined }; +const result = spawnSync(npm, ["exec", "--offline", "--", "node-gyp@11.2.0", "rebuild", "--offline"], { cwd: root, stdio: "inherit", env }); if (result.error) throw result.error; if (result.status !== 0) process.exit(result.status ?? 1); const output = resolve(root, "build/Release/workspace_fs_at.node"); diff --git a/backend/src/native/workspace-fs-at-binding.d.ts b/backend/src/native/workspace-fs-at-binding.d.ts index e832d894..1d82718f 100644 --- a/backend/src/native/workspace-fs-at-binding.d.ts +++ b/backend/src/native/workspace-fs-at-binding.d.ts @@ -11,6 +11,4 @@ export interface WorkspaceFsAtBindingV1 { fstatat(parent:NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1):NativeWorkspaceFsAtStatV1; fsyncDirectory(handle:NativeWorkspaceFsAtHandleV1):void; close(handle:NativeWorkspaceFsAtHandleV1):void; - fdNumberForSynchronousBorrow(handle: NativeWorkspaceFsAtHandleV1): number; - fcntl(handle: NativeWorkspaceFsAtHandleV1, operation: "probe-exclusive-nonblocking" | "hold-exclusive" | "unlock"): "held" | "available"; } diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index 3df1bfdb..c908b060 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -15,6 +15,8 @@ import { type CanonicalWorkspaceId, type WorkspaceLockRootIdentityV1, Revision40, + workspaceRootInternals, + makeBorrowedRootLease, } from "./workspace-lock-root-lease.js"; export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; } @@ -72,14 +74,14 @@ async function durableJson(path: string, value: unknown): Promise { } export class PreprocessingStateStore { - private readonly rootLease: Pick; + 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 = rootLease; + } else this.rootLease = { assertLive: () => workspaceRootInternals.assertLive(rootLease), anchoredPath: () => { throw fail(); } }; } private root(): string { this.rootLease.assertLive(); return this.rootLease.anchoredPath(); } private paths(input: { runId: string }) { @@ -158,13 +160,13 @@ 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(); } + assertWriterPath(): void { this.assertLive(); workspaceRootInternals.assertLive(this.root); this.writer.assertPath(); } invalidateForSettlement(): void { this.settled = true; } static [INTERNAL_STATE](id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); } async runUnderSessionReadersExclusive(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise): Promise { this.assertLive(); if (this.readerExclusive || this.spawnActive) throw fail(); this.readerExclusive = true; let lock: WorkspaceRootLock; - try { lock = await this.root.acquireSessionReadersExclusive(); } catch { this.readerExclusive = false; this.poisoned = true; throw fail(); } + 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 { @@ -182,7 +184,7 @@ export class WorkspaceWriterLockCapability { 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 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) }); + 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; } } 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(); } @@ -199,7 +201,7 @@ export class OrderedWorkspaceWriterCapabilitySet { static [INTERNAL_STATE](caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } invalidate(): void { this.live = false; for (const cap of this.caps.values()) cap.invalidateForSettlement(); } get workspaceIds(): readonly CanonicalWorkspaceId[] { if (!this.live) throw fail(); return [...this.caps.keys()]; } - async forWorkspace(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise): Promise { if (!this.live) throw fail(); const cap = this.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: BorrowedVerifiedWorkspaceLockRootLease.make(cap.rootIdentity, () => { if (!this.live) throw fail(); }), writerCapability: cap }); } + 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 }); } async forEachWorkspace(action: (lease: OrderedWorkspaceCapability) => Promise): Promise { return Promise.all(this.workspaceIds.map(id => this.forWorkspace(id, action))); } } export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise): Promise { @@ -210,7 +212,7 @@ export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonl const root = source.transfer(); let writer: WorkspaceRootLock | undefined; try { - writer = await root.acquireWriterLock(); + writer = await workspaceRootInternals.acquireWriter(root); caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer)); } catch (error) { try { writer?.close(); } catch {} diff --git a/backend/src/workspaces/workspace-fs-at.ts b/backend/src/workspaces/workspace-fs-at.ts index 8cd3f2e2..c6cff631 100644 --- a/backend/src/workspaces/workspace-fs-at.ts +++ b/backend/src/workspaces/workspace-fs-at.ts @@ -50,7 +50,7 @@ function wrapDirectory(result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: 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"); + if ((st.mode & 0o170000) !== 0o100000 || (st.mode & 0o7777) !== 0o600 || st.uid !== (process.getuid?.() ?? st.uid) || st.nlink !== 1n) throw new Error("invalid lock identity"); return OwnedWorkspaceFsAtRegularFile[INTERNAL](result.handle, st); } catch (error) { try { binding.close(result.handle); } catch { /* preserve conversion error */ } throw error; } } @@ -60,50 +60,47 @@ function withLockFd(value: OwnedWorkspaceFsAtRegularFile, action: (fd: number const count = borrowing.get(value) ?? 0; borrowing.set(value, count + 1); try { return binding.withFd(rawHandles.get(value)!, action as (fd: number) => void) as T; } finally { borrowing.set(value, count); } } -function anchoredDirectoryPath(value: OwnedWorkspaceFsAtDirectory): string { - if (!rawHandles.has(value)) throw new Error("workspace descriptor is closed"); - const count = borrowing.get(value) ?? 0; borrowing.set(value, count + 1); - try { return binding.withFd(rawHandles.get(value)!, fd => `${process.platform === "darwin" ? "/dev/fd" : "/proc/self/fd"}/${fd}`); } - finally { borrowing.set(value, count); } -} export class WorkspaceFsAtV1 { openRoot(): OwnedWorkspaceFsAtDirectory { return wrapDirectory(binding.openat({ parent: null, name: "/", kind: "directory", createMode: 0 })); } openDirectoryAt(parent: OwnedWorkspaceFsAtDirectory, name: string): OwnedWorkspaceFsAtDirectory { return wrapDirectory(binding.openat({ parent: rawDirectory(parent), name: component(name), kind: "directory", createMode: 0 })); } - openOrCreateLockAt(parent: OwnedWorkspaceFsAtDirectory, name: LockFileName, mode: 0o600): OwnedWorkspaceFsAtRegularFile { if ((name !== "writer.lock" && name !== "session-readers.lock") || mode !== 0o600) throw new Error("invalid lock"); return wrapLock(binding.openat({ parent: rawDirectory(parent), name: component(name), kind: "regular_lock", createMode: 0o600 })); } + openOrCreateLockAt(parent: OwnedWorkspaceFsAtDirectory, name: LockFileName, mode: 0o600): OwnedWorkspaceFsAtRegularFile { + if ((name !== "writer.lock" && name !== "session-readers.lock") || mode !== 0o600) throw new Error("invalid lock"); + return wrapLock(binding.openat({ parent: rawDirectory(parent), name: component(name), kind: "regular_lock", createMode: 0o600 })); + } 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 spawnChildFromOwnedCapability( - writer: OwnedWorkspaceFsAtRegularFile, - root: OwnedWorkspaceFsAtDirectory, - executable: string, - args: readonly string[], - environment: NodeJS.ProcessEnv = process.env, -): Promise<{ exitCode: number; stdout: Uint8Array; stderr: Uint8Array }> { - if (!rawHandles.has(writer) || !rawHandles.has(root)) throw new Error("preprocessing_conflict"); - const writerRaw = rawHandles.get(writer)!; const rootRaw = rawHandles.get(root)!; - const opened: number[] = []; - try { - const writerFd = openSync("/dev/null", "r"); opened.push(writerFd); - const rootFd = openSync("/dev/null", "r"); opened.push(rootFd); - binding.duplicateForChildStdio(writerRaw, rootRaw, writerFd, rootFd); - const child = spawn(executable, [...args], { stdio: ["ignore", "pipe", "pipe", writerFd, rootFd], env: { ...environment, THOTH_WORKSPACE_CAPABILITY_REQUIRED: "1" } }); - const out: Buffer[] = []; const err: Buffer[] = []; - child.stdout?.on("data", (chunk: Buffer) => out.push(chunk)); child.stderr?.on("data", (chunk: Buffer) => err.push(chunk)); - return new Promise((resolve, reject) => { - child.once("error", reject); child.once("close", code => resolve({ exitCode: code ?? 1, stdout: Buffer.concat(out), stderr: Buffer.concat(err) })); - }); - } finally { for (const fd of opened) { try { closeSync(fd); } catch {} } } -} +// Friend operations are deliberately not methods on the frozen public classes. They are +// exported only for the preprocessing-state module; callers cannot obtain handles or paths. +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 }>; +}; +export type WorkspaceLockRootIdentityV1Like = { readonly device: bigint; readonly inode: bigint }; +export const workspaceFsAtFriend: WorkspaceRootFriend = { + 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"); + }, + spawn(writer, root, executable, args, environment = process.env) { + if (!rawHandles.has(writer) || !rawHandles.has(root)) return Promise.reject(new Error("preprocessing_conflict")); + const writerRaw = rawHandles.get(writer)!; const rootRaw = rawHandles.get(root)!; const opened: number[] = []; + try { + const writerFd = openSync("/dev/null", "r"); opened.push(writerFd); const rootFd = openSync("/dev/null", "r"); opened.push(rootFd); + binding.duplicateForChildStdio(writerRaw, rootRaw, writerFd, rootFd); + const child = spawn(executable, [...args], { stdio: ["ignore", "pipe", "pipe", writerFd, rootFd], env: { ...environment, THOTH_WORKSPACE_CAPABILITY_REQUIRED: "1" } }); + const out: Buffer[] = []; const err: Buffer[] = []; + child.stdout?.on("data", (chunk: Buffer) => out.push(chunk)); child.stderr?.on("data", (chunk: Buffer) => err.push(chunk)); + return new Promise((resolve, reject) => { child.once("error", reject); child.once("close", code => resolve({ exitCode: code ?? 1, stdout: Buffer.concat(out), stderr: Buffer.concat(err) })); }); + } catch (error) { return Promise.reject(normalizeError(error)); } + finally { for (const fd of opened) { try { closeSync(fd); } catch {} } } + }, +}; +function fsStatAtNoFollow(parent: OwnedWorkspaceFsAtDirectory, name: LockFileName): WorkspaceFsAtStatV1 { return binding.fstatat(rawDirectory(parent), component(name)); } diff --git a/backend/src/workspaces/workspace-lock-root-lease.ts b/backend/src/workspaces/workspace-lock-root-lease.ts index 6aecb80b..347f538c 100644 --- a/backend/src/workspaces/workspace-lock-root-lease.ts +++ b/backend/src/workspaces/workspace-lock-root-lease.ts @@ -1,6 +1,5 @@ -import { lstatSync, realpathSync } from "node:fs"; import { resolve } from "node:path"; -import { spawnChildFromOwnedCapability, WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, type OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js"; +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; declare const sha256HexBrand: unique symbol; @@ -8,99 +7,97 @@ export type CanonicalWorkspaceId = string & { readonly [canonicalWorkspaceIdBran 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 INPUT_BRAND = new WeakMap(); const ROOT_INTERNAL = Symbol("workspace-root-internal"); -const activeWriterLocks = new Set(); +const conflict = (message = "preprocessing_conflict"): Error => { const e = new Error(message); e.name = "PreprocessingConflictError"; return e; }; +function modeExact(mode: number, type: number, permissions: number): boolean { return (mode & 0o170000) === type && (mode & 0o7777) === permissions; } +function exactRoot(stat: WorkspaceFsAtStatV1, uid: number): boolean { return modeExact(stat.mode, 0o040000, 0o700) && stat.uid === uid && stat.nlink >= 2n; } +function exactParent(stat: WorkspaceFsAtStatV1, uid: number): boolean { return exactRoot(stat, uid); } +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; } -/** Internal opaque lock returned only after the root has been checked. */ -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 { +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 }>; }; +class RootLock implements InternalLock { private live = true; - 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(); } + 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) {} + 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); } } } -export class BorrowedVerifiedWorkspaceLockRootLease { - private constructor(readonly identity: WorkspaceLockRootIdentityV1, private readonly check: () => void) {} - static make(identity: WorkspaceLockRootIdentityV1, check: () => void): BorrowedVerifiedWorkspaceLockRootLease { return new BorrowedVerifiedWorkspaceLockRootLease(identity, check); } - assertLive(): void { this.check(); } +const activeWriterLocks = new Set(); +type RootState = { fs: WorkspaceFsAtV1; parent: OwnedWorkspaceFsAtDirectory; root: OwnedWorkspaceFsAtDirectory; serviceUid: number; owner: symbol; identity: WorkspaceLockRootIdentityV1; live: boolean }; +const rootState = new WeakMap(); + +export class CanonicalWorkspaceLockRootInput { + private constructor(readonly workspaceId: CanonicalWorkspaceId) {} } +function makeInput(id: CanonicalWorkspaceId, owner: symbol): CanonicalWorkspaceLockRootInput { const value = Object.create(CanonicalWorkspaceLockRootInput.prototype) as CanonicalWorkspaceLockRootInput; Object.defineProperty(value, "workspaceId", { value: id, enumerable: true, writable: false }); INPUT_BRAND.set(value, owner); return value; } + +export class BorrowedVerifiedWorkspaceLockRootLease { + private constructor(readonly identity: WorkspaceLockRootIdentityV1) {} +} +function makeBorrowed(identity: WorkspaceLockRootIdentityV1): BorrowedVerifiedWorkspaceLockRootLease { return Reflect.construct(BorrowedVerifiedWorkspaceLockRootLease, [identity]) as BorrowedVerifiedWorkspaceLockRootLease; } +export function makeBorrowedRootLease(identity: WorkspaceLockRootIdentityV1, _check?: () => void): BorrowedVerifiedWorkspaceLockRootLease { return makeBorrowed(identity); } + export class WorkspaceSessionReadersLockLease { private live = true; - 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); } + private constructor(private readonly lock: InternalLock, private readonly root: OwnedWorkspaceFsAtDirectory, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} + transfer(): WorkspaceSessionReadersLockLease { if (!this.live) throw conflict(); this.live = false; return Reflect.construct(WorkspaceSessionReadersLockLease, [this.lock, this.root, this.rootIdentity]) as WorkspaceSessionReadersLockLease; } async close(): Promise { if (!this.live) return; this.live = false; let failure: unknown; try { this.lock.close(); } catch (e) { failure = e; } try { this.root.close(); } catch (e) { failure ??= e; } if (failure) throw conflict(); } - static [ROOT_INTERNAL](lock: InternalWorkspaceRootLock, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1): WorkspaceSessionReadersLockLease { return new WorkspaceSessionReadersLockLease(lock, root, identity); } } + +function assertRootLive(root: VerifiedWorkspaceLockRootLease): void { const state = rootState.get(root); if (!state || state.live !== true) throw conflict(); try { const retained = state.root.stat(); const named = state.fs.statAtNoFollow(state.parent, state.identity.workspaceId); if (!sameIdentity(retained, state.identity) || !sameIdentity(named, state.identity) || !exactRoot(retained, state.serviceUid) || !exactRoot(named, state.serviceUid)) throw conflict("preprocessing_conflict: workspace root identity changed"); } catch (e) { if ((e as Error).name === "PreprocessingConflictError") throw e; throw conflict("preprocessing_conflict: workspace root identity changed"); } } export class VerifiedWorkspaceLockRootLease { - private live = true; - private borrowed = 0; - private constructor(private readonly owner: symbol, private readonly fs: WorkspaceFsAtV1, private readonly root: OwnedWorkspaceFsAtDirectory, readonly identity: WorkspaceLockRootIdentityV1, private readonly rootPath: string, private readonly serviceUid: number) {} - private assertPath(): void { - if (!this.live) throw conflict(); - try { const st = lstatSync(this.rootPath); if (!st.isDirectory() || st.uid !== this.serviceUid || (st.mode & 0o777) !== 0o700 || !sameIdentity(st, this.identity)) throw conflict("preprocessing_conflict: workspace root identity changed"); } - catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("preprocessing_conflict: workspace root identity changed"); } - const retained = this.root.stat(); if (!sameIdentity(retained, this.identity) || !exactRoot(retained, this.serviceUid)) throw conflict("preprocessing_conflict: workspace root identity changed"); - } - assertLive(): void { this.assertPath(); } - /** Internal capability-scoped path; it resolves through the retained root FD. */ - anchoredPath(): string { this.assertPath(); return this.fs.anchoredDirectoryPath(this.root); } - async borrow(action: (lease: { readonly identity: WorkspaceLockRootIdentityV1; assertLive(): void }) => Promise): Promise { this.assertPath(); this.borrowed++; try { return await action({ identity: this.identity, assertLive: () => this.assertPath() }); } finally { this.borrowed--; } } - transfer(): VerifiedWorkspaceLockRootLease { this.assertPath(); if (this.borrowed) throw conflict("workspace root is borrowed"); this.live = false; return new VerifiedWorkspaceLockRootLease(this.owner, this.fs, this.root, this.identity, this.rootPath, this.serviceUid); } + private live = true; private borrowed = 0; + private constructor(fs: WorkspaceFsAtV1, parent: OwnedWorkspaceFsAtDirectory, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, serviceUid: number, owner: symbol) { this.fs = fs; this.parent = parent; this.root = root; this.serviceUid = serviceUid; this.owner = owner; this.identity = identity; rootState.set(this, { fs, parent, root, serviceUid, owner, identity, live: true }); } + private readonly fs: WorkspaceFsAtV1; private readonly parent: OwnedWorkspaceFsAtDirectory; private readonly root: OwnedWorkspaceFsAtDirectory; private readonly serviceUid: number; private readonly owner: symbol; + readonly identity: WorkspaceLockRootIdentityV1; + borrow(action: (borrowed: BorrowedVerifiedWorkspaceLockRootLease) => Promise): Promise { assertRootLive(this); this.borrowed++; try { return Promise.resolve(action(makeBorrowed(this.identity))).finally(() => { this.borrowed--; }); } catch (e) { this.borrowed--; return Promise.reject(e); } } async acquireSessionReadersShared(): Promise { - this.assertPath(); let lock: WorkspaceRootLock | undefined; - try { const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); lock = new WorkspaceRootLock(this.fs, 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(); } + assertRootLive(this); let lock: RootLock | undefined; + 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(); } } - async acquireWriterLock(): Promise { - this.assertPath(); - const key = `${this.identity.device}:${this.identity.inode}`; - if (activeWriterLocks.has(key)) throw conflict(); - const owned = this.fs.openOrCreateLockAt(this.root, "writer.lock", 0o600); - const lock = new WorkspaceRootLock(this.fs, this.root, owned, "writer.lock", key); - try { lock.assertPath(); lock.flock("exclusive", "nonblocking"); lock.assertPath(); activeWriterLocks.add(key); return lock; } - catch (error) { try { lock.close(); } catch {} throw conflict(); } - } - async acquireSessionReadersExclusive(): Promise { this.assertPath(); const owned = this.fs.openOrCreateLockAt(this.root, "session-readers.lock", 0o600); const lock = new WorkspaceRootLock(this.fs, this.root, owned, "session-readers.lock"); try { lock.flock("exclusive", "nonblocking"); this.assertPath(); return lock; } catch { try { lock.close(); } catch {} throw conflict(); } } - async spawnChild(lock: InternalWorkspaceRootLock, executable: string, args: readonly string[], environment?: NodeJS.ProcessEnv) { this.assertPath(); return lock.spawn(this.root, executable, args, environment); } - fsync(): void { this.assertPath(); this.fs.fsyncDirectory(this.root); this.assertPath(); } - async close(): Promise { if (!this.live) return; while (this.borrowed) await new Promise(resolve => setTimeout(resolve, 1)); this.live = false; try { this.root.close(); } catch { throw conflict(); } } - static [ROOT_INTERNAL](owner: symbol, fs: WorkspaceFsAtV1, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, rootPath: string, uid: number): VerifiedWorkspaceLockRootLease { return new VerifiedWorkspaceLockRootLease(owner, fs, root, identity, rootPath, uid); } + 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); } + 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); } + export class VerifiedWorkspaceLockRootLeaseFactory { - private readonly owner = Symbol("workspace-root-factory"); - private readonly parent: OwnedWorkspaceFsAtDirectory; - private readonly parentPath: string; + private readonly owner = Symbol("workspace-root-factory"); private readonly parent: OwnedWorkspaceFsAtDirectory; private provisionTail: Promise = Promise.resolve(); constructor(private readonly input: { readonly workspaceFsAt: WorkspaceFsAtV1; readonly installationId: string; readonly sessionsRootFromValidatedInstallationConfig: string; readonly serviceUid: number; readonly provisionedWorkspaceMode: 0o700 }) { if (!Number.isInteger(input.serviceUid) || input.serviceUid < 0 || input.provisionedWorkspaceMode !== 0o700) throw new Error("invalid workspace root policy"); - try { const configured = resolve(input.sessionsRootFromValidatedInstallationConfig); this.parentPath = realpathSync(configured); } catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("invalid sessions root"); } - if (!this.parentPath.startsWith("/") || this.parentPath.split("/").includes("..")) throw conflict("invalid sessions root"); + const configured = input.sessionsRootFromValidatedInstallationConfig; + if (!configured.startsWith("/") || configured.split("/").some(c => c === "" ? false : c === "." || c === "..")) throw conflict("invalid sessions root"); let d = input.workspaceFsAt.openRoot(); - try { for (const c of this.parentPath.split("/").filter(Boolean)) { const n = input.workspaceFsAt.openDirectoryAt(d, c); d.close(); d = n; } this.parent = d; this.checkParent(); } - catch (error) { try { d.close(); } catch {} throw error; } + 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(); } + catch (e) { try { d.close(); } catch {} throw conflict("invalid sessions root"); } } - private checkParent(): void { try { const s = lstatSync(this.parentPath); const p = this.parent.stat(); if (!s.isDirectory() || !sameIdentity(s, p) || s.uid !== this.input.serviceUid || (s.mode & 0o777) !== 0o700) throw conflict("installation root identity changed"); } catch (error) { if ((error as Error).name === "PreprocessingConflictError") throw error; throw conflict("installation root identity changed"); } } - canonicalInput(workspaceId: string): CanonicalWorkspaceLockRootInput { if (!/^[a-z][a-z0-9-]{2,62}$/.test(workspaceId)) throw conflict(); return Object.assign(Object.create(CanonicalWorkspaceLockRootInput.prototype), { workspaceId: workspaceId as CanonicalWorkspaceId, owner: this.owner }) as CanonicalWorkspaceLockRootInput; } - private async open(input: CanonicalWorkspaceLockRootInput): Promise { if (input.owner !== this.owner) throw conflict(); this.checkParent(); let root: OwnedWorkspaceFsAtDirectory; try { root = this.input.workspaceFsAt.openDirectoryAt(this.parent, input.workspaceId); this.checkParent(); } catch { throw conflict(); } if (!exactRoot(root.stat(), this.input.serviceUid)) { root.close(); throw conflict(); } const identity = { schemaVersion: 1 as const, workspaceId: input.workspaceId, device: root.stat().device, inode: root.stat().inode }; const lease = VerifiedWorkspaceLockRootLease[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; } + 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"); } } + canonicalInput(workspaceId: string): CanonicalWorkspaceLockRootInput { if (!/^[a-z][a-z0-9-]{2,62}$/.test(workspaceId)) throw conflict(); return makeInput(workspaceId as CanonicalWorkspaceId, this.owner); } + private validInput(input: CanonicalWorkspaceLockRootInput): boolean { return typeof input === "object" && input !== null && INPUT_BRAND.get(input) === this.owner; } + private async open(input: CanonicalWorkspaceLockRootInput): Promise { if (!this.validInput(input)) throw conflict(); this.checkParent(); let root: OwnedWorkspaceFsAtDirectory; try { root = this.input.workspaceFsAt.openDirectoryAt(this.parent, input.workspaceId); } catch (e) { throw e; } try { this.checkParent(); const st = root.stat(); if (!exactRoot(st, this.input.serviceUid)) throw conflict(); const identity = { schemaVersion: 1 as const, workspaceId: input.workspaceId, device: st.device, inode: st.inode }; const lease = makeRoot(this.input.workspaceFsAt, this.parent, root, identity, this.input.serviceUid, this.owner); this.assertRoot(lease); return lease; } catch (e) { try { root.close(); } catch {} throw ((e as Error).name === "PreprocessingConflictError" ? e : conflict()); } } + private assertRoot(root: VerifiedWorkspaceLockRootLease): void { assertRootLive(root) } acquire(input: CanonicalWorkspaceLockRootInput): Promise { return this.open(input); } async acquireOrProvision(input: CanonicalWorkspaceLockRootInput): Promise { - if (input.owner !== this.owner) throw conflict(); - try { return await this.open(input); } catch (error) { if (!String((error as Error).message).includes("identity") && (error as {code?: string}).code !== "ENOENT") { /* open errors are normalized; probe the anchored parent below */ } } - this.checkParent(); - try { this.input.workspaceFsAt.mkdirAt(this.parent, input.workspaceId, 0o700); } catch (error) { if ((error as {code?: string}).code !== "EEXIST") throw conflict(); } - this.checkParent(); - try { this.input.workspaceFsAt.fsyncDirectory(this.parent); } catch { throw conflict(); } - return this.open(input); + if (!this.validInput(input)) throw conflict(); + 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; + } finally { release(); } } } + +// Internal friend boundary for writer/quiescence code. The root lease class itself has only +// 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()); } }, + 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 3e5b014b..af97f2c3 100644 --- a/backend/test/workspace-lock-root-lease.test.ts +++ b/backend/test/workspace-lock-root-lease.test.ts @@ -2,7 +2,7 @@ import { describe, expect, it, afterEach } from "vitest"; 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"; +import { CanonicalWorkspaceLockRootInput, VerifiedWorkspaceLockRootLeaseFactory } from "../src/workspaces/workspace-lock-root-lease.js"; const roots: string[] = []; afterEach(() => { for (const root of roots.splice(0)) rmSync(root, { recursive: true, force: true }); }); function factory(sessionsRootFromValidatedInstallationConfig: string) { return new VerifiedWorkspaceLockRootLeaseFactory({ workspaceFsAt: new WorkspaceFsAtV1(), installationId: "i", sessionsRootFromValidatedInstallationConfig, serviceUid: process.getuid!(), provisionedWorkspaceMode: 0o700 }); } @@ -41,4 +41,14 @@ describe("retained canonical workspace root", () => { await expect(outcome).rejects.toThrow(/preprocessing/); }); + it("rejects forged canonical inputs even when the prototype is copied", async () => { + const parent = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(parent); const f = factory(parent); + const forged = Object.assign(Object.create(CanonicalWorkspaceLockRootInput.prototype), { workspaceId: "abc-workspace" }); + await expect(f.acquireOrProvision(forged as CanonicalWorkspaceLockRootInput)).rejects.toThrow(/preprocessing/); + }); + it("rejects symlink sessions roots and special permission bits", () => { + const base = mkdtempSync(join(process.cwd(), "thoth-root-")); roots.push(base); const target = mkdtempSync(join(process.cwd(), "thoth-target-")); roots.push(target); + const link = join(base, "sessions"); symlinkSync(target, link); expect(() => factory(link)).toThrow(); + chmodSync(base, 0o1700); expect(() => factory(base)).toThrow(); + }); });