From 6fcc86037d1bb46a0e68688bd9d8a76db7d16560 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Iv=C3=A1n=20Ovejero?= Date: Fri, 20 Mar 2026 10:19:41 +0100 Subject: [PATCH] fix(core): Fix unhandled rejection in task broker on runner disconnect (#27278) --- .../__tests__/task-broker.service.test.ts | 22 +++++++++++++++++++ .../task-broker/task-broker.service.ts | 4 ++++ .../task-managers/local-task-requester.ts | 4 +++- .../task-managers/task-requester.ts | 2 +- 4 files changed, 30 insertions(+), 2 deletions(-) diff --git a/packages/cli/src/task-runners/task-broker/__tests__/task-broker.service.test.ts b/packages/cli/src/task-runners/task-broker/__tests__/task-broker.service.test.ts index 1adfb1be6d0..b7d2a3c37cb 100644 --- a/packages/cli/src/task-runners/task-broker/__tests__/task-broker.service.test.ts +++ b/packages/cli/src/task-runners/task-broker/__tests__/task-broker.service.test.ts @@ -714,6 +714,28 @@ describe('TaskBroker', () => { nodeTypes, }); }); + + it('should discard `requester:rpcresponse` for an already-cleaned-up task', async () => { + await expect( + taskBroker.handleRequesterRpcResponse('nonexistent', 'call1', 'success', {}), + ).resolves.toBeUndefined(); + }); + + it('should discard `requester:taskdataresponse` for an already-cleaned-up task', async () => { + await expect( + taskBroker.handleRequesterDataResponse('nonexistent', 'req1', {}), + ).resolves.toBeUndefined(); + }); + + it('should discard `requester:nodetypesresponse` for an already-cleaned-up task', async () => { + await expect( + taskBroker.handleRequesterNodeTypesResponse('nonexistent', 'req1', []), + ).resolves.toBeUndefined(); + }); + + it('should discard `sendTaskSettings` for an already-cleaned-up task', async () => { + await expect(taskBroker.sendTaskSettings('nonexistent', {})).resolves.toBeUndefined(); + }); }); describe('task execution timeouts', () => { diff --git a/packages/cli/src/task-runners/task-broker/task-broker.service.ts b/packages/cli/src/task-runners/task-broker/task-broker.service.ts index 049dc8aa74a..abd1d92b783 100644 --- a/packages/cli/src/task-runners/task-broker/task-broker.service.ts +++ b/packages/cli/src/task-runners/task-broker/task-broker.service.ts @@ -363,6 +363,7 @@ export class TaskBroker { status: RequesterMessage.ToBroker.RPCResponse['status'], data: unknown, ) { + if (!this.tasks.has(taskId)) return; const runner = await this.getRunnerOrFailTask(taskId); await this.messageRunner(runner.id, { type: 'broker:rpcresponse', @@ -374,6 +375,7 @@ export class TaskBroker { } async handleRequesterDataResponse(taskId: Task['id'], requestId: string, data: unknown) { + if (!this.tasks.has(taskId)) return; const runner = await this.getRunnerOrFailTask(taskId); await this.messageRunner(runner.id, { @@ -389,6 +391,7 @@ export class TaskBroker { requestId: RequesterMessage.ToBroker.NodeTypesResponse['requestId'], nodeTypes: RequesterMessage.ToBroker.NodeTypesResponse['nodeTypes'], ) { + if (!this.tasks.has(taskId)) return; const runner = await this.getRunnerOrFailTask(taskId); await this.messageRunner(runner.id, { @@ -463,6 +466,7 @@ export class TaskBroker { } async sendTaskSettings(taskId: Task['id'], settings: unknown) { + if (!this.tasks.has(taskId)) return; const runner = await this.getRunnerOrFailTask(taskId); const task = this.tasks.get(taskId); diff --git a/packages/cli/src/task-runners/task-managers/local-task-requester.ts b/packages/cli/src/task-runners/task-managers/local-task-requester.ts index c9b666bcec9..e5fd3d7c52c 100644 --- a/packages/cli/src/task-runners/task-managers/local-task-requester.ts +++ b/packages/cli/src/task-runners/task-managers/local-task-requester.ts @@ -37,6 +37,8 @@ export class LocalTaskRequester extends TaskRequester { } sendMessage(message: RequesterMessage.ToBroker.All) { - void this.taskBroker.onRequesterMessage(this.id, message); + void this.taskBroker.onRequesterMessage(this.id, message).catch((error) => { + this.errorReporter.error(error); + }); } } diff --git a/packages/cli/src/task-runners/task-managers/task-requester.ts b/packages/cli/src/task-runners/task-managers/task-requester.ts index a38dbae9c38..e2288d9286e 100644 --- a/packages/cli/src/task-runners/task-managers/task-requester.ts +++ b/packages/cli/src/task-runners/task-managers/task-requester.ts @@ -72,7 +72,7 @@ export abstract class TaskRequester { private readonly eventService: EventService, private readonly taskRunnersConfig: TaskRunnersConfig, private readonly globalConfig: GlobalConfig, - private readonly errorReporter: ErrorReporter, + protected readonly errorReporter: ErrorReporter, ) {} setRunnerUnavailable(taskType: string, reason: string) {