feat: add npm update checks and manual installation

Add profile-scoped npm updates for Desktop and Web, protect source installations, and keep restarts manual. Cover stale status recovery, request races, and dialog focus with regression tests.

Verified with npm run check (1626 tests) and isolated Desktop/Web npm installation flows.

Refs #61
This commit is contained in:
xmanrui 2026-08-28 06:37:40 +08:00
parent 5888641297
commit f3cefd8a84
23 changed files with 4058 additions and 271 deletions

View file

@ -10,6 +10,7 @@ import { apply as applyWeixin } from './channels/weixin/index.mjs';
import { apply as applyWhatsapp } from './channels/whatsapp/index.mjs';
import { installOutboundArtifactTool } from '../../src/channels/shared/semantic/artifact.mjs';
import { setImHostLanguage } from '../../src/channels/shared/i18n.mjs';
import { installUpdateRpc } from './update-rpc.mjs';
export const name = 'dsh-im-host';
export const inject = [
@ -27,6 +28,7 @@ function channelConfig(config, name) {
}
export function createImHostPlugin(internals = {}) {
const startUpdate = internals.installUpdateRpc ?? installUpdateRpc;
const startFeishu = internals.applyFeishu ?? applyFeishu;
const startWeixin = internals.applyWeixin ?? applyWeixin;
const startDingtalk = internals.applyDingtalk ?? applyDingtalk;
@ -64,6 +66,13 @@ export function createImHostPlugin(internals = {}) {
const logger = typeof ctx?.logger === 'function'
? ctx.logger(name)
: (ctx?.logger ?? console);
if (ctx?.connection?.rpc) {
try {
startUpdate(ctx);
} catch (error) {
logger.error?.('[dsh-im] failed to activate update management; continuing with channels', error);
}
}
const failures = [];
for (const [channel, start] of channels) {
try {

View file

@ -0,0 +1,50 @@
import { createUpdateRuntime } from './update-runtime.mjs';
import { createUpdateService } from './update-service.mjs';
export const UPDATE_RPC_CHANNEL = '/dsh-im';
export const UPDATE_ENDPOINTS = Object.freeze(['update.status', 'update.check', 'update.install']);
function validPayload(endpoint, payload) {
if (payload === null || typeof payload !== 'object' || Array.isArray(payload)) return false;
const keys = Object.keys(payload);
if (endpoint !== 'update.install') return keys.length === 0;
return keys.length === 2 && keys.every((key) => key === 'checkId' || key === 'requestId')
&& ['checkId', 'requestId'].every((key) => typeof payload[key] === 'string'
&& /^[A-Za-z0-9_-]{1,128}$/.test(payload[key]));
}
const PUBLIC_ERRORS = new Set([
'check-failed', 'invalid-release', 'check-expired', 'installation-changed', 'update-busy',
'install-failed', 'verify-failed', 'state-unavailable', 'interrupted', 'disposed',
'source-install', 'unknown-profile', 'unsupported-runtime', 'registry-conflict',
'incompatible-node', 'pending-restart', 'recovery-required', 'executor-unavailable',
'invalid-installation', 'registry-check-failed', 'install-timeout', 'install-interrupted', 'invalid-version',
]);
export function createUpdateRpcHandler(service) {
return async (endpoint, payload, signal) => {
if (!UPDATE_ENDPOINTS.includes(endpoint) || !validPayload(endpoint, payload)) {
return { ok: false, error: { code: 'bad-request', message: 'Invalid update request.' } };
}
if (signal?.aborted) return { ok: false, error: { code: 'cancelled', message: 'Request cancelled.' } };
try {
// The submitted install belongs to the Host, not the lifetime of this browser request.
const value = endpoint === 'update.install' ? await service.install(payload)
: endpoint === 'update.check' ? await service.check() : await service.status();
return { ok: true, value };
} catch (error) {
const code = PUBLIC_ERRORS.has(error?.code) ? error.code : 'update-failed';
return { ok: false, error: { code, message: code } };
}
};
}
export function installUpdateRpc(ctx, options = {}) {
const runtime = options.runtime ?? createUpdateRuntime({ ctx, moduleUrl: import.meta.url });
const service = options.service ?? createUpdateService({ runtime });
const dispose = ctx.connection.rpc.handle(UPDATE_RPC_CHANNEL, createUpdateRpcHandler(service), {
authority: 'loopback',
});
ctx.effect(() => () => service.close(), 'dsh-im: close update installer');
return dispose;
}

View file

@ -0,0 +1,338 @@
import { createHash } from 'node:crypto';
import { readFile, realpath, stat } from 'node:fs/promises';
import { homedir } from 'node:os';
import { basename, dirname, isAbsolute, join, relative, resolve, sep } from 'node:path';
import { fileURLToPath } from 'node:url';
import semver from 'semver';
export const PACKAGE_NAME = '@xmanrui/dsh-im';
export const NPM_REGISTRY = 'https://registry.npmjs.org/';
const INSTALL_TIMEOUT_MS = 15 * 60_000;
const CONFIG_TIMEOUT_MS = 10_000;
const OUTPUT_LIMIT = 16 * 1024;
function failure(code, extra = {}) {
return Object.assign(new Error(code), { code, ...extra });
}
// Cordis 4's get() is the supported optional-service lookup. Putting these
// Desktop-only services in the plugin's inject array would disable web Hosts.
function service(ctx, name) {
return typeof ctx?.get === 'function' ? ctx.get(name) : ctx?.[name];
}
function inside(directory, filename) {
const suffix = relative(directory, filename);
return suffix !== '..' && !suffix.startsWith(`..${sep}`) && !isAbsolute(suffix);
}
async function readOptional(filename) {
try {
return await readFile(filename, 'utf8');
} catch (error) {
if (error.code === 'ENOENT') return '';
throw error;
}
}
async function packageAt(directory) {
const contents = await readFile(join(directory, 'package.json'), 'utf8');
return { directory: await realpath(directory), manifest: JSON.parse(contents), contents };
}
async function containingPackage(filename, name) {
let directory = dirname(await realpath(filename));
while (true) {
try {
const found = await packageAt(directory);
return found.manifest?.name === name ? found : null;
} catch (error) {
if (error.code !== 'ENOENT') throw error;
}
const parent = dirname(directory);
if (parent === directory) return null;
directory = parent;
}
}
function profileNameValid(name) {
return typeof name === 'string' && name.length > 0 && Buffer.byteLength(name) <= 255
&& !name.startsWith('-') && !['.', '..', 'node_modules'].includes(name)
&& !/[\\/\x00-\x1f\x7f<>:"|?*]/u.test(name);
}
function cliProfile(args) {
if (args[0] === 'web') return 'web';
let name;
for (let index = 0; index < args.length; index += 1) {
const token = args[index];
if (token === '--profile') name = args[++index];
else if (token.startsWith('--profile=')) name = token.slice('--profile='.length);
else if (token === '--patch') index += 1;
else if (!token.startsWith('--patch=')) break;
}
return name;
}
function dshHome(env, osHome) {
let selected = env.DSH_HOME?.trim() ? env.DSH_HOME : join(osHome, '.dsh');
if (selected === '~') selected = osHome;
else if (/^~[\\/]/u.test(selected)) selected = join(osHome, selected.slice(2));
return resolve(selected);
}
function registrySpec(spec) {
return typeof spec === 'string' && spec.trim().length > 0
&& (semver.validRange(spec) !== null || /^[A-Za-z][A-Za-z0-9._-]*$/u.test(spec));
}
async function validPackage(pkg) {
if (pkg.manifest?.name !== PACKAGE_NAME || semver.valid(pkg.manifest.version) === null) return false;
const entries = [
pkg.manifest.main,
pkg.manifest.exports?.['./client'],
pkg.manifest.dsh?.bundle?.patch,
];
for (const entry of entries) {
if (typeof entry !== 'string' || !entry || isAbsolute(entry) || entry.includes('\0')) return false;
const filename = resolve(pkg.directory, entry);
if (!inside(pkg.directory, filename) || !inside(pkg.directory, await realpath(filename))) return false;
if (!(await stat(filename)).isFile()) return false;
}
return true;
}
function officialRegistry(value) {
if (value === undefined || value === null || value === '') return true;
if (typeof value !== 'string') return false;
try {
const url = new URL(value);
return url.href === NPM_REGISTRY;
} catch {
return false;
}
}
/** Drain bounded output; raw subprocess diagnostics never cross the RPC boundary. */
async function run(start, { signal, timeoutMs, errorCode, capture = false }) {
if (signal?.aborted) throw failure('install-interrupted');
const controller = new AbortController();
let interrupted;
let operation;
const cancel = (code) => {
interrupted ??= code;
controller.abort();
operation?.cancel?.();
};
const onAbort = () => cancel('install-interrupted');
signal?.addEventListener('abort', onAbort, { once: true });
const timer = setTimeout(() => cancel('install-timeout'), timeoutMs);
let stdout = '';
let stderr = '';
try {
operation = start(controller.signal);
// The Host services own process-tree termination and wait for actual exit.
// Do not race their completion: a timed-out installer may still hold files.
operation.stdout?.on('data', (chunk) => {
if (capture) stdout = (stdout + chunk.toString()).slice(-OUTPUT_LIMIT);
});
operation.stderr?.on('data', (chunk) => {
stderr = (stderr + chunk.toString()).slice(-OUTPUT_LIMIT);
});
operation.stdout?.on('error', () => cancel(errorCode));
operation.stderr?.on('error', () => cancel(errorCode));
const outcome = await operation.done;
if (interrupted) throw failure(interrupted);
if (outcome.exitCode !== 0 || outcome.signal) {
throw failure(errorCode, {
exitCode: outcome.exitCode,
diagnosticCode: stderr.match(/\bERR_PNPM_[A-Z0-9_]+\b/u)?.[0],
});
}
return { exitCode: 0, signal: null, ...(capture ? { stdout } : {}) };
} catch (error) {
if (interrupted) throw failure(interrupted);
if (error.code === errorCode) throw error;
throw failure(errorCode);
} finally {
clearTimeout(timer);
signal?.removeEventListener('abort', onAbort);
}
}
/**
* Adapt the running Host's existing package-management capabilities. Options
* replace process facts in tests only; no runtime path comes from RPC input.
*/
export function createUpdateRuntime(options = {}) {
const ctx = options.ctx;
const env = options.env ?? process.env;
const argv = options.argv ?? process.argv;
const execArgv = options.execArgv ?? process.execArgv;
const execPath = options.execPath ?? process.execPath;
const cwd = options.cwd ?? process.cwd();
const platform = options.platform ?? process.platform;
const osHome = options.osHome ?? homedir();
const electron = options.electron ?? process.versions.electron;
const moduleUrl = options.moduleUrl ?? import.meta.url;
const loadedPackage = containingPackage(fileURLToPath(moduleUrl), PACKAGE_NAME).catch(() => null);
let boundProfile;
async function environment() {
const profiles = service(ctx, 'desktopProfiles');
const desktop = service(ctx, 'desktopPnpm');
const bootstrap = service(ctx, 'desktopPnpmBootstrap');
const isDesktop = Boolean(electron || profiles || desktop || bootstrap);
const current = profiles?.current;
const profileName = isDesktop ? current?.name : cliProfile(argv.slice(2));
const homeDir = isDesktop && bootstrap?.homeDir ? resolve(bootstrap.homeDir) : dshHome(env, osHome);
const result = { environmentKind: isDesktop ? 'desktop' : 'cli', homeDir, profileName };
if (!profileNameValid(profileName)) return { ...result, blockedReason: 'unknown-profile' };
const profileDir = await realpath(join(homeDir, 'profiles', profileName));
const home = await realpath(homeDir);
const base = { ...result, homeDir: home, profileDir };
if (isDesktop) {
try {
if (typeof desktop?.runPlugin !== 'function' || typeof desktop?.run !== 'function' || !bootstrap
|| !current?.dir || bootstrap.activeProfileName !== profileName
|| await realpath(current.dir) !== profileDir
|| await realpath(bootstrap.activeProfileDir) !== profileDir
|| await realpath(bootstrap.appExecutable) !== await realpath(execPath)) {
return { ...base, blockedReason: 'executor-unavailable' };
}
const desktopPackage = await containingPackage(bootstrap.dshBootstrapPath, 'dsh-plugin-desktop');
if (!desktopPackage || basename(bootstrap.dshBootstrapPath) !== 'desktop-cli.js'
|| dirname(await realpath(bootstrap.dshBootstrapPath)) !== join(desktopPackage.directory, 'lib')) {
return { ...base, blockedReason: 'executor-unavailable' };
}
return { ...base, desktop, executable: execPath, cliEntry: bootstrap.dshBootstrapPath };
} catch {
return { ...base, blockedReason: 'executor-unavailable' };
}
}
const subprocess = service(ctx, 'subprocess');
if (typeof subprocess?.spawn !== 'function' || typeof argv[1] !== 'string' || !isAbsolute(argv[1]) || platform === 'win32') {
return { ...base, blockedReason: 'executor-unavailable' };
}
try {
const cli = await containingPackage(argv[1], '@deepseek-ai/dsh');
if (!cli) return { ...base, blockedReason: 'executor-unavailable' };
const cliEntry = await realpath(argv[1]);
const declared = typeof cli.manifest.bin === 'string' ? cli.manifest.bin : cli.manifest.bin?.dsh;
const publishedEntry = typeof declared === 'string' ? resolve(cli.directory, declared) : '';
if (cliEntry !== publishedEntry && cliEntry !== join(cli.directory, 'src', 'bin.ts')) {
return { ...base, blockedReason: 'executor-unavailable' };
}
return { ...base, subprocess, executable: execPath, cliEntry };
} catch {
return { ...base, blockedReason: 'executor-unavailable' };
}
}
function cliOperation(runtime, args, signal, directPnpm = false) {
const child = runtime.subprocess.spawn({
argv: directPnpm ? ['pnpm', ...args] : [execPath, ...execArgv, runtime.cliEntry, ...args],
cwd: directPnpm ? runtime.profileDir : cwd,
env: { DSH_HOME: runtime.homeDir, CI: 'true' },
stdio: { stdin: 'ignore', stdout: 'pipe', stderr: 'pipe' },
graceMs: 3_000,
signal,
});
return {
stdout: child.stdout,
stderr: child.stderr,
cancel: () => child.terminate(),
done: (async () => {
try { return await child.done; }
finally { await child.waitForExit(); }
})(),
};
}
async function checkRegistry(runtime, signal) {
const args = ['config', 'get', '@xmanrui:registry', '--json'];
const result = await run(
(childSignal) => runtime.desktop
? runtime.desktop.run(args, childSignal)
: cliOperation(runtime, args, childSignal, true),
{ signal, timeoutMs: options.configTimeoutMs ?? CONFIG_TIMEOUT_MS, errorCode: 'registry-check-failed', capture: true },
);
const output = result.stdout.trim();
let value;
try { value = output === '' || output === 'undefined' ? undefined : JSON.parse(output); }
catch { throw failure('registry-check-failed'); }
if (!officialRegistry(value)) throw failure('registry-conflict');
}
async function inspect({ preflight = false } = {}) {
let runtime;
const result = { installedVersion: null, packageValid: false, eligible: false, installationKey: null };
try {
runtime = await environment();
for (const key of ['homeDir', 'profileDir', 'profileName', 'environmentKind', 'executable', 'cliEntry']) {
result[key] = runtime[key] ?? null;
}
if (!runtime.profileDir) return { ...result, blockedReason: runtime.blockedReason ?? 'unknown-profile' };
const profile = await packageAt(runtime.profileDir);
const installed = await packageAt(join(runtime.profileDir, 'node_modules', PACKAGE_NAME));
result.installedVersion = typeof installed.manifest.version === 'string' ? installed.manifest.version : null;
result.packageValid = await validPackage(installed);
const loaded = await loadedPackage;
const identity = `${runtime.homeDir}\0${runtime.profileDir}\0${runtime.profileName}`;
const sameLoadedPackage = loaded?.directory === installed.directory
&& loaded?.manifest.version === installed.manifest.version;
if (boundProfile === undefined && sameLoadedPackage && result.packageValid) boundProfile = identity;
const stateFiles = await Promise.all(['pnpm-lock.yaml', 'pnpm-workspace.yaml', 'package-lock.json']
.map((filename) => readOptional(join(runtime.profileDir, filename))));
result.installationKey = createHash('sha256').update(JSON.stringify([
identity, profile.contents, installed.directory, installed.contents,
runtime.cliEntry, runtime.executable, ...stateFiles,
])).digest('hex');
if (boundProfile !== undefined && boundProfile !== identity) result.blockedReason = 'installation-changed';
else if (!sameLoadedPackage && boundProfile === undefined) result.blockedReason = 'installation-changed';
else if (!result.packageValid) result.blockedReason = 'invalid-installation';
else if (!registrySpec(profile.manifest.dependencies?.[PACKAGE_NAME])
|| !inside(join(runtime.profileDir, 'node_modules'), installed.directory)) result.blockedReason = 'source-install';
else if (runtime.blockedReason) result.blockedReason = runtime.blockedReason;
else if (!sameLoadedPackage) result.blockedReason = 'pending-restart';
else if (preflight) await checkRegistry(runtime);
result.eligible = !result.blockedReason;
return { ...result, blockedReason: result.blockedReason ?? null };
} catch (error) {
const known = ['registry-conflict', 'registry-check-failed', 'install-timeout', 'install-interrupted'];
return { ...result, blockedReason: known.includes(error.code) ? error.code : runtime?.profileDir ? 'invalid-installation' : 'unknown-profile' };
}
}
async function install(version, { signal, expectedInstallationKey } = {}) {
if (semver.valid(version) !== version || semver.prerelease(version)) throw failure('invalid-version');
const before = await inspect();
if (expectedInstallationKey !== undefined && before.installationKey !== expectedInstallationKey) {
throw failure('installation-changed');
}
if (!before.eligible) throw failure(before.blockedReason);
const runtime = await environment();
await checkRegistry(runtime, signal);
// Config validation is asynchronous: do not replace a package another
// package manager changed while it ran.
const checked = await inspect();
if (!checked.eligible || checked.installationKey !== before.installationKey) throw failure('installation-changed');
const args = ['add', '-w', '--save-exact', `${PACKAGE_NAME}@${version}`, `--registry=${NPM_REGISTRY}`];
return run(
(childSignal) => runtime.desktop
? runtime.desktop.runPlugin(args, runtime.profileDir, childSignal)
: cliOperation(runtime, ['plugin', '--profile', runtime.profileName, ...args], childSignal),
{ signal, timeoutMs: options.installTimeoutMs ?? INSTALL_TIMEOUT_MS, errorCode: 'install-failed' },
);
}
return Object.freeze({ inspect, install });
}

View file

@ -0,0 +1,385 @@
import { createHash, randomUUID } from 'node:crypto';
import { mkdir, open, readFile, rename, stat, unlink, writeFile } from 'node:fs/promises';
import { isAbsolute, join } from 'node:path';
import semver from 'semver';
import manifest from '../../package.json' with { type: 'json' };
import { withSessionBindingLock } from '../../src/channels/shared/session-binding-lock.mjs';
import { NPM_REGISTRY, PACKAGE_NAME } from './update-runtime.mjs';
const ACTIVE_STATES = new Set(['installing', 'verifying']);
const JOB_STATES = new Set([...ACTIVE_STATES, 'restart-required', 'completed', 'failed', 'interrupted']);
const MAX_METADATA_BYTES = 256 * 1024;
const SNAPSHOT_FILES = ['package.json', 'pnpm-lock.yaml', 'pnpm-workspace.yaml'];
const NO_LOCK = Symbol('no update lock');
export function updateError(code) {
return Object.assign(new Error(code), { code });
}
/** Read only the fixed npm package; neither RPC callers nor registry metadata choose a command. */
export async function fetchNpmRelease(fetchImpl = globalThis.fetch, timeoutMs = 10_000) {
let response;
try {
response = await fetchImpl(`${NPM_REGISTRY}${encodeURIComponent(PACKAGE_NAME)}/latest`, {
headers: { accept: 'application/json' },
redirect: 'error',
signal: AbortSignal.timeout(timeoutMs),
});
if (!response.ok) throw updateError('check-failed');
if (Number(response.headers.get('content-length')) > MAX_METADATA_BYTES) {
throw updateError('invalid-release');
}
const chunks = [];
let length = 0;
for await (const chunk of response.body) {
length += chunk.byteLength;
if (length > MAX_METADATA_BYTES) throw updateError('invalid-release');
chunks.push(Buffer.from(chunk));
}
const value = JSON.parse(Buffer.concat(chunks).toString('utf8'));
const version = value.version;
if (value.name !== PACKAGE_NAME || typeof version !== 'string'
|| semver.valid(version) !== version || semver.prerelease(version)) {
throw updateError('invalid-release');
}
const nodeRange = value.engines?.node;
if (nodeRange !== undefined && (typeof nodeRange !== 'string' || !semver.validRange(nodeRange))) {
throw updateError('invalid-release');
}
const tarball = new URL(value.dist?.tarball);
const integrity = value.dist?.integrity;
if (tarball.origin !== new URL(NPM_REGISTRY).origin || tarball.username || tarball.password
|| tarball.search || tarball.hash
|| tarball.pathname !== `/@xmanrui/dsh-im/-/dsh-im-${version}.tgz`
|| typeof integrity !== 'string' || !/^sha512-[A-Za-z0-9+/]{86}==$/.test(integrity)) {
throw updateError('invalid-release');
}
return { version, nodeRange: nodeRange ?? '*', integrity, tarball: tarball.href };
} catch (error) {
if (error?.code === 'invalid-release' || error instanceof SyntaxError || error instanceof TypeError && response?.ok) {
throw updateError('invalid-release');
}
throw updateError('check-failed');
}
}
function pathsFor(environment) {
if (!isAbsolute(environment.homeDir ?? '') || !isAbsolute(environment.profileDir ?? '')) return null;
const key = createHash('sha256').update(environment.profileDir).digest('hex').slice(0, 24);
const directory = join(environment.homeDir, 'updates', 'dsh-im', key);
return {
directory,
state: join(directory, 'state.json'),
lock: join(directory, 'install.lock'),
backup: join(directory, 'before.json'),
};
}
async function readJson(path, missing = null) {
try {
if ((await stat(path)).size > 10 * 1024 * 1024) throw updateError('state-unavailable');
return JSON.parse(await readFile(path, 'utf8'));
} catch (error) {
if (error.code === 'ENOENT') return missing;
throw updateError('state-unavailable');
}
}
async function writeJson(path, value) {
const temporary = `${path}.${randomUUID()}.tmp`;
try {
await writeFile(temporary, `${JSON.stringify(value, null, 2)}\n`, { mode: 0o600, flag: 'wx' });
await rename(temporary, path);
} catch {
throw updateError('state-unavailable');
} finally {
await unlink(temporary).catch((error) => {
if (error.code !== 'ENOENT') throw updateError('state-unavailable');
});
}
}
function processAlive(pid) {
if (!Number.isSafeInteger(pid) || pid <= 0) return false;
try {
process.kill(pid, 0);
return true;
} catch (error) {
return error.code !== 'ESRCH';
}
}
function publicJob(job) {
if (!job) return null;
const { id, state, targetVersion, message } = job;
return { id, state, targetVersion, message };
}
/** One update job per profile. Persist intent before starting pnpm, and never apply a restart here. */
export function createUpdateService({
runtime,
runningVersion = manifest.version,
nodeVersion = process.versions.node,
fetchImpl = globalThis.fetch,
now = Date.now,
checkTimeoutMs = 10_000,
confirmationTtlMs = 10 * 60_000,
installTimeoutMs = 15 * 60_000,
} = {}) {
const queue = {};
let checked = null;
let checking = null;
let lastCheckAt = -Infinity;
let activeJob = null;
let activeTask = null;
let abortController = null;
let unsavedJob = null;
let disposed = false;
function assertActive() {
if (disposed) throw updateError('disposed');
}
async function readJob(environment) {
const paths = pathsFor(environment);
if (!paths) return null;
// A failed final write must not make this Host display the old "installing"
// record forever. Keep the lock and expose the outcome we actually observed.
if (unsavedJob?.statePath === paths.state) return unsavedJob.job;
let job = await readJson(paths.state);
const lock = await readJson(paths.lock, NO_LOCK);
if (!job) {
if (lock !== NO_LOCK) {
return { id: 'locked', state: 'interrupted', message: 'recovery-required', targetVersion: null };
}
return null;
}
if (!JOB_STATES.has(job.state) || typeof job.id !== 'string' || !semver.valid(job.targetVersion)) {
throw updateError('state-unavailable');
}
if (ACTIVE_STATES.has(job.state) && job.id !== activeJob?.id) {
if (lock === NO_LOCK || lock?.id !== job.id || !processAlive(lock?.pid)) {
job = { ...job, state: 'interrupted', message: 'recovery-required' };
}
}
if (job.id !== activeJob?.id && !ACTIVE_STATES.has(job.state)) {
if (lock !== NO_LOCK) return { ...job, state: 'interrupted', message: 'recovery-required' };
// A later Host can retry after a verified manual repair. Never infer this
// from version equality while an old process lock is still present.
if (['failed', 'interrupted'].includes(job.state) && environment.packageValid === true
&& environment.installedVersion === runningVersion && !environment.blockedReason) {
return job.targetVersion === runningVersion
? { ...job, state: 'completed', message: 'recovered' }
: { ...job, recoverable: true };
}
}
if (job.state === 'restart-required' || job.state === 'completed') {
if (environment.installedVersion !== job.targetVersion || environment.packageValid !== true) {
return { ...job, state: 'interrupted', message: 'installation-changed' };
}
return { ...job, state: runningVersion === job.targetVersion ? 'completed' : 'restart-required' };
}
return job;
}
function snapshot(environment, job) {
let blockedReason = environment.blockedReason ?? checked?.blockedReason ?? null;
if (job?.state === 'interrupted' && !job.recoverable
|| job?.state === 'failed' && environment.installedVersion !== runningVersion) {
blockedReason = 'recovery-required';
} else if (job?.state === 'restart-required' || environment.installedVersion && environment.installedVersion !== runningVersion) {
blockedReason = 'pending-restart';
}
const busy = ACTIVE_STATES.has(job?.state);
const canInstall = Boolean(environment.eligible && !blockedReason && !busy && checked?.checkId
&& now() < checked.expiresAt && checked.installationKey === environment.installationKey
&& semver.valid(runningVersion) && semver.gt(checked.release.version, runningVersion));
return {
runningVersion,
installedVersion: environment.installedVersion ?? null,
latestVersion: checked?.release.version ?? null,
profileName: environment.profileName ?? null,
environmentKind: environment.environmentKind ?? 'cli',
canInstall,
blockedReason,
checkedAt: checked?.checkedAt ?? null,
checkId: canInstall ? checked.checkId : null,
job: publicJob(job),
};
}
async function status() {
assertActive();
const environment = await runtime.inspect();
return snapshot(environment, await readJob(environment));
}
async function check() {
assertActive();
if (checking) return checking;
// Collapse double clicks without turning a failed request into a cached success.
if (checked?.checkId && now() - lastCheckAt < 2_000) return status();
lastCheckAt = now();
checking = (async () => {
try {
const environment = await runtime.inspect({ preflight: true });
const release = await fetchNpmRelease(fetchImpl, checkTimeoutMs);
assertActive();
checked = {
release,
checkId: randomUUID(),
checkedAt: now(),
expiresAt: now() + confirmationTtlMs,
installationKey: environment.installationKey,
profileDir: environment.profileDir,
blockedReason: environment.blockedReason
?? (!semver.satisfies(nodeVersion, release.nodeRange) ? 'incompatible-node' : null),
};
return snapshot(environment, await readJob(environment));
} catch (error) {
if (checked) checked = { ...checked, checkId: null, expiresAt: 0 };
throw error;
} finally {
checking = null;
}
})();
return checking;
}
async function releaseLock(paths, id) {
const lock = await readJson(paths.lock);
if (lock?.id === id) await unlink(paths.lock);
}
async function backupProfile(environment, paths, job) {
const files = {};
for (const filename of SNAPSHOT_FILES) {
try {
const path = join(environment.profileDir, filename);
if ((await stat(path)).size > 3 * 1024 * 1024) throw updateError('state-unavailable');
files[filename] = await readFile(path, 'utf8');
} catch (error) {
if (error.code !== 'ENOENT') throw updateError('state-unavailable');
}
}
// Keep only the previous attempt's small manifests, never credentials or the entire home.
await writeJson(paths.backup, { jobId: job.id, previousVersion: job.previousVersion, files });
}
function rememberUnsavedJob(paths) {
activeJob = { ...activeJob, state: 'interrupted', message: 'state-unavailable' };
unsavedJob = { statePath: paths.state, job: activeJob };
}
async function execute(paths, job, originalEnvironment) {
const deadline = AbortSignal.timeout(installTimeoutMs);
try {
await runtime.install(job.targetVersion, {
signal: AbortSignal.any([abortController.signal, deadline]),
expectedInstallationKey: originalEnvironment.installationKey,
});
if (disposed) throw updateError('interrupted');
activeJob = { ...activeJob, state: 'verifying', updatedAt: now() };
await writeJson(paths.state, activeJob);
const environment = await runtime.inspect();
if (environment.homeDir !== originalEnvironment.homeDir || environment.profileDir !== originalEnvironment.profileDir
|| environment.profileName !== originalEnvironment.profileName || environment.blockedReason === 'installation-changed') {
throw updateError('installation-changed');
}
if (environment.installedVersion !== job.targetVersion || environment.packageValid !== true) {
throw updateError('verify-failed');
}
if (environment.blockedReason && environment.blockedReason !== 'pending-restart') throw updateError('verify-failed');
activeJob = { ...activeJob, state: 'restart-required', message: null, updatedAt: now() };
await writeJson(paths.state, activeJob);
} catch (error) {
const timedOut = deadline.aborted || error.code === 'install-timeout';
const interrupted = disposed || abortController.signal.aborted
|| !timedOut && ['interrupted', 'install-interrupted'].includes(error.code);
activeJob = {
...activeJob,
state: interrupted ? 'interrupted' : 'failed',
message: interrupted ? 'interrupted' : timedOut ? 'install-timeout'
: ['verify-failed', 'state-unavailable', 'installation-changed', 'registry-conflict'].includes(error.code)
? error.code : 'install-failed',
updatedAt: now(),
};
try {
await writeJson(paths.state, activeJob);
} catch {
// Retain the lock when the final record cannot be saved; the next Host must inspect it.
rememberUnsavedJob(paths);
return;
}
}
await releaseLock(paths, job.id);
}
function install({ checkId, requestId }) {
return withSessionBindingLock(queue, 'install', async () => {
assertActive();
let environment = await runtime.inspect({ preflight: true });
const previous = await readJob(environment);
if (previous?.requestId === requestId) return snapshot(environment, previous);
if (ACTIVE_STATES.has(previous?.state) || previous?.state === 'restart-required') throw updateError('update-busy');
const confirmation = checked;
if (!confirmation || !checkId || confirmation.checkId !== checkId || now() >= confirmation.expiresAt) {
throw updateError('check-expired');
}
if (environment.profileDir !== confirmation.profileDir || environment.installationKey !== confirmation.installationKey) {
throw updateError('installation-changed');
}
if (!snapshot(environment, previous).canInstall) throw updateError(environment.blockedReason ?? 'update-busy');
const paths = pathsFor(environment);
if (!paths) throw updateError('state-unavailable');
const job = {
id: randomUUID(), requestId, state: 'installing', message: null,
targetVersion: confirmation.release.version, previousVersion: environment.installedVersion,
startedAt: now(), updatedAt: now(),
};
let locked = false;
try {
await mkdir(paths.directory, { recursive: true, mode: 0o700 });
const lock = await open(paths.lock, 'wx', 0o600);
locked = true;
try {
await lock.writeFile(JSON.stringify({ id: job.id, pid: process.pid, startedAt: now() }));
} finally {
await lock.close();
}
const currentRelease = await fetchNpmRelease(fetchImpl, checkTimeoutMs);
if (JSON.stringify(currentRelease) !== JSON.stringify(confirmation.release)) throw updateError('check-expired');
environment = await runtime.inspect({ preflight: true });
if (environment.profileDir !== confirmation.profileDir || environment.installationKey !== confirmation.installationKey) {
throw updateError('installation-changed');
}
if (!environment.eligible || environment.blockedReason) throw updateError(environment.blockedReason ?? 'update-busy');
assertActive();
await backupProfile(environment, paths, job);
await writeJson(paths.state, job);
assertActive();
} catch (error) {
if (locked) await releaseLock(paths, job.id);
if (error.code === 'EEXIST') throw updateError('update-busy');
if (error.code?.startsWith('E')) throw updateError('state-unavailable');
throw error;
}
activeJob = job;
abortController = new AbortController();
activeTask = execute(paths, job, environment).catch(() => {
// Keep unknown cleanup failures local and keep the on-disk lock for manual recovery.
rememberUnsavedJob(paths);
});
return snapshot(environment, job);
});
}
async function close() {
disposed = true;
abortController?.abort();
if (activeTask) await activeTask;
}
return Object.freeze({ status, check, install, close });
}