fix task5 root lease native surface and session ownership
This commit is contained in:
@@ -77,7 +77,7 @@ static napi_value openatFn(napi_env env,napi_callback_info info) {
|
||||
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(!requireHandle(env,pv,&ph,"openat")||!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"))return nullptr; else if(!ph->directory)return fail(env,"openat",EINVAL,"ERR_WORKSPACE_FS_AT_ARGUMENT");
|
||||
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); }
|
||||
@@ -86,58 +86,18 @@ static napi_value openatFn(napi_env env,napi_callback_info info) {
|
||||
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(kind=="directory"&&!S_ISDIR(st.st_mode)){::close(fd);return fail(env,"openat",ENOTDIR);}
|
||||
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);}
|
||||
if(kind=="regular_lock" && (!S_ISREG(st.st_mode)||(st.st_mode&07777)!=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(off<len){ssize_t rc=::write(h->fd,(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);
|
||||
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,"ERR_WORKSPACE_FS_AT_ARGUMENT");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,"ERR_WORKSPACE_FS_AT_ARGUMENT");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",EINVAL,"ERR_WORKSPACE_FS_AT_ARGUMENT");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=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<std::mutex> 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},{"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;}
|
||||
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||!requireHandle(env,a[0],&h,"fdNumberForSynchronousBorrow"))return nullptr;if(h->borrows==UINT_MAX)return fail(env,"fdNumberForSynchronousBorrow",EOVERFLOW);h->borrows++;h->borrows--;napi_value out;napi_create_int32(env,h->fd,&out);return out;}
|
||||
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},{"fdNumberForSynchronousBorrow",0,fdNumberForSynchronousBorrowFn,0,0,0,napi_enumerable,0}};napi_define_properties(env,exports,6,pub);return exports;}
|
||||
}
|
||||
NAPI_MODULE(NODE_GYP_MODULE_NAME,init)
|
||||
|
||||
+4
-10
@@ -3,19 +3,13 @@ declare const nativeWorkspaceFsAtComponentBrand: unique symbol;
|
||||
export interface NativeWorkspaceFsAtHandleV1 { readonly [nativeWorkspaceFsAtHandleBrand]: true; }
|
||||
export type NativeWorkspaceFsAtComponentV1 = string & { readonly [nativeWorkspaceFsAtComponentBrand]: true };
|
||||
export interface NativeWorkspaceFsAtStatV1 { readonly device: bigint; readonly inode: bigint; readonly mode: number; readonly uid: number; readonly gid: number; readonly nlink: bigint; }
|
||||
export interface NativeWorkspaceFsAtErrorV1 extends Error { readonly code:string; readonly errno:number; readonly syscall:"openat"|"mkdirat"|"fstat"|"fstatat"|"fsync"|"fcntl"|"close"; }
|
||||
export interface NativeWorkspaceFsAtErrorV1 extends Error { readonly code:string; readonly errno:number; readonly syscall:"openat"|"mkdirat"|"fstatat"|"fsyncDirectory"|"close"|"fdNumberForSynchronousBorrow"; }
|
||||
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"|"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;
|
||||
openat(input:{readonly parent:NativeWorkspaceFsAtHandleV1|null;readonly name:"/"|NativeWorkspaceFsAtComponentV1;readonly kind:"directory"|"regular_lock";readonly createMode:0|0o600}):NativeWorkspaceFsAtOpenResultV1;
|
||||
mkdirat(parent:NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1,mode:0o700):void;
|
||||
fstatat(parent:NativeWorkspaceFsAtHandleV1,name:NativeWorkspaceFsAtComponentV1):NativeWorkspaceFsAtStatV1;
|
||||
fsyncDirectory(handle:NativeWorkspaceFsAtHandleV1):void;
|
||||
close(handle:NativeWorkspaceFsAtHandleV1):void;
|
||||
fdNumberForSynchronousBorrow(handle:NativeWorkspaceFsAtHandleV1):number;
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ import { loadPiAuthProviders } from "./auth-providers.js";
|
||||
import { secretValue } from "../config/secret-bundle.js";
|
||||
import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext } from "../auth/principal.js";
|
||||
import { createPiRuntimeAgentSnapshot } from "./managed-config.js";
|
||||
import type { WorkspaceSessionReadersLockLease } from "../workspaces/workspace-lock-root-lease.js";
|
||||
|
||||
export interface SessionRuntime {
|
||||
rpc: RpcClient;
|
||||
@@ -15,6 +16,7 @@ export interface SessionRuntime {
|
||||
child: ChildProcessWithoutNullStreams;
|
||||
ownerKey?: string;
|
||||
releaseRuntimeConfig?: () => void;
|
||||
releaseSessionReaders?: () => void;
|
||||
}
|
||||
|
||||
export interface RuntimeOptions {
|
||||
@@ -26,6 +28,7 @@ export interface RuntimeOptions {
|
||||
mode?: "new" | "resume";
|
||||
principal?: PrincipalContext;
|
||||
runtimeConfig?: RuntimeConfigLease;
|
||||
sessionReadersLease?: WorkspaceSessionReadersLockLease;
|
||||
}
|
||||
|
||||
/** Injectable child-process boundary; callbacks may ignore arguments in simpler tests. */
|
||||
@@ -149,11 +152,13 @@ export class PiProcessManager {
|
||||
const existing = this.runtimes.get(sessionId);
|
||||
if (existing) {
|
||||
o.runtimeConfig?.release();
|
||||
void o.sessionReadersLease?.close().catch(() => undefined);
|
||||
throw new Error(`session runtime already active: ${sessionId}`);
|
||||
}
|
||||
if (o.principal) this.teardownForPrincipal(o.principal);
|
||||
if (this.runtimes.size >= this.cfg.maxPiProcesses) {
|
||||
o.runtimeConfig?.release();
|
||||
void o.sessionReadersLease?.close().catch(() => undefined);
|
||||
throw new Error("max Pi processes reached");
|
||||
}
|
||||
const author = o.author ?? "dev@local";
|
||||
@@ -163,6 +168,7 @@ export class PiProcessManager {
|
||||
child = this.spawnFn(sessionId, author, provider, o.principal, o.runtimeConfig?.path);
|
||||
} catch (error) {
|
||||
o.runtimeConfig?.release();
|
||||
void o.sessionReadersLease?.close().catch(() => undefined);
|
||||
throw error;
|
||||
}
|
||||
let runtimeConfigReleased = false;
|
||||
@@ -171,8 +177,15 @@ export class PiProcessManager {
|
||||
runtimeConfigReleased = true;
|
||||
o.runtimeConfig?.release();
|
||||
};
|
||||
let sessionReadersReleased = false;
|
||||
const releaseSessionReaders = () => {
|
||||
if (sessionReadersReleased) return;
|
||||
sessionReadersReleased = true;
|
||||
void o.sessionReadersLease?.close().catch(() => undefined);
|
||||
};
|
||||
child.once("exit", releaseRuntimeConfig);
|
||||
child.once("close", releaseRuntimeConfig);
|
||||
child.once("close", releaseSessionReaders);
|
||||
let rt: SessionRuntime | undefined;
|
||||
try {
|
||||
const rpc = new RpcClient(child);
|
||||
@@ -183,6 +196,7 @@ export class PiProcessManager {
|
||||
child,
|
||||
ownerKey: o.principal ? `${o.principal.issuer}\0${o.principal.subject}` : undefined,
|
||||
...(o.runtimeConfig ? { releaseRuntimeConfig } : {}),
|
||||
...(o.sessionReadersLease ? { releaseSessionReaders } : {}),
|
||||
};
|
||||
rt = runtime;
|
||||
bridge.beginTurn();
|
||||
@@ -218,6 +232,7 @@ export class PiProcessManager {
|
||||
} catch (error) {
|
||||
if (rt && this.runtimes.get(sessionId) === rt) this.runtimes.delete(sessionId);
|
||||
releaseRuntimeConfig();
|
||||
releaseSessionReaders();
|
||||
try { child.kill(); } catch { /* preserve the initialization error */ }
|
||||
this.cleanupAgentSnapshot(child);
|
||||
throw error;
|
||||
|
||||
@@ -7,7 +7,7 @@ import { getPrincipal } from "../auth/auth.js";
|
||||
import type { PrincipalContext } from "../auth/principal.js";
|
||||
import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
||||
import type { ListModelsFn } from "./meta.js";
|
||||
import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
import { workspaceRegistryRecoveryIdentity, workspaceRegistrySnapshotReader, acquireWorkspaceSessionReadersShared, type WorkspaceRegistry } from "../workspaces/registry.js";
|
||||
import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js";
|
||||
import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js";
|
||||
|
||||
@@ -322,6 +322,8 @@ export function sessionRoutes(
|
||||
});
|
||||
}
|
||||
let revisionLease: Awaited<ReturnType<ReturnType<typeof snapshotReader>["acquireSessionRevision"]>> | undefined;
|
||||
let sessionReadersLease: Awaited<ReturnType<typeof acquireWorkspaceSessionReadersShared>> | undefined;
|
||||
let sessionReadersHandedOff = false;
|
||||
let manifestPersisted = false;
|
||||
try {
|
||||
let workspaceConfigPath: string | undefined;
|
||||
@@ -356,6 +358,10 @@ export function sessionRoutes(
|
||||
});
|
||||
}
|
||||
}
|
||||
if (workspaceId && "rootLeaseFactory" in (d.workspaceRegistry as object)) {
|
||||
try { sessionReadersLease = await acquireWorkspaceSessionReadersShared(d.workspaceRegistry, workspaceId); }
|
||||
catch { return storageFailure(reply); }
|
||||
}
|
||||
const provider = b.provider ?? s.provider;
|
||||
const model = b.model ?? s.model;
|
||||
const thinking = b.thinking ?? s.thinking;
|
||||
@@ -427,12 +433,15 @@ export function sessionRoutes(
|
||||
author: principal.displayName ?? principal.subject,
|
||||
principal,
|
||||
question: b.question,
|
||||
...(sessionReadersLease ? { sessionReadersLease } : {}),
|
||||
};
|
||||
let runtimeOptions = options;
|
||||
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
||||
try {
|
||||
runtimeOptions = await optionsWithRuntimeConfig(runner, workspaceConfigPath, options);
|
||||
runtimeOptions = await optionsWithRuntimeConfig(runner, workspaceConfigPath, { ...options, ...(sessionReadersLease ? { sessionReadersLease } : {}) });
|
||||
rt = d.mgr.createFor(id, runtimeOptions);
|
||||
sessionReadersHandedOff = true;
|
||||
sessionReadersHandedOff = true;
|
||||
bindRuntime(id, rt, runner, workspaceConfigPath);
|
||||
} catch (error) {
|
||||
if (rt) d.mgr.teardownIfCurrent(id, rt);
|
||||
@@ -453,6 +462,7 @@ export function sessionRoutes(
|
||||
);
|
||||
return { id };
|
||||
} finally {
|
||||
if (sessionReadersLease && !sessionReadersHandedOff) await sessionReadersLease.close().catch(() => undefined);
|
||||
if (revisionLease && !manifestPersisted) {
|
||||
await revisionLease.abort().catch((error: unknown) => {
|
||||
console.error(
|
||||
@@ -609,6 +619,10 @@ export function sessionRoutes(
|
||||
boundRuntimes.delete(stoppedId);
|
||||
}
|
||||
|
||||
let sessionReadersLease: Awaited<ReturnType<typeof acquireWorkspaceSessionReadersShared>> | undefined;
|
||||
let sessionReadersHandedOff = false;
|
||||
try { if (saved.workspace_id && "rootLeaseFactory" in (d.workspaceRegistry as object)) sessionReadersLease = await acquireWorkspaceSessionReadersShared(d.workspaceRegistry, saved.workspace_id); }
|
||||
catch { return storageFailure(reply); }
|
||||
let rt: ReturnType<PiProcessManager["createFor"]> | undefined;
|
||||
try {
|
||||
if (current) {
|
||||
@@ -625,6 +639,7 @@ export function sessionRoutes(
|
||||
if (boundRuntimes.get(id) === rt) boundRuntimes.delete(id);
|
||||
d.mgr.teardownIfCurrent(id, rt);
|
||||
}
|
||||
if (sessionReadersLease && !sessionReadersHandedOff) await sessionReadersLease.close().catch(() => undefined);
|
||||
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
||||
}
|
||||
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
import { fsAtInternal } from "./workspace-fs-at-internal.js";
|
||||
import { rootLeaseInternals } from "./workspace-lock-root-lease-runtime.js";
|
||||
import { createHash, randomBytes } from "node:crypto";
|
||||
import {
|
||||
WorkspaceFsAtV1,
|
||||
type OwnedWorkspaceFsAtRegularFile,
|
||||
type OwnedWorkspaceFsAtDirectory,
|
||||
workspaceFsAtFriend,
|
||||
} from "./workspace-fs-at.js";
|
||||
import type { RuntimeConfigLease } from "./runtime-config-lease.js";
|
||||
import {
|
||||
@@ -12,8 +13,7 @@ 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; }
|
||||
@@ -59,29 +59,29 @@ function validFile(file: OwnedWorkspaceFsAtRegularFile, max = MAX_FILE_BYTES): v
|
||||
export class PreprocessingStateStore {
|
||||
constructor(private readonly rootLease: VerifiedWorkspaceLockRootLease) {}
|
||||
private dirs<T>(action: (fs: WorkspaceFsAtV1, base: OwnedWorkspaceFsAtDirectory, jobs: OwnedWorkspaceFsAtDirectory, candidates: OwnedWorkspaceFsAtDirectory, reviews: OwnedWorkspaceFsAtDirectory) => T): T {
|
||||
return workspaceRootInternals.withRoot(this.rootLease, (fs, root) => {
|
||||
return rootLeaseInternals.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); }
|
||||
const mkdir = (parent: OwnedWorkspaceFsAtDirectory, name: string): OwnedWorkspaceFsAtDirectory => { try { return fsAtInternal.openDirectory(parent, name); } catch { try { fsAtInternal.mkdir(parent, name); } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; } return fsAtInternal.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]) fsAtInternal.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 file(directory: OwnedWorkspaceFsAtDirectory, name: string, access: "read"|"create"): OwnedWorkspaceFsAtRegularFile {
|
||||
try { const f = workspaceFsAtFriend.openFile(directory, name, access); validFile(f); return f; } catch { throw fail(); }
|
||||
try { const f = fsAtInternal.openFile(directory, name, access); validFile(f); return f; } catch { 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 {} }
|
||||
const f = this.file(directory, name, "read"); try { const bytes = fsAtInternal.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 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);
|
||||
try { fsAtInternal.writeFile(f, bytes); fsAtInternal.fsyncFile(f); } catch { throw fail(); } finally { try { f.close(); } catch {} }
|
||||
fsAtInternal.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(); }
|
||||
const tmp = `.${name}.tmp-${process.pid}-${randomBytes(8).toString("hex")}`; try { this.writeExclusive(directory, tmp, bytes); fsAtInternal.rename(directory, tmp, name, true); fsAtInternal.fsyncDirectory(directory); } catch { try { fsAtInternal.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 rawState(jobs: OwnedWorkspaceFsAtDirectory, runId: string): PreprocessingRunStateV1 { id(runId); const names = fsAtInternal.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<PreprocessingRunStateV1> {
|
||||
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();
|
||||
@@ -122,7 +122,7 @@ export class WorkspaceWriterLockCapability {
|
||||
static [INTERNAL_STATE](id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, writer: WorkspaceRootLock) { return new WorkspaceWriterLockCapability(id, identity, root, writer); }
|
||||
async runUnderSessionReadersExclusive<T>(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise<T>): Promise<T> {
|
||||
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(); }
|
||||
try { lock = await rootLeaseInternals.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; }
|
||||
@@ -136,11 +136,11 @@ export class WorkspaceWriterLockCapability {
|
||||
try {
|
||||
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) }));
|
||||
return await rootLeaseInternals.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; }
|
||||
}
|
||||
}
|
||||
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 assertCapability(cap: WorkspaceWriterLockCapability): WriterState { const state = writerState.get(cap); if (!state || !state.live || state.settled || state.poisoned) throw fail(); rootLeaseInternals.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<void> { 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); }
|
||||
@@ -156,7 +156,7 @@ export class OrderedWorkspaceWriterCapabilitySet {
|
||||
private constructor(caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { orderedState.set(this, { live: true, caps }); }
|
||||
static [INTERNAL_STATE](caps: Map<CanonicalWorkspaceId, WorkspaceWriterLockCapability>) { return new OrderedWorkspaceWriterCapabilitySet(caps); }
|
||||
get workspaceIds(): readonly CanonicalWorkspaceId[] { const state = orderedState.get(this); if (!state?.live) throw fail(); return [...state.caps.keys()]; }
|
||||
async forWorkspace<T>(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T> { 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 forWorkspace<T>(workspaceId: CanonicalWorkspaceId, action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<T> { const state = orderedState.get(this); if (!state?.live) throw fail(); const cap = state.caps.get(workspaceId); if (!cap) throw fail(); return action({ workspaceId, rootLease: Reflect.construct(BorrowedVerifiedWorkspaceLockRootLease, [cap.rootIdentity]) as BorrowedVerifiedWorkspaceLockRootLease, writerCapability: cap }); }
|
||||
async forEachWorkspace<T>(action: (lease: OrderedWorkspaceCapability) => Promise<T>): Promise<readonly T[]> { 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); }
|
||||
@@ -168,7 +168,7 @@ export async function runUnderOrderedWorkspaceWriterLocks<T>(rootLeases: readonl
|
||||
const root = source.transfer();
|
||||
let writer: WorkspaceRootLock | undefined;
|
||||
try {
|
||||
writer = await workspaceRootInternals.acquireWriter(root);
|
||||
writer = await rootLeaseInternals.acquireWriter(root);
|
||||
caps.push(makeWriterCapability(root.identity.workspaceId, root.identity, root, writer));
|
||||
} catch (error) {
|
||||
try { writer?.close(); } catch {}
|
||||
|
||||
@@ -145,6 +145,14 @@ function registryContext(registry: WorkspaceRegistry): WorkspaceRegistryContext
|
||||
return context;
|
||||
}
|
||||
|
||||
export async function acquireWorkspaceSessionReadersShared(registry: WorkspaceRegistry, workspaceId: string) {
|
||||
const context = registryContext(registry);
|
||||
const input = context.rootLeaseFactory.canonicalInput(workspaceId);
|
||||
const root = await context.rootLeaseFactory.acquireOrProvision(input);
|
||||
try { return await root.acquireSessionReadersShared(); }
|
||||
catch (error) { await root.close().catch(() => undefined); throw error; }
|
||||
}
|
||||
|
||||
export function workspaceRegistryRecoveryIdentity(registry: WorkspaceRegistry): RegistryBootstrapRecoveryIdentityV1 {
|
||||
return registryContext(registry).installationIdentity;
|
||||
}
|
||||
@@ -165,7 +173,7 @@ export function createWorkspaceRegistry(config: WorkspaceRegistryConfig, deps: W
|
||||
mkdirSync(sessionsRoot, { recursive: true, mode: 0o700 });
|
||||
const rootLeaseFactory = deps.rootLeaseFactory ?? new VerifiedWorkspaceLockRootLeaseFactory({
|
||||
workspaceFsAt: new WorkspaceFsAtV1(), installationId: config.installationId,
|
||||
sessionsRootFromValidatedInstallationConfig: realpathSync(sessionsRoot), serviceUid: process.getuid?.() ?? 0,
|
||||
sessionsRootFromValidatedInstallationConfig: sessionsRoot, serviceUid: process.getuid?.() ?? 0,
|
||||
provisionedWorkspaceMode: 0o700,
|
||||
});
|
||||
const hash = (value: string) => createHash("sha256").update(value).digest("hex");
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
const fdPath = (fd:number): string => `${process.platform === "linux" ? "/proc/self/fd" : "/dev/fd"}/${fd}`;
|
||||
import { createRequire } from "node:module";
|
||||
import { closeSync, openSync, readSync, writeSync, fsyncSync, fstatSync, readdirSync, renameSync, unlinkSync } from "node:fs";
|
||||
import { spawn } from "node:child_process";
|
||||
import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1, NativeWorkspaceFsAtComponentV1 } from "../native/workspace-fs-at-binding.js";
|
||||
const require = createRequire(import.meta.url);
|
||||
const binding = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node") as WorkspaceFsAtBindingV1;
|
||||
export interface WorkspaceFsAtStatV1 extends NativeWorkspaceFsAtStatV1 {}
|
||||
export type LockFileName = "writer.lock" | "session-readers.lock";
|
||||
export type WorkspaceFlockKindV1 = "shared" | "exclusive";
|
||||
export type WorkspaceFlockWaitV1 = "blocking" | "nonblocking";
|
||||
const component = (value: string): NativeWorkspaceFsAtComponentV1 => {
|
||||
if (typeof value !== "string" || value.length === 0 || Buffer.byteLength(value, "utf8") > 255 || value !== value.trim() || value === "." || value === ".." || value.includes("/") || value.includes("\\") || value.includes("\0")) throw new Error("invalid path component");
|
||||
return value as NativeWorkspaceFsAtComponentV1;
|
||||
};
|
||||
const normalizeError = (e: unknown): Error => e instanceof Error ? e : new Error(String(e));
|
||||
const rawHandles = new WeakMap<object, NativeWorkspaceFsAtHandleV1>();
|
||||
const borrowing = new WeakMap<object, number>();
|
||||
const live = new WeakMap<object, boolean>();
|
||||
const rawOf = (value: object): NativeWorkspaceFsAtHandleV1 => {
|
||||
if (!rawHandles.has(value) || live.get(value) !== true) throw Object.assign(new Error("workspace descriptor is closed"), { code: "ERR_WORKSPACE_FS_AT_HANDLE_CLOSED" });
|
||||
return rawHandles.get(value)!;
|
||||
};
|
||||
abstract class Owned {
|
||||
constructor(raw: NativeWorkspaceFsAtHandleV1, readonly opened: WorkspaceFsAtStatV1) { rawHandles.set(this, raw); borrowing.set(this, 0); live.set(this, true); }
|
||||
stat(): WorkspaceFsAtStatV1 { if (live.get(this) !== true) throw Object.assign(new Error("workspace descriptor is closed"), { code: "ERR_WORKSPACE_FS_AT_HANDLE_CLOSED" }); return this.opened; }
|
||||
close(): void { if (live.get(this) !== true) return; if ((borrowing.get(this) ?? 0) !== 0) throw Object.assign(new Error("workspace descriptor is borrowed"), { code: "ERR_WORKSPACE_FS_AT_BORROWED" }); live.set(this, false); binding.close(rawHandles.get(this)!); }
|
||||
}
|
||||
export class OwnedWorkspaceFsAtDirectory extends Owned {}
|
||||
export class OwnedWorkspaceFsAtRegularFile extends Owned {}
|
||||
const wrapDirectory = (result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtDirectory => {
|
||||
try { if ((result.openedStat.mode & 0o170000) !== 0o040000) throw new Error("not a directory"); return new OwnedWorkspaceFsAtDirectory(result.handle, result.openedStat); }
|
||||
catch (e) { try { binding.close(result.handle); } catch {} throw e; }
|
||||
};
|
||||
const wrapFile = (result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtRegularFile => {
|
||||
try { const st=result.openedStat; if ((st.mode&0o170000)!==0o100000 || (st.mode&0o7777)!==0o600 || st.uid!==(process.getuid?.()??st.uid) || st.nlink!==1n) throw new Error("invalid regular file"); return new OwnedWorkspaceFsAtRegularFile(result.handle,st); }
|
||||
catch (e) { try { binding.close(result.handle); } catch {} throw e; }
|
||||
};
|
||||
const wrapLock = wrapFile;
|
||||
const withBorrowedFd = <T>(value: object, action: (fd:number)=>T): T => {
|
||||
const raw=rawOf(value), n=borrowing.get(value)??0; borrowing.set(value,n+1);
|
||||
try { return action(binding.fdNumberForSynchronousBorrow(raw)); }
|
||||
finally { borrowing.set(value,n); }
|
||||
};
|
||||
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:rawOf(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:rawOf(parent),name:component(name),kind:"regular_lock",createMode:0o600}));
|
||||
}
|
||||
mkdirAt(parent: OwnedWorkspaceFsAtDirectory, name:string, mode:0o700):void { binding.mkdirat(rawOf(parent),component(name),mode); }
|
||||
statAtNoFollow(parent:OwnedWorkspaceFsAtDirectory,name:string):WorkspaceFsAtStatV1 { return binding.fstatat(rawOf(parent),component(name)); }
|
||||
fsyncDirectory(directory:OwnedWorkspaceFsAtDirectory):void { binding.fsyncDirectory(rawOf(directory)); }
|
||||
flockOwnedLock(owned:OwnedWorkspaceFsAtRegularFile,kind:WorkspaceFlockKindV1,wait:WorkspaceFlockWaitV1):void { const ext:{flockSync(fd:number,operation:string):void}=require("fs-ext"); withBorrowedFd(owned,fd=>ext.flockSync(fd,kind==="shared"?(wait==="blocking"?"sh":"shnb"):(wait==="blocking"?"ex":"exnb"))); }
|
||||
}
|
||||
export const fsAtInternal = {
|
||||
raw: (v: object) => rawOf(v),
|
||||
withFd: withBorrowedFd,
|
||||
openDirectory(parent: OwnedWorkspaceFsAtDirectory,name:string):OwnedWorkspaceFsAtDirectory { return wrapDirectory(binding.openat({parent:rawOf(parent),name:component(name),kind:"directory",createMode:0})); },
|
||||
mkdir(parent:OwnedWorkspaceFsAtDirectory,name:string):void { binding.mkdirat(rawOf(parent),component(name),0o700); },
|
||||
openFile(parent:OwnedWorkspaceFsAtDirectory,name:string,access:"read"|"create"):OwnedWorkspaceFsAtRegularFile {
|
||||
// regular_lock is the sole native regular-file operation. Every state file is 0600,
|
||||
// and the no-follow open still occurs relative to the retained directory handle.
|
||||
if(access==="create") { try { binding.fstatat(rawOf(parent),component(name)); throw Object.assign(new Error("exists"),{code:"EEXIST"}); } catch(e) { if((e as NodeJS.ErrnoException).code!=="ENOENT") throw e; } }
|
||||
else binding.fstatat(rawOf(parent), component(name));
|
||||
return wrapFile(binding.openat({parent:rawOf(parent),name:component(name),kind:"regular_lock",createMode:0o600}));
|
||||
},
|
||||
readFile(file:OwnedWorkspaceFsAtRegularFile,max:number):Uint8Array { return withBorrowedFd(file,fd=>{const st=fstatSync(fd);if(st.size>max)throw new Error("file too large");const out=Buffer.alloc(st.size);let off=0;while(off<out.length){const n=readSync(fd,out,off,out.length-off,off);if(n<=0)throw new Error("short read");off+=n;}return out;}); },
|
||||
writeFile(file:OwnedWorkspaceFsAtRegularFile,bytes:Uint8Array):void { withBorrowedFd(file,fd=>{let off=0;while(off<bytes.length){const n=writeSync(fd,bytes,off,bytes.length-off,off);if(n<=0)throw new Error("short write");off+=n;}}); },
|
||||
fsyncFile(file:OwnedWorkspaceFsAtRegularFile):void { withBorrowedFd(file,fd=>fsyncSync(fd)); },
|
||||
rename(parent:OwnedWorkspaceFsAtDirectory,from:string,to:string,replace:boolean):void { if(!replace){ try{ withBorrowedFd(parent,fd=>{ readdirSync(fdPath(fd)); }); }catch{} } withBorrowedFd(parent,fd=>renameSync(`${fdPath(fd)}/${component(from)}`,`${fdPath(fd)}/${component(to)}`)); },
|
||||
unlink(parent:OwnedWorkspaceFsAtDirectory,name:string):void { withBorrowedFd(parent,fd=>unlinkSync(`${fdPath(fd)}/${component(name)}`)); },
|
||||
listDirectory(directory:OwnedWorkspaceFsAtDirectory):readonly string[] { return withBorrowedFd(directory,fd=>readdirSync(fdPath(fd))); },
|
||||
fsyncDirectory(directory:OwnedWorkspaceFsAtDirectory):void { binding.fsyncDirectory(rawOf(directory)); },
|
||||
assertPath(root:OwnedWorkspaceFsAtDirectory,lock:OwnedWorkspaceFsAtRegularFile,name:LockFileName,identity:{device:bigint;inode:bigint}):void { const st=binding.fstatat(rawOf(root),component(name)), opened=lock.stat(); if(st.device!==opened.device||st.inode!==opened.inode||st.mode!==opened.mode||st.uid!==opened.uid||st.nlink!==opened.nlink) throw new Error("preprocessing_conflict: lock pathname identity changed"); },
|
||||
spawn(writer:OwnedWorkspaceFsAtRegularFile,root:OwnedWorkspaceFsAtDirectory,executable:string,args:readonly string[],environment?:NodeJS.ProcessEnv):Promise<{exitCode:number;stdout:Uint8Array;stderr:Uint8Array}> {
|
||||
return withBorrowedFd(writer,wfd=>withBorrowedFd(root,rfd=>new Promise((resolve,reject)=>{const child=spawn(executable,[...args],{stdio:["ignore","pipe","pipe",wfd,rfd],env:{...(environment??process.env),THOTH_WORKSPACE_CAPABILITY_REQUIRED:"1"}});const out:Buffer[]=[];const err:Buffer[]=[];child.stdout?.on("data",(x:Buffer)=>out.push(x));child.stderr?.on("data",(x:Buffer)=>err.push(x));child.once("error",reject);child.once("close",code=>resolve({exitCode:code??1,stdout:Buffer.concat(out),stderr:Buffer.concat(err)}));})));
|
||||
},
|
||||
};
|
||||
@@ -1,130 +1,2 @@
|
||||
import { createRequire } from "node:module";
|
||||
import { closeSync, openSync } from "node:fs";
|
||||
import { spawn } from "node:child_process";
|
||||
import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1, NativeWorkspaceFsAtComponentV1 } from "../native/workspace-fs-at-binding.js";
|
||||
const require = createRequire(import.meta.url);
|
||||
const binding = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node") as WorkspaceFsAtBindingV1 & {
|
||||
withFd<T>(handle: NativeWorkspaceFsAtHandleV1, action: (fd: number) => T): T;
|
||||
duplicateForChildStdio(writer: NativeWorkspaceFsAtHandleV1, root: NativeWorkspaceFsAtHandleV1, writerFd: number, rootFd: number): void;
|
||||
};
|
||||
export interface WorkspaceFsAtStatV1 extends NativeWorkspaceFsAtStatV1 {}
|
||||
export type LockFileName = "writer.lock" | "session-readers.lock";
|
||||
export type WorkspaceFlockKindV1 = "shared" | "exclusive";
|
||||
export type WorkspaceFlockWaitV1 = "blocking" | "nonblocking";
|
||||
|
||||
function component(value: string): NativeWorkspaceFsAtComponentV1 {
|
||||
if (typeof value !== "string" || value.length === 0 || Buffer.byteLength(value, "utf8") > 255 || value !== value.trim() || value === "." || value === ".." || value.includes("/") || value.includes("\0")) throw new Error("invalid path component");
|
||||
return value as NativeWorkspaceFsAtComponentV1;
|
||||
}
|
||||
function normalizeError(error: unknown): Error {
|
||||
if (error instanceof Error) return error;
|
||||
return new Error(String(error));
|
||||
}
|
||||
const INTERNAL = Symbol("workspace-fs-at-owned");
|
||||
const rawHandles = new WeakMap<object, NativeWorkspaceFsAtHandleV1>();
|
||||
const borrowing = new WeakMap<object, number>();
|
||||
abstract class Owned {
|
||||
private live = true;
|
||||
private borrowing = 0;
|
||||
protected constructor(raw: NativeWorkspaceFsAtHandleV1, readonly opened: WorkspaceFsAtStatV1, token: symbol) { if (token !== INTERNAL) throw new TypeError("private workspace descriptor"); rawHandles.set(this, raw); borrowing.set(this, 0); }
|
||||
stat(): WorkspaceFsAtStatV1 { if (!this.live) throw new Error("workspace descriptor is closed"); return this.opened; }
|
||||
close(): void {
|
||||
if (!this.live) return;
|
||||
if ((borrowing.get(this) ?? 0) !== 0) throw Object.assign(new Error("workspace descriptor is borrowed"), { code: "ERR_WORKSPACE_FS_AT_BORROWED" });
|
||||
this.live = false;
|
||||
try { binding.close(rawHandles.get(this)!); } catch (error) { throw normalizeError(error); }
|
||||
}
|
||||
}
|
||||
export class OwnedWorkspaceFsAtDirectory extends Owned {
|
||||
private constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); }
|
||||
static [INTERNAL](raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1): OwnedWorkspaceFsAtDirectory { return new OwnedWorkspaceFsAtDirectory(raw, stat, INTERNAL); }
|
||||
}
|
||||
export class OwnedWorkspaceFsAtRegularFile extends Owned {
|
||||
private constructor(raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1, token: symbol) { super(raw, stat, token); }
|
||||
static [INTERNAL](raw: NativeWorkspaceFsAtHandleV1, stat: WorkspaceFsAtStatV1): OwnedWorkspaceFsAtRegularFile { return new OwnedWorkspaceFsAtRegularFile(raw, stat, INTERNAL); }
|
||||
}
|
||||
function wrapDirectory(result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtDirectory {
|
||||
try { if ((result.openedStat.mode & 0o170000) !== 0o040000) throw new Error("not a directory"); return OwnedWorkspaceFsAtDirectory[INTERNAL](result.handle, result.openedStat); }
|
||||
catch (error) { try { binding.close(result.handle); } catch { /* preserve conversion error */ } throw error; }
|
||||
}
|
||||
function wrapLock(result: {handle: NativeWorkspaceFsAtHandleV1; openedStat: WorkspaceFsAtStatV1}): OwnedWorkspaceFsAtRegularFile {
|
||||
try {
|
||||
const st = result.openedStat;
|
||||
if ((st.mode & 0o170000) !== 0o100000 || (st.mode & 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; }
|
||||
}
|
||||
function rawDirectory(value: OwnedWorkspaceFsAtDirectory): NativeWorkspaceFsAtHandleV1 { if (!rawHandles.has(value)) throw new Error("workspace descriptor is closed"); return rawHandles.get(value)!; }
|
||||
function withLockFd<T>(value: OwnedWorkspaceFsAtRegularFile, action: (fd: number) => T): T {
|
||||
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)!, action as (fd: number) => void) as T; } 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 }));
|
||||
}
|
||||
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)); }
|
||||
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));
|
||||
}
|
||||
}
|
||||
|
||||
// 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 }>;
|
||||
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");
|
||||
},
|
||||
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)); }
|
||||
export { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, OwnedWorkspaceFsAtRegularFile } from "./workspace-fs-at-internal.js";
|
||||
export type { WorkspaceFsAtStatV1, LockFileName, WorkspaceFlockKindV1, WorkspaceFlockWaitV1 } from "./workspace-fs-at-internal.js";
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
import { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, OwnedWorkspaceFsAtRegularFile, fsAtInternal, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at-internal.js";
|
||||
export type WorkspaceLockRootIdentityRuntime = { readonly device: bigint; readonly inode: bigint; readonly workspaceId: string };
|
||||
export 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}> };
|
||||
export type RootState = { fs:WorkspaceFsAtV1; parent:OwnedWorkspaceFsAtDirectory; root:OwnedWorkspaceFsAtDirectory; serviceUid:number; owner:symbol; identity:WorkspaceLockRootIdentityRuntime; assertParent:()=>void; live:boolean };
|
||||
export const rootState = new WeakMap<object,RootState>();
|
||||
const writerAdmissions = new Set<string>();
|
||||
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:WorkspaceLockRootIdentityRuntime,private readonly fs:WorkspaceFsAtV1,private readonly admissionKey?:string){}
|
||||
assertPath():void { if(!this.live) throw new Error("preprocessing_conflict"); fsAtInternal.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 fsAtInternal.spawn(this.lock,root,executable,args,environment); }
|
||||
close():void { if(!this.live)return; this.live=false; try{this.lock.close();}finally{if(this.admissionKey)writerAdmissions.delete(this.admissionKey);} }
|
||||
}
|
||||
const conflict=()=>Object.assign(new Error("preprocessing_conflict"),{name:"PreprocessingConflictError"});
|
||||
export const rootLeaseInternals = {
|
||||
assertLive(root:object):void { const s=rootState.get(root); if(!s||!s.live)throw conflict(); s.assertParent(); const retained=s.root.stat(); const named=s.fs.statAtNoFollow(s.parent,s.identity.workspaceId as string); if(retained.device!==s.identity.device||retained.inode!==s.identity.inode||named.device!==s.identity.device||named.inode!==s.identity.inode)throw conflict(); },
|
||||
withRoot<T>(root:object,action:(fs:WorkspaceFsAtV1,directory:OwnedWorkspaceFsAtDirectory)=>T):T { this.assertLive(root); const s=rootState.get(root);if(!s)throw conflict();return action(s.fs,s.root); },
|
||||
acquireShared(root:object):InternalLock { this.assertLive(root); const s=rootState.get(root)!; const owned=s.fs.openOrCreateLockAt(s.root,"session-readers.lock",0o600); const lock=new RootLock(s.root,owned,"session-readers.lock",s.identity,s.fs); try{lock.flock("shared","nonblocking");return lock;}catch{try{lock.close();}catch{}throw conflict();}},
|
||||
async acquireWriter(root:object):Promise<InternalLock>{this.assertLive(root);const s=rootState.get(root)!;const key=`${s.identity.device}:${s.identity.inode}`;if(writerAdmissions.has(key))throw conflict();const owned=s.fs.openOrCreateLockAt(s.root,"writer.lock",0o600);const lock=new RootLock(s.root,owned,"writer.lock",s.identity,s.fs,key);try{lock.flock("exclusive","nonblocking");writerAdmissions.add(key);return lock;}catch{try{lock.close();}catch{}throw conflict();}},
|
||||
async acquireReadersExclusive(root:object):Promise<InternalLock>{this.assertLive(root);const s=rootState.get(root)!;const owned=s.fs.openOrCreateLockAt(s.root,"session-readers.lock",0o600);const lock=new RootLock(s.root,owned,"session-readers.lock",s.identity,s.fs);try{lock.flock("exclusive","nonblocking");return lock;}catch{try{lock.close();}catch{}throw conflict();}},
|
||||
};
|
||||
@@ -1,5 +1,6 @@
|
||||
import { lstatSync, realpathSync } from "node:fs";
|
||||
import { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, OwnedWorkspaceFsAtRegularFile, workspaceFsAtFriend, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js";
|
||||
import { fsAtInternal } from "./workspace-fs-at-internal.js";
|
||||
import { WorkspaceFsAtV1, OwnedWorkspaceFsAtDirectory, OwnedWorkspaceFsAtRegularFile, type WorkspaceFsAtStatV1, type WorkspaceFlockKindV1, type WorkspaceFlockWaitV1 } from "./workspace-fs-at.js";
|
||||
import { rootState, rootLeaseInternals, type InternalLock } from "./workspace-lock-root-lease-runtime.js";
|
||||
declare const canonicalWorkspaceIdBrand: unique symbol;
|
||||
declare const revision40Brand: unique symbol;
|
||||
declare const sha256HexBrand: unique symbol;
|
||||
@@ -9,26 +10,12 @@ export type Sha256Hex = string & { readonly [sha256HexBrand]: true };
|
||||
export interface WorkspaceLockRootIdentityV1 { readonly schemaVersion: 1; readonly workspaceId: CanonicalWorkspaceId; readonly device: bigint; readonly inode: bigint; }
|
||||
|
||||
const INPUT_BRAND = new WeakMap<object, symbol>();
|
||||
const ROOT_INTERNAL = Symbol("workspace-root-internal");
|
||||
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; }
|
||||
|
||||
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<string>();
|
||||
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 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.admissionKey) writerAdmissions.delete(this.admissionKey); } }
|
||||
}
|
||||
type RootState = { fs: WorkspaceFsAtV1; parent: OwnedWorkspaceFsAtDirectory; root: OwnedWorkspaceFsAtDirectory; serviceUid: number; owner: symbol; identity: WorkspaceLockRootIdentityV1; live: boolean };
|
||||
const rootState = new WeakMap<object, RootState>();
|
||||
|
||||
export class CanonicalWorkspaceLockRootInput {
|
||||
private constructor(readonly workspaceId: CanonicalWorkspaceId) {}
|
||||
}
|
||||
@@ -38,7 +25,6 @@ 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;
|
||||
@@ -47,40 +33,53 @@ export class WorkspaceSessionReadersLockLease {
|
||||
async close(): Promise<void> { if (!this.live) return; this.live = false; let failure: unknown; try { this.lock.close(); } catch (e) { failure = e; } try { this.root.close(); } catch (e) { failure ??= e; } if (failure) throw conflict(); }
|
||||
}
|
||||
|
||||
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"); } }
|
||||
function assertRootLive(root: VerifiedWorkspaceLockRootLease): void { const state = rootState.get(root); if (!state || state.live !== true) throw conflict(); try { state.assertParent(); 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(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;
|
||||
private constructor(fs: WorkspaceFsAtV1, parent: OwnedWorkspaceFsAtDirectory, root: OwnedWorkspaceFsAtDirectory, identity: WorkspaceLockRootIdentityV1, serviceUid: number, owner: symbol, assertParent: () => void) { this.fs = fs; this.parent = parent; this.root = root; this.serviceUid = serviceUid; this.owner = owner; this.identity = identity; this.assertParent = assertParent; rootState.set(this, { fs, parent, root, serviceUid, owner, identity, assertParent, live: true }); }
|
||||
private readonly assertParent: () => void; private readonly fs: WorkspaceFsAtV1; private readonly parent: OwnedWorkspaceFsAtDirectory; private readonly root: OwnedWorkspaceFsAtDirectory; private readonly serviceUid: number; private readonly owner: symbol;
|
||||
readonly identity: WorkspaceLockRootIdentityV1;
|
||||
borrow<T>(action: (borrowed: BorrowedVerifiedWorkspaceLockRootLease) => Promise<T>): Promise<T> { 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<WorkspaceSessionReadersLockLease> {
|
||||
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; }
|
||||
assertRootLive(this); let lock: InternalLock | undefined;
|
||||
try { lock = rootLeaseInternals.acquireShared(this); assertRootLive(this); this.live = false; rootState.get(this)!.live = false; return Reflect.construct(WorkspaceSessionReadersLockLease, [lock, rootState.get(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 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, this.assertParent]) as VerifiedWorkspaceLockRootLease; }
|
||||
async close(): Promise<void> { if (!this.live) return; while (this.borrowed) await new Promise<void>(r => setTimeout(r, 1)); this.live = false; rootState.get(this)!.live = false; try { this.root.close(); } catch { throw conflict(); } }
|
||||
|
||||
}
|
||||
|
||||
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; }
|
||||
function makeRoot(fs: WorkspaceFsAtV1, parent: OwnedWorkspaceFsAtDirectory, root: OwnedWorkspaceFsAtDirectory, id: WorkspaceLockRootIdentityV1, uid: number, owner: symbol, assertParent: () => void): VerifiedWorkspaceLockRootLease { return Reflect.construct(VerifiedWorkspaceLockRootLease, [fs, parent, root, id, uid, owner, assertParent]) as VerifiedWorkspaceLockRootLease; }
|
||||
|
||||
export class VerifiedWorkspaceLockRootLeaseFactory {
|
||||
private readonly owner = Symbol("workspace-root-factory"); private readonly parent: OwnedWorkspaceFsAtDirectory; private provisionTail: Promise<void> = Promise.resolve();
|
||||
private readonly owner = Symbol("workspace-root-factory"); private readonly parent: OwnedWorkspaceFsAtDirectory; private readonly components: readonly string[]; private readonly parentIdentities: readonly {device: bigint; inode: bigint}[]; private provisionTail: Promise<void> = 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");
|
||||
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 canonical.split("/").filter(Boolean)) { const n = input.workspaceFsAt.openDirectoryAt(d, c); d.close(); d = n; } this.parent = d; this.checkParent(); }
|
||||
if (!configured.startsWith("/") || configured.split("/").some(c => c === "." || c === "..")) throw conflict("invalid sessions root");
|
||||
const components = configured.split("/").filter(Boolean);
|
||||
if (!components.length) throw conflict("invalid sessions root");
|
||||
let d = input.workspaceFsAt.openRoot(); const identities: {device: bigint; inode: bigint}[] = [];
|
||||
try { for (const c of components) { const n = input.workspaceFsAt.openDirectoryAt(d, c); const st = n.stat(); if ((st.mode & 0o170000) !== 0o040000) throw conflict("invalid sessions root"); identities.push({device: st.device, inode: st.inode}); d.close(); d = n; } this.parent = d; this.components = components; this.parentIdentities = identities; 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"); } }
|
||||
private checkParent(): void {
|
||||
let d: OwnedWorkspaceFsAtDirectory | undefined;
|
||||
try {
|
||||
d = this.input.workspaceFsAt.openRoot();
|
||||
for (let i = 0; i < this.components.length; i++) {
|
||||
const name = this.components[i]; const st = this.input.workspaceFsAt.statAtNoFollow(d, name);
|
||||
if (!sameIdentity(st, this.parentIdentities[i]) || (st.mode & 0o170000) !== 0o040000) throw conflict("installation root identity changed");
|
||||
const next = this.input.workspaceFsAt.openDirectoryAt(d, name); d.close(); d = next;
|
||||
}
|
||||
const retained = this.parent.stat(); if (!sameIdentity(retained, this.parentIdentities[this.parentIdentities.length - 1]) || !exactParent(retained, 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"); }
|
||||
finally { try { d?.close(); } catch {} }
|
||||
}
|
||||
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<VerifiedWorkspaceLockRootLease> { 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 async open(input: CanonicalWorkspaceLockRootInput): Promise<VerifiedWorkspaceLockRootLease> { 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.checkParent()); 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<VerifiedWorkspaceLockRootLease> { return this.open(input); }
|
||||
async acquireOrProvision(input: CanonicalWorkspaceLockRootInput): Promise<VerifiedWorkspaceLockRootLease> {
|
||||
@@ -92,14 +91,3 @@ export class VerifiedWorkspaceLockRootLeaseFactory {
|
||||
} 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) },
|
||||
withRoot<T>(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<InternalLock> { 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<InternalLock> { 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()); } },
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user