mirror of
https://github.com/n8n-io/n8n.git
synced 2025-03-05 20:50:17 -08:00
refactor(core): Simplify postExecutePromises
This commit is contained in:
parent
cef64329a9
commit
2b06cb91c7
|
@ -96,7 +96,7 @@ describe('ActiveExecutions', () => {
|
||||||
test('Should remove an existing execution', async () => {
|
test('Should remove an existing execution', async () => {
|
||||||
const newExecution = mockExecutionData();
|
const newExecution = mockExecutionData();
|
||||||
const executionId = await activeExecutions.add(newExecution);
|
const executionId = await activeExecutions.add(newExecution);
|
||||||
activeExecutions.remove(executionId);
|
activeExecutions.finishExecution(executionId);
|
||||||
|
|
||||||
expect(activeExecutions.getActiveExecutions().length).toBe(0);
|
expect(activeExecutions.getActiveExecutions().length).toBe(0);
|
||||||
});
|
});
|
||||||
|
@ -110,7 +110,7 @@ describe('ActiveExecutions', () => {
|
||||||
setTimeout(res, 100);
|
setTimeout(res, 100);
|
||||||
});
|
});
|
||||||
const fakeOutput = mockFullRunData();
|
const fakeOutput = mockFullRunData();
|
||||||
activeExecutions.remove(executionId, fakeOutput);
|
activeExecutions.finishExecution(executionId, fakeOutput);
|
||||||
|
|
||||||
await expect(postExecutePromise).resolves.toEqual(fakeOutput);
|
await expect(postExecutePromise).resolves.toEqual(fakeOutput);
|
||||||
});
|
});
|
||||||
|
@ -126,7 +126,7 @@ describe('ActiveExecutions', () => {
|
||||||
const cancellablePromise = mockCancelablePromise();
|
const cancellablePromise = mockCancelablePromise();
|
||||||
cancellablePromise.cancel = cancelExecution;
|
cancellablePromise.cancel = cancelExecution;
|
||||||
activeExecutions.attachWorkflowExecution(executionId, cancellablePromise);
|
activeExecutions.attachWorkflowExecution(executionId, cancellablePromise);
|
||||||
void activeExecutions.stopExecution(executionId);
|
activeExecutions.stopExecution(executionId);
|
||||||
|
|
||||||
expect(cancelExecution).toHaveBeenCalledTimes(1);
|
expect(cancelExecution).toHaveBeenCalledTimes(1);
|
||||||
});
|
});
|
||||||
|
|
|
@ -95,13 +95,23 @@ export class ActiveExecutions {
|
||||||
await this.executionRepository.updateExistingExecution(executionId, execution);
|
await this.executionRepository.updateExistingExecution(executionId, execution);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const postExecutePromise = createDeferredPromise<IRun | undefined>();
|
||||||
|
|
||||||
this.activeExecutions[executionId] = {
|
this.activeExecutions[executionId] = {
|
||||||
executionData,
|
executionData,
|
||||||
startedAt: new Date(),
|
startedAt: new Date(),
|
||||||
postExecutePromises: [],
|
postExecutePromise,
|
||||||
status: executionStatus,
|
status: executionStatus,
|
||||||
};
|
};
|
||||||
|
|
||||||
|
// Automatically remove execution once the postExecutePromise settles
|
||||||
|
void postExecutePromise.promise.finally(() => {
|
||||||
|
this.concurrencyControl.release({ mode: executionData.executionMode });
|
||||||
|
setImmediate(() => {
|
||||||
|
delete this.activeExecutions[executionId];
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
return executionId;
|
return executionId;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -125,68 +135,24 @@ export class ActiveExecutions {
|
||||||
execution?.responsePromise?.resolve(response);
|
execution?.responsePromise?.resolve(response);
|
||||||
}
|
}
|
||||||
|
|
||||||
getPostExecutePromiseCount(executionId: string): number {
|
/** Forces an execution to stop */
|
||||||
return this.activeExecutions[executionId]?.postExecutePromises.length ?? 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Remove an active execution
|
|
||||||
*/
|
|
||||||
remove(executionId: string, fullRunData?: IRun): void {
|
|
||||||
const execution = this.activeExecutions[executionId];
|
|
||||||
if (execution === undefined) {
|
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Resolve all the waiting promises
|
|
||||||
for (const promise of execution.postExecutePromises) {
|
|
||||||
promise.resolve(fullRunData);
|
|
||||||
}
|
|
||||||
|
|
||||||
this.postExecuteCleanup(executionId);
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Forces an execution to stop
|
|
||||||
*/
|
|
||||||
stopExecution(executionId: string): void {
|
stopExecution(executionId: string): void {
|
||||||
const execution = this.activeExecutions[executionId];
|
const execution = this.getExecution(executionId);
|
||||||
if (execution === undefined) {
|
execution.workflowExecution?.cancel();
|
||||||
// There is no execution running with that id
|
execution.postExecutePromise.reject(new ExecutionCancelledError(executionId));
|
||||||
return;
|
|
||||||
}
|
|
||||||
|
|
||||||
execution.workflowExecution!.cancel();
|
|
||||||
|
|
||||||
// Reject all the waiting promises
|
|
||||||
const reason = new ExecutionCancelledError(executionId);
|
|
||||||
for (const promise of execution.postExecutePromises) {
|
|
||||||
promise.reject(reason);
|
|
||||||
}
|
|
||||||
|
|
||||||
this.postExecuteCleanup(executionId);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
private postExecuteCleanup(executionId: string) {
|
/** Mark an execution as completed */
|
||||||
const execution = this.activeExecutions[executionId];
|
finishExecution(executionId: string, fullRunData?: IRun) {
|
||||||
if (execution === undefined) {
|
const execution = this.getExecution(executionId);
|
||||||
return;
|
execution.postExecutePromise.resolve(fullRunData);
|
||||||
}
|
|
||||||
|
|
||||||
// Remove from the list of active executions
|
|
||||||
delete this.activeExecutions[executionId];
|
|
||||||
|
|
||||||
this.concurrencyControl.release({ mode: execution.executionData.executionMode });
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Returns a promise which will resolve with the data of the execution with the given id
|
* Returns a promise which will resolve with the data of the execution with the given id
|
||||||
*/
|
*/
|
||||||
async getPostExecutePromise(executionId: string): Promise<IRun | undefined> {
|
async getPostExecutePromise(executionId: string): Promise<IRun | undefined> {
|
||||||
// Create the promise which will be resolved when the execution finished
|
return await this.getExecution(executionId).postExecutePromise.promise;
|
||||||
const waitPromise = createDeferredPromise<IRun | undefined>();
|
|
||||||
this.getExecution(executionId).postExecutePromises.push(waitPromise);
|
|
||||||
return await waitPromise.promise;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
@ -193,7 +193,7 @@ export interface IExecutionsCurrentSummary {
|
||||||
export interface IExecutingWorkflowData {
|
export interface IExecutingWorkflowData {
|
||||||
executionData: IWorkflowExecutionDataProcess;
|
executionData: IWorkflowExecutionDataProcess;
|
||||||
startedAt: Date;
|
startedAt: Date;
|
||||||
postExecutePromises: Array<IDeferredPromise<IRun | undefined>>;
|
postExecutePromise: IDeferredPromise<IRun | undefined>;
|
||||||
responsePromise?: IDeferredPromise<IExecuteResponsePromiseData>;
|
responsePromise?: IDeferredPromise<IExecuteResponsePromiseData>;
|
||||||
workflowExecution?: PCancelable<IRun>;
|
workflowExecution?: PCancelable<IRun>;
|
||||||
status: ExecutionStatus;
|
status: ExecutionStatus;
|
||||||
|
|
|
@ -880,7 +880,7 @@ async function executeWorkflow(
|
||||||
}
|
}
|
||||||
|
|
||||||
// remove execution from active executions
|
// remove execution from active executions
|
||||||
activeExecutions.remove(executionId, fullRunData);
|
activeExecutions.finishExecution(executionId, fullRunData);
|
||||||
|
|
||||||
await executionRepository.updateExistingExecution(executionId, fullExecutionData);
|
await executionRepository.updateExistingExecution(executionId, fullExecutionData);
|
||||||
throw objectToError(
|
throw objectToError(
|
||||||
|
@ -906,11 +906,11 @@ async function executeWorkflow(
|
||||||
if (data.finished === true || data.status === 'waiting') {
|
if (data.finished === true || data.status === 'waiting') {
|
||||||
// Workflow did finish successfully
|
// Workflow did finish successfully
|
||||||
|
|
||||||
activeExecutions.remove(executionId, data);
|
activeExecutions.finishExecution(executionId, data);
|
||||||
const returnData = WorkflowHelpers.getDataLastExecutedNodeData(data);
|
const returnData = WorkflowHelpers.getDataLastExecutedNodeData(data);
|
||||||
return returnData!.data!.main;
|
return returnData!.data!.main;
|
||||||
}
|
}
|
||||||
activeExecutions.remove(executionId, data);
|
activeExecutions.finishExecution(executionId, data);
|
||||||
|
|
||||||
// Workflow did fail
|
// Workflow did fail
|
||||||
const { error } = data.data.resultData;
|
const { error } = data.data.resultData;
|
||||||
|
|
|
@ -26,7 +26,6 @@ import { ActiveExecutions } from '@/active-executions';
|
||||||
import config from '@/config';
|
import config from '@/config';
|
||||||
import { ExecutionRepository } from '@/databases/repositories/execution.repository';
|
import { ExecutionRepository } from '@/databases/repositories/execution.repository';
|
||||||
import { ExternalHooks } from '@/external-hooks';
|
import { ExternalHooks } from '@/external-hooks';
|
||||||
import type { IExecutionResponse } from '@/interfaces';
|
|
||||||
import { Logger } from '@/logger';
|
import { Logger } from '@/logger';
|
||||||
import { NodeTypes } from '@/node-types';
|
import { NodeTypes } from '@/node-types';
|
||||||
import type { ScalingService } from '@/scaling/scaling.service';
|
import type { ScalingService } from '@/scaling/scaling.service';
|
||||||
|
@ -101,7 +100,7 @@ export class WorkflowRunner {
|
||||||
|
|
||||||
// Remove from active execution with empty data. That will
|
// Remove from active execution with empty data. That will
|
||||||
// set the execution to failed.
|
// set the execution to failed.
|
||||||
this.activeExecutions.remove(executionId, fullRunData);
|
this.activeExecutions.finishExecution(executionId, fullRunData);
|
||||||
|
|
||||||
if (hooks) {
|
if (hooks) {
|
||||||
await hooks.executeHookFunctions('workflowExecuteAfter', [fullRunData]);
|
await hooks.executeHookFunctions('workflowExecuteAfter', [fullRunData]);
|
||||||
|
@ -129,7 +128,7 @@ export class WorkflowRunner {
|
||||||
await workflowHooks.executeHookFunctions('workflowExecuteBefore', []);
|
await workflowHooks.executeHookFunctions('workflowExecuteBefore', []);
|
||||||
await workflowHooks.executeHookFunctions('workflowExecuteAfter', [runData]);
|
await workflowHooks.executeHookFunctions('workflowExecuteAfter', [runData]);
|
||||||
responsePromise?.reject(error);
|
responsePromise?.reject(error);
|
||||||
this.activeExecutions.remove(executionId);
|
this.activeExecutions.finishExecution(executionId);
|
||||||
return executionId;
|
return executionId;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -321,7 +320,7 @@ export class WorkflowRunner {
|
||||||
fullRunData.finished = false;
|
fullRunData.finished = false;
|
||||||
}
|
}
|
||||||
fullRunData.status = this.activeExecutions.getStatus(executionId);
|
fullRunData.status = this.activeExecutions.getStatus(executionId);
|
||||||
this.activeExecutions.remove(executionId, fullRunData);
|
this.activeExecutions.finishExecution(executionId, fullRunData);
|
||||||
})
|
})
|
||||||
.catch(
|
.catch(
|
||||||
async (error) =>
|
async (error) =>
|
||||||
|
@ -490,45 +489,27 @@ export class WorkflowRunner {
|
||||||
reject(error);
|
reject(error);
|
||||||
}
|
}
|
||||||
|
|
||||||
// optimization: only pull and unflatten execution data from the Db when it is needed
|
|
||||||
const executionHasPostExecutionPromises =
|
|
||||||
this.activeExecutions.getPostExecutePromiseCount(executionId) > 0;
|
|
||||||
|
|
||||||
if (executionHasPostExecutionPromises) {
|
|
||||||
this.logger.debug(
|
|
||||||
`Reading execution data for execution ${executionId} from db for PostExecutionPromise.`,
|
|
||||||
);
|
|
||||||
} else {
|
|
||||||
this.logger.debug(
|
|
||||||
`Skipping execution data for execution ${executionId} since there are no PostExecutionPromise.`,
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
const fullExecutionData = await this.executionRepository.findSingleExecution(executionId, {
|
const fullExecutionData = await this.executionRepository.findSingleExecution(executionId, {
|
||||||
includeData: executionHasPostExecutionPromises,
|
includeData: true,
|
||||||
unflattenData: executionHasPostExecutionPromises,
|
unflattenData: true,
|
||||||
});
|
});
|
||||||
if (!fullExecutionData) {
|
if (!fullExecutionData) {
|
||||||
return reject(new Error(`Could not find execution with id "${executionId}"`));
|
return reject(new Error(`Could not find execution with id "${executionId}"`));
|
||||||
}
|
}
|
||||||
|
|
||||||
const runData: IRun = {
|
const runData: IRun = {
|
||||||
data: {},
|
|
||||||
finished: fullExecutionData.finished,
|
finished: fullExecutionData.finished,
|
||||||
mode: fullExecutionData.mode,
|
mode: fullExecutionData.mode,
|
||||||
startedAt: fullExecutionData.startedAt,
|
startedAt: fullExecutionData.startedAt,
|
||||||
stoppedAt: fullExecutionData.stoppedAt,
|
stoppedAt: fullExecutionData.stoppedAt,
|
||||||
status: fullExecutionData.status,
|
status: fullExecutionData.status,
|
||||||
} as IRun;
|
data: fullExecutionData.data,
|
||||||
|
};
|
||||||
if (executionHasPostExecutionPromises) {
|
|
||||||
runData.data = (fullExecutionData as IExecutionResponse).data;
|
|
||||||
}
|
|
||||||
|
|
||||||
// NOTE: due to the optimization of not loading the execution data from the db when no post execution promises are present,
|
// NOTE: due to the optimization of not loading the execution data from the db when no post execution promises are present,
|
||||||
// the execution data in runData.data MAY not be available here.
|
// the execution data in runData.data MAY not be available here.
|
||||||
// This means that any function expecting with runData has to check if the runData.data defined from this point
|
// This means that any function expecting with runData has to check if the runData.data defined from this point
|
||||||
this.activeExecutions.remove(executionId, runData);
|
this.activeExecutions.finishExecution(executionId, runData);
|
||||||
|
|
||||||
// Normally also static data should be supplied here but as it only used for sending
|
// Normally also static data should be supplied here but as it only used for sending
|
||||||
// data to editor-UI is not needed.
|
// data to editor-UI is not needed.
|
||||||
|
|
Loading…
Reference in a new issue