diff --git a/CHANGELOG.md b/CHANGELOG.md index 038ce141f..9db7a473d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -22,6 +22,10 @@ Versioning](https://semver.org/spec/v2.0.0.html). ### Fixed +- Fixed issue where the local cache could restore corrupt output. If Wireit was + killed or crashed while writing a cache entry, later runs restored the partial + entry, with files missing or truncated and no error. Entries are now written + to a temp folder and renamed into place. - GitHub Actions caching now uses `http` or `https` based on the scheme of `ACTIONS_RESULTS_URL`. Always calling `https.request` broke `http://` cache proxies used by some third-party runners. diff --git a/README.md b/README.md index b34c607cf..f2a106034 100644 --- a/README.md +++ b/README.md @@ -370,7 +370,7 @@ Wireit reminds you once a day with the folders to delete. Note the limit is applied per script, so a package with many cached scripts will still use a multiple of this space. To free all of it at once, use -`rm -rf .wireit/*/cache .wireit/trash`. +`rm -rf .wireit/*/cache .wireit/*/temp .wireit/trash`. ### GitHub Actions caching diff --git a/src/caching/local-cache.ts b/src/caching/local-cache.ts index 3238c8e79..ef543be84 100644 --- a/src/caching/local-cache.ts +++ b/src/caching/local-cache.ts @@ -53,11 +53,16 @@ const REMIND_OVER_LIMIT_EVERY_MS = 24 * 60 * 60 * 1000; * an entry and evicts up to {@link MAX_EVICTIONS_PER_WRITE}. Evicted entries * move to the package's ".wireit/trash", which {@link sweepTrash} empties. * + * Entries are copied into the script's "temp" folder and then renamed into + * place, so a killed Wireit can't leave a partial entry. + * * Eviction needs no lock of its own: it touches only the calling script's cache * folder, and StandardScriptExecution#acquireSystemLockIfNeeded already holds * that script's lock, except for an empty "output", where the entries are empty - * directories. Sweeping is deliberately unlocked, so any number of Wireit - * processes can empty the same trash at once and a vanished entry is expected. + * directories. The lock also means anything in the script's temp folder was + * left by a killed or failed write, not one in progress. Sweeping is + * deliberately unlocked, so any number of Wireit processes can empty the same + * trash at once and a vanished entry is expected. */ export class LocalCache implements Cache { readonly #maxEntries: number; @@ -119,20 +124,67 @@ export class LocalCache implements Cache { ): Promise { this.#packageDirs.add(script.packageDir); const absCacheDir = this.#getCacheDir(script, fingerprint); - // Note fs.mkdir returns the first created directory, or undefined if no - // directory was created. - const existed = - (await fs.mkdir(absCacheDir, {recursive: true})) === undefined; - if (existed) { - // This is an unexpected error because the Executor should already have - // checked for an existing cache hit. - throw new Error(`Did not expect ${absCacheDir} to already exist.`); + if (absoluteFiles.length === 0) { + // No temp folder, because an empty "output" runs without the lock. + // + // Note fs.mkdir returns the first created directory, or undefined if no + // directory was created. + const existed = + (await fs.mkdir(absCacheDir, {recursive: true})) === undefined; + if (existed) { + // This is an unexpected error because the Executor should already have + // checked for an existing cache hit. + throw new Error(`Did not expect ${absCacheDir} to already exist.`); + } + await this.#evictLeastRecentlyUsed(script, pathlib.basename(absCacheDir)); + return true; } - await copyEntries(absoluteFiles, script.packageDir, absCacheDir); - await this.#evictLeastRecentlyUsed(script, pathlib.basename(absCacheDir)); + await this.#writeThroughTemp(script, absoluteFiles, absCacheDir); + await Promise.all([ + this.#evictLeastRecentlyUsed(script, pathlib.basename(absCacheDir)), + this.#trashLeftoverTemp(script), + ]); return true; } + async #writeThroughTemp( + script: ScriptReference, + absoluteFiles: AbsoluteEntry[], + absCacheDir: string, + ): Promise { + // Short, so a path that fits the Windows limit in the cache fits here. + const tempDir = pathlib.join( + this.#getScriptTempDir(script), + randomBytes(8).toString('hex'), + ); + await fs.mkdir(this.#getScriptCacheDir(script), {recursive: true}); + try { + await copyEntries(absoluteFiles, script.packageDir, tempDir); + await fs.rename(tempDir, absCacheDir); + } catch (error) { + // Not moved to the trash, because creating the trash folder fails on a + // full disk. + await fs.rmTree(tempDir).catch(() => {}); + throw error; + } + } + + /** Needs the script's lock, and must run after this write's rename. */ + async #trashLeftoverTemp(script: ScriptReference): Promise { + const tempDir = this.#getScriptTempDir(script); + let leftovers; + try { + leftovers = await fs.readdir(tempDir, {withFileTypes: true}); + } catch { + return; + } + await Promise.allSettled( + leftovers.map((entry) => + this.#moveToTrash(script.packageDir, pathlib.join(tempDir, entry.name)), + ), + ); + } + async sweepTrash({ signal, background = false, @@ -302,6 +354,10 @@ export class LocalCache implements Cache { return pathlib.join(getScriptDataDir(script), 'cache'); } + #getScriptTempDir(script: ScriptReference): string { + return pathlib.join(getScriptDataDir(script), 'temp'); + } + #getCacheDir(script: ScriptReference, fingerprint: Fingerprint): string { return pathlib.join( this.#getScriptCacheDir(script), diff --git a/src/test/cache-local.test.ts b/src/test/cache-local.test.ts index 872a6a977..635a8094e 100644 --- a/src/test/cache-local.test.ts +++ b/src/test/cache-local.test.ts @@ -72,6 +72,20 @@ async function startA( return exec; } +async function readdirIfExists( + rig: WireitTestRig, + path: string, +): Promise { + try { + return (await fs.readdir(rig.resolve(path))).sort(); + } catch (error) { + if ((error as {code?: string}).code === 'ENOENT') { + return []; + } + throw error; + } +} + /** The names of the entries in script "a"'s cache folder. */ async function cacheEntries(rig: WireitTestRig): Promise { const cacheDir = pathlib.join( @@ -81,6 +95,15 @@ async function cacheEntries(rig: WireitTestRig): Promise { return (await fs.readdir(cacheDir)).sort(); } +const dataDir = (rig: WireitTestRig, script: string) => + getScriptDataDir({packageDir: rig.resolve('.'), name: script}); + +const cacheEntriesIfAny = (rig: WireitTestRig, script: string) => + readdirIfExists(rig, pathlib.join(dataDir(rig, script), 'cache')); + +const tempEntries = (rig: WireitTestRig) => + readdirIfExists(rig, pathlib.join(dataDir(rig, 'a'), 'temp')); + /** Writes entries into the trash, as an interrupted sweep leaves them. */ async function writeTrash( rig: WireitTestRig, @@ -98,19 +121,9 @@ async function writeTrash( /** The number of files in the trash, as written by {@link writeTrash}. */ async function countTrashFiles(rig: WireitTestRig): Promise { - const readdirIfExists = async (path: string) => { - try { - return await fs.readdir(rig.resolve(path)); - } catch (error) { - if ((error as {code?: string}).code === 'ENOENT') { - return []; - } - throw error; - } - }; let count = 0; - for (const entry of await readdirIfExists(TRASH)) { - count += (await readdirIfExists(pathlib.join(TRASH, entry))).length; + for (const entry of await readdirIfExists(rig, TRASH)) { + count += (await readdirIfExists(rig, pathlib.join(TRASH, entry))).length; } return count; } @@ -134,6 +147,10 @@ const TRASH_DELETIONS = { const holdTrashDeletions = (rig: WireitTestRig) => gateWireitFs(rig, TRASH_DELETIONS); +/** Holds script "a"'s writes to the cache. Restores would be held too. */ +const holdCacheWrites = (rig: WireitTestRig) => + gateWireitFs(rig, {functions: ['copyFile'], path: /[\\/]output$/}); + void test( 'WIREIT_CACHE_MAX_ENTRIES caps the cache directory end to end', rigTest( @@ -533,3 +550,108 @@ for (const signal of ['SIGINT', 'SIGTERM'] as const) { ); } } + +void test( + 'a crash while writing an entry leaves no entry, and the next run cleans up', + {timeout: DEFAULT_TIMEOUT}, + rigTest(async ({rig}) => { + const cmdA = await writePackage(rig); + await using gate = await holdCacheWrites(rig); + const crashed = await startA(rig, cmdA, 'v0'); + await gate.firstCall(crashed); + crashed.kill('SIGKILL'); + await crashed.exit; + assert.deepEqual(await cacheEntriesIfAny(rig, 'a'), []); + assert.equal((await tempEntries(rig)).length, 1); + + rig.env = {...rig.env, WIREIT_TEST_FS_GATE: undefined}; + // Otherwise the killed run's lock takes 10 seconds to go stale. + await rig.delete(pathlib.join(dataDir(rig, 'a'), 'lock.lock')); + // Otherwise the run is fresh and never looks in the cache. + await rig.delete('output'); + assert.equal((await (await startA(rig, cmdA, 'v0')).exit).code, 0); + assert.equal(cmdA.numInvocations, 2); + assert.deepEqual(await tempEntries(rig), []); + await assertNoTrash(rig); + + assert.equal((await cacheEntries(rig)).length, 1); + await rig.delete('output'); + assert.equal((await rig.exec('npm run a').exit).code, 0); + assert.equal(cmdA.numInvocations, 2); + assert.equal(await rig.read('output'), 'v0'); + }), +); + +void test( + 'a run of another script leaves alone an entry still being written', + {timeout: DEFAULT_TIMEOUT}, + rigTest(async ({rig}) => { + const cmdA = await rig.newCommand(); + const cmdB = await rig.newCommand(); + await rig.write({ + 'package.json': { + scripts: {a: 'wireit', b: 'wireit'}, + wireit: { + a: {command: cmdA.command, files: ['input'], output: ['output']}, + b: {command: cmdB.command, files: ['input'], output: ['outputB']}, + }, + }, + }); + await using gate = await holdCacheWrites(rig); + const execA = await startA(rig, cmdA, 'v0'); + await gate.firstCall(execA); + + rig.env = {...rig.env, WIREIT_TEST_FS_GATE: undefined}; + const execB = rig.exec('npm run b'); + const invB = await cmdB.nextInvocation(); + await rig.write({outputB: 'v0'}); + invB.exit(0); + assert.equal((await execB.exit).code, 0); + // Only a write cleans up the temp folder. + assert.equal((await cacheEntriesIfAny(rig, 'b')).length, 1); + + await gate.release(); + assert.equal((await execA.exit).code, 0); + assert.deepEqual(await tempEntries(rig), []); + await rig.delete('output'); + assert.equal((await rig.exec('npm run a').exit).code, 0); + assert.equal(cmdA.numInvocations, 1); + assert.equal(await rig.read('output'), 'v0'); + }), +); + +void test( + 'a run of the same script waits for an entry still being written', + {timeout: DEFAULT_TIMEOUT}, + rigTest(async ({rig}) => { + const cmdA = await writePackage(rig); + await using gate = await holdCacheWrites(rig); + const first = await startA(rig, cmdA, 'v0'); + await gate.firstCall(first); + + // The quiet logger doesn't log waiting for a lock. + rig.env = { + ...rig.env, + WIREIT_TEST_FS_GATE: undefined, + WIREIT_LOGGER: 'simple', + }; + await rig.write({input: 'v1'}); + const second = rig.exec('npm run a'); + await waitForLog(second, /Waiting for another process/); + + await gate.release(); + assert.equal((await first.exit).code, 0); + const inv = await withTimeout('the second run', cmdA.nextInvocation()); + await rig.write({output: 'v1'}); + inv.exit(0); + assert.equal((await second.exit).code, 0); + assert.equal((await cacheEntries(rig)).length, 2); + assert.deepEqual(await tempEntries(rig), []); + + await rig.write({input: 'v0'}); + await rig.delete('output'); + assert.equal((await rig.exec('npm run a').exit).code, 0); + assert.equal(cmdA.numInvocations, 2); + assert.equal(await rig.read('output'), 'v0'); + }), +); diff --git a/src/test/local-cache.test.ts b/src/test/local-cache.test.ts index 855d3b757..1a7f0279b 100644 --- a/src/test/local-cache.test.ts +++ b/src/test/local-cache.test.ts @@ -47,6 +47,12 @@ async function setup(maxEntries: number): Promise< /** The names of the evicted entries waiting to be swept. */ trashEntries: () => Promise; + + /** The "output" file, as {@link LocalCache.set} takes it. */ + outputEntry: AbsoluteEntry; + + /** The names in the script's folder of entries still being written. */ + tempEntries: () => Promise; } & AsyncDisposable > { const rig = new FilesystemTestRig(); @@ -76,32 +82,17 @@ async function setup(maxEntries: number): Promise< ); }; - const entryHashes = async () => { - try { - return (await fs.readdir(cacheDir)).sort(); - } catch (error) { - if ((error as {code?: string}).code === 'ENOENT') { - return []; - } - throw error; - } - }; + const entryHashes = () => readdirIfExists(cacheDir); const setRecency = async (name: string, secondsSinceEpoch: number) => { const when = new Date(secondsSinceEpoch * 1000); await fs.utimes(pathlib.join(cacheDir, hashOf(name)), when, when); }; - const trashEntries = async () => { - try { - return (await fs.readdir(trashDir)).sort(); - } catch (error) { - if ((error as {code?: string}).code === 'ENOENT') { - return []; - } - throw error; - } - }; + const trashEntries = () => readdirIfExists(trashDir); + + const tempEntries = () => + readdirIfExists(pathlib.join(getScriptDataDir(script), 'temp')); return { rig, @@ -113,10 +104,23 @@ async function setup(maxEntries: number): Promise< setRecency, trashDir, trashEntries, + outputEntry, + tempEntries, [Symbol.asyncDispose]: () => rig.cleanup(), }; } +async function readdirIfExists(dir: string): Promise { + try { + return (await fs.readdir(dir)).sort(); + } catch (error) { + if ((error as {code?: string}).code === 'ENOENT') { + return []; + } + throw error; + } +} + /** The cache only keys off the string form, so any distinct string works. */ const fingerprint = (name: string) => Fingerprint.fromString(name as FingerprintString); @@ -459,3 +463,25 @@ void test('get returns undefined for an evicted entry', async () => { await ctx.cacheOutput('v1'); assert.equal(await ctx.cache.get(ctx.script, fingerprint('v0')), undefined); }); + +void test('a write that fails leaves no entry, and deletes its temp copy', async () => { + await using ctx = await setup(1); + await ctx.rig.write({output: 'v0'}); + { + using _gate = new FsGate({ + functions: ['copyFile'], + path: /[\\/]output$/, + failWith: 'ENOSPC', + }); + await assert.rejects( + ctx.cache.set(ctx.script, fingerprint('v0'), [ctx.outputEntry]), + {code: 'ENOSPC'}, + ); + } + assert.deepEqual(await ctx.entryHashes(), []); + assert.deepEqual(await ctx.tempEntries(), []); + assert.deepEqual(await ctx.trashEntries(), []); + + await ctx.cacheOutput('v0'); + assert.deepEqual(await ctx.entryHashes(), [hashOf('v0')]); +});