From 568dba2c8f84f79ddf6a34ef647382f80ae0e47e Mon Sep 17 00:00:00 2001 From: Marc Littlemore Date: Mon, 15 Dec 2025 13:06:22 +0000 Subject: [PATCH] feat(core): Add workflow cancellation events to log streaming (#23151) --- packages/cli/src/active-executions.ts | 9 ++++++- .../eventbus/event-message-classes/index.ts | 1 + .../log-streaming-event-relay.test.ts | 24 +++++++++++++++++++ .../cli/src/events/maps/relay.event-map.ts | 4 ++++ .../relays/log-streaming.event-relay.ts | 18 ++++++++++++++ packages/cli/src/scaling/job-processor.ts | 9 +++++-- .../src/errors/execution-cancelled.error.ts | 13 ++++++---- packages/workflow/src/errors/index.ts | 1 + 8 files changed, 72 insertions(+), 7 deletions(-) diff --git a/packages/cli/src/active-executions.ts b/packages/cli/src/active-executions.ts index 73641be3d0d..3ed36a280e8 100644 --- a/packages/cli/src/active-executions.ts +++ b/packages/cli/src/active-executions.ts @@ -178,7 +178,14 @@ export class ActiveExecutions { // There is no execution running with that id return; } - this.eventService.emit('execution-cancelled', { executionId }); + + const workflowData = execution.executionData.workflowData; + this.eventService.emit('execution-cancelled', { + executionId, + workflowId: workflowData?.id, + workflowName: workflowData?.name, + reason: cancellationError.reason, + }); execution.responsePromise?.reject(cancellationError); if (execution.status === 'waiting') { // A waiting execution will not have a valid workflowExecution or postExecutePromise diff --git a/packages/cli/src/eventbus/event-message-classes/index.ts b/packages/cli/src/eventbus/event-message-classes/index.ts index 059b90686db..e02f07bc010 100644 --- a/packages/cli/src/eventbus/event-message-classes/index.ts +++ b/packages/cli/src/eventbus/event-message-classes/index.ts @@ -47,6 +47,7 @@ export const eventNamesWorkflow = [ 'n8n.workflow.started', 'n8n.workflow.success', 'n8n.workflow.failed', + 'n8n.workflow.cancelled', ] as const; export const eventNamesGeneric = ['n8n.worker.started', 'n8n.worker.stopped'] as const; export const eventNamesNode = ['n8n.node.started', 'n8n.node.finished'] as const; diff --git a/packages/cli/src/events/__tests__/log-streaming-event-relay.test.ts b/packages/cli/src/events/__tests__/log-streaming-event-relay.test.ts index 9eb4f2f2a70..af4b567144f 100644 --- a/packages/cli/src/events/__tests__/log-streaming-event-relay.test.ts +++ b/packages/cli/src/events/__tests__/log-streaming-event-relay.test.ts @@ -1083,6 +1083,30 @@ describe('LogStreamingEventRelay', () => { }, }); }); + + it.each(['manual', 'timeout', 'shutdown'] as const)( + 'should log on `execution-cancelled` event with %s reason', + (reason) => { + const event: RelayEventMap['execution-cancelled'] = { + executionId: 'exec-cancelled-123', + workflowId: 'wf-456', + workflowName: 'Cancelled Workflow', + reason, + }; + + eventService.emit('execution-cancelled', event); + + expect(eventBus.sendWorkflowEvent).toHaveBeenCalledWith({ + eventName: 'n8n.workflow.cancelled', + payload: { + executionId: 'exec-cancelled-123', + workflowId: 'wf-456', + workflowName: 'Cancelled Workflow', + reason, + }, + }); + }, + ); }); describe('AI events', () => { diff --git a/packages/cli/src/events/maps/relay.event-map.ts b/packages/cli/src/events/maps/relay.event-map.ts index 26a2813f7ef..6231a973473 100644 --- a/packages/cli/src/events/maps/relay.event-map.ts +++ b/packages/cli/src/events/maps/relay.event-map.ts @@ -1,6 +1,7 @@ import type { AuthenticationMethod, ProjectRelation } from '@n8n/api-types'; import type { AuthProviderType, User, IWorkflowDb } from '@n8n/db'; import type { + CancellationReason, IPersonalizationSurveyAnswersV4, IRun, IWorkflowBase, @@ -395,6 +396,9 @@ export type RelayEventMap = { 'execution-cancelled': { executionId: string; + workflowId?: string; + workflowName?: string; + reason: CancellationReason; }; // #endregion diff --git a/packages/cli/src/events/relays/log-streaming.event-relay.ts b/packages/cli/src/events/relays/log-streaming.event-relay.ts index afd108f2a17..c68034f8db5 100644 --- a/packages/cli/src/events/relays/log-streaming.event-relay.ts +++ b/packages/cli/src/events/relays/log-streaming.event-relay.ts @@ -51,6 +51,7 @@ export class LogStreamingEventRelay extends EventRelay { 'community-package-deleted': (event) => this.communityPackageDeleted(event), 'execution-throttled': (event) => this.executionThrottled(event), 'execution-started-during-bootup': (event) => this.executionStartedDuringBootup(event), + 'execution-cancelled': (event) => this.executionCancelled(event), 'ai-messages-retrieved-from-memory': (event) => this.aiMessagesRetrievedFromMemory(event), 'ai-message-added-to-memory': (event) => this.aiMessageAddedToMemory(event), 'ai-output-parsed': (event) => this.aiOutputParsed(event), @@ -465,6 +466,23 @@ export class LogStreamingEventRelay extends EventRelay { }); } + private executionCancelled({ + executionId, + workflowId, + workflowName, + reason, + }: RelayEventMap['execution-cancelled']) { + void this.eventBus.sendWorkflowEvent({ + eventName: 'n8n.workflow.cancelled', + payload: { + executionId, + workflowId, + workflowName, + reason, + }, + }); + } + // #endregion // #region AI diff --git a/packages/cli/src/scaling/job-processor.ts b/packages/cli/src/scaling/job-processor.ts index cff52c55bd0..a1741a18481 100644 --- a/packages/cli/src/scaling/job-processor.ts +++ b/packages/cli/src/scaling/job-processor.ts @@ -287,8 +287,13 @@ export class JobProcessor { const runningJob = this.runningJobs[jobId]; if (!runningJob) return; - const executionId = runningJob.executionId; - this.eventService.emit('execution-cancelled', { executionId }); + const { executionId, workflowId, workflowName } = runningJob; + this.eventService.emit('execution-cancelled', { + executionId, + workflowId, + workflowName, + reason: 'manual', // Job stops via scaling service are always user-initiated + }); runningJob.run.cancel(); delete this.runningJobs[jobId]; diff --git a/packages/workflow/src/errors/execution-cancelled.error.ts b/packages/workflow/src/errors/execution-cancelled.error.ts index 4f11348e3e4..deb949d332a 100644 --- a/packages/workflow/src/errors/execution-cancelled.error.ts +++ b/packages/workflow/src/errors/execution-cancelled.error.ts @@ -1,32 +1,37 @@ import { ExecutionBaseError } from './abstract/execution-base.error'; +export type CancellationReason = 'manual' | 'timeout' | 'shutdown'; + export abstract class ExecutionCancelledError extends ExecutionBaseError { + readonly reason: CancellationReason; + // NOTE: prefer one of the more specific - constructor(executionId: string) { + constructor(executionId: string, reason: CancellationReason) { super('The execution was cancelled', { level: 'warning', extra: { executionId }, }); + this.reason = reason; } } export class ManualExecutionCancelledError extends ExecutionCancelledError { constructor(executionId: string) { - super(executionId); + super(executionId, 'manual'); this.message = 'The execution was cancelled manually'; } } export class TimeoutExecutionCancelledError extends ExecutionCancelledError { constructor(executionId: string) { - super(executionId); + super(executionId, 'timeout'); this.message = 'The execution was cancelled because it timed out'; } } export class SystemShutdownExecutionCancelledError extends ExecutionCancelledError { constructor(executionId: string) { - super(executionId); + super(executionId, 'shutdown'); this.message = 'The execution was cancelled because the system is shutting down'; } } diff --git a/packages/workflow/src/errors/index.ts b/packages/workflow/src/errors/index.ts index a0824e88e6f..4345f4fdfec 100644 --- a/packages/workflow/src/errors/index.ts +++ b/packages/workflow/src/errors/index.ts @@ -9,6 +9,7 @@ export { ManualExecutionCancelledError, SystemShutdownExecutionCancelledError, TimeoutExecutionCancelledError, + type CancellationReason, } from './execution-cancelled.error'; export { NodeApiError } from './node-api.error'; export { NodeOperationError } from './node-operation.error';