implemented p-queue for use as a taskQueue

Signed-off-by: Hao Luo <howlowck@gmail.com>
This commit is contained in:
Hao Luo
2022-10-06 12:46:47 -07:00
parent 0c55e847a4
commit 19ecf55817
3 changed files with 25 additions and 37 deletions
+1
View File
@@ -72,6 +72,7 @@
"octokit": "^2.0.0",
"octokit-plugin-create-pull-request": "^3.10.0",
"p-limit": "^3.1.0",
"p-queue": "^7.3.0",
"uuid": "^8.2.0",
"vm2": "^3.9.11",
"winston": "^3.2.1",
@@ -15,6 +15,7 @@
*/
import { TaskContext, TaskBroker, WorkflowRunner } from './types';
import PQueue from 'p-queue';
import { NunjucksWorkflowRunner } from './NunjucksWorkflowRunner';
import { Logger } from 'winston';
import { TemplateActionRegistry } from '../actions';
@@ -54,15 +55,6 @@ export type CreateWorkerOptions = {
additionalTemplateGlobals?: Record<string, TemplateGlobal>;
};
// Same implementation as StorageTaskBroker
function makeLock() {
let unlock = () => {};
const promise = new Promise<void>(_ => {
unlock = _;
});
return { promise, unlock };
}
/**
* TaskWorker
*
@@ -71,28 +63,9 @@ function makeLock() {
export class TaskWorker {
private constructor(private readonly options: TaskWorkerOptions) {}
private runningTasks: number = 0;
private workerLock = makeLock();
private get isWorkerAvailable() {
return this.runningTasks < this.options.concurrentTasksLimit;
}
private completeTask() {
this.runningTasks--;
if (this.runningTasks === this.options.concurrentTasksLimit - 1) {
this.workerLock.unlock();
this.workerLock = makeLock();
}
}
private async waitWorkerToBeAvailable() {
if (this.isWorkerAvailable) {
return;
}
await this.workerLock.promise;
}
private taskQueue: PQueue = new PQueue({
concurrency: this.options.concurrentTasksLimit,
});
static async create(options: CreateWorkerOptions): Promise<TaskWorker> {
const {
@@ -125,10 +98,8 @@ export class TaskWorker {
start() {
(async () => {
for (;;) {
await this.waitWorkerToBeAvailable();
const task = await this.options.taskBroker.claim();
this.runningTasks++;
this.runOneTask(task);
this.taskQueue.add(() => this.runOneTask(task));
}
})();
}
@@ -151,8 +122,6 @@ export class TaskWorker {
await task.complete('failed', {
error: { name: error.name, message: error.message },
});
} finally {
this.completeTask();
}
}
}
+19 -1
View File
@@ -6680,6 +6680,7 @@ __metadata:
octokit: ^2.0.0
octokit-plugin-create-pull-request: ^3.10.0
p-limit: ^3.1.0
p-queue: ^7.3.0
supertest: ^6.1.3
uuid: ^8.2.0
vm2: ^3.9.11
@@ -21960,7 +21961,7 @@ __metadata:
languageName: node
linkType: hard
"eventemitter3@npm:^4.0.0, eventemitter3@npm:^4.0.1, eventemitter3@npm:^4.0.4":
"eventemitter3@npm:^4.0.0, eventemitter3@npm:^4.0.1, eventemitter3@npm:^4.0.4, eventemitter3@npm:^4.0.7":
version: 4.0.7
resolution: "eventemitter3@npm:4.0.7"
checksum: 1875311c42fcfe9c707b2712c32664a245629b42bb0a5a84439762dd0fd637fc54d078155ea83c2af9e0323c9ac13687e03cfba79b03af9f40c89b4960099374
@@ -31128,6 +31129,16 @@ __metadata:
languageName: node
linkType: hard
"p-queue@npm:^7.3.0":
version: 7.3.0
resolution: "p-queue@npm:7.3.0"
dependencies:
eventemitter3: ^4.0.7
p-timeout: ^5.0.2
checksum: eda8b8d4dac9456d1722c8ee0892bbb2e7f7b79d1f90b6cc1bac5fb9dcf3120285f90c28a68dc2f66bc9e46c7d7d42ad206a2c4d83e7aad27a0226ae042081ee
languageName: node
linkType: hard
"p-retry@npm:^4.5.0":
version: 4.5.0
resolution: "p-retry@npm:4.5.0"
@@ -31147,6 +31158,13 @@ __metadata:
languageName: node
linkType: hard
"p-timeout@npm:^5.0.2":
version: 5.1.0
resolution: "p-timeout@npm:5.1.0"
checksum: f5cd4e17301ff1ff1d8dbf2817df0ad88c6bba99349fc24d8d181827176ad4f8aca649190b8a5b1a428dfd6ddc091af4606835d3e0cb0656e04045da5c9e270c
languageName: node
linkType: hard
"p-transform@npm:^1.3.0":
version: 1.3.0
resolution: "p-transform@npm:1.3.0"