mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
fix(core): Persist execution context before writing to db (#28973)
Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.7
parent
4358f1d51c
commit
c4bb5ae8df
@@ -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<IWorkflowExecutionDataProcess>({
|
||||
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<string, string> }).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<core.ExecutionLifecycleHooks>({ 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
|
||||
|
||||
@@ -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<IExecuteResponsePromiseData>,
|
||||
): Promise<void> {
|
||||
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<IExecuteResponsePromiseData>,
|
||||
): Promise<string> {
|
||||
// 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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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<INode>({
|
||||
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<string, string>;
|
||||
expect(headers.authorization).toBe('original-header-value');
|
||||
});
|
||||
});
|
||||
|
||||
describe('sub-workflow context inheritance', () => {
|
||||
|
||||
@@ -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<void> => {
|
||||
assertExecutionDataExists(runExecutionData.executionData, workflow, additionalData, mode);
|
||||
|
||||
@@ -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';
|
||||
|
||||
@@ -9,14 +9,14 @@ import {
|
||||
export function assertExecutionDataExists(
|
||||
executionData: IRunExecutionData['executionData'],
|
||||
workflow: Workflow,
|
||||
additionalData: IWorkflowExecuteAdditionalData,
|
||||
additionalData: IWorkflowExecuteAdditionalData | undefined,
|
||||
mode: WorkflowExecuteMode,
|
||||
): asserts executionData is NonNullable<IRunExecutionData['executionData']> {
|
||||
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,
|
||||
},
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user