From c4f945f4700fa95148496e17c5aa2c4f96c05f41 Mon Sep 17 00:00:00 2001 From: Charlie <5764343+charlielye@users.noreply.github.com> Date: Thu, 1 Oct 2026 10:55:04 +0100 Subject: [PATCH] fix(bb.js): report a dead bb process as retryable, and replace it on request (#25548) Addresses https://github.com/AztecProtocol/barretenberg-claude/issues/4427 from the bb.js side. The aztec-node side is [aztec-node#354](https://github.com/aztec-labs-eng/aztec-node/pull/354). A pooled bb verifier whose process dies is handed back to the pool and handed out again on every later borrow. The pool cannot tell, because the socket backend marks itself permanently unusable and every later call fails with a bare `Socket not connected`. The node therefore reports a dead helper process as an invalid transaction proof, persistently, and books it into the metric that means users are submitting bad proofs. This gives an owner the two things it needs, in the shape the AVM simulator pool already relies on. ### A failed call says whether retrying can help An environmental failure now carries `retry: true`: the bb process died, its connection broke, or it could not be started. Only a binary that cannot be executed stays non-retryable, because retrying cannot fix it. The bare property is the whole contract, feature-detected rather than imported: ```ts if (err instanceof Error && (err as Error & { retry?: unknown }).retry === true) { ... } ``` That is deliberately the same convention `ipc-runtime`'s transports use, so a caller written against one works unchanged against the other. It is also what lets a verifier distinguish a bad proof from a dead helper, which is the misattribution half of the issue. ### The socket backend can replace a dead process, on request `respawn` is a new `BackendOptions` flag, off by default. With it on, the bb process and its connection are one unit swapped as a whole: the next call starts a replacement, while calls already in flight still fail retryably. Concurrent callers that find the connection down share one replacement, so a single death costs a single process, and a replacement that arrives after `destroy()` is not left running. It is opt-in because a replacement 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 it off, so a death fails loudly instead of the next call quietly running against a process that has forgotten everything. The verifier pool qualifies because every verification carries its proof and key in the call. With this, a pool needs no liveness check and no maintenance loop: returning an instance unconditionally becomes correct, exactly as it already is for the AVM pool. ### Also `BarretenbergSync.initSingleton()` cached a failed initialization for the life of the process, so one bad spawn was permanent. It now clears the failure, as the asynchronous singleton already did. ### Testing `native_socket.test.ts`, against the fake bb the existing tests use: a killed bb fails the call retryably and keeps failing without the option; with it on the next call is served by a replacement process, a different pid; four concurrent calls that find the connection down start exactly one replacement; `destroy()` during a replacement leaves no process running; a replacement that cannot start fails retryably too; a connection that breaks while the process keeps running leaves no bb behind. `singleton.test.ts` covers the initialization fix and fails without it. Checked against a real bb as well: with the option off a killed bb gives a retryable error and keeps doing so, and with it on the next call transparently returns the same hash from a fresh process. ### What this does not cover `BarretenbergSync` runs on the shared-memory backend, which has no respawn and no way to report a death mid-call: the NAPI receive loop retries without a deadline, and because the call blocks the event loop the process exit is never even observed, so the caller wedges rather than fails. That path is untouched here and is deliberately left alone: the synchronous bb runs one thread doing hashes and signatures, so it is the least likely process on the machine to be killed, and the fix would mean threading a liveness check into the C++ client for a case that may never happen. #25546 covers the idle-death half of it by replacing a dead singleton, so the two PRs cover different backends rather than one superseding the other. ### Note on direction bb.js's hand-written backends are replaced by `ipc-runtime`'s in the codegen migration (#25362). Nothing above is lost in that move: `ipc-runtime`'s spawned backend already carries the same `retry` contract and the same opt-in respawn, so the callers written against this keep working and the implementation here is deleted. That is why this is expressed as the retry contract rather than as a liveness query, which would have to become part of the generated client's interface. --------- Co-authored-by: Claude Opus 5 (1M context) (cherry picked from commit d5b4b8ae96be70eabd6e9bc058737aebff9f70db) --- .../msgpack_client/msgpack_client_wrapper.cpp | 16 +- .../ts/bb.js/src/barretenberg/index.ts | 12 +- .../bb.js/src/barretenberg/singleton.test.ts | 30 +++ .../ts/bb.js/src/bb_backends/errors.ts | 29 +++ .../ts/bb.js/src/bb_backends/index.ts | 19 ++ .../ts/bb.js/src/bb_backends/node/index.ts | 8 +- .../bb_backends/node/native_socket.test.ts | 148 +++++++++++++ .../src/bb_backends/node/native_socket.ts | 205 +++++++++++++++--- 8 files changed, 437 insertions(+), 30 deletions(-) create mode 100644 barretenberg/ts/bb.js/src/barretenberg/singleton.test.ts create mode 100644 barretenberg/ts/bb.js/src/bb_backends/errors.ts 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 b72114a00abf..39606d309d80 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 217d2e5ff0f1..e79c5ccc6f8f 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 000000000000..8e68c4691360 --- /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 000000000000..1c5bf976b9dd --- /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 760f51fe9da0..a538a9b614aa 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 d02b4d2c76c5..9ef091fb4e1e 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 c1c169968d65..f8e856f8de5f 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 c09bb7cfd224..17c22e816223 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. + } +}