Skip to content

Commit ad37e01

Browse files
authored
Merge branch 'current' into typescript-entrypoint
2 parents 0977853 + ec213ad commit ad37e01

7 files changed

Lines changed: 66 additions & 48 deletions

File tree

docs/docs/api-reference/class.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -223,7 +223,7 @@ If the `PiscinaLoadBalancer` returns `null`, `Piscina` will attempt to spawn a n
223223
224224
```ts
225225
interface PiscinaTask {
226-
taskId: number; // Unique identifier for the task
226+
taskId: string; // Unique identifier for the task
227227
filename: string; // Filename of the worker module
228228
name: string; // Name of the worker function
229229
created: number; // Timestamp when the task was created

src/abort.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,5 @@
1+
import type { EventEmitter } from 'node:events';
2+
13
interface AbortSignalEventTargetAddOptions {
24
once: boolean;
35
}
@@ -32,10 +34,12 @@ export class AbortError extends Error {
3234
}
3335
}
3436

35-
export function onabort (abortSignal: AbortSignalAny, listener: () => void) {
37+
export function onabort (abortSignal: AbortSignalAny, listener: () => void): () => void {
3638
if ('addEventListener' in abortSignal) {
3739
abortSignal.addEventListener('abort', listener, { once: true });
40+
return () => abortSignal.removeEventListener('abort', listener);
3841
} else {
3942
abortSignal.once('abort', listener);
43+
return () => (abortSignal as EventEmitter).removeListener('abort', listener);
4044
}
4145
}

src/index.ts

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,6 @@ import {
4040
AbortSignalAny,
4141
AbortSignalEventTarget,
4242
AbortError,
43-
onabort
4443
} from './abort';
4544
import {
4645
PiscinaHistogram,
@@ -507,12 +506,15 @@ class ThreadPool {
507506
}
508507

509508
const { promise: ret, resolve, reject } = promiseResolvers();
510-
const taskInfo = new TaskInfo(
509+
const taskInfo = new TaskInfo({
511510
task,
512511
transferList,
513512
filename,
514513
name,
515-
(err : Error | null, result : any) => {
514+
abortSignal: signal,
515+
triggerAsyncId: this.publicInterface.asyncResource.asyncId()
516+
},
517+
(err : Error | null, result : any) => {
516518
this.completed++;
517519
if (taskInfo.started) {
518520
this.histogram?.recordRunTime(performance.now() - taskInfo.started);
@@ -524,9 +526,7 @@ class ThreadPool {
524526
}
525527

526528
queueMicrotask(this._maybeDrain.bind(this))
527-
},
528-
signal,
529-
this.publicInterface.asyncResource.asyncId());
529+
});
530530

531531
if (signal != null) {
532532
// If the AbortSignal has an aborted property and it's truthy,
@@ -536,7 +536,7 @@ class ThreadPool {
536536
return ret;
537537
}
538538

539-
taskInfo.abortListener = () => {
539+
taskInfo.onAbort(() => {
540540
// Call reject() first to make sure we always reject with the AbortError
541541
// if the task is aborted, not with an Error from the possible
542542
// thread termination below.
@@ -551,9 +551,9 @@ class ThreadPool {
551551
// Call should be idempotent
552552
this.taskQueue.remove(taskInfo);
553553
}
554-
};
554+
});
555555

556-
onabort(signal, taskInfo.abortListener);
556+
taskInfo.setAbortListener(signal);
557557
}
558558

559559
if (this.taskQueue.size > 0) {

src/task_queue/common.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ export interface TaskQueue {
99

1010
// Public Interface
1111
export interface PiscinaTask extends Task {
12-
taskId: number;
12+
taskId: string;
1313
filename: string;
1414
name: string;
1515
created: number;

src/task_queue/index.ts

Lines changed: 46 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -3,11 +3,12 @@ import { performance } from 'node:perf_hooks';
33
import { AsyncResource } from 'node:async_hooks';
44

55
import type { WorkerInfo } from '../worker_pool';
6-
import type { AbortSignalAny, AbortSignalEventEmitter } from '../abort';
6+
import type { Task, TaskQueue, PiscinaTask } from './common';
7+
8+
import { onabort, type AbortSignalAny } from '../abort';
79
import { isMovable } from '../common';
810
import { kTransferable, kValue, kQueueOptions } from '../symbols';
911

10-
import type { Task, TaskQueue, PiscinaTask } from './common';
1112

1213
export { ArrayTaskQueue } from './array_queue';
1314
export { FixedQueue } from './fixed_queue';
@@ -23,6 +24,15 @@ export type TransferList = MessagePort extends {
2324
: never
2425
export type TransferListItem = TransferList extends Array<infer T> ? T : never
2526

27+
type TaskInfoParameters = {
28+
task : any;
29+
transferList : TransferList;
30+
filename : string;
31+
name : string;
32+
abortSignal : AbortSignalAny | null;
33+
triggerAsyncId : number;
34+
}
35+
2636
/**
2737
* Verifies if a given TaskQueue is valid
2838
*
@@ -41,32 +51,43 @@ export function isTaskQueue (value: TaskQueue): boolean {
4151
);
4252
}
4353

44-
let taskIdCounter = 0;
54+
55+
function taskIdFactory() {
56+
let taskIdCounter = 0;
57+
const maxint = 2147483647;
58+
return () => {
59+
taskIdCounter = (taskIdCounter + 1) & maxint;
60+
return taskIdCounter.toString(36);
61+
}
62+
}
63+
4564
// Extend AsyncResource so that async relations between posting a task and
4665
// receiving its result are visible to diagnostic tools.
4766
export class TaskInfo extends AsyncResource implements Task {
67+
static getTaskId: () => string = taskIdFactory();
4868
callback : TaskCallback;
4969
task : any;
5070
transferList : TransferList;
5171
filename : string;
5272
name : string;
53-
taskId : number;
73+
taskId : string;
5474
abortSignal : AbortSignalAny | null;
55-
// abortListener : (() => void) | null = null;
5675
workerInfo : WorkerInfo | null = null;
5776
created : number;
5877
started : number;
5978
aborted = false;
60-
_abortListener: (() => void) | null = null;
61-
62-
constructor (
63-
task : any,
64-
transferList : TransferList,
65-
filename : string,
66-
name : string,
67-
callback : TaskCallback,
68-
abortSignal : AbortSignalAny | null,
69-
triggerAsyncId : number) {
79+
_abortListener: (() => void) = () => { this.aborted = true; };
80+
_abortCleaner: (() => void) | null = null;
81+
82+
constructor ({
83+
task,
84+
transferList,
85+
filename,
86+
name,
87+
abortSignal,
88+
triggerAsyncId,
89+
}: TaskInfoParameters,
90+
callback: TaskCallback) {
7091
super('Piscina.Task', { requireManualDestroy: true, triggerAsyncId });
7192
this.callback = callback;
7293
this.task = task;
@@ -78,34 +99,31 @@ export class TaskInfo extends AsyncResource implements Task {
7899
if (isMovable(task)) {
79100
// This condition should never be hit but typescript
80101
// complains if we dont do the check.
81-
/* istanbul ignore if */
102+
/* c8 ignore next */
82103
if (this.transferList == null) {
83104
this.transferList = [];
84105
}
85-
this.transferList =
86-
this.transferList.concat(task[kTransferable]);
106+
this.transferList = this.transferList.concat(task[kTransferable]);
87107
this.task = task[kValue];
88108
}
89109

90110
this.filename = filename;
91111
this.name = name;
92-
// TODO: This should not be global
93-
this.taskId = taskIdCounter++;
112+
this.taskId = TaskInfo.getTaskId();
94113
this.abortSignal = abortSignal;
95114
this.created = performance.now();
96115
this.started = 0;
97116
}
98117

99-
// TODO: improve this handling - ideally should be extended
100-
set abortListener (value: (() => void)) {
118+
onAbort (value: (() => void)) {
101119
this._abortListener = () => {
102120
this.aborted = true;
103121
value();
104122
};
105123
}
106124

107-
get abortListener (): (() => void) | null {
108-
return this._abortListener;
125+
setAbortListener(signal: AbortSignalAny) : void {
126+
this._abortCleaner = onabort(signal, this._abortListener);
109127
}
110128

111129
releaseTask () : any {
@@ -114,19 +132,15 @@ export class TaskInfo extends AsyncResource implements Task {
114132
return ret;
115133
}
116134

135+
// TODO: implement - helpful for streaming chunks of data from worker to parent
136+
onResponse(_result: any) {}
137+
117138
done (err : Error | null, result? : any) : void {
118139
this.runInAsyncScope(this.callback, null, err, result);
119140
this.emitDestroy(); // `TaskInfo`s are used only once.
120141
// If an abort signal was used, remove the listener from it when
121142
// done to make sure we do not accidentally leak.
122-
if (this.abortSignal && this.abortListener) {
123-
if ('removeEventListener' in this.abortSignal && this.abortListener) {
124-
this.abortSignal.removeEventListener('abort', this.abortListener);
125-
} else {
126-
(this.abortSignal as AbortSignalEventEmitter).off(
127-
'abort', this.abortListener);
128-
}
129-
}
143+
this._abortCleaner?.();
130144
}
131145

132146
get [kQueueOptions] () : {} | null {

src/types.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ export interface StartupMessage {
1313
}
1414

1515
export interface RequestMessage {
16-
taskId: number
16+
taskId: string
1717
task: any
1818
filename: string
1919
name: string
@@ -25,7 +25,7 @@ export interface ReadyMessage {
2525
}
2626

2727
export interface ResponseMessage {
28-
taskId: number
28+
taskId: string
2929
result: any
3030
error: Error | null
3131
time: number | null

src/worker_pool/index.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ type WorkerInfoParams = {
3333

3434
export class WorkerInfo extends AsynchronouslyCreatedResource {
3535
worker : Worker;
36-
taskInfos : Map<number, TaskInfo>;
36+
taskInfos : Map<string, TaskInfo>;
3737
idleTimeout : NodeJS.Timeout | null = null;
3838
port : MessagePort;
3939
sharedBuffer : Int32Array;
@@ -184,7 +184,7 @@ export class WorkerInfo extends AsynchronouslyCreatedResource {
184184
return this.taskInfos.size;
185185
}
186186

187-
popTask (taskId: number) : TaskInfo | null {
187+
popTask (taskId: string) : TaskInfo | null {
188188
const task = this.taskInfos.get(taskId) ?? null;
189189

190190
this.taskInfos.delete(taskId);

0 commit comments

Comments
 (0)