diff --git a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts index 5e116922de..cd7dc06a91 100644 --- a/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts +++ b/plugins/scaffolder-backend/src/scaffolder/tasks/TaskWorker.ts @@ -76,8 +76,10 @@ export type CreateWorkerOptions = { export class TaskWorker { private taskQueue: PQueue; private logger: Logger | undefined; + private stopWorkers: boolean; private constructor(private readonly options: TaskWorkerOptions) { + this.stopWorkers = false; this.logger = options.logger; this.taskQueue = new PQueue({ concurrency: options.concurrentTasksLimit, @@ -120,26 +122,31 @@ export class TaskWorker { await this.options.taskBroker.recoverTasks?.(); } catch (err) { this.logger?.error(stringifyError(err)); - // ignore } } start() { (async () => { - for (;;) { + while (!this.stopWorkers) { await new Promise(resolve => setTimeout(resolve, 10000)); await this.recoverTasks(); } })(); (async () => { - for (;;) { + while (!this.stopWorkers) { await this.onReadyToClaimTask(); - const task = await this.options.taskBroker.claim(); - this.taskQueue.add(() => this.runOneTask(task)); + if (!this.stopWorkers) { + const task = await this.options.taskBroker.claim(); + void this.taskQueue.add(() => this.runOneTask(task)); + } } })(); } + stop() { + this.stopWorkers = true; + } + protected onReadyToClaimTask(): Promise { if (this.taskQueue.pending < this.options.concurrentTasksLimit) { return Promise.resolve(); diff --git a/plugins/scaffolder-backend/src/service/router.ts b/plugins/scaffolder-backend/src/service/router.ts index 365ec8748d..c0322af283 100644 --- a/plugins/scaffolder-backend/src/service/router.ts +++ b/plugins/scaffolder-backend/src/service/router.ts @@ -336,8 +336,13 @@ export async function createRouter( const launchWorkers = () => workers.forEach(worker => worker.start()); + const shutdownWorkers = () => { + workers.forEach(worker => worker.stop()); + }; + if (options.lifecycle) { options.lifecycle.addStartupHook(launchWorkers); + options.lifecycle.addShutdownHook(shutdownWorkers); } else { launchWorkers(); }