Allow to stop running workflow (#15270)

https://github.com/user-attachments/assets/599154e4-8743-471b-b05a-721b635bcf4e

On stoppage:
- if no running steps, mark pending as failed and end the workflow
- if running steps, set as stopping and exit. Going to the next step, as
the workflow is not running anymore, it will naturally stop
This commit is contained in:
Thomas Trompette
2025-10-23 13:01:15 +00:00
committed by GitHub
parent 1cf442966d
commit 4effa5351c
23 changed files with 532 additions and 70 deletions
@@ -18,7 +18,10 @@ export class WorkflowRunUpdateOnePreQueryHook
_objectName: string,
payload: UpdateOneResolverArgs<WorkflowRunWorkspaceEntity>,
): Promise<UpdateOneResolverArgs<WorkflowRunWorkspaceEntity>> {
if (Object.keys(payload.data).length === 1 && payload.data.name) {
const allowedFields = ['name'];
const payloadKeys = Object.keys(payload.data);
if (payloadKeys.every((key) => allowedFields.includes(key))) {
return payload;
}
@@ -1,3 +1,5 @@
import { registerEnumType } from '@nestjs/graphql';
import { msg } from '@lingui/core/macro';
import { FieldMetadataType } from 'twenty-shared/types';
import { type WorkflowRunStepInfos } from 'twenty-shared/workflow';
@@ -40,8 +42,15 @@ export enum WorkflowRunStatus {
COMPLETED = 'COMPLETED',
FAILED = 'FAILED',
ENQUEUED = 'ENQUEUED',
STOPPING = 'STOPPING',
STOPPED = 'STOPPED',
}
registerEnumType(WorkflowRunStatus, {
name: 'WorkflowRunStatusEnum',
description: 'Status of the workflow run',
});
export type StepOutput = {
id: string;
output: WorkflowActionOutput;
@@ -159,6 +168,18 @@ export class WorkflowRunWorkspaceEntity extends BaseWorkspaceEntity {
position: 4,
color: 'blue',
},
{
value: WorkflowRunStatus.STOPPING,
label: 'Stopping',
position: 5,
color: 'orange',
},
{
value: WorkflowRunStatus.STOPPED,
label: 'Stopped',
position: 6,
color: 'gray',
},
],
defaultValue: "'NOT_STARTED'",
})
@@ -0,0 +1,15 @@
import { StepStatus, type WorkflowRunStepInfos } from 'twenty-shared/workflow';
import { type WorkflowAction } from 'src/modules/workflow/workflow-executor/workflow-actions/types/workflow-action.type';
export const workflowHasRunningSteps = ({
stepInfos,
steps,
}: {
stepInfos: WorkflowRunStepInfos;
steps: WorkflowAction[];
}) => {
return steps.some(
(step) => stepInfos[step.id]?.status === StepStatus.RUNNING,
);
};
@@ -20,6 +20,7 @@ import { MessageQueue } from 'src/engine/core-modules/message-queue/message-queu
import { MessageQueueService } from 'src/engine/core-modules/message-queue/services/message-queue.service';
import { WorkspaceEventEmitter } from 'src/engine/workspace-event-emitter/workspace-event-emitter';
import { WorkflowRunStatus } from 'src/modules/workflow/common/standard-objects/workflow-run.workspace-entity';
import { workflowHasRunningSteps } from 'src/modules/workflow/common/utils/workflow-has-running-steps.util';
import { WorkflowActionFactory } from 'src/modules/workflow/workflow-executor/factories/workflow-action.factory';
import { type WorkflowActionOutput } from 'src/modules/workflow/workflow-executor/types/workflow-action-output.type';
import {
@@ -224,6 +225,18 @@ export class WorkflowExecutorWorkspaceService {
const steps = workflowRun.state.flow.steps;
if (workflowRun.status === WorkflowRunStatus.STOPPING) {
if (!workflowHasRunningSteps({ stepInfos, steps })) {
await this.workflowRunWorkspaceService.endWorkflowRun({
workflowRunId,
workspaceId,
status: WorkflowRunStatus.STOPPED,
});
}
return;
}
if (workflowShouldFail({ stepInfos, steps })) {
await this.workflowRunWorkspaceService.endWorkflowRun({
workflowRunId,
@@ -190,7 +190,7 @@ export class WorkflowRunWorkspaceService {
}: {
workflowRunId: string;
workspaceId: string;
status: WorkflowRunStatus;
status: Extract<WorkflowRunStatus, 'COMPLETED' | 'FAILED' | 'STOPPED'>;
error?: string;
}) {
const workflowRunToUpdate = await this.getWorkflowRunOrFail({
@@ -199,13 +199,10 @@ export class WorkflowRunWorkspaceService {
});
let updatedStepInfos = {};
const shouldUpdateStepInfos = status === WorkflowRunStatus.FAILED;
if (shouldUpdateStepInfos) {
updatedStepInfos = this.markRunningStepsAsFailed({
stepInfosToUpdate: workflowRunToUpdate.state?.stepInfos ?? {},
});
}
updatedStepInfos = this.markRunningStepsAsFailed({
stepInfosToUpdate: workflowRunToUpdate.state?.stepInfos ?? {},
});
const partialUpdate = {
status,
@@ -213,7 +210,7 @@ export class WorkflowRunWorkspaceService {
state: {
...workflowRunToUpdate.state,
workflowRunError: error,
...(shouldUpdateStepInfos && { stepInfos: updatedStepInfos }),
stepInfos: updatedStepInfos,
},
};
@@ -223,7 +220,9 @@ export class WorkflowRunWorkspaceService {
key:
status === WorkflowRunStatus.COMPLETED
? MetricsKeys.WorkflowRunCompleted
: MetricsKeys.WorkflowRunFailed,
: status === WorkflowRunStatus.STOPPED
? MetricsKeys.WorkflowRunStopped
: MetricsKeys.WorkflowRunFailed,
eventId: workflowRunId,
});
}
@@ -378,35 +377,7 @@ export class WorkflowRunWorkspaceService {
return workflowRun;
}
private getInitState(
workflowVersion: WorkflowVersionWorkspaceEntity,
triggerPayload: object,
): WorkflowRunState | undefined {
if (
!isDefined(workflowVersion.trigger) ||
!isDefined(workflowVersion.steps)
) {
return undefined;
}
return {
flow: {
trigger: workflowVersion.trigger,
steps: workflowVersion.steps,
},
stepInfos: {
trigger: { status: StepStatus.NOT_STARTED, result: triggerPayload },
...Object.fromEntries(
workflowVersion.steps.map((step) => [
step.id,
{ status: StepStatus.NOT_STARTED },
]),
),
},
};
}
private async updateWorkflowRun({
async updateWorkflowRun({
workflowRunId,
workspaceId,
partialUpdate,
@@ -441,6 +412,34 @@ export class WorkflowRunWorkspaceService {
);
}
private getInitState(
workflowVersion: WorkflowVersionWorkspaceEntity,
triggerPayload: object,
): WorkflowRunState | undefined {
if (
!isDefined(workflowVersion.trigger) ||
!isDefined(workflowVersion.steps)
) {
return undefined;
}
return {
flow: {
trigger: workflowVersion.trigger,
steps: workflowVersion.steps,
},
stepInfos: {
trigger: { status: StepStatus.NOT_STARTED, result: triggerPayload },
...Object.fromEntries(
workflowVersion.steps.map((step) => [
step.id,
{ status: StepStatus.NOT_STARTED },
]),
),
},
};
}
private markRunningStepsAsFailed({
stepInfosToUpdate,
}: {
@@ -14,9 +14,14 @@ import {
WorkflowVersionStepExceptionCode,
} from 'src/modules/workflow/common/exceptions/workflow-version-step.exception';
import { WorkflowRunStatus } from 'src/modules/workflow/common/standard-objects/workflow-run.workspace-entity';
import { workflowHasRunningSteps } from 'src/modules/workflow/common/utils/workflow-has-running-steps.util';
import { WorkflowCommonWorkspaceService } from 'src/modules/workflow/common/workspace-services/workflow-common.workspace-service';
import { WorkflowVersionStepOperationsWorkspaceService } from 'src/modules/workflow/workflow-builder/workflow-version-step/workflow-version-step-operations.workspace-service';
import { isWorkflowFormAction } from 'src/modules/workflow/workflow-executor/workflow-actions/form/guards/is-workflow-form-action.guard';
import {
WorkflowRunException,
WorkflowRunExceptionCode,
} from 'src/modules/workflow/workflow-runner/exceptions/workflow-run.exception';
import { RunWorkflowJob } from 'src/modules/workflow/workflow-runner/jobs/run-workflow.job';
import { type RunWorkflowJobData } from 'src/modules/workflow/workflow-runner/types/run-workflow-job-data.type';
import { WorkflowRunQueueWorkspaceService } from 'src/modules/workflow/workflow-runner/workflow-run-queue/workspace-services/workflow-run-queue.workspace-service';
@@ -191,4 +196,52 @@ export class WorkflowRunnerWorkspaceService {
workspaceId,
);
}
async stopWorkflowRun(workspaceId: string, workflowRunId: string) {
const workflowRun =
await this.workflowRunWorkspaceService.getWorkflowRunOrFail({
workflowRunId,
workspaceId,
});
if (workflowRun.status !== WorkflowRunStatus.RUNNING) {
throw new WorkflowRunException(
'Workflow run is not running',
WorkflowRunExceptionCode.INVALID_OPERATION,
{
userFriendlyMessage: msg`Workflow run is not running`,
},
);
}
let newStatus: WorkflowRunStatus;
if (
workflowHasRunningSteps({
stepInfos: workflowRun.state.stepInfos,
steps: workflowRun.state.flow.steps,
})
) {
await this.workflowRunWorkspaceService.updateWorkflowRun({
workflowRunId,
workspaceId,
partialUpdate: {
status: WorkflowRunStatus.STOPPING,
},
});
newStatus = WorkflowRunStatus.STOPPING;
} else {
await this.workflowRunWorkspaceService.endWorkflowRun({
workflowRunId,
workspaceId,
status: WorkflowRunStatus.STOPPED,
});
newStatus = WorkflowRunStatus.STOPPED;
}
return {
id: workflowRun.id,
status: newStatus,
};
}
}
@@ -142,6 +142,13 @@ export class WorkflowTriggerWorkspaceService {
return true;
}
async stopWorkflowRun(workflowRunId: string) {
return this.workflowRunnerWorkspaceService.stopWorkflowRun(
this.getWorkspaceId(),
workflowRunId,
);
}
private async performActivationSteps(
workflow: WorkflowWorkspaceEntity,
workflowVersion: WorkflowVersionWorkspaceEntity,