From 87ee7644755ba27144f3de3c501d85fa1137351f Mon Sep 17 00:00:00 2001 From: TippyFlits Date: Thu, 4 Jun 2026 20:14:39 +0100 Subject: [PATCH 1/2] perf(backup-helper): compute prepare CommP across a worker_threads pool prepare computes a piece CID for every shard missing one via the pure-JS @filoz/synapse-core/piece hash. That hash is CPU-bound and single-threaded, so on a multi-core node prepare pins one core and leaves the rest idle. Move the hashing into a worker_threads pool sized to --concurrency. Each worker calls the same calculateFromIterable, so output piece CIDs are byte-identical to the single-thread path; only throughput changes. DB writes and shard->piece renames stay on the main thread (sqlite is not shared with workers). Measured on a 64-core machine (~10 MB avg CARs): ~20 MB/s at -c 8, ~69 MB/s at -c 24, ~84 MB/s at -c 32 (knee ~32 on that box). No change to default concurrency or any existing flag. (cherry picked from commit 1de3e21824afc88cf466508378154de25d04d00d) --- scripts/backup-helper/commands/prepare.mjs | 92 ++++++++++++++++--- .../backup-helper/lib/piece-cid-worker.mjs | 29 ++++++ 2 files changed, 110 insertions(+), 11 deletions(-) create mode 100644 scripts/backup-helper/lib/piece-cid-worker.mjs diff --git a/scripts/backup-helper/commands/prepare.mjs b/scripts/backup-helper/commands/prepare.mjs index 04f2428..0e7db76 100644 --- a/scripts/backup-helper/commands/prepare.mjs +++ b/scripts/backup-helper/commands/prepare.mjs @@ -7,10 +7,10 @@ * CAR bytes and persisted back into tracking.db. */ -import { createReadStream } from 'node:fs' import fs from 'node:fs/promises' +import { fileURLToPath } from 'node:url' +import { Worker } from 'node:worker_threads' -import { calculateFromIterable } from '@filoz/synapse-core/piece' import pMap from 'p-map' import { pathExists, renderProgressLine } from '../../utils.js' @@ -28,15 +28,83 @@ const PREPARE_BATCH_SIZE = 1_000 */ /** - * @param {string} carPath + * Pool of worker threads that each compute a piece CID from a CAR file, so the + * CPU-bound (pure-JS, single-threaded) CommP hashing runs in parallel across + * cores. Each worker calls the same `@filoz/synapse-core/piece` hash as before, + * so output piece CIDs are identical — only throughput changes. DB writes and + * file renames stay on the main thread (sqlite is not shared with workers). */ -async function calculateLocalPieceCid(carPath) { - const stream = createReadStream(carPath) - try { - const pieceCid = await calculateFromIterable(stream) - return pieceCid.toString() - } finally { - stream.destroy() +class PieceCidWorkerPool { + /** @param {number} size */ + constructor(size) { + this.workerPath = fileURLToPath(new URL('../lib/piece-cid-worker.mjs', import.meta.url)) + /** @type {import('node:worker_threads').Worker[]} */ + this.workers = [] + /** @type {import('node:worker_threads').Worker[]} */ + this.idle = [] + /** @type {Array<{carPath: string, resolve: (v: string) => void, reject: (e: Error) => void}>} */ + this.queue = [] + /** @type {Map void, reject: (e: Error) => void}>} */ + this.busy = new Map() + for (let i = 0; i < Math.max(1, size); i++) { + const worker = new Worker(this.workerPath) + worker.on('message', (msg) => this._onMessage(worker, msg)) + worker.on('error', (err) => this._onError(worker, err)) + this.workers.push(worker) + this.idle.push(worker) + } + } + + /** + * @param {string} carPath + * @returns {Promise} + */ + compute(carPath) { + return new Promise((resolve, reject) => { + const worker = this.idle.pop() + if (worker) this._assign(worker, { carPath, resolve, reject }) + else this.queue.push({ carPath, resolve, reject }) + }) + } + + /** + * @param {import('node:worker_threads').Worker} worker + * @param {{carPath: string, resolve: (v: string) => void, reject: (e: Error) => void}} job + */ + _assign(worker, job) { + this.busy.set(worker, { resolve: job.resolve, reject: job.reject }) + worker.postMessage({ carPath: job.carPath }) + } + + /** + * @param {import('node:worker_threads').Worker} worker + * @param {{pieceCid?: string, error?: string}} msg + */ + _onMessage(worker, msg) { + const job = this.busy.get(worker) + this.busy.delete(worker) + if (job) { + if (msg.error) job.reject(new Error(msg.error)) + else job.resolve(/** @type {string} */ (msg.pieceCid)) + } + const next = this.queue.shift() + if (next) this._assign(worker, next) + else this.idle.push(worker) + } + + /** + * @param {import('node:worker_threads').Worker} worker + * @param {Error} err + */ + _onError(worker, err) { + const job = this.busy.get(worker) + this.busy.delete(worker) + if (job) job.reject(err) + // a crashed worker leaves the pool; remaining workers continue + } + + async close() { + await Promise.all(this.workers.map((worker) => worker.terminate())) } } @@ -127,6 +195,7 @@ function renderPrepareProgress(summary) { export async function runPrepare({ dir, concurrency }) { const tracking = openTrackingDb(dir) const workerConcurrency = concurrency ?? DEFAULT_PREPARE_CONCURRENCY + const pool = new PieceCidWorkerPool(workerConcurrency) try { const total = tracking.getDownloadStats().complete @@ -171,7 +240,7 @@ export async function runPrepare({ dir, concurrency }) { let pieceCid = candidate.pieceCid if (!pieceCid) { - pieceCid = await calculateLocalPieceCid(workItem.carPath) + pieceCid = await pool.compute(workItem.carPath) tracking.setPieceCid(workItem.shardCid, pieceCid) computedPieceCid = true } else { @@ -211,6 +280,7 @@ export async function runPrepare({ dir, concurrency }) { `prepare: done. total=${summary.total} done=${summary.done}/${summary.total} computed=${summary.computed} failed=${summary.failed}`, ) } finally { + await pool.close() tracking.close() } } diff --git a/scripts/backup-helper/lib/piece-cid-worker.mjs b/scripts/backup-helper/lib/piece-cid-worker.mjs new file mode 100644 index 0000000..8fd3e36 --- /dev/null +++ b/scripts/backup-helper/lib/piece-cid-worker.mjs @@ -0,0 +1,29 @@ +/** + * Worker thread for `backup-helper prepare`: computes a piece CID from a local + * CAR file using the same `@filoz/synapse-core/piece` hash as the main thread, + * so the CPU-bound (pure-JS) CommP hashing can run in parallel across cores. + * + * Local addition (not upstream) — parallelises prepare's hashing. The hash + * itself is unchanged, so output piece CIDs are identical to the single-thread + * path; only the throughput differs. + */ +import { createReadStream } from 'node:fs' +import { parentPort } from 'node:worker_threads' + +import { calculateFromIterable } from '@filoz/synapse-core/piece' + +if (!parentPort) { + throw new Error('piece-cid-worker.mjs must be run as a worker thread') +} + +parentPort.on('message', async ({ carPath }) => { + const stream = createReadStream(carPath) + try { + const pieceCid = await calculateFromIterable(stream) + parentPort.postMessage({ pieceCid: pieceCid.toString() }) + } catch (err) { + parentPort.postMessage({ error: String(err?.message || err) }) + } finally { + stream.destroy() + } +}) From 01f40fd54f850879f828790c9527512ec2f72753 Mon Sep 17 00:00:00 2001 From: bravonatalie Date: Fri, 5 Jun 2026 09:09:31 -0300 Subject: [PATCH 2/2] fix(backup-helper): harden prepare worker pool failure handling --- scripts/backup-helper/commands/prepare.mjs | 81 +++++++++++++++++++++- 1 file changed, 78 insertions(+), 3 deletions(-) diff --git a/scripts/backup-helper/commands/prepare.mjs b/scripts/backup-helper/commands/prepare.mjs index 0e7db76..6c1ccd5 100644 --- a/scripts/backup-helper/commands/prepare.mjs +++ b/scripts/backup-helper/commands/prepare.mjs @@ -46,10 +46,16 @@ class PieceCidWorkerPool { this.queue = [] /** @type {Map void, reject: (e: Error) => void}>} */ this.busy = new Map() + /** @type {Error | null} */ + this.fatalError = null + this.closing = false + /** @type {Set} */ + this.failedWorkers = new Set() for (let i = 0; i < Math.max(1, size); i++) { const worker = new Worker(this.workerPath) worker.on('message', (msg) => this._onMessage(worker, msg)) worker.on('error', (err) => this._onError(worker, err)) + worker.on('exit', (code) => this._onExit(worker, code)) this.workers.push(worker) this.idle.push(worker) } @@ -61,6 +67,11 @@ class PieceCidWorkerPool { */ compute(carPath) { return new Promise((resolve, reject) => { + if (this.fatalError) { + reject(this.fatalError) + return + } + const worker = this.idle.pop() if (worker) this._assign(worker, { carPath, resolve, reject }) else this.queue.push({ carPath, resolve, reject }) @@ -97,13 +108,73 @@ class PieceCidWorkerPool { * @param {Error} err */ _onError(worker, err) { + this._handleWorkerFailure(worker, err, `prepare: piece CID worker error: ${err?.message || err}`) + } + + /** + * @param {import('node:worker_threads').Worker} worker + * @param {number} code + */ + _onExit(worker, code) { + if (this.closing) return + if (code !== 0) { + const err = new Error(`piece CID worker exited with code ${code}`) + this._handleWorkerFailure(worker, err, `prepare: ${err.message}`) + return + } + + if (this.busy.has(worker)) { + const err = new Error('piece CID worker exited unexpectedly while processing a job') + this._handleWorkerFailure(worker, err, `prepare: ${err.message}`) + } + } + + /** + * @param {import('node:worker_threads').Worker} worker + */ + _removeWorker(worker) { + this.workers = this.workers.filter((candidate) => candidate !== worker) + this.idle = this.idle.filter((candidate) => candidate !== worker) + } + + /** + * @param {import('node:worker_threads').Worker} worker + * @param {Error} err + * @param {string} logMessage + */ + _handleWorkerFailure(worker, err, logMessage) { + if (this.failedWorkers.has(worker)) return + this.failedWorkers.add(worker) + const job = this.busy.get(worker) this.busy.delete(worker) if (job) job.reject(err) - // a crashed worker leaves the pool; remaining workers continue + this._removeWorker(worker) + console.error(logMessage) + this._drainOrFailQueue() + } + + _drainOrFailQueue() { + while (this.idle.length > 0 && this.queue.length > 0) { + const worker = this.idle.pop() + const next = this.queue.shift() + if (!worker || !next) break + this._assign(worker, next) + } + + if (this.workers.length > 0) return + + this.fatalError = new Error('all piece CID workers failed') + console.error(`prepare: ${this.fatalError.message}`) + while (this.queue.length > 0) { + const next = this.queue.shift() + if (!next) break + next.reject(this.fatalError) + } } async close() { + this.closing = true await Promise.all(this.workers.map((worker) => worker.terminate())) } } @@ -195,7 +266,8 @@ function renderPrepareProgress(summary) { export async function runPrepare({ dir, concurrency }) { const tracking = openTrackingDb(dir) const workerConcurrency = concurrency ?? DEFAULT_PREPARE_CONCURRENCY - const pool = new PieceCidWorkerPool(workerConcurrency) + /** @type {PieceCidWorkerPool | null} */ + let pool = null try { const total = tracking.getDownloadStats().complete @@ -204,6 +276,8 @@ export async function runPrepare({ dir, concurrency }) { return } + pool = new PieceCidWorkerPool(workerConcurrency) + const summary = { total, done: 0, @@ -240,6 +314,7 @@ export async function runPrepare({ dir, concurrency }) { let pieceCid = candidate.pieceCid if (!pieceCid) { + if (!pool) throw new Error('prepare: piece CID worker pool was not initialized') pieceCid = await pool.compute(workItem.carPath) tracking.setPieceCid(workItem.shardCid, pieceCid) computedPieceCid = true @@ -280,7 +355,7 @@ export async function runPrepare({ dir, concurrency }) { `prepare: done. total=${summary.total} done=${summary.done}/${summary.total} computed=${summary.computed} failed=${summary.failed}`, ) } finally { - await pool.close() + if (pool) await pool.close() tracking.close() } }