From 6eb19e1668fd403b1bb7de9f31be9f31bebcdee2 Mon Sep 17 00:00:00 2001 From: Tomi Turtiainen <10324676+tomi@users.noreply.github.com> Date: Thu, 27 Aug 2026 15:32:25 +0000 Subject: [PATCH] feat(engine): Push engine 2.0 execution progress to the editor (#37142) --- .../node-engine-compatibility/src/index.ts | 2 +- .../engine-control-plane-server.test.ts | 15 +- .../engine-lifecycle-event-push-relay.test.ts | 379 ++++++++++++++++++ .../engine-lifecycle-event.controller.test.ts | 42 +- .../engine-lifecycle-event-push-relay.ts | 235 +++++++++++ .../engine-lifecycle-event.controller.ts | 31 +- .../engine-v2-dispatcher.service.test.ts | 47 ++- .../engine-v2-push-registry.service.test.ts | 125 ++++++ .../services/engine-v2-dispatcher.service.ts | 48 ++- .../engine-v2-push-registry.service.ts | 64 +++ .../src/services/engine-v2-push-session.ts | 39 ++ 11 files changed, 955 insertions(+), 72 deletions(-) create mode 100644 packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event-push-relay.test.ts create mode 100644 packages/cli/src/modules/engine-v2/engine-lifecycle-event-push-relay.ts create mode 100644 packages/cli/src/services/__tests__/engine-v2-push-registry.service.test.ts create mode 100644 packages/cli/src/services/engine-v2-push-registry.service.ts create mode 100644 packages/cli/src/services/engine-v2-push-session.ts diff --git a/packages/@n8n/node-engine-compatibility/src/index.ts b/packages/@n8n/node-engine-compatibility/src/index.ts index b4e5e7cf30a..ea712b3758d 100644 --- a/packages/@n8n/node-engine-compatibility/src/index.ts +++ b/packages/@n8n/node-engine-compatibility/src/index.ts @@ -1,6 +1,6 @@ export { V1WorkflowConverter } from './v1-workflow-converter'; export { V1StepExecutor } from './v1-step-executor'; export { createEngineStepDataLoader } from './engine-step-data-loader'; -export { toStepOutputs } from './io'; +export { fromStepInputs, toStepOutputs } from './io'; export { UnsupportedTriggerError, UnsupportedWorkflowError } from './errors'; export type { StepData, StepDataLoader, V1StepExecutorDeps } from './types'; diff --git a/packages/cli/src/modules/engine-v2/__tests__/engine-control-plane-server.test.ts b/packages/cli/src/modules/engine-v2/__tests__/engine-control-plane-server.test.ts index fb3c7d34a8f..b8d182e9977 100644 --- a/packages/cli/src/modules/engine-v2/__tests__/engine-control-plane-server.test.ts +++ b/packages/cli/src/modules/engine-v2/__tests__/engine-control-plane-server.test.ts @@ -5,6 +5,7 @@ import request from 'supertest'; import { mock } from 'vitest-mock-extended'; import { EngineControlPlaneServer } from '../engine-control-plane-server'; +import type { EngineLifecycleEventPushRelay } from '../engine-lifecycle-event-push-relay'; import { EngineLifecycleEventController } from '../engine-lifecycle-event.controller'; const authSecret = 'a'.repeat(32); @@ -21,8 +22,8 @@ const events: LifecycleEvent[] = [ /** Binds for real, so these exercise the wiring rather than a mock app. */ describe('EngineControlPlaneServer', () => { let server: EngineControlPlaneServer; - let logger: Logger; let serverLogger: Logger; + let pushRelay: EngineLifecycleEventPushRelay; let baseUrl: string; const engineConfig = (overrides: Partial = {}) => @@ -35,11 +36,9 @@ describe('EngineControlPlaneServer', () => { }); beforeEach(async () => { - logger = mock(); serverLogger = mock(); - const controller = new EngineLifecycleEventController( - mock({ scoped: vi.fn().mockReturnValue(logger) }), - ); + pushRelay = mock(); + const controller = new EngineLifecycleEventController(pushRelay); server = new EngineControlPlaneServer( engineConfig(), controller, @@ -92,10 +91,8 @@ describe('EngineControlPlaneServer', () => { const response = await post({ events }, mintActionToken(authSecret, 'lifecycle-events:write')); expect(response.status).toBe(204); - expect(logger.debug).toHaveBeenCalledExactlyOnceWith( - 'Engine lifecycle event: execution:completed', - events[0], - ); + // Confirms the request reached the relay. + expect(pushRelay.relay).toHaveBeenCalledExactlyOnceWith(events); }); it.each([ diff --git a/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event-push-relay.test.ts b/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event-push-relay.test.ts new file mode 100644 index 00000000000..956573ef988 --- /dev/null +++ b/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event-push-relay.test.ts @@ -0,0 +1,379 @@ +import type { Logger } from '@n8n/backend-common'; +import type { LifecycleEvent } from '@n8n/engine'; +import type { PushMessage } from '@n8n/api-types'; +import { mock } from 'vitest-mock-extended'; + +import type { Push } from '@/push'; +import { EngineV2PushRegistry } from '@/services/engine-v2-push-registry.service'; + +import { EngineLifecycleEventPushRelay } from '../engine-lifecycle-event-push-relay'; + +const EXECUTION_ID = 'exec-1'; +const PUSH_REF = 'push-1'; +const WORKFLOW_ID = 'wf-1'; +const TRIGGER_NAME = 'When clicking Execute'; + +const stepFields = { + executionId: EXECUTION_ID, + stepId: 'step-1', + nodeId: 'node-a', + nodeName: 'Edit Fields', + iteration: 0, +}; + +const executionStarted: LifecycleEvent = { + type: 'execution:started', + executionId: EXECUTION_ID, + workflowId: WORKFLOW_ID, + mode: 'manual', + at: '2026-08-25T10:00:00.000Z', +}; + +const stepStarted: LifecycleEvent = { + ...stepFields, + type: 'step:started', + at: '2026-08-25T10:00:01.000Z', +}; + +const stepCompleted: LifecycleEvent = { + ...stepFields, + type: 'step:completed', + outputs: [[{ json: { greeting: 'hi' } }]], + at: '2026-08-25T10:00:01.500Z', +}; + +const stepFailed: LifecycleEvent = { + ...stepFields, + type: 'step:failed', + at: '2026-08-25T10:00:01.500Z', +}; + +describe('EngineLifecycleEventPushRelay', () => { + let push: Push; + let logger: Logger; + let registry: EngineV2PushRegistry; + let relay: EngineLifecycleEventPushRelay; + + /** Sent push messages, in order. */ + const sent = () => vi.mocked(push.send).mock.calls.map(([message]) => message); + + const sentOfType = (type: T) => + sent().filter((message): message is Extract => message.type === type); + + const register = (overrides: { trigger?: { nodeName: string; outputs: never } } = {}) => + registry.register(EXECUTION_ID, { + pushRef: PUSH_REF, + workflowId: WORKFLOW_ID, + ...overrides, + }); + + beforeEach(() => { + push = mock(); + logger = mock(); + registry = new EngineV2PushRegistry(); + relay = new EngineLifecycleEventPushRelay( + registry, + push, + mock({ scoped: vi.fn().mockReturnValue(logger) }), + ); + }); + + it('sends nothing for an execution it has no session for', () => { + relay.relay([executionStarted, stepStarted, stepCompleted]); + + expect(push.send).not.toHaveBeenCalled(); + }); + + it('routes every message to the session that started the run', () => { + register(); + + relay.relay([stepStarted, stepCompleted]); + + for (const [, pushRef] of vi.mocked(push.send).mock.calls) { + expect(pushRef).toBe(PUSH_REF); + } + }); + + it('never sends executionStarted', () => { + // Would overwrite the editor's existing run data with an empty scaffold. + register(); + + relay.relay([executionStarted, stepStarted, stepCompleted]); + + expect(sentOfType('executionStarted')).toHaveLength(0); + }); + + describe('the trigger', () => { + const triggerOutputs = [[{ json: { first: true } }]]; + + const registerWithTrigger = () => + registry.register(EXECUTION_ID, { + pushRef: PUSH_REF, + workflowId: WORKFLOW_ID, + trigger: { nodeName: TRIGGER_NAME, outputs: triggerOutputs }, + }); + + it('reports its run, because the engine never announces it', () => { + registerWithTrigger(); + + relay.relay([executionStarted]); + + expect(sent().map((message) => message.type)).toEqual([ + 'nodeExecuteBefore', + 'nodeExecuteAfter', + 'nodeExecuteAfterData', + ]); + const [before] = sentOfType('nodeExecuteBefore'); + expect(before.data.nodeName).toBe(TRIGGER_NAME); + expect(before.data.data.executionIndex).toBe(0); + expect(sentOfType('nodeExecuteAfterData')[0].data.data.data).toEqual({ + main: triggerOutputs, + }); + }); + + it('reports it only once, however often the update is redelivered', () => { + registerWithTrigger(); + + relay.relay([executionStarted, executionStarted]); + + expect(sentOfType('nodeExecuteBefore')).toHaveLength(1); + }); + + it('reports nothing when the run named no trigger', () => { + register(); + + relay.relay([executionStarted]); + + expect(push.send).not.toHaveBeenCalled(); + }); + }); + + describe('step:started', () => { + it('sends nodeExecuteBefore with the step start time', () => { + register(); + + relay.relay([stepStarted]); + + const [before] = sentOfType('nodeExecuteBefore'); + expect(before.data).toEqual({ + executionId: EXECUTION_ID, + nodeName: 'Edit Fields', + sequenceNumber: 0, + data: { + startTime: Date.parse(stepStarted.at), + executionIndex: 0, + source: [], + }, + }); + }); + + it('ignores a redelivered start', () => { + register(); + + relay.relay([stepStarted, stepStarted]); + + expect(sentOfType('nodeExecuteBefore')).toHaveLength(1); + }); + }); + + describe('step:completed', () => { + it('sends nodeExecuteAfter without the data, then nodeExecuteAfterData with it', () => { + register(); + + relay.relay([stepStarted, stepCompleted]); + + const [after] = sentOfType('nodeExecuteAfter'); + expect(after.data.data).not.toHaveProperty('data'); + expect(after.data.data.executionStatus).toBe('success'); + expect(after.data.itemCountByConnectionType).toEqual({ main: [1] }); + + const [afterData] = sentOfType('nodeExecuteAfterData'); + expect(afterData.data.data.data).toEqual({ main: [[{ json: { greeting: 'hi' } }]] }); + }); + + it('sends the output data as a binary frame', () => { + register(); + + relay.relay([stepStarted, stepCompleted]); + + const call = vi + .mocked(push.send) + .mock.calls.find(([message]) => message.type === 'nodeExecuteAfterData'); + expect(call?.[2]).toBe(true); + }); + + it('reuses the executionIndex allocated at the start, so the two messages pair up', () => { + register(); + + relay.relay([stepStarted, stepCompleted]); + + const [before] = sentOfType('nodeExecuteBefore'); + const [after] = sentOfType('nodeExecuteAfter'); + const [afterData] = sentOfType('nodeExecuteAfterData'); + expect(after.data.data.executionIndex).toBe(before.data.data.executionIndex); + expect(afterData.data.data.executionIndex).toBe(before.data.data.executionIndex); + }); + + it('reports the time between the two updates', () => { + register(); + + relay.relay([stepStarted, stepCompleted]); + + expect(sentOfType('nodeExecuteAfter')[0].data.data.executionTime).toBe(500); + }); + + it('turns a slot the step did not fire into an empty branch', () => { + register(); + + relay.relay([stepStarted, { ...stepCompleted, outputs: [[{ json: { a: 1 } }], null] }]); + + const [after] = sentOfType('nodeExecuteAfter'); + expect(after.data.itemCountByConnectionType).toEqual({ main: [1, 0] }); + expect(sentOfType('nodeExecuteAfterData')[0].data.data.data).toEqual({ + main: [[{ json: { a: 1 } }], []], + }); + }); + + it('still reports the outcome when the start was lost', () => { + register(); + + relay.relay([stepCompleted]); + + expect(sentOfType('nodeExecuteAfter')).toHaveLength(1); + expect(sentOfType('nodeExecuteAfterData')).toHaveLength(1); + }); + + it('ignores a redelivered completion', () => { + register(); + + relay.relay([stepStarted, stepCompleted, stepCompleted]); + + expect(sentOfType('nodeExecuteAfter')).toHaveLength(1); + }); + + it('reports the outcome on redelivery when the first send failed', () => { + register(); + vi.mocked(push.send).mockImplementationOnce(() => { + throw new Error('socket gone'); + }); + + relay.relay([stepCompleted, stepCompleted]); + + // The failed attempt sent nothing, so only the redelivery reported the run. + expect(sentOfType('nodeExecuteAfterData')).toHaveLength(1); + // The retry reuses the run, so the editor replaces it instead of appending. + const indexes = sentOfType('nodeExecuteAfter').map((m) => m.data.data.executionIndex); + expect(indexes).toEqual([0, 0]); + }); + + it('gives each step its own executionIndex and a rising sequenceNumber', () => { + register(); + const second = { ...stepFields, stepId: 'step-2', nodeName: 'Edit Fields 2' }; + + relay.relay([ + stepStarted, + stepCompleted, + { ...second, type: 'step:started', at: '2026-08-25T10:00:02.000Z' }, + { ...second, type: 'step:completed', outputs: [[]], at: '2026-08-25T10:00:03.000Z' }, + ]); + + expect(sentOfType('nodeExecuteBefore').map((m) => m.data.data.executionIndex)).toEqual([ + 0, 1, + ]); + expect( + [ + ...sentOfType('nodeExecuteBefore').map((m) => m.data.sequenceNumber), + ...sentOfType('nodeExecuteAfter').map((m) => m.data.sequenceNumber), + ].sort(), + ).toEqual([0, 1, 2, 3]); + }); + }); + + describe('step:failed', () => { + it('marks the node failed and sends no data message', () => { + register(); + + relay.relay([stepStarted, stepFailed]); + + const [after] = sentOfType('nodeExecuteAfter'); + expect(after.data.data.executionStatus).toBe('error'); + // The editor keys failure display on `error`, not the status. + expect(after.data.data.error).toBeDefined(); + expect(after.data.itemCountByConnectionType).toEqual({}); + expect(sentOfType('nodeExecuteAfterData')).toHaveLength(0); + }); + + it('sends an error that survives the wire', () => { + // A plain `Error` would serialize to `{}`. + register(); + + relay.relay([stepStarted, stepFailed]); + + const { error } = sentOfType('nodeExecuteAfter')[0].data.data; + expect(JSON.stringify(error)).toContain('Node execution failed'); + }); + }); + + describe('the end of an execution', () => { + it.each([ + ['execution:completed', 'success'], + ['execution:failed', 'error'], + ] as const)('maps %s to executionFinished %s', (type, status) => { + register(); + + relay.relay([ + { + type, + executionId: EXECUTION_ID, + workflowId: WORKFLOW_ID, + at: '2026-08-25T10:00:04.000Z', + }, + ]); + + expect(sentOfType('executionFinished')[0].data).toEqual({ + executionId: EXECUTION_ID, + workflowId: WORKFLOW_ID, + status, + }); + }); + + it('releases the session, so later updates are ignored', () => { + register(); + + relay.relay([ + { + type: 'execution:completed', + executionId: EXECUTION_ID, + workflowId: WORKFLOW_ID, + at: '2026-08-25T10:00:04.000Z', + }, + stepStarted, + ]); + + expect(registry.get(EXECUTION_ID)).toBeUndefined(); + expect(sentOfType('nodeExecuteBefore')).toHaveLength(0); + }); + }); + + it('keeps two concurrent runs apart', () => { + register(); + registry.register('exec-2', { pushRef: 'push-2', workflowId: 'wf-2' }); + + relay.relay([stepStarted, { ...stepStarted, executionId: 'exec-2', stepId: 'other-step' }]); + + const refs = vi.mocked(push.send).mock.calls.map(([, pushRef]) => pushRef); + expect(refs).toEqual([PUSH_REF, 'push-2']); + }); + + it('logs a failing update and relays the rest of the batch', () => { + register(); + vi.mocked(push.send).mockImplementationOnce(() => { + throw new Error('socket gone'); + }); + + relay.relay([stepStarted, stepCompleted]); + + expect(logger.error).toHaveBeenCalledTimes(1); + expect(sentOfType('nodeExecuteAfter')).toHaveLength(1); + }); +}); diff --git a/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event.controller.test.ts b/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event.controller.test.ts index cc88a3ebe79..8393f9d87c8 100644 --- a/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event.controller.test.ts +++ b/packages/cli/src/modules/engine-v2/__tests__/engine-lifecycle-event.controller.test.ts @@ -1,4 +1,3 @@ -import type { Logger } from '@n8n/backend-common'; import type { LifecycleEvent } from '@n8n/engine'; import type { Request, Response } from 'express'; import type { Mocked } from 'vitest'; @@ -6,6 +5,7 @@ import { mock } from 'vitest-mock-extended'; import { BadRequestError } from '@/errors/response-errors/bad-request.error'; +import type { EngineLifecycleEventPushRelay } from '../engine-lifecycle-event-push-relay'; import { EngineLifecycleEventController } from '../engine-lifecycle-event.controller'; const events: LifecycleEvent[] = [ @@ -29,8 +29,7 @@ const events: LifecycleEvent[] = [ ]; describe('EngineLifecycleEventController', () => { - // The controller scopes its logger, so assert on the scoped one. - let logger: Logger; + let pushRelay: EngineLifecycleEventPushRelay; let controller: EngineLifecycleEventController; const newResponse = () => { @@ -47,10 +46,8 @@ describe('EngineLifecycleEventController', () => { const newRequest = (body: unknown = { events }) => ({ body }) as unknown as Request; beforeEach(() => { - logger = mock(); - controller = new EngineLifecycleEventController( - mock({ scoped: vi.fn().mockReturnValue(logger) }), - ); + pushRelay = mock(); + controller = new EngineLifecycleEventController(pushRelay); }); describe('receiveLifecycleEvents', () => { @@ -63,35 +60,10 @@ describe('EngineLifecycleEventController', () => { expect(res.end).toHaveBeenCalled(); }); - it('logs every event in the batch, not just a count', async () => { + it('hands the whole batch to the push relay, in order', async () => { await controller.receiveLifecycleEvents(newRequest(), newResponse()); - expect(logger.debug).toHaveBeenCalledTimes(2); - expect(logger.debug).toHaveBeenNthCalledWith( - 1, - 'Engine lifecycle event: execution:started', - events[0], - ); - }); - - it('logs a completed step by its output slot count, never its contents', async () => { - // A log must not become a copy of a user's execution data. - await controller.receiveLifecycleEvents(newRequest(), newResponse()); - - const [message, metadata] = vi.mocked(logger.debug).mock.calls[1]; - - expect(message).toBe('Engine lifecycle event: step:completed'); - expect(metadata).toEqual({ - type: 'step:completed', - executionId: 'exec-1', - stepId: 'step-1', - nodeId: 'node-a', - nodeName: 'Edit Fields', - iteration: 0, - at: '2026-08-24T10:00:01.000Z', - outputSlots: 1, - }); - expect(JSON.stringify(metadata)).not.toContain('greeting'); + expect(pushRelay.relay).toHaveBeenCalledExactlyOnceWith(events); }); it.each([ @@ -105,6 +77,8 @@ describe('EngineLifecycleEventController', () => { BadRequestError, ); expect(res.status).not.toHaveBeenCalled(); + // Unvalidated input must never reach the relay. + expect(pushRelay.relay).not.toHaveBeenCalled(); }); }); }); diff --git a/packages/cli/src/modules/engine-v2/engine-lifecycle-event-push-relay.ts b/packages/cli/src/modules/engine-v2/engine-lifecycle-event-push-relay.ts new file mode 100644 index 00000000000..6a2af1aa702 --- /dev/null +++ b/packages/cli/src/modules/engine-v2/engine-lifecycle-event-push-relay.ts @@ -0,0 +1,235 @@ +import { Logger } from '@n8n/backend-common'; +import { Service } from '@n8n/di'; +import type { LifecycleEvent } from '@n8n/engine'; +import { fromStepInputs } from '@n8n/node-engine-compatibility'; +import type { ExecutionStatus, INodeExecutionData, ITaskData } from 'n8n-workflow'; +import { WorkflowOperationError } from 'n8n-workflow'; + +import { Push } from '@/push'; +import { EngineV2PushRegistry } from '@/services/engine-v2-push-registry.service'; +import type { EngineV2PushSession } from '@/services/engine-v2-push-session'; +import { EngineV2StepRun } from '@/services/engine-v2-push-session'; +import { getItemCountByConnectionType } from '@/utils/get-item-count-by-connection-type'; + +/** A lifecycle event scoped to one step. */ +type StepUpdate = Extract; + +/** + * Placeholder error: no failure detail is available yet, and without an + * `error` the editor would show a failed step as successful. + * + * TODO(CAT-2878 follow-up): carry the real failure through and drop this. + */ +const STEP_FAILURE_DESCRIPTION = 'Engine 2.0 does not report error detail yet.'; + +/** + * Relays engine lifecycle events to the editor as push messages. + * + * Reuses the push messages v1 runs already send, so the frontend needs no + * engine-specific code. + */ +@Service() +export class EngineLifecycleEventPushRelay { + constructor( + private readonly registry: EngineV2PushRegistry, + private readonly push: Push, + private readonly logger: Logger, + ) { + this.logger = this.logger.scoped('engine-v2'); + } + + relay(events: LifecycleEvent[]): void { + for (const update of events) { + const session = this.registry.get(update.executionId); + // No session means nothing is watching this run — safe to drop. + if (!session) continue; + + try { + this.relayOne(update, session); + } catch (error) { + // Isolate failures so one bad event doesn't drop the rest of the batch. + this.logger.error('Failed to relay engine lifecycle event to the editor', { + executionId: update.executionId, + type: update.type, + error, + }); + } + } + } + + private relayOne(update: LifecycleEvent, session: EngineV2PushSession): void { + switch (update.type) { + case 'execution:started': + return this.onExecutionStarted(update.executionId, update.at, session); + case 'step:started': + return this.onStepStarted(update, session); + case 'step:completed': + return this.onStepSettled(update, session, fromStepInputs(update.outputs)); + case 'step:failed': + return this.onStepSettled(update, session, undefined); + case 'execution:completed': + return this.onExecutionFinished(update.executionId, update.workflowId, 'success', session); + case 'execution:failed': + return this.onExecutionFinished(update.executionId, update.workflowId, 'error', session); + } + } + + /** + * Reports the trigger's run, since the engine never announces it as a step. + * + * Sends no `executionStarted`: the editor already promoted the run from the + * dispatch response, and the message would overwrite its run data with an + * empty scaffold. Runs the editor did not start need it — TODO(CAT-4258). + */ + private onExecutionStarted(executionId: string, at: string, session: EngineV2PushSession): void { + const { trigger } = session; + // Clear before use so a redelivery can't re-emit the run or hold onto + // the (possibly large) pinned data. + session.trigger = undefined; + // No trigger, no node to report outputs for. + if (!trigger) return; + + // The trigger has no step id, so it isn't tracked in `steps`. + const run = new EngineV2StepRun(session.nextExecutionIndex++, Date.parse(at)); + + this.sendNodeExecuteBefore(executionId, trigger.nodeName, run, session); + this.sendNodeExecuteAfter(executionId, trigger.nodeName, run, session, { + executionTime: 0, + outputs: trigger.outputs, + }); + } + + private onStepStarted(update: StepUpdate, session: EngineV2PushSession): void { + if (session.steps.has(update.stepId)) return; + + const run = new EngineV2StepRun(session.nextExecutionIndex++, Date.parse(update.at)); + session.steps.set(update.stepId, run); + + this.sendNodeExecuteBefore(update.executionId, update.nodeName, run, session); + } + + /** + * @param outputs The step's output slots, or `undefined` when it failed. + */ + private onStepSettled( + update: StepUpdate, + session: EngineV2PushSession, + outputs: INodeExecutionData[][] | undefined, + ): void { + const started = session.steps.get(update.stepId); + if (started?.settled) return; + + // A missing `step:started` must not lose the outcome too. + const run = started ?? new EngineV2StepRun(session.nextExecutionIndex++, Date.parse(update.at)); + session.steps.set(update.stepId, run); + + this.sendNodeExecuteAfter(update.executionId, update.nodeName, run, session, { + // Same clock, same process — the difference is safe to trust. + executionTime: Math.max(0, Date.parse(update.at) - run.startTime), + outputs, + }); + + run.settled = true; + } + + private onExecutionFinished( + executionId: string, + workflowId: string, + status: ExecutionStatus, + session: EngineV2PushSession, + ): void { + this.push.send( + { type: 'executionFinished', data: { executionId, workflowId, status } }, + session.pushRef, + ); + + // Releasing here makes a redelivered terminal event a no-op. + this.registry.release(executionId); + } + + private sendNodeExecuteBefore( + executionId: string, + nodeName: string, + run: EngineV2StepRun, + session: EngineV2PushSession, + ): void { + this.push.send( + { + type: 'nodeExecuteBefore', + data: { + executionId, + nodeName, + sequenceNumber: session.sequenceNumber++, + // No input lineage to report, so `source` is empty. See CAT-4265. + data: { startTime: run.startTime, executionIndex: run.executionIndex, source: [] }, + }, + }, + session.pushRef, + ); + } + + private sendNodeExecuteAfter( + executionId: string, + nodeName: string, + run: EngineV2StepRun, + session: EngineV2PushSession, + result: { executionTime: number; outputs: INodeExecutionData[][] | undefined }, + ): void { + const { executionTime, outputs } = result; + // Every output slot is `main`: no other connection type exists here. + const data = outputs ? { main: outputs } : undefined; + + const taskData: ITaskData = { + startTime: run.startTime, + executionIndex: run.executionIndex, + // No input lineage to report, so `source` is empty. See CAT-4265. + source: [], + executionTime, + executionStatus: outputs ? 'success' : 'error', + data, + ...(outputs + ? {} + : { + error: new WorkflowOperationError( + 'Node execution failed', + undefined, + STEP_FAILURE_DESCRIPTION, + ), + }), + }; + + const itemCountByConnectionType = getItemCountByConnectionType(data); + const { data: _withheld, ...trimmed } = taskData; + + this.push.send( + { + type: 'nodeExecuteAfter', + data: { + executionId, + nodeName, + sequenceNumber: session.sequenceNumber++, + data: trimmed, + itemCountByConnectionType, + }, + }, + session.pushRef, + ); + + // A failed step has no data, so there's no second message to send. + if (!data) return; + + // TODO(CAT-2878 follow-up): redact this before sending; needs a resolved + // user, which isn't available here yet. + + // Binary avoids a copy: the editor hands it straight to a worker. + const asBinary = true; + this.push.send( + { + type: 'nodeExecuteAfterData', + data: { executionId, nodeName, data: taskData, itemCountByConnectionType }, + }, + session.pushRef, + asBinary, + ); + } +} diff --git a/packages/cli/src/modules/engine-v2/engine-lifecycle-event.controller.ts b/packages/cli/src/modules/engine-v2/engine-lifecycle-event.controller.ts index f62f176c7c7..fb370ac0663 100644 --- a/packages/cli/src/modules/engine-v2/engine-lifecycle-event.controller.ts +++ b/packages/cli/src/modules/engine-v2/engine-lifecycle-event.controller.ts @@ -1,17 +1,15 @@ -import { Logger } from '@n8n/backend-common'; import { Service } from '@n8n/di'; -import type { LifecycleEvent } from '@n8n/engine'; import { lifecycleEventBatchSchema } from '@n8n/engine'; import type { Request, Response } from 'express'; import { BadRequestError } from '@/errors/response-errors/bad-request.error'; +import { EngineLifecycleEventPushRelay } from './engine-lifecycle-event-push-relay'; + /** Handles `LifecycleEvent` batches from the engine 2.0 data plane. */ @Service() export class EngineLifecycleEventController { - constructor(private readonly logger: Logger) { - this.logger = this.logger.scoped('engine-v2'); - } + constructor(private readonly pushRelay: EngineLifecycleEventPushRelay) {} async receiveLifecycleEvents(req: Request, res: Response): Promise { const parsed = lifecycleEventBatchSchema.safeParse(req.body); @@ -19,23 +17,14 @@ export class EngineLifecycleEventController { // stays out of the response: only a data plane calls this. if (!parsed.success) throw new BadRequestError('Invalid lifecycle event batch'); - // TODO(CAT-2878): forward these to the editor and dispatch the error workflow. - for (const event of parsed.data.events) { - this.logger.debug(`Engine lifecycle event: ${event.type}`, toLogMetadata(event)); - } + // TODO(CAT-2877 follow-up): dispatch the error workflow on + // `execution:failed` once the failure reason and run data are available. - // Nothing to return, and re-delivery is harmless while this only logs. + // Never throws: a bad event is logged and skipped inside. + this.pushRelay.relay(parsed.data.events); + + // Re-delivering a batch is harmless: the relay ignores events it has + // already reported. res.status(204).end(); } } - -/** - * An event's identifiers, ready to log. Outputs become a slot count, so a log - * never becomes a copy of a user's execution data. - */ -function toLogMetadata(event: LifecycleEvent): Record { - if (event.type !== 'step:completed') return { ...event }; - - const { outputs, ...rest } = event; - return { ...rest, outputSlots: outputs.length }; -} diff --git a/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts b/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts index 3edc3dd2e87..561cffb924b 100644 --- a/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts +++ b/packages/cli/src/services/__tests__/engine-v2-dispatcher.service.test.ts @@ -17,6 +17,7 @@ import type { CredentialsPermissionChecker } from '@/executions/pre-execution-ch import type { ResumableExecution } from '@/interfaces'; import type { EngineDataPlaneProxyService } from '@/services/engine-data-plane-proxy.service'; import { EngineV2Dispatcher } from '@/services/engine-v2-dispatcher.service'; +import type { EngineV2PushRegistry } from '@/services/engine-v2-push-registry.service'; const node = (id: string, name: string, type: string): INode => ({ id, @@ -72,6 +73,7 @@ function runData( describe('EngineV2Dispatcher', () => { const proxy = mock(); const credentialsPermissionChecker = mock(); + const pushRegistry = mock(); let dispatcher: EngineV2Dispatcher; @@ -79,7 +81,7 @@ describe('EngineV2Dispatcher', () => { vi.clearAllMocks(); proxy.isAvailable.mockReturnValue(true); proxy.startExecution.mockResolvedValue({ executionId: 'dp-uuid' }); - dispatcher = new EngineV2Dispatcher(proxy, credentialsPermissionChecker); + dispatcher = new EngineV2Dispatcher(proxy, credentialsPermissionChecker, pushRegistry); }); describe('routesToEngineV2', () => { @@ -317,5 +319,48 @@ describe('EngineV2Dispatcher', () => { expect(startedWith()).toBeNull(); }); }); + + describe('the push session', () => { + it('records the run against the data plane execution id', async () => { + await dispatcher.start(runData({ pushRef: 'push-1' })); + + expect(pushRegistry.register).toHaveBeenCalledExactlyOnceWith('dp-uuid', { + pushRef: 'push-1', + workflowId: 'wf-1', + trigger: { nodeName: MANUAL_TRIGGER.name, outputs: [[{ json: {} }]] }, + }); + }); + + it('records the trigger payload the engine was given', async () => { + const data = runData({ + pushRef: 'push-1', + triggerToStartFrom: { + name: MANUAL_TRIGGER.name, + data: taskData([[{ json: { from: 'trigger' } }]]), + }, + }); + + await dispatcher.start(data); + + expect(pushRegistry.register.mock.calls[0][1].trigger).toEqual({ + nodeName: MANUAL_TRIGGER.name, + outputs: [[{ json: { from: 'trigger' } }]], + }); + }); + + it('records nothing when nothing is watching the run', async () => { + await dispatcher.start(runData()); + + expect(pushRegistry.register).not.toHaveBeenCalled(); + }); + + it('records nothing when the data plane refused the run', async () => { + proxy.startExecution.mockRejectedValueOnce(new Error('down')); + + await expect(dispatcher.start(runData({ pushRef: 'push-1' }))).rejects.toThrow('down'); + + expect(pushRegistry.register).not.toHaveBeenCalled(); + }); + }); }); }); diff --git a/packages/cli/src/services/__tests__/engine-v2-push-registry.service.test.ts b/packages/cli/src/services/__tests__/engine-v2-push-registry.service.test.ts new file mode 100644 index 00000000000..0c506ec4a31 --- /dev/null +++ b/packages/cli/src/services/__tests__/engine-v2-push-registry.service.test.ts @@ -0,0 +1,125 @@ +import { EngineV2PushRegistry } from '@/services/engine-v2-push-registry.service'; + +const HOUR_MS = 60 * 60 * 1000; +const TTL_MS = 12 * HOUR_MS; + +describe('EngineV2PushRegistry', () => { + let registry: EngineV2PushRegistry; + + beforeEach(() => { + registry = new EngineV2PushRegistry(); + }); + + it('returns nothing for an execution it never saw', () => { + expect(registry.get('unknown')).toBeUndefined(); + }); + + it('starts a session with zeroed counters and no steps', () => { + registry.register('exec-1', { pushRef: 'push-1', workflowId: 'wf-1' }); + + expect(registry.get('exec-1')).toMatchObject({ + pushRef: 'push-1', + workflowId: 'wf-1', + sequenceNumber: 0, + nextExecutionIndex: 0, + }); + expect(registry.get('exec-1')?.steps.size).toBe(0); + }); + + it('keeps the trigger the run started from', () => { + const outputs = [[{ json: { x: 1 } }]]; + registry.register('exec-1', { + pushRef: 'push-1', + workflowId: 'wf-1', + trigger: { nodeName: 'When clicking Execute', outputs }, + }); + + expect(registry.get('exec-1')?.trigger).toEqual({ + nodeName: 'When clicking Execute', + outputs, + }); + }); + + it('releases a session, and releasing twice is a no-op', () => { + registry.register('exec-1', { pushRef: 'push-1', workflowId: 'wf-1' }); + + registry.release('exec-1'); + registry.release('exec-1'); + + expect(registry.get('exec-1')).toBeUndefined(); + }); + + it('keeps sessions apart', () => { + registry.register('exec-1', { pushRef: 'push-1', workflowId: 'wf-1' }); + registry.register('exec-2', { pushRef: 'push-2', workflowId: 'wf-2' }); + + registry.release('exec-1'); + + expect(registry.get('exec-2')?.pushRef).toBe('push-2'); + }); + + describe('eviction', () => { + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it('drops a session whose terminal update never arrived', () => { + // No `cancelled` event exists, so an unreleased session must expire. + vi.setSystemTime(new Date('2026-08-25T10:00:00.000Z')); + registry.register('stale', { pushRef: 'push-1', workflowId: 'wf-1' }); + + vi.advanceTimersByTime(TTL_MS + 1000); + registry.register('fresh', { pushRef: 'push-2', workflowId: 'wf-2' }); + + expect(registry.get('stale')).toBeUndefined(); + expect(registry.get('fresh')).toBeDefined(); + }); + + it('keeps a session that is still inside the retention window', () => { + vi.setSystemTime(new Date('2026-08-25T10:00:00.000Z')); + registry.register('exec-1', { pushRef: 'push-1', workflowId: 'wf-1' }); + + vi.advanceTimersByTime(TTL_MS - 1000); + registry.register('exec-2', { pushRef: 'push-2', workflowId: 'wf-2' }); + + expect(registry.get('exec-1')).toBeDefined(); + }); + + it('measures the window from the last event, not from registration', () => { + vi.setSystemTime(new Date('2026-08-25T10:00:00.000Z')); + registry.register('long-run', { pushRef: 'push-1', workflowId: 'wf-1' }); + + // A run that keeps reporting stays alive past the raw TTL. + for (let i = 0; i < 3; i++) { + vi.advanceTimersByTime(TTL_MS - 1000); + expect(registry.get('long-run')).toBeDefined(); + } + + vi.advanceTimersByTime(TTL_MS + 1000); + registry.register('other', { pushRef: 'push-2', workflowId: 'wf-2' }); + + expect(registry.get('long-run')).toBeUndefined(); + }); + + it('caps the number of sessions, dropping the least recently seen', () => { + vi.setSystemTime(new Date('2026-08-25T10:00:00.000Z')); + for (let i = 0; i < 1000; i++) { + registry.register(`exec-${i}`, { pushRef: `push-${i}`, workflowId: 'wf-1' }); + } + + vi.advanceTimersByTime(1000); + // Touching the oldest session makes the next one the eviction target. + expect(registry.get('exec-0')).toBeDefined(); + + registry.register('exec-1000', { pushRef: 'push-1000', workflowId: 'wf-1' }); + + expect(registry.get('exec-1')).toBeUndefined(); + expect(registry.get('exec-0')).toBeDefined(); + expect(registry.get('exec-1000')).toBeDefined(); + }); + }); +}); diff --git a/packages/cli/src/services/engine-v2-dispatcher.service.ts b/packages/cli/src/services/engine-v2-dispatcher.service.ts index ededb2b0eca..0af0e3b229f 100644 --- a/packages/cli/src/services/engine-v2-dispatcher.service.ts +++ b/packages/cli/src/services/engine-v2-dispatcher.service.ts @@ -10,6 +10,7 @@ import { getChildNodes, NodeConnectionTypes, UserError } from 'n8n-workflow'; import { CredentialsPermissionChecker } from '@/executions/pre-execution-checks'; import type { ResumableExecution } from '@/interfaces'; import { EngineDataPlaneProxyService } from '@/services/engine-data-plane-proxy.service'; +import { EngineV2PushRegistry } from '@/services/engine-v2-push-registry.service'; type ToStepOutputs = (outputs: INodeExecutionData[][]) => StepSlots; @@ -33,6 +34,7 @@ export class EngineV2Dispatcher { constructor( private readonly proxy: EngineDataPlaneProxyService, private readonly credentialsPermissionChecker: CredentialsPermissionChecker, + private readonly pushRegistry: EngineV2PushRegistry, ) {} /** @@ -64,17 +66,47 @@ export class EngineV2Dispatcher { const { V1WorkflowConverter, toStepOutputs } = await import('@n8n/node-engine-compatibility'); const graph = new V1WorkflowConverter().convert(this.selectTriggerSubgraph(data)); + const triggerMain = this.triggerMainOutputs(data); const { executionId } = await this.proxy.startExecution({ workflowId: workflowData.id, graph, - triggerOutputs: this.toTriggerOutputs(data, toStepOutputs), + triggerOutputs: this.toTriggerOutputs(triggerMain, toStepOutputs), mode: 'manual', }); + // TODO(CAT-4255): the engine can publish lifecycle events before this line + // runs, and the relay drops them because no session exists yet. Let the + // control plane mint the execution id so this can register before dispatch. + this.registerPushSession(executionId, data, triggerMain); + return executionId; } + /** + * Lifecycle events carry no session id, so the push ref is recorded here, + * keyed by execution id, before any events can arrive. + */ + private registerPushSession( + executionId: string, + data: IWorkflowExecutionDataProcess, + triggerMain: INodeExecutionData[][], + ): void { + const { pushRef, workflowData, triggerToStartFrom } = data; + // No push ref means nothing is watching this run. + if (!pushRef) return; + + this.pushRegistry.register(executionId, { + pushRef, + workflowId: workflowData.id, + // The engine never announces the trigger, so save its outputs for the relay. + trigger: triggerToStartFrom && { + nodeName: triggerToStartFrom.name, + outputs: triggerMain, + }, + }); + } + /** Keep only the branch that starts at the trigger selected for this manual run. */ private selectTriggerSubgraph(data: IWorkflowExecutionDataProcess): IWorkflowBase { const { triggerToStartFrom, workflowData } = data; @@ -152,10 +184,7 @@ export class EngineV2Dispatcher { * record the trigger with no slots, so every successor edge reads as dead and * the execution completes having run nothing. */ - private toTriggerOutputs( - data: IWorkflowExecutionDataProcess, - toStepOutputs: ToStepOutputs, - ): TriggerOutputs | null { + private triggerMainOutputs(data: IWorkflowExecutionDataProcess): INodeExecutionData[][] { const triggerName = data.triggerToStartFrom?.name; // `IPinData` values are a flat item array; the Manual Trigger has one output. const pinned = triggerName ? data.pinData?.[triggerName] : undefined; @@ -166,7 +195,14 @@ export class EngineV2Dispatcher { // v1 uses `null` for a slot it has no data for; an empty slot says the same // thing to the engine, which `toStepOutputs` collapses back to a dead edge. - const slots = toStepOutputs(main.map((slot) => slot ?? [])); + return main.map((slot) => slot ?? []); + } + + private toTriggerOutputs( + main: INodeExecutionData[][], + toStepOutputs: ToStepOutputs, + ): TriggerOutputs | null { + const slots = toStepOutputs(main); // The wire schema rejects an empty array; `null` is how "no slots" is sent. return slots.length === 0 ? null : slots; diff --git a/packages/cli/src/services/engine-v2-push-registry.service.ts b/packages/cli/src/services/engine-v2-push-registry.service.ts new file mode 100644 index 00000000000..ad3ea700962 --- /dev/null +++ b/packages/cli/src/services/engine-v2-push-registry.service.ts @@ -0,0 +1,64 @@ +import { Time } from '@n8n/constants'; +import { Service } from '@n8n/di'; + +import { EngineV2PushSession } from '@/services/engine-v2-push-session'; + +/** Long enough to outlive a run that idles between steps, e.g. on a wait. */ +const SESSION_TTL_MS = 12 * Time.hours.toMilliseconds; + +/** Hard ceiling on the map, in case many runs stall inside the TTL. */ +const MAX_SESSIONS = 1000; + +/** + * Correlates a data-plane execution id with the editor session that started it. + * + * Lifecycle events carry no session id, so the push ref is recorded here at + * dispatch and read back as events arrive. + */ +@Service() +export class EngineV2PushRegistry { + private readonly sessions = new Map(); + + register( + executionId: string, + init: Pick, + ): void { + this.evict(); + this.sessions.set( + executionId, + new EngineV2PushSession(init.pushRef, init.workflowId, init.trigger), + ); + } + + get(executionId: string): EngineV2PushSession | undefined { + const session = this.sessions.get(executionId); + if (session) session.lastSeenAt = Date.now(); + + return session; + } + + release(executionId: string): void { + this.sessions.delete(executionId); + } + + /** + * No `cancelled` event exists, so a session whose terminal event never + * arrives would live forever. Swept on write instead of on a timer, so + * there's no interval to manage. + */ + private evict(): void { + const cutoff = Date.now() - SESSION_TTL_MS; + for (const [executionId, session] of this.sessions) { + if (session.lastSeenAt < cutoff) this.sessions.delete(executionId); + } + + // Leave room for the caller's session, so the cap holds after the insert. + const excess = this.sessions.size - MAX_SESSIONS + 1; + if (excess <= 0) return; + + const oldestFirst = [...this.sessions].sort((a, b) => a[1].lastSeenAt - b[1].lastSeenAt); + for (const [executionId] of oldestFirst.slice(0, excess)) { + this.sessions.delete(executionId); + } + } +} diff --git a/packages/cli/src/services/engine-v2-push-session.ts b/packages/cli/src/services/engine-v2-push-session.ts new file mode 100644 index 00000000000..e8a1fcfb335 --- /dev/null +++ b/packages/cli/src/services/engine-v2-push-session.ts @@ -0,0 +1,39 @@ +import type { INodeExecutionData } from 'n8n-workflow'; + +/** + * A step run reported to the editor. Kept after it settles so a redelivered + * event is ignored instead of appending a duplicate. + */ +export class EngineV2StepRun { + /** Whether the step's outcome has been reported. */ + settled = false; + + constructor( + /** Pairs this run's `nodeExecuteAfter` with its `nodeExecuteAfterData`. */ + readonly executionIndex: number, + readonly startTime: number, + ) {} +} + +/** State needed to relay one execution's lifecycle events to the editor. */ +export class EngineV2PushSession { + /** Ordering counter for `nodeExecuteBefore`/`nodeExecuteAfter`; starts at 0. */ + sequenceNumber = 0; + /** Next `ITaskData.executionIndex` to hand out. */ + nextExecutionIndex = 0; + /** Step runs keyed by the engine's step id. */ + readonly steps = new Map(); + /** When the last lifecycle event for this execution arrived. */ + lastSeenAt = Date.now(); + + constructor( + /** The only routing key {@link Push.send} accepts. */ + readonly pushRef: string, + readonly workflowId: string, + /** + * The trigger's outputs, since the engine never announces it as a step. + * Cleared once emitted — pinned data can be large. + */ + public trigger?: { nodeName: string; outputs: INodeExecutionData[][] }, + ) {} +}