版博士V2.0程序
25'ten fazla konu seçemezsiniz Konular bir harf veya rakamla başlamalı, kısa çizgiler ('-') içerebilir ve en fazla 35 karakter uzunluğunda olabilir.

index.js 25 KiB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835
  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(timeout) {
  366. let resolve2;
  367. let reject;
  368. const ret = new Promise((res, rej) => {
  369. resolve2 = res;
  370. reject = rej;
  371. });
  372. const timer = timeout ? setTimeout(() => reject(new Error("Failed to terminate worker")), timeout) : null;
  373. this.worker.terminate().then(() => {
  374. if (timer !== null) {
  375. clearTimeout(timer);
  376. }
  377. this.port.close();
  378. this.clearIdleTimeout();
  379. for (const taskInfo of this.taskInfos.values()) {
  380. taskInfo.done(Errors.ThreadTermination());
  381. }
  382. this.taskInfos.clear();
  383. resolve2();
  384. });
  385. return ret;
  386. }
  387. clearIdleTimeout() {
  388. if (this.idleTimeout !== null) {
  389. clearTimeout(this.idleTimeout);
  390. this.idleTimeout = null;
  391. }
  392. }
  393. ref() {
  394. this.port.ref();
  395. return this;
  396. }
  397. unref() {
  398. this.port.unref();
  399. return this;
  400. }
  401. _handleResponse(message) {
  402. this.onMessage(message);
  403. if (this.taskInfos.size === 0) {
  404. this.unref();
  405. }
  406. }
  407. postTask(taskInfo) {
  408. assert(!this.taskInfos.has(taskInfo.taskId));
  409. const message = {
  410. task: taskInfo.releaseTask(),
  411. taskId: taskInfo.taskId,
  412. filename: taskInfo.filename,
  413. name: taskInfo.name
  414. };
  415. try {
  416. this.port.postMessage(message, taskInfo.transferList);
  417. } catch (err) {
  418. taskInfo.done(err);
  419. return;
  420. }
  421. taskInfo.workerInfo = this;
  422. this.taskInfos.set(taskInfo.taskId, taskInfo);
  423. this.ref();
  424. this.clearIdleTimeout();
  425. Atomics.add(this.sharedBuffer, kRequestCountField, 1);
  426. Atomics.notify(this.sharedBuffer, kRequestCountField, 1);
  427. }
  428. processPendingMessages() {
  429. const actualResponseCount = Atomics.load(this.sharedBuffer, kResponseCountField);
  430. if (actualResponseCount !== this.lastSeenResponseCount) {
  431. this.lastSeenResponseCount = actualResponseCount;
  432. let entry;
  433. while ((entry = receiveMessageOnPort(this.port)) !== void 0) {
  434. this._handleResponse(entry.message);
  435. }
  436. }
  437. }
  438. isRunningAbortableTask() {
  439. if (this.taskInfos.size !== 1)
  440. return false;
  441. const [[, task]] = this.taskInfos;
  442. return task.abortSignal !== null;
  443. }
  444. currentUsage() {
  445. if (this.isRunningAbortableTask())
  446. return Infinity;
  447. return this.taskInfos.size;
  448. }
  449. };
  450. var ThreadPool = class {
  451. constructor(publicInterface, options) {
  452. __publicField(this, "publicInterface");
  453. __publicField(this, "workers");
  454. __publicField(this, "workerIds");
  455. __publicField(this, "options");
  456. __publicField(this, "taskQueue");
  457. __publicField(this, "skipQueue", []);
  458. __publicField(this, "completed", 0);
  459. __publicField(this, "start", performance.now());
  460. __publicField(this, "inProcessPendingMessages", false);
  461. __publicField(this, "startingUp", false);
  462. __publicField(this, "workerFailsDuringBootstrap", false);
  463. this.publicInterface = publicInterface;
  464. this.taskQueue = options.taskQueue || new ArrayTaskQueue();
  465. const filename = options.filename ? maybeFileURLToPath(options.filename) : null;
  466. this.options = { ...kDefaultOptions, ...options, filename, maxQueue: 0 };
  467. if (options.maxThreads !== void 0 && this.options.minThreads >= options.maxThreads) {
  468. this.options.minThreads = options.maxThreads;
  469. }
  470. if (options.minThreads !== void 0 && this.options.maxThreads <= options.minThreads) {
  471. this.options.maxThreads = options.minThreads;
  472. }
  473. if (options.maxQueue === "auto") {
  474. this.options.maxQueue = this.options.maxThreads ** 2;
  475. } else {
  476. this.options.maxQueue = options.maxQueue ?? kDefaultOptions.maxQueue;
  477. }
  478. this.workerIds = new Map(new Array(this.options.maxThreads).fill(0).map((_, i) => [i + 1, true]));
  479. this.workers = new AsynchronouslyCreatedResourcePool(this.options.concurrentTasksPerWorker);
  480. this.workers.onAvailable((w) => this._onWorkerAvailable(w));
  481. this.startingUp = true;
  482. this._ensureMinimumWorkers();
  483. this.startingUp = false;
  484. }
  485. _ensureEnoughWorkersForTaskQueue() {
  486. while (this.workers.size < this.taskQueue.size && this.workers.size < this.options.maxThreads) {
  487. this._addNewWorker();
  488. }
  489. }
  490. _ensureMaximumWorkers() {
  491. while (this.workers.size < this.options.maxThreads) {
  492. this._addNewWorker();
  493. }
  494. }
  495. _ensureMinimumWorkers() {
  496. while (this.workers.size < this.options.minThreads) {
  497. this._addNewWorker();
  498. }
  499. }
  500. _addNewWorker() {
  501. const pool = this;
  502. const workerIds = this.workerIds;
  503. const __dirname2 = dirname(fileURLToPath(import.meta.url));
  504. let workerId;
  505. workerIds.forEach((isIdAvailable, _workerId2) => {
  506. if (isIdAvailable && !workerId) {
  507. workerId = _workerId2;
  508. workerIds.set(_workerId2, false);
  509. }
  510. });
  511. const tinypoolPrivateData = { workerId };
  512. const worker = new Worker(resolve(__dirname2, "./worker.js"), {
  513. env: this.options.env,
  514. argv: this.options.argv,
  515. execArgv: this.options.execArgv,
  516. resourceLimits: this.options.resourceLimits,
  517. workerData: [
  518. tinypoolPrivateData,
  519. this.options.workerData
  520. ],
  521. trackUnmanagedFds: this.options.trackUnmanagedFds
  522. });
  523. const onMessage = (message2) => {
  524. const { taskId, result } = message2;
  525. const taskInfo = workerInfo.taskInfos.get(taskId);
  526. workerInfo.taskInfos.delete(taskId);
  527. if (!this.options.isolateWorkers)
  528. pool.workers.maybeAvailable(workerInfo);
  529. if (taskInfo === void 0) {
  530. const err = new Error(`Unexpected message from Worker: ${inspect(message2)}`);
  531. pool.publicInterface.emit("error", err);
  532. } else {
  533. taskInfo.done(message2.error, result);
  534. }
  535. pool._processPendingMessages();
  536. };
  537. const { port1, port2 } = new MessageChannel();
  538. const workerInfo = new WorkerInfo(worker, port1, workerId, () => workerIds.set(workerId, true), onMessage);
  539. if (this.startingUp) {
  540. workerInfo.markAsReady();
  541. }
  542. const message = {
  543. filename: this.options.filename,
  544. name: this.options.name,
  545. port: port2,
  546. sharedBuffer: workerInfo.sharedBuffer,
  547. useAtomics: this.options.useAtomics
  548. };
  549. worker.postMessage(message, [port2]);
  550. worker.on("message", (message2) => {
  551. if (message2.ready === true) {
  552. if (workerInfo.currentUsage() === 0) {
  553. workerInfo.unref();
  554. }
  555. if (!workerInfo.isReady()) {
  556. workerInfo.markAsReady();
  557. }
  558. return;
  559. }
  560. worker.emit("error", new Error(`Unexpected message on Worker: ${inspect(message2)}`));
  561. });
  562. worker.on("error", (err) => {
  563. worker.ref = () => {
  564. };
  565. const taskInfos = [...workerInfo.taskInfos.values()];
  566. workerInfo.taskInfos.clear();
  567. this._removeWorker(workerInfo);
  568. if (workerInfo.isReady() && !this.workerFailsDuringBootstrap) {
  569. this._ensureMinimumWorkers();
  570. } else {
  571. this.workerFailsDuringBootstrap = true;
  572. }
  573. if (taskInfos.length > 0) {
  574. for (const taskInfo of taskInfos) {
  575. taskInfo.done(err, null);
  576. }
  577. } else {
  578. this.publicInterface.emit("error", err);
  579. }
  580. });
  581. worker.unref();
  582. port1.on("close", () => {
  583. worker.ref();
  584. });
  585. this.workers.add(workerInfo);
  586. }
  587. _processPendingMessages() {
  588. if (this.inProcessPendingMessages || !this.options.useAtomics) {
  589. return;
  590. }
  591. this.inProcessPendingMessages = true;
  592. try {
  593. for (const workerInfo of this.workers) {
  594. workerInfo.processPendingMessages();
  595. }
  596. } finally {
  597. this.inProcessPendingMessages = false;
  598. }
  599. }
  600. _removeWorker(workerInfo) {
  601. workerInfo.freeWorkerId();
  602. this.workers.delete(workerInfo);
  603. return workerInfo.destroy(this.options.terminateTimeout);
  604. }
  605. _onWorkerAvailable(workerInfo) {
  606. while ((this.taskQueue.size > 0 || this.skipQueue.length > 0) && workerInfo.currentUsage() < this.options.concurrentTasksPerWorker) {
  607. const taskInfo = this.skipQueue.shift() || this.taskQueue.shift();
  608. if (taskInfo.abortSignal && workerInfo.taskInfos.size > 0) {
  609. this.skipQueue.push(taskInfo);
  610. break;
  611. }
  612. const now = performance.now();
  613. taskInfo.started = now;
  614. workerInfo.postTask(taskInfo);
  615. this._maybeDrain();
  616. return;
  617. }
  618. if (workerInfo.taskInfos.size === 0 && this.workers.size > this.options.minThreads) {
  619. workerInfo.idleTimeout = setTimeout(() => {
  620. assert.strictEqual(workerInfo.taskInfos.size, 0);
  621. if (this.workers.size > this.options.minThreads) {
  622. this._removeWorker(workerInfo);
  623. }
  624. }, this.options.idleTimeout).unref();
  625. }
  626. }
  627. runTask(task, options) {
  628. let { filename, name } = options;
  629. const { transferList = [], signal = null } = options;
  630. if (filename == null) {
  631. filename = this.options.filename;
  632. }
  633. if (name == null) {
  634. name = this.options.name;
  635. }
  636. if (typeof filename !== "string") {
  637. return Promise.reject(Errors.FilenameNotProvided());
  638. }
  639. filename = maybeFileURLToPath(filename);
  640. let resolve2;
  641. let reject;
  642. const ret = new Promise((res, rej) => {
  643. resolve2 = res;
  644. reject = rej;
  645. });
  646. const taskInfo = new TaskInfo(task, transferList, filename, name, (err, result) => {
  647. this.completed++;
  648. if (err !== null) {
  649. reject(err);
  650. }
  651. if (this.options.isolateWorkers && taskInfo.workerInfo) {
  652. this._removeWorker(taskInfo.workerInfo).then(() => this._ensureEnoughWorkersForTaskQueue()).then(() => resolve2(result)).catch(reject);
  653. } else {
  654. resolve2(result);
  655. }
  656. }, signal, this.publicInterface.asyncResource.asyncId());
  657. if (signal !== null) {
  658. if (signal.aborted) {
  659. return Promise.reject(new AbortError());
  660. }
  661. taskInfo.abortListener = () => {
  662. reject(new AbortError());
  663. if (taskInfo.workerInfo !== null) {
  664. this._removeWorker(taskInfo.workerInfo);
  665. this._ensureMinimumWorkers();
  666. } else {
  667. this.taskQueue.remove(taskInfo);
  668. }
  669. };
  670. onabort(signal, taskInfo.abortListener);
  671. }
  672. if (this.taskQueue.size > 0) {
  673. const totalCapacity = this.options.maxQueue + this.pendingCapacity();
  674. if (this.taskQueue.size >= totalCapacity) {
  675. if (this.options.maxQueue === 0) {
  676. return Promise.reject(Errors.NoTaskQueueAvailable());
  677. } else {
  678. return Promise.reject(Errors.TaskQueueAtLimit());
  679. }
  680. } else {
  681. if (this.workers.size < this.options.maxThreads) {
  682. this._addNewWorker();
  683. }
  684. this.taskQueue.push(taskInfo);
  685. }
  686. return ret;
  687. }
  688. let workerInfo = this.workers.findAvailable();
  689. if (workerInfo !== null && workerInfo.currentUsage() > 0 && signal) {
  690. workerInfo = null;
  691. }
  692. let waitingForNewWorker = false;
  693. if ((workerInfo === null || workerInfo.currentUsage() > 0) && this.workers.size < this.options.maxThreads) {
  694. this._addNewWorker();
  695. waitingForNewWorker = true;
  696. }
  697. if (workerInfo === null) {
  698. if (this.options.maxQueue <= 0 && !waitingForNewWorker) {
  699. return Promise.reject(Errors.NoTaskQueueAvailable());
  700. } else {
  701. this.taskQueue.push(taskInfo);
  702. }
  703. return ret;
  704. }
  705. const now = performance.now();
  706. taskInfo.started = now;
  707. workerInfo.postTask(taskInfo);
  708. this._maybeDrain();
  709. return ret;
  710. }
  711. pendingCapacity() {
  712. return this.workers.pendingItems.size * this.options.concurrentTasksPerWorker;
  713. }
  714. _maybeDrain() {
  715. if (this.taskQueue.size === 0 && this.skipQueue.length === 0) {
  716. this.publicInterface.emit("drain");
  717. }
  718. }
  719. async destroy() {
  720. while (this.skipQueue.length > 0) {
  721. const taskInfo = this.skipQueue.shift();
  722. taskInfo.done(new Error("Terminating worker thread"));
  723. }
  724. while (this.taskQueue.size > 0) {
  725. const taskInfo = this.taskQueue.shift();
  726. taskInfo.done(new Error("Terminating worker thread"));
  727. }
  728. const exitEvents = [];
  729. while (this.workers.size > 0) {
  730. const [workerInfo] = this.workers;
  731. exitEvents.push(once(workerInfo.worker, "exit"));
  732. this._removeWorker(workerInfo);
  733. }
  734. await Promise.all(exitEvents);
  735. }
  736. };
  737. var _pool;
  738. var Tinypool = class extends EventEmitterAsyncResource_default {
  739. constructor(options = {}) {
  740. if (options.minThreads !== void 0 && options.minThreads > 0 && options.minThreads < 1) {
  741. options.minThreads = Math.max(1, Math.floor(options.minThreads * cpuCount));
  742. }
  743. if (options.maxThreads !== void 0 && options.maxThreads > 0 && options.maxThreads < 1) {
  744. options.maxThreads = Math.max(1, Math.floor(options.maxThreads * cpuCount));
  745. }
  746. super({ ...options, name: "Tinypool" });
  747. __privateAdd(this, _pool, void 0);
  748. if (options.minThreads !== void 0 && options.maxThreads !== void 0 && options.minThreads > options.maxThreads) {
  749. throw new RangeError("options.minThreads and options.maxThreads must not conflict");
  750. }
  751. __privateSet(this, _pool, new ThreadPool(this, options));
  752. }
  753. run(task, options = kDefaultRunOptions) {
  754. const { transferList, filename, name, signal } = options;
  755. return __privateGet(this, _pool).runTask(task, { transferList, filename, name, signal });
  756. }
  757. destroy() {
  758. return __privateGet(this, _pool).destroy();
  759. }
  760. get options() {
  761. return __privateGet(this, _pool).options;
  762. }
  763. get threads() {
  764. const ret = [];
  765. for (const workerInfo of __privateGet(this, _pool).workers) {
  766. ret.push(workerInfo.worker);
  767. }
  768. return ret;
  769. }
  770. get queueSize() {
  771. const pool = __privateGet(this, _pool);
  772. return Math.max(pool.taskQueue.size - pool.pendingCapacity(), 0);
  773. }
  774. get completed() {
  775. return __privateGet(this, _pool).completed;
  776. }
  777. get duration() {
  778. return performance.now() - __privateGet(this, _pool).start;
  779. }
  780. static get isWorkerThread() {
  781. return process.__tinypool_state__?.isWorkerThread || false;
  782. }
  783. static get workerData() {
  784. return process.__tinypool_state__?.workerData || void 0;
  785. }
  786. static get version() {
  787. const { version } = JSON.parse(readFileSync(join(__dirname, "../package.json"), "utf-8"));
  788. return version;
  789. }
  790. static move(val) {
  791. if (val != null && typeof val === "object" && typeof val !== "function") {
  792. if (!isTransferable(val)) {
  793. if (types.isArrayBufferView(val)) {
  794. val = new ArrayBufferViewTransferable(val);
  795. } else {
  796. val = new DirectlyTransferable(val);
  797. }
  798. }
  799. markMovable(val);
  800. }
  801. return val;
  802. }
  803. static get transferableSymbol() {
  804. return kTransferable;
  805. }
  806. static get valueSymbol() {
  807. return kValue;
  808. }
  809. static get queueOptionsSymbol() {
  810. return kQueueOptions;
  811. }
  812. };
  813. _pool = new WeakMap();
  814. var _workerId = process.__tinypool_state__?.workerId;
  815. var src_default = Tinypool;
  816. export {
  817. Tinypool,
  818. src_default as default,
  819. isMovable,
  820. isTaskQueue,
  821. isTransferable,
  822. kFieldCount,
  823. kQueueOptions,
  824. kRequestCountField,
  825. kResponseCountField,
  826. kTransferable,
  827. kValue,
  828. markMovable,
  829. _workerId as workerId
  830. };