Skip to content

Commit ec213ad

Browse files
metcoder95Copilot
andauthored
refactor: TaskInfo (#857)
Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com>
1 parent e3e2d39 commit ec213ad

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
*
@@ -62,21 +72,22 @@ export class TaskInfo extends AsyncResource implements Task {
6272
name : string;
6373
taskId : string;
6474
abortSignal : AbortSignalAny | null;
65-
// abortListener : (() => void) | null = null;
6675
workerInfo : WorkerInfo | null = null;
6776
created : number;
6877
started : number;
6978
aborted = false;
70-
_abortListener: (() => void) | null = null;
71-
72-
constructor (
73-
task : any,
74-
transferList : TransferList,
75-
filename : string,
76-
name : string,
77-
callback : TaskCallback,
78-
abortSignal : AbortSignalAny | null,
79-
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) {
8091
super('Piscina.Task', { requireManualDestroy: true, triggerAsyncId });
8192
this.callback = callback;
8293
this.task = task;
@@ -88,12 +99,11 @@ export class TaskInfo extends AsyncResource implements Task {
8899
if (isMovable(task)) {
89100
// This condition should never be hit but typescript
90101
// complains if we dont do the check.
91-
/* istanbul ignore if */
102+
/* c8 ignore next */
92103
if (this.transferList == null) {
93104
this.transferList = [];
94105
}
95-
this.transferList =
96-
this.transferList.concat(task[kTransferable]);
106+
this.transferList = this.transferList.concat(task[kTransferable]);
97107
this.task = task[kValue];
98108
}
99109

@@ -105,16 +115,15 @@ export class TaskInfo extends AsyncResource implements Task {
105115
this.started = 0;
106116
}
107117

108-
// TODO: improve this handling - ideally should be extended
109-
set abortListener (value: (() => void)) {
118+
onAbort (value: (() => void)) {
110119
this._abortListener = () => {
111120
this.aborted = true;
112121
value();
113122
};
114123
}
115124

116-
get abortListener (): (() => void) | null {
117-
return this._abortListener;
125+
setAbortListener(signal: AbortSignalAny) : void {
126+
this._abortCleaner = onabort(signal, this._abortListener);
118127
}
119128

120129
releaseTask () : any {
@@ -123,19 +132,15 @@ export class TaskInfo extends AsyncResource implements Task {
123132
return ret;
124133
}
125134

135+
// TODO: implement - helpful for streaming chunks of data from worker to parent
136+
onResponse(_result: any) {}
137+
126138
done (err : Error | null, result? : any) : void {
127139
this.runInAsyncScope(this.callback, null, err, result);
128140
this.emitDestroy(); // `TaskInfo`s are used only once.
129141
// If an abort signal was used, remove the listener from it when
130142
// done to make sure we do not accidentally leak.
131-
if (this.abortSignal && this.abortListener) {
132-
if ('removeEventListener' in this.abortSignal && this.abortListener) {
133-
this.abortSignal.removeEventListener('abort', this.abortListener);
134-
} else {
135-
(this.abortSignal as AbortSignalEventEmitter).off(
136-
'abort', this.abortListener);
137-
}
138-
}
143+
this._abortCleaner?.();
139144
}
140145

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

0 commit comments

Comments
 (0)