|
1 | 1 | import type { TaskRunnersConfig } from '@n8n/config';
|
2 | 2 | import type { RunnerMessage, TaskResultData } from '@n8n/task-runner';
|
3 | 3 | import { mock } from 'jest-mock-extended';
|
| 4 | +import type { Logger } from 'n8n-core'; |
4 | 5 | import { ApplicationError, type INodeTypeBaseDescription } from 'n8n-workflow';
|
5 | 6 |
|
6 | 7 | import { Time } from '@/constants';
|
7 | 8 | import type { TaskRunnerLifecycleEvents } from '@/task-runners/task-runner-lifecycle-events';
|
8 | 9 |
|
9 | 10 | import { TaskRejectError } from '../errors/task-reject.error';
|
10 |
| -import { TaskRunnerTimeoutError } from '../errors/task-runner-timeout.error'; |
| 11 | +import { TaskRunnerExecutionTimeoutError } from '../errors/task-runner-execution-timeout.error'; |
11 | 12 | import { TaskBroker } from '../task-broker.service';
|
12 | 13 | import type { TaskOffer, TaskRequest, TaskRunner } from '../task-broker.service';
|
13 | 14 |
|
@@ -715,7 +716,7 @@ describe('TaskBroker', () => {
|
715 | 716 | });
|
716 | 717 | });
|
717 | 718 |
|
718 |
| - describe('task timeouts', () => { |
| 719 | + describe('task execution timeouts', () => { |
719 | 720 | let taskBroker: TaskBroker;
|
720 | 721 | let config: TaskRunnersConfig;
|
721 | 722 | let runnerLifecycleEvents = mock<TaskRunnerLifecycleEvents>();
|
@@ -879,11 +880,54 @@ describe('TaskBroker', () => {
|
879 | 880 | expect(requesterCallback).toHaveBeenCalledWith({
|
880 | 881 | type: 'broker:taskerror',
|
881 | 882 | taskId,
|
882 |
| - error: expect.any(TaskRunnerTimeoutError), |
| 883 | + error: expect.any(TaskRunnerExecutionTimeoutError), |
883 | 884 | });
|
884 | 885 |
|
885 | 886 | expect(clearTimeout).toHaveBeenCalled();
|
886 | 887 | expect(taskBroker.getTasks().get(taskId)).toBeUndefined();
|
887 | 888 | });
|
888 | 889 | });
|
| 890 | + |
| 891 | + describe('task runner accept timeout', () => { |
| 892 | + it('broker should handle timeout when waiting for acknowledgment of offer accept', async () => { |
| 893 | + const runnerId = 'runner1'; |
| 894 | + const runner = mock<TaskRunner>({ id: runnerId }); |
| 895 | + const messageCallback = jest.fn(); |
| 896 | + const loggerMock = mock<Logger>(); |
| 897 | + |
| 898 | + taskBroker = new TaskBroker(loggerMock, mock(), mock()); |
| 899 | + taskBroker.registerRunner(runner, messageCallback); |
| 900 | + |
| 901 | + const offer: TaskOffer = { |
| 902 | + offerId: 'offer1', |
| 903 | + runnerId, |
| 904 | + taskType: 'taskType1', |
| 905 | + validFor: 1000, |
| 906 | + validUntil: createValidUntil(1000), |
| 907 | + }; |
| 908 | + |
| 909 | + const request: TaskRequest = { |
| 910 | + requestId: 'request1', |
| 911 | + requesterId: 'requester1', |
| 912 | + taskType: 'taskType1', |
| 913 | + }; |
| 914 | + |
| 915 | + jest.useFakeTimers(); |
| 916 | + |
| 917 | + const acceptPromise = taskBroker.acceptOffer(offer, request); |
| 918 | + |
| 919 | + jest.advanceTimersByTime(2100); |
| 920 | + |
| 921 | + await acceptPromise; |
| 922 | + |
| 923 | + expect(request.acceptInProgress).toBe(false); |
| 924 | + expect(loggerMock.warn).toHaveBeenCalledWith( |
| 925 | + expect.stringContaining( |
| 926 | + `Runner (${runnerId}) took too long to acknowledge acceptance of task`, |
| 927 | + ), |
| 928 | + ); |
| 929 | + |
| 930 | + jest.useRealTimers(); |
| 931 | + }); |
| 932 | + }); |
889 | 933 | });
|
0 commit comments