diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 9defd71c49..b18f06f1f6 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -32,6 +32,7 @@ export type TaskWorkerOptions = { runners: { workflowRunner: WorkflowRunner; }; + concurrentTasksLimit: number; }; /** @@ -46,8 +47,18 @@ export type CreateWorkerOptions = { workingDirectory: string; logger: Logger; additionalTemplateFilters?: Record; + concurrentTasksLimit?: number; }; +// Same implementation as StorageTaskBroker +function makeLock() { + let unlock = () => {}; + const promise = new Promise(_ => { + unlock = _; + }); + return { promise, unlock }; +} + /** * TaskWorker * @@ -56,6 +67,29 @@ export type CreateWorkerOptions = { 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; + } + static async create(options: CreateWorkerOptions): Promise { const { taskBroker, @@ -64,6 +98,7 @@ export class TaskWorker { integrations, workingDirectory, additionalTemplateFilters, + concurrentTasksLimit = 10, // Or Infinity } = options; const workflowRunner = new NunjucksWorkflowRunner({ @@ -77,14 +112,17 @@ export class TaskWorker { return new TaskWorker({ taskBroker: taskBroker, runners: { workflowRunner }, + concurrentTasksLimit, }); } start() { (async () => { for (;;) { + await this.waitWorkerToBeAvailable(); const task = await this.options.taskBroker.claim(); - await this.runOneTask(task); + this.runningTasks++; + this.runOneTask(task); } })(); } @@ -107,6 +145,8 @@ export class TaskWorker { await task.complete('failed', { error: { name: error.name, message: error.message }, }); + } finally { + this.completeTask(); } } } diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 28dd1ca413..3fbfe24074 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -199,7 +199,7 @@ export async function createRouter( const actionRegistry = new TemplateActionRegistry(); const workers = []; - for (let i = 0; i < (taskWorkers || 3); i++) { + for (let i = 0; i < (taskWorkers || 1); i++) { const worker = await TaskWorker.create({ taskBroker, actionRegistry,