feat: delay node

This commit is contained in:
Fu Diwei
2025-08-22 22:31:56 +08:00
parent 752fc72d6c
commit 625e13ece3
21 changed files with 312 additions and 33 deletions
+29 -18
View File
@@ -93,6 +93,7 @@ const (
WorkflowNodeTypeTryCatch = WorkflowNodeType("tryCatch")
WorkflowNodeTypeTryBlock = WorkflowNodeType("tryBlock")
WorkflowNodeTypeCatchBlock = WorkflowNodeType("catchBlock")
WorkflowNodeTypeDelay = WorkflowNodeType("delay")
WorkflowNodeTypeBizApply = WorkflowNodeType("bizApply")
WorkflowNodeTypeBizUpload = WorkflowNodeType("bizUpload")
WorkflowNodeTypeBizMonitor = WorkflowNodeType("bizMonitor")
@@ -107,6 +108,29 @@ type WorkflowNodeData struct {
type WorkflowNodeConfig map[string]any
func (c WorkflowNodeConfig) AsDelay() WorkflowNodeConfigForDelay {
return WorkflowNodeConfigForDelay{
Wait: xmaps.GetInt32(c, "wait"),
}
}
func (c WorkflowNodeConfig) AsBranchBlock() WorkflowNodeConfigForBranchBlock {
expression := c["expression"]
if expression == nil {
return WorkflowNodeConfigForBranchBlock{}
}
exprRaw, _ := json.Marshal(expression)
expr, err := expr.UnmarshalExpr([]byte(exprRaw))
if err != nil {
return WorkflowNodeConfigForBranchBlock{}
}
return WorkflowNodeConfigForBranchBlock{
Expression: expr,
}
}
func (c WorkflowNodeConfig) AsBizApply() WorkflowNodeConfigForBizApply {
return WorkflowNodeConfigForBizApply{
Domains: xmaps.GetString(c, "domains"),
@@ -169,21 +193,12 @@ func (c WorkflowNodeConfig) AsBizNotify() WorkflowNodeConfigForBizNotify {
}
}
func (c WorkflowNodeConfig) AsBranchBlock() WorkflowNodeConfigForBranchBlock {
expression := c["expression"]
if expression == nil {
return WorkflowNodeConfigForBranchBlock{}
}
type WorkflowNodeConfigForDelay struct {
Wait int32 `json:"wait"` // 等待时间
}
exprRaw, _ := json.Marshal(expression)
expr, err := expr.UnmarshalExpr([]byte(exprRaw))
if err != nil {
return WorkflowNodeConfigForBranchBlock{}
}
return WorkflowNodeConfigForBranchBlock{
Expression: expr,
}
type WorkflowNodeConfigForBranchBlock struct {
Expression expr.Expr `json:"expression"` // 条件表达式
}
type WorkflowNodeConfigForBizApply struct {
@@ -236,7 +251,3 @@ type WorkflowNodeConfigForBizNotify struct {
Message string `json:"message"` // 通知内容
SkipOnAllPrevSkipped bool `json:"skipOnAllPrevSkipped"` // 前序节点均已跳过时是否跳过
}
type WorkflowNodeConfigForBranchBlock struct {
Expression expr.Expr `json:"expression"` // 条件表达式
}
+5
View File
@@ -2,10 +2,12 @@ package handlers
import (
"context"
"errors"
"github.com/pocketbase/pocketbase/core"
"github.com/pocketbase/pocketbase/tools/router"
"github.com/certimate-go/certimate/internal/domain"
"github.com/certimate-go/certimate/internal/domain/dtos"
"github.com/certimate-go/certimate/internal/rest/resp"
)
@@ -36,6 +38,9 @@ func (handler *WorkflowHandler) run(e *core.RequestEvent) error {
if err := e.BindBody(req); err != nil {
return resp.Err(e, err)
}
if req.RunTrigger != domain.WorkflowTriggerTypeManual {
return resp.Err(e, errors.New("invalid parameters: the value of 'trigger' must be 'manual'"))
}
res, err := handler.service.StartRun(e.Request.Context(), req)
if err != nil {
+1
View File
@@ -282,6 +282,7 @@ func NewWorkflowEngine() WorkflowEngine {
}
engine.executors[NodeTypeStart] = newStartNodeExecutor()
engine.executors[NodeTypeEnd] = newEndNodeExecutor()
engine.executors[NodeTypeDelay] = newDelayNodeExecutor()
engine.executors[NodeTypeCondition] = newConditionNodeExecutor()
engine.executors[NodeTypeBranchBlock] = newBranchBlockNodeExecutor()
engine.executors[NodeTypeTryCatch] = newTryCatchNodeExecutor()
@@ -0,0 +1,28 @@
package engine
import (
"fmt"
"log/slog"
"time"
)
type delayNodeExecutor struct {
nodeExecutor
}
func (ne *delayNodeExecutor) Execute(execCtx *NodeExecutionContext) (*NodeExecutionResult, error) {
execRes := newNodeExecutionResult(execCtx.Node)
nodeCfg := execCtx.Node.Data.Config.AsDelay()
ne.logger.Info(fmt.Sprintf("delay for %d second(s) before continuing ...", nodeCfg.Wait))
time.Sleep(time.Duration(nodeCfg.Wait) * time.Second)
return execRes, nil
}
func newDelayNodeExecutor() NodeExecutor {
return &delayNodeExecutor{
nodeExecutor: nodeExecutor{logger: slog.Default()},
}
}
+1
View File
@@ -16,6 +16,7 @@ const (
NodeTypeTryCatch = domain.WorkflowNodeTypeTryCatch
NodeTypeTryBlock = domain.WorkflowNodeTypeTryBlock
NodeTypeCatchBlock = domain.WorkflowNodeTypeCatchBlock
NodeTypeDelay = domain.WorkflowNodeTypeDelay
NodeTypeBizApply = domain.WorkflowNodeTypeBizApply
NodeTypeBizUpload = domain.WorkflowNodeTypeBizUpload
NodeTypeBizMonitor = domain.WorkflowNodeTypeBizMonitor
+1 -1
View File
@@ -81,7 +81,7 @@ func (s *WorkflowService) StartRun(ctx context.Context, req *dtos.WorkflowStartR
return nil, err
}
if workflow.LastRunStatus == domain.WorkflowRunStatusTypePending || workflow.LastRunStatus == domain.WorkflowRunStatusTypeProcessing {
if req.RunTrigger == domain.WorkflowTriggerTypeManual && (workflow.LastRunStatus == domain.WorkflowRunStatusTypePending || workflow.LastRunStatus == domain.WorkflowRunStatusTypeProcessing) {
return nil, errors.New("workflow is already pending or processing")
} else if workflow.GraphContent == nil {
return nil, errors.New("workflow graph content is empty")
@@ -10,6 +10,7 @@ import BizMonitorNodeConfigDrawer from "./forms/BizMonitorNodeConfigDrawer";
import BizNotifyNodeConfigDrawer from "./forms/BizNotifyNodeConfigDrawer";
import BizUploadNodeConfigDrawer from "./forms/BizUploadNodeConfigDrawer";
import BranchBlockNodeConfigDrawer from "./forms/BranchBlockNodeConfigDrawer";
import DelayNodeConfigDrawer from "./forms/DelayNodeConfigDrawer";
import StartNodeConfigDrawer from "./forms/StartNodeConfigDrawer";
import { NodeType } from "./nodes/typings";
@@ -50,6 +51,10 @@ const NodeDrawer = ({ node, trigger, ...props }: NodeDrawerProps) => {
{node?.flowNodeType === NodeType.Start ? (
<StartNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.Delay ? (
<DelayNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.BranchBlock ? (
<BranchBlockNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.BizApply ? (
<BizApplyNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.BizUpload ? (
@@ -60,8 +65,6 @@ const NodeDrawer = ({ node, trigger, ...props }: NodeDrawerProps) => {
<BizDeployNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.BizNotify ? (
<BizNotifyNodeConfigDrawer {...drawerProps} />
) : node?.flowNodeType === NodeType.BranchBlock ? (
<BranchBlockNodeConfigDrawer {...drawerProps} />
) : (
<></>
)}
@@ -0,0 +1,40 @@
import { useTranslation } from "react-i18next";
import { type FlowNodeEntity } from "@flowgram.ai/fixed-layout-editor";
import { Form } from "antd";
import { NodeConfigDrawer } from "./_shared";
import DelayNodeConfigForm from "./DelayNodeConfigForm";
import { NodeType } from "../nodes/typings";
export interface DelayNodeConfigDrawerProps {
afterClose?: () => void;
loading?: boolean;
node: FlowNodeEntity;
open?: boolean;
onOpenChange?: (open: boolean) => void;
}
const DelayNodeConfigDrawer = ({ node, ...props }: DelayNodeConfigDrawerProps) => {
if (node.flowNodeType !== NodeType.Delay) {
console.warn(`[certimate] current workflow node type is not: ${NodeType.Delay}`);
}
const { i18n } = useTranslation();
const [formInst] = Form.useForm();
return (
<NodeConfigDrawer
anchor={{
items: DelayNodeConfigForm.getAnchorItems({ i18n }),
}}
form={formInst}
node={node}
{...props}
>
<DelayNodeConfigForm form={formInst} node={node} />
</NodeConfigDrawer>
);
};
export default DelayNodeConfigDrawer;
@@ -0,0 +1,92 @@
import { useMemo } from "react";
import { getI18n, useTranslation } from "react-i18next";
import { type FlowNodeEntity, getNodeForm } from "@flowgram.ai/fixed-layout-editor";
import { type AnchorProps, Form, type FormInstance, InputNumber } from "antd";
import { createSchemaFieldRule } from "antd-zod";
import { z } from "zod";
import { type WorkflowNodeConfigForDelay, defaultNodeConfigForDelay } from "@/domain/workflow";
import { useAntdForm } from "@/hooks";
import { NodeFormContextProvider } from "./_context";
import { NodeType } from "../nodes/typings";
export interface DelayNodeConfigFormProps {
form: FormInstance;
node: FlowNodeEntity;
}
const DelayNodeConfigForm = ({ node, ...props }: DelayNodeConfigFormProps) => {
if (node.flowNodeType !== NodeType.Delay) {
console.warn(`[certimate] current workflow node type is not: ${NodeType.Delay}`);
}
const { i18n, t } = useTranslation();
const initialValues = useMemo(() => {
return getNodeForm(node)?.getValueIn("config") as WorkflowNodeConfigForDelay | undefined;
}, [node]);
const formSchema = getSchema({ i18n });
const formRule = createSchemaFieldRule(formSchema);
const { form: formInst, formProps } = useAntdForm({
form: props.form,
name: "workflowNodeDelayConfigForm",
initialValues: initialValues ?? getInitialValues(),
});
return (
<NodeFormContextProvider value={{ node }}>
<Form {...formProps} clearOnDestroy={true} form={formInst} layout="vertical" preserve={false} scrollToFirstError>
<div id="parameters" data-anchor="parameters">
<Form.Item name="wait" label={t("workflow_node.delay.form.wait.label")} rules={[formRule]}>
<InputNumber
style={{ width: "100%" }}
min={0}
max={3600}
placeholder={t("workflow_node.delay.form.wait.placeholder")}
addonAfter={t("workflow_node.delay.form.wait.unit")}
/>
</Form.Item>
</div>
</Form>
</NodeFormContextProvider>
);
};
const getAnchorItems = ({ i18n = getI18n() }: { i18n?: ReturnType<typeof getI18n> }): Required<AnchorProps>["items"] => {
const { t } = i18n;
return ["parameters"].map((key) => ({
key: key,
title: t(`workflow_node.delay.form_anchor.${key}.tab`),
href: "#" + key,
}));
};
const getInitialValues = (): Nullish<z.infer<ReturnType<typeof getSchema>>> => {
return {
...defaultNodeConfigForDelay(),
};
};
const getSchema = ({ i18n = getI18n() }: { i18n?: ReturnType<typeof getI18n> }) => {
const { t } = i18n;
return z.object({
wait: z.preprocess(
(v) => Number(v),
z
.number(t("workflow_node.delay.form.wait.placeholder"))
.int(t("workflow_node.delay.form.wait.placeholder"))
.gte(1, t("workflow_node.delay.form.wait.placeholder"))
),
});
};
const _default = Object.assign(DelayNodeConfigForm, {
getAnchorItems,
getSchema,
});
export default _default;
@@ -24,6 +24,7 @@ export const BizApplyNodeRegistry: NodeRegistry = {
iconBgColor: "#5b65f5",
clickable: true,
expandable: false,
},
formMeta: {
@@ -25,6 +25,7 @@ export const BizDeployNodeRegistry: NodeRegistry = {
iconBgColor: "#5b65f5",
clickable: true,
expandable: false,
},
formMeta: {
@@ -22,6 +22,7 @@ export const BizMonitorNodeRegistry: NodeRegistry = {
iconBgColor: "#5b65f5",
clickable: true,
expandable: false,
},
formMeta: {
@@ -24,6 +24,7 @@ export const BizNotifyNodeRegistry: NodeRegistry = {
iconBgColor: "#0693d4",
clickable: true,
expandable: false,
},
formMeta: {
@@ -22,6 +22,7 @@ export const BizUploadNodeRegistry: NodeRegistry = {
iconBgColor: "#5b65f5",
clickable: true,
expandable: false,
},
formMeta: {
@@ -0,0 +1,63 @@
import { getI18n } from "react-i18next";
import { FeedbackLevel, Field } from "@flowgram.ai/fixed-layout-editor";
import { IconHourglassHigh } from "@tabler/icons-react";
import { newNode } from "@/domain/workflow";
import { BaseNode } from "./_shared";
import { NodeKindType, type NodeRegistry, NodeType } from "./typings";
import DelayNodeConfigForm from "../forms/DelayNodeConfigForm";
export const DelayNodeRegistry: NodeRegistry = {
type: NodeType.Delay,
kind: NodeKindType.Basis,
meta: {
helpText: getI18n().t("workflow_node.delay.help"),
labelText: getI18n().t("workflow_node.delay.label"),
icon: IconHourglassHigh,
iconColor: "#2a354c",
iconBgColor: "#fed421",
clickable: true,
expandable: false,
},
formMeta: {
validate: {
["config"]: ({ value }) => {
const res = DelayNodeConfigForm.getSchema({}).safeParse(value);
if (!res.success) {
return {
message: res.error.message,
level: FeedbackLevel.Error,
};
}
},
},
render: () => {
const { t } = getI18n();
return (
<BaseNode
description={
<Field name="config.wait">
{({ field: { value } }) => (
<>
<div>{value != null ? `${value} ${t("workflow_node.delay.form.wait.unit")}` : t("workflow.detail.design.editor.placeholder")}</div>
</>
)}
</Field>
}
/>
);
},
},
onAdd() {
return newNode(NodeType.Delay, { i18n: getI18n() });
},
};
@@ -4,6 +4,7 @@ import { BizMonitorNodeRegistry } from "./BizMonitorNodeRegistry";
import { BizNotifyNodeRegistry } from "./BizNotifyNodeRegistry";
import { BizUploadNodeRegistry } from "./BizUploadNodeRegistry";
import { BranchBlockNodeRegistry, ConditionNodeRegistry } from "./ConditionNode";
import { DelayNodeRegistry } from "./DelayNode";
import { EndNodeRegistry } from "./EndNode";
import { StartNodeRegistry } from "./StartNode";
import { CatchBlockNodeRegistry, TryCatchNodeRegistry } from "./TryCatchNode";
@@ -12,6 +13,7 @@ export const getAllNodeRegistries = () => {
return [
StartNodeRegistry,
EndNodeRegistry,
DelayNodeRegistry,
BizApplyNodeRegistry,
BizUploadNodeRegistry,
BizMonitorNodeRegistry,
@@ -13,6 +13,7 @@ import { WORKFLOW_NODE_TYPES, type WorkflowNode } from "@/domain/workflow";
export enum NodeType {
Start = "start",
End = "end",
Delay = "delay",
Condition = "condition",
BranchBlock = "branchBlock",
TryCatch = "tryCatch",
@@ -28,6 +29,7 @@ export enum NodeType {
/* TYPE GUARD, PLEASE DO NOT REMOVE THESE! */
console.assert(NodeType.Start === WORKFLOW_NODE_TYPES.START);
console.assert(NodeType.End === WORKFLOW_NODE_TYPES.END);
console.assert(NodeType.Delay === WORKFLOW_NODE_TYPES.DELAY);
console.assert(NodeType.Condition === WORKFLOW_NODE_TYPES.CONDITION);
console.assert(NodeType.BranchBlock === WORKFLOW_NODE_TYPES.BRANCHBLOCK);
console.assert(NodeType.TryCatch === WORKFLOW_NODE_TYPES.TRYCATCH);
+19
View File
@@ -37,6 +37,7 @@ export type WorkflowTriggerType = (typeof WORKFLOW_TRIGGERS)[keyof typeof WORKFL
export const WORKFLOW_NODE_TYPES = Object.freeze({
START: "start",
END: "end",
DELAY: "delay",
CONDITION: "condition",
BRANCHBLOCK: "branchBlock",
TRYCATCH: "tryCatch",
@@ -73,6 +74,14 @@ export const defaultNodeConfigForStart = (): Partial<WorkflowNodeConfigForStart>
};
};
export type WorkflowNodeConfigForDelay = {
wait?: number;
};
export const defaultNodeConfigForDelay = (): Partial<WorkflowNodeConfigForDelay> => {
return {};
};
export type WorkflowNodeConfigForBranchBlock = {
expression?: Expr;
};
@@ -189,6 +198,16 @@ export const newNode = (type: WorkflowNodeType, { i18n = getI18n() }: { i18n?: R
},
};
case WORKFLOW_NODE_TYPES.DELAY:
return {
id: newNodeId(),
type: type,
data: {
name: t("workflow_node.delay.default_name"),
config: defaultNodeConfigForDelay(),
},
};
case WORKFLOW_NODE_TYPES.CONDITION: {
const branch1 = newNode(WORKFLOW_NODE_TYPES.BRANCHBLOCK, { i18n });
branch1.data.name = `${branch1.data.name} 1`;
@@ -1025,6 +1025,14 @@
"workflow_node.notify.form.skip_on_all_prev_skipped.switch.on": "skip",
"workflow_node.notify.form.skip_on_all_prev_skipped.switch.off": "not skip",
"workflow_node.delay.label": "Delay",
"workflow_node.delay.help": "Pause the execution, and wait for a certain amount of time.",
"workflow_node.delay.default_name": "Delay",
"workflow_node.delay.form_anchor.parameters.tab": "Parameters",
"workflow_node.delay.form.wait.label": "Waiting time",
"workflow_node.delay.form.wait.placeholder": "Please enter waiting time",
"workflow_node.delay.form.wait.unit": "seconds",
"workflow_node.condition.label": "Parallel/Conditional branch",
"workflow_node.condition.help": "When the specified conditions are met, enter the corresponding branch. The failure of a node in a certain branch does not affect the continuation of parallel branch execution.",
"workflow_node.condition.default_name": "Parallel",
@@ -1023,6 +1023,14 @@
"workflow_node.notify.form.skip_on_all_prev_skipped.switch.on": "跳过",
"workflow_node.notify.form.skip_on_all_prev_skipped.switch.off": "不跳过",
"workflow_node.delay.label": "延迟等待",
"workflow_node.delay.help": "暂停执行工作流并等待一段时间。这在某些需要降低调用执行频率的场景很有用。",
"workflow_node.delay.default_name": "延迟",
"workflow_node.delay.form_anchor.parameters.tab": "参数设置",
"workflow_node.delay.form.wait.label": "等待时间",
"workflow_node.delay.form.wait.placeholder": "请输入等待时间",
"workflow_node.delay.form.wait.unit": "秒",
"workflow_node.condition.label": "并行/条件分支",
"workflow_node.condition.help": "当满足指定的条件时,进入相应分支。某一分支中的节点执行失败不影响平行分支继续执行。",
"workflow_node.condition.default_name": "并行",
@@ -1,4 +1,4 @@
import { useEffect, useMemo, useRef, useState } from "react";
import { useMemo, useRef, useState } from "react";
import { useTranslation } from "react-i18next";
import { IconArrowBackUp, IconDots } from "@tabler/icons-react";
import { useDeepCompareEffect } from "ahooks";
@@ -7,7 +7,6 @@ import { debounce } from "radash";
import Show from "@/components/Show";
import { WorkflowDesigner, type WorkflowDesignerInstance, WorkflowNodeDrawer, WorkflowToolbar } from "@/components/workflow/designer";
import { WORKFLOW_RUN_STATUSES } from "@/domain/workflowRun";
import { useZustandShallowSelector } from "@/hooks";
import { useWorkflowStore } from "@/stores/workflow";
import { getErrMsg } from "@/utils/error";
@@ -20,16 +19,8 @@ const WorkflowDetailDesign = () => {
const { workflow, ...workflowStore } = useWorkflowStore(useZustandShallowSelector(["workflow", "orchestrate", "publish", "rollback"]));
const [workflowRunDisabled, setWorkflowRunDisabled] = useState(false);
const workflowRollbackDisabled = useMemo(
() => workflowRunDisabled || !workflow.hasDraft || !workflow.hasContent,
[workflowRunDisabled, workflow.hasDraft, workflow.hasContent]
);
const workflowPublishDisabled = useMemo(() => workflowRunDisabled || !workflow.hasDraft, [workflowRunDisabled, workflow.hasDraft]);
useEffect(() => {
const running = workflow.lastRunStatus === WORKFLOW_RUN_STATUSES.PENDING || workflow.lastRunStatus === WORKFLOW_RUN_STATUSES.PROCESSING;
setWorkflowRunDisabled(running);
}, [workflow.lastRunStatus]);
const workflowRollbackDisabled = useMemo(() => !workflow.hasDraft || !workflow.hasContent, [workflow.hasDraft, workflow.hasContent]);
const workflowPublishDisabled = useMemo(() => !workflow.hasDraft, [workflow.hasDraft]);
const designerRef = useRef<WorkflowDesignerInstance>(null);
const designerPending = useRef(false); // 保存中时阻止刷新画布