Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<const uint8_t> response;
while ((response = client_->receive(TIMEOUT_NS)).empty()) {
// Response not ready yet, server is processing - retry
Expand Down
12 changes: 10 additions & 2 deletions barretenberg/ts/bb.js/src/barretenberg/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
30 changes: 30 additions & 0 deletions barretenberg/ts/bb.js/src/barretenberg/singleton.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
29 changes: 29 additions & 0 deletions barretenberg/ts/bb.js/src/bb_backends/errors.ts
Original file line number Diff line number Diff line change
@@ -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;
}
19 changes: 19 additions & 0 deletions barretenberg/ts/bb.js/src/bb_backends/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 7 additions & 1 deletion barretenberg/ts/bb.js/src/bb_backends/node/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: {
Expand Down
148 changes: 148 additions & 0 deletions barretenberg/ts/bb.js/src/bb_backends/node/native_socket.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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');
Expand Down Expand Up @@ -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<void> {
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;
}
}
Loading
Loading