fix(core): Fix unhandled rejection in task broker on runner disconnect (#27278)

This commit is contained in:
Iván Ovejero
2026-03-20 09:19:41 +00:00
committed by GitHub
parent e60d9e7f39
commit 6fcc86037d
4 changed files with 30 additions and 2 deletions
@@ -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', () => {
@@ -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);
@@ -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);
});
}
}
@@ -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) {