feat(core): Cancel runner task on timeout in external mode (#12101)

This commit is contained in:
Iván Ovejero authored and GitHub committed 2024-12-10 12:50:22 +01:00
1 parent a63f0e878e
commit addb4fa352
12 files changed
+283 -34

No files matched your search

@@ -6,6 +6,7 @@ import { ApplicationError, type INodeTypeBaseDescription } from 'n8n-workflow';
import { Time } from '@/constants';
import { TaskRejectError } from '../errors';
import { TaskRunnerTimeoutError } from '../errors/task-runner-timeout.error';
import type { RunnerLifecycleEvents } from '../runner-lifecycle-events';
import { TaskBroker } from '../task-broker.service';
import type { TaskOffer, TaskRequest, TaskRunner } from '../task-broker.service';
@@ -721,7 +722,7 @@ describe('TaskBroker', () => {
beforeAll(() => {
jest.useFakeTimers();
config = mock<TaskRunnersConfig>({ taskTimeout: 30 });
config = mock<TaskRunnersConfig>({ taskTimeout: 30, mode: 'internal' });
taskBroker = new TaskBroker(mock(), config, runnerLifecycleEvents);
});
@@ -800,7 +801,7 @@ describe('TaskBroker', () => {
expect(taskBroker.getTasks().get(taskId)).toBeUndefined();
});
it('on timeout, we should emit `runner:timed-out-during-task` event and send error to requester', async () => {
it('[internal mode] on timeout, we should emit `runner:timed-out-during-task` event and send error to requester', async () => {
jest.spyOn(global, 'clearTimeout');
const taskId = 'task1';
@@ -839,5 +840,50 @@ describe('TaskBroker', () => {
expect(taskBroker.getTasks().get(taskId)).toBeUndefined();
});
it('[external mode] on timeout, we should instruct the runner to cancel and send error to requester', async () => {
const config = mock<TaskRunnersConfig>({ taskTimeout: 30, mode: 'external' });
taskBroker = new TaskBroker(mock(), config, runnerLifecycleEvents);
jest.spyOn(global, 'clearTimeout');
const taskId = 'task1';
const runnerId = 'runner1';
const requesterId = 'requester1';
const runner = mock<TaskRunner>({ id: runnerId });
const runnerCallback = jest.fn();
const requesterCallback = jest.fn();
taskBroker.registerRunner(runner, runnerCallback);
taskBroker.registerRequester(requesterId, requesterCallback);
taskBroker.setTasks({
[taskId]: { id: taskId, runnerId, requesterId, taskType: 'test' },
});
await taskBroker.sendTaskSettings(taskId, {});
runnerCallback.mockClear();
jest.runAllTimers();
await Promise.resolve(); // for timeout callback
await Promise.resolve(); // for sending messages to runner and requester
await Promise.resolve(); // for task cleanup and removal
expect(runnerCallback).toHaveBeenLastCalledWith({
type: 'broker:taskcancel',
taskId,
reason: 'Task execution timed out',
});
expect(requesterCallback).toHaveBeenCalledWith({
type: 'broker:taskerror',
taskId,
error: expect.any(TaskRunnerTimeoutError),
});
expect(clearTimeout).toHaveBeenCalled();
expect(taskBroker.getTasks().get(taskId)).toBeUndefined();
});
});
});
@@ -1,15 +1,23 @@
import type { TaskRunnerMode } from '@n8n/config/src/configs/runners.config';
import { ApplicationError } from 'n8n-workflow';
export class TaskRunnerTimeoutError extends ApplicationError {
description: string;
constructor(taskTimeout: number, isSelfHosted: boolean) {
constructor({
taskTimeout,
isSelfHosted,
mode,
}: { taskTimeout: number; isSelfHosted: boolean; mode: TaskRunnerMode }) {
super(
`Task execution timed out after ${taskTimeout} ${taskTimeout === 1 ? 'second' : 'seconds'}`,
);
const subtitle =
'The task runner was taking too long on this task, so it was suspected of being unresponsive and restarted, and the task was aborted. You can try the following:';
const subtitles = {
internal:
'The task runner was taking too long on this task, so it was suspected of being unresponsive and restarted, and the task was aborted.',
external: 'The task runner was taking too long on this task, so the task was aborted.',
};
const fixes = {
optimizeScript:
@@ -27,7 +35,7 @@ export class TaskRunnerTimeoutError extends ApplicationError {
.map((suggestion, index) => `${index + 1}. ${suggestion}`)
.join('<br/>');
const description = `${subtitle}<br/><br/>${suggestionsText}`;
const description = `${mode === 'internal' ? subtitles.internal : subtitles.external} You can try the following:<br/><br/>${suggestionsText}`;
this.description = description;
}
@@ -459,14 +459,25 @@ export class TaskBroker {
const task = this.tasks.get(taskId);
if (!task) return;
this.runnerLifecycleEvents.emit('runner:timed-out-during-task');
if (this.taskRunnersConfig.mode === 'internal') {
this.runnerLifecycleEvents.emit('runner:timed-out-during-task');
} else if (this.taskRunnersConfig.mode === 'external') {
await this.messageRunner(task.runnerId, {
type: 'broker:taskcancel',
taskId,
reason: 'Task execution timed out',
});
}
const { taskTimeout, mode } = this.taskRunnersConfig;
await this.taskErrorHandler(
taskId,
new TaskRunnerTimeoutError(
this.taskRunnersConfig.taskTimeout,
config.getEnv('deployment.type') !== 'cloud',
),
new TaskRunnerTimeoutError({
taskTimeout,
isSelfHosted: config.getEnv('deployment.type') !== 'cloud',
mode,
}),
);
}