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) {