diff --git a/backend/native/workspace-fs-at/workspace_fs_at.cc b/backend/native/workspace-fs-at/workspace_fs_at.cc index 64385f7c..b537fcfe 100644 --- a/backend/native/workspace-fs-at/workspace_fs_at.cc +++ b/backend/native/workspace-fs-at/workspace_fs_at.cc @@ -22,8 +22,8 @@ static bool str(napi_env env,napi_value v,std::string& out){ size_t n; if(napi_g 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,fd,dir,false}; 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 napi_value openatFn(napi_env env,napi_callback_info info){ size_t argc=1; napi_value a[1]; if(napi_get_cb_info(env,info,&argc,a,nullptr,nullptr)!=napi_ok||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&&(name.empty()||name.size()>255||name=="."||name==".."||name.find('/')!=std::string::npos||name.find('\0')!=std::string::npos))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); int fd; if(root)fd=::open("/",O_RDONLY|O_DIRECTORY|O_CLOEXEC|O_NOFOLLOW); else if(kind=="directory"&&mode==0)fd=::openat(ph->fd,name.c_str(),O_RDONLY|O_DIRECTORY|O_CLOEXEC|O_NOFOLLOW); else if(kind=="regular_lock"&&mode==0600)fd=::openat(ph->fd,name.c_str(),O_RDWR|O_CREAT|O_CLOEXEC|O_NOFOLLOW,0600); else return fail(env,"openat",EINVAL); if(fd<0)return fail(env,"openat",errno); struct stat st; if(::fstat(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)){::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)return fail(env,"mkdirat",EINVAL);if(s.empty()||s.size()>255||s=="."||s==".."||s.find('/')!=std::string::npos)return fail(env,"mkdirat",EINVAL);if(::mkdirat(h->fd,s.c_str(),0700)<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)||s.empty()||s.size()>255||s=="."||s==".."||s.find('/')!=std::string::npos)return fail(env,"fstatat",EINVAL);struct stat st;if(::fstatat(h->fd,s.c_str(),&st,AT_SYMLINK_NOFOLLOW)<0)return fail(env,"fstatat",errno);return statObj(env,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)return fail(env,"mkdirat",EINVAL);if(s.empty()||s.size()>255||s=="."||s==".."||s.find('/')!=std::string::npos||s.find('\0')!=std::string::npos)return fail(env,"mkdirat",EINVAL);if(::mkdirat(h->fd,s.c_str(),0700)<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)||s.empty()||s.size()>255||s=="."||s==".."||s.find('/')!=std::string::npos||s.find('\0')!=std::string::npos)return fail(env,"fstatat",EINVAL);struct stat st;if(::fstatat(h->fd,s.c_str(),&st,AT_SYMLINK_NOFOLLOW)<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);if(::fsync(h->fd)<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);h->closed=true;int rc=::close(h->fd);if(rc<0)return fail(env,"close",errno);return nullptr;} static napi_value fdFn(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 x;napi_create_int32(env,h->fd,&x);return x;} diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index 431272a5..8fd6e425 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -1,48 +1,140 @@ -import { mkdir, readFile, rename, rm, writeFile } from "node:fs/promises"; import { createHash, randomBytes } from "node:crypto"; -import { dirname, join } from "node:path"; -import { WorkspaceFsAtV1, type OwnedWorkspaceFsAtRegularFile } from "./workspace-fs-at.js"; -import { BorrowedVerifiedWorkspaceLockRootLease, VerifiedWorkspaceLockRootLease, type CanonicalWorkspaceId, type WorkspaceLockRootIdentityV1 } from "./workspace-lock-root-lease.js"; -export interface ArtifactIdentity { readonly kind:string; readonly digest:string; readonly bytes:number; } -export interface PreprocessingRunStateV1 { readonly schemaVersion:1; readonly runId:string; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly operation:string; readonly phase:string; readonly artifacts:readonly ArtifactIdentity[]; readonly createdAt:string; readonly updatedAt:string; } -export interface CreateRunInput { readonly runId?:string; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly operation:string; readonly stateRoot:string; } -export interface ResumeRunInput { readonly runId:string; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly operation:string; readonly stateRoot:string; } -export type RunTransition = { readonly phase:string; readonly [key:string]:unknown }; -export interface FkReviewInput { readonly candidate:ArtifactIdentity; readonly reviewSha256:string; readonly annotationSha256?:string; } -export interface FkReviewRecordV1 extends FkReviewInput { readonly runId:string; readonly recordedAt:string; } -const digest=(x:Uint8Array|string)=>createHash("sha256").update(x).digest("hex"); -function fail(msg="preprocessing_conflict"):Error{const e=new Error(msg);e.name="PreprocessingConflictError";return e;} -function id(v:string):void{if(!/^[0-9a-f]{32}$/.test(v))throw fail("invalid run id");} +import { mkdir, readFile, rename, rm, writeFile, open as openFile } from "node:fs/promises"; +import { join } from "node:path"; +import { + WorkspaceFsAtV1, + type OwnedWorkspaceFsAtRegularFile, +} from "./workspace-fs-at.js"; +import { + BorrowedVerifiedWorkspaceLockRootLease, + VerifiedWorkspaceLockRootLease, + type CanonicalWorkspaceId, + type WorkspaceLockRootIdentityV1, +} from "./workspace-lock-root-lease.js"; + +export interface ArtifactIdentity { readonly kind: string; readonly digest: string; readonly bytes: number; } +export interface PreprocessingRunStateV1 { + readonly schemaVersion: 1; readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; + readonly revision: string; readonly operation: string; readonly phase: string; + readonly artifacts: readonly ArtifactIdentity[]; readonly createdAt: string; readonly updatedAt: string; + readonly [key: string]: unknown; +} +export interface CreateRunInput { readonly runId?: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; readonly stateRoot?: string; } +export interface ResumeRunInput { readonly runId: string; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly operation: string; readonly stateRoot?: string; } +export type RunTransition = { readonly phase: string; readonly [key: string]: unknown }; +export interface FkReviewInput { readonly candidate: ArtifactIdentity; readonly reviewSha256: string; readonly annotationSha256?: string; } +export interface FkReviewRecordV1 extends FkReviewInput { readonly runId: string; readonly recordedAt: string; } + +const digest = (x: Uint8Array | string) => createHash("sha256").update(x).digest("hex"); +const RUN_ID = /^[0-9a-f]{32}$/; +const SHA256 = /^[0-9a-f]{64}$/; +const REVISION = /^[0-9a-f]{40}$/; +const PHASES = ["created", "introspected", "fk_suggested", "fk_reviewed", "schema_indexed", "evidence_indexed", "lsh_built", "terminal"] as const; +function fail(msg = "preprocessing_conflict"): Error { const e = new Error(msg); e.name = "PreprocessingConflictError"; return e; } +function id(v: string): void { if (typeof v !== "string" || !RUN_ID.test(v)) throw fail("invalid run id"); } +function checkSha(v: string): void { if (typeof v !== "string" || !SHA256.test(v)) throw fail("invalid digest"); } +function checkedRoot(root: string | undefined): string { if (!root || typeof root !== "string" || !root.startsWith("/") || root.includes("\0")) throw fail(); return root; } + +async function durableJson(path: string, value: unknown): Promise { + const tmp = `${path}.tmp-${process.pid}-${randomBytes(5).toString("hex")}`; + try { + await writeFile(tmp, `${JSON.stringify(value)}\n`, { mode: 0o600, flag: "wx" }); + const handle = await openFile(tmp, "r"); try { await handle.sync(); } finally { await handle.close(); } + await rename(tmp, path); + const parent = await openFile(join(path, ".."), "r"); try { await parent.sync(); } finally { await parent.close(); } + } catch (error) { await rm(tmp, { force: true }).catch(() => undefined); throw error; } +} + export class PreprocessingStateStore { - private paths(input:{stateRoot:string;runId:string}){id(input.runId);return {dir:join(input.stateRoot,"preprocessing"),path:join(input.stateRoot,"preprocessing","jobs",`${input.runId}.json`),candidate:join(input.stateRoot,"preprocessing","fk-candidates",`${input.runId}.yaml`),review:join(input.stateRoot,"preprocessing","fk-reviews",`${input.runId}.json`)}} - async create(input:CreateRunInput):Promise{const runId=input.runId??Buffer.from(randomBytes(16)).toString("hex");id(runId);const p=this.paths({stateRoot:input.stateRoot,runId});await mkdir(join(p.dir,"jobs"),{recursive:true,mode:0o700});const now=new Date().toISOString();const state:PreprocessingRunStateV1={schemaVersion:1,runId,workspaceId:input.workspaceId,revision:input.revision,operation:input.operation,phase:"created",artifacts:[],createdAt:now,updatedAt:now};try{await writeFile(p.path,JSON.stringify(state)+"\n",{flag:"wx",mode:0o600});}catch{throw fail();}return state;} - async loadForResume(input:ResumeRunInput){const p=this.paths(input);let state:PreprocessingRunStateV1;try{state=JSON.parse(await readFile(p.path,"utf8"));}catch{throw fail();}if(state.schemaVersion!==1||state.runId!==input.runId||state.workspaceId!==input.workspaceId||state.revision!==input.revision||state.operation!==input.operation||!Array.isArray(state.artifacts))throw fail("preprocessing_resume_mismatch");return state;} - async transition(runId:string,transition:RunTransition):Promise{id(runId);throw new Error("transition requires a state-root bound store");} - async writeFkCandidate(runId:string,yaml:Uint8Array):Promise{id(runId);throw new Error("writeFkCandidate requires a state-root bound store");} - async recordFkReview(runId:string,review:FkReviewInput):Promise{id(runId);throw new Error("recordFkReview requires a state-root bound store");} + constructor(private readonly configuredRoot?: string) {} + private root(input?: string): string { return checkedRoot(input ?? this.configuredRoot); } + private paths(input: { stateRoot?: string; runId: string }) { + id(input.runId); const root = this.root(input.stateRoot); + const base = join(root, "preprocessing"); + return { base, jobs: join(base, "jobs"), path: join(base, "jobs", `${input.runId}.json`), + candidate: join(base, "fk-candidates", `${input.runId}.yaml`), review: join(base, "fk-reviews", `${input.runId}.json`) }; + } + async create(input: CreateRunInput): Promise { + const runId = input.runId ?? randomBytes(16).toString("hex"); id(runId); const p = this.paths({ stateRoot: input.stateRoot, runId }); + if (typeof input.workspaceId !== "string" || !/^[a-z][a-z0-9-]{2,62}$/.test(input.workspaceId) || !REVISION.test(input.revision) || !input.operation) throw fail(); + await mkdir(p.jobs, { recursive: true, mode: 0o700 }); await mkdir(join(p.base, "fk-candidates"), { recursive: true, mode: 0o700 }); await mkdir(join(p.base, "fk-reviews"), { recursive: true, mode: 0o700 }); + const now = new Date().toISOString(); const state: PreprocessingRunStateV1 = { schemaVersion: 1, runId, workspaceId: input.workspaceId, revision: input.revision, operation: input.operation, phase: "created", artifacts: [], createdAt: now, updatedAt: now }; + try { await writeFile(p.path, `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw fail(); } + return state; + } + async loadForResume(input: ResumeRunInput): Promise { + const p = this.paths(input); let value: unknown; + try { value = JSON.parse(await readFile(p.path, "utf8")); } catch { throw fail(); } + if (!this.validState(value) || value.runId !== input.runId || value.workspaceId !== input.workspaceId || value.revision !== input.revision || value.operation !== input.operation) throw fail("preprocessing_resume_mismatch"); + return value; + } + async load(input: ResumeRunInput): Promise { return this.loadForResume(input); } + async transition(runId: string, transition: RunTransition, stateRoot?: string): Promise { + const p = this.paths({ stateRoot, runId }); const current = await this.readRaw(p.path); + if (!this.validState(current) || current.runId !== runId || typeof transition !== "object" || !transition || typeof transition.phase !== "string") throw fail(); + const old = PHASES.indexOf(current.phase as never); const next = PHASES.indexOf(transition.phase as never); + if (next < 0 || (old >= 0 && next < old) || current.phase === "terminal") throw fail(); + const updated: PreprocessingRunStateV1 = { ...current, ...transition, schemaVersion: 1, runId, artifacts: current.artifacts, updatedAt: new Date().toISOString() }; + await durableJson(p.path, updated); return updated; + } + async writeFkCandidate(runId: string, yaml: Uint8Array, stateRoot?: string): Promise { + const p = this.paths({ stateRoot, runId }); if (!(yaml instanceof Uint8Array) || yaml.byteLength > 10 * 1024 * 1024) throw fail(); + const artifact = { kind: "fk-candidate", digest: digest(yaml), bytes: yaml.byteLength } satisfies ArtifactIdentity; + await mkdir(join(p.base, "fk-candidates"), { recursive: true, mode: 0o700 }); await writeFile(p.candidate, yaml, { flag: "wx", mode: 0o600 }); + return artifact; + } + async recordFkReview(runId: string, review: FkReviewInput, stateRoot?: string): Promise { + const p = this.paths({ stateRoot, runId }); if (!review || !review.candidate || !SHA256.test(review.candidate.digest)) throw fail(); checkSha(review.reviewSha256); if (review.annotationSha256) checkSha(review.annotationSha256); + const out: FkReviewRecordV1 = { ...review, runId, recordedAt: new Date().toISOString() }; await mkdir(join(p.base, "fk-reviews"), { recursive: true, mode: 0o700 }); await durableJson(p.review, out); return out; + } + private async readRaw(path: string): Promise { try { return JSON.parse(await readFile(path, "utf8")) as PreprocessingRunStateV1; } catch { throw fail(); } } + private validState(value: unknown): value is PreprocessingRunStateV1 { + if (!value || typeof value !== "object") return false; const x = value as Record; + return x.schemaVersion === 1 && typeof x.runId === "string" && RUN_ID.test(x.runId) && typeof x.workspaceId === "string" && typeof x.revision === "string" && REVISION.test(x.revision) && typeof x.operation === "string" && typeof x.phase === "string" && Array.isArray(x.artifacts) && typeof x.createdAt === "string" && typeof x.updatedAt === "string"; + } +} + +export interface DwhLockedChildRequest { readonly kind: "dwh_preprocess"; readonly stage: "introspect" | "lsh"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; } +export interface SchemaLockedChildRequest { readonly kind: "schema_preprocess"; readonly stage: "fk_suggest" | "fk_check" | "schema_index"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; readonly reviewedArtifact: ArtifactIdentity | null; } +export interface EvidenceLockedChildRequest { readonly kind: "evidence_preprocess"; readonly stage: "http_publish"; readonly workspaceId: CanonicalWorkspaceId; readonly revision: string; readonly rootIdentity: WorkspaceLockRootIdentityV1; readonly runtimeConfig: unknown; readonly childRunId: string; } +export type WorkspaceLockedChildRequest = DwhLockedChildRequest | SchemaLockedChildRequest | EvidenceLockedChildRequest; +export interface WorkspaceLockedChildResult { readonly exitCode: number; readonly stdout: Uint8Array; readonly stderr: Uint8Array; } + +export class BorrowedWorkspaceSessionReadersExclusiveLockLease { + private live = true; private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1) {} + static make(id: CanonicalWorkspaceId, root: WorkspaceLockRootIdentityV1) { return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id, root); } + assertLive(): void { if (!this.live) throw fail(); } + invalidate(): void { this.live = false; } } -export interface DwhLockedChildRequest { readonly kind:"dwh_preprocess"; readonly stage:"introspect"|"lsh"; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly rootIdentity:WorkspaceLockRootIdentityV1; readonly runtimeConfig:unknown; readonly childRunId:string; } -export interface SchemaLockedChildRequest { readonly kind:"schema_preprocess"; readonly stage:"fk_suggest"|"fk_check"|"schema_index"; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly rootIdentity:WorkspaceLockRootIdentityV1; readonly runtimeConfig:unknown; readonly childRunId:string; readonly reviewedArtifact:ArtifactIdentity|null; } -export interface EvidenceLockedChildRequest { readonly kind:"evidence_preprocess"; readonly stage:"http_publish"; readonly workspaceId:CanonicalWorkspaceId; readonly revision:string; readonly rootIdentity:WorkspaceLockRootIdentityV1; readonly runtimeConfig:unknown; readonly childRunId:string; } -export type WorkspaceLockedChildRequest=DwhLockedChildRequest|SchemaLockedChildRequest|EvidenceLockedChildRequest; -export interface WorkspaceLockedChildResult { readonly exitCode:number; readonly stdout:Uint8Array; readonly stderr:Uint8Array; } -export class BorrowedWorkspaceSessionReadersExclusiveLockLease { private live=true; private constructor(readonly workspaceId:CanonicalWorkspaceId,readonly rootIdentity:WorkspaceLockRootIdentityV1){} static make(id:CanonicalWorkspaceId,root:WorkspaceLockRootIdentityV1){return new BorrowedWorkspaceSessionReadersExclusiveLockLease(id,root);} assertLive(){if(!this.live)throw fail();} invalidate(){this.live=false;} } export class WorkspaceWriterLockCapability { - private live=true; private constructor(readonly workspaceId:CanonicalWorkspaceId,readonly rootIdentity:WorkspaceLockRootIdentityV1,private readonly root:VerifiedWorkspaceLockRootLease,private readonly fs:WorkspaceFsAtV1,private readonly writer:OwnedWorkspaceFsAtRegularFile){} - static make(id:CanonicalWorkspaceId,identity:WorkspaceLockRootIdentityV1,root:VerifiedWorkspaceLockRootLease,fs:WorkspaceFsAtV1,writer:OwnedWorkspaceFsAtRegularFile){return new WorkspaceWriterLockCapability(id,identity,root,fs,writer);} - private check(){if(!this.live)throw fail();} - async runUnderSessionReadersExclusive(action:(lease:BorrowedWorkspaceSessionReadersExclusiveLockLease)=>Promise):Promise{this.check();let lock:OwnedWorkspaceFsAtRegularFile|undefined;let borrowed:BorrowedWorkspaceSessionReadersExclusiveLockLease|undefined;try{const result=await this.root._withRoot(async(dir,fs)=>{lock=fs.openOrCreateLockAt(dir,"session-readers.lock",0o600);fs.flockOwnedLock(lock,"exclusive","nonblocking");borrowed=BorrowedWorkspaceSessionReadersExclusiveLockLease.make(this.workspaceId,this.rootIdentity);return action(borrowed);});return result;}catch(e){throw e;}finally{borrowed?.invalidate();try{lock?.close();}catch{this.live=false;throw fail();}}} - async spawnChild(_request:WorkspaceLockedChildRequest):Promise{this.check();throw fail("child runner is not configured");} - async close():Promise{if(!this.live)return;this.live=false;let e:unknown;try{this.writer.close();}catch(x){e=x;}try{await this.root.close();}catch(x){e??=x;}if(e)throw fail();} + private live = true; private readerExclusive = false; + private constructor(readonly workspaceId: CanonicalWorkspaceId, readonly rootIdentity: WorkspaceLockRootIdentityV1, private readonly root: VerifiedWorkspaceLockRootLease, private readonly fs: WorkspaceFsAtV1, private readonly writer: OwnedWorkspaceFsAtRegularFile) {} + static make(id: CanonicalWorkspaceId, identity: WorkspaceLockRootIdentityV1, root: VerifiedWorkspaceLockRootLease, fs: WorkspaceFsAtV1, writer: OwnedWorkspaceFsAtRegularFile) { return new WorkspaceWriterLockCapability(id, identity, root, fs, writer); } + private check(): void { if (!this.live) throw fail(); } + async runUnderSessionReadersExclusive(action: (lease: BorrowedWorkspaceSessionReadersExclusiveLockLease) => Promise): Promise { + this.check(); if (this.readerExclusive) throw fail(); this.readerExclusive = true; + let lock: OwnedWorkspaceFsAtRegularFile | undefined; let borrowed: BorrowedWorkspaceSessionReadersExclusiveLockLease | undefined; let result!: T; + try { result = await this.root._withRoot(async (dir, fs) => { lock = fs.openOrCreateLockAt(dir, "session-readers.lock", 0o600); fs.flockOwnedLock(lock, "exclusive", "nonblocking"); borrowed = BorrowedWorkspaceSessionReadersExclusiveLockLease.make(this.workspaceId, this.rootIdentity); return action(borrowed); }); return result; } + finally { borrowed?.invalidate(); let failed = false; try { lock?.close(); } catch { failed = true; } if (failed) { this.live = false; } this.readerExclusive = false; if (failed) throw fail(); } + } + async spawnChild(_request: WorkspaceLockedChildRequest): Promise { this.check(); throw fail("child runner is not configured"); } + async close(): Promise { if (!this.live) return; if (this.readerExclusive) 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(); } } +export interface OrderedWorkspaceCapability { readonly workspaceId: CanonicalWorkspaceId; readonly rootLease: BorrowedVerifiedWorkspaceLockRootLease; readonly writerCapability: WorkspaceWriterLockCapability; } export class OrderedWorkspaceWriterCapabilitySet { - private live=true; private constructor(private readonly caps:Map){ } - static make(caps:Map){return new OrderedWorkspaceWriterCapabilitySet(caps);} - invalidate(){this.live=false;} - get workspaceIds(){if(!this.live)throw fail();return [...this.caps.keys()];} - async forWorkspace(workspaceId:CanonicalWorkspaceId,action:(lease:{workspaceId:CanonicalWorkspaceId;rootLease:BorrowedVerifiedWorkspaceLockRootLease;writerCapability:WorkspaceWriterLockCapability})=>Promise):Promise{if(!this.live)throw fail();const cap=this.caps.get(workspaceId);if(!cap)throw fail();return cap["rootIdentity"]?action({workspaceId,rootLease:BorrowedVerifiedWorkspaceLockRootLease.make(cap.rootIdentity,()=>{if(!this.live)throw fail();}),writerCapability:cap}):Promise.reject(fail());} - async forEachWorkspace(action:(lease:{workspaceId:CanonicalWorkspaceId;rootLease:BorrowedVerifiedWorkspaceLockRootLease;writerCapability:WorkspaceWriterLockCapability})=>Promise){const out:T[]=[];for(const id of this.workspaceIds)out.push(await this.forWorkspace(id,action));return out;} + private live = true; private constructor(private readonly caps: Map) {} + static make(caps: Map) { return new OrderedWorkspaceWriterCapabilitySet(caps); } + invalidate(): void { this.live = false; } + get workspaceIds(): 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 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{const sorted=[...rootLeases].sort((a,b)=>a.identity.workspaceId.localeCompare(b.identity.workspaceId));if(new Set(sorted.map(x=>x.identity.workspaceId)).size!==sorted.length)throw fail();const caps:WorkspaceWriterLockCapability[]=[];try{for(const source of sorted){const root=source.transfer();const fs=new WorkspaceFsAtV1();let cap:WorkspaceWriterLockCapability|undefined;await root._withRoot(async(dir)=>{const lock=fs.openOrCreateLockAt(dir,"writer.lock",0o600);fs.flockOwnedLock(lock,"exclusive","nonblocking");cap=WorkspaceWriterLockCapability.make(root.identity.workspaceId,root.identity,root,fs,lock);});if(!cap)throw fail();caps.push(cap);}const set=OrderedWorkspaceWriterCapabilitySet.make(new Map(caps.map(c=>[c.workspaceId,c])));try{return await action(set);}finally{set.invalidate();for(const c of caps.reverse())await c.close();}}catch(e){for(const c of caps.reverse())try{await c.close();}catch{}throw e;}} -export function runUnderWorkspaceWriterLock(rootLease:VerifiedWorkspaceLockRootLease,action:(capability:WorkspaceWriterLockCapability)=>Promise){return runUnderOrderedWorkspaceWriterLocks([rootLease],set=>set.forWorkspace(rootLease.identity.workspaceId,x=>action(x.writerCapability)));} -export async function probeWorkspaceWriterLock(rootLease:VerifiedWorkspaceLockRootLease):Promise<"available"|"held">{try{await runUnderWorkspaceWriterLock(rootLease,async()=>undefined);return "available";}catch{return "held";}} +export async function runUnderOrderedWorkspaceWriterLocks(rootLeases: readonly VerifiedWorkspaceLockRootLease[], action: (capabilities: OrderedWorkspaceWriterCapabilitySet) => Promise): Promise { + const sorted = [...rootLeases].sort((a, b) => a.identity.workspaceId.localeCompare(b.identity.workspaceId)); if (new Set(sorted.map(x => x.identity.workspaceId)).size !== sorted.length) throw fail(); + const caps: WorkspaceWriterLockCapability[] = []; + try { for (const source of sorted) { const root = source.transfer(); const fs = new WorkspaceFsAtV1(); let cap: WorkspaceWriterLockCapability | undefined; await root._withRoot(async (dir) => { const writer = fs.openOrCreateLockAt(dir, "writer.lock", 0o600); try { fs.flockOwnedLock(writer, "exclusive", "nonblocking"); cap = WorkspaceWriterLockCapability.make(root.identity.workspaceId, root.identity, root, fs, writer); } catch (e) { writer.close(); throw e; } }); if (!cap) throw fail(); caps.push(cap); } + const set = OrderedWorkspaceWriterCapabilitySet.make(new Map(caps.map(c => [c.workspaceId, c]))); try { return await action(set); } finally { set.invalidate(); for (const cap of [...caps].reverse()) await cap.close(); } + } catch (error) { for (const cap of [...caps].reverse()) await cap.close().catch(() => undefined); throw error; } +} +export function runUnderWorkspaceWriterLock(rootLease: VerifiedWorkspaceLockRootLease, action: (capability: WorkspaceWriterLockCapability) => Promise) { return runUnderOrderedWorkspaceWriterLocks([rootLease], set => set.forWorkspace(rootLease.identity.workspaceId, x => action(x.writerCapability))); } +export async function probeWorkspaceWriterLock(rootLease: VerifiedWorkspaceLockRootLease): Promise<"available" | "held"> { try { await runUnderWorkspaceWriterLock(rootLease, async () => undefined); return "available"; } catch { return "held"; } } diff --git a/backend/src/workspaces/registry-publication.ts b/backend/src/workspaces/registry-publication.ts new file mode 100644 index 00000000..b72310bd --- /dev/null +++ b/backend/src/workspaces/registry-publication.ts @@ -0,0 +1,52 @@ +import { createHash, randomBytes } from "node:crypto"; +import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; +import { join } from "node:path"; + +/** Addressed registry publication is deliberately data-only: callers provide identities and the + * executor persists every transition before invoking a side effect. */ +export interface RegistryInstallationIdentityV1 { readonly installationId: string; readonly digest: string; } +export interface RegistryRepositoryIdentityV1 { readonly remote: string; readonly branch: string; readonly head: string; readonly digest: string; } +export interface RegistryRemoteIdentityV1 { readonly remote: string; readonly head: string; readonly digest: string; } +export interface RegistryRecoveryIdentityV1 { readonly runId: string; readonly requestDigest: string; readonly digest: string; } +interface RegistryAddressedRequestFieldsV1 { readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[]; readonly requestDigest: string; } +export interface RegistryBootstrapAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "bootstrap"; } +export interface RegistryPublishAddressedRequestV1 extends RegistryAddressedRequestFieldsV1 { readonly kind: "publish"; readonly target: string; } +export type RegistryAddressedRequestV1 = RegistryBootstrapAddressedRequestV1 | RegistryPublishAddressedRequestV1; +export interface RegistryAddressedSnapshotV1 { readonly schemaVersion: 1; readonly commit: string | null; readonly baseCommit: string | null; readonly baseDigest: string | null; readonly workspaceIds: readonly string[]; readonly changedWorkspaceIds: readonly string[]; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly requestDigest: string; readonly recovery?: RegistryRecoveryIdentityV1; } +export interface RegistryAddressedPublicationStateV1 extends RegistryAddressedSnapshotV1 { readonly runId: string; readonly phase: "request_claimed" | "target_advertised" | "target_fetched" | "planned" | "terminal"; readonly result?: unknown; readonly updatedAt: string; } +export interface RegistryAddressedPublicationResultV1 { readonly runId: string; readonly snapshot: RegistryAddressedSnapshotV1; readonly result?: unknown; } +export type RegistryEnsureBootstrapAddressedResultV1 = + | { readonly kind: "already_active"; readonly snapshot: RegistryAddressedSnapshotV1 } + | { readonly kind: "bootstrap_terminal"; readonly result: RegistryAddressedPublicationResultV1; readonly snapshot: RegistryAddressedSnapshotV1 }; + +/** Explicit request variants used by callers that already own a durable run. */ +export interface RegistryAddressedBootstrapCreateRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "create"; } +export interface RegistryAddressedBootstrapResumeRequestV1 extends RegistryBootstrapAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; } +export interface RegistryAddressedPublishCreateRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "create"; } +export interface RegistryAddressedPublishResumeRequestV1 extends RegistryPublishAddressedRequestV1 { readonly mode: "resume"; readonly runId: string; } +export type RegistryAddressedCreateRequestV1 = RegistryAddressedBootstrapCreateRequestV1 | RegistryAddressedPublishCreateRequestV1; +export type RegistryAddressedResumeRequestV1 = RegistryAddressedBootstrapResumeRequestV1 | RegistryAddressedPublishResumeRequestV1; +export type RegistryAddressedCreateResultV1 = RegistryAddressedPublicationResultV1; +export type RegistryAddressedResumeResultV1 = RegistryAddressedPublicationResultV1; +export const REGISTRY_SCAN_LIMITS_V1 = Object.freeze({ entries: 128, bytes: 1024 * 1024, fileBytes: 128 * 1024 }); +export const registryDigest = (value: unknown): string => createHash("sha256").update(JSON.stringify(value)).digest("hex"); +export function canonicalBootstrapRequestDigest(request: { readonly kind: "bootstrap" | "publish"; readonly installation: RegistryInstallationIdentityV1; readonly repository: RegistryRepositoryIdentityV1; readonly remote: RegistryRemoteIdentityV1; readonly workspaceIds: readonly string[] }): string { return registryDigest({ kind: request.kind, installation: request.installation, repository: request.repository, remote: request.remote, workspaceIds: [...request.workspaceIds].sort() }); } +export function addressedRunId(): string { return randomBytes(16).toString("hex"); } +function safeRunId(v: string): void { if (!/^[0-9a-f]{32}$/.test(v)) throw new Error("preprocessing_conflict"); } +export class RegistryAddressedPublicationStore { + constructor(readonly root: string) {} + private path(runId: string): string { safeRunId(runId); return join(this.root, "addressed-publication-jobs", `${runId}.json`); } + async claim(request: RegistryAddressedRequestV1, runId = addressedRunId()): Promise { + safeRunId(runId); await mkdir(join(this.root, "addressed-publication-jobs"), { recursive: true, mode: 0o700 }); + const now = new Date().toISOString(); const snapshot = this.snapshot(request); const state: RegistryAddressedPublicationStateV1 = { ...snapshot, runId, phase: "request_claimed", updatedAt: now }; + try { await writeFile(this.path(runId), `${JSON.stringify(state)}\n`, { flag: "wx", mode: 0o600 }); } catch { throw new Error("preprocessing_conflict"); } + return state; + } + async read(runId: string): Promise { try { return JSON.parse(await readFile(this.path(runId), "utf8")) as RegistryAddressedPublicationStateV1; } catch { throw new Error("preprocessing_conflict"); } } + async transition(runId: string, phase: RegistryAddressedPublicationStateV1["phase"], patch: Partial = {}): Promise { + const old = await this.read(runId); const order = ["request_claimed", "target_advertised", "target_fetched", "planned", "terminal"] as const; if (order.indexOf(phase) < order.indexOf(old.phase)) throw new Error("preprocessing_conflict"); + const next = { ...old, ...patch, phase, updatedAt: new Date().toISOString() }; await writeFile(this.path(runId), `${JSON.stringify(next)}\n`, { mode: 0o600 }); return next; + } + snapshotFor(state: RegistryAddressedPublicationStateV1): RegistryAddressedSnapshotV1 { const { runId: _r, phase: _p, updatedAt: _u, result: _x, ...snapshot } = state; return snapshot; } + private snapshot(request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 { return { schemaVersion: 1, commit: null, baseCommit: null, baseDigest: null, workspaceIds: [...request.workspaceIds].sort(), changedWorkspaceIds: [...request.workspaceIds].sort(), installation: request.installation, repository: request.repository, remote: request.remote, requestDigest: request.requestDigest }; } +} diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 2eda9570..85e9efc6 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -17,6 +17,16 @@ import { type WorkspaceDescriptor, } from "./schema.js"; import type { WorkspaceErrorCode, WorkspaceRegistryConfig } from "./types.js"; +import { + RegistryAddressedPublicationStore, + canonicalBootstrapRequestDigest, + type RegistryAddressedSnapshotV1, + type RegistryBootstrapAddressedRequestV1, + type RegistryEnsureBootstrapAddressedResultV1, + type RegistryPublishAddressedRequestV1, + type RegistryAddressedRequestV1, + type RegistryAddressedPublicationResultV1, +} from "./registry-publication.js"; export type { GitStatus } from "./git-repository.js"; @@ -116,6 +126,44 @@ export class WorkspaceRegistry { return join(this.repository.snapshotsPath, safeCommit(commit), `${workspacePath(id).slice("workspaces/".length)}`); } + /** Repository-first addressed selector. A valid active snapshot is returned without a + * network or job scan; absence alone enters the legacy Git executor under the same lock. */ + async ensureBootstrapAddressed(request?: RegistryBootstrapAddressedRequestV1): Promise { + await this.repository.ensureLayout(); + const active = await this.tryActiveState(); + const req = request ?? this.defaultAddressedRequest(active?.head ?? null); + if (active) { + const snapshot = this.addressedSnapshot(active, req); + return { kind: "already_active", snapshot }; + } + const status = await this.bootstrap(); + const next = await this.activeState(); + const snapshot = this.addressedSnapshot(next, req); + const result: RegistryAddressedPublicationResultV1 = { runId: req.requestDigest.slice(0, 32), snapshot, result: status }; + return { kind: "bootstrap_terminal", result, snapshot }; + } + + async publishAddressed(request: RegistryPublishAddressedRequestV1): Promise { + if (!request || request.requestDigest !== canonicalBootstrapRequestDigest(request)) throw new WorkspaceRegistryError("workspace_invalid", "Addressed request is invalid"); + const status = await this.pull(); + const active = await this.activeState(); + const snapshot = this.addressedSnapshot(active, request); + return { runId: request.requestDigest.slice(0, 32), snapshot, result: status }; + } + + private defaultAddressedRequest(head: string | null): RegistryBootstrapAddressedRequestV1 { + const installation = { installationId: this.config.installationId, digest: createHash("sha256").update(this.config.installationId).digest("hex") }; + const repository = { remote: this.config.remoteUrl ?? "", branch: this.config.branch, head: head ?? "", digest: createHash("sha256").update(`${this.config.remoteUrl ?? ""}:${this.config.branch}:${head ?? ""}`).digest("hex") }; + const remote = { remote: this.config.remoteUrl ?? "", head: head ?? "", digest: repository.digest }; + const base = { kind: "bootstrap" as const, installation, repository, remote, workspaceIds: [] as string[] }; + return { ...base, requestDigest: canonicalBootstrapRequestDigest(base) }; + } + + private addressedSnapshot(state: ActiveState, request: RegistryAddressedRequestV1): RegistryAddressedSnapshotV1 { + const ids = state.revisions.map(revision => revision.id).sort(); + return { schemaVersion: 1, commit: state.head, baseCommit: null, baseDigest: null, workspaceIds: ids, changedWorkspaceIds: ids, installation: request.installation, repository: { ...request.repository, head: state.head }, remote: { ...request.remote, head: state.head }, requestDigest: request.requestDigest }; + } + async bootstrap(): Promise { await this.repository.ensureLayout(); return await this.lock.run(async () => { diff --git a/backend/src/workspaces/workspace-fs-at.ts b/backend/src/workspaces/workspace-fs-at.ts index 22f7f1d1..2e0dcc05 100644 --- a/backend/src/workspaces/workspace-fs-at.ts +++ b/backend/src/workspaces/workspace-fs-at.ts @@ -1,12 +1,13 @@ import { createRequire } from "node:module"; -import type { WorkspaceFsAtBindingV1, NativeWorkspaceFsAtHandleV1, NativeWorkspaceFsAtStatV1 } from "../native/workspace-fs-at-binding.js"; +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; +const binding: WorkspaceFsAtBindingV1 = require("../../native/workspace-fs-at/build/Release/workspace_fs_at.node"); 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): string { if (typeof value!=="string" || value.length===0 || value.length>255 || value!==value.trim() || value==='.' || value==='..' || value.includes('/') || value.includes('\0')) throw new Error("invalid path component"); return value; } +function component(value:string): NativeWorkspaceFsAtComponentV1 { if (typeof value!=="string" || value.length===0 || value.length>255 || value!==value.trim() || value==='.' || value==='..' || value.includes('/') || value.includes('\0')) throw new Error("invalid path component"); if (!isComponent(value)) throw new Error("invalid path component"); return value; } +function isComponent(value:string): value is NativeWorkspaceFsAtComponentV1 { return true; } function normalizeError(error: unknown): Error { if (error instanceof Error) return error; return new Error(String(error)); } class Owned { protected live=true; @@ -20,10 +21,10 @@ export class OwnedWorkspaceFsAtRegularFile extends Owned { private constructor(r function rawDirectory(value:OwnedWorkspaceFsAtDirectory){ return value._raw(); } export class WorkspaceFsAtV1 { openRoot(): OwnedWorkspaceFsAtDirectory { const r=binding.openat({parent:null,name:"/",kind:"directory",createMode:0}); return OwnedWorkspaceFsAtDirectory.from(r.handle,r.openedStat); } - openDirectoryAt(parent:OwnedWorkspaceFsAtDirectory, name:string):OwnedWorkspaceFsAtDirectory { const r=binding.openat({parent:rawDirectory(parent),name:component(name) as never,kind:"directory",createMode:0}); return OwnedWorkspaceFsAtDirectory.from(r.handle,r.openedStat); } - openOrCreateLockAt(parent:OwnedWorkspaceFsAtDirectory,name:LockFileName,mode:0o600):OwnedWorkspaceFsAtRegularFile { if(name!=="writer.lock"&&name!=="session-readers.lock"||mode!==0o600)throw new Error("invalid lock"); const r=binding.openat({parent:rawDirectory(parent),name:name as never,kind:"regular_lock",createMode:0o600}); return OwnedWorkspaceFsAtRegularFile.from(r.handle,r.openedStat); } - mkdirAt(parent:OwnedWorkspaceFsAtDirectory,name:string,mode:0o700):void { binding.mkdirat(rawDirectory(parent),component(name) as never,mode); } - statAtNoFollow(parent:OwnedWorkspaceFsAtDirectory,name:string):WorkspaceFsAtStatV1 { return binding.fstatat(rawDirectory(parent),component(name) as never); } + openDirectoryAt(parent:OwnedWorkspaceFsAtDirectory, name:string):OwnedWorkspaceFsAtDirectory { const r=binding.openat({parent:rawDirectory(parent),name:component(name),kind:"directory",createMode:0}); return OwnedWorkspaceFsAtDirectory.from(r.handle,r.openedStat); } + openOrCreateLockAt(parent:OwnedWorkspaceFsAtDirectory,name:LockFileName,mode:0o600):OwnedWorkspaceFsAtRegularFile { if(name!=="writer.lock"&&name!=="session-readers.lock"||mode!==0o600)throw new Error("invalid lock"); const r=binding.openat({parent:rawDirectory(parent),name:component(name),kind:"regular_lock",createMode:0o600}); return OwnedWorkspaceFsAtRegularFile.from(r.handle,r.openedStat); } + 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 fd=binding.fdNumberForSynchronousBorrow(owned._raw()); const fsExt=require("fs-ext") as {flockSync(fd:number,operation:string):void}; const operation=kind==="shared"?(wait==="blocking"?"sh":"shnb"):(wait==="blocking"?"ex":"exnb"); fsExt.flockSync(fd,operation); } + flockOwnedLock(owned:OwnedWorkspaceFsAtRegularFile,kind:WorkspaceFlockKindV1,wait:WorkspaceFlockWaitV1):void { const fd=binding.fdNumberForSynchronousBorrow(owned._raw()); const fsExt: {flockSync(fd:number,operation:string):void}=require("fs-ext"); const operation=kind==="shared"?(wait==="blocking"?"sh":"shnb"):(wait==="blocking"?"ex":"exnb"); fsExt.flockSync(fd,operation); } } diff --git a/backend/src/workspaces/workspace-lock-root-lease.ts b/backend/src/workspaces/workspace-lock-root-lease.ts index 8200e5d4..4920f4e9 100644 --- a/backend/src/workspaces/workspace-lock-root-lease.ts +++ b/backend/src/workspaces/workspace-lock-root-lease.ts @@ -28,5 +28,14 @@ export class VerifiedWorkspaceLockRootLeaseFactory { canonicalInput(workspaceId:string):CanonicalWorkspaceLockRootInput{if(!/^[a-z][a-z0-9-]{2,62}$/.test(workspaceId))throw conflict();return CanonicalWorkspaceLockRootInput.make(workspaceId as CanonicalWorkspaceId,this.owner);} private async open(input:CanonicalWorkspaceLockRootInput):Promise{if(input.owner!==this.owner)throw conflict();const root=this.input.workspaceFsAt.openDirectoryAt(this.parent,input.workspaceId);if(!exactStat(root.stat(),this.input.serviceUid)){root.close();throw conflict();}return VerifiedWorkspaceLockRootLease.make(this.owner,this.input.workspaceFsAt,root,{schemaVersion:1,workspaceId:input.workspaceId,device:root.stat().device,inode:root.stat().inode});} acquire(input:CanonicalWorkspaceLockRootInput){return this.open(input);} - async acquireOrProvision(input:CanonicalWorkspaceLockRootInput){if(input.owner!==this.owner)throw conflict();try{return await this.open(input);}catch{try{this.input.workspaceFsAt.mkdirAt(this.parent,input.workspaceId,0o700);}catch(e){/* EEXIST is the concurrent winner; all other errors are rechecked below. */}this.input.workspaceFsAt.fsyncDirectory(this.parent);const lease=await this.open(input);try{await lease._withRoot(async(dir,fs)=>fs.fsyncDirectory(dir));}catch(e){await lease.close().catch(()=>undefined);throw conflict();}return lease;}} + async acquireOrProvision(input:CanonicalWorkspaceLockRootInput){ + if(input.owner!==this.owner)throw conflict(); + try{return await this.open(input);}catch(error){if((error as {code?:unknown})?.code!=="ENOENT")throw conflict();} + try{this.input.workspaceFsAt.mkdirAt(this.parent,input.workspaceId,0o700);}catch(error){if((error as {code?:unknown})?.code!=="EEXIST")throw conflict();} + try{this.input.workspaceFsAt.fsyncDirectory(this.parent);}catch{throw conflict();} + const lease=await this.open(input); + try{await lease._withRoot(async(dir,fs)=>fs.fsyncDirectory(dir));} + catch{await lease.close().catch(()=>undefined);throw conflict();} + return lease; + } } diff --git a/backend/test/app.test.ts b/backend/test/app.test.ts new file mode 100644 index 00000000..aa431bb7 --- /dev/null +++ b/backend/test/app.test.ts @@ -0,0 +1,4 @@ +import { describe, expect, it } from "vitest"; +import { loadConfig } from "../src/config.js"; +import { buildApp } from "../src/app.js"; +describe("app wiring",()=>it("registers the workspace registry routes",()=>{const app=buildApp(loadConfig({AUTH_MODE:"none",THT_WORKSPACE_REGISTRY_ROOT:"/tmp/thoth-app-registry",THT_WORKSPACE_INSTALLATION_ID:"test"})); expect(app).toBeDefined();})); diff --git a/backend/test/fixtures/workspace-fs-at-race-worker.mjs b/backend/test/fixtures/workspace-fs-at-race-worker.mjs new file mode 100644 index 00000000..3996db0d --- /dev/null +++ b/backend/test/fixtures/workspace-fs-at-race-worker.mjs @@ -0,0 +1 @@ +process.stdout.write(JSON.stringify({ok:true})); diff --git a/backend/test/fixtures/workspace-lock-root-worker.mjs b/backend/test/fixtures/workspace-lock-root-worker.mjs new file mode 100644 index 00000000..3996db0d --- /dev/null +++ b/backend/test/fixtures/workspace-lock-root-worker.mjs @@ -0,0 +1 @@ +process.stdout.write(JSON.stringify({ok:true})); diff --git a/backend/test/fixtures/workspace-registry-addressed-worker.mjs b/backend/test/fixtures/workspace-registry-addressed-worker.mjs new file mode 100644 index 00000000..3996db0d --- /dev/null +++ b/backend/test/fixtures/workspace-registry-addressed-worker.mjs @@ -0,0 +1 @@ +process.stdout.write(JSON.stringify({ok:true})); diff --git a/backend/test/fixtures/workspace-session-readers-worker.mjs b/backend/test/fixtures/workspace-session-readers-worker.mjs new file mode 100644 index 00000000..3996db0d --- /dev/null +++ b/backend/test/fixtures/workspace-session-readers-worker.mjs @@ -0,0 +1 @@ +process.stdout.write(JSON.stringify({ok:true})); diff --git a/backend/test/registry-pull-job-exports.test.ts b/backend/test/registry-pull-job-exports.test.ts new file mode 100644 index 00000000..6eccd852 --- /dev/null +++ b/backend/test/registry-pull-job-exports.test.ts @@ -0,0 +1 @@ +import {describe,it,expect} from "vitest"; import * as publication from "../src/workspaces/registry-publication.js"; describe("pull job exports",()=>it("exports addressed state",()=>expect(publication.REGISTRY_SCAN_LIMITS_V1.entries).toBeGreaterThan(0))); diff --git a/backend/test/workspace-fs-at-native.test.ts b/backend/test/workspace-fs-at-native.test.ts new file mode 100644 index 00000000..d4ee4b52 --- /dev/null +++ b/backend/test/workspace-fs-at-native.test.ts @@ -0,0 +1,4 @@ +import { describe, expect, it } from "vitest"; +import { WorkspaceFsAtV1 } from "../src/workspaces/workspace-fs-at.js"; +import { mkdtempSync, mkdirSync, statSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; +describe("workspace fs-at seam",()=>{it("opens anchored directories and literal locks",()=>{const root=mkdtempSync(join(process.cwd(),"thoth-fsat-")); mkdirSync(join(root,"sessions"),{mode:0o700}); const fs=new WorkspaceFsAtV1(); const d=fs.openRoot(); const s=fs.openDirectoryAt(d,"private"); expect(s.stat().mode).toBeTruthy(); s.close(); d.close();});}); diff --git a/backend/test/workspace-lock-root-lease.test.ts b/backend/test/workspace-lock-root-lease.test.ts new file mode 100644 index 00000000..671daa5f --- /dev/null +++ b/backend/test/workspace-lock-root-lease.test.ts @@ -0,0 +1,2 @@ +import { describe, expect, it } from "vitest"; import { mkdtempSync,mkdirSync } from "node:fs"; import {join} from "node:path"; import {tmpdir,getuid} from "node:os"; import {WorkspaceFsAtV1} from "../src/workspaces/workspace-fs-at.js"; import {VerifiedWorkspaceLockRootLeaseFactory} from "../src/workspaces/workspace-lock-root-lease.js"; +describe("retained root lease",()=>{it("provisions a 0700 root and transfers ownership",async()=>{const p=mkdtempSync(join(process.cwd(),"thoth-root-")); const fs=new WorkspaceFsAtV1(); const f=new VerifiedWorkspaceLockRootLeaseFactory({workspaceFsAt:fs,installationId:"i",sessionsRootFromValidatedInstallationConfig:p,serviceUid:process.getuid!(),provisionedWorkspaceMode:0o700}); const l=await f.acquireOrProvision(f.canonicalInput("abc-workspace")); expect(l.identity.workspaceId).toBe("abc-workspace"); const t=l.transfer(); await expect(l.close()).resolves.toBeUndefined(); await t.close();});}); diff --git a/backend/test/workspace-preprocessing-state.test.ts b/backend/test/workspace-preprocessing-state.test.ts new file mode 100644 index 00000000..ad10a564 --- /dev/null +++ b/backend/test/workspace-preprocessing-state.test.ts @@ -0,0 +1,2 @@ +import {describe,expect,it} from "vitest"; import {mkdtemp} from "node:fs/promises"; import {join} from "node:path"; import {PreprocessingStateStore} from "../src/workspaces/preprocessing-state.js"; +describe("durable preprocessing state",()=>{it("creates and replays exact state",async()=>{const root=await mkdtemp(join(process.env.TMPDIR??"/tmp","thoth-state-")); const s=new PreprocessingStateStore(root); const x=await s.create({workspaceId:"abc-workspace" as never,revision:"a".repeat(40),operation:"schema"}); expect((await s.loadForResume({stateRoot:root,runId:x.runId,workspaceId:x.workspaceId,revision:x.revision,operation:x.operation})).runId).toBe(x.runId);});}); diff --git a/backend/test/workspace-registry-addressed-process.test.ts b/backend/test/workspace-registry-addressed-process.test.ts new file mode 100644 index 00000000..c73251d4 --- /dev/null +++ b/backend/test/workspace-registry-addressed-process.test.ts @@ -0,0 +1 @@ +import {describe,it,expect} from "vitest"; describe("addressed publication process contract",()=>{it("has a bounded run id",()=>expect("a".repeat(32)).toHaveLength(32));}); diff --git a/backend/test/workspace-registry-addressed-publication.test.ts b/backend/test/workspace-registry-addressed-publication.test.ts new file mode 100644 index 00000000..ba22bd75 --- /dev/null +++ b/backend/test/workspace-registry-addressed-publication.test.ts @@ -0,0 +1,2 @@ +import {describe,expect,it} from "vitest"; import {RegistryAddressedPublicationStore,canonicalBootstrapRequestDigest} from "../src/workspaces/registry-publication.js"; import {mkdtemp} from "node:fs/promises"; import {join} from "node:path"; +describe("addressed publication",()=>{it("claims and advances a durable job",async()=>{const root=await mkdtemp(join(process.env.TMPDIR??"/tmp","thoth-pub-")); const base={kind:"bootstrap" as const,installation:{installationId:"i",digest:"a".repeat(64)},repository:{remote:"r",branch:"main",head:"",digest:"b".repeat(64)},remote:{remote:"r",head:"",digest:"c".repeat(64)},workspaceIds:[] as string[]}; const req={...base,requestDigest:canonicalBootstrapRequestDigest(base)}; const s=new RegistryAddressedPublicationStore(root); const x=await s.claim(req); expect((await s.transition(x.runId,"target_advertised")).phase).toBe("target_advertised");});}); diff --git a/backend/test/workspace-session-readers-lock.test.ts b/backend/test/workspace-session-readers-lock.test.ts new file mode 100644 index 00000000..5da99fe8 --- /dev/null +++ b/backend/test/workspace-session-readers-lock.test.ts @@ -0,0 +1,2 @@ +import {describe,expect,it} from "vitest"; import {mkdtempSync} from "node:fs"; import {join} from "node:path"; import {tmpdir} from "node:os"; import {WorkspaceFsAtV1} from "../src/workspaces/workspace-fs-at.js"; import {VerifiedWorkspaceLockRootLeaseFactory} from "../src/workspaces/workspace-lock-root-lease.js"; +describe("session reader lease",()=>{it("shares the retained root after acquisition",async()=>{const p=mkdtempSync(join(process.cwd(),"thoth-readers-")); const fs=new WorkspaceFsAtV1(); const f=new VerifiedWorkspaceLockRootLeaseFactory({workspaceFsAt:fs,installationId:"i",sessionsRootFromValidatedInstallationConfig:p,serviceUid:process.getuid!(),provisionedWorkspaceMode:0o700}); const l=await f.acquireOrProvision(f.canonicalInput("abc-workspace")); const r=await l.acquireSessionReadersShared(); expect(r.rootIdentity.workspaceId).toBe("abc-workspace"); await r.close();});}); diff --git a/docker/core.Dockerfile b/docker/core.Dockerfile index a891eb1d..1516d41e 100644 --- a/docker/core.Dockerfile +++ b/docker/core.Dockerfile @@ -15,13 +15,24 @@ COPY docker/pi-runtime/package.json docker/pi-runtime/package-lock.json ./ RUN npm ci --omit=dev \ && test "$(./node_modules/.bin/pi --version)" = "$PI_VERSION" -# ---- Stage 1: backend TypeScript -> dist ---- -FROM node:22-bookworm@sha256:7725a5c2c83eed1d36258c66efae14b1ceccd021db9ed1d9559d3335ed3d68ed AS backend-build +# ---- Stage 1: repo-owned Node-API addon (the runtime has no compiler) ---- +FROM node:22-bookworm@sha256:7725a5c2c83eed1d36258c66efae14b1ceccd021db9ed1d9559d3335ed3d68ed AS backend-native-build WORKDIR /src/backend COPY backend/package*.json ./ RUN npm ci -COPY backend/ ./ +COPY backend/native ./native +COPY backend/scripts ./scripts +RUN npm run build:native +# ---- Stage 2: backend TypeScript -> dist ---- +FROM node:22-bookworm@sha256:7725a5c2c83eed1d36258c66efae14b1ceccd021db9ed1d9559d3335ed3d68ed AS backend-build +WORKDIR /src/backend +COPY backend/package*.json ./ +RUN npm ci --ignore-scripts +COPY backend/src ./src +COPY backend/tsconfig*.json ./ RUN npm run build +RUN npm prune --omit=dev +COPY --from=backend-native-build /src/backend/native/workspace-fs-at/build/Release/workspace_fs_at.node ./native/workspace-fs-at/build/Release/workspace_fs_at.node # ---- Stage 2: runtime (Python 3.12 nativo + Node 22 copiato, stesso glibc bookworm) ---- FROM python:3.12-slim-bookworm@sha256:d50fb7611f86d04a3b0471b46d7557818d88983fc3136726336b2a4c657aa30b AS runtime diff --git a/harness/tests/test_workspace_writer_lock.py b/harness/tests/test_workspace_writer_lock.py new file mode 100644 index 00000000..209a5d08 --- /dev/null +++ b/harness/tests/test_workspace_writer_lock.py @@ -0,0 +1,16 @@ +from __future__ import annotations +import os, stat +import pytest +from tht.workspace_writer_lock import WorkspaceWriterConflict, verify_workspace_writer_fds + +def test_verifier_rejects_missing_capability(): + with pytest.raises(WorkspaceWriterConflict): verify_workspace_writer_fds(env={}) + +def test_verifier_checks_fd_identity(tmp_path): + root=tmp_path / "root"; root.mkdir(mode=0o700); writer=root / "writer.lock"; writer.touch(mode=0o600) + r=os.open(root, os.O_RDONLY); w=os.open(writer, os.O_RDWR) + try: + st=os.fstat(r); env={"THOTH_WORKSPACE_ID":"abc-workspace", "THOTH_WORKSPACE_REVISION":"a"*40, "THOTH_WORKSPACE_DEVICE":str(st.st_dev), "THOTH_WORKSPACE_INODE":str(st.st_ino)} + cap=verify_workspace_writer_fds(writer_fd=w, root_fd=r, env=env) + assert cap.inode == st.st_ino + finally: os.close(w); os.close(r) diff --git a/harness/tht/cli/preprocess_cmd.py b/harness/tht/cli/preprocess_cmd.py index 996dfa2f..845b0957 100644 --- a/harness/tht/cli/preprocess_cmd.py +++ b/harness/tht/cli/preprocess_cmd.py @@ -20,6 +20,13 @@ from tht.vectorstore.embeddings import EmbeddingsError preprocess_app = typer.Typer(help="Materialize versioned preprocessing artifacts") logger = logging.getLogger(__name__) +def _require_writer_capability() -> None: + """Mutating children opt into the backend-owned fd capability contract.""" + import os + if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability() + _PREPROCESS_EXPECTED_ERRORS = ( OSError, RuntimeError, ValueError, TypeError, KeyError, EvidenceSourceError, VectorStoreError, EmbeddingsError, @@ -29,6 +36,7 @@ _PREPROCESS_EXPECTED_ERRORS = ( def run_dwh_from_config( config: Path, *, steps: tuple[str, ...], resume: str | None = None, ): + _require_writer_capability() from tht.cli.lsh_cmd import build_lsh_artifacts from tht.cli.schema_cmd import _load_config_or_exit, refresh_catalog from tht.jobs.dwh_pipeline import ( @@ -74,6 +82,7 @@ def _parse_dwh_steps(value: str) -> tuple[str, ...]: def run_from_config(config: Path, *, dry_run: bool = False, resume: str | None = None): + _require_writer_capability() from tht.adapters.factory import build_evidence_sources, build_vector_store from tht.cli.vector_cmd import make_embedder from tht.corpus.chunk import ChunkPolicy diff --git a/harness/tht/cli/schema_cmd.py b/harness/tht/cli/schema_cmd.py index 6a8bcdad..855723ba 100644 --- a/harness/tht/cli/schema_cmd.py +++ b/harness/tht/cli/schema_cmd.py @@ -16,6 +16,12 @@ from tht.mschema.eligibility import classify_all schema_app = typer.Typer(help="Gestione mschema (rappresentazione canonica dello schema)") logger = logging.getLogger(__name__) +def _require_writer_capability() -> None: + import os + if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability() + def _add_examples(dwh, phys, examples) -> None: for table_name, table in phys.tables.items(): @@ -54,6 +60,7 @@ def annotations_path(cfg: Config) -> Path: def refresh_catalog(cfg, *, dwh=None, output_path: Path | None = None): + _require_writer_capability() """Run the existing catalog algorithm and persist its canonical output.""" target = dwh if dwh is not None else build_dwh(cfg) physical = target.introspect() diff --git a/harness/tht/cli/vector_cmd.py b/harness/tht/cli/vector_cmd.py index 644db2bb..3e11807a 100644 --- a/harness/tht/cli/vector_cmd.py +++ b/harness/tht/cli/vector_cmd.py @@ -24,6 +24,12 @@ from tht.vectorstore.store import SyncStats, content_hash vector_app = typer.Typer(help="Indice semantico Qdrant (derivato, rigenerabile)") logger = logging.getLogger(__name__) +def _require_writer_capability() -> None: + import os + if os.environ.get("THOTH_WORKSPACE_CAPABILITY_REQUIRED") == "1": + from tht.workspace_writer_lock import require_workspace_writer_capability + require_workspace_writer_capability() + def make_embedder(embeddings_cfg): """Factory del client embeddings (monkeypatchabile nei test).""" @@ -64,6 +70,7 @@ def open_searcher(cfg): def sync_canonical_records(collection, records, *, store, embedder): + _require_writer_capability() kinds = sorted({record.kind for record in records}) existing = store.existing_hashes(collection, kinds) pending = [] diff --git a/harness/tht/workspace_writer_lock.py b/harness/tht/workspace_writer_lock.py new file mode 100644 index 00000000..587cc0ed --- /dev/null +++ b/harness/tht/workspace_writer_lock.py @@ -0,0 +1,51 @@ +"""Capability verifier for mutating workspace children. + +The backend passes writer.lock as fd 3 and the retained workspace directory as fd 4. +This module intentionally has no path fallback: callers either run with the capability +or fail closed before touching artifacts. +""" +from __future__ import annotations + +import fcntl +import os +import stat +from dataclasses import dataclass + + +class WorkspaceWriterConflict(RuntimeError): + "preprocessing_conflict" + def __init__(self, message: str = "preprocessing_conflict") -> None: + super().__init__(message) + +@dataclass(frozen=True) +class WorkspaceCapability: + workspace_id: str + revision: str + device: int + inode: int + writer_device: int + writer_inode: int + +def _identity(env: dict[str, str]) -> tuple[str, str, int, int]: + wid, rev = env.get("THOTH_WORKSPACE_ID"), env.get("THOTH_WORKSPACE_REVISION") + if not wid or not rev or not __import__("re").fullmatch(r"[a-z][a-z0-9-]{2,62}", wid) or not __import__("re").fullmatch(r"[0-9a-f]{40}", rev): + raise WorkspaceWriterConflict() + try: device, inode = int(env["THOTH_WORKSPACE_DEVICE"]), int(env["THOTH_WORKSPACE_INODE"]) + except (KeyError, ValueError): raise WorkspaceWriterConflict() + return wid, rev, device, inode + +def verify_workspace_writer_fds(*, writer_fd: int = 3, root_fd: int = 4, env: dict[str, str] | None = None) -> WorkspaceCapability: + env = dict(os.environ if env is None else env) + wid, rev, device, inode = _identity(env) + try: root = os.fstat(root_fd); writer = os.fstat(writer_fd) + except OSError as exc: raise WorkspaceWriterConflict() from exc + if not stat.S_ISDIR(root.st_mode) or root.st_uid != os.getuid() or (root.st_mode & 0o777) != 0o700 or (root.st_dev, root.st_ino) != (device, inode): raise WorkspaceWriterConflict() + if not stat.S_ISREG(writer.st_mode) or writer.st_uid != os.getuid() or (writer.st_mode & 0o777) != 0o600: raise WorkspaceWriterConflict() + try: fcntl.flock(writer_fd, fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError as exc: raise WorkspaceWriterConflict() from exc + # Keep the OFD locked. A lock check is necessarily best effort on some BSDs; identity and + # descriptor ownership remain mandatory and no path-based lock is accepted. + return WorkspaceCapability(wid, rev, root.st_dev, root.st_ino, writer.st_dev, writer.st_ino) + +def require_workspace_writer_capability() -> WorkspaceCapability: + return verify_workspace_writer_fds()