77 * CAR bytes and persisted back into tracking.db.
88 */
99
10- import { createReadStream } from 'node:fs'
1110import fs from 'node:fs/promises'
11+ import { fileURLToPath } from 'node:url'
12+ import { Worker } from 'node:worker_threads'
1213
13- import { calculateFromIterable } from '@filoz/synapse-core/piece'
1414import pMap from 'p-map'
1515
1616import { pathExists , renderProgressLine } from '../../utils.js'
@@ -28,15 +28,154 @@ const PREPARE_BATCH_SIZE = 1_000
2828 */
2929
3030/**
31- * @param {string } carPath
31+ * Pool of worker threads that each compute a piece CID from a CAR file, so the
32+ * CPU-bound (pure-JS, single-threaded) CommP hashing runs in parallel across
33+ * cores. Each worker calls the same `@filoz/synapse-core/piece` hash as before,
34+ * so output piece CIDs are identical — only throughput changes. DB writes and
35+ * file renames stay on the main thread (sqlite is not shared with workers).
3236 */
33- async function calculateLocalPieceCid ( carPath ) {
34- const stream = createReadStream ( carPath )
35- try {
36- const pieceCid = await calculateFromIterable ( stream )
37- return pieceCid . toString ( )
38- } finally {
39- stream . destroy ( )
37+ class PieceCidWorkerPool {
38+ /** @param {number } size */
39+ constructor ( size ) {
40+ this . workerPath = fileURLToPath ( new URL ( '../lib/piece-cid-worker.mjs' , import . meta. url ) )
41+ /** @type {import('node:worker_threads').Worker[] } */
42+ this . workers = [ ]
43+ /** @type {import('node:worker_threads').Worker[] } */
44+ this . idle = [ ]
45+ /** @type {Array<{carPath: string, resolve: (v: string) => void, reject: (e: Error) => void}> } */
46+ this . queue = [ ]
47+ /** @type {Map<import('node:worker_threads').Worker, {resolve: (v: string) => void, reject: (e: Error) => void}> } */
48+ this . busy = new Map ( )
49+ /** @type {Error | null } */
50+ this . fatalError = null
51+ this . closing = false
52+ /** @type {Set<import('node:worker_threads').Worker> } */
53+ this . failedWorkers = new Set ( )
54+ for ( let i = 0 ; i < Math . max ( 1 , size ) ; i ++ ) {
55+ const worker = new Worker ( this . workerPath )
56+ worker . on ( 'message' , ( msg ) => this . _onMessage ( worker , msg ) )
57+ worker . on ( 'error' , ( err ) => this . _onError ( worker , err ) )
58+ worker . on ( 'exit' , ( code ) => this . _onExit ( worker , code ) )
59+ this . workers . push ( worker )
60+ this . idle . push ( worker )
61+ }
62+ }
63+
64+ /**
65+ * @param {string } carPath
66+ * @returns {Promise<string> }
67+ */
68+ compute ( carPath ) {
69+ return new Promise ( ( resolve , reject ) => {
70+ if ( this . fatalError ) {
71+ reject ( this . fatalError )
72+ return
73+ }
74+
75+ const worker = this . idle . pop ( )
76+ if ( worker ) this . _assign ( worker , { carPath, resolve, reject } )
77+ else this . queue . push ( { carPath, resolve, reject } )
78+ } )
79+ }
80+
81+ /**
82+ * @param {import('node:worker_threads').Worker } worker
83+ * @param {{carPath: string, resolve: (v: string) => void, reject: (e: Error) => void} } job
84+ */
85+ _assign ( worker , job ) {
86+ this . busy . set ( worker , { resolve : job . resolve , reject : job . reject } )
87+ worker . postMessage ( { carPath : job . carPath } )
88+ }
89+
90+ /**
91+ * @param {import('node:worker_threads').Worker } worker
92+ * @param {{pieceCid?: string, error?: string} } msg
93+ */
94+ _onMessage ( worker , msg ) {
95+ const job = this . busy . get ( worker )
96+ this . busy . delete ( worker )
97+ if ( job ) {
98+ if ( msg . error ) job . reject ( new Error ( msg . error ) )
99+ else job . resolve ( /** @type {string } */ ( msg . pieceCid ) )
100+ }
101+ const next = this . queue . shift ( )
102+ if ( next ) this . _assign ( worker , next )
103+ else this . idle . push ( worker )
104+ }
105+
106+ /**
107+ * @param {import('node:worker_threads').Worker } worker
108+ * @param {Error } err
109+ */
110+ _onError ( worker , err ) {
111+ this . _handleWorkerFailure ( worker , err , `prepare: piece CID worker error: ${ err ?. message || err } ` )
112+ }
113+
114+ /**
115+ * @param {import('node:worker_threads').Worker } worker
116+ * @param {number } code
117+ */
118+ _onExit ( worker , code ) {
119+ if ( this . closing ) return
120+ if ( code !== 0 ) {
121+ const err = new Error ( `piece CID worker exited with code ${ code } ` )
122+ this . _handleWorkerFailure ( worker , err , `prepare: ${ err . message } ` )
123+ return
124+ }
125+
126+ if ( this . busy . has ( worker ) ) {
127+ const err = new Error ( 'piece CID worker exited unexpectedly while processing a job' )
128+ this . _handleWorkerFailure ( worker , err , `prepare: ${ err . message } ` )
129+ }
130+ }
131+
132+ /**
133+ * @param {import('node:worker_threads').Worker } worker
134+ */
135+ _removeWorker ( worker ) {
136+ this . workers = this . workers . filter ( ( candidate ) => candidate !== worker )
137+ this . idle = this . idle . filter ( ( candidate ) => candidate !== worker )
138+ }
139+
140+ /**
141+ * @param {import('node:worker_threads').Worker } worker
142+ * @param {Error } err
143+ * @param {string } logMessage
144+ */
145+ _handleWorkerFailure ( worker , err , logMessage ) {
146+ if ( this . failedWorkers . has ( worker ) ) return
147+ this . failedWorkers . add ( worker )
148+
149+ const job = this . busy . get ( worker )
150+ this . busy . delete ( worker )
151+ if ( job ) job . reject ( err )
152+ this . _removeWorker ( worker )
153+ console . error ( logMessage )
154+ this . _drainOrFailQueue ( )
155+ }
156+
157+ _drainOrFailQueue ( ) {
158+ while ( this . idle . length > 0 && this . queue . length > 0 ) {
159+ const worker = this . idle . pop ( )
160+ const next = this . queue . shift ( )
161+ if ( ! worker || ! next ) break
162+ this . _assign ( worker , next )
163+ }
164+
165+ if ( this . workers . length > 0 ) return
166+
167+ this . fatalError = new Error ( 'all piece CID workers failed' )
168+ console . error ( `prepare: ${ this . fatalError . message } ` )
169+ while ( this . queue . length > 0 ) {
170+ const next = this . queue . shift ( )
171+ if ( ! next ) break
172+ next . reject ( this . fatalError )
173+ }
174+ }
175+
176+ async close ( ) {
177+ this . closing = true
178+ await Promise . all ( this . workers . map ( ( worker ) => worker . terminate ( ) ) )
40179 }
41180}
42181
@@ -127,6 +266,8 @@ function renderPrepareProgress(summary) {
127266export async function runPrepare ( { dir, concurrency } ) {
128267 const tracking = openTrackingDb ( dir )
129268 const workerConcurrency = concurrency ?? DEFAULT_PREPARE_CONCURRENCY
269+ /** @type {PieceCidWorkerPool | null } */
270+ let pool = null
130271
131272 try {
132273 const total = tracking . getDownloadStats ( ) . complete
@@ -135,6 +276,8 @@ export async function runPrepare({ dir, concurrency }) {
135276 return
136277 }
137278
279+ pool = new PieceCidWorkerPool ( workerConcurrency )
280+
138281 const summary = {
139282 total,
140283 done : 0 ,
@@ -171,7 +314,8 @@ export async function runPrepare({ dir, concurrency }) {
171314
172315 let pieceCid = candidate . pieceCid
173316 if ( ! pieceCid ) {
174- pieceCid = await calculateLocalPieceCid ( workItem . carPath )
317+ if ( ! pool ) throw new Error ( 'prepare: piece CID worker pool was not initialized' )
318+ pieceCid = await pool . compute ( workItem . carPath )
175319 tracking . setPieceCid ( workItem . shardCid , pieceCid )
176320 computedPieceCid = true
177321 } else {
@@ -211,6 +355,7 @@ export async function runPrepare({ dir, concurrency }) {
211355 `prepare: done. total=${ summary . total } done=${ summary . done } /${ summary . total } computed=${ summary . computed } failed=${ summary . failed } ` ,
212356 )
213357 } finally {
358+ if ( pool ) await pool . close ( )
214359 tracking . close ( )
215360 }
216361}
0 commit comments