From fac005b1654530f7162d448ba7bde00b50a22392 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Iv=C3=A1n=20Ovejero?= Date: Thu, 25 Sep 2025 11:24:03 +0200 Subject: [PATCH] fix(core): Ensure cancellation interrupts runner tasks in worker (#19864) --- .../scaling/__tests__/job-processor.service.test.ts | 4 ++++ packages/cli/src/scaling/job-processor.ts | 10 +++++++++- 2 files changed, 13 insertions(+), 1 deletion(-) diff --git a/packages/cli/src/scaling/__tests__/job-processor.service.test.ts b/packages/cli/src/scaling/__tests__/job-processor.service.test.ts index 38c854d6cc4..0a7663adae2 100644 --- a/packages/cli/src/scaling/__tests__/job-processor.service.test.ts +++ b/packages/cli/src/scaling/__tests__/job-processor.service.test.ts @@ -72,6 +72,7 @@ describe('JobProcessor', () => { mock(), mock(), executionsConfig, + mock(), ); const result = await jobProcessor.processJob(mock()); @@ -102,6 +103,7 @@ describe('JobProcessor', () => { mock(), manualExecutionService, executionsConfig, + mock(), ); await jobProcessor.processJob(mock()); @@ -141,6 +143,7 @@ describe('JobProcessor', () => { mock(), manualExecutionService, executionsConfig, + mock(), ); const executionId = 'execution-id'; @@ -204,6 +207,7 @@ describe('JobProcessor', () => { mock(), manualExecutionService, executionsConfig, + mock(), ); await jobProcessor.processJob(mock()); diff --git a/packages/cli/src/scaling/job-processor.ts b/packages/cli/src/scaling/job-processor.ts index aacb809cbaf..25aaab17c87 100644 --- a/packages/cli/src/scaling/job-processor.ts +++ b/packages/cli/src/scaling/job-processor.ts @@ -14,6 +14,7 @@ import type { import { BINARY_ENCODING, Workflow, UnexpectedError } from 'n8n-workflow'; import type PCancelable from 'p-cancelable'; +import { EventService } from '@/events/event.service'; import { getLifecycleHooksForScalingWorker } from '@/execution-lifecycle/execution-lifecycle-hooks'; import { ManualExecutionService } from '@/manual-execution.service'; import { NodeTypes } from '@/node-types'; @@ -44,6 +45,7 @@ export class JobProcessor { private readonly instanceSettings: InstanceSettings, private readonly manualExecutionService: ManualExecutionService, private readonly executionsConfig: ExecutionsConfig, + private readonly eventService: EventService, ) { this.logger = this.logger.scoped('scaling'); } @@ -269,7 +271,13 @@ export class JobProcessor { } stopJob(jobId: JobId) { - this.runningJobs[jobId]?.run.cancel(); + const runningJob = this.runningJobs[jobId]; + if (!runningJob) return; + + const executionId = runningJob.executionId; + this.eventService.emit('execution-cancelled', { executionId }); + + runningJob.run.cancel(); delete this.runningJobs[jobId]; }