index.js 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823
  1. "use strict";
  2. import {
  3. __dirname,
  4. __privateAdd,
  5. __privateGet,
  6. __privateSet,
  7. __publicField,
  8. isMovable,
  9. isTaskQueue,
  10. isTransferable,
  11. kFieldCount,
  12. kQueueOptions,
  13. kRequestCountField,
  14. kResponseCountField,
  15. kTransferable,
  16. kValue,
  17. markMovable
  18. } from "./chunk-QYFJIXNO.js";
  19. // src/index.ts
  20. import {
  21. Worker,
  22. MessageChannel,
  23. receiveMessageOnPort
  24. } from "worker_threads";
  25. import { once } from "events";
  26. // src/EventEmitterAsyncResource.ts
  27. import { EventEmitter } from "events";
  28. import { AsyncResource } from "async_hooks";
  29. var kEventEmitter = Symbol("kEventEmitter");
  30. var kAsyncResource = Symbol("kAsyncResource");
  31. var _a;
  32. var EventEmitterReferencingAsyncResource = class extends AsyncResource {
  33. constructor(ee, type, options) {
  34. super(type, options);
  35. __publicField(this, _a);
  36. this[kEventEmitter] = ee;
  37. }
  38. get eventEmitter() {
  39. return this[kEventEmitter];
  40. }
  41. };
  42. _a = kEventEmitter;
  43. var _a2;
  44. var _EventEmitterAsyncResource = class extends EventEmitter {
  45. constructor(options) {
  46. let name;
  47. if (typeof options === "string") {
  48. name = options;
  49. options = void 0;
  50. } else {
  51. name = options?.name || new.target.name;
  52. }
  53. super(options);
  54. __publicField(this, _a2);
  55. this[kAsyncResource] = new EventEmitterReferencingAsyncResource(this, name, options);
  56. }
  57. emit(event, ...args) {
  58. return this.asyncResource.runInAsyncScope(super.emit, this, event, ...args);
  59. }
  60. emitDestroy() {
  61. this.asyncResource.emitDestroy();
  62. }
  63. asyncId() {
  64. return this.asyncResource.asyncId();
  65. }
  66. triggerAsyncId() {
  67. return this.asyncResource.triggerAsyncId();
  68. }
  69. get asyncResource() {
  70. return this[kAsyncResource];
  71. }
  72. static get EventEmitterAsyncResource() {
  73. return _EventEmitterAsyncResource;
  74. }
  75. };
  76. var EventEmitterAsyncResource = _EventEmitterAsyncResource;
  77. _a2 = kAsyncResource;
  78. var EventEmitterAsyncResource_default = EventEmitterAsyncResource;
  79. // src/index.ts
  80. import { AsyncResource as AsyncResource2 } from "async_hooks";
  81. import { fileURLToPath, URL } from "url";
  82. import { dirname, join, resolve } from "path";
  83. import { inspect, types } from "util";
  84. import assert from "assert";
  85. import { performance } from "perf_hooks";
  86. import { readFileSync } from "fs";
  87. // src/physicalCpuCount.ts
  88. import os from "os";
  89. import childProcess from "child_process";
  90. function exec(command) {
  91. const output = childProcess.execSync(command, {
  92. encoding: "utf8",
  93. stdio: [null, null, null]
  94. });
  95. return output;
  96. }
  97. var amount;
  98. try {
  99. const platform = os.platform();
  100. if (platform === "linux") {
  101. const output1 = exec('cat /proc/cpuinfo | grep "physical id" | sort |uniq | wc -l');
  102. const output2 = exec('cat /proc/cpuinfo | grep "core id" | sort | uniq | wc -l');
  103. const physicalCpuAmount = parseInt(output1.trim(), 10);
  104. const physicalCoreAmount = parseInt(output2.trim(), 10);
  105. amount = physicalCpuAmount * physicalCoreAmount;
  106. } else if (platform === "darwin") {
  107. const output = exec("sysctl -n hw.physicalcpu_max");
  108. amount = parseInt(output.trim(), 10);
  109. } else if (platform === "win32") {
  110. throw new Error();
  111. } else {
  112. const cores = os.cpus().filter(function(cpu, index) {
  113. const hasHyperthreading = cpu.model.includes("Intel");
  114. const isOdd = index % 2 === 1;
  115. return !hasHyperthreading || isOdd;
  116. });
  117. amount = cores.length;
  118. }
  119. } catch {
  120. amount = os.cpus().length;
  121. }
  122. // src/index.ts
  123. var cpuCount = amount;
  124. function onabort(abortSignal, listener) {
  125. if ("addEventListener" in abortSignal) {
  126. abortSignal.addEventListener("abort", listener, { once: true });
  127. } else {
  128. abortSignal.once("abort", listener);
  129. }
  130. }
  131. var AbortError = class extends Error {
  132. constructor() {
  133. super("The task has been aborted");
  134. }
  135. get name() {
  136. return "AbortError";
  137. }
  138. };
  139. var ArrayTaskQueue = class {
  140. constructor() {
  141. __publicField(this, "tasks", []);
  142. }
  143. get size() {
  144. return this.tasks.length;
  145. }
  146. shift() {
  147. return this.tasks.shift();
  148. }
  149. push(task) {
  150. this.tasks.push(task);
  151. }
  152. remove(task) {
  153. const index = this.tasks.indexOf(task);
  154. assert.notStrictEqual(index, -1);
  155. this.tasks.splice(index, 1);
  156. }
  157. };
  158. var kDefaultOptions = {
  159. filename: null,
  160. name: "default",
  161. minThreads: Math.max(cpuCount / 2, 1),
  162. maxThreads: cpuCount,
  163. idleTimeout: 0,
  164. maxQueue: Infinity,
  165. concurrentTasksPerWorker: 1,
  166. useAtomics: true,
  167. taskQueue: new ArrayTaskQueue(),
  168. trackUnmanagedFds: true
  169. };
  170. var kDefaultRunOptions = {
  171. transferList: void 0,
  172. filename: null,
  173. signal: null,
  174. name: null
  175. };
  176. var _value;
  177. var DirectlyTransferable = class {
  178. constructor(value) {
  179. __privateAdd(this, _value, void 0);
  180. __privateSet(this, _value, value);
  181. }
  182. get [kTransferable]() {
  183. return __privateGet(this, _value);
  184. }
  185. get [kValue]() {
  186. return __privateGet(this, _value);
  187. }
  188. };
  189. _value = new WeakMap();
  190. var _view;
  191. var ArrayBufferViewTransferable = class {
  192. constructor(view) {
  193. __privateAdd(this, _view, void 0);
  194. __privateSet(this, _view, view);
  195. }
  196. get [kTransferable]() {
  197. return __privateGet(this, _view).buffer;
  198. }
  199. get [kValue]() {
  200. return __privateGet(this, _view);
  201. }
  202. };
  203. _view = new WeakMap();
  204. var taskIdCounter = 0;
  205. function maybeFileURLToPath(filename) {
  206. return filename.startsWith("file:") ? fileURLToPath(new URL(filename)) : filename;
  207. }
  208. var TaskInfo = class extends AsyncResource2 {
  209. constructor(task, transferList, filename, name, callback, abortSignal, triggerAsyncId) {
  210. super("Tinypool.Task", { requireManualDestroy: true, triggerAsyncId });
  211. __publicField(this, "callback");
  212. __publicField(this, "task");
  213. __publicField(this, "transferList");
  214. __publicField(this, "filename");
  215. __publicField(this, "name");
  216. __publicField(this, "taskId");
  217. __publicField(this, "abortSignal");
  218. __publicField(this, "abortListener", null);
  219. __publicField(this, "workerInfo", null);
  220. __publicField(this, "created");
  221. __publicField(this, "started");
  222. this.callback = callback;
  223. this.task = task;
  224. this.transferList = transferList;
  225. if (isMovable(task)) {
  226. if (this.transferList == null) {
  227. this.transferList = [];
  228. }
  229. this.transferList = this.transferList.concat(task[kTransferable]);
  230. this.task = task[kValue];
  231. }
  232. this.filename = filename;
  233. this.name = name;
  234. this.taskId = taskIdCounter++;
  235. this.abortSignal = abortSignal;
  236. this.created = performance.now();
  237. this.started = 0;
  238. }
  239. releaseTask() {
  240. const ret = this.task;
  241. this.task = null;
  242. return ret;
  243. }
  244. done(err, result) {
  245. this.emitDestroy();
  246. this.runInAsyncScope(this.callback, null, err, result);
  247. if (this.abortSignal && this.abortListener) {
  248. if ("removeEventListener" in this.abortSignal && this.abortListener) {
  249. this.abortSignal.removeEventListener("abort", this.abortListener);
  250. } else {
  251. ;
  252. this.abortSignal.off("abort", this.abortListener);
  253. }
  254. }
  255. }
  256. get [kQueueOptions]() {
  257. return kQueueOptions in this.task ? this.task[kQueueOptions] : null;
  258. }
  259. };
  260. var AsynchronouslyCreatedResource = class {
  261. constructor() {
  262. __publicField(this, "onreadyListeners", []);
  263. }
  264. markAsReady() {
  265. const listeners = this.onreadyListeners;
  266. assert(listeners !== null);
  267. this.onreadyListeners = null;
  268. for (const listener of listeners) {
  269. listener();
  270. }
  271. }
  272. isReady() {
  273. return this.onreadyListeners === null;
  274. }
  275. onReady(fn) {
  276. if (this.onreadyListeners === null) {
  277. fn();
  278. return;
  279. }
  280. this.onreadyListeners.push(fn);
  281. }
  282. };
  283. var AsynchronouslyCreatedResourcePool = class {
  284. constructor(maximumUsage) {
  285. __publicField(this, "pendingItems", /* @__PURE__ */ new Set());
  286. __publicField(this, "readyItems", /* @__PURE__ */ new Set());
  287. __publicField(this, "maximumUsage");
  288. __publicField(this, "onAvailableListeners");
  289. this.maximumUsage = maximumUsage;
  290. this.onAvailableListeners = [];
  291. }
  292. add(item) {
  293. this.pendingItems.add(item);
  294. item.onReady(() => {
  295. if (this.pendingItems.has(item)) {
  296. this.pendingItems.delete(item);
  297. this.readyItems.add(item);
  298. this.maybeAvailable(item);
  299. }
  300. });
  301. }
  302. delete(item) {
  303. this.pendingItems.delete(item);
  304. this.readyItems.delete(item);
  305. }
  306. findAvailable() {
  307. let minUsage = this.maximumUsage;
  308. let candidate = null;
  309. for (const item of this.readyItems) {
  310. const usage = item.currentUsage();
  311. if (usage === 0)
  312. return item;
  313. if (usage < minUsage) {
  314. candidate = item;
  315. minUsage = usage;
  316. }
  317. }
  318. return candidate;
  319. }
  320. *[Symbol.iterator]() {
  321. yield* this.pendingItems;
  322. yield* this.readyItems;
  323. }
  324. get size() {
  325. return this.pendingItems.size + this.readyItems.size;
  326. }
  327. maybeAvailable(item) {
  328. if (item.currentUsage() < this.maximumUsage) {
  329. for (const listener of this.onAvailableListeners) {
  330. listener(item);
  331. }
  332. }
  333. }
  334. onAvailable(fn) {
  335. this.onAvailableListeners.push(fn);
  336. }
  337. };
  338. var Errors = {
  339. ThreadTermination: () => new Error("Terminating worker thread"),
  340. FilenameNotProvided: () => new Error("filename must be provided to run() or in options object"),
  341. TaskQueueAtLimit: () => new Error("Task queue is at limit"),
  342. NoTaskQueueAvailable: () => new Error("No task queue available and all Workers are busy")
  343. };
  344. var WorkerInfo = class extends AsynchronouslyCreatedResource {
  345. constructor(worker, port, workerId, freeWorkerId, onMessage) {
  346. super();
  347. __publicField(this, "worker");
  348. __publicField(this, "workerId");
  349. __publicField(this, "freeWorkerId");
  350. __publicField(this, "taskInfos");
  351. __publicField(this, "idleTimeout", null);
  352. __publicField(this, "port");
  353. __publicField(this, "sharedBuffer");
  354. __publicField(this, "lastSeenResponseCount", 0);
  355. __publicField(this, "onMessage");
  356. this.worker = worker;
  357. this.workerId = workerId;
  358. this.freeWorkerId = freeWorkerId;
  359. this.port = port;
  360. this.port.on("message", (message) => this._handleResponse(message));
  361. this.onMessage = onMessage;
  362. this.taskInfos = /* @__PURE__ */ new Map();
  363. this.sharedBuffer = new Int32Array(new SharedArrayBuffer(kFieldCount * Int32Array.BYTES_PER_ELEMENT));
  364. }
  365. async destroy() {
  366. await this.worker.terminate();
  367. this.port.close();
  368. this.clearIdleTimeout();
  369. for (const taskInfo of this.taskInfos.values()) {
  370. taskInfo.done(Errors.ThreadTermination());
  371. }
  372. this.taskInfos.clear();
  373. }
  374. clearIdleTimeout() {
  375. if (this.idleTimeout !== null) {
  376. clearTimeout(this.idleTimeout);
  377. this.idleTimeout = null;
  378. }
  379. }
  380. ref() {
  381. this.port.ref();
  382. return this;
  383. }
  384. unref() {
  385. this.port.unref();
  386. return this;
  387. }
  388. _handleResponse(message) {
  389. this.onMessage(message);
  390. if (this.taskInfos.size === 0) {
  391. this.unref();
  392. }
  393. }
  394. postTask(taskInfo) {
  395. assert(!this.taskInfos.has(taskInfo.taskId));
  396. const message = {
  397. task: taskInfo.releaseTask(),
  398. taskId: taskInfo.taskId,
  399. filename: taskInfo.filename,
  400. name: taskInfo.name
  401. };
  402. try {
  403. this.port.postMessage(message, taskInfo.transferList);
  404. } catch (err) {
  405. taskInfo.done(err);
  406. return;
  407. }
  408. taskInfo.workerInfo = this;
  409. this.taskInfos.set(taskInfo.taskId, taskInfo);
  410. this.ref();
  411. this.clearIdleTimeout();
  412. Atomics.add(this.sharedBuffer, kRequestCountField, 1);
  413. Atomics.notify(this.sharedBuffer, kRequestCountField, 1);
  414. }
  415. processPendingMessages() {
  416. const actualResponseCount = Atomics.load(this.sharedBuffer, kResponseCountField);
  417. if (actualResponseCount !== this.lastSeenResponseCount) {
  418. this.lastSeenResponseCount = actualResponseCount;
  419. let entry;
  420. while ((entry = receiveMessageOnPort(this.port)) !== void 0) {
  421. this._handleResponse(entry.message);
  422. }
  423. }
  424. }
  425. isRunningAbortableTask() {
  426. if (this.taskInfos.size !== 1)
  427. return false;
  428. const [[, task]] = this.taskInfos;
  429. return task.abortSignal !== null;
  430. }
  431. currentUsage() {
  432. if (this.isRunningAbortableTask())
  433. return Infinity;
  434. return this.taskInfos.size;
  435. }
  436. };
  437. var ThreadPool = class {
  438. constructor(publicInterface, options) {
  439. __publicField(this, "publicInterface");
  440. __publicField(this, "workers");
  441. __publicField(this, "workerIds");
  442. __publicField(this, "options");
  443. __publicField(this, "taskQueue");
  444. __publicField(this, "skipQueue", []);
  445. __publicField(this, "completed", 0);
  446. __publicField(this, "start", performance.now());
  447. __publicField(this, "inProcessPendingMessages", false);
  448. __publicField(this, "startingUp", false);
  449. __publicField(this, "workerFailsDuringBootstrap", false);
  450. this.publicInterface = publicInterface;
  451. this.taskQueue = options.taskQueue || new ArrayTaskQueue();
  452. const filename = options.filename ? maybeFileURLToPath(options.filename) : null;
  453. this.options = { ...kDefaultOptions, ...options, filename, maxQueue: 0 };
  454. if (options.maxThreads !== void 0 && this.options.minThreads >= options.maxThreads) {
  455. this.options.minThreads = options.maxThreads;
  456. }
  457. if (options.minThreads !== void 0 && this.options.maxThreads <= options.minThreads) {
  458. this.options.maxThreads = options.minThreads;
  459. }
  460. if (options.maxQueue === "auto") {
  461. this.options.maxQueue = this.options.maxThreads ** 2;
  462. } else {
  463. this.options.maxQueue = options.maxQueue ?? kDefaultOptions.maxQueue;
  464. }
  465. this.workerIds = new Map(new Array(this.options.maxThreads).fill(0).map((_, i) => [i + 1, true]));
  466. this.workers = new AsynchronouslyCreatedResourcePool(this.options.concurrentTasksPerWorker);
  467. this.workers.onAvailable((w) => this._onWorkerAvailable(w));
  468. this.startingUp = true;
  469. this._ensureMinimumWorkers();
  470. this.startingUp = false;
  471. }
  472. _ensureEnoughWorkersForTaskQueue() {
  473. while (this.workers.size < this.taskQueue.size && this.workers.size < this.options.maxThreads) {
  474. this._addNewWorker();
  475. }
  476. }
  477. _ensureMaximumWorkers() {
  478. while (this.workers.size < this.options.maxThreads) {
  479. this._addNewWorker();
  480. }
  481. }
  482. _ensureMinimumWorkers() {
  483. while (this.workers.size < this.options.minThreads) {
  484. this._addNewWorker();
  485. }
  486. }
  487. _addNewWorker() {
  488. const pool = this;
  489. const workerIds = this.workerIds;
  490. const __dirname2 = dirname(fileURLToPath(import.meta.url));
  491. let workerId;
  492. workerIds.forEach((isIdAvailable, _workerId2) => {
  493. if (isIdAvailable && !workerId) {
  494. workerId = _workerId2;
  495. workerIds.set(_workerId2, false);
  496. }
  497. });
  498. const tinypoolPrivateData = { workerId };
  499. const worker = new Worker(resolve(__dirname2, "./worker.js"), {
  500. env: this.options.env,
  501. argv: this.options.argv,
  502. execArgv: this.options.execArgv,
  503. resourceLimits: this.options.resourceLimits,
  504. workerData: [
  505. tinypoolPrivateData,
  506. this.options.workerData
  507. ],
  508. trackUnmanagedFds: this.options.trackUnmanagedFds
  509. });
  510. const onMessage = (message2) => {
  511. const { taskId, result } = message2;
  512. const taskInfo = workerInfo.taskInfos.get(taskId);
  513. workerInfo.taskInfos.delete(taskId);
  514. if (!this.options.isolateWorkers)
  515. pool.workers.maybeAvailable(workerInfo);
  516. if (taskInfo === void 0) {
  517. const err = new Error(`Unexpected message from Worker: ${inspect(message2)}`);
  518. pool.publicInterface.emit("error", err);
  519. } else {
  520. taskInfo.done(message2.error, result);
  521. }
  522. pool._processPendingMessages();
  523. };
  524. const { port1, port2 } = new MessageChannel();
  525. const workerInfo = new WorkerInfo(worker, port1, workerId, () => workerIds.set(workerId, true), onMessage);
  526. if (this.startingUp) {
  527. workerInfo.markAsReady();
  528. }
  529. const message = {
  530. filename: this.options.filename,
  531. name: this.options.name,
  532. port: port2,
  533. sharedBuffer: workerInfo.sharedBuffer,
  534. useAtomics: this.options.useAtomics
  535. };
  536. worker.postMessage(message, [port2]);
  537. worker.on("message", (message2) => {
  538. if (message2.ready === true) {
  539. if (workerInfo.currentUsage() === 0) {
  540. workerInfo.unref();
  541. }
  542. if (!workerInfo.isReady()) {
  543. workerInfo.markAsReady();
  544. }
  545. return;
  546. }
  547. worker.emit("error", new Error(`Unexpected message on Worker: ${inspect(message2)}`));
  548. });
  549. worker.on("error", (err) => {
  550. worker.ref = () => {
  551. };
  552. const taskInfos = [...workerInfo.taskInfos.values()];
  553. workerInfo.taskInfos.clear();
  554. this._removeWorker(workerInfo);
  555. if (workerInfo.isReady() && !this.workerFailsDuringBootstrap) {
  556. this._ensureMinimumWorkers();
  557. } else {
  558. this.workerFailsDuringBootstrap = true;
  559. }
  560. if (taskInfos.length > 0) {
  561. for (const taskInfo of taskInfos) {
  562. taskInfo.done(err, null);
  563. }
  564. } else {
  565. this.publicInterface.emit("error", err);
  566. }
  567. });
  568. worker.unref();
  569. port1.on("close", () => {
  570. worker.ref();
  571. });
  572. this.workers.add(workerInfo);
  573. }
  574. _processPendingMessages() {
  575. if (this.inProcessPendingMessages || !this.options.useAtomics) {
  576. return;
  577. }
  578. this.inProcessPendingMessages = true;
  579. try {
  580. for (const workerInfo of this.workers) {
  581. workerInfo.processPendingMessages();
  582. }
  583. } finally {
  584. this.inProcessPendingMessages = false;
  585. }
  586. }
  587. _removeWorker(workerInfo) {
  588. workerInfo.freeWorkerId();
  589. workerInfo.destroy();
  590. this.workers.delete(workerInfo);
  591. }
  592. _onWorkerAvailable(workerInfo) {
  593. while ((this.taskQueue.size > 0 || this.skipQueue.length > 0) && workerInfo.currentUsage() < this.options.concurrentTasksPerWorker) {
  594. const taskInfo = this.skipQueue.shift() || this.taskQueue.shift();
  595. if (taskInfo.abortSignal && workerInfo.taskInfos.size > 0) {
  596. this.skipQueue.push(taskInfo);
  597. break;
  598. }
  599. const now = performance.now();
  600. taskInfo.started = now;
  601. workerInfo.postTask(taskInfo);
  602. this._maybeDrain();
  603. return;
  604. }
  605. if (workerInfo.taskInfos.size === 0 && this.workers.size > this.options.minThreads) {
  606. workerInfo.idleTimeout = setTimeout(() => {
  607. assert.strictEqual(workerInfo.taskInfos.size, 0);
  608. if (this.workers.size > this.options.minThreads) {
  609. this._removeWorker(workerInfo);
  610. }
  611. }, this.options.idleTimeout).unref();
  612. }
  613. }
  614. runTask(task, options) {
  615. let { filename, name } = options;
  616. const { transferList = [], signal = null } = options;
  617. if (filename == null) {
  618. filename = this.options.filename;
  619. }
  620. if (name == null) {
  621. name = this.options.name;
  622. }
  623. if (typeof filename !== "string") {
  624. return Promise.reject(Errors.FilenameNotProvided());
  625. }
  626. filename = maybeFileURLToPath(filename);
  627. let resolve2;
  628. let reject;
  629. const ret = new Promise((res, rej) => {
  630. resolve2 = res;
  631. reject = rej;
  632. });
  633. const taskInfo = new TaskInfo(task, transferList, filename, name, (err, result) => {
  634. this.completed++;
  635. if (err !== null) {
  636. reject(err);
  637. } else {
  638. resolve2(result);
  639. }
  640. if (this.options.isolateWorkers && taskInfo.workerInfo) {
  641. this._removeWorker(taskInfo.workerInfo);
  642. this._ensureEnoughWorkersForTaskQueue();
  643. }
  644. }, signal, this.publicInterface.asyncResource.asyncId());
  645. if (signal !== null) {
  646. if (signal.aborted) {
  647. return Promise.reject(new AbortError());
  648. }
  649. taskInfo.abortListener = () => {
  650. reject(new AbortError());
  651. if (taskInfo.workerInfo !== null) {
  652. this._removeWorker(taskInfo.workerInfo);
  653. this._ensureMinimumWorkers();
  654. } else {
  655. this.taskQueue.remove(taskInfo);
  656. }
  657. };
  658. onabort(signal, taskInfo.abortListener);
  659. }
  660. if (this.taskQueue.size > 0) {
  661. const totalCapacity = this.options.maxQueue + this.pendingCapacity();
  662. if (this.taskQueue.size >= totalCapacity) {
  663. if (this.options.maxQueue === 0) {
  664. return Promise.reject(Errors.NoTaskQueueAvailable());
  665. } else {
  666. return Promise.reject(Errors.TaskQueueAtLimit());
  667. }
  668. } else {
  669. if (this.workers.size < this.options.maxThreads) {
  670. this._addNewWorker();
  671. }
  672. this.taskQueue.push(taskInfo);
  673. }
  674. return ret;
  675. }
  676. let workerInfo = this.workers.findAvailable();
  677. if (workerInfo !== null && workerInfo.currentUsage() > 0 && signal) {
  678. workerInfo = null;
  679. }
  680. let waitingForNewWorker = false;
  681. if ((workerInfo === null || workerInfo.currentUsage() > 0) && this.workers.size < this.options.maxThreads) {
  682. this._addNewWorker();
  683. waitingForNewWorker = true;
  684. }
  685. if (workerInfo === null) {
  686. if (this.options.maxQueue <= 0 && !waitingForNewWorker) {
  687. return Promise.reject(Errors.NoTaskQueueAvailable());
  688. } else {
  689. this.taskQueue.push(taskInfo);
  690. }
  691. return ret;
  692. }
  693. const now = performance.now();
  694. taskInfo.started = now;
  695. workerInfo.postTask(taskInfo);
  696. this._maybeDrain();
  697. return ret;
  698. }
  699. pendingCapacity() {
  700. return this.workers.pendingItems.size * this.options.concurrentTasksPerWorker;
  701. }
  702. _maybeDrain() {
  703. if (this.taskQueue.size === 0 && this.skipQueue.length === 0) {
  704. this.publicInterface.emit("drain");
  705. }
  706. }
  707. async destroy() {
  708. while (this.skipQueue.length > 0) {
  709. const taskInfo = this.skipQueue.shift();
  710. taskInfo.done(new Error("Terminating worker thread"));
  711. }
  712. while (this.taskQueue.size > 0) {
  713. const taskInfo = this.taskQueue.shift();
  714. taskInfo.done(new Error("Terminating worker thread"));
  715. }
  716. const exitEvents = [];
  717. while (this.workers.size > 0) {
  718. const [workerInfo] = this.workers;
  719. exitEvents.push(once(workerInfo.worker, "exit"));
  720. this._removeWorker(workerInfo);
  721. }
  722. await Promise.all(exitEvents);
  723. }
  724. };
  725. var _pool;
  726. var Tinypool = class extends EventEmitterAsyncResource_default {
  727. constructor(options = {}) {
  728. if (options.minThreads !== void 0 && options.minThreads > 0 && options.minThreads < 1) {
  729. options.minThreads = Math.max(1, Math.floor(options.minThreads * cpuCount));
  730. }
  731. if (options.maxThreads !== void 0 && options.maxThreads > 0 && options.maxThreads < 1) {
  732. options.maxThreads = Math.max(1, Math.floor(options.maxThreads * cpuCount));
  733. }
  734. super({ ...options, name: "Tinypool" });
  735. __privateAdd(this, _pool, void 0);
  736. if (options.minThreads !== void 0 && options.maxThreads !== void 0 && options.minThreads > options.maxThreads) {
  737. throw new RangeError("options.minThreads and options.maxThreads must not conflict");
  738. }
  739. __privateSet(this, _pool, new ThreadPool(this, options));
  740. }
  741. run(task, options = kDefaultRunOptions) {
  742. const { transferList, filename, name, signal } = options;
  743. return __privateGet(this, _pool).runTask(task, { transferList, filename, name, signal });
  744. }
  745. destroy() {
  746. return __privateGet(this, _pool).destroy();
  747. }
  748. get options() {
  749. return __privateGet(this, _pool).options;
  750. }
  751. get threads() {
  752. const ret = [];
  753. for (const workerInfo of __privateGet(this, _pool).workers) {
  754. ret.push(workerInfo.worker);
  755. }
  756. return ret;
  757. }
  758. get queueSize() {
  759. const pool = __privateGet(this, _pool);
  760. return Math.max(pool.taskQueue.size - pool.pendingCapacity(), 0);
  761. }
  762. get completed() {
  763. return __privateGet(this, _pool).completed;
  764. }
  765. get duration() {
  766. return performance.now() - __privateGet(this, _pool).start;
  767. }
  768. static get isWorkerThread() {
  769. return process.__tinypool_state__?.isWorkerThread || false;
  770. }
  771. static get workerData() {
  772. return process.__tinypool_state__?.workerData || void 0;
  773. }
  774. static get version() {
  775. const { version } = JSON.parse(readFileSync(join(__dirname, "../package.json"), "utf-8"));
  776. return version;
  777. }
  778. static move(val) {
  779. if (val != null && typeof val === "object" && typeof val !== "function") {
  780. if (!isTransferable(val)) {
  781. if (types.isArrayBufferView(val)) {
  782. val = new ArrayBufferViewTransferable(val);
  783. } else {
  784. val = new DirectlyTransferable(val);
  785. }
  786. }
  787. markMovable(val);
  788. }
  789. return val;
  790. }
  791. static get transferableSymbol() {
  792. return kTransferable;
  793. }
  794. static get valueSymbol() {
  795. return kValue;
  796. }
  797. static get queueOptionsSymbol() {
  798. return kQueueOptions;
  799. }
  800. };
  801. _pool = new WeakMap();
  802. var _workerId = process.__tinypool_state__?.workerId;
  803. var src_default = Tinypool;
  804. export {
  805. Tinypool,
  806. src_default as default,
  807. isMovable,
  808. isTaskQueue,
  809. isTransferable,
  810. kFieldCount,
  811. kQueueOptions,
  812. kRequestCountField,
  813. kResponseCountField,
  814. kTransferable,
  815. kValue,
  816. markMovable,
  817. _workerId as workerId
  818. };