import * as plugins from './plugins.js'; import { controllerPackageName, type IControllerStatus, type IReq_ControllerStatus, type IReq_ControllerUpgradeLaunch, type IReq_ControllerUpgradePrepareBegin, type IReq_ControllerUpgradeFinalize, } from '../ts_interfaces/index.js'; import { controllerProcessIdentityIsLive, signalVerifiedControllerProcess, } from './classes.processinspection.js'; const currentCliName = String(controllerPackageName) === 'agl' ? 'agl' : 'hcon'; const currentRuntimeName = currentCliName === 'agl' ? 'AGL' : currentCliName; export interface IControllerStatusExpectation { packageName: string; packageVersion: string; protocolVersion?: number; upgradeManagementVersion: number; } const harnessStates = ['starting', 'ready', 'stopping', 'stopped', 'failed'] as const; export const controllerUpgradePrepareBeginRequestTimeoutMs = 30_000; const isControllerHarnessStatus = (statusArg: unknown): boolean => { if (typeof statusArg !== 'object' || statusArg === null) return false; const status = statusArg as Record; return ( (status.harnessId === 'opencode' || status.harnessId === 'flex' || status.harnessId === 'codex') && harnessStates.some((stateArg) => status.state === stateArg) && typeof status.healthy === 'boolean' && (status.pid === undefined || Number.isSafeInteger(status.pid)) && (status.version === undefined || typeof status.version === 'string') && (status.supportsLocalAttachments === undefined || typeof status.supportsLocalAttachments === 'boolean') ); }; const isControllerUpgradeStatus = (statusArg: unknown): boolean => { if (!statusArg || typeof statusArg !== 'object' || Array.isArray(statusArg)) return false; const status = statusArg as Record; return Object.keys(status).length === 3 && typeof status.fromVersion === 'string' && status.fromVersion.length > 0 && Buffer.byteLength(status.fromVersion, 'utf8') <= 128 && typeof status.toVersion === 'string' && status.toVersion.length > 0 && Buffer.byteLength(status.toVersion, 'utf8') <= 128 && typeof status.phase === 'string' && ['preparing', 'pausing', 'installing', 'restarting', 'continuing', 'completed', 'failed'] .includes(status.phase); }; export const assertControllerStatus = ( statusArg: unknown, expectationArg: IControllerStatusExpectation, ): IControllerStatus => { if (typeof statusArg !== 'object' || statusArg === null) { throw new Error(`The service on the management port is not an ${currentRuntimeName} endpoint.`); } const status = statusArg as Partial; const harnesses = Array.isArray(status.harnesses) ? status.harnesses : []; const openCode = harnesses.find((entryArg) => entryArg?.harnessId === 'opencode'); const flex = harnesses.find((entryArg) => entryArg?.harnessId === 'flex'); if ( status.packageName !== expectationArg.packageName || status.packageVersion !== expectationArg.packageVersion || (expectationArg.protocolVersion !== undefined && status.protocolVersion !== expectationArg.protocolVersion) || status.upgradeManagementVersion !== expectationArg.upgradeManagementVersion || !Number.isSafeInteger(status.controllerPid) || !Number.isSafeInteger(status.processGroupId) || typeof status.processFingerprint !== 'string' || status.processFingerprint.length < 8 || status.processFingerprint.length > 256 || (status.processMode !== 'detached' && status.processMode !== 'foreground') || (status.lifecycleState !== 'starting' && status.lifecycleState !== 'ready' && status.lifecycleState !== 'stopping' && status.lifecycleState !== 'stopped') || typeof status.startedAt !== 'number' || typeof status.setupRequired !== 'boolean' || (status.upgrade !== undefined && !isControllerUpgradeStatus(status.upgrade)) || harnesses.length !== (Number(status.protocolVersion) >= 28 ? 3 : 2) || (Number(status.protocolVersion) >= 28 && !harnesses.some((entry) => entry.harnessId === 'codex')) || !openCode || !flex || !harnesses.every((entryArg) => isControllerHarnessStatus(entryArg)) ) { throw new Error('The service on the management port has an incompatible controller identity.'); } return status as IControllerStatus; }; export const requestControllerUpgradePrepareBegin = async ( portArg: number, tokenArg: string, targetVersionArg: string, gracePeriodMsArg: number, ): Promise => { const abortController = new AbortController(); const timeout = setTimeout( () => abortController.abort(), controllerUpgradePrepareBeginRequestTimeoutMs, ); let socket: plugins.typedsocket.TypedSocket | undefined; try { socket = await plugins.typedsocket.TypedSocket.createClient( new plugins.typedrequest.TypedRouter(), `http://127.0.0.1:${portArg}`, { autoReconnect: false, maxRetries: 0, abortSignal: abortController.signal, }, ); const response = await socket.createTypedRequest( 'controller.upgrade.prepare.begin', undefined, { timeoutMs: controllerUpgradePrepareBeginRequestTimeoutMs, abortSignal: abortController.signal, }, ).fire({ token: tokenArg, targetVersion: targetVersionArg, gracePeriodMs: gracePeriodMsArg }); if ( !response || typeof response !== 'object' || Array.isArray(response) || Object.keys(response).length !== 1 || response.accepted !== true ) throw new Error('The controller returned an invalid upgrade preparation response.'); return response; } finally { clearTimeout(timeout); await socket?.stop(); } }; export const requestControllerUpgradeFinalize = async ( portArg: number, tokenArg: string, modeArg: IReq_ControllerUpgradeFinalize['request']['mode'], timeoutMsArg = 60_000, ): Promise => { const socket = await plugins.typedsocket.TypedSocket.createClient( new plugins.typedrequest.TypedRouter(), `http://127.0.0.1:${portArg}`, { autoReconnect: false, maxRetries: 0 }, ); try { return await socket.createTypedRequest( 'controller.upgrade.finalize', undefined, { timeoutMs: timeoutMsArg }, ).fire({ token: tokenArg, mode: modeArg }); } finally { await socket.stop(); } }; export const queryControllerStatus = async ( portArg: number, expectationArg: IControllerStatusExpectation, timeoutMsArg = 2_500, ): Promise => { const abortController = new AbortController(); const timeout = setTimeout(() => abortController.abort(), timeoutMsArg); let socket: plugins.typedsocket.TypedSocket | undefined; try { const router = new plugins.typedrequest.TypedRouter(); socket = await plugins.typedsocket.TypedSocket.createClient( router, `http://127.0.0.1:${portArg}`, { autoReconnect: false, maxRetries: 0, abortSignal: abortController.signal, }, ); const request = socket.createTypedRequest( 'controller.status', undefined, { timeoutMs: timeoutMsArg, abortSignal: abortController.signal }, ); return assertControllerStatus(await request.fire({}), expectationArg); } finally { clearTimeout(timeout); await socket?.stop(); } }; export const tryQueryControllerStatus = async ( portArg: number, expectationArg: IControllerStatusExpectation, ): Promise => { try { return await queryControllerStatus(portArg, expectationArg); } catch { return undefined; } }; export const isLoopbackPortListening = async ( portArg: number, timeoutMsArg = 500, ): Promise => { return await new Promise((resolve) => { const socket = plugins.net.createConnection({ host: '127.0.0.1', port: portArg }); let settled = false; const finish = (resultArg: boolean) => { if (settled) return; settled = true; clearTimeout(timeout); socket.destroy(); resolve(resultArg); }; const timeout = setTimeout(() => finish(false), timeoutMsArg); socket.once('connect', () => finish(true)); socket.once('error', () => finish(false)); }); }; export const requestControllerUpgradeLaunch = async ( portArg: number, tokenArg: string, timeoutMsArg = 15_000, ): Promise => { const abortController = new AbortController(); const timeout = setTimeout(() => abortController.abort(), timeoutMsArg); let socket: plugins.typedsocket.TypedSocket | undefined; try { const router = new plugins.typedrequest.TypedRouter(); socket = await plugins.typedsocket.TypedSocket.createClient( router, `http://127.0.0.1:${portArg}`, { autoReconnect: false, maxRetries: 0, abortSignal: abortController.signal, }, ); const response = await socket.createTypedRequest( 'controller.upgrade.launch', undefined, { timeoutMs: timeoutMsArg, abortSignal: abortController.signal }, ).fire({ token: tokenArg }); if ( !response || typeof response !== 'object' || response.accepted !== true || !Number.isSafeInteger(response.workerPid) || response.workerPid < 2 || typeof response.logFilePath !== 'string' || !plugins.path.isAbsolute(response.logFilePath) ) { throw new Error(`The controller returned an invalid ${currentCliName} upgrade launch response.`); } return response; } finally { clearTimeout(timeout); await socket?.stop(); } }; export const waitUntilControllerExited = async ( optionsArg: { pid: number; cliPath: string; port: number; expectedProcessGroupId: number; expectedFingerprint: string; command: '__serve' | 'foreground'; }, timeoutMsArg: number, ): Promise => { const deadline = Date.now() + timeoutMsArg; while (Date.now() < deadline) { if (!await controllerProcessIdentityIsLive(optionsArg.pid, optionsArg.expectedFingerprint)) return true; await new Promise((resolve) => setTimeout(resolve, 150)); } return !await controllerProcessIdentityIsLive(optionsArg.pid, optionsArg.expectedFingerprint); }; export const resolveControllerStopSignal = ( statusArg: IControllerStatus, phaseArg: 'cooperative' | 'hard', ): { signal: NodeJS.Signals; signalProcessGroup: boolean } => { if (phaseArg === 'cooperative') { return { signal: 'SIGTERM', signalProcessGroup: false }; } const signalProcessGroup = statusArg.processMode === 'detached'; if (signalProcessGroup && statusArg.processGroupId !== statusArg.controllerPid) { throw new Error('The detached controller is not the verified leader of its process group.'); } return { signal: 'SIGKILL', signalProcessGroup }; }; export const stopVerifiedController = async (optionsArg: { status: IControllerStatus; cliPath: string; port: number; }): Promise => { const command: '__serve' | 'foreground' = optionsArg.status.processMode === 'detached' ? '__serve' : 'foreground'; const processIdentity = { pid: optionsArg.status.controllerPid, cliPath: optionsArg.cliPath, port: optionsArg.port, expectedProcessGroupId: optionsArg.status.processGroupId, expectedFingerprint: optionsArg.status.processFingerprint, command, }; const cooperativeSignal = resolveControllerStopSignal(optionsArg.status, 'cooperative'); await signalVerifiedControllerProcess({ ...processIdentity, ...cooperativeSignal, }); if (await waitUntilControllerExited(processIdentity, 30_000)) return; const hardSignal = resolveControllerStopSignal(optionsArg.status, 'hard'); await signalVerifiedControllerProcess({ ...processIdentity, ...hardSignal, }); if (!await waitUntilControllerExited(processIdentity, 15_000)) { throw new Error( 'The exact controller process survived verified hard-stop escalation.', ); } };