mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
fix(core): Move subworfklow binary duplication to workflowExecuteAfter before execution cleaning (#22390)
This commit is contained in:
@@ -858,5 +858,90 @@ describe('Execution Lifecycle Hooks', () => {
|
||||
expect(handlers.sendResponse).toHaveLength(0);
|
||||
expect(handlers.sendChunk).toHaveLength(0);
|
||||
});
|
||||
|
||||
describe('when parentExecution is provided', () => {
|
||||
const parentWorkflowId = 'parent-workflow-id';
|
||||
const parentExecutionId = 'parent-execution-id';
|
||||
const parentExecution = {
|
||||
workflowId: parentWorkflowId,
|
||||
executionId: parentExecutionId,
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
lifecycleHooks = getLifecycleHooksForSubExecutions(
|
||||
'integrated',
|
||||
executionId,
|
||||
workflowData,
|
||||
undefined,
|
||||
parentExecution,
|
||||
);
|
||||
});
|
||||
|
||||
it('should duplicate binary data to parent execution', async () => {
|
||||
const binaryDataId = `filesystem:workflows/${workflowId}/executions/${executionId}/binary_data/123`;
|
||||
const duplicatedBinaryDataId = `filesystem:workflows/${parentWorkflowId}/executions/${parentExecutionId}/binary_data/456`;
|
||||
|
||||
const mainOutputData = [
|
||||
[
|
||||
{
|
||||
json: {},
|
||||
binary: {
|
||||
data: {
|
||||
id: binaryDataId,
|
||||
data: '',
|
||||
mimeType: 'text/plain',
|
||||
},
|
||||
},
|
||||
},
|
||||
],
|
||||
];
|
||||
|
||||
successfulRun.data.resultData.runData = {
|
||||
[nodeName]: [
|
||||
{
|
||||
startTime: 1,
|
||||
executionIndex: 0,
|
||||
executionTime: 1,
|
||||
source: [],
|
||||
data: {
|
||||
main: mainOutputData,
|
||||
},
|
||||
},
|
||||
],
|
||||
};
|
||||
successfulRun.data.resultData.lastNodeExecuted = nodeName;
|
||||
|
||||
binaryDataService.duplicateBinaryData.mockResolvedValue([
|
||||
[
|
||||
{
|
||||
json: {},
|
||||
binary: {
|
||||
data: {
|
||||
id: duplicatedBinaryDataId,
|
||||
data: '',
|
||||
mimeType: 'text/plain',
|
||||
},
|
||||
},
|
||||
},
|
||||
],
|
||||
]);
|
||||
|
||||
await lifecycleHooks.runHook('workflowExecuteAfter', [successfulRun, {}]);
|
||||
|
||||
expect(binaryDataService.duplicateBinaryData).toHaveBeenCalledWith(
|
||||
{ type: 'execution', workflowId: parentWorkflowId, executionId: parentExecutionId },
|
||||
mainOutputData,
|
||||
);
|
||||
});
|
||||
|
||||
it('should not duplicate binary data when there is no output data', async () => {
|
||||
successfulRun.data.resultData.runData = {};
|
||||
successfulRun.data.resultData.lastNodeExecuted = undefined;
|
||||
|
||||
await lifecycleHooks.runHook('workflowExecuteAfter', [successfulRun, {}]);
|
||||
|
||||
expect(binaryDataService.duplicateBinaryData).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3,9 +3,17 @@ import { ExecutionRepository } from '@n8n/db';
|
||||
import { LifecycleMetadata } from '@n8n/decorators';
|
||||
import { Container, Service } from '@n8n/di';
|
||||
import { stringify } from 'flatted';
|
||||
import { ErrorReporter, InstanceSettings, ExecutionLifecycleHooks } from 'n8n-core';
|
||||
import {
|
||||
BinaryDataService,
|
||||
ErrorReporter,
|
||||
FileLocation,
|
||||
InstanceSettings,
|
||||
ExecutionLifecycleHooks,
|
||||
} from 'n8n-core';
|
||||
import type {
|
||||
IRun,
|
||||
IWorkflowBase,
|
||||
RelatedExecution,
|
||||
WorkflowExecuteMode,
|
||||
IWorkflowExecutionDataProcess,
|
||||
} from 'n8n-workflow';
|
||||
@@ -28,6 +36,7 @@ import {
|
||||
} from './shared/shared-hook-functions';
|
||||
import { type ExecutionSaveSettings, toSaveSettings } from './to-save-settings';
|
||||
import { getItemCountByConnectionType } from '@/utils/get-item-count-by-connection-type';
|
||||
import { getDataLastExecutedNodeData } from '@/workflow-helpers';
|
||||
|
||||
@Service()
|
||||
class ModulesHooksRegistry {
|
||||
@@ -99,6 +108,7 @@ type HooksSetupParameters = {
|
||||
saveSettings: ExecutionSaveSettings;
|
||||
pushRef?: string;
|
||||
retryOf?: string;
|
||||
parentExecution?: RelatedExecution;
|
||||
};
|
||||
|
||||
function hookFunctionsWorkflowEvents(hooks: ExecutionLifecycleHooks, userId?: string) {
|
||||
@@ -305,16 +315,39 @@ function hookFunctionsStatistics(hooks: ExecutionLifecycleHooks) {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Duplicates binary data from a subworkflow execution to the parent execution.
|
||||
* This ensures the parent can access the binary data after the subworkflow
|
||||
* execution is cleaned up. The duplicateBinaryData method also updates
|
||||
* the binary data IDs in the data to point to the new location.
|
||||
*/
|
||||
async function duplicateBinaryDataToParent(
|
||||
fullRunData: IRun,
|
||||
parentExecution: RelatedExecution,
|
||||
binaryDataService: BinaryDataService,
|
||||
) {
|
||||
const outputData = getDataLastExecutedNodeData(fullRunData);
|
||||
if (outputData?.data?.main) {
|
||||
const duplicatedData = await binaryDataService.duplicateBinaryData(
|
||||
FileLocation.ofExecution(parentExecution.workflowId, parentExecution.executionId),
|
||||
outputData.data.main,
|
||||
);
|
||||
// Update the run data with the new binary data IDs
|
||||
outputData.data.main = duplicatedData;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns hook functions to save workflow execution and call error workflow
|
||||
*/
|
||||
function hookFunctionsSave(
|
||||
hooks: ExecutionLifecycleHooks,
|
||||
{ pushRef, retryOf, saveSettings }: HooksSetupParameters,
|
||||
{ pushRef, retryOf, saveSettings, parentExecution }: HooksSetupParameters,
|
||||
) {
|
||||
const logger = Container.get(Logger);
|
||||
const errorReporter = Container.get(ErrorReporter);
|
||||
const executionRepository = Container.get(ExecutionRepository);
|
||||
const binaryDataService = Container.get(BinaryDataService);
|
||||
const workflowStaticDataService = Container.get(WorkflowStaticDataService);
|
||||
const workflowStatisticsService = Container.get(WorkflowStatisticsService);
|
||||
hooks.addHandler('workflowExecuteAfter', async function (fullRunData, newStaticData) {
|
||||
@@ -325,6 +358,13 @@ function hookFunctionsSave(
|
||||
|
||||
await restoreBinaryDataId(fullRunData, this.executionId, this.mode);
|
||||
|
||||
// If this is a subworkflow execution, duplicate binary data to the parent's
|
||||
// execution. This must happen before any potential deletion of this execution's
|
||||
// data, and updates the binary data IDs in fullRunData to point to the parent location.
|
||||
if (parentExecution) {
|
||||
await duplicateBinaryDataToParent(fullRunData, parentExecution, binaryDataService);
|
||||
}
|
||||
|
||||
const isManualMode = this.mode === 'manual';
|
||||
|
||||
try {
|
||||
@@ -481,13 +521,14 @@ export function getLifecycleHooksForSubExecutions(
|
||||
executionId: string,
|
||||
workflowData: IWorkflowBase,
|
||||
userId?: string,
|
||||
parentExecution?: RelatedExecution,
|
||||
): ExecutionLifecycleHooks {
|
||||
const hooks = new ExecutionLifecycleHooks(mode, executionId, workflowData);
|
||||
const saveSettings = toSaveSettings(workflowData.settings);
|
||||
hookFunctionsWorkflowEvents(hooks, userId);
|
||||
hookFunctionsNodeEvents(hooks);
|
||||
hookFunctionsFinalizeExecutionStatus(hooks);
|
||||
hookFunctionsSave(hooks, { saveSettings });
|
||||
hookFunctionsSave(hooks, { saveSettings, parentExecution });
|
||||
hookFunctionsSaveProgress(hooks, { saveSettings });
|
||||
hookFunctionsStatistics(hooks);
|
||||
hookFunctionsExternalHooks(hooks);
|
||||
|
||||
@@ -230,6 +230,7 @@ async function startExecution(
|
||||
executionId,
|
||||
workflowData,
|
||||
additionalData.userId,
|
||||
options.parentExecution,
|
||||
);
|
||||
additionalDataIntegrated.executionId = executionId;
|
||||
additionalDataIntegrated.parentCallbackManager = options.parentCallbackManager;
|
||||
|
||||
@@ -269,7 +269,7 @@ export const describeCommonTests = (
|
||||
|
||||
describe('executeWorkflow', () => {
|
||||
const data = [[{ json: { test: true } }]];
|
||||
const executeWorkflowData = mock<ExecuteWorkflowData>();
|
||||
const executeWorkflowData = mock<ExecuteWorkflowData>({ data });
|
||||
const workflowInfo = mock<IExecuteWorkflowInfo>();
|
||||
const parentExecution: RelatedExecution = {
|
||||
executionId: 'parent_execution_id',
|
||||
@@ -278,23 +278,18 @@ export const describeCommonTests = (
|
||||
|
||||
it('should execute workflow and return data', async () => {
|
||||
additionalData.executeWorkflow.mockResolvedValue(executeWorkflowData);
|
||||
binaryDataService.duplicateBinaryData.mockResolvedValue(data);
|
||||
|
||||
const result = await context.executeWorkflow(workflowInfo, undefined, undefined, {
|
||||
parentExecution,
|
||||
});
|
||||
|
||||
expect(result.data).toEqual(data);
|
||||
expect(binaryDataService.duplicateBinaryData).toHaveBeenCalledWith(
|
||||
{ type: 'execution', workflowId: workflow.id, executionId: additionalData.executionId },
|
||||
executeWorkflowData.data,
|
||||
);
|
||||
expect(result).toBe(executeWorkflowData);
|
||||
});
|
||||
|
||||
it('should put execution to wait if waitTill is returned', async () => {
|
||||
const waitTill = new Date();
|
||||
additionalData.executeWorkflow.mockResolvedValue({ ...executeWorkflowData, waitTill });
|
||||
binaryDataService.duplicateBinaryData.mockResolvedValue(data);
|
||||
|
||||
const result = await context.executeWorkflow(workflowInfo, undefined, undefined, {
|
||||
parentExecution,
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
import { Container } from '@n8n/di';
|
||||
import get from 'lodash/get';
|
||||
import type {
|
||||
Workflow,
|
||||
@@ -33,14 +32,9 @@ import {
|
||||
createEnvProviderState,
|
||||
} from 'n8n-workflow';
|
||||
|
||||
import { BinaryDataService } from '@/binary-data/binary-data.service';
|
||||
import { FileLocation } from '@/binary-data/utils';
|
||||
|
||||
import { NodeExecutionContext } from './node-execution-context';
|
||||
|
||||
export class BaseExecuteContext extends NodeExecutionContext {
|
||||
protected readonly binaryDataService = Container.get(BinaryDataService);
|
||||
|
||||
constructor(
|
||||
workflow: Workflow,
|
||||
node: INode,
|
||||
@@ -155,11 +149,7 @@ export class BaseExecuteContext extends NodeExecutionContext {
|
||||
await this.putExecutionToWait(WAIT_INDEFINITELY);
|
||||
}
|
||||
|
||||
const data = await this.binaryDataService.duplicateBinaryData(
|
||||
FileLocation.ofExecution(this.workflow.id, this.additionalData.executionId!),
|
||||
result.data,
|
||||
);
|
||||
return { ...result, data };
|
||||
return result;
|
||||
}
|
||||
|
||||
async getExecutionDataById(executionId: string): Promise<IRunExecutionData | undefined> {
|
||||
|
||||
Reference in New Issue
Block a user