diff --git a/barretenberg/cpp/src/barretenberg/nodejs_module/msgpack_client/msgpack_client_wrapper.cpp b/barretenberg/cpp/src/barretenberg/nodejs_module/msgpack_client/msgpack_client_wrapper.cpp index b72114a00ab..39606d309d8 100644 --- a/barretenberg/cpp/src/barretenberg/nodejs_module/msgpack_client/msgpack_client_wrapper.cpp +++ b/barretenberg/cpp/src/barretenberg/nodejs_module/msgpack_client/msgpack_client_wrapper.cpp @@ -61,7 +61,21 @@ Napi::Value MsgpackClientWrapper::call(const Napi::CallbackInfo& info) } // Receive response with retry (1s timeout per attempt) - // Loop until response is ready - handles case where server is processing + // Loop until response is ready - handles case where server is processing. + // + // This waits forever by design, and there is no deadline to tune: a slow bb and a dead one look + // identical from here, so any deadline would eventually fail a legitimately slow call. The cost + // is that a bb process dying mid-call wedges the caller rather than failing it. This call blocks + // the JS thread, so on the main thread the event loop stops, the child's exit is never observed + // (it stays an unreaped zombie, which still answers kill(pid, 0)), and nothing recovers. + // + // Deliberately not addressed. Detecting it needs a liveness signal that survives the zombie + // window — an inherited pipe whose read end sees EOF is the usual answer, since descriptors + // close on death before anything reaps — which means plumbing a descriptor through this client + // and bb. The process on the other end of a synchronous client runs one thread doing hashes and + // signatures, so it is the smallest and least likely thing on the machine to be killed, and the + // fix costs more than the case is worth. The socket backend, which carries the workloads that do + // get killed, detects death and can replace the process; see bb.js BackendOptions.respawn. std::span response; while ((response = client_->receive(TIMEOUT_NS)).empty()) { // Response not ready yet, server is processing - retry diff --git a/barretenberg/ts/bb.js/src/barretenberg/index.ts b/barretenberg/ts/bb.js/src/barretenberg/index.ts index 217d2e5ff0f..e79c5ccc6f8 100644 --- a/barretenberg/ts/bb.js/src/barretenberg/index.ts +++ b/barretenberg/ts/bb.js/src/barretenberg/index.ts @@ -225,8 +225,16 @@ export class BarretenbergSync extends SyncApi { barretenbergSyncSingletonPromise = BarretenbergSync.new(options); } - barretenbergSyncSingleton = await barretenbergSyncSingletonPromise; - return barretenbergSyncSingleton; + try { + barretenbergSyncSingleton = await barretenbergSyncSingletonPromise; + return barretenbergSyncSingleton; + } catch (error) { + // Clear the failure so the next call can try again, as the asynchronous singleton does. + // Caching it would make one bad spawn permanent for the life of the process. + barretenbergSyncSingleton = undefined; + barretenbergSyncSingletonPromise = undefined; + throw error; + } } static destroySingleton() { diff --git a/barretenberg/ts/bb.js/src/barretenberg/singleton.test.ts b/barretenberg/ts/bb.js/src/barretenberg/singleton.test.ts new file mode 100644 index 00000000000..8e68c469136 --- /dev/null +++ b/barretenberg/ts/bb.js/src/barretenberg/singleton.test.ts @@ -0,0 +1,30 @@ +import { jest } from '@jest/globals'; + +import { BackendType, Barretenberg, BarretenbergSync } from './index.js'; + +jest.setTimeout(30_000); + +// A backend that cannot be created, so initialization fails without touching native code. +const UNUSABLE = { backend: BackendType.Wasm, wasmPath: '/nonexistent/wasm-directory' } as const; + +describe.each([ + ['Barretenberg', Barretenberg], + ['BarretenbergSync', BarretenbergSync], +])('%s.initSingleton', (_name, Api) => { + afterEach(async () => { + await Api.destroySingleton(); + }); + + it('tries again after a failed initialization instead of caching the failure', async () => { + const first = await Api.initSingleton(UNUSABLE).catch((err: unknown) => err); + const second = await Api.initSingleton(UNUSABLE).catch((err: unknown) => err); + + // Not toBeInstanceOf: jest runs the test in its own vm context, so an Error thrown by node's + // own fs is not an instance of this realm's Error. + expect((first as Error).message).toMatch(/no such file or directory/); + expect((second as Error).message).toMatch(/no such file or directory/); + // A second attempt, not the first one's cached rejection. Caching it would make one bad spawn + // permanent for the life of the process, with no way back short of a restart. + expect(second).not.toBe(first); + }); +}); diff --git a/barretenberg/ts/bb.js/src/bb_backends/errors.ts b/barretenberg/ts/bb.js/src/bb_backends/errors.ts new file mode 100644 index 00000000000..1c5bf976b9d --- /dev/null +++ b/barretenberg/ts/bb.js/src/bb_backends/errors.ts @@ -0,0 +1,29 @@ +/** + * The cross-layer contract for a failed call is the bare `retry` property, not this class. + * + * An error with `retry === true` failed for environmental reasons — the bb process died, its + * connection broke, the machine was too loaded to start one — and the operation may be retried. + * `retry === false`, or no such property, means retrying cannot help: the command itself is the + * problem, or the backend was destroyed by its owner. + * + * Consumers feature-detect the property rather than importing this class, so the convention + * survives package boundaries: + * + * if (err instanceof Error && (err as Error & { retry?: unknown }).retry === true) { ... } + * + * This is deliberately the same convention the ipc-runtime transports use, so a caller written + * against one works unchanged against the other. + */ +export class BackendUnavailableError extends Error { + readonly retry = true; + + constructor(message: string, options?: { cause?: unknown }) { + super(message, options); + this.name = 'BackendUnavailableError'; + } +} + +/** Whether a failure was environmental, and so may be retried. */ +export function isRetryable(err: unknown): boolean { + return err instanceof Error && (err as Error & { retry?: unknown }).retry === true; +} diff --git a/barretenberg/ts/bb.js/src/bb_backends/index.ts b/barretenberg/ts/bb.js/src/bb_backends/index.ts index 760f51fe9da..a538a9b614a 100644 --- a/barretenberg/ts/bb.js/src/bb_backends/index.ts +++ b/barretenberg/ts/bb.js/src/bb_backends/index.ts @@ -46,6 +46,25 @@ export type BackendOptions = { */ maxClients?: number; + /** + * @description Replace the bb process when it dies, on the next call (NativeUnixSocket only). + * + * There is no equivalent for the shared-memory backends, and no error either: a call to a bb + * process that has died never returns, because the receive loop cannot tell a dead server from a + * slow one. That is documented where it happens, in msgpack_client_wrapper.cpp. A caller that + * needs to survive a bb death should be on this backend. + * + * Calls that were in flight still fail, with an error carrying `retry: true`; only later calls + * see the replacement. + * + * Safe only for an owner whose every call stands alone, carrying what it needs in its arguments. + * A replacement process has none of the state a command sequence establishes: no SRS loaded over + * the connection, no Chonk accumulation, no batch-verifier session with its registered keys. An + * owner that runs any such sequence must leave this off, so a death fails loudly instead of the + * next call quietly running against a process that has forgotten everything. + */ + respawn?: boolean; + /** * @description Specify exact backend to use * - If unset: tries backends in default order with fallback diff --git a/barretenberg/ts/bb.js/src/bb_backends/node/index.ts b/barretenberg/ts/bb.js/src/bb_backends/node/index.ts index d02b4d2c76c..9ef091fb4e1 100644 --- a/barretenberg/ts/bb.js/src/bb_backends/node/index.ts +++ b/barretenberg/ts/bb.js/src/bb_backends/node/index.ts @@ -26,7 +26,13 @@ export async function createAsyncBackend( throw new Error('Native backend requires bb binary.'); } logger(`Using native Unix socket backend: ${bbPath}`); - return await BarretenbergNativeSocketAsyncBackend.new(bbPath, options.threads, options.logger, options.unref); + return await BarretenbergNativeSocketAsyncBackend.new( + bbPath, + options.threads, + options.logger, + options.unref, + options.respawn, + ); } case BackendType.NativeSharedMemory: { diff --git a/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.test.ts b/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.test.ts index c1c169968d6..f8e856f8de5 100644 --- a/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.test.ts +++ b/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.test.ts @@ -3,6 +3,7 @@ import * as fs from 'fs'; import * as os from 'os'; import * as path from 'path'; +import { isRetryable } from '../errors.js'; import { BarretenbergNativeSocketAsyncBackend } from './native_socket.js'; jest.setTimeout(30_000); @@ -43,6 +44,21 @@ function writeFakeBb(startupDelaySecs: number): string { return file; } +// A fake bb that writes its pid next to the script before serving, so a replacement process can +// be told apart from the one it replaced. +function writeFakeBbRecordingPid(): { path: string; pids: () => number[] } { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'fake-bb-')); + const serverJs = path.join(dir, 'echo_server.cjs'); + fs.writeFileSync(serverJs, ECHO_SERVER_JS); + const pidLog = path.join(dir, 'pids'); + const file = path.join(dir, 'bb'); + fs.writeFileSync(file, `#!/bin/bash\necho $$ >> ${pidLog}\nexec node ${serverJs} "$4"\n`, { mode: 0o755 }); + return { + path: file, + pids: () => fs.readFileSync(pidLog, 'utf-8').split('\n').filter(Boolean).map(Number), + }; +} + function writeFakeBbScript(script: string): string { const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'fake-bb-')); const file = path.join(dir, 'bb'); @@ -79,4 +95,136 @@ describe('BarretenbergNativeSocketAsyncBackend', () => { /Native backend process error/, ); }); + + describe('when bb dies', () => { + it('fails the call as retryable, and keeps failing without respawn', async () => { + const fake = writeFakeBbRecordingPid(); + const backend = await BarretenbergNativeSocketAsyncBackend.new(fake.path); + expect(await backend.call(new Uint8Array([1]))).toEqual(new Uint8Array([1])); + + process.kill(fake.pids()[0], 'SIGKILL'); + await waitUntil(() => !backend.isConnected()); + + const err = await backend.call(new Uint8Array([2])).catch(e => e); + expect(isRetryable(err)).toBe(true); + // No replacement: the same retryable failure, and no second process. + expect(isRetryable(await backend.call(new Uint8Array([3])).catch(e => e))).toBe(true); + expect(fake.pids()).toHaveLength(1); + await backend.destroy(); + }); + + it('serves the next call from a replacement process when respawn is on', async () => { + const fake = writeFakeBbRecordingPid(); + const backend = await BarretenbergNativeSocketAsyncBackend.new(fake.path, undefined, undefined, undefined, true); + expect(await backend.call(new Uint8Array([1]))).toEqual(new Uint8Array([1])); + + process.kill(fake.pids()[0], 'SIGKILL'); + await waitUntil(() => !backend.isConnected()); + + expect(await backend.call(new Uint8Array([2]))).toEqual(new Uint8Array([2])); + const pids = fake.pids(); + expect(pids).toHaveLength(2); + expect(pids[1]).not.toEqual(pids[0]); + await backend.destroy(); + }); + + it('starts one replacement however many calls find the connection down', async () => { + const fake = writeFakeBbRecordingPid(); + const backend = await BarretenbergNativeSocketAsyncBackend.new(fake.path, undefined, undefined, undefined, true); + await backend.call(new Uint8Array([1])); + + process.kill(fake.pids()[0], 'SIGKILL'); + await waitUntil(() => !backend.isConnected()); + + const results = await Promise.all([1, 2, 3, 4].map(n => backend.call(new Uint8Array([n])))); + expect(results).toEqual([1, 2, 3, 4].map(n => new Uint8Array([n]))); + expect(fake.pids()).toHaveLength(2); + await backend.destroy(); + }); + + it('does not leave a replacement running when destroyed while it starts', async () => { + const fake = writeFakeBbRecordingPid(); + const backend = await BarretenbergNativeSocketAsyncBackend.new(fake.path, undefined, undefined, undefined, true); + await backend.call(new Uint8Array([1])); + + process.kill(fake.pids()[0], 'SIGKILL'); + await waitUntil(() => !backend.isConnected()); + + const call = backend.call(new Uint8Array([2])).catch(e => e); + await backend.destroy(); + expect(String(await call)).toMatch(/Backend connection closed/); + + // destroy() does not wait for the replacement, so watch for it to arrive and then go. + await waitUntil(() => fake.pids().length === 2); + await waitUntil(() => !isProcessAlive(fake.pids()[1])); + }); + + it('fails retryably when the replacement cannot start either', async () => { + // The scenario the option exists for, gone wrong: bb is killed, and its replacement dies + // under the same pressure before it can connect. The caller must still be told to retry. + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'fake-bb-')); + const serverJs = path.join(dir, 'echo_server.cjs'); + fs.writeFileSync(serverJs, ECHO_SERVER_JS); + const pidLog = path.join(dir, 'pids'); + const bb = path.join(dir, 'bb'); + // Serves on the first start; every later start exits before creating its socket. + fs.writeFileSync( + bb, + `#!/bin/bash\nif [ -e ${pidLog} ]; then exit 9; fi\necho $$ >> ${pidLog}\nexec node ${serverJs} "$4"\n`, + { mode: 0o755 }, + ); + + const backend = await BarretenbergNativeSocketAsyncBackend.new(bb, undefined, undefined, undefined, true); + await backend.call(new Uint8Array([1])); + + process.kill(Number(fs.readFileSync(pidLog, 'utf-8').trim()), 'SIGKILL'); + await waitUntil(() => !backend.isConnected()); + + const err = await backend.call(new Uint8Array([2])).catch(e => e); + expect(isRetryable(err)).toBe(true); + expect(String(err)).toMatch(/exited before socket connection was established/); + await backend.destroy(); + }); + + it('leaves no bb running when the connection breaks but the process does not exit', async () => { + // A server that hangs up on the first request and keeps running, as bb does: its client is + // disconnected, its serve loop is not. + const dir = fs.mkdtempSync(path.join(os.tmpdir(), 'fake-bb-')); + const serverJs = path.join(dir, 'hangup_server.cjs'); + fs.writeFileSync( + serverJs, + `const net = require('net');\nnet.createServer(s => s.on('data', () => s.end())).listen(process.argv[2]);\nsetInterval(() => {}, 1000);\n`, + ); + const pidLog = path.join(dir, 'pids'); + const bb = path.join(dir, 'bb'); + fs.writeFileSync(bb, `#!/bin/bash\necho $$ >> ${pidLog}\nexec node ${serverJs} "$4"\n`, { mode: 0o755 }); + + const backend = await BarretenbergNativeSocketAsyncBackend.new(bb); + const pid = Number(fs.readFileSync(pidLog, 'utf-8').trim()); + await expect(backend.call(new Uint8Array([1]))).rejects.toThrow(); + + await waitUntil(() => !isProcessAlive(pid)); + await backend.destroy(); + }); + }); }); + +/** Poll until `predicate` holds, so a test never depends on when an event lands. */ +async function waitUntil(predicate: () => boolean, timeoutMs = 10_000): Promise { + const deadline = Date.now() + timeoutMs; + while (!predicate()) { + if (Date.now() > deadline) { + throw new Error('timed out waiting for condition'); + } + await new Promise(resolve => setTimeout(resolve, 10)); + } +} + +function isProcessAlive(pid: number): boolean { + try { + process.kill(pid, 0); + return true; + } catch { + return false; + } +} diff --git a/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.ts b/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.ts index c09bb7cfd22..17c22e81622 100644 --- a/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.ts +++ b/barretenberg/ts/bb.js/src/bb_backends/node/native_socket.ts @@ -7,6 +7,7 @@ import * as path from 'path'; import readline from 'readline'; import { threadId } from 'worker_threads'; +import { BackendUnavailableError } from '../errors.js'; import { IMsgpackBackendAsync } from '../interface.js'; let instanceCounter = 0; @@ -16,6 +17,29 @@ let instanceCounter = 0; // useful upper bound, so this must only ever fire when bb is genuinely stuck, never under load. const STARTUP_TIMEOUT_MS = 60_000; +/** What it takes to start a bb process, kept so a dead one can be replaced by an identical one. */ +interface SpawnOptions { + bbBinaryPath: string; + threads?: number; + logger?: (msg: string) => void; + unref?: boolean; + /** + * Replace the bb process when it dies, on the next call. Only safe where a bb process holds no + * state between calls: a replacement has no Chonk accumulation and no batch-verifier session. + * Off by default, so a caller that does hold such state keeps failing loudly instead of + * silently continuing against a fresh process. + */ + respawn?: boolean; +} + +/** A bb process and the connection to it. Replaced as a unit. */ +interface Incarnation { + proc: ChildProcess; + socket: net.Socket; + /** Kept so the path can be removed when the process is gone without having unlinked it. */ + socketPath: string; +} + /** * Asynchronous native backend that communicates with bb binary via Unix Domain Socket. * Uses event-based I/O with a state machine to handle partial reads. @@ -46,25 +70,45 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn private responseBuffer: Buffer | null = null; private responseBytesRead: number = 0; - private constructor( - private process: ChildProcess, - socket: net.Socket, - private logger: (msg: string) => void, - ) { - this.socket = socket; + private proc: ChildProcess | null = null; + /** Set when the process died; cleared when a replacement is adopted. */ + private death: Error | null = null; + /** Shared by every caller that finds the connection down, so one death causes one replacement. */ + private starting: Promise | null = null; + private destroyed = false; - this.process.on('error', err => { - this.failAllPending(new Error(`Native backend process error: ${err.message}`)); + private constructor(private opts: SpawnOptions) { + this.socket = null; + this.logger = opts.logger ?? (() => {}); + } + + private logger: (msg: string) => void; + + /** Take ownership of a started process and its connection, and watch for its death. */ + private adopt(incarnation: Incarnation): void { + this.proc = incarnation.proc; + this.socket = incarnation.socket; + this.death = null; + this.readingLength = true; + this.lengthBytesRead = 0; + this.responseBuffer = null; + this.responseBytesRead = 0; + + const { proc, socket } = incarnation; + // Every listener is scoped to this incarnation and removed with it, so a death observed after + // a replacement was adopted cannot tear the replacement down. + proc.on('error', err => { + this.onDeath(incarnation, `Native backend process error: ${err.message}`); }); - this.process.on('exit', (code, signal) => { - const errorMsg = + proc.on('exit', (code, signal) => { + const reason = code !== null && code !== 0 ? `Native backend process exited with code ${code}` : signal && signal !== 'SIGTERM' ? `Native backend process killed with signal ${signal}` : 'Native backend process exited unexpectedly'; - this.failAllPending(new Error(errorMsg)); + this.onDeath(incarnation, reason); }); socket.on('data', (chunk: Buffer) => { @@ -72,14 +116,32 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn }); socket.on('error', err => { - this.failAllPending(new Error(`Socket error: ${err.message}`)); + this.onDeath(incarnation, `Socket error: ${err.message}`); }); socket.on('end', () => { - this.failAllPending(new Error('Socket connection ended unexpectedly')); + this.onDeath(incarnation, 'Socket connection ended unexpectedly'); }); } + /** + * This backend has lost the bb process it was talking to. In-flight calls fail as retryable, and + * the connection is dropped so the next call either starts a replacement or reports the death. + * + * Losing the connection does not mean the process is gone: bb's server keeps serving after a + * client disconnects, so it is killed here rather than left to outlive the backend that spawned + * it. Killing a process that has already exited is a no-op. + */ + private onDeath(incarnation: Incarnation, reason: string): void { + if (this.proc !== incarnation.proc) { + return; // A later incarnation is already in charge. + } + this.death = new BackendUnavailableError(reason); + this.proc = null; + retire(incarnation); + this.failAllPending(this.death); + } + /** * Spawn a bb process and wait until a socket connection to it is established. * Waits as long as the bb process is alive (bb startup has no useful upper bound on a loaded @@ -91,7 +153,17 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn threads?: number, logger?: (msg: string) => void, unref?: boolean, + respawn?: boolean, ): Promise { + const opts: SpawnOptions = { bbBinaryPath, threads, logger, unref, respawn }; + const backend = new BarretenbergNativeSocketAsyncBackend(opts); + backend.adopt(await this.start(opts)); + return backend; + } + + /** Start a bb process and connect to it. The whole of what a replacement has to repeat. */ + private static async start(opts: SpawnOptions): Promise { + const { bbBinaryPath, threads, logger, unref } = opts; // Create a unique socket path in temp directory const socketPath = path.join(os.tmpdir(), `bb-${process.pid}-${threadId}-${instanceCounter++}.sock`); @@ -132,15 +204,21 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn try { await once(proc, 'spawn'); } catch (err) { - throw new Error(`Native backend process error: ${(err as Error).message}`); + // A missing or non-executable binary is a configuration fault: retrying cannot fix it. + const code = (err as NodeJS.ErrnoException).code; + const message = `Native backend process error: ${(err as Error).message}`; + throw code === 'ENOENT' || code === 'EACCES' ? new Error(message) : new BackendUnavailableError(message); } try { const socket = await this.waitForSocketAndConnect(socketPath, proc); - return new BarretenbergNativeSocketAsyncBackend(proc, socket, logger ?? (() => {})); + return { proc, socket, socketPath }; } catch (err) { proc.kill('SIGKILL'); - throw err; + cleanUpSocketPath(socketPath); + // A bb that died starting up, never accepted a connection, or refused one failed for + // environmental reasons — a loaded machine, memory pressure — so it is worth another try. + throw new BackendUnavailableError((err as Error).message, { cause: err }); } } @@ -207,6 +285,11 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn } } + /** Whether a call can be made without first starting a replacement bb process. */ + isConnected(): boolean { + return this.socket !== null; + } + private handleData(chunk: Buffer): void { let offset = 0; @@ -258,15 +341,13 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn } } - call(inputBuffer: Uint8Array): Promise { - if (!this.socket) { - return Promise.reject(new Error('Socket not connected')); - } + async call(inputBuffer: Uint8Array): Promise { + const socket = await this.ensureConnected(); return new Promise((resolve, reject) => { // If this is the first pending callback, ref the socket to keep event loop alive if (this.pendingCallbacks.length === 0) { - this.socket!.ref(); + socket.ref(); } // Enqueue this promise's callbacks (FIFO order) @@ -276,16 +357,88 @@ export class BarretenbergNativeSocketAsyncBackend implements IMsgpackBackendAsyn // Socket will buffer these if needed, maintaining order const lengthBuf = Buffer.alloc(4); lengthBuf.writeUInt32LE(inputBuffer.length, 0); - this.socket!.write(lengthBuf); - this.socket!.write(inputBuffer); + socket.write(lengthBuf); + socket.write(inputBuffer); }); } + /** + * The connection to use for the next call, replacing a dead bb process when that is allowed. + * Concurrent callers share one replacement, so a single death costs a single process. + */ + private async ensureConnected(): Promise { + if (this.socket) { + return this.socket; + } + if (this.destroyed) { + throw new Error('Backend connection closed'); + } + if (!this.opts.respawn) { + // A fresh error each time: callers attach to what they are given, and one shared instance + // would carry one caller's stack and annotations to every other. + const death = this.death; + throw death + ? new BackendUnavailableError(death.message, { cause: death }) + : new BackendUnavailableError('Socket not connected'); + } + this.starting ??= this.replace().finally(() => { + this.starting = null; + }); + await this.starting; + // destroy() may have run while the replacement was starting. + if (!this.socket) { + throw new Error('Backend connection closed'); + } + return this.socket; + } + + private async replace(): Promise { + this.logger('bb process died; starting a replacement'); + const incarnation = await BarretenbergNativeSocketAsyncBackend.start(this.opts); + if (this.destroyed) { + // destroy() ran while this was starting: the owner is gone, so neither is this process. + retire(incarnation); + throw new Error('Backend connection closed'); + } + this.adopt(incarnation); + } + destroy(): Promise { + this.destroyed = true; this.failAllPending(new Error('Backend connection closed')); - // Don't try to unlink socket - bb owns it and will clean it up - this.process.kill('SIGTERM'); - this.process.removeAllListeners(); + // A replacement still starting is not waited for: replace() kills whatever it produces once it + // sees the backend destroyed. Waiting here would hold up shutdown for as long as a bb can take + // to come up, which has no useful upper bound. + const proc = this.proc; + this.proc = null; + // bb unlinks its own socket path when it shuts down cleanly, which SIGTERM gives it a chance to. + proc?.kill('SIGTERM'); + proc?.removeAllListeners(); return Promise.resolve(); } } + +/** + * Stop watching an incarnation, drop its connection and kill its process. + * + * A bb killed this way never gets to unlink its socket path, so this does it instead; otherwise a + * long-lived backend that replaces its process leaves one file in the temp directory per death. + */ +function retire(incarnation: Incarnation): void { + incarnation.proc.removeAllListeners(); + incarnation.socket.removeAllListeners(); + incarnation.socket.destroy(); + incarnation.proc.kill('SIGKILL'); + const socketPath = incarnation.socketPath; + incarnation.proc.once('exit', () => cleanUpSocketPath(socketPath)); + cleanUpSocketPath(socketPath); +} + +/** Remove a socket path a bb process is no longer listening on. */ +function cleanUpSocketPath(socketPath: string): void { + try { + fs.unlinkSync(socketPath); + } catch { + // Already gone, which is the usual case: bb unlinks its own path when it exits cleanly. + } +}