Skip to content

Commit ce998b2

Browse files
feat: add stricterFIFO option for queue fairness (#1100)
Co-authored-by: Cursor <cursoragent@cursor.com> Signed-off-by: huntiez <huntieezz@outlook.com>
1 parent 15033ef commit ce998b2

8 files changed

Lines changed: 144 additions & 5 deletions

File tree

README.md

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -463,6 +463,10 @@ This class extends [`EventEmitter`][] from Node.js.
463463
complete all in-flight tasks when `close()` is called. The default is `30000`
464464
- `recordTiming`: (`boolean`) By default, run and wait time will be recorded
465465
for the pool. To disable, set to `false`.
466+
- `stricterFIFO`: (`boolean`) When `true`, tasks that cannot be dispatched
467+
immediately are returned to the front of the queue instead of the back.
468+
This avoids head-of-line blocking under sustained load (especially with a
469+
single worker). Defaults to `false`.
466470

467471
Use caution when setting resource limits. Setting limits that are too low may
468472
result in the `Piscina` worker threads being unusable.

src/index.ts

Lines changed: 25 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ interface Options {
7070
trackUnmanagedFds? : boolean,
7171
closeTimeout?: number,
7272
recordTiming?: boolean,
73+
stricterFIFO?: boolean,
7374
loadBalancer?: PiscinaLoadBalancer,
7475
workerHistogram?: boolean,
7576
}
@@ -87,6 +88,7 @@ interface FilledOptions extends Options {
8788
niceIncrement : number,
8889
closeTimeout : number,
8990
recordTiming : boolean,
91+
stricterFIFO : boolean,
9092
workerHistogram: boolean,
9193
}
9294

@@ -122,6 +124,7 @@ const kDefaultOptions : FilledOptions = {
122124
trackUnmanagedFds: true,
123125
closeTimeout: 30000,
124126
recordTiming: true,
127+
stricterFIFO: false,
125128
workerHistogram: false
126129
};
127130

@@ -407,8 +410,8 @@ class ThreadPool {
407410
if (distributed) {
408411
// If task was distributed, we should continue to distribute more tasks
409412
continue;
410-
}
411-
413+
}
414+
412415
if (this.workers.size < this.options.maxThreads) {
413416
// We spawn if possible
414417
// TODO: scheduler will intercept this.
@@ -458,13 +461,27 @@ class ThreadPool {
458461
return true;
459462
}
460463

464+
this._returnUndistributedTask(task, !this.options.stricterFIFO);
465+
return false;
466+
}
467+
468+
_returnUndistributedTask (task: TaskInfo, toTail = false) : void {
461469
if (task.abortSignal != null) {
462-
this.skipQueue.push(task);
470+
if (toTail) {
471+
this.skipQueue.push(task);
472+
} else {
473+
this.skipQueue.unshift(task);
474+
}
475+
return;
476+
}
477+
478+
if (toTail) {
479+
this.taskQueue.push(task);
480+
} else if (this.taskQueue.unshift != null) {
481+
this.taskQueue.unshift(task);
463482
} else {
464483
this.taskQueue.push(task);
465484
}
466-
467-
return false;
468485
}
469486

470487
runTask (
@@ -787,6 +804,9 @@ export default class Piscina<Exports extends Record<string, (payload: any) => an
787804
if (opts.workerHistogram != null && (typeof opts.workerHistogram !== 'boolean')) {
788805
throw Errors.ValidationError('options.workerHistogram must be a boolean');
789806
}
807+
if (opts.stricterFIFO !== undefined && (typeof opts.stricterFIFO !== 'boolean')) {
808+
throw new TypeError('options.stricterFIFO must be a boolean');
809+
}
790810

791811
this.#pool = new ThreadPool(this, opts);
792812
}

src/task_queue/array_queue.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,10 @@ export class ArrayTaskQueue implements TaskQueue {
1717
this.tasks.push(task);
1818
}
1919

20+
unshift (task: Task): void {
21+
this.tasks.unshift(task);
22+
}
23+
2024
remove (task: Task): void {
2125
const index = this.tasks.indexOf(task);
2226
assert.notStrictEqual(index, -1);

src/task_queue/common.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ export interface TaskQueue {
55
shift(): Task | null;
66
remove(task: Task): void;
77
push(task: Task): void;
8+
unshift?(task: Task): void;
89
}
910

1011
// Public Interface

src/task_queue/fixed_queue.ts

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,11 @@ class FixedCircularBuffer {
7777
this.top = (this.top + 1) & kMask;
7878
}
7979

80+
unshift (data:Task) {
81+
this.bottom = (this.bottom - 1) & kMask;
82+
this.list[this.bottom] = data;
83+
}
84+
8085
shift () {
8186
const nextItem = this.list[this.bottom];
8287
if (nextItem === undefined) { return null; }
@@ -129,6 +134,16 @@ export class FixedQueue implements TaskQueue {
129134
this.#size++;
130135
}
131136

137+
unshift (data:Task) {
138+
if (this.tail.isFull()) {
139+
const newTail = new FixedCircularBuffer();
140+
newTail.next = this.tail;
141+
this.tail = newTail;
142+
}
143+
this.tail.unshift(data);
144+
this.#size++;
145+
}
146+
132147
shift (): Task | null {
133148
const tail = this.tail;
134149
const next = tail.shift();

test/fixtures/task-logger.js

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,8 @@
1+
const tasks = [];
2+
3+
module.exports = {
4+
run: task => {
5+
tasks.push(task);
6+
},
7+
getLog: () => tasks
8+
};

test/queues.test.ts

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
import * as assert from 'node:assert/strict';
2+
import { test } from 'node:test';
3+
import { kQueueOptions } from '../dist/symbols';
4+
import { FixedQueue, ArrayTaskQueue, PiscinaTask as Task } from '..';
5+
6+
// @ts-expect-error - it misses several properties, but it's enough for the test
7+
class NumberedTask implements Task {
8+
constructor(readonly index: number) {}
9+
get [kQueueOptions] () {
10+
return null;
11+
}
12+
}
13+
14+
for (const QueueClass of [FixedQueue, ArrayTaskQueue]) {
15+
test(`${QueueClass.name} - unshift`, () => {
16+
const bufferSize = 2048;
17+
let queue = new QueueClass();
18+
19+
queue.unshift(new NumberedTask(1));
20+
assert.equal(queue.size, 1);
21+
22+
for (let i = 2; i <= bufferSize + 10; ++i) {
23+
queue.unshift(new NumberedTask(i));
24+
}
25+
assert.equal(queue.size, bufferSize + 10);
26+
27+
for (let i = 0; i < 5; ++i) {
28+
assert.equal((queue.shift() as NumberedTask).index, bufferSize + 10 - i);
29+
}
30+
assert.equal(queue.size, bufferSize + 5);
31+
32+
for (let i = 0; i < bufferSize + 5; ++i) {
33+
assert.equal((queue.shift() as NumberedTask).index, bufferSize + 5 - i);
34+
}
35+
assert.equal(queue.size, 0);
36+
37+
queue = new QueueClass();
38+
queue.push(new NumberedTask(1));
39+
queue.push(new NumberedTask(2));
40+
queue.push(new NumberedTask(3));
41+
queue.push(new NumberedTask(4));
42+
queue.push(new NumberedTask(5));
43+
assert.equal((queue.shift() as NumberedTask).index, 1);
44+
assert.equal((queue.shift() as NumberedTask).index, 2);
45+
assert.equal((queue.shift() as NumberedTask).index, 3);
46+
assert.equal(queue.size, 2);
47+
queue.unshift(new NumberedTask(6));
48+
queue.unshift(new NumberedTask(7));
49+
queue.unshift(new NumberedTask(8));
50+
queue.unshift(new NumberedTask(9));
51+
assert.equal((queue.shift() as NumberedTask).index, 9);
52+
assert.equal((queue.shift() as NumberedTask).index, 8);
53+
assert.equal((queue.shift() as NumberedTask).index, 7);
54+
assert.equal((queue.shift() as NumberedTask).index, 6);
55+
assert.equal((queue.shift() as NumberedTask).index, 4);
56+
assert.equal((queue.shift() as NumberedTask).index, 5);
57+
});
58+
}

test/stricter-fifo.test.ts

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,29 @@
1+
import * as assert from 'node:assert/strict';
2+
import { resolve } from 'node:path';
3+
import { test } from 'node:test';
4+
import Piscina from '..';
5+
6+
test('stricterFIFO keeps queued tasks in submission order', async () => {
7+
const pool = new Piscina({
8+
filename: resolve(__dirname, 'fixtures/task-logger.js'),
9+
minThreads: 1,
10+
maxThreads: 1,
11+
stricterFIFO: true
12+
});
13+
14+
const tasks = new Array(20).fill(0).map((_, i) => ({ i }));
15+
const pending = tasks.map((task) => pool.run(task, { name: 'run' }));
16+
await Promise.all(pending);
17+
18+
const log = await pool.run({}, { name: 'getLog' });
19+
assert.deepEqual(log, tasks);
20+
21+
await pool.destroy();
22+
});
23+
24+
test('stricterFIFO must be a boolean', () => {
25+
assert.throws(
26+
() => new Piscina({ stricterFIFO: 1 } as any),
27+
/options.stricterFIFO must be a boolean/
28+
);
29+
});

0 commit comments

Comments
 (0)