Files
ThothII/backend/src/workspaces/preprocessing-state.ts
T

49 lines
9.7 KiB
TypeScript

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");}
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<PreprocessingRunStateV1>{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<PreprocessingRunStateV1>{id(runId);throw new Error("transition requires a state-root bound store");}
async writeFkCandidate(runId:string,yaml:Uint8Array):Promise<ArtifactIdentity>{id(runId);throw new Error("writeFkCandidate requires a state-root bound store");}
async recordFkReview(runId:string,review:FkReviewInput):Promise<FkReviewRecordV1>{id(runId);throw new Error("recordFkReview requires a state-root bound store");}
}
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<T>(action:(lease:BorrowedWorkspaceSessionReadersExclusiveLockLease)=>Promise<T>):Promise<T>{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<WorkspaceLockedChildResult>{this.check();throw fail("child runner is not configured");}
async close():Promise<void>{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();}
}
export class OrderedWorkspaceWriterCapabilitySet {
private live=true; private constructor(private readonly caps:Map<CanonicalWorkspaceId,WorkspaceWriterLockCapability>){ }
static make(caps:Map<CanonicalWorkspaceId,WorkspaceWriterLockCapability>){return new OrderedWorkspaceWriterCapabilitySet(caps);}
invalidate(){this.live=false;}
get workspaceIds(){if(!this.live)throw fail();return [...this.caps.keys()];}
async forWorkspace<T>(workspaceId:CanonicalWorkspaceId,action:(lease:{workspaceId:CanonicalWorkspaceId;rootLease:BorrowedVerifiedWorkspaceLockRootLease;writerCapability:WorkspaceWriterLockCapability})=>Promise<T>):Promise<T>{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<T>(action:(lease:{workspaceId:CanonicalWorkspaceId;rootLease:BorrowedVerifiedWorkspaceLockRootLease;writerCapability:WorkspaceWriterLockCapability})=>Promise<T>){const out:T[]=[];for(const id of this.workspaceIds)out.push(await this.forWorkspace(id,action));return out;}
}
export async function runUnderOrderedWorkspaceWriterLocks<T>(rootLeases:readonly VerifiedWorkspaceLockRootLease[],action:(capabilities:OrderedWorkspaceWriterCapabilitySet)=>Promise<T>):Promise<T>{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<T>(rootLease:VerifiedWorkspaceLockRootLease,action:(capability:WorkspaceWriterLockCapability)=>Promise<T>){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";}}