import * as plugins from './plugins.js'; import { commitinfo } from './00_commitinfo_data.js'; import { controllerPackageName, controllerProtocolVersion, controllerUpgradeManagementVersion, type IControllerStatus, } from '../ts_interfaces/index.js'; import { resolveAGLHomePaths } from './classes.aglhome.js'; import { isLoopbackPortListening, queryControllerStatus, requestControllerUpgradeFinalize, requestControllerUpgradePrepareBegin, stopVerifiedController, } from './classes.controllermanagement.js'; import { inspectControllerProcess, listControllerDataWriterProcessesForCli, listControllerDataWriterProcessesForCliPaths, listControllerProcessesForCli, processIdentityHasCliCommand, readControllerProcessIdentity, readProcessGroupMemberPids, type IControllerProcessIdentity, } from './classes.processinspection.js'; import { createUpgradeToken, compareUpgradeSemver, type TUpgradeTransaction, type IUpgradeTransactionInventoryEntry, type IUpgradeTransactionV3, type IUpgradeExpectedController, type IUpgradeWorkerPayload, normalizeUpgradeRegistryUrl, UpgradeCoordinator, type UpgradeInstallationLock, upgradeCoordinationVersion, upgradePackageTransitionTransactionVersion, upgradePackageTransitionSource, upgradePackageTransitionTarget, upgradeTransactionIsStalled, upgradeTransactionRequiresForwardRecovery, upgradeTransactionTargetVersion, serializeUpgradeWorkerPayload, terminateVerifiedUpgradeWorkerProcessGroup, } from './classes.upgradecoordinator.js'; import { upgradeCoordinationRootEnvironmentVariable, upgradeTokenEnvironmentVariable, } from './constants.upgradeenvironment.js'; const commandOutputByteLimit = 1024 * 1024; const loggedCommandOutputByteLimit = 16 * 1024 * 1024; const packageManagerTimeoutMs = 10 * 60 * 1000; const controllerRestartTimeoutMs = 90_000; const commandGroupDrainTimeoutMs = 5_000; const retainedUpgradeLogCount = 10; const minimumPrunableUpgradeLogAgeMs = 30 * 60 * 1000; const upgradePreparationPollMs = 250; const upgradePreparationSourceCheckMs = 500; const upgradePreparationAcceptanceReconciliationMs = 5_000; const upgradePreparationCleanupAllowanceMs = 60_000; const scrubUpgradeCoordinationEnvironment = (): void => { delete process.env[upgradeTokenEnvironmentVariable]; delete process.env[upgradeCoordinationRootEnvironmentVariable]; }; const upgradeLogPattern = /^upgrade-[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{6}\.[0-9]{3}Z-[A-Za-z0-9_-]{10}\.log$/; let upgradeWorkerCliPath: string | undefined; const currentPackageIdentity = String(controllerPackageName) === upgradePackageTransitionSource.packageName ? upgradePackageTransitionSource : upgradePackageTransitionTarget; const currentCliName = currentPackageIdentity.cliName; const registryArguments = (registryUrlArg?: string): string[] => registryUrlArg ? [`--registry=${registryUrlArg}`] : []; const packageManagerEnvironmentAllowlist = new Set([ 'CI', 'COREPACK_HOME', 'FORCE_COLOR', 'HOME', 'HTTP_PROXY', 'HTTPS_PROXY', 'LANG', 'LOGNAME', 'NODE_EXTRA_CA_CERTS', 'NO_PROXY', 'NPM_CONFIG_REGISTRY', 'NPM_CONFIG_USERCONFIG', 'PATH', 'PNPM_CONFIG_REGISTRY', 'PNPM_HOME', 'SHELL', 'SSL_CERT_DIR', 'SSL_CERT_FILE', 'TEMP', 'TERM', 'TMP', 'TMPDIR', 'TZ', 'USER', 'XDG_CACHE_HOME', 'XDG_CONFIG_HOME', 'XDG_DATA_HOME', 'XDG_RUNTIME_DIR', 'XDG_STATE_HOME', ]); class UpgradeCommandCleanupError extends Error {} interface IGlobalListDependency { version?: unknown; path?: unknown; } interface IGlobalListResult { path?: unknown; dependencies?: Record; } export interface IPnpmGlobalInstallation { packageManagerPath: string; globalRoot: string; globalBinDirectory: string; packageDirectory: string; cliPath: string; packageVersion: string; } export interface IDetachedUpgradeWorker { pid: number; processGroupId: number; fingerprint: string; cliPath: string; logFilePath: string; } export interface IUpgradeWorkerLaunchCandidate { pid: number; cliPath: string; logFilePath: string; processGroupId?: number; fingerprint?: string; } export class UpgradeWorkerHandoffError extends Error { constructor( messageArg: string, causeArg: unknown, public readonly recoveryLock: UpgradeInstallationLock, ) { super(messageArg, { cause: causeArg }); this.name = 'UpgradeWorkerHandoffError'; } } export class UpgradeWorkerCandidateDrainageError extends Error { constructor(messageArg: string, causeArg: unknown) { super(messageArg, { cause: causeArg }); this.name = 'UpgradeWorkerCandidateDrainageError'; } } export interface ITerminalizeUpgradeWorkerHandoffFailureOptions { coordinator: UpgradeCoordinator; token: string; error: UpgradeWorkerHandoffError; } export const terminalizeUpgradeWorkerHandoffFailure = async ( optionsArg: ITerminalizeUpgradeWorkerHandoffFailureOptions, ): Promise => { const handoffFailure = optionsArg.error.cause ?? optionsArg.error; let terminalizationError: unknown; let terminalized = false; try { await optionsArg.error.recoveryLock.assertOwned(); await optionsArg.coordinator.finishTransaction( optionsArg.token, false, 'The upgrade worker handoff failed before package work began.', handoffFailure instanceof Error ? handoffFailure.message : String(handoffFailure), ); terminalized = true; } catch (errorArg) { terminalizationError = errorArg; } let releaseError: unknown; if (terminalized) { try { await optionsArg.error.recoveryLock.release(); } catch (errorArg) { releaseError = errorArg; } } const failures = [optionsArg.error, terminalizationError, releaseError] .filter((failureArg) => failureArg !== undefined); throw failures.length > 1 ? new AggregateError( failures, 'The upgrade worker handoff failed and its transaction cleanup was incomplete.', ) : optionsArg.error; }; const proxyEnvironmentValueIsSafe = (valueArg: string): boolean => { try { const parsed = new URL(valueArg); return parsed.username.length === 0 && parsed.password.length === 0; } catch { return false; } }; export const createPackageManagerEnvironment = ( sourceArg: NodeJS.ProcessEnv = process.env, ): NodeJS.ProcessEnv => { const environment: NodeJS.ProcessEnv = {}; for (const [name, value] of Object.entries(sourceArg)) { if (value === undefined) continue; const normalizedName = name.toUpperCase(); if ( !packageManagerEnvironmentAllowlist.has(normalizedName) && !normalizedName.startsWith('LC_') ) continue; if ( (normalizedName === 'HTTP_PROXY' || normalizedName === 'HTTPS_PROXY') && !proxyEnvironmentValueIsSafe(value) ) continue; if (normalizedName === 'NPM_CONFIG_REGISTRY' || normalizedName === 'PNPM_CONFIG_REGISTRY') { try { normalizeUpgradeRegistryUrl(value); } catch { continue; } } environment[name] = value; } return environment; }; const findExecutable = async (nameArg: string, environmentArg: NodeJS.ProcessEnv): Promise => { const pathValue = environmentArg.PATH; if (!pathValue) throw new Error(`Unable to locate ${nameArg}: PATH is empty.`); for (const directory of pathValue.split(plugins.path.delimiter)) { if (!directory) continue; const candidate = plugins.path.join(directory, nameArg); try { await plugins.fs.promises.access(candidate, plugins.fs.constants.X_OK); return await plugins.fs.promises.realpath(candidate); } catch { continue; } } throw new Error(`Unable to locate ${nameArg} on PATH.`); }; interface IManagedCommandOptions { executable: string; arguments: string[]; environment: NodeJS.ProcessEnv; timeoutMs: number; captureOutput: boolean; outputByteLimit?: number; } interface IPreparedManagedCommand extends IManagedCommandOptions { cwd: string; workerIdentity: NonNullable>>; ownsWorkerProcessGroup: boolean; } const prepareManagedCommand = async ( optionsArg: IManagedCommandOptions, ): Promise => { const workerIdentity = await readControllerProcessIdentity(process.pid); if (!workerIdentity) throw new Error('Unable to establish upgrade command ownership.'); const ownsWorkerProcessGroup = workerIdentity.processGroupLeader && await processIdentityHasCliCommand( workerIdentity, upgradeWorkerCliPath ?? await resolveCurrentCliPath(), '__upgrade-worker', { allowMissingAbsolutePath: true }, ); return { ...optionsArg, cwd: plugins.os.homedir(), workerIdentity, ownsWorkerProcessGroup, }; }; const runPreparedManagedCommand = async ( optionsArg: IPreparedManagedCommand, ): Promise<{ stdout: string; stderr: string }> => { const child = plugins.childProcess.spawn(optionsArg.executable, optionsArg.arguments, { cwd: optionsArg.cwd, detached: !optionsArg.ownsWorkerProcessGroup, env: optionsArg.environment, shell: false, stdio: ['ignore', 'pipe', 'pipe'], }); const outcomePromise = new Promise<{ code: number | null; signal: NodeJS.Signals | null; error?: Error; }>((resolve) => { child.once('error', (errorArg) => resolve({ code: null, signal: null, error: errorArg })); child.once('close', (codeArg, signalArg) => resolve({ code: codeArg, signal: signalArg })); }); const commandProcessGroupId = optionsArg.ownsWorkerProcessGroup ? optionsArg.workerIdentity.processGroupId : child.pid; let stdout = ''; let stderr = ''; let outputBytes = 0; let outputError: Error | undefined; const recordOutput = ( streamArg: 'stdout' | 'stderr', chunkArg: Buffer, sourceArg?: NodeJS.ReadableStream, ) => { outputBytes += chunkArg.byteLength; if (outputBytes > (optionsArg.outputByteLimit ?? commandOutputByteLimit)) { outputError ??= new Error('Upgrade command output exceeded its size limit.'); return; } if (optionsArg.captureOutput) { if (streamArg === 'stdout') stdout += chunkArg.toString('utf8'); else stderr += chunkArg.toString('utf8'); return; } const destination = streamArg === 'stdout' ? process.stdout : process.stderr; if (!destination.write(chunkArg) && sourceArg && 'pause' in sourceArg && 'resume' in sourceArg) { sourceArg.pause(); destination.once('drain', () => sourceArg.resume()); } }; child.stdout?.on('data', (chunkArg: Buffer) => recordOutput('stdout', chunkArg, child.stdout!)); child.stderr?.on('data', (chunkArg: Buffer) => recordOutput('stderr', chunkArg, child.stderr!)); const listCommandPids = async (): Promise => { if (!commandProcessGroupId) { return child.pid && child.exitCode === null && child.signalCode === null ? [child.pid] : []; } const members = await readProcessGroupMemberPids(commandProcessGroupId); return optionsArg.ownsWorkerProcessGroup ? members.filter((processIdArg) => processIdArg !== process.pid) : members; }; const signalDirectChild = (signalArg: NodeJS.Signals): void => { if (!child.pid || child.exitCode !== null || child.signalCode !== null) return; try { process.kill(child.pid, signalArg); } catch (errorArg) { if ((errorArg as NodeJS.ErrnoException).code !== 'ESRCH') throw errorArg; } }; const signalCommand = async (signalArg: NodeJS.Signals): Promise => { let processIds: number[]; try { processIds = await listCommandPids(); } catch (errorArg) { signalDirectChild(signalArg); throw errorArg; } for (const processId of processIds) { try { process.kill(processId, signalArg); } catch (errorArg) { if ((errorArg as NodeJS.ErrnoException).code !== 'ESRCH') throw errorArg; } } }; const waitForCommandDrain = async (timeoutMsArg: number): Promise => { const deadline = Date.now() + timeoutMsArg; while (Date.now() < deadline) { if ((await listCommandPids()).length === 0) return true; await new Promise((resolve) => setTimeout(resolve, 50)); } return (await listCommandPids()).length === 0; }; let interruption: string | undefined; let terminating: Promise | undefined; let terminationError: unknown; const terminate = (): Promise => { if (terminating) return terminating; terminating = (async () => { try { await signalCommand('SIGTERM'); if (await waitForCommandDrain(commandGroupDrainTimeoutMs)) return; await signalCommand('SIGKILL'); if (!await waitForCommandDrain(commandGroupDrainTimeoutMs)) { throw new Error('The upgrade command process group did not terminate.'); } } catch (errorArg) { terminationError = errorArg; signalDirectChild('SIGKILL'); throw errorArg; } })(); return terminating; }; const beginTermination = (reasonArg: string): void => { interruption ??= reasonArg; void terminate().catch(() => undefined); }; const onSigint = () => beginTermination('SIGINT'); const onSigterm = () => beginTermination('SIGTERM'); process.on('SIGINT', onSigint); process.on('SIGTERM', onSigterm); const timeout = setTimeout(() => beginTermination('timeout'), optionsArg.timeoutMs); const outputLimitPoll = setInterval(() => { if (outputError) beginTermination('output limit'); }, 25); try { const outcome = await outcomePromise; try { if (!await waitForCommandDrain(2_000)) { beginTermination('surviving child processes'); } } catch (errorArg) { terminationError = errorArg; signalDirectChild('SIGKILL'); } if (terminating) await terminating.catch(() => undefined); if (terminationError) { throw new UpgradeCommandCleanupError('Upgrade command cleanup failed.', { cause: terminationError, }); } if (interruption) { throw new Error(`${plugins.path.basename(optionsArg.executable)} was interrupted by ${interruption}.`); } if (outputError) throw outputError; if (outcome.error) throw outcome.error; if (outcome.code !== 0) { throw new Error( `${plugins.path.basename(optionsArg.executable)} exited with ${outcome.signal ?? outcome.code ?? 'an unknown result'}${stderr.trim() ? `: ${stderr.trim()}` : '.'}`, ); } return { stdout, stderr }; } finally { clearTimeout(timeout); clearInterval(outputLimitPoll); process.removeListener('SIGINT', onSigint); process.removeListener('SIGTERM', onSigterm); } }; const runManagedCommand = async ( optionsArg: IManagedCommandOptions, ): Promise<{ stdout: string; stderr: string }> => await runPreparedManagedCommand( await prepareManagedCommand(optionsArg), ); const runCapturedCommand = async ( executableArg: string, argumentsArg: string[], environmentArg: NodeJS.ProcessEnv, timeoutMsArg = 30_000, outputByteLimitArg = commandOutputByteLimit, ): Promise => { const result = await runManagedCommand({ executable: executableArg, arguments: argumentsArg, environment: environmentArg, timeoutMs: timeoutMsArg, captureOutput: true, outputByteLimit: outputByteLimitArg, }); return result.stdout; }; const assertSingleLineAbsolutePath = async (valueArg: string, nameArg: string): Promise => { const value = valueArg.trim(); if (!value || value.includes('\n') || !plugins.path.isAbsolute(value)) { throw new Error(`${nameArg} did not return one absolute path.`); } return await plugins.fs.promises.realpath(value); }; const readGlobalListDependency = ( outputArg: string, packageNameArg: string, ): { version: string; path: string } => { let value: unknown; try { value = JSON.parse(outputArg) as unknown; } catch { throw new Error('pnpm list --global returned malformed JSON.'); } if (!Array.isArray(value)) { throw new Error('pnpm list --global returned an unexpected result.'); } const dependencies = value.flatMap((entryArg) => { if (!entryArg || typeof entryArg !== 'object' || Array.isArray(entryArg)) return []; const entry = entryArg as IGlobalListResult; const dependency = entry.dependencies?.[packageNameArg]; return dependency ? [dependency] : []; }); if (dependencies.length !== 1) { throw new Error(`Expected exactly one pnpm-global installation of ${packageNameArg}.`); } const dependency = dependencies[0]; if ( typeof dependency.version !== 'string' || typeof dependency.path !== 'string' || !plugins.path.isAbsolute(dependency.path) ) { throw new Error('pnpm reported invalid global package metadata.'); } return { version: dependency.version, path: dependency.path }; }; export const resolvePnpmGlobalInstallation = async (): Promise => { if (process.platform !== 'linux' && process.platform !== 'darwin') { throw new Error(`${currentCliName} upgrade is currently supported on Linux and macOS only.`); } const environment = createPackageManagerEnvironment(); const packageManagerPath = await findExecutable('pnpm', environment); const globalRootOutput = await runCapturedCommand( packageManagerPath, ['root', '--global'], environment, ); const globalBinOutput = await runCapturedCommand( packageManagerPath, ['bin', '--global'], environment, ); const listOutput = await runCapturedCommand( packageManagerPath, ['list', '--global', controllerPackageName, '--depth=0', '--json'], environment, ); const globalRoot = await assertSingleLineAbsolutePath(globalRootOutput, 'pnpm root --global'); const globalBinDirectory = await assertSingleLineAbsolutePath( globalBinOutput, 'pnpm bin --global', ); const dependency = readGlobalListDependency(listOutput, controllerPackageName); const dependencyRelativePath = plugins.path.relative(globalRoot, dependency.path); if ( dependencyRelativePath.startsWith(`..${plugins.path.sep}`) || dependencyRelativePath === '..' || plugins.path.isAbsolute(dependencyRelativePath) ) { throw new Error('pnpm reported a global package path outside its global root.'); } const packageDirectory = await plugins.fs.promises.realpath(dependency.path); const packageJsonPath = plugins.path.join(packageDirectory, 'package.json'); const packageJsonStats = await plugins.fs.promises.stat(packageJsonPath); if (!packageJsonStats.isFile() || packageJsonStats.size < 2 || packageJsonStats.size > 64 * 1024) { throw new Error(`The pnpm-global ${controllerPackageName} package metadata is not a bounded file.`); } const packageJson = JSON.parse(await plugins.fs.promises.readFile(packageJsonPath, 'utf8')) as { name?: unknown; version?: unknown; bin?: unknown; controllerUpgradeManagementVersion?: unknown; }; const currentIdentity: { cliName: string; managementVersion: number } = currentPackageIdentity; if ( packageJson.name !== controllerPackageName || packageJson.version !== dependency.version || packageJson.controllerUpgradeManagementVersion !== currentIdentity.managementVersion || Number(controllerUpgradeManagementVersion) !== currentIdentity.managementVersion || !packageJson.bin || typeof packageJson.bin !== 'object' || Array.isArray(packageJson.bin) || !Object.keys(packageJson.bin).length || Object.keys(packageJson.bin).length !== 1 || (packageJson.bin as Record)[currentIdentity.cliName] !== './cli.js' ) { throw new Error(`The pnpm-global ${controllerPackageName} package metadata is incompatible.`); } const cliPath = await plugins.fs.promises.realpath(plugins.path.join(packageDirectory, 'cli.js')); if (!(await plugins.fs.promises.stat(cliPath)).isFile()) { throw new Error('The pnpm-global controller direct CLI is not a file.'); } return { packageManagerPath, globalRoot, globalBinDirectory, packageDirectory, cliPath, packageVersion: dependency.version, }; }; interface IPackageTransitionContext { packageManagerPath: string; globalRoot: string; globalBinDirectory: string; } interface IExactGlobalPackageInstallation extends IPnpmGlobalInstallation { packageName: string; cliName: string; cliRelativePath: string; listedPackagePath: string; groupDirectory: string; } interface IPackageTransitionTopology { source?: IExactGlobalPackageInstallation; target?: IExactGlobalPackageInstallation; } type TPackageTransitionTopology = 'source-only' | 'grouped' | 'target-only'; const parseGlobalListDependencies = ( outputArg: string, ): { path: string; dependencies: Record } => { let value: unknown; try { value = JSON.parse(outputArg) as unknown; } catch { throw new Error('pnpm list --global returned malformed JSON.'); } if (!Array.isArray(value) || value.length !== 1) { throw new Error('pnpm list --global returned an unexpected package-transition result.'); } const entry = value[0]; if (!entry || typeof entry !== 'object' || Array.isArray(entry)) { throw new Error('pnpm list --global returned invalid package-transition metadata.'); } const raw = entry as IGlobalListResult; if (typeof raw.path !== 'string' || !plugins.path.isAbsolute(raw.path)) { throw new Error('pnpm list --global omitted its global package root.'); } const dependencies: Record = {}; for (const packageName of [ upgradePackageTransitionSource.packageName, upgradePackageTransitionTarget.packageName, ]) { const dependency = raw.dependencies?.[packageName]; if (!dependency) continue; if ( typeof dependency.version !== 'string' || typeof dependency.path !== 'string' || !plugins.path.isAbsolute(dependency.path) ) throw new Error(`pnpm reported invalid global metadata for ${packageName}.`); dependencies[packageName] = { version: dependency.version, path: dependency.path, }; } return { path: raw.path, dependencies }; }; const packageGroupDirectory = ( globalRootArg: string, packagePathArg: string, packageNameArg: string, ): string => { const relative = plugins.path.relative(globalRootArg, packagePathArg); if ( !relative || relative === '..' || relative.startsWith(`..${plugins.path.sep}`) || plugins.path.isAbsolute(relative) ) throw new Error(`pnpm reported ${packageNameArg} outside its global root.`); const expectedSuffix = plugins.path.join('node_modules', ...packageNameArg.split('/')); const suffix = `${plugins.path.sep}${expectedSuffix}`; if (!relative.endsWith(suffix)) { throw new Error(`pnpm reported an unexpected global path for ${packageNameArg}.`); } const groupRelativePath = relative.slice(0, -suffix.length); if (!groupRelativePath || groupRelativePath.includes(plugins.path.sep)) { throw new Error(`pnpm reported an invalid isolated group for ${packageNameArg}.`); } return plugins.path.join(globalRootArg, groupRelativePath); }; const readExactGlobalPackage = async ( contextArg: IPackageTransitionContext, dependencyArg: { version: string; path: string }, identityArg: { packageName: string; version: string; managementVersion: number; cliName: string; cliRelativePath: string; }, ): Promise => { if (dependencyArg.version !== identityArg.version) { throw new Error( `The pnpm-global ${identityArg.packageName} version is ${dependencyArg.version}, not ${identityArg.version}.`, ); } const groupDirectory = packageGroupDirectory( contextArg.globalRoot, dependencyArg.path, identityArg.packageName, ); const packageDirectory = await plugins.fs.promises.realpath(dependencyArg.path); const packageJsonPath = plugins.path.join(packageDirectory, 'package.json'); const packageJsonStats = await plugins.fs.promises.stat(packageJsonPath); if (!packageJsonStats.isFile() || packageJsonStats.size < 2 || packageJsonStats.size > 64 * 1024) { throw new Error(`The pnpm-global ${identityArg.packageName} metadata is not a bounded file.`); } let packageJson: unknown; try { packageJson = JSON.parse(await plugins.fs.promises.readFile(packageJsonPath, 'utf8')) as unknown; } catch (errorArg) { throw new Error(`The pnpm-global ${identityArg.packageName} metadata is malformed.`, { cause: errorArg, }); } if (!packageJson || typeof packageJson !== 'object' || Array.isArray(packageJson)) { throw new Error(`The pnpm-global ${identityArg.packageName} metadata is invalid.`); } const manifest = packageJson as Record; const bin = manifest.bin; if ( manifest.name !== identityArg.packageName || manifest.version !== identityArg.version || manifest.controllerUpgradeManagementVersion !== identityArg.managementVersion || !bin || typeof bin !== 'object' || Array.isArray(bin) || Object.keys(bin).length !== 1 || (bin as Record)[identityArg.cliName] !== identityArg.cliRelativePath ) throw new Error(`The pnpm-global ${identityArg.packageName} identity is incompatible.`); const cliPath = await plugins.fs.promises.realpath( plugins.path.join(packageDirectory, identityArg.cliRelativePath), ); if (!(await plugins.fs.promises.stat(cliPath)).isFile()) { throw new Error(`The pnpm-global ${identityArg.packageName} direct CLI is not a file.`); } return { packageManagerPath: contextArg.packageManagerPath, globalRoot: contextArg.globalRoot, globalBinDirectory: contextArg.globalBinDirectory, packageDirectory, cliPath, packageVersion: identityArg.version, packageName: identityArg.packageName, cliName: identityArg.cliName, cliRelativePath: identityArg.cliRelativePath, listedPackagePath: dependencyArg.path, groupDirectory, }; }; const readGlobalShimTarget = async ( binDirectoryArg: string, cliNameArg: string, ): Promise => { const shimPath = plugins.path.join(binDirectoryArg, cliNameArg); let stats: plugins.fs.Stats; try { stats = await plugins.fs.promises.lstat(shimPath); } catch (errorArg) { if ((errorArg as NodeJS.ErrnoException).code === 'ENOENT') return undefined; throw errorArg; } if ( !stats.isFile() || stats.isSymbolicLink() || stats.size < 2 || stats.size > 1024 * 1024 || (stats.mode & 0o111) === 0 ) throw new Error(`The pnpm-global ${cliNameArg} shim is not a bounded executable file.`); const contents = await plugins.fs.promises.readFile(shimPath, 'utf8'); const markers = [...contents.matchAll(/^# cmd-shim-target=(.+)$/gm)]; if (markers.length !== 1 || !plugins.path.isAbsolute(markers[0][1])) { throw new Error(`The pnpm-global ${cliNameArg} shim has no exact ownership marker.`); } try { return await plugins.fs.promises.realpath(markers[0][1]); } catch (errorArg) { throw new Error(`The pnpm-global ${cliNameArg} shim target is unavailable.`, { cause: errorArg, }); } }; const assertGlobalShimOwnership = async ( binDirectoryArg: string, cliNameArg: string, cliPathArg?: string, ): Promise => { const markerTarget = await readGlobalShimTarget(binDirectoryArg, cliNameArg); if (markerTarget === undefined) { if (cliPathArg === undefined) return; throw new Error(`The pnpm-global ${cliNameArg} shim is missing.`); } if (cliPathArg === undefined) { throw new Error(`The pnpm-global ${cliNameArg} shim survived without its owning package.`); } if (markerTarget !== cliPathArg) { throw new Error(`The pnpm-global ${cliNameArg} shim is owned by a different package.`); } }; const maximumPnpmGlobalGroups = 1_024; const findExactGlobalPackageLink = async ( globalRootArg: string, packageNameArg: string, packageDirectoryArg: string, ): Promise<{ listedPackagePath: string; groupDirectory: string }> => { const entries = await plugins.fs.promises.readdir(globalRootArg, { withFileTypes: true }); if (entries.length > maximumPnpmGlobalGroups) { throw new Error('The pnpm global root exceeds the recovery scan bound.'); } const matches: string[] = []; for (const entry of entries) { if (!entry.isDirectory()) continue; const candidate = plugins.path.join( globalRootArg, entry.name, 'node_modules', ...packageNameArg.split('/'), ); try { if (await plugins.fs.promises.realpath(candidate) === packageDirectoryArg) { matches.push(candidate); } } catch (errorArg) { const code = (errorArg as NodeJS.ErrnoException).code; if (code !== 'ENOENT' && code !== 'ENOTDIR') throw errorArg; } } if (matches.length !== 1) { throw new Error(`Expected exactly one filesystem-owned installation of ${packageNameArg}.`); } return { listedPackagePath: matches[0], groupDirectory: packageGroupDirectory(globalRootArg, matches[0], packageNameArg), }; }; const resolveAuthoritativeRecoveryInstallation = async ( transactionArg: TUpgradeTransaction, capturedInstallationArg: IPnpmGlobalInstallation, ): Promise => { const [globalRoot, globalBinDirectory, packageManagerPath] = await Promise.all([ plugins.fs.promises.realpath(capturedInstallationArg.globalRoot), plugins.fs.promises.realpath(capturedInstallationArg.globalBinDirectory), plugins.fs.promises.realpath(capturedInstallationArg.packageManagerPath), ]); if ( globalRoot !== capturedInstallationArg.globalRoot || globalBinDirectory !== capturedInstallationArg.globalBinDirectory || packageManagerPath !== capturedInstallationArg.packageManagerPath ) throw new Error('The captured pnpm recovery paths changed unexpectedly.'); await plugins.fs.promises.access(packageManagerPath, plugins.fs.constants.X_OK); const currentIdentity: { packageName: string; versions: string[]; managementVersion: number; cliName: string; cliRelativePath: string; } = transactionArg.version === upgradePackageTransitionTransactionVersion ? (() => { const identity = transactionArg.targetPackageCommitStarted === true ? upgradePackageTransitionTarget : upgradePackageTransitionSource; return { packageName: identity.packageName, versions: transactionArg.targetPackageCommitStarted === true ? [...new Set([ identity.version, transactionArg.recoveryTargetVersion ?? capturedInstallationArg.packageVersion, ])] : [identity.version], managementVersion: identity.managementVersion, cliName: identity.cliName, cliRelativePath: identity.cliRelativePath, }; })() : (() => { const identity = String(controllerPackageName) === upgradePackageTransitionTarget.packageName ? upgradePackageTransitionTarget : upgradePackageTransitionSource; return { packageName: String(controllerPackageName), versions: [ transactionArg.sourceVersion, ...(transactionArg.targetVersion ? [transactionArg.targetVersion] : []), ], managementVersion: Number(controllerUpgradeManagementVersion), cliName: identity.cliName, cliRelativePath: identity.cliRelativePath, }; })(); const cliPath = await readGlobalShimTarget(globalBinDirectory, currentIdentity.cliName); if (!cliPath) { throw new Error(`The authoritative pnpm-global ${currentIdentity.cliName} shim is missing.`); } const packageDirectory = plugins.path.dirname(cliPath); if ( await plugins.fs.promises.realpath( plugins.path.join(packageDirectory, currentIdentity.cliRelativePath), ) !== cliPath ) throw new Error('The authoritative recovery CLI path is incompatible.'); const packageJsonPath = plugins.path.join(packageDirectory, 'package.json'); const packageJsonStats = await plugins.fs.promises.stat(packageJsonPath); if (!packageJsonStats.isFile() || packageJsonStats.size < 2 || packageJsonStats.size > 64 * 1024) { throw new Error('The authoritative recovery package metadata is not a bounded file.'); } let packageJson: unknown; try { packageJson = JSON.parse(await plugins.fs.promises.readFile(packageJsonPath, 'utf8')) as unknown; } catch (errorArg) { throw new Error('The authoritative recovery package metadata is malformed.', { cause: errorArg, }); } if (!packageJson || typeof packageJson !== 'object' || Array.isArray(packageJson)) { throw new Error('The authoritative recovery package metadata is invalid.'); } const manifest = packageJson as Record; const bin = manifest.bin; if ( manifest.name !== currentIdentity.packageName || typeof manifest.version !== 'string' || !currentIdentity.versions.includes(manifest.version) || manifest.controllerUpgradeManagementVersion !== currentIdentity.managementVersion || !bin || typeof bin !== 'object' || Array.isArray(bin) || Object.keys(bin).length !== 1 || (bin as Record)[currentIdentity.cliName] !== currentIdentity.cliRelativePath ) throw new Error('The authoritative recovery package identity is incompatible.'); if (!(await plugins.fs.promises.stat(cliPath)).isFile()) { throw new Error('The authoritative recovery CLI is not a file.'); } await findExactGlobalPackageLink(globalRoot, currentIdentity.packageName, packageDirectory); return { packageManagerPath, globalRoot, globalBinDirectory, packageDirectory, cliPath, packageVersion: manifest.version, }; }; const verifyPackageTransitionTopology = async ( contextArg: IPackageTransitionContext, expectedArg: TPackageTransitionTopology, targetVersionArg: string = upgradePackageTransitionTarget.version, ): Promise => { const output = await runCapturedCommand( contextArg.packageManagerPath, [ 'list', '--global', upgradePackageTransitionSource.packageName, upgradePackageTransitionTarget.packageName, '--depth=0', '--json', ], createPackageManagerEnvironment(), ); const listed = parseGlobalListDependencies(output); if (await plugins.fs.promises.realpath(listed.path) !== contextArg.globalRoot) { throw new Error('The pnpm global root changed during the package transition.'); } const sourceDependency = listed.dependencies[upgradePackageTransitionSource.packageName]; const targetDependency = listed.dependencies[upgradePackageTransitionTarget.packageName]; const sourceExpected = expectedArg !== 'target-only'; const targetExpected = expectedArg !== 'source-only'; if (Boolean(sourceDependency) !== sourceExpected || Boolean(targetDependency) !== targetExpected) { throw new Error(`The pnpm-global package topology is not ${expectedArg}.`); } const [source, target] = await Promise.all([ sourceDependency ? readExactGlobalPackage(contextArg, sourceDependency, upgradePackageTransitionSource) : undefined, targetDependency ? readExactGlobalPackage(contextArg, targetDependency, { ...upgradePackageTransitionTarget, version: targetVersionArg, }) : undefined, ]); if (expectedArg === 'grouped' && source?.groupDirectory !== target?.groupDirectory) { throw new Error('The source and target packages are not installed in one pnpm group.'); } await Promise.all([ assertGlobalShimOwnership( contextArg.globalBinDirectory, upgradePackageTransitionSource.cliName, source?.cliPath, ), assertGlobalShimOwnership( contextArg.globalBinDirectory, upgradePackageTransitionTarget.cliName, target?.cliPath, ), ]); return { ...(source ? { source } : {}), ...(target ? { target } : {}), }; }; const resolvePackageManagerContext = async (): Promise => { if (process.platform !== 'linux' && process.platform !== 'darwin') { throw new Error(`${currentCliName} upgrade is currently supported on Linux and macOS only.`); } const environment = createPackageManagerEnvironment(); const packageManagerPath = await findExecutable('pnpm', environment); const globalRoot = await assertSingleLineAbsolutePath( await runCapturedCommand(packageManagerPath, ['root', '--global'], environment), 'pnpm root --global', ); const globalBinDirectory = await assertSingleLineAbsolutePath( await runCapturedCommand(packageManagerPath, ['bin', '--global'], environment), 'pnpm bin --global', ); return { packageManagerPath, globalRoot, globalBinDirectory }; }; const canRestartUpgradeSource = (transactionArg: TUpgradeTransaction): boolean => !transactionArg.terminal && transactionArg.controllerWasRunning && transactionArg.targetStartupInvoked !== true && !upgradeTransactionRequiresForwardRecovery(transactionArg); const isUpgradeSourceRestart = (transactionArg: TUpgradeTransaction): boolean => canRestartUpgradeSource(transactionArg) && transactionArg.phase === 'restarting'; export const resolveInheritedUpgradeTargetInstallation = async ( transactionArg: TUpgradeTransaction, ): Promise => { if (transactionArg.version === upgradeCoordinationVersion) { return await resolvePnpmGlobalInstallation(); } const restoringSource = isUpgradeSourceRestart(transactionArg); const topology = await verifyPackageTransitionTopology( await resolvePackageManagerContext(), restoringSource ? 'source-only' : 'target-only', upgradeTransactionTargetVersion(transactionArg)!, ); if (restoringSource) { if (!topology.source) throw new Error('The restored source package is unavailable.'); return topology.source; } if (!topology.target) throw new Error('The committed target package is unavailable.'); return topology.target; }; export const resolveCurrentCliPath = async (): Promise => { return await plugins.fs.promises.realpath( plugins.url.fileURLToPath(new URL('../cli.js', import.meta.url)), ); }; const currentPackageIsDevelopmentCheckout = async (): Promise => { const packageRoot = plugins.path.dirname(await resolveCurrentCliPath()); try { const stats = await plugins.fs.promises.lstat(plugins.path.join(packageRoot, '.git')); return stats.isDirectory() || stats.isFile(); } catch (errorArg) { if ((errorArg as NodeJS.ErrnoException).code === 'ENOENT') return false; throw errorArg; } }; export const resolveCurrentPnpmGlobalInstallation = async (): Promise => { const [installation, cliPath] = await Promise.all([ resolvePnpmGlobalInstallation(), resolveCurrentCliPath(), ]); if (installation.cliPath !== cliPath || installation.packageVersion !== commitinfo.version) { throw new Error(`${currentCliName} upgrade must be run from the active pnpm-global installation.`); } return installation; }; export const tryResolveCurrentPnpmGlobalInstallation = async ( ): Promise => { if (await currentPackageIsDevelopmentCheckout()) return undefined; return await resolveCurrentPnpmGlobalInstallation(); }; export const compareSemver = compareUpgradeSemver; export interface IInspectOrphanedUpgradeForAdoptionOptions { coordinator: UpgradeCoordinator; installation: IPnpmGlobalInstallation; port: number; registryUrl?: string; } export const inspectOrphanedUpgradeForAdoption = async ( optionsArg: IInspectOrphanedUpgradeForAdoptionOptions, ): Promise => { const inventory = await optionsArg.coordinator.inspectCanonicalTransactionInventory(); const nonterminal = inventory.transactions.filter((entryArg) => !entryArg.transaction.terminal); if (nonterminal.length === 0) return undefined; if (nonterminal.length !== 1) { throw new Error('Multiple nonterminal upgrade transactions block a fresh upgrade.'); } const entry = nonterminal[0]; const transaction = entry.transaction; const eligibleV2 = transaction.version === upgradeCoordinationVersion && transaction.controllerWasRunning === true && transaction.targetStartupInvoked === true && transaction.targetVersion !== undefined && transaction.preparationCompletedAt !== undefined && (transaction.phase === 'restarting' || transaction.phase === 'continuing') && compareSemver(transaction.sourceVersion, upgradePackageTransitionTarget.version) >= 0 && compareSemver(transaction.sourceVersion, transaction.targetVersion) < 0 && compareSemver(transaction.targetVersion, optionsArg.installation.packageVersion) <= 0; const eligibleV3 = transaction.version === upgradePackageTransitionTransactionVersion && currentPackageIdentity.packageName === upgradePackageTransitionTarget.packageName && transaction.targetPackageCommitStarted === true && (!transaction.controllerWasRunning || transaction.preparationCompletedAt !== undefined) && ( transaction.phase === 'installing' || transaction.phase === 'restarting' || transaction.phase === 'continuing' ) && compareSemver( upgradeTransactionTargetVersion(transaction)!, optionsArg.installation.packageVersion, ) <= 0; if ( transaction.port !== optionsArg.port || (!eligibleV2 && !eligibleV3) || !upgradeTransactionIsStalled(transaction) ) throw new Error('A nonterminal upgrade transaction blocks a fresh upgrade.'); if (transaction.version === upgradePackageTransitionTransactionVersion && optionsArg.registryUrl) { throw new Error( 'Legacy package-transition recovery uses pnpm registry configuration and cannot retain --registry.', ); } if ( transaction.version === upgradeCoordinationVersion && transaction.registryUrl !== undefined && optionsArg.registryUrl !== undefined && transaction.registryUrl !== optionsArg.registryUrl ) throw new Error('The orphaned upgrade registry conflicts with the requested registry.'); if (await optionsArg.coordinator.transactionWorkerIsLive(transaction)) { throw new Error('A nonterminal upgrade worker is still running.'); } if ( transaction.controller && await optionsArg.coordinator.transactionControllerIsLive(transaction) ) throw new Error('A nonterminal upgrade controller is still running.'); const adoptionMatch = /^transaction-([a-f0-9]{20})\.json\.adopting-([a-f0-9]{20})$/ .exec(entry.fileName); const prefixes = new Set([ transaction.tokenHash.slice(0, 20), ...(adoptionMatch ? [adoptionMatch[1], adoptionMatch[2]] : []), ]); if (inventory.tokenMetadataFileNames.some((fileNameArg) => ( [...prefixes].some((prefixArg) => fileNameArg.includes(`-${prefixArg}.json`)) ))) throw new Error('The orphaned upgrade retains token-bound coordination metadata.'); return entry; }; export interface IAdoptOrphanedUpgradeUnderLockOptions extends IInspectOrphanedUpgradeForAdoptionOptions { lock: UpgradeInstallationLock; token: string; expected: IUpgradeTransactionInventoryEntry; } export interface IAdoptOrphanedUpgradeUnderLockDependencies { listDataWriterProcesses?: typeof listControllerDataWriterProcessesForCliPaths; } export const adoptOrphanedUpgradeUnderLock = async ( optionsArg: IAdoptOrphanedUpgradeUnderLockOptions, dependenciesArg: IAdoptOrphanedUpgradeUnderLockDependencies = {}, ): Promise => { await optionsArg.lock.assertOwned(); const candidate = await inspectOrphanedUpgradeForAdoption(optionsArg); if ( !candidate || candidate.transaction.tokenHash !== optionsArg.expected.transaction.tokenHash || candidate.transaction.revision !== optionsArg.expected.transaction.revision ) throw new Error('The orphaned upgrade transaction changed before adoption.'); const cliPaths = [ optionsArg.installation.cliPath, ...(candidate.transaction.worker ? [candidate.transaction.worker.cliPath] : []), ...(candidate.transaction.controller ? [candidate.transaction.controller.cliPath] : []), ]; const listDataWriterProcesses = dependenciesArg.listDataWriterProcesses ?? listControllerDataWriterProcessesForCliPaths; const assertNoRuntimeWriters = async (): Promise => { if (await isLoopbackPortListening(optionsArg.port)) { throw new Error('The orphaned upgrade port became occupied before recovery.'); } const writers = await listDataWriterProcesses(cliPaths, { includeInstalledPackagePaths: true, }); if (writers.length > 0) { throw new Error(`A ${writers[0].kind} data writer prevents orphaned upgrade recovery.`); } }; await assertNoRuntimeWriters(); await optionsArg.lock.assertOwned(); const adopted = await optionsArg.coordinator.adoptOrphanedTransaction({ lock: optionsArg.lock, token: optionsArg.token, expectedTokenHash: candidate.transaction.tokenHash, expectedRevision: candidate.transaction.revision, port: optionsArg.port, installedVersion: optionsArg.installation.packageVersion, ...(optionsArg.registryUrl === undefined ? {} : { registryUrl: optionsArg.registryUrl }), }); await assertNoRuntimeWriters(); await optionsArg.lock.assertOwned(); return adopted; }; const readLatestVersion = async ( installationArg: IPnpmGlobalInstallation, registryUrlArg?: string, ): Promise => { const output = await runCapturedCommand( installationArg.packageManagerPath, [ 'view', controllerPackageName, 'dist-tags.latest', ...registryArguments(registryUrlArg), '--json', ], createPackageManagerEnvironment(), ); let version: unknown; try { version = JSON.parse(output) as unknown; } catch { throw new Error('pnpm view returned malformed JSON for the registry latest tag.'); } if (typeof version !== 'string') { throw new Error('The registry latest tag did not resolve to one package version.'); } compareUpgradeSemver(version, version); return version; }; const resolveExactPackageTransitionTarget = async ( contextArg: IPackageTransitionContext, registryUrlArg?: string, ): Promise => { const output = await runCapturedCommand( contextArg.packageManagerPath, [ 'view', `${upgradePackageTransitionTarget.packageName}@${upgradePackageTransitionTarget.version}`, ...registryArguments(registryUrlArg), '--json', ], createPackageManagerEnvironment(), ); let metadata: unknown; try { metadata = JSON.parse(output) as unknown; } catch (errorArg) { throw new Error('pnpm view returned malformed target package metadata.', { cause: errorArg }); } if (!metadata || typeof metadata !== 'object' || Array.isArray(metadata)) { throw new Error('pnpm view returned invalid target package metadata.'); } const manifest = metadata as Record; const bin = manifest.bin; if ( manifest.name !== upgradePackageTransitionTarget.packageName || manifest.version !== upgradePackageTransitionTarget.version || manifest.controllerUpgradeManagementVersion !== upgradePackageTransitionTarget.managementVersion || !bin || typeof bin !== 'object' || Array.isArray(bin) || Object.keys(bin).length !== 1 || (bin as Record)[upgradePackageTransitionTarget.cliName] !== upgradePackageTransitionTarget.cliRelativePath ) throw new Error('The exact target package metadata is incompatible.'); }; const openUpgradeLog = async (tokenArg: string): Promise<{ handle: plugins.fs.promises.FileHandle; path: string; }> => { const directory = resolveAGLHomePaths().logs; await plugins.fs.promises.mkdir(directory, { recursive: true, mode: 0o700 }); const directoryStats = await plugins.fs.promises.lstat(directory); if ( !directoryStats.isDirectory() || directoryStats.isSymbolicLink() || (typeof process.getuid === 'function' && directoryStats.uid !== process.getuid()) || (directoryStats.mode & 0o077) !== 0 ) { throw new Error(`The ${currentCliName} upgrade log directory is unsafe: ${directory}`); } const existingLogs: Array<{ path: string; mtimeMs: number }> = []; for (const entry of await plugins.fs.promises.readdir(directory, { withFileTypes: true })) { if (!entry.isFile() || !upgradeLogPattern.test(entry.name)) continue; const existingPath = plugins.path.join(directory, entry.name); const stats = await plugins.fs.promises.lstat(existingPath); if ( !stats.isFile() || stats.isSymbolicLink() || (typeof process.getuid === 'function' && stats.uid !== process.getuid()) || (stats.mode & 0o077) !== 0 ) { throw new Error(`An existing ${currentCliName} upgrade log is unsafe: ${existingPath}`); } existingLogs.push({ path: existingPath, mtimeMs: stats.mtimeMs }); } existingLogs.sort((leftArg, rightArg) => rightArg.mtimeMs - leftArg.mtimeMs); for (const staleLog of existingLogs.slice(retainedUpgradeLogCount - 1)) { if (Date.now() - staleLog.mtimeMs < minimumPrunableUpgradeLogAgeMs) continue; await plugins.fs.promises.unlink(staleLog.path).catch((errorArg) => { if ((errorArg as NodeJS.ErrnoException).code !== 'ENOENT') throw errorArg; }); } const filePath = plugins.path.join( directory, `upgrade-${new Date().toISOString().replaceAll(':', '')}-${plugins.crypto .createHash('sha256') .update(tokenArg, 'utf8') .digest('hex') .slice(0, 10)}.log`, ); const handle = await plugins.fs.promises.open( filePath, plugins.fs.constants.O_WRONLY | plugins.fs.constants.O_APPEND | plugins.fs.constants.O_CREAT | plugins.fs.constants.O_EXCL | plugins.fs.constants.O_NOFOLLOW, 0o600, ); return { handle, path: filePath }; }; interface IUpgradeRecoveryWorkerContext { version: 1; mode: 'recovery'; canonicalGlobalRoot: string; packageManagerPath: string; globalBinDirectory: string; cliPath: string; } const serializeUpgradeRecoveryWorkerContext = ( contextArg: IUpgradeRecoveryWorkerContext, ): string => Buffer.from(JSON.stringify(contextArg), 'utf8').toString('base64url'); const parseUpgradeRecoveryWorkerContext = ( valueArg: unknown, ): IUpgradeRecoveryWorkerContext => { if (typeof valueArg !== 'string' || valueArg.length < 16 || valueArg.length > 16 * 1024) { throw new Error('The internal upgrade recovery worker context is invalid.'); } let parsed: unknown; try { parsed = JSON.parse(Buffer.from(valueArg, 'base64url').toString('utf8')) as unknown; } catch (errorArg) { throw new Error('The internal upgrade recovery worker context is malformed.', { cause: errorArg, }); } if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) { throw new Error('The internal upgrade recovery worker context is malformed.'); } const context = parsed as Record; if ( Object.keys(context).length !== 6 || ![ 'version', 'mode', 'canonicalGlobalRoot', 'packageManagerPath', 'globalBinDirectory', 'cliPath', ].every((keyArg) => Object.hasOwn(context, keyArg)) || context.version !== 1 || context.mode !== 'recovery' || typeof context.canonicalGlobalRoot !== 'string' || typeof context.packageManagerPath !== 'string' || typeof context.globalBinDirectory !== 'string' || typeof context.cliPath !== 'string' || !plugins.path.isAbsolute(context.canonicalGlobalRoot) || !plugins.path.isAbsolute(context.packageManagerPath) || !plugins.path.isAbsolute(context.globalBinDirectory) || !plugins.path.isAbsolute(context.cliPath) ) throw new Error('The internal upgrade recovery worker context is incompatible.'); return context as unknown as IUpgradeRecoveryWorkerContext; }; const terminateSpawnedDetachedWorker = async ( identityArg: IControllerProcessIdentity, cliPathArg: string, ): Promise => await terminateVerifiedUpgradeWorkerProcessGroup({ pid: identityArg.pid, processGroupId: identityArg.processGroupId, fingerprint: identityArg.fingerprint, cliPath: cliPathArg, }); export interface ICandidateCleanupDependencies { beforeAttempt?: (processIdArg?: number) => Promise; } const cleanupCandidateUntilDrained = async ( cleanupArg: () => Promise, dependenciesArg: ICandidateCleanupDependencies = {}, processIdArg?: number, ): Promise => { await dependenciesArg.beforeAttempt?.(processIdArg); await cleanupArg(); }; export const terminateUpgradeWorkerCandidateUntilDrained = async ( coordinatorArg: UpgradeCoordinator, candidateArg: IDetachedUpgradeWorker, ): Promise => await cleanupCandidateUntilDrained( async () => await coordinatorArg.terminateUpgradeWorkerCandidate(candidateArg), {}, candidateArg.pid, ); interface IDetachedUpgradeWorkerProcessDependencies { readProcessIdentity?: typeof readControllerProcessIdentity; processHasCliCommand?: typeof processIdentityHasCliCommand; candidateCleanup?: ICandidateCleanupDependencies; } const launchDetachedUpgradeWorkerProcess = async (optionsArg: { installation: IPnpmGlobalInstallation; payload: IUpgradeWorkerPayload; environment?: NodeJS.ProcessEnv; internalArguments?: string[]; }, dependenciesArg: IDetachedUpgradeWorkerProcessDependencies = {}): Promise => { if (process.platform !== 'linux' && process.platform !== 'darwin') { throw new Error(`${currentCliName} upgrade is currently supported on Linux and macOS only.`); } const log = await openUpgradeLog(optionsArg.payload.token); let child: plugins.childProcess.ChildProcess | undefined; let spawnedIdentity: IControllerProcessIdentity | undefined; let logClosed = false; try { const workerEnvironment = { ...(optionsArg.environment ?? process.env) }; const inheritedToken = workerEnvironment[upgradeTokenEnvironmentVariable]; const useLegacyCoordination = workerEnvironment[upgradeCoordinationRootEnvironmentVariable] === undefined && inheritedToken === optionsArg.payload.token; if (useLegacyCoordination) { delete workerEnvironment[upgradeCoordinationRootEnvironmentVariable]; } else if (workerEnvironment[upgradeCoordinationRootEnvironmentVariable] === undefined) { workerEnvironment[upgradeCoordinationRootEnvironmentVariable] = new UpgradeCoordinator( optionsArg.installation.globalRoot, ).baseDirectory; } const launchedChild = plugins.childProcess.spawn( process.execPath, [ optionsArg.installation.cliPath, '__upgrade-worker', serializeUpgradeWorkerPayload(optionsArg.payload), ...(optionsArg.internalArguments ?? []), ], { cwd: plugins.os.homedir(), detached: true, env: { ...workerEnvironment, [upgradeTokenEnvironmentVariable]: optionsArg.payload.token, }, shell: false, stdio: ['ignore', log.handle.fd, log.handle.fd], }, ); child = launchedChild; await new Promise((resolve, reject) => { launchedChild.once('spawn', resolve); launchedChild.once('error', reject); }); if (!launchedChild.pid) { throw new Error(`The detached ${currentCliName} upgrade worker has no process ID.`); } spawnedIdentity = await readControllerProcessIdentity(launchedChild.pid) ?? undefined; await log.handle.close(); logClosed = true; const identity = await (dependenciesArg.readProcessIdentity ? dependenciesArg.readProcessIdentity(launchedChild.pid) : spawnedIdentity); if ( !identity || !identity.processGroupLeader || identity.processGroupId !== launchedChild.pid || !await (dependenciesArg.processHasCliCommand ?? processIdentityHasCliCommand)( identity, optionsArg.installation.cliPath, '__upgrade-worker', ) ) throw new Error('The detached upgrade worker candidate identity is invalid.'); launchedChild.unref(); return { pid: identity.pid, processGroupId: identity.processGroupId, fingerprint: identity.fingerprint, cliPath: optionsArg.installation.cliPath, logFilePath: log.path, }; } catch (errorArg) { const errors: unknown[] = [errorArg]; let candidateCleanupError: unknown; if (!logClosed) { try { await log.handle.close(); logClosed = true; } catch (logCleanupErrorArg) { errors.push(logCleanupErrorArg); } } if (child?.pid) { try { if (!spawnedIdentity) { throw new Error('The detached upgrade worker cleanup identity was not established.'); } await cleanupCandidateUntilDrained( async () => terminateSpawnedDetachedWorker( spawnedIdentity!, optionsArg.installation.cliPath, ), dependenciesArg.candidateCleanup, child.pid, ); } catch (cleanupErrorArg) { candidateCleanupError = cleanupErrorArg; errors.push(cleanupErrorArg); } } if (candidateCleanupError !== undefined) { throw new UpgradeWorkerCandidateDrainageError( 'The detached upgrade worker launch failed and candidate drainage is unproven.', new AggregateError(errors, 'Detached upgrade worker launch cleanup was incomplete.'), ); } if (errors.length > 1) { throw new AggregateError( errors, 'The detached upgrade worker launch failed and cleanup was incomplete.', ); } throw errorArg; } }; interface ILaunchDetachedUpgradeWorkerOptions { installation: IPnpmGlobalInstallation; payload: IUpgradeWorkerPayload; environment?: NodeJS.ProcessEnv; } export function launchDetachedUpgradeWorker( optionsArg: ILaunchDetachedUpgradeWorkerOptions, ): Promise; export async function launchDetachedUpgradeWorker( optionsArg: ILaunchDetachedUpgradeWorkerOptions, dependenciesArg: IDetachedUpgradeWorkerProcessDependencies = {}, ): Promise { return await launchDetachedUpgradeWorkerProcess(optionsArg, dependenciesArg); } const launchDetachedRecoveryWorker = async (optionsArg: { installation: IPnpmGlobalInstallation; payload: IUpgradeWorkerPayload; environment?: NodeJS.ProcessEnv; }): Promise => await launchDetachedUpgradeWorkerProcess({ ...optionsArg, internalArguments: [serializeUpgradeRecoveryWorkerContext({ version: 1, mode: 'recovery', canonicalGlobalRoot: optionsArg.installation.globalRoot, packageManagerPath: optionsArg.installation.packageManagerPath, globalBinDirectory: optionsArg.installation.globalBinDirectory, cliPath: optionsArg.installation.cliPath, })], }); const runLoggedCommand = async (optionsArg: { executable: string; arguments: string[]; environment: NodeJS.ProcessEnv; timeoutMs: number; }): Promise => { await runManagedCommand({ ...optionsArg, captureOutput: false, outputByteLimit: loggedCommandOutputByteLimit, }); }; const installVersion = async ( installationArg: IPnpmGlobalInstallation, versionArg: string, registryUrlArg?: string, ): Promise => { await runLoggedCommand({ executable: installationArg.packageManagerPath, arguments: packageAddArguments(`${controllerPackageName}@${versionArg}`, registryUrlArg), environment: createPackageManagerEnvironment(), timeoutMs: packageManagerTimeoutMs, }); }; const packageAddArguments = (packageSpecArg: string, registryUrlArg?: string): string[] => [ 'add', '--global', '--save-exact', '--yes', '--reporter=append-only', ...registryArguments(registryUrlArg), packageSpecArg, ]; const preparePackageAdd = async ( contextArg: IPackageTransitionContext, packageSpecArg: string, registryUrlArg?: string, ): Promise => await prepareManagedCommand({ executable: contextArg.packageManagerPath, arguments: packageAddArguments(packageSpecArg, registryUrlArg), environment: createPackageManagerEnvironment(), timeoutMs: packageManagerTimeoutMs, captureOutput: false, outputByteLimit: loggedCommandOutputByteLimit, }); const addExactPackage = async ( contextArg: IPackageTransitionContext, packageSpecArg: string, registryUrlArg?: string, ): Promise => { await runPreparedManagedCommand( await preparePackageAdd(contextArg, packageSpecArg, registryUrlArg), ); }; const assertExpectedController = ( statusArg: IControllerStatus, expectedArg: IUpgradeExpectedController, ): void => { if ( statusArg.controllerPid !== expectedArg.pid || statusArg.processGroupId !== expectedArg.processGroupId || statusArg.processFingerprint !== expectedArg.processFingerprint || statusArg.lifecycleState !== 'ready' ) { throw new Error('The running controller changed after the upgrade request was authorized.'); } }; const prepareControllerStart = async ( installationArg: IPnpmGlobalInstallation, portArg: number, tokenArg: string, coordinatorArg: UpgradeCoordinator, ): Promise => await prepareManagedCommand({ executable: process.execPath, arguments: [installationArg.cliPath, 'start', '--port', String(portArg)], environment: { ...process.env, [upgradeTokenEnvironmentVariable]: tokenArg, [upgradeCoordinationRootEnvironmentVariable]: coordinatorArg.baseDirectory, }, timeoutMs: controllerRestartTimeoutMs, captureOutput: false, outputByteLimit: loggedCommandOutputByteLimit, }); const startController = async ( installationArg: IPnpmGlobalInstallation, portArg: number, tokenArg: string, coordinatorArg: UpgradeCoordinator, ): Promise => await runPreparedManagedCommand( await prepareControllerStart(installationArg, portArg, tokenArg, coordinatorArg), ).then(() => undefined); export const authorizeUpgradeSourceStartup = async ( coordinatorArg: UpgradeCoordinator, tokenArg: string, portArg: number, sourceVersionArg: string, ): Promise => coordinatorArg.mutateTransaction(tokenArg, (current) => { if ( !canRestartUpgradeSource(current) || current.port !== portArg || current.sourceVersion !== sourceVersionArg ) throw new Error('The upgrade transaction cannot restart the restored source package.'); return { ...current, phase: 'restarting', message: `Restarting restored source ${current.sourceVersion}.`, }; }); const startRestoredController = async ( installationArg: IPnpmGlobalInstallation, portArg: number, tokenArg: string, coordinatorArg: UpgradeCoordinator, ): Promise => { const start = await prepareControllerStart(installationArg, portArg, tokenArg, coordinatorArg); await authorizeUpgradeSourceStartup( coordinatorArg, tokenArg, portArg, installationArg.packageVersion, ); await runPreparedManagedCommand(start); }; const retainedTargetTransactionMatches = ( transactionArg: TUpgradeTransaction, installationArg: IPnpmGlobalInstallation, ): transactionArg is TUpgradeTransaction & { targetVersion: string; targetStartupInvoked: true; controller: NonNullable; } => Boolean( !transactionArg.terminal && transactionArg.controllerWasRunning && transactionArg.targetStartupInvoked === true && (transactionArg.version === upgradeCoordinationVersion || transactionArg.targetPackageCommitted === true) && (transactionArg.phase === 'restarting' || transactionArg.phase === 'continuing') && upgradeTransactionTargetVersion(transactionArg) === installationArg.packageVersion && transactionArg.controller?.command === '__serve' && transactionArg.controller.cliPath === installationArg.cliPath && transactionArg.controller.pid === transactionArg.controller.processGroupId && /^linux:[1-9][0-9]*:[1-9][0-9]*$/.test(transactionArg.controller.fingerprint) ); export interface IRetainedUpgradeTargetAuthorityOptions { coordinator: UpgradeCoordinator; token: string; } export interface IRetainedUpgradeTargetAuthorityDependencies { platform?: NodeJS.Platform; resolveInstallation?: ( transactionArg: TUpgradeTransaction, ) => Promise; inspectProcess?: typeof inspectControllerProcess; } export const retainedUpgradeTargetIsAuthoritative = async ( optionsArg: IRetainedUpgradeTargetAuthorityOptions, dependenciesArg: IRetainedUpgradeTargetAuthorityDependencies = {}, ): Promise => { if ((dependenciesArg.platform ?? process.platform) !== 'linux') return false; const resolveInstallation = dependenciesArg.resolveInstallation ?? resolveInheritedUpgradeTargetInstallation; const inspectProcess = dependenciesArg.inspectProcess ?? inspectControllerProcess; const transaction = await optionsArg.coordinator.readTransaction(optionsArg.token); const installation = await resolveInstallation(transaction); if (!retainedTargetTransactionMatches(transaction, installation)) return false; const controller = transaction.controller; const identity = await inspectProcess({ pid: controller.pid, cliPath: controller.cliPath, port: transaction.port, expectedProcessGroupId: controller.processGroupId, expectedFingerprint: controller.fingerprint, command: controller.command, }); if (!identity?.processGroupLeader || identity.fingerprint !== controller.fingerprint) return false; const currentTransaction = await optionsArg.coordinator.readTransaction(optionsArg.token); const currentInstallation = await resolveInstallation(currentTransaction); if ( !retainedTargetTransactionMatches(currentTransaction, currentInstallation) || currentTransaction.controller.pid !== controller.pid || currentTransaction.controller.processGroupId !== controller.processGroupId || currentTransaction.controller.fingerprint !== controller.fingerprint || currentTransaction.controller.cliPath !== controller.cliPath ) return false; const currentIdentity = await inspectProcess({ pid: currentTransaction.controller.pid, cliPath: currentTransaction.controller.cliPath, port: currentTransaction.port, expectedProcessGroupId: currentTransaction.controller.processGroupId, expectedFingerprint: currentTransaction.controller.fingerprint, command: currentTransaction.controller.command, }); return Boolean( currentIdentity?.processGroupLeader && currentIdentity.fingerprint === currentTransaction.controller.fingerprint ); }; const updateWorkerProgress = async ( coordinatorArg: UpgradeCoordinator, tokenArg: string, phaseArg: TUpgradeTransaction['phase'], messageArg: string, ): Promise => coordinatorArg.mutateTransaction(tokenArg, (current) => ({ ...current, phase: phaseArg, message: messageArg, })); const grantControllerAction = async ( coordinatorArg: UpgradeCoordinator, tokenArg: string, statusArg: IControllerStatus, actionArg: 'prepare' | 'finalize', modeArg?: 'continue' | 'compensate' | 'reopen', ): Promise => coordinatorArg.createControllerActionGrant({ token: tokenArg, action: actionArg, ...(modeArg === undefined ? {} : { mode: modeArg }), packageVersion: statusArg.packageVersion, controller: { pid: statusArg.controllerPid, processGroupId: statusArg.processGroupId, processFingerprint: statusArg.processFingerprint, }, }); const queryUpgradeManagedController = async ( portArg: number, identityArg: { packageName: string; packageVersion: string; managementVersion: number; }, ): Promise => queryControllerStatus(portArg, { packageName: identityArg.packageName, packageVersion: identityArg.packageVersion, upgradeManagementVersion: identityArg.managementVersion, }); const finalizeControllerUpgrade = async ( coordinatorArg: UpgradeCoordinator, tokenArg: string, statusArg: IControllerStatus, modeArg: 'continue' | 'compensate' | 'reopen', ): Promise => { await grantControllerAction(coordinatorArg, tokenArg, statusArg, 'finalize', modeArg); let response: Awaited>; try { response = await requestControllerUpgradeFinalize(statusArg.controllerPid > 0 ? (await coordinatorArg.readTransaction(tokenArg)).port : 0, tokenArg, modeArg); } finally { await coordinatorArg.removeControllerActionGrant('finalize', tokenArg); } if (response.failed > 0) { throw new Error(`${response.failed} paused session${response.failed === 1 ? '' : 's'} could not be continued.`); } }; export interface IWaitForControllerUpgradePreparationOptions { coordinator: UpgradeCoordinator; token: string; targetVersion: string; gracePeriodMs: number; controllerIsLive: () => Promise; now?: () => number; delay?: (millisecondsArg: number) => Promise; } export const waitForControllerUpgradePreparation = async ( optionsArg: IWaitForControllerUpgradePreparationOptions, ): Promise => { const now = optionsArg.now ?? Date.now; const delay = optionsArg.delay ?? (async (millisecondsArg: number) => new Promise((resolve) => ( setTimeout(resolve, millisecondsArg) ))); const acceptanceReconciliationDeadline = now() + upgradePreparationAcceptanceReconciliationMs; let nextSourceCheckAt = 0; while (true) { const currentTime = now(); if (currentTime >= nextSourceCheckAt) { if (!await optionsArg.controllerIsLive()) { throw new Error('The source controller exited before upgrade preparation completed.'); } nextSourceCheckAt = currentTime + upgradePreparationSourceCheckMs; } const transaction = await optionsArg.coordinator.readTransaction(optionsArg.token); if (transaction.terminal) { throw new Error( transaction.terminal.error ?? transaction.message ?? 'The source controller rejected upgrade preparation.', ); } if ( transaction.preparationAcceptedAt === undefined || transaction.preparationDeadlineAt === undefined ) { if (currentTime >= acceptanceReconciliationDeadline) { throw new Error('The source controller did not durably accept upgrade preparation.'); } await delay(upgradePreparationPollMs); continue; } if ( transaction.targetVersion !== optionsArg.targetVersion || transaction.gracePeriodMs !== optionsArg.gracePeriodMs || transaction.preparationDeadlineAt !== transaction.preparationAcceptedAt + optionsArg.gracePeriodMs ) throw new Error('The durable upgrade preparation does not match the worker request.'); if (transaction.phase === 'stopping') { if ( transaction.preparationCompletedAt === undefined || transaction.preparationCompletedAt < transaction.preparationAcceptedAt || transaction.preparationCompletedAt > transaction.preparationDeadlineAt ) throw new Error('The durable upgrade preparation completion is invalid.'); return transaction; } if (currentTime > transaction.preparationDeadlineAt + upgradePreparationCleanupAllowanceMs) { throw new Error('Timed out waiting for upgrade preparation cleanup to settle.'); } await delay(upgradePreparationPollMs); } }; export interface IPrepareAndStopRunningControllerOptions { coordinator: UpgradeCoordinator; token: string; port: number; targetVersion: string; gracePeriodMs: number; currentInstallation: IPnpmGlobalInstallation; runningStatus: IControllerStatus; } export interface IPrepareAndStopRunningControllerDependencies { requestPrepareBegin?: typeof requestControllerUpgradePrepareBegin; stopController?: typeof stopVerifiedController; controllerIsLive?: () => Promise; now?: () => number; delay?: (millisecondsArg: number) => Promise; } export const prepareAndStopRunningController = async ( optionsArg: IPrepareAndStopRunningControllerOptions, dependenciesArg: IPrepareAndStopRunningControllerDependencies = {}, ): Promise => { await updateWorkerProgress( optionsArg.coordinator, optionsArg.token, 'preparing', `Preparing sessions for ${optionsArg.currentInstallation.packageVersion} to ${optionsArg.targetVersion}.`, ); await grantControllerAction( optionsArg.coordinator, optionsArg.token, optionsArg.runningStatus, 'prepare', ); let beginError: unknown; try { await (dependenciesArg.requestPrepareBegin ?? requestControllerUpgradePrepareBegin)( optionsArg.port, optionsArg.token, optionsArg.targetVersion, optionsArg.gracePeriodMs, ); } catch (errorArg) { beginError = errorArg; } finally { await optionsArg.coordinator.removeControllerActionGrant('prepare', optionsArg.token); } try { await waitForControllerUpgradePreparation({ coordinator: optionsArg.coordinator, token: optionsArg.token, targetVersion: optionsArg.targetVersion, gracePeriodMs: optionsArg.gracePeriodMs, controllerIsLive: dependenciesArg.controllerIsLive ?? (async () => Boolean( await inspectControllerProcess({ pid: optionsArg.runningStatus.controllerPid, cliPath: optionsArg.currentInstallation.cliPath, port: optionsArg.port, expectedProcessGroupId: optionsArg.runningStatus.processGroupId, expectedFingerprint: optionsArg.runningStatus.processFingerprint, command: optionsArg.runningStatus.processMode === 'detached' ? '__serve' : 'foreground', }), )), ...(dependenciesArg.now ? { now: dependenciesArg.now } : {}), ...(dependenciesArg.delay ? { delay: dependenciesArg.delay } : {}), }); } catch (errorArg) { if (beginError === undefined) throw errorArg; throw new AggregateError( [beginError, errorArg], 'The source controller did not complete asynchronous upgrade preparation.', ); } await (dependenciesArg.stopController ?? stopVerifiedController)({ status: optionsArg.runningStatus, cliPath: optionsArg.currentInstallation.cliPath, port: optionsArg.port, }); }; const restorePreviousInstallation = async ( installationArg: IPnpmGlobalInstallation, registryUrlArg?: string, ): Promise => { try { const current = await resolvePnpmGlobalInstallation(); if (current.packageVersion === installationArg.packageVersion) return current; } catch (errorArg) { if (errorContainsCommandCleanupFailure(errorArg)) throw errorArg; // Reinstalling the exact prior version is the authoritative rollback. } await installVersion(installationArg, installationArg.packageVersion, registryUrlArg); const restored = await resolvePnpmGlobalInstallation(); if (restored.packageVersion !== installationArg.packageVersion) { throw new Error('The previous pnpm-global package version was not restored.'); } return restored; }; const errorContainsCommandCleanupFailure = (errorArg: unknown): boolean => { if (errorArg instanceof UpgradeCommandCleanupError) return true; if (errorArg instanceof AggregateError) { return errorArg.errors.some((nestedErrorArg) => ( errorContainsCommandCleanupFailure(nestedErrorArg) )); } return errorArg instanceof Error && errorArg.cause !== undefined && errorContainsCommandCleanupFailure(errorArg.cause); }; const sourceControllerIdentity = (transactionArg: TUpgradeTransaction): { packageName: string; packageVersion: string; managementVersion: number; } => transactionArg.version === upgradeCoordinationVersion ? { packageName: controllerPackageName, packageVersion: transactionArg.sourceVersion, managementVersion: controllerUpgradeManagementVersion, } : { packageName: transactionArg.sourcePackageName, packageVersion: transactionArg.sourceVersion, managementVersion: transactionArg.sourceManagementVersion, }; const requirePackageTransitionTransaction = ( transactionArg: TUpgradeTransaction, ): IUpgradeTransactionV3 => { if (transactionArg.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return transactionArg; }; const targetControllerIdentity = (transactionArg: TUpgradeTransaction): { packageName: string; packageVersion: string; managementVersion: number; } => transactionArg.version === upgradeCoordinationVersion ? { packageName: controllerPackageName, packageVersion: transactionArg.targetVersion ?? (() => { throw new Error('The upgrade target version is unavailable.'); })(), managementVersion: controllerUpgradeManagementVersion, } : { packageName: transactionArg.targetPackageName, packageVersion: upgradeTransactionTargetVersion(transactionArg)!, managementVersion: transactionArg.targetManagementVersion, }; export const assertInheritedUpgradeTransaction = ( transactionArg: TUpgradeTransaction, currentIdentityArg: { packageName: string; packageVersion: string; managementVersion: number; cliName: string; cliRelativePath: string; }, ): void => { if (transactionArg.terminal) { throw new Error('A terminal upgrade transaction cannot authorize controller startup.'); } if (isUpgradeSourceRestart(transactionArg)) { const sourceIdentity = transactionArg.version === upgradeCoordinationVersion ? { ...sourceControllerIdentity(transactionArg), cliName: currentCliName, cliRelativePath: './cli.js', } : { ...sourceControllerIdentity(transactionArg), cliName: transactionArg.sourceCliName, cliRelativePath: transactionArg.sourceCliRelativePath, }; if ( currentIdentityArg.packageName !== sourceIdentity.packageName || currentIdentityArg.packageVersion !== sourceIdentity.packageVersion || currentIdentityArg.managementVersion !== sourceIdentity.managementVersion || currentIdentityArg.cliName !== sourceIdentity.cliName || currentIdentityArg.cliRelativePath !== sourceIdentity.cliRelativePath ) throw new Error('The inherited upgrade source identity is incompatible.'); return; } if ( !transactionArg.controllerWasRunning || transactionArg.targetStartupInvoked !== true || upgradeTransactionTargetVersion(transactionArg) !== currentIdentityArg.packageVersion ) throw new Error('The inherited upgrade transaction did not authorize target startup.'); if (transactionArg.version === upgradeCoordinationVersion) { if ( currentIdentityArg.packageName !== String(controllerPackageName) || currentIdentityArg.managementVersion !== controllerUpgradeManagementVersion ) throw new Error('The inherited v2 upgrade target identity is incompatible.'); return; } if ( transactionArg.targetPackageCommitted !== true || transactionArg.targetPackageName !== currentIdentityArg.packageName || transactionArg.targetManagementVersion !== currentIdentityArg.managementVersion || transactionArg.targetCliName !== currentIdentityArg.cliName || transactionArg.targetCliRelativePath !== currentIdentityArg.cliRelativePath ) throw new Error('The inherited v3 upgrade target identity is incompatible.'); }; const runPackageTransitionUpgradeWorker = async (optionsArg: { coordinator: UpgradeCoordinator; payload: IUpgradeWorkerPayload; transaction: IUpgradeTransactionV3; currentInstallation: IPnpmGlobalInstallation; registryUrl?: string; }): Promise => { const { coordinator, payload, transaction: initialTransaction, currentInstallation, registryUrl, } = optionsArg; if ( commitinfo.name !== initialTransaction.sourcePackageName || commitinfo.version !== initialTransaction.sourceVersion || currentInstallation.packageVersion !== initialTransaction.sourceVersion ) throw new Error('The package-transition worker is not running from the exact bridge package.'); const context: IPackageTransitionContext = { packageManagerPath: currentInstallation.packageManagerPath, globalRoot: currentInstallation.globalRoot, globalBinDirectory: currentInstallation.globalBinDirectory, }; // Registry validation is deliberately complete before controller preparation or package mutation. await resolveExactPackageTransitionTarget(context, registryUrl); await coordinator.mutateTransaction(payload.token, (current) => { if (current.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...current, phase: 'checking', message: `Resolved exact target ${current.targetPackageName}@${current.targetVersion}.`, }; }); process.stdout.write( `${initialTransaction.sourceCliName} upgrade: resolved ${initialTransaction.targetPackageName}@${initialTransaction.targetVersion}.\n`, ); const controllerProcesses = await listControllerProcessesForCli(currentInstallation.cliPath); let runningStatus: IControllerStatus | undefined; if (payload.expectedController) { runningStatus = await queryControllerStatus(payload.port, { packageName: initialTransaction.sourcePackageName, packageVersion: initialTransaction.sourceVersion, protocolVersion: controllerProtocolVersion, upgradeManagementVersion: initialTransaction.sourceManagementVersion, }); assertExpectedController(runningStatus, payload.expectedController); if ( controllerProcesses.length !== 1 || controllerProcesses[0].identity.pid !== runningStatus.controllerPid || controllerProcesses[0].port !== payload.port ) throw new Error('Additional or mismatched pnpm-global controller instances prevent upgrade.'); } else if (controllerProcesses.length > 0) { throw new Error('A pnpm-global controller started after the upgrade request; upgrade aborted.'); } else if (await isLoopbackPortListening(payload.port)) { throw new Error( `The requested controller port became occupied while ${initialTransaction.sourceCliName} upgrade was starting.`, ); } if (runningStatus) { await prepareAndStopRunningController({ coordinator, token: payload.token, port: payload.port, targetVersion: initialTransaction.targetVersion, gracePeriodMs: payload.gracePeriodMs, currentInstallation, runningStatus, }); process.stdout.write(`${initialTransaction.sourceCliName} upgrade: bridge controller stopped cleanly.\n`); } let targetInstallation: IPnpmGlobalInstallation; try { await updateWorkerProgress( coordinator, payload.token, 'installing', `Transitioning ${initialTransaction.sourcePackageName} to ${initialTransaction.targetPackageName}.`, ); if (await isLoopbackPortListening(payload.port)) { throw new Error('The controller port became occupied before package replacement.'); } const sourceDataWriters = await listControllerDataWriterProcessesForCli( currentInstallation.cliPath, ); if (sourceDataWriters.length > 0) { const writer = sourceDataWriters[0]; throw new Error( `Package replacement is blocked by ${writer.kind} source writer PID ${writer.identity.pid}.`, ); } const sourceSpec = `${initialTransaction.sourcePackageName}@${initialTransaction.sourceVersion}`; const targetSpec = `${initialTransaction.targetPackageName}@${initialTransaction.targetVersion}`; const groupedSpec = `${sourceSpec},${targetSpec}`; let current = requirePackageTransitionTransaction( await coordinator.readTransaction(payload.token), ); if (current.targetPackageCommitStarted !== true) { if (current.groupedTargetVerified !== true) { if (current.packageTransitionStarted === true) { await addExactPackage(context, sourceSpec, registryUrl); await verifyPackageTransitionTopology(context, 'source-only'); } const groupedAdd = await preparePackageAdd(context, groupedSpec, registryUrl); if (current.packageTransitionStarted !== true) { current = requirePackageTransitionTransaction(await coordinator.mutateTransaction(payload.token, (transaction) => { if (transaction.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...transaction, packageTransitionStarted: true as const, message: `Installing grouped ${sourceSpec},${targetSpec}.`, }; })); } await runPreparedManagedCommand(groupedAdd); await verifyPackageTransitionTopology(context, 'grouped'); current = requirePackageTransitionTransaction(await coordinator.mutateTransaction(payload.token, (transaction) => { if (transaction.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...transaction, groupedTargetVerified: true as const, message: 'Verified the grouped bridge and target package topology.', }; })); } else { await verifyPackageTransitionTopology(context, 'grouped'); } const targetAdd = await preparePackageAdd(context, targetSpec, registryUrl); current = requirePackageTransitionTransaction(await coordinator.mutateTransaction(payload.token, (transaction) => { if (transaction.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...transaction, targetPackageCommitStarted: true as const, message: `Committing exact target ${targetSpec}.`, }; })); await runPreparedManagedCommand(targetAdd); } else if (current.targetPackageCommitted !== true) { await addExactPackage(context, targetSpec, registryUrl); } if (current.targetPackageCommitted !== true) { const topology = await verifyPackageTransitionTopology(context, 'target-only'); if (!topology.target) throw new Error('The exact target package installation is unavailable.'); targetInstallation = topology.target; await coordinator.mutateTransaction(payload.token, (transaction) => { if (transaction.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...transaction, targetPackageCommitted: true as const, message: `Committed exact target ${targetSpec}.`, }; }); } else { const topology = await verifyPackageTransitionTopology(context, 'target-only'); if (!topology.target) throw new Error('The exact target package installation is unavailable.'); targetInstallation = topology.target; } } catch (upgradeErrorArg) { if (errorContainsCommandCleanupFailure(upgradeErrorArg)) throw upgradeErrorArg; const current = await coordinator.readTransaction(payload.token); if ( current.version === upgradePackageTransitionTransactionVersion && current.targetPackageCommitStarted === true ) throw upgradeErrorArg; const recoveryErrors: unknown[] = [upgradeErrorArg]; try { const sourceSpec = `${initialTransaction.sourcePackageName}@${initialTransaction.sourceVersion}`; if ( current.version === upgradePackageTransitionTransactionVersion && current.packageTransitionStarted === true ) await addExactPackage(context, sourceSpec, registryUrl); const sourceTopology = await verifyPackageTransitionTopology(context, 'source-only'); if (!sourceTopology.source) throw new Error('The exact bridge package was not restored.'); if (runningStatus) { await startRestoredController(sourceTopology.source, payload.port, payload.token, coordinator); const restoredStatus = await queryUpgradeManagedController( payload.port, sourceControllerIdentity(initialTransaction), ); await finalizeControllerUpgrade( coordinator, payload.token, restoredStatus, 'compensate', ); process.stdout.write( `${initialTransaction.sourceCliName} upgrade: bridge package restored and restarted.\n`, ); } } catch (recoveryErrorArg) { recoveryErrors.push(recoveryErrorArg); } throw new AggregateError( recoveryErrors, `${initialTransaction.sourceCliName} upgrade failed before target package commit began.`, ); } process.stdout.write( `${initialTransaction.sourceCliName} upgrade: committed ${initialTransaction.targetPackageName}@${initialTransaction.targetVersion}.\n`, ); if (runningStatus) { await updateWorkerProgress( coordinator, payload.token, 'restarting', `Restarting ${initialTransaction.targetPackageName} ${initialTransaction.targetVersion}.`, ); const preparedTargetStart = await prepareControllerStart( targetInstallation, payload.port, payload.token, coordinator, ); await coordinator.mutateTransaction(payload.token, (current) => { if (current.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...current, targetStartupInvoked: true as const, message: `Invoking ${current.targetPackageName} ${current.targetVersion}.`, }; }); await runPreparedManagedCommand(preparedTargetStart); const targetStatus = await queryUpgradeManagedController( payload.port, targetControllerIdentity(initialTransaction), ); await updateWorkerProgress( coordinator, payload.token, payload.continueSessions ? 'continuing' : 'restarting', payload.continueSessions ? 'Continuing sessions paused for the upgrade.' : 'Reopening harness prompt admission.', ); await finalizeControllerUpgrade( coordinator, payload.token, targetStatus, payload.continueSessions ? 'continue' : 'reopen', ); process.stdout.write( `${initialTransaction.sourceCliName} upgrade: target controller restarted on port ${payload.port}.\n`, ); } await coordinator.finishTransaction( payload.token, true, `Upgrade to ${initialTransaction.targetPackageName}@${initialTransaction.targetVersion} completed.`, ); }; interface IUpgradeWorkerInternalOptions { registryUrl?: string; recoveryContext?: string; recoveryBody?: (optionsArg: { coordinator: UpgradeCoordinator; token: string; transaction: TUpgradeTransaction; installation: IPnpmGlobalInstallation; }) => Promise; isCommandCleanupFailure?: (errorArg: unknown) => boolean; } const runRecoveryUpgradeWorker = async ( payloadArg: IUpgradeWorkerPayload, internalOptionsArg: IUpgradeWorkerInternalOptions, ): Promise => { const [context, coordinator] = (() => { try { const parsedContext = parseUpgradeRecoveryWorkerContext(internalOptionsArg.recoveryContext); return [parsedContext, new UpgradeCoordinator(parsedContext.canonicalGlobalRoot)] as const; } finally { scrubUpgradeCoordinationEnvironment(); } })(); let lock: Awaited> | undefined; let releaseLock = true; try { await coordinator.waitForWorkerAcknowledgement(payloadArg.token, 10_000); lock = await coordinator.acquireUpgradeLock({ token: payloadArg.token, cliPath: context.cliPath, }); await coordinator.waitForStartLeasesToDrain(lock); const transaction = await coordinator.readTransaction(payloadArg.token); if (transaction.terminal) return; if ( transaction.port !== payloadArg.port || transaction.gracePeriodMs !== payloadArg.gracePeriodMs || transaction.continueSessions !== payloadArg.continueSessions ) throw new Error('The recovery worker payload does not match its upgrade transaction.'); const identity = await readControllerProcessIdentity(process.pid); if ( !identity || !identity.processGroupLeader || identity.processGroupId !== process.pid || transaction.worker?.pid !== identity.pid || transaction.worker.processGroupId !== identity.processGroupId || transaction.worker.fingerprint !== identity.fingerprint || transaction.worker.cliPath !== context.cliPath ) throw new Error('The recovery worker transaction identity is incompatible.'); const installation = await resolveAuthoritativeRecoveryInstallation(transaction, { packageManagerPath: context.packageManagerPath, globalRoot: context.canonicalGlobalRoot, globalBinDirectory: context.globalBinDirectory, packageDirectory: plugins.path.dirname(context.cliPath), cliPath: context.cliPath, packageVersion: '', }); if (installation.cliPath !== context.cliPath) { throw new Error('The authoritative recovery CLI changed during handoff.'); } upgradeWorkerCliPath = installation.cliPath; await (internalOptionsArg.recoveryBody ?? recoverStalledUpgradeUnderLock)({ coordinator, token: payloadArg.token, transaction, installation, }); } catch (errorArg) { const cleanupFailed = (internalOptionsArg.isCommandCleanupFailure ?? errorContainsCommandCleanupFailure)(errorArg); if (cleanupFailed) releaseLock = false; await coordinator.mutateTransaction(payloadArg.token, (current) => ({ ...current, message: errorArg instanceof Error ? errorArg.message : String(errorArg), })).catch(() => undefined); throw errorArg; } finally { if (releaseLock) await lock?.release(); } }; export function runUpgradeWorker(payloadArg: IUpgradeWorkerPayload): Promise; export async function runUpgradeWorker( payloadArg: IUpgradeWorkerPayload, internalOptionsArg: IUpgradeWorkerInternalOptions = {}, ): Promise { const environmentToken = process.env[upgradeTokenEnvironmentVariable]; if (environmentToken !== payloadArg.token) { scrubUpgradeCoordinationEnvironment(); throw new Error(`The detached ${currentCliName} upgrade worker token binding is invalid.`); } if (internalOptionsArg.recoveryContext !== undefined) { try { await runRecoveryUpgradeWorker(payloadArg, internalOptionsArg); } finally { scrubUpgradeCoordinationEnvironment(); } return; } let installation: IPnpmGlobalInstallation; let coordinator: UpgradeCoordinator; try { installation = await resolveCurrentPnpmGlobalInstallation(); upgradeWorkerCliPath = installation.cliPath; coordinator = new UpgradeCoordinator(installation.globalRoot); } finally { scrubUpgradeCoordinationEnvironment(); } let lock: Awaited> | undefined; let releaseLock = true; try { await coordinator.registerWorker(payloadArg.token, installation.cliPath); await coordinator.waitForWorkerAcknowledgement(payloadArg.token, 10_000); lock = await coordinator.acquireUpgradeLock({ token: payloadArg.token, cliPath: installation.cliPath, }); await coordinator.waitForStartLeasesToDrain(lock); const currentInstallation = await resolveCurrentPnpmGlobalInstallation(); if (currentInstallation.globalRoot !== installation.globalRoot) { throw new Error(`The pnpm global root changed while ${currentCliName} upgrade was starting.`); } const transaction = await coordinator.readTransaction(payloadArg.token); if (transaction.version === upgradePackageTransitionTransactionVersion) { await runPackageTransitionUpgradeWorker({ coordinator, payload: payloadArg, transaction, currentInstallation, registryUrl: internalOptionsArg.registryUrl === undefined ? undefined : normalizeUpgradeRegistryUrl(internalOptionsArg.registryUrl), }); return; } const latestVersion = await readLatestVersion(currentInstallation, transaction.registryUrl); await coordinator.mutateTransaction(payloadArg.token, (current) => { if (current.version !== upgradeCoordinationVersion) { throw new Error('The same-package upgrade transaction changed format.'); } return { ...current, targetVersion: latestVersion, phase: 'checking', message: `Installed ${currentInstallation.packageVersion}; registry latest is ${latestVersion}.`, }; }); const versionComparison = compareSemver(latestVersion, currentInstallation.packageVersion); process.stdout.write( `${currentCliName} upgrade: installed ${currentInstallation.packageVersion}, registry latest ${latestVersion}.\n`, ); if (versionComparison <= 0) { process.stdout.write( versionComparison === 0 ? `${currentCliName} upgrade: already up to date.\n` : `${currentCliName} upgrade: installed version is newer than registry latest; refusing to downgrade.\n`, ); await coordinator.finishTransaction( payloadArg.token, true, versionComparison === 0 ? 'Already up to date.' : 'Installed version is newer than registry latest; no downgrade was performed.', ); return; } const controllerProcesses = await listControllerProcessesForCli(currentInstallation.cliPath); let runningStatus: IControllerStatus | undefined; if (payloadArg.expectedController) { runningStatus = await queryControllerStatus(payloadArg.port, { packageName: controllerPackageName, packageVersion: currentInstallation.packageVersion, protocolVersion: controllerProtocolVersion, upgradeManagementVersion: controllerUpgradeManagementVersion, }); assertExpectedController(runningStatus, payloadArg.expectedController); if ( controllerProcesses.length !== 1 || controllerProcesses[0].identity.pid !== runningStatus.controllerPid || controllerProcesses[0].port !== payloadArg.port ) { throw new Error('Additional or mismatched pnpm-global controller instances prevent upgrade.'); } } else if (controllerProcesses.length > 0) { throw new Error('A pnpm-global controller started after the upgrade request; upgrade aborted.'); } else if (await isLoopbackPortListening(payloadArg.port)) { throw new Error( `The requested controller port became occupied while ${currentCliName} upgrade was starting.`, ); } if (runningStatus) { await prepareAndStopRunningController({ coordinator, token: payloadArg.token, port: payloadArg.port, targetVersion: latestVersion, gracePeriodMs: payloadArg.gracePeriodMs, currentInstallation, runningStatus, }); process.stdout.write(`${currentCliName} upgrade: previous controller stopped cleanly.\n`); } let upgradedInstallation: IPnpmGlobalInstallation; try { await updateWorkerProgress( coordinator, payloadArg.token, 'installing', `Installing ${controllerPackageName} ${latestVersion}.`, ); if (await isLoopbackPortListening(payloadArg.port)) { throw new Error('The controller port became occupied before package replacement.'); } const sourceDataWriters = await listControllerDataWriterProcessesForCli( currentInstallation.cliPath, ); if (sourceDataWriters.length > 0) { const writer = sourceDataWriters[0]; throw new Error( `Package replacement is blocked by ${writer.kind} source writer PID ${writer.identity.pid}.`, ); } await installVersion(currentInstallation, latestVersion, transaction.registryUrl); upgradedInstallation = await resolvePnpmGlobalInstallation(); if (upgradedInstallation.packageVersion !== latestVersion) { throw new Error('pnpm completed without installing the requested registry version.'); } } catch (upgradeErrorArg) { if (errorContainsCommandCleanupFailure(upgradeErrorArg)) throw upgradeErrorArg; const recoveryErrors: unknown[] = [upgradeErrorArg]; try { const restored = await restorePreviousInstallation( currentInstallation, transaction.registryUrl, ); if (runningStatus) { await startRestoredController(restored, payloadArg.port, payloadArg.token, coordinator); const restoredStatus = await queryUpgradeManagedController( payloadArg.port, { packageName: controllerPackageName, packageVersion: restored.packageVersion, managementVersion: controllerUpgradeManagementVersion, }, ); await finalizeControllerUpgrade( coordinator, payloadArg.token, restoredStatus, 'compensate', ); process.stdout.write( `${currentCliName} upgrade: previous controller version restored and restarted.\n`, ); } } catch (recoveryErrorArg) { recoveryErrors.push(recoveryErrorArg); } throw new AggregateError( recoveryErrors, `${currentCliName} upgrade failed before the new controller started.`, ); } process.stdout.write(`${currentCliName} upgrade: installed ${latestVersion}.\n`); if (runningStatus) { // From this point onward rollback is unsafe: startup may run data migrations. await updateWorkerProgress( coordinator, payloadArg.token, 'restarting', `Restarting ${controllerPackageName} ${latestVersion}.`, ); await coordinator.mutateTransaction(payloadArg.token, (current) => ({ ...current, targetStartupInvoked: true, message: `Invoking ${controllerPackageName} ${latestVersion}.`, })); await startController(upgradedInstallation, payloadArg.port, payloadArg.token, coordinator); const upgradedStatus = await queryUpgradeManagedController(payloadArg.port, { packageName: controllerPackageName, packageVersion: latestVersion, managementVersion: controllerUpgradeManagementVersion, }); await updateWorkerProgress( coordinator, payloadArg.token, payloadArg.continueSessions ? 'continuing' : 'restarting', payloadArg.continueSessions ? 'Continuing sessions paused for the upgrade.' : 'Reopening harness prompt admission.', ); await finalizeControllerUpgrade( coordinator, payloadArg.token, upgradedStatus, payloadArg.continueSessions ? 'continue' : 'reopen', ); process.stdout.write( `${currentCliName} upgrade: controller restarted on port ${payloadArg.port}.\n`, ); } await coordinator.finishTransaction( payloadArg.token, true, `Upgrade to ${latestVersion} completed.`, ); } catch (errorArg) { if (errorContainsCommandCleanupFailure(errorArg)) releaseLock = false; const failedTransaction = await coordinator.readTransaction(payloadArg.token).catch(() => undefined); if ( (failedTransaction?.version === upgradePackageTransitionTransactionVersion ? failedTransaction.targetPackageCommitStarted === true : failedTransaction?.targetStartupInvoked === true) || errorContainsCommandCleanupFailure(errorArg) ) { await coordinator.mutateTransaction(payloadArg.token, (current) => ({ ...current, message: errorArg instanceof Error ? errorArg.message : String(errorArg), })).catch(() => undefined); throw errorArg; } await coordinator.finishTransaction( payloadArg.token, false, 'Upgrade failed.', errorArg instanceof Error ? errorArg.message : String(errorArg), ).catch(() => undefined); throw errorArg; } finally { if (releaseLock) await lock?.release(); } } const recoverPackageTransition = async (optionsArg: { coordinator: UpgradeCoordinator; token: string; transaction: IUpgradeTransactionV3; installation: IPnpmGlobalInstallation; }): Promise => { const context: IPackageTransitionContext = { packageManagerPath: optionsArg.installation.packageManagerPath, globalRoot: optionsArg.installation.globalRoot, globalBinDirectory: optionsArg.installation.globalBinDirectory, }; const registryUrl = undefined; const transaction = optionsArg.transaction; if (transaction.targetPackageCommitStarted === true) { const targetVersion = upgradeTransactionTargetVersion(transaction)!; const targetSpec = `${transaction.targetPackageName}@${targetVersion}`; let targetTopology: IPackageTransitionTopology; if (transaction.targetPackageCommitted === true) { try { targetTopology = await verifyPackageTransitionTopology( context, 'target-only', targetVersion, ); } catch (verificationErrorArg) { try { await addExactPackage(context, targetSpec, registryUrl); targetTopology = await verifyPackageTransitionTopology( context, 'target-only', targetVersion, ); } catch (normalizationErrorArg) { throw new AggregateError( [verificationErrorArg, normalizationErrorArg], 'The committed target package could not be normalized without restoring the bridge.', ); } } } else { await addExactPackage(context, targetSpec, registryUrl); targetTopology = await verifyPackageTransitionTopology( context, 'target-only', targetVersion, ); await optionsArg.coordinator.mutateTransaction(optionsArg.token, (current) => { if (current.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...current, targetPackageCommitted: true as const, message: `Committed exact target ${targetSpec} during bounded recovery.`, }; }); } const targetInstallation = targetTopology.target; if (!targetInstallation) throw new Error('The exact target package is unavailable.'); let current = requirePackageTransitionTransaction( await optionsArg.coordinator.readTransaction(optionsArg.token), ); if (current.controllerWasRunning) { if (current.targetStartupInvoked !== true) { const candidates = (await listControllerProcessesForCli(targetInstallation.cliPath)) .filter((candidateArg) => candidateArg.port === current.port); if (candidates.length > 0 || await isLoopbackPortListening(current.port)) { throw new Error( 'A target controller exists before its durable startup marker; refusing to adopt it.', ); } const preparedTargetStart = await prepareControllerStart( targetInstallation, current.port, optionsArg.token, optionsArg.coordinator, ); current = requirePackageTransitionTransaction(await optionsArg.coordinator.mutateTransaction(optionsArg.token, (candidate) => { if (candidate.version !== upgradePackageTransitionTransactionVersion) { throw new Error('The package-transition transaction changed format.'); } return { ...candidate, phase: 'restarting', targetStartupInvoked: true as const, message: `Invoking ${candidate.targetPackageName} during bounded recovery.`, }; })); await runPreparedManagedCommand(preparedTargetStart); } let status = await queryControllerStatus(current.port, { packageName: current.targetPackageName, packageVersion: upgradeTransactionTargetVersion(current)!, upgradeManagementVersion: current.targetManagementVersion, }).catch(() => undefined); if (!status) { if (await retainedUpgradeTargetIsAuthoritative({ coordinator: optionsArg.coordinator, token: optionsArg.token, })) return false; const targetCandidates = (await listControllerProcessesForCli(targetInstallation.cliPath)) .filter((candidateArg) => candidateArg.port === current.port); if (targetCandidates.length > 0 || await isLoopbackPortListening(current.port)) { throw new Error( 'A target controller is already starting or listening but cannot be adopted safely; refusing to launch a replacement.', ); } await startController( targetInstallation, current.port, optionsArg.token, optionsArg.coordinator, ); status = await queryUpgradeManagedController( current.port, targetControllerIdentity(current), ); } await finalizeControllerUpgrade( optionsArg.coordinator, optionsArg.token, status, current.continueSessions ? 'continue' : 'reopen', ); } await optionsArg.coordinator.finishTransaction( optionsArg.token, true, `Upgrade to ${current.targetPackageName}@${upgradeTransactionTargetVersion(current)} completed during bounded recovery.`, ); return true; } const sourceSpec = `${transaction.sourcePackageName}@${transaction.sourceVersion}`; let sourceTopology: IPackageTransitionTopology; if (transaction.packageTransitionStarted === true) { await addExactPackage(context, sourceSpec, registryUrl); sourceTopology = await verifyPackageTransitionTopology(context, 'source-only'); } else { sourceTopology = await verifyPackageTransitionTopology(context, 'source-only'); } const sourceInstallation = sourceTopology.source; if (!sourceInstallation) throw new Error('The exact bridge package is unavailable.'); if (transaction.controllerWasRunning) { let status = await queryControllerStatus(transaction.port, { packageName: transaction.sourcePackageName, packageVersion: transaction.sourceVersion, upgradeManagementVersion: transaction.sourceManagementVersion, }).catch(() => undefined); if (!status) { await startRestoredController( sourceInstallation, transaction.port, optionsArg.token, optionsArg.coordinator, ); status = await queryUpgradeManagedController( transaction.port, sourceControllerIdentity(transaction), ); } await finalizeControllerUpgrade( optionsArg.coordinator, optionsArg.token, status, 'compensate', ); } await optionsArg.coordinator.finishTransaction( optionsArg.token, false, 'The stalled package transition restored the exact bridge package.', `The upgrade stopped making progress during phase ${transaction.phase}.`, ); return true; }; const recoverStalledUpgradeUnderLock = async (optionsArg: { coordinator: UpgradeCoordinator; token: string; transaction: TUpgradeTransaction; installation: IPnpmGlobalInstallation; }): Promise => { const transaction = optionsArg.transaction; if (transaction.terminal) return true; if (transaction.targetStartupInvoked === true && await retainedUpgradeTargetIsAuthoritative({ coordinator: optionsArg.coordinator, token: optionsArg.token, }, { resolveInstallation: async () => optionsArg.installation, })) return false; if (transaction.version === upgradePackageTransitionTransactionVersion) { return await recoverPackageTransition({ coordinator: optionsArg.coordinator, token: optionsArg.token, transaction, installation: optionsArg.installation, }); } const afterTargetStartupBoundary = transaction.targetStartupInvoked === true; if (afterTargetStartupBoundary) { if (!transaction.targetVersion) { throw new Error('The stalled upgrade has no retained target version.'); } const targetInstallation = await resolvePnpmGlobalInstallation(); if (targetInstallation.packageVersion !== transaction.targetVersion) { throw new Error( 'The package installation does not match the retained post-startup upgrade target.', ); } if (transaction.controllerWasRunning) { let status = await queryControllerStatus(transaction.port, { packageName: controllerPackageName, packageVersion: targetInstallation.packageVersion, upgradeManagementVersion: controllerUpgradeManagementVersion, }).catch(() => undefined); if (!status) { if (await retainedUpgradeTargetIsAuthoritative({ coordinator: optionsArg.coordinator, token: optionsArg.token, })) return false; const targetCandidates = (await listControllerProcessesForCli(targetInstallation.cliPath)) .filter((candidateArg) => candidateArg.port === transaction.port); if (targetCandidates.length > 0 || await isLoopbackPortListening(transaction.port)) { throw new Error( 'A target controller is already starting or listening but cannot be adopted safely; refusing to launch a replacement.', ); } await startController( targetInstallation, transaction.port, optionsArg.token, optionsArg.coordinator, ); status = await queryUpgradeManagedController( transaction.port, { packageName: controllerPackageName, packageVersion: targetInstallation.packageVersion, managementVersion: controllerUpgradeManagementVersion, }, ); } await finalizeControllerUpgrade( optionsArg.coordinator, optionsArg.token, status, transaction.continueSessions ? 'continue' : 'reopen', ); } await optionsArg.coordinator.finishTransaction( optionsArg.token, true, `Upgrade to ${transaction.targetVersion} completed during bounded recovery.`, ); return true; } const sourceInstallation = await restorePreviousInstallation({ ...optionsArg.installation, packageVersion: transaction.sourceVersion, }, transaction.version === upgradeCoordinationVersion ? transaction.registryUrl : undefined); if (transaction.controllerWasRunning) { let status = await queryControllerStatus(transaction.port, { packageName: controllerPackageName, packageVersion: sourceInstallation.packageVersion, upgradeManagementVersion: controllerUpgradeManagementVersion, }).catch(() => undefined); if (!status) { await startRestoredController( sourceInstallation, transaction.port, optionsArg.token, optionsArg.coordinator, ); status = await queryUpgradeManagedController( transaction.port, { packageName: controllerPackageName, packageVersion: sourceInstallation.packageVersion, managementVersion: controllerUpgradeManagementVersion, }, ); } await finalizeControllerUpgrade( optionsArg.coordinator, optionsArg.token, status, 'compensate', ); } await optionsArg.coordinator.finishTransaction( optionsArg.token, false, 'The stalled upgrade restored the source package.', `The upgrade stopped making progress during phase ${transaction.phase}.`, ); return true; }; interface IRecoverStalledUpgradeOptions { coordinator: UpgradeCoordinator; token: string; transaction: TUpgradeTransaction; installation: IPnpmGlobalInstallation; command: '__serve' | 'foreground' | 'upgrade'; } export interface IUpgradeWorkerHandoffDependencies { terminateWorkerCandidate?: ( candidateArg: IDetachedUpgradeWorker, ) => Promise; candidateCleanup?: ICandidateCleanupDependencies; } interface IRecoverStalledUpgradeDependencies extends IUpgradeWorkerHandoffDependencies { launchRecoveryWorker?: typeof launchDetachedRecoveryWorker; } export interface IHandoffLaunchedUpgradeWorkerUnderLockOptions { coordinator: UpgradeCoordinator; token: string; installation: IPnpmGlobalInstallation; command: '__serve' | 'foreground' | 'upgrade'; lock: UpgradeInstallationLock; candidate: IUpgradeWorkerLaunchCandidate; retainLockOnFailure?: boolean; } export const handoffLaunchedUpgradeWorkerUnderLock = async ( optionsArg: IHandoffLaunchedUpgradeWorkerUnderLockOptions, dependenciesArg: IUpgradeWorkerHandoffDependencies = {}, ): Promise => { let worker: IDetachedUpgradeWorker | undefined; try { await optionsArg.coordinator.removeWorkerAcknowledgement(optionsArg.token); const launchedTransaction = await optionsArg.coordinator.recordWorkerLaunch( optionsArg.token, optionsArg.candidate.pid, optionsArg.candidate.cliPath, optionsArg.candidate.logFilePath, ); const retainedWorker = launchedTransaction.worker; if ( !retainedWorker || retainedWorker.pid !== optionsArg.candidate.pid || retainedWorker.cliPath !== optionsArg.candidate.cliPath || retainedWorker.logFilePath !== optionsArg.candidate.logFilePath || (optionsArg.candidate.processGroupId !== undefined && retainedWorker.processGroupId !== optionsArg.candidate.processGroupId) || (optionsArg.candidate.fingerprint !== undefined && retainedWorker.fingerprint !== optionsArg.candidate.fingerprint) ) throw new Error('The upgrade worker launch identity was not retained exactly.'); worker = { pid: retainedWorker.pid, processGroupId: retainedWorker.processGroupId, fingerprint: retainedWorker.fingerprint, cliPath: retainedWorker.cliPath, logFilePath: optionsArg.candidate.logFilePath, }; await optionsArg.coordinator.transferUpgradeLockToWorker({ lock: optionsArg.lock, token: optionsArg.token, worker, }); await optionsArg.coordinator.acknowledgeWorker(optionsArg.token, worker); return { pid: worker.pid, processGroupId: worker.processGroupId, fingerprint: worker.fingerprint, cliPath: worker.cliPath, logFilePath: worker.logFilePath!, }; } catch (errorArg) { const errors: unknown[] = [errorArg]; if (!worker) { const current = await optionsArg.coordinator.readTransaction(optionsArg.token).catch(() => undefined); if ( current?.worker && current.worker.pid === optionsArg.candidate.pid && current.worker.cliPath === optionsArg.candidate.cliPath && current.worker.logFilePath === optionsArg.candidate.logFilePath ) worker = { pid: current.worker.pid, processGroupId: current.worker.processGroupId, fingerprint: current.worker.fingerprint, cliPath: current.worker.cliPath, logFilePath: optionsArg.candidate.logFilePath, }; } if (!worker && ( optionsArg.candidate.processGroupId !== undefined && optionsArg.candidate.fingerprint !== undefined )) { worker = { pid: optionsArg.candidate.pid, processGroupId: optionsArg.candidate.processGroupId, fingerprint: optionsArg.candidate.fingerprint, cliPath: optionsArg.candidate.cliPath, logFilePath: optionsArg.candidate.logFilePath, }; } if (!worker) { throw new UpgradeWorkerCandidateDrainageError( 'The upgrade worker handoff failed before its cleanup identity could be retained.', errorArg, ); } const cleanupWorker = worker; try { await cleanupCandidateUntilDrained( async () => await (dependenciesArg.terminateWorkerCandidate ?? ((candidateArg) => optionsArg.coordinator.terminateUpgradeWorkerCandidate(candidateArg)))( cleanupWorker, ), dependenciesArg.candidateCleanup, cleanupWorker.pid, ); } catch (cleanupErrorArg) { errors.push(cleanupErrorArg); } if (errors.length > 1) { throw new AggregateError( errors, 'Upgrade worker handoff failed and candidate drainage is unproven; ownership was retained.', ); } const acquireCallerLock = async (): Promise => { try { await optionsArg.lock.assertOwned(); return optionsArg.lock; } catch { return await optionsArg.coordinator.acquireUpgradeLock({ token: optionsArg.token, cliPath: optionsArg.installation.cliPath, command: optionsArg.command, }); } }; if (optionsArg.retainLockOnFailure) { try { const recoveryLock = await acquireCallerLock(); await recoveryLock.assertOwned(); throw new UpgradeWorkerHandoffError( 'The upgrade worker handoff failed; exclusive ownership was retained for recovery.', errors.length === 1 ? errorArg : new AggregateError(errors, 'Upgrade worker handoff and candidate cleanup failed.'), recoveryLock, ); } catch (lockRetentionErrorArg) { if (lockRetentionErrorArg instanceof UpgradeWorkerHandoffError) { throw lockRetentionErrorArg; } errors.push(lockRetentionErrorArg); } } try { await (await acquireCallerLock()).release(); } catch (lockCleanupErrorArg) { errors.push(lockCleanupErrorArg); } if (errors.length > 1) { throw new AggregateError( errors, 'Upgrade worker handoff failed and cleanup was incomplete.', ); } throw errorArg; } }; export interface IHandoffStalledUpgradeRecoveryUnderLockOptions { coordinator: UpgradeCoordinator; token: string; installation: IPnpmGlobalInstallation; command: '__serve' | 'foreground' | 'upgrade'; lock: UpgradeInstallationLock; } export const handoffStalledUpgradeRecoveryUnderLock = async ( optionsArg: IHandoffStalledUpgradeRecoveryUnderLockOptions, dependenciesArg: IRecoverStalledUpgradeDependencies = {}, ): Promise => { let candidate: IDetachedUpgradeWorker; try { const transaction = await optionsArg.coordinator.readTransaction(optionsArg.token); if (transaction.terminal) { await optionsArg.lock.release(); return undefined; } const recoveryInstallation = await resolveAuthoritativeRecoveryInstallation( transaction, optionsArg.installation, ); const payload: IUpgradeWorkerPayload = { version: upgradeCoordinationVersion, token: optionsArg.token, port: transaction.port, gracePeriodMs: transaction.gracePeriodMs, continueSessions: transaction.continueSessions, }; await optionsArg.coordinator.removeWorkerAcknowledgement(optionsArg.token); candidate = await (dependenciesArg.launchRecoveryWorker ?? launchDetachedRecoveryWorker)({ installation: recoveryInstallation, payload, }); } catch (errorArg) { if (errorArg instanceof UpgradeWorkerCandidateDrainageError) throw errorArg; try { await optionsArg.lock.release(); } catch (lockCleanupErrorArg) { throw new AggregateError( [errorArg, lockCleanupErrorArg], 'Stalled upgrade recovery launch failed and its lock could not be released.', ); } throw errorArg; } return await handoffLaunchedUpgradeWorkerUnderLock({ ...optionsArg, candidate, }, dependenciesArg); }; export function recoverStalledUpgrade( optionsArg: IRecoverStalledUpgradeOptions, ): Promise; export async function recoverStalledUpgrade( optionsArg: IRecoverStalledUpgradeOptions, dependenciesArg: IRecoverStalledUpgradeDependencies = {}, ): Promise { await optionsArg.coordinator.terminateTransactionWorker(optionsArg.transaction); let lock: Awaited>; try { lock = await optionsArg.coordinator.acquireUpgradeLock({ token: optionsArg.token, cliPath: optionsArg.installation.cliPath, command: optionsArg.command, }); } catch (errorArg) { if ( errorArg instanceof Error && (errorArg.message.includes('already running') || errorArg.message.includes('initializing')) ) return false; throw errorArg; } await handoffStalledUpgradeRecoveryUnderLock({ coordinator: optionsArg.coordinator, token: optionsArg.token, installation: optionsArg.installation, command: optionsArg.command, lock, }, dependenciesArg); return true; } export const createUpgradeWorkerPayload = (optionsArg: { port: number; gracePeriodMs: number; continueSessions: boolean; expectedController?: IUpgradeExpectedController; }): IUpgradeWorkerPayload => ({ version: upgradeCoordinationVersion, token: createUpgradeToken(), port: optionsArg.port, gracePeriodMs: optionsArg.gracePeriodMs, continueSessions: optionsArg.continueSessions, ...(optionsArg.expectedController ? { expectedController: optionsArg.expectedController } : {}), });