Skip to content

Commit 3c1abe9

Browse files
github-actions[bot]metcoder95Copilot
authored
refactor: TaskInfo (#859)
Co-authored-by: Carlos Fuentes <me@metcoder.dev> Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
1 parent 3b73e34 commit 3c1abe9

3 files changed

Lines changed: 47 additions & 38 deletions

File tree

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/index.ts

Lines changed: 33 additions & 28 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
*
@@ -52,21 +62,22 @@ export class TaskInfo extends AsyncResource implements Task {
5262
name : string;
5363
taskId : number;
5464
abortSignal : AbortSignalAny | null;
55-
// abortListener : (() => void) | null = null;
5665
workerInfo : WorkerInfo | null = null;
5766
created : number;
5867
started : number;
5968
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) {
69+
_abortListener: (() => void) = () => { this.aborted = true; };
70+
_abortCleaner: (() => void) | null = null;
71+
72+
constructor ({
73+
task,
74+
transferList,
75+
filename,
76+
name,
77+
abortSignal,
78+
triggerAsyncId,
79+
}: TaskInfoParameters,
80+
callback: TaskCallback) {
7081
super('Piscina.Task', { requireManualDestroy: true, triggerAsyncId });
7182
this.callback = callback;
7283
this.task = task;
@@ -78,12 +89,11 @@ export class TaskInfo extends AsyncResource implements Task {
7889
if (isMovable(task)) {
7990
// This condition should never be hit but typescript
8091
// complains if we dont do the check.
81-
/* istanbul ignore if */
92+
/* c8 ignore next */
8293
if (this.transferList == null) {
8394
this.transferList = [];
8495
}
85-
this.transferList =
86-
this.transferList.concat(task[kTransferable]);
96+
this.transferList = this.transferList.concat(task[kTransferable]);
8797
this.task = task[kValue];
8898
}
8999

@@ -96,16 +106,15 @@ export class TaskInfo extends AsyncResource implements Task {
96106
this.started = 0;
97107
}
98108

99-
// TODO: improve this handling - ideally should be extended
100-
set abortListener (value: (() => void)) {
109+
onAbort (value: (() => void)) {
101110
this._abortListener = () => {
102111
this.aborted = true;
103112
value();
104113
};
105114
}
106115

107-
get abortListener (): (() => void) | null {
108-
return this._abortListener;
116+
setAbortListener(signal: AbortSignalAny) : void {
117+
this._abortCleaner = onabort(signal, this._abortListener);
109118
}
110119

111120
releaseTask () : any {
@@ -114,19 +123,15 @@ export class TaskInfo extends AsyncResource implements Task {
114123
return ret;
115124
}
116125

126+
// TODO: implement - helpful for streaming chunks of data from worker to parent
127+
onResponse(_result: any) {}
128+
117129
done (err : Error | null, result? : any) : void {
118130
this.runInAsyncScope(this.callback, null, err, result);
119131
this.emitDestroy(); // `TaskInfo`s are used only once.
120132
// If an abort signal was used, remove the listener from it when
121133
// 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-
}
134+
this._abortCleaner?.();
130135
}
131136

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

0 commit comments

Comments
 (0)