import * as path from "node:path"; import { MonitorSocketServer } from "../lib/hermes-monitor-socket.ts"; import { createOwnedHandleRegistry, recordWaitOnlyCancellation } from "./monitor-control.ts"; export interface MonitorLifecycleConfig { runtimeDir:string; profileId:string; profilePath:string; } export function monitorLifecycleConfig(env:NodeJS.ProcessEnv):MonitorLifecycleConfig|null{const profileId=env.AGENT_FLEET_PROFILE_ID,runtimeDir=env.AGENT_FLEET_MONITOR_RUNTIME_DIR;if(!profileId||!runtimeDir||!path.isAbsolute(runtimeDir)||!/^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$/.test(profileId)||profileId.includes(".."))return null;return{runtimeDir,profileId,profilePath:path.join(runtimeDir,"profiles",profileId)};} export function createMonitorLifecycle(deps:any){const control=createOwnedHandleRegistry(deps),handles=new Map(),completedCancels=new Map(),setTimer=deps.setInterval??setInterval,clearTimer=deps.clearInterval??clearInterval;let socket:MonitorSocketServer|null=null,registration:any=null,renewTimer:any=null,reconcileTimer:any=null,retryTimer:any=null,unavailable=false,bridgeRef:any=null,bridgeInput:any=null,epoch=0;const clearRenew=()=>{if(renewTimer!==null){clearTimer(renewTimer);renewTimer=null;}if(reconcileTimer!==null){clearTimer(reconcileTimer);reconcileTimer=null;}if(retryTimer!==null){(deps.clearTimeout??clearTimeout)(retryTimer);retryTimer=null;}};const cleanup=()=>{clearRenew();registration?.cleanup?.();registration=null;};const api:any={ async startBridge(bridge:any,input:any){const startEpoch=epoch;bridgeRef=bridge;bridgeInput=input;const value=await api.start({...input,snapshot:()=>{bridge.prune();return bridge.snapshot();},output:(request:any)=>bridge.readOutput(request),events:input.events,invoke:input.invoke});if(!value||startEpoch!==epoch||bridgeRef!==bridge){if(value){try{await socket?.close();}catch{}value.cleanup?.();}return null;}bridge.setEventIdentity?.({profileKey:value.profileKey,hubInstanceId:input.hubInstanceId});bridge.setCurrentOwner?.({ownerSessionId:value.ownerId,ownerLeaseExpiresAt:value.leaseExpiresAt,updateActive:true});bridge.publishHubEvent?.("hub.capability_changed",{capabilities:{events:!!input.events,invoke:!!input.invoke}});reconcileTimer=setTimer(()=>{const probeEpoch=epoch;const reconcile=(task:any,e:any)=>{if(probeEpoch!==epoch||bridgeRef!==bridge)return;bridge.reconcile({...e,taskId:`${task.id}:${task.generation}`});};try{for(const task of bridge.snapshot?.().tasks??[]){if(probeEpoch!==epoch||bridgeRef!==bridge)break;Promise.resolve(deps.getRecoveryEvidence?.(task)??deps.reconcileEvidence?.()??{}).then(e=>reconcile(task,e)).catch(()=>reconcile(task,{transient:true}));}}catch{if(probeEpoch===epoch&&bridgeRef===bridge)bridge.reconcile({transient:true});}},30_000);value.cancel=(request:any)=>bridge.cancelTask(request);return value;}, lowLevelCancelOwnedGeneration(request:any){return api.cancelOwnedGeneration(request);}, async start(input:any){const startEpoch=epoch;try{registration=deps.registry.register(input);registration.output=input.output;registration.events=input.events;registration.invoke=input.invoke;registration.cancel=(request:any)=>api.cancelOwnedGeneration(request);socket=(deps.createSocketServer??((r:any)=>new MonitorSocketServer(r)))(registration);await socket.listen();if(startEpoch!==epoch){try{await socket.close();}catch{}registration?.cleanup?.();registration=null;socket=null;return null;}if(registration.renew&®istration.leaseMs){renewTimer=setTimer(()=>{try{Promise.resolve(registration?.renew()).then((renewed:any)=>{if(renewed?.leaseExpiresAt){registration.leaseExpiresAt=renewed.leaseExpiresAt;bridgeRef?.setCurrentOwner?.({ownerSessionId:registration.ownerId,ownerLeaseExpiresAt:renewed.leaseExpiresAt,updateActive:true});}}).catch(()=>api.disable());}catch{void api.disable();}},Math.max(1,Math.floor(registration.leaseMs/2)));}return registration;}catch(error:any){try{await socket?.close();}catch{}socket=null;cleanup();if(error?.leaseExpiresAt){const delay=Math.max(0,new Date(error.leaseExpiresAt).getTime()-(deps.now?.()??Date.now()));const retryEpoch=epoch;retryTimer=(deps.setTimeout??setTimeout)(()=>{retryTimer=null;if(retryEpoch!==epoch)return;void (bridgeRef?api.startBridge(bridgeRef,bridgeInput):api.start(input)).catch(()=>{});},delay);return null;}throw error;}}, cancelOwnedGeneration(request:any){const key=`${request.taskId}:${request.generation}`;if(completedCancels.has(key))return Promise.resolve(completedCancels.get(key));const handle=handles.get(key);if(!handle)return Promise.resolve({cancelled:false,reason:"unsupported"});return control.cancel({...request,profileKey:registration?.profileKey,token:registration?.token,handle}).then((result:any)=>{completedCancels.set(key,result);return result;});}, registerOwnedGeneration(value:any){const proc:any=value.process;const exitPromise=value.exitPromise??(proc?.exitCode!==null&&proc?.exitCode!==undefined?Promise.resolve(true):new Promise(resolve=>proc?.once?.("close",()=>resolve(true))));const owned={...value,exitPromise,profileKey:registration?.profileKey??value.profileKey,token:registration?.token??value.token};const handle=control.register(owned);handles.set(`${owned.taskId}:${owned.generation}`,handle);return handle;},async disable(){unavailable=true;await api.stop();},isUnavailable(){return unavailable;},isAlive(){return !!socket&&!unavailable;},recordWaitOnlyCancellation, async stop(){epoch++;bridgeRef=null;bridgeInput=null;clearRenew();try{await socket?.close();}finally{socket=null;registration?.cleanup?.();registration=null;}} };return api;}