From c4bb5ae8df8e7de4c7b919a82d3cf2f492edcc5b Mon Sep 17 00:00:00 2001 From: phyllis-noester <102315132+phyllis-noester@users.noreply.github.com> Date: Wed, 29 Apr 2026 16:15:13 +0200 Subject: [PATCH] fix(core): Persist execution context before writing to db (#28973) Co-authored-by: Claude Opus 4.7 (1M context) --- .../cli/src/__tests__/workflow-runner.test.ts | 199 ++++++++++++++++++ packages/cli/src/workflow-runner.ts | 92 ++++++-- .../__tests__/execution-context.test.ts | 67 ++++++ .../src/execution-engine/execution-context.ts | 2 +- packages/core/src/execution-engine/index.ts | 1 + packages/core/src/utils/assertions.ts | 4 +- 6 files changed, 350 insertions(+), 15 deletions(-) diff --git a/packages/cli/src/__tests__/workflow-runner.test.ts b/packages/cli/src/__tests__/workflow-runner.test.ts index 23f6c8ce69f..35382fe424f 100644 --- a/packages/cli/src/__tests__/workflow-runner.test.ts +++ b/packages/cli/src/__tests__/workflow-runner.test.ts @@ -16,6 +16,7 @@ import { type IExecuteData, type INode, type IRun, + type IRunExecutionData, type ITaskData, type IWaitingForExecution, type IWaitingForExecutionSource, @@ -23,9 +24,11 @@ import { type IWorkflowExecutionDataProcess, type StartNodeData, type IWorkflowExecuteAdditionalData, + type WorkflowExecuteMode, Workflow, ExecutionError, TimeoutExecutionCancelledError, + createRunExecutionData, } from 'n8n-workflow'; import PCancelable from 'p-cancelable'; @@ -630,6 +633,202 @@ describe('needsFullExecutionData', () => { }); }); +describe('pre-persist context establishment', () => { + const callOrder: string[] = []; + let establishSpy: jest.SpyInstance; + let addSpy: jest.SpyInstance; + let capturedAddData: IWorkflowExecutionDataProcess | undefined; + + const buildRunData = ( + executionData: IRunExecutionData | undefined, + executionMode: WorkflowExecuteMode = 'webhook', + ): IWorkflowExecutionDataProcess => + mock({ + executionMode, + workflowData: { + id: 'wf-1', + name: 'Test', + nodes: [], + connections: {}, + settings: undefined, + }, + executionData, + userId: 'u1', + }); + + const buildExecutionDataWithHeader = (): IRunExecutionData => + createRunExecutionData({ + executionData: { + contextData: {}, + nodeExecutionStack: [ + { + node: { + id: 'n1', + name: 'Webhook', + type: 'n8n-nodes-base.webhook', + typeVersion: 2, + position: [0, 0], + parameters: {}, + }, + data: { + main: [[{ json: { headers: { authorization: 'Bearer eyJtest' } } }]], + }, + source: null, + }, + ], + metadata: {}, + waitingExecution: {}, + waitingExecutionSource: {}, + }, + }); + + beforeEach(() => { + callOrder.length = 0; + capturedAddData = undefined; + + const permissionChecker = Container.get(CredentialsPermissionChecker); + jest.spyOn(permissionChecker, 'check').mockResolvedValue(); + + establishSpy = jest + .spyOn(core, 'establishExecutionContext') + .mockImplementation(async (_workflow, runExecutionData) => { + callOrder.push('establishExecutionContext'); + // Simulate real behaviour: mask header, set runtimeData. + const stack = runExecutionData.executionData?.nodeExecutionStack; + const item = stack?.[0]?.data?.main?.[0]?.[0]; + const headers = item && (item.json as { headers?: Record }).headers; + if (headers && typeof headers.authorization === 'string') { + headers.authorization = '**********'; + } + if (runExecutionData.executionData) { + (runExecutionData.executionData as { runtimeData?: unknown }).runtimeData = { + version: 1, + establishedAt: Date.now(), + source: 'webhook', + }; + } + }); + + const activeExecutions = Container.get(ActiveExecutions); + addSpy = jest + .spyOn(activeExecutions, 'add') + .mockImplementation(async (data: IWorkflowExecutionDataProcess) => { + capturedAddData = data; + callOrder.push('activeExecutions.add'); + throw new Error('short-circuit for test'); + }); + }); + + afterEach(() => { + establishSpy.mockRestore(); + addSpy.mockRestore(); + }); + + it('calls establishExecutionContext before activeExecutions.add', async () => { + const data = buildRunData(buildExecutionDataWithHeader()); + + await expect(runner.run(data)).rejects.toThrow('short-circuit for test'); + + expect(callOrder).toEqual(['establishExecutionContext', 'activeExecutions.add']); + }); + + it('passes masked executionData to activeExecutions.add', async () => { + const data = buildRunData(buildExecutionDataWithHeader()); + + await expect(runner.run(data)).rejects.toThrow('short-circuit for test'); + + expect(capturedAddData).toBeDefined(); + const stack = capturedAddData!.executionData!.executionData!.nodeExecutionStack; + expect(stack[0].data.main[0]![0].json).toMatchObject({ + headers: { authorization: '**********' }, + }); + expect(capturedAddData!.executionData!.executionData!.runtimeData).toBeDefined(); + }); + + it('skips establishExecutionContext when data.executionData is undefined', async () => { + const data = buildRunData(undefined); + + await expect(runner.run(data)).rejects.toThrow('short-circuit for test'); + + expect(establishSpy).not.toHaveBeenCalled(); + expect(callOrder).toEqual(['activeExecutions.add']); + }); + + it('skips establishExecutionContext when nodeExecutionStack has not been populated yet', async () => { + // Queue mode with OFFLOAD_MANUAL_EXECUTIONS_TO_WORKERS=true creates + // the outer IRunExecutionData with `executionData: null`, which + // `createRunExecutionData` normalises to an object whose inner + // `.executionData` is undefined. The worker establishes context + // later, once it populates the trigger-item stack. + const data = buildRunData(createRunExecutionData({ executionData: null }), 'manual'); + + await expect(runner.run(data)).rejects.toThrow('short-circuit for test'); + + expect(establishSpy).not.toHaveBeenCalled(); + expect(callOrder).toEqual(['activeExecutions.add']); + }); + + describe('when establishExecutionContext throws', () => { + const lifecycleRunHook = jest.fn().mockResolvedValue(undefined); + const responseReject = jest.fn(); + + beforeEach(() => { + establishSpy.mockReset(); + establishSpy.mockImplementation(async () => { + callOrder.push('establishExecutionContext'); + throw new Error('hook augmentation failed'); + }); + + addSpy.mockReset(); + addSpy.mockImplementation(async (data: IWorkflowExecutionDataProcess) => { + capturedAddData = data; + callOrder.push('activeExecutions.add'); + return 'exec-1'; + }); + + lifecycleRunHook.mockClear(); + responseReject.mockClear(); + jest + .spyOn(ExecutionLifecycleHooks, 'getLifecycleHooksForRegularMain') + .mockReturnValue(mock({ runHook: lifecycleRunHook })); + jest.spyOn(Container.get(ActiveExecutions), 'finalizeExecution').mockReturnValue(); + }); + + it('clears the trigger-item stack so raw headers do not get persisted', async () => { + const data = buildRunData(buildExecutionDataWithHeader()); + + await runner.run(data, undefined, undefined, undefined, { + reject: responseReject, + resolve: jest.fn(), + promise: Promise.resolve() as never, + } as never); + + expect(capturedAddData).toBeDefined(); + expect(capturedAddData!.executionData!.executionData!.nodeExecutionStack).toEqual([]); + }); + + it('creates a failed-execution record and rejects the responsePromise', async () => { + const data = buildRunData(buildExecutionDataWithHeader()); + + const executionId = await runner.run(data, undefined, undefined, undefined, { + reject: responseReject, + resolve: jest.fn(), + promise: Promise.resolve() as never, + } as never); + + expect(executionId).toBe('exec-1'); + expect(lifecycleRunHook).toHaveBeenCalledWith('workflowExecuteBefore', [ + undefined, + data.executionData, + ]); + expect(lifecycleRunHook).toHaveBeenCalledWith('workflowExecuteAfter', [expect.any(Object)]); + expect(responseReject).toHaveBeenCalledWith( + expect.objectContaining({ message: 'hook augmentation failed' }), + ); + }); + }); +}); + describe('streaming functionality', () => { it('should setup heartbeat interval and sendChunk handler when streaming is enabled', async () => { // ARRANGE diff --git a/packages/cli/src/workflow-runner.ts b/packages/cli/src/workflow-runner.ts index 9127570c4ae..c831e678acb 100644 --- a/packages/cli/src/workflow-runner.ts +++ b/packages/cli/src/workflow-runner.ts @@ -7,11 +7,18 @@ import { ExecutionsConfig } from '@n8n/config'; import { ExecutionRepository } from '@n8n/db'; import { Container, Service } from '@n8n/di'; import type { ExecutionLifecycleHooks } from 'n8n-core'; -import { ErrorReporter, InstanceSettings, StorageConfig, WorkflowExecute } from 'n8n-core'; +import { + ErrorReporter, + establishExecutionContext, + InstanceSettings, + StorageConfig, + WorkflowExecute, +} from 'n8n-core'; import type { ExecutionError, IDeferredPromise, IExecuteResponsePromiseData, + INode, IPinData, IRun, WorkflowExecuteMode, @@ -149,6 +156,28 @@ export class WorkflowRunner { await hooks?.runHook('workflowExecuteAfter', [fullRunData]); } + /** + * Persist a failed execution record for an already-registered execution and + * finalize it, so a pre-flight failure surfaces as a normal failed run. + */ + private async failExecution( + data: IWorkflowExecutionDataProcess, + executionId: string, + error: ExecutionError & { node?: INode }, + responsePromise?: IDeferredPromise, + ): Promise { + const runData = this.failedRunFactory.generateFailedExecutionFromError( + data.executionMode, + error, + error.node, + ); + const lifecycleHooks = getLifecycleHooksForRegularMain(data, executionId); + await lifecycleHooks.runHook('workflowExecuteBefore', [undefined, data.executionData]); + await lifecycleHooks.runHook('workflowExecuteAfter', [runData]); + responsePromise?.reject(error); + this.activeExecutions.finalizeExecution(executionId); + } + /** Run the workflow * @param realtime This is used in queue mode to change the priority of an execution, making sure they are picked up quicker. */ @@ -159,25 +188,64 @@ export class WorkflowRunner { restartExecutionId?: string, responsePromise?: IDeferredPromise, ): Promise { + // Establish the execution context before persisting to the DB. + // activeExecutions.add() -> executionPersistence.create() writes + // data.executionData to the DB; any header masking or runtimeData + // population must happen before that write so the persisted record + // does not contain raw trigger-item data (e.g. Authorization headers). + // The runtimeData early-exit guard in establishExecutionContext keeps + // the subsequent worker-side call at workflow-execute.ts idempotent. + // Guard on the inner executionData: in queue mode with manual offload + // the outer IRunExecutionData is created with `executionData: null` + // so the trigger-item stack is undefined here; nothing to mask yet, + // the worker will establish context once it populates the stack. + let establishContextError: (ExecutionError & { node?: INode }) | undefined; + if (data.executionData?.executionData) { + // Deliberately lightweight: no pinData, no staticData loading, + // no additionalData. establishExecutionContext only needs the + // workflow's settings (for redactionPolicy) and node lookups. + // runMainProcess() builds its own fully-configured Workflow for + // actual execution. + const contextWorkflow = new Workflow({ + id: data.workflowData.id, + name: data.workflowData.name, + nodes: data.workflowData.nodes, + connections: data.workflowData.connections, + active: data.workflowData.activeVersionId !== null, + nodeTypes: this.nodeTypes, + staticData: data.workflowData.staticData, + settings: data.workflowData.settings ?? {}, + }); + try { + await establishExecutionContext( + contextWorkflow, + data.executionData, + undefined, + data.executionMode, + ); + } catch (error) { + // Masking may have failed partway through, so the trigger-item + // stack can still contain raw header data. Drop it before + // activeExecutions.add() persists the execution row. + data.executionData.executionData.nodeExecutionStack = []; + establishContextError = error as ExecutionError & { node?: INode }; + } + } + // Register a new execution const executionId = await this.activeExecutions.add(data, restartExecutionId); + if (establishContextError) { + await this.failExecution(data, executionId, establishContextError, responsePromise); + return executionId; + } + const { id: workflowId, nodes } = data.workflowData; try { await this.credentialsPermissionChecker.check(workflowId, nodes); } catch (error) { - // Create a failed execution with the data for the node, save it and abort execution - const runData = this.failedRunFactory.generateFailedExecutionFromError( - data.executionMode, - error, - error.node, - ); - const lifecycleHooks = getLifecycleHooksForRegularMain(data, executionId); - await lifecycleHooks.runHook('workflowExecuteBefore', [undefined, data.executionData]); - await lifecycleHooks.runHook('workflowExecuteAfter', [runData]); - responsePromise?.reject(error); - this.activeExecutions.finalizeExecution(executionId); + await this.failExecution(data, executionId, error, responsePromise); return executionId; } diff --git a/packages/core/src/execution-engine/__tests__/execution-context.test.ts b/packages/core/src/execution-engine/__tests__/execution-context.test.ts index 677be2f5409..aaa8f0c38c9 100644 --- a/packages/core/src/execution-engine/__tests__/execution-context.test.ts +++ b/packages/core/src/execution-engine/__tests__/execution-context.test.ts @@ -456,6 +456,73 @@ describe('establishExecutionContext', () => { 'parent-exec-123', ); }); + + it('should skip hook augmentation when runtimeData is already set (queue-mode worker resume)', async () => { + // When the main process established the context before persisting, + // the worker fetches the execution with runtimeData already populated. + // The function must return immediately without re-running hook augmentation. + const existingContext: IExecutionContext = { + version: 1, + establishedAt: 1234567890, + source: 'webhook', + credentials: 'encrypted-credentials-blob', + }; + + const webhookNode = mock({ + name: 'Webhook', + type: 'n8n-nodes-base.webhook', + parameters: {}, + }); + + const runExecutionData = createRunExecutionData({ + startData: {}, + resultData: { runData: {} }, + executionData: { + contextData: {}, + nodeExecutionStack: [ + { + node: webhookNode, + data: { + main: [ + [ + { + json: { + headers: { authorization: 'original-header-value' }, + }, + }, + ], + ], + }, + source: null, + }, + ], + metadata: {}, + waitingExecution: {}, + waitingExecutionSource: {}, + runtimeData: existingContext, + }, + }); + + await establishExecutionContext( + mockWorkflow, + runExecutionData, + mockAdditionalData, + 'webhook', + ); + + // runtimeData reference is preserved — same object, same values + expect(runExecutionData.executionData!.runtimeData).toBe(existingContext); + expect(runExecutionData.executionData!.runtimeData!.credentials).toBe( + 'encrypted-credentials-blob', + ); + + // Trigger items are left untouched (hook augmentation was skipped). + // The main process already applied any transformations before persisting; + // this assertion just verifies the worker-side call is a no-op. + const headers = runExecutionData.executionData!.nodeExecutionStack[0].data.main[0]![0].json + .headers as Record; + expect(headers.authorization).toBe('original-header-value'); + }); }); describe('sub-workflow context inheritance', () => { diff --git a/packages/core/src/execution-engine/execution-context.ts b/packages/core/src/execution-engine/execution-context.ts index d4641cb3022..4991d8f214b 100644 --- a/packages/core/src/execution-engine/execution-context.ts +++ b/packages/core/src/execution-engine/execution-context.ts @@ -109,7 +109,7 @@ import { ExecutionContextService } from './execution-context.service'; export const establishExecutionContext = async ( workflow: Workflow, runExecutionData: IRunExecutionData, - additionalData: IWorkflowExecuteAdditionalData, + additionalData: IWorkflowExecuteAdditionalData | undefined, mode: WorkflowExecuteMode, ): Promise => { assertExecutionDataExists(runExecutionData.executionData, workflow, additionalData, mode); diff --git a/packages/core/src/execution-engine/index.ts b/packages/core/src/execution-engine/index.ts index 0cfbd8b51f5..ca7a64f4c90 100644 --- a/packages/core/src/execution-engine/index.ts +++ b/packages/core/src/execution-engine/index.ts @@ -91,4 +91,5 @@ export * from './execution-context-hook-registry.service'; export { ExecutionLifecycleHooks } from './execution-lifecycle-hooks'; export { ExternalSecretsProxy, type IExternalSecretsManager } from './external-secrets-proxy'; export { ExecutionContextService } from './execution-context.service'; +export { establishExecutionContext } from './execution-context'; export { isEngineRequest } from './requests-response'; diff --git a/packages/core/src/utils/assertions.ts b/packages/core/src/utils/assertions.ts index 74fabed1d51..450913ac68b 100644 --- a/packages/core/src/utils/assertions.ts +++ b/packages/core/src/utils/assertions.ts @@ -9,14 +9,14 @@ import { export function assertExecutionDataExists( executionData: IRunExecutionData['executionData'], workflow: Workflow, - additionalData: IWorkflowExecuteAdditionalData, + additionalData: IWorkflowExecuteAdditionalData | undefined, mode: WorkflowExecuteMode, ): asserts executionData is NonNullable { if (!executionData) { throw new UnexpectedError('Failed to run workflow due to missing execution data', { extra: { workflowId: workflow.id, - executionId: additionalData.executionId, + executionId: additionalData?.executionId, mode, }, });