diff --git a/docs/DEEPSEEK.md b/docs/DEEPSEEK.md index 19a8d3f1..c27be37a 100644 --- a/docs/DEEPSEEK.md +++ b/docs/DEEPSEEK.md @@ -265,7 +265,9 @@ catalog the phone reads: normal connection. The compose file keeps it in the `agents` volume rather than in `./data`: a harness process reads the whole tree every time it boots, and from a Docker Desktop host directory that is some 25 seconds - against about one from a volume. (A host directory can also refuse a - rename for a moment, which the installer waits out.) + against about one from a volume. (The installer never renames into + place, which a host directory can refuse for a moment: each package is + extracted where it belongs and counts as installed once its + `package.json` — written last — names the pinned version.) - A container recreated with an existing `agents` volume keeps whatever agents that volume already has; the bundled copy reaches a *new* volume. diff --git a/packages/agent-host/src/__tests__/agentInstall.test.ts b/packages/agent-host/src/__tests__/agentInstall.test.ts index 7839e38b..ba986d37 100644 --- a/packages/agent-host/src/__tests__/agentInstall.test.ts +++ b/packages/agent-host/src/__tests__/agentInstall.test.ts @@ -9,7 +9,6 @@ import { extractTree, installBinary, installPackageTree, - renameIntoPlace, type PackagedBinary, withoutAgentBin, } from '../agentInstall'; @@ -120,6 +119,21 @@ describe('installBinary', () => { expect(fs.readdirSync(path.join(cache, '@scope+tool-linux-x64@2.0.0'), { recursive: true })).toEqual(['bin']); }); + it('installs again over a binary an interrupted install left without its marker', async () => { + const target = path.join(cache, '@scope+tool-linux-x64@2.0.0', 'bin', 'tool'); + fs.mkdirSync(path.dirname(target), { recursive: true }); + fs.writeFileSync(target, 'truncated'); + const fetchFn = vi.fn(async () => new Response(tgz)); + const options = { cacheDir: cache, pins, log, fetchFn: fetchFn as unknown as typeof fetch }; + expect(await installBinary(binary, options)).toBe(target); + expect(fetchFn).toHaveBeenCalledTimes(1); + expect(fs.readFileSync(target).equals(payload)).toBe(true); + // A marker for another pin (the same version republished) does not count either. + fs.writeFileSync(path.join(cache, '@scope+tool-linux-x64@2.0.0', '.installed'), 'sha512-other'); + await installBinary(binary, options); + expect(fetchFn).toHaveBeenCalledTimes(2); + }); + it('refuses a package this build does not pin', async () => { await expect(installBinary({ ...binary, pkg: 'unpinned' }, { cacheDir: cache, pins, log })).rejects.toThrow(/not pinned/); }); @@ -198,6 +212,38 @@ describe('extractTree', () => { } }); + it('writes the file named last after every other one', async () => { + const tar = Buffer.concat([ + tarEntry('package/package.json', '{"version":"1.0.0"}', '0', 0o644), + tarEntry('package/lib/bin.js', 'console.log(1)'), + Buffer.alloc(1024), + ]); + // Every block of the archive has been read, the other file written — + // and still no package.json. + let seenBeforeEnd: boolean | undefined; + async function* watched(): AsyncGenerator { + yield* chunks(tar.subarray(0, tar.length - 1024), 512); + seenBeforeEnd = fs.existsSync(path.join(dir, 'package.json')); + yield tar.subarray(tar.length - 1024); + } + await extractTree(watched(), dir, { last: 'package.json' }); + expect(seenBeforeEnd).toBe(false); + expect(fs.readFileSync(path.join(dir, 'package.json'), 'utf8')).toBe('{"version":"1.0.0"}'); + expect(fs.readFileSync(path.join(dir, 'lib', 'bin.js'), 'utf8')).toBe('console.log(1)'); + }); + + it('strips whatever top-level directory the package ships under, as npm does', async () => { + // DefinitelyTyped's packages are packed under their own name, not package/. + const tar = Buffer.concat([ + tarEntry('node/package.json', '{"version":"22.0.0"}', '0', 0o644), + tarEntry('node/fs.d.ts', 'declare module "fs";', '0', 0o644), + Buffer.alloc(1024), + ]); + await extractTree(chunks(tar, 512), dir, { last: 'package.json' }); + expect(fs.readFileSync(path.join(dir, 'package.json'), 'utf8')).toBe('{"version":"22.0.0"}'); + expect(fs.readFileSync(path.join(dir, 'fs.d.ts'), 'utf8')).toBe('declare module "fs";'); + }); + it('skips an entry that would leave the package, and keeps going', async () => { const tar = Buffer.concat([ tarEntry('package/../evil', 'no'), @@ -230,39 +276,6 @@ describe('extractTree', () => { }); }); -describe('renameIntoPlace', () => { - it('retries a rename a filesystem refuses for a moment, and then takes it', async () => { - const codes: Array = ['EACCES', 'EPERM', 'EBUSY', undefined]; - let calls = 0; - const rename = (): void => { - const code = codes[calls++]; - if (code !== undefined) throw Object.assign(new Error(`${code}: permission denied`), { code }); - }; - const logs: string[] = []; - await renameIntoPlace('/staging/a', '/tree/a', { rename, log: (line) => logs.push(line) }); - expect(calls).toBe(4); - // Said once, so a slow install on such a filesystem is explicable. - expect(logs).toEqual(['[install] /tree/a was busy (EPERM); retrying']); - }); - - it('fails at once for a rename that is a real error', async () => { - const failing = (): void => { - throw Object.assign(new Error('ENOTEMPTY: directory not empty'), { code: 'ENOTEMPTY' }); - }; - await expect(renameIntoPlace('/a', '/b', { rename: failing })).rejects.toThrow(/ENOTEMPTY/); - }); - - it('gives up after enough attempts on one that never takes', async () => { - let attempts = 0; - const never = (): void => { - attempts++; - throw Object.assign(new Error('EACCES: still busy'), { code: 'EACCES' }); - }; - await expect(renameIntoPlace('/a', '/b', { rename: never })).rejects.toThrow(/still busy/); - expect(attempts).toBe(8); - }); -}); - describe('installPackageTree', () => { let cache: string; const log = vi.fn(); @@ -329,6 +342,27 @@ describe('installPackageTree', () => { expect(JSON.parse(fs.readFileSync(path.join(stale, 'package.json'), 'utf8')).version).toBe('15.0.0'); }); + it('lays down again a package an interrupted install left without its package.json', async () => { + const fetchFn = registryOf(tgzOf); + await installPackageTree('@deepseek-ai/dsh', entries, { cacheDir: cache, log, fetchFn, label: 'DSH' }); + const commander = path.join(cache, '@deepseek-ai+dsh@1.0.0', 'node_modules', 'commander'); + fs.rmSync(path.join(commander, 'package.json')); + fs.rmSync(path.join(commander, 'index.js')); + await installPackageTree('@deepseek-ai/dsh', entries, { cacheDir: cache, log, fetchFn, label: 'DSH' }); + expect(fs.readFileSync(path.join(commander, 'index.js'), 'utf8')).toBe('ok'); + expect(JSON.parse(fs.readFileSync(path.join(commander, 'package.json'), 'utf8')).version).toBe('15.0.0'); + }); + + it('puts back what was nested inside a package it lays down again', async () => { + const fetchFn = registryOf(tgzOf); + await installPackageTree('@deepseek-ai/dsh', entries, { cacheDir: cache, log, fetchFn, label: 'DSH' }); + const commander = path.join(cache, '@deepseek-ai+dsh@1.0.0', 'node_modules', 'commander'); + fs.writeFileSync(path.join(commander, 'package.json'), '{"version":"14.0.0"}'); + await installPackageTree('@deepseek-ai/dsh', entries, { cacheDir: cache, log, fetchFn, label: 'DSH' }); + // shared-old was at its version, but lived inside commander's directory. + expect(JSON.parse(fs.readFileSync(path.join(commander, 'node_modules', 'shared-old', 'package.json'), 'utf8')).version).toBe('0.6.4'); + }); + it('skips an optional package the registry will not serve', async () => { const withOptional = [ ...entries.slice(0, 3), diff --git a/packages/agent-host/src/agentInstall.ts b/packages/agent-host/src/agentInstall.ts index 57241509..061f869f 100644 --- a/packages/agent-host/src/agentInstall.ts +++ b/packages/agent-host/src/agentInstall.ts @@ -89,6 +89,17 @@ function realPath(p: string): string { } } +/** + * The file beside an installed binary that says it is complete: written, + * holding the pin's sha512, only once the binary is whole. Installs never + * rename anything into place — some filesystems the cache lives on (a + * Windows host's bind mount in a container, a Windows disk an antivirus is + * scanning) refuse a rename of a file they have just seen written, for + * longer than any retry should wait — so completeness is this marker's to + * say instead. + */ +const INSTALLED_MARKER = '.installed'; + /** Install `binary` if it is not cached yet, and return its path. */ export async function installBinary(binary: PackagedBinary, options: InstallOptions): Promise { const pin = (options.pins ?? PLATFORM_PACKAGES)[binary.pkg]; @@ -97,14 +108,17 @@ export async function installBinary(binary: PackagedBinary, options: InstallOpti const prefix = `${binary.pkg.replace('/', '+')}@`; const dir = path.join(options.cacheDir, `${prefix}${pin.version}`); const target = path.join(dir, ...binary.file.split('/')); - if (isFile(target)) { + const marker = path.join(dir, INSTALLED_MARKER); + if (isFile(target) && readQuietly(marker) === pin.integrity) { linkIntoBin(options.cacheDir, target); return target; } + // Whatever is there is incomplete: it stops counting as installed before + // a byte of the new one is written. + removeQuietly(marker); fs.mkdirSync(path.dirname(target), { recursive: true }); const tarball = path.join(dir, `.download-${process.pid}.tgz`); - const partial = `${target}.part-${process.pid}`; try { const url = `${options.registry ?? DEFAULT_REGISTRY}/${binary.pkg}/-/${binary.pkg.split('/').pop()}-${pin.version}.tgz`; const res = await (options.fetchFn ?? fetch)(url); @@ -136,21 +150,18 @@ export async function installBinary(binary: PackagedBinary, options: InstallOpti const source = fs.createReadStream(tarball); let found: boolean; try { - found = await extractFile(source.pipe(createGunzip()), `package/${binary.file}`, partial); + found = await extractFile(source.pipe(createGunzip()), `package/${binary.file}`, target); } finally { source.destroy(); if (!source.closed) await once(source, 'close'); } if (!found) throw new Error(`${binary.pkg}@${pin.version} has no ${binary.file}`); - fs.chmodSync(partial, 0o755); - // A file's rename is refused by the same filesystem in the same way (see - // renameIntoPlace), so it waits the same way. - await renameIntoPlace(partial, target, { log: options.log }); + fs.chmodSync(target, 0o755); + fs.writeFileSync(marker, pin.integrity); } catch (err) { throw new Error(`${binary.label} could not be installed: ${err instanceof Error ? err.message : String(err)}`); } finally { removeQuietly(tarball); - removeQuietly(partial); } pruneOtherVersions(options.cacheDir, prefix, path.basename(dir)); @@ -192,6 +203,15 @@ function removeQuietly(target: string): void { } } +/** A small text file's contents, or undefined when it cannot be read. */ +function readQuietly(file: string): string | undefined { + try { + return fs.readFileSync(file, 'utf8'); + } catch { + return undefined; + } +} + function pruneOtherVersions(cacheDir: string, prefix: string, keep: string): void { let names: string[]; try { @@ -352,23 +372,32 @@ async function* oneBuffer(data: Buffer): AsyncGenerator { * extractFile; non-file entries other than directories (symlinks, devices) * do not occur in npm tarballs and are skipped. The executable bit of a * file's tar mode survives as 0o755. + * + * `last` names one file (relative to `dest`, `/`-separated) that is held + * back and written only after every other one: a reader that trusts that + * file to mean "complete" never sees it beside a half-written package. */ -export async function extractTree(tar: AsyncIterable, dest: string): Promise { +export async function extractTree(tar: AsyncIterable, dest: string, options: { last?: string } = {}): Promise { let pending: Buffer = Buffer.alloc(0); let inBody = false; let bodyLeft = 0; let padLeft = 0; - let sink: 'skip' | 'file' | 'dir' | 'pax' | 'longname' = 'skip'; + let sink: 'skip' | 'file' | 'held' | 'dir' | 'pax' | 'longname' = 'skip'; let meta: Buffer[] = []; let nextName: string | undefined; let out: fs.WriteStream | null = null; let outPath = ''; let mode = 0; + const lastPath = options.last !== undefined ? path.join(dest, ...options.last.split('/')) : undefined; + let held: Buffer[] | undefined; - /** npm packs under `package/`; anything else is not an npm tarball, and a - * `..` or empty segment would write outside `dest` — refuse both. */ + /** Everything sits under one top-level directory, stripped as npm strips + * it: `package/` for what `npm pack` makes, the package's own name for + * some (DefinitelyTyped's `@types/*` ship under `node/`, `retry/`, …). + * A `..` or empty segment would write outside `dest` — refused. */ const target = (name: string): string | null => { - const rel = name.startsWith('package/') ? name.slice('package/'.length) : name; + const slash = name.indexOf('/'); + const rel = slash < 0 ? '' : name.slice(slash + 1); if (rel.length === 0) return null; const parts = rel.split('/'); if (parts.some((p) => p === '' || p === '.' || p === '..')) return null; @@ -389,6 +418,8 @@ export async function extractTree(tar: AsyncIterable, dest: string): Pro /* best effort — the bit is not worth failing an install */ } } + } else if (sink === 'held') { + held = meta; } else if (sink === 'dir') { fs.mkdirSync(outPath, { recursive: true }); } @@ -421,6 +452,7 @@ export async function extractTree(tar: AsyncIterable, dest: string): Pro const where = sink === 'file' || sink === 'dir' ? target(name) : null; if (sink === 'file' && where === null) sink = 'skip'; // not an npm path: skip, keep parsing if (sink === 'dir' && where === null) sink = 'skip'; + if (sink === 'file' && where === lastPath) sink = 'held'; if (sink === 'file') { outPath = where!; fs.mkdirSync(path.dirname(outPath), { recursive: true }); @@ -438,7 +470,7 @@ export async function extractTree(tar: AsyncIterable, dest: string): Pro offset += n; bodyLeft -= n; if (sink === 'file' && out && !out.write(piece)) await once(out, 'drain'); - else if (sink === 'pax' || sink === 'longname') meta.push(Buffer.from(piece)); + else if (sink === 'pax' || sink === 'longname' || sink === 'held') meta.push(Buffer.from(piece)); if (bodyLeft === 0) await endBody(); } pending = pending.subarray(offset); @@ -446,56 +478,15 @@ export async function extractTree(tar: AsyncIterable, dest: string): Pro } finally { if (out && !out.writableFinished) out.destroy(); } + if (held && lastPath) { + fs.mkdirSync(path.dirname(lastPath), { recursive: true }); + fs.writeFileSync(lastPath, Buffer.concat(held)); + } } /** How many packages the installer downloads at once. */ const TREE_CONCURRENCY = 4; -/** How many times a directory rename is retried before it is a real failure, - * and how long the longest wait between two attempts is. Long enough to sit - * out a foreign filesystem's moment of busyness (~2.5s in all), short enough - * that a filesystem that is simply not going to allow it fails the install - * rather than hanging it. */ -const RENAME_ATTEMPTS = 8; -const RENAME_BACKOFF_MS = 500; - -/** - * Move a staged package into place. - * - * Retried, because one filesystem this runs on refuses a directory rename - * that the same filesystem accepts a moment later: a bind mount from a - * Windows host (the container's `/data`, Docker Desktop) answers EACCES while - * a file just written inside the directory is still being let go of by the - * host side. It is not a race this code could remove — the rename is - * serialised, the destination is recreated first, and the same tree of a few - * hundred packages installs on the first attempt on an ordinary filesystem — - * and it is not a failure the caller can fix, since which package it lands on - * differs from run to run. Waiting is the whole of the answer; a rename that - * keeps failing, or fails for any other reason, is still a failure. - */ -export async function renameIntoPlace( - from: string, - to: string, - options: { log?: (message: string) => void; rename?: (from: string, to: string) => void } = {}, -): Promise { - const rename = options.rename ?? fs.renameSync; - for (let attempt = 1; ; attempt++) { - try { - rename(from, to); - return; - } catch (error) { - const code = (error as NodeJS.ErrnoException).code ?? ''; - if (attempt >= RENAME_ATTEMPTS || !['EACCES', 'EPERM', 'EBUSY'].includes(code)) throw error; - if (attempt === 2) options.log?.(`[install] ${to} was busy (${code}); retrying`); - await delay(Math.min(attempt * 100, RENAME_BACKOFF_MS)); - } - } -} - -function delay(ms: number): Promise { - return new Promise((resolve) => setTimeout(resolve, ms)); -} - /** * Install one runtime's pinned dependency closure — `entries` come from a * generated pins module (src/generated/dshPackages.ts today; any pure-JS @@ -507,6 +498,12 @@ function delay(ms: number): Promise { * skipped, so an interrupted install resumes; an optional package that * cannot be fetched is skipped (npm's own semantics), everything else fails * the install. + * + * Each package is extracted straight into its destination with its + * `package.json` written last, and a package counts as installed only when + * that file names the pinned version: a package cut off mid-extraction has + * none and is laid down again by the next install. Nothing is renamed into + * place — see INSTALLED_MARKER for why. */ export async function installPackageTree( root: string, @@ -518,7 +515,6 @@ export async function installPackageTree( const label = options.label ?? root; const dir = path.join(options.cacheDir, `${root.replace('/', '+')}@${rootEntry.version}`); - const nodeModules = path.join(dir, 'node_modules'); const target = (dest: string): string => path.join(dir, ...dest.split('/')); // A package whose platform gates exclude this machine: optional ones are @@ -534,83 +530,74 @@ export async function installPackageTree( wanted.push(entry); } - const atVersion = (entry: TreePackageEntry): boolean => installedVersion(target(entry.dest)) === entry.version; - if (wanted.every(atVersion)) return target(rootEntry.dest); + // What must be laid down: every package not at its pinned version, and + // everything nested inside one, since laying a package down replaces its + // whole directory — its own node_modules included. + const stale = wanted.filter((entry) => installedVersion(target(entry.dest)) !== entry.version); + if (stale.length === 0) return target(rootEntry.dest); + const replaced = stale.map((entry) => `${entry.dest}/`); + const redo = new Set(wanted.filter((entry) => stale.includes(entry) || replaced.some((dest) => entry.dest.startsWith(dest)))); // The same tarball serves every spot its package occupies (a nested // variant may sit under several consumers): fetch once, extract per dest. const byTarball = new Map(); - for (const entry of wanted) { + for (const entry of redo) { const key = `${entry.name}@${entry.version}`; byTarball.set(key, [...(byTarball.get(key) ?? []), entry]); } - const staging = path.join(dir, `.staging-${process.pid}`); const registry = options.registry ?? DEFAULT_REGISTRY; const fetchFn = options.fetchFn ?? fetch; const started = Date.now(); options.log(`[install] downloading ${label} (${byTarball.size} packages) from ${registry}`); - try { - fs.mkdirSync(staging, { recursive: true }); - - // Phase 1: fetch every distinct tarball once, in parallel — this is the - // network-bound part. An optional package that cannot be fetched is - // dropped here (npm's own semantics); anything else fails the install. - // The tarballs are kept as they came, gzipped: a whole runtime inflated - // at once is the better part of a gigabyte held in memory, on a machine - // that may be a small VPS. - const tarballs = new Map(); - const jobs = [...byTarball.keys()]; - let next = 0; - const download = async (): Promise => { - for (;;) { - const index = next++; - if (index >= jobs.length) return; - const group = byTarball.get(jobs[index]!)!; - const first = group[0]!; - if (group.every(atVersion)) { - tarballs.set(jobs[index]!, null); // already on disk + // Phase 1: fetch every distinct tarball once, in parallel — this is the + // network-bound part. An optional package that cannot be fetched is + // dropped here (npm's own semantics); anything else fails the install. + // The tarballs are kept as they came, gzipped: a whole runtime inflated + // at once is the better part of a gigabyte held in memory, on a machine + // that may be a small VPS. + const tarballs = new Map(); + const jobs = [...byTarball.keys()]; + let next = 0; + const download = async (): Promise => { + for (;;) { + const index = next++; + if (index >= jobs.length) return; + const group = byTarball.get(jobs[index]!)!; + const first = group[0]!; + try { + const url = `${registry}/${first.name}/-/${first.name.split('/').pop()}-${first.version}.tgz`; + tarballs.set(jobs[index]!, await fetchVerifiedTarball(url, first.integrity, fetchFn)); + } catch (err) { + const reason = err instanceof Error ? err.message : String(err); + if (first.optional) { + options.log(`[install] skipping optional ${first.name}@${first.version}: ${reason}`); + tarballs.set(jobs[index]!, null); continue; } - try { - const url = `${registry}/${first.name}/-/${first.name.split('/').pop()}-${first.version}.tgz`; - tarballs.set(jobs[index]!, await fetchVerifiedTarball(url, first.integrity, fetchFn)); - } catch (err) { - const reason = err instanceof Error ? err.message : String(err); - if (first.optional) { - options.log(`[install] skipping optional ${first.name}@${first.version}: ${reason}`); - tarballs.set(jobs[index]!, null); - continue; - } - throw new Error(`${first.name} could not be installed: ${reason}`); - } + throw new Error(`${first.name} could not be installed: ${reason}`); } - }; - await Promise.all(Array.from({ length: TREE_CONCURRENCY }, () => download())); - - // Phase 2: lay the tarballs down serially, shallow dest first. A nested - // entry installs inside its parent's directory, so extracting in - // parallel would race a parent's replace against its children. - const depth = (dest: string): number => dest.split('/').length; - const ordered = wanted - .filter((entry) => tarballs.get(`${entry.name}@${entry.version}`)) - .sort((a, b) => depth(a.dest) - depth(b.dest)); - for (const entry of ordered) { - if (atVersion(entry)) continue; - const tar = gunzipSync(tarballs.get(`${entry.name}@${entry.version}`)!); - const pkgDir = target(entry.dest); - const stage = path.join(staging, entry.dest); - fs.rmSync(stage, { recursive: true, force: true }); - await extractTree(oneBuffer(tar), stage); - fs.rmSync(pkgDir, { recursive: true, force: true }); - fs.mkdirSync(path.dirname(pkgDir), { recursive: true }); - await renameIntoPlace(stage, pkgDir, { log: options.log }); } - - pruneOtherVersions(options.cacheDir, `${root.replace('/', '+')}@`, path.basename(dir)); - options.log(`[install] ${label} installed at ${dir} (${((Date.now() - started) / 1000) | 0}s)`); - return target(rootEntry.dest); - } finally { - removeQuietly(staging); + }; + await Promise.all(Array.from({ length: TREE_CONCURRENCY }, () => download())); + + // Phase 2: lay the tarballs down serially, shallow dest first. A nested + // entry installs inside its parent's directory, so extracting in + // parallel would race a parent's replace against its children. + const depth = (dest: string): number => dest.split('/').length; + const ordered = [...redo] + .filter((entry) => tarballs.get(`${entry.name}@${entry.version}`)) + .sort((a, b) => depth(a.dest) - depth(b.dest)); + for (const entry of ordered) { + const tar = gunzipSync(tarballs.get(`${entry.name}@${entry.version}`)!); + const pkgDir = target(entry.dest); + // Retried by Node itself on the transient EBUSY / EPERM a scanner or + // a foreign filesystem answers with. + fs.rmSync(pkgDir, { recursive: true, force: true, maxRetries: 10, retryDelay: 200 }); + await extractTree(oneBuffer(tar), pkgDir, { last: 'package.json' }); } + + pruneOtherVersions(options.cacheDir, `${root.replace('/', '+')}@`, path.basename(dir)); + options.log(`[install] ${label} installed at ${dir} (${((Date.now() - started) / 1000) | 0}s)`); + return target(rootEntry.dest); }