// 工作流引擎核心逻辑 package workflow import ( "encoding/json" "errors" "fmt" "strconv" "strings" "time" "com.sclktx/m/v2/internal/common/model" "gorm.io/gorm" ) // Engine 工作流引擎 type Engine struct { db *gorm.DB } // NewEngine 创建工作流引擎 func NewEngine(db *gorm.DB) *Engine { return &Engine{db: db} } // Node 节点配置(从 nodes_json 解析) type Node struct { NodeID string `json:"node_id"` NodeName string `json:"node_name"` NodeType int `json:"node_type"` // 0-开始 1-审批 2-网关 3-结束 PrevNodeIDs []string `json:"prev_node_ids"` UserIDs []string `json:"user_ids"` Roles []string `json:"roles"` GwConfig *GatewayConfig `json:"gw_config"` IsCosigned int `json:"is_cosigned"` // 0-否 1-会签 NodeStartEvents []string `json:"node_start_events"` NodeEndEvents []string `json:"node_end_events"` TaskFinishEvents []string `json:"task_finish_events"` } // GatewayConfig 网关配置 type GatewayConfig struct { Conditions []Condition `json:"conditions"` InevitableNodes []string `json:"inevitable_nodes"` WaitForAllPrevNode int `json:"wait_for_all_prev_node"` // 0-否 1-是 } // Condition 条件分支 type Condition struct { Expression string `json:"expression"` NodeID string `json:"node_id"` } // parseNodes 解析流程定义的节点配置 func (e *Engine) parseNodes(def *model.ProcessDefinition) ([]Node, error) { if def.NodesJSON == "" { return nil, errors.New("流程节点配置为空") } var nodes []Node if err := json.Unmarshal([]byte(def.NodesJSON), &nodes); err != nil { return nil, fmt.Errorf("解析节点配置失败: %w", err) } return nodes, nil } // findStartNode 查找开始节点 func findStartNode(nodes []Node) *Node { for i := range nodes { if nodes[i].NodeType == 0 { return &nodes[i] } } return nil } // findNextNodes 查找下一个节点 func findNextNodes(nodes []Node, currentNodeID string) []Node { var next []Node for _, n := range nodes { for _, prevID := range n.PrevNodeIDs { if prevID == currentNodeID { next = append(next, n) break } } } return next } // findNodeByID 根据节点 ID 查找节点 func findNodeByID(nodes []Node, nodeID string) *Node { for i := range nodes { if nodes[i].NodeID == nodeID { return &nodes[i] } } return nil } // parseVariables 解析流程实例中的变量 JSON func (e *Engine) parseVariables(variablesJSON string) map[string]interface{} { var vars map[string]interface{} if variablesJSON != "" { json.Unmarshal([]byte(variablesJSON), &vars) } if vars == nil { vars = make(map[string]interface{}) } return vars } // StartInstanceByCode 启动流程实例(按流程编码) func (e *Engine) StartInstanceByCode(processCode string, businessType string, businessID uint, starterID uint, variables map[string]interface{}) (*model.ProcessInstance, error) { var def model.ProcessDefinition if err := e.db.Where("process_code = ?", processCode).First(&def).Error; err != nil { return nil, fmt.Errorf("流程定义不存在: %s", processCode) } return e.StartInstance(def.ID, businessType, businessID, starterID, variables) } func (e *Engine) StartInstance(defID uint, businessType string, businessID uint, starterID uint, variables map[string]interface{}) (*model.ProcessInstance, error) { // 查询流程定义 var def model.ProcessDefinition if err := e.db.First(&def, defID).Error; err != nil { return nil, errors.New("流程定义不存在") } if !def.IsActive { return nil, errors.New("流程定义已禁用") } // 解析节点配置 nodes, err := e.parseNodes(&def) if err != nil { return nil, err } // 查找开始节点 startNode := findStartNode(nodes) if startNode == nil { return nil, errors.New("流程定义缺少开始节点") } // 序列化变量 variablesJSON := "{}" if variables != nil { b, _ := json.Marshal(variables) variablesJSON = string(b) } // 创建流程实例 instance := &model.ProcessInstance{ ProcessDefID: def.ID, BusinessType: businessType, BusinessID: businessID, Status: 0, // 运行中 StarterID: starterID, CurrentNodeID: startNode.NodeID, Variables: variablesJSON, } if err := e.db.Create(instance).Error; err != nil { return nil, fmt.Errorf("创建流程实例失败: %w", err) } // 流转到下一个节点 if err := e.moveToNextNodes(instance, nodes, startNode.NodeID); err != nil { return nil, err } return instance, nil } // moveToNextNodes 流转到下一个节点,创建任务(使用默认 DB) func (e *Engine) moveToNextNodes(instance *model.ProcessInstance, nodes []Node, currentNodeID string) error { return e.moveToNextNodesTx(e.db, instance, nodes, currentNodeID) } // moveToNextNodesTx 流转到下一个节点(支持外部事务) func (e *Engine) moveToNextNodesTx(tx *gorm.DB, instance *model.ProcessInstance, nodes []Node, currentNodeID string) error { nextNodes := findNextNodes(nodes, currentNodeID) // 如果没有下一个节点(到达结束节点),标记流程完成 if len(nextNodes) == 0 { now := time.Now() if err := tx.Model(instance).Updates(map[string]interface{}{ "status": 1, // 已完成 "finished_at": now, }).Error; err != nil { return err } return e.updateAppointmentStatus(tx, instance, 1) } // 为每个下一个节点创建任务 for _, node := range nextNodes { switch node.NodeType { case 1: // 审批节点 if err := e.createApprovalTasksTx(tx, instance, nodes, &node); err != nil { return err } case 3: // 结束节点 now := time.Now() if err := tx.Model(instance).Updates(map[string]interface{}{ "status": 1, "finished_at": now, }).Error; err != nil { return err } return e.updateAppointmentStatus(tx, instance, 1) case 2: // 网关节点 - 自动判断条件 if err := e.handleGatewayNodeTx(tx, instance, nodes, &node); err != nil { return err } } } // 更新当前节点 if len(nextNodes) > 0 { tx.Model(instance).Update("current_node_id", nextNodes[0].NodeID) } return nil } // handleGatewayNodeTx 处理网关节点(支持外部事务) func (e *Engine) handleGatewayNodeTx(tx *gorm.DB, instance *model.ProcessInstance, nodes []Node, gwNode *Node) error { if gwNode.GwConfig == nil { return e.moveToNextNodesTx(tx, instance, nodes, gwNode.NodeID) } // 解析流程变量 var variables map[string]interface{} if instance.Variables != "" { json.Unmarshal([]byte(instance.Variables), &variables) } if variables == nil { variables = make(map[string]interface{}) } // 评估条件 for _, cond := range gwNode.GwConfig.Conditions { if evaluateCondition(cond.Expression, variables) { return e.moveToNextNodesTx(tx, instance, nodes, cond.NodeID) } } // 没有条件满足,走必经节点 if len(gwNode.GwConfig.InevitableNodes) > 0 { return e.moveToNextNodesTx(tx, instance, nodes, gwNode.GwConfig.InevitableNodes[0]) } return e.moveToNextNodesTx(tx, instance, nodes, gwNode.NodeID) } // createApprovalTasks 为审批节点创建任务(使用默认 DB) func (e *Engine) createApprovalTasks(instance *model.ProcessInstance, nodes []Node, node *Node) error { return e.createApprovalTasksTx(e.db, instance, nodes, node) } // createApprovalTasksTx 为审批节点创建任务(支持外部事务) func (e *Engine) createApprovalTasksTx(tx *gorm.DB, instance *model.ProcessInstance, nodes []Node, node *Node) error { // 解析审批人:优先使用 user_ids,其次使用 roles var assigneeIDs []uint variables := e.parseVariables(instance.Variables) if len(node.UserIDs) > 0 { for _, uidStr := range node.UserIDs { if strings.HasPrefix(uidStr, "$") { // $变量名 → 从流程变量中取值 varName := uidStr[1:] if val, ok := variables[varName]; ok { switch v := val.(type) { case float64: assigneeIDs = append(assigneeIDs, uint(v)) case uint: assigneeIDs = append(assigneeIDs, v) case int: assigneeIDs = append(assigneeIDs, uint(v)) case string: // 先尝试解析为单个 ID if id, err := strconv.ParseUint(v, 10, 64); err == nil { assigneeIDs = append(assigneeIDs, uint(id)) } else { // 再尝试解析为 JSON 数组(如 "[1,2,3]") var ids []uint if err := json.Unmarshal([]byte(v), &ids); err == nil { assigneeIDs = append(assigneeIDs, ids...) } } } } } else { // 静态数字 ID var uid uint if _, err := fmt.Sscanf(uidStr, "%d", &uid); err == nil { assigneeIDs = append(assigneeIDs, uid) } } } } if len(assigneeIDs) == 0 && len(node.Roles) > 0 { // 根据角色编码查找用户 var users []model.User tx.Where("role IN ?", node.Roles).Find(&users) for _, u := range users { assigneeIDs = append(assigneeIDs, u.ID) } } if len(assigneeIDs) == 0 { // 没有找到审批人 — 跳过此节点,自动流转到下一节点(适用于未配置分管领导等情况) return e.moveToNextNodesTx(tx, instance, nodes, node.NodeID) } // 为每个审批人创建任务 for _, assigneeID := range assigneeIDs { task := &model.Task{ InstanceID: instance.ID, NodeID: node.NodeID, NodeName: node.NodeName, Status: 0, // 待处理 AssigneeID: &assigneeID, } if err := tx.Create(task).Error; err != nil { return fmt.Errorf("创建任务失败: %w", err) } } return nil } // handleGatewayNode 处理网关节点 func (e *Engine) handleGatewayNode(instance *model.ProcessInstance, nodes []Node, gwNode *Node) error { if gwNode.GwConfig == nil { return e.moveToNextNodes(instance, nodes, gwNode.NodeID) } // 解析流程变量 var variables map[string]interface{} if instance.Variables != "" { json.Unmarshal([]byte(instance.Variables), &variables) } if variables == nil { variables = make(map[string]interface{}) } // 评估条件 for _, cond := range gwNode.GwConfig.Conditions { if evaluateCondition(cond.Expression, variables) { return e.moveToNextNodes(instance, nodes, cond.NodeID) } } // 没有条件满足,走必经节点 if len(gwNode.GwConfig.InevitableNodes) > 0 { return e.moveToNextNodes(instance, nodes, gwNode.GwConfig.InevitableNodes[0]) } return e.moveToNextNodes(instance, nodes, gwNode.NodeID) } // evaluateCondition 简单条件评估(支持 $days>=3 这样的表达式) func evaluateCondition(expression string, variables map[string]interface{}) bool { // 简单实现:解析 "$key op value" 格式 var key string var op string var expected float64 n, _ := fmt.Sscanf(expression, "$%s %s %f", &key, &op, &expected) if n != 3 { return false } val, ok := variables[key] if !ok { return false } var numVal float64 switch v := val.(type) { case float64: numVal = v case int: numVal = float64(v) case json.Number: numVal, _ = v.Float64() default: return false } switch op { case ">=": return numVal >= expected case ">": return numVal > expected case "<=": return numVal <= expected case "<": return numVal < expected case "==": return numVal == expected case "!=": return numVal != expected } return false } // ApproveTask 审批通过任务 func (e *Engine) ApproveTask(taskID uint, userID uint, comment string) error { var task model.Task if err := e.db.First(&task, taskID).Error; err != nil { return errors.New("任务不存在") } if task.Status != 0 { return errors.New("任务已处理") } // 验证审批人身份 if task.AssigneeID == nil || *task.AssigneeID != userID { return errors.New("您不是该任务的审批人") } now := time.Now() // 开启事务,保证一致性 tx := e.db.Begin() if err := tx.Model(&task).Updates(map[string]interface{}{ "status": 1, // 已完成 "comment": comment, "handled_at": now, }).Error; err != nil { tx.Rollback() return fmt.Errorf("更新任务失败: %w", err) } // 获取实例和流程定义 var instance model.ProcessInstance if err := tx.First(&instance, task.InstanceID).Error; err != nil { tx.Rollback() return err } var def model.ProcessDefinition if err := tx.First(&def, instance.ProcessDefID).Error; err != nil { tx.Rollback() return err } nodes, err := e.parseNodes(&def) if err != nil { tx.Rollback() return err } // 查找当前节点的定义 currentNode := findNodeByID(nodes, task.NodeID) needAllApprovers := currentNode != nil && currentNode.IsCosigned == 1 if needAllApprovers { // 会签模式(is_cosigned=1):需要所有任务都完成 var pendingCount int64 tx.Model(&model.Task{}). Where("instance_id = ? AND node_id = ? AND status = 0", task.InstanceID, task.NodeID). Count(&pendingCount) if pendingCount > 0 { tx.Commit() return nil // 等待其他审批人 } } // 非会签模式(is_cosigned=0 或不设置):一人通过即流转下一节点 // 取消本节点其他待处理任务,防止并发重复创建下一节点 if !needAllApprovers { tx.Model(&model.Task{}). Where("instance_id = ? AND node_id = ? AND status = 0 AND id != ?", task.InstanceID, task.NodeID, task.ID). Update("status", 3) // 已取消 } // 流转到下一个节点(在同一事务中) if err := e.moveToNextNodesTx(tx, &instance, nodes, task.NodeID); err != nil { tx.Rollback() return err } return tx.Commit().Error } // RejectTask 审批拒绝任务 func (e *Engine) RejectTask(taskID uint, userID uint, comment string) error { var task model.Task if err := e.db.First(&task, taskID).Error; err != nil { return errors.New("任务不存在") } if task.Status != 0 { return errors.New("任务已处理") } // 验证审批人身份 if task.AssigneeID == nil || *task.AssigneeID != userID { return errors.New("您不是该任务的审批人") } now := time.Now() tx := e.db.Begin() if err := tx.Model(&task).Updates(map[string]interface{}{ "status": 2, // 已拒绝 "comment": comment, "handled_at": now, }).Error; err != nil { tx.Rollback() return fmt.Errorf("更新任务失败: %w", err) } // 终止流程实例 var instance model.ProcessInstance if err := tx.First(&instance, task.InstanceID).Error; err != nil { tx.Rollback() return err } if err := tx.Model(&instance).Updates(map[string]interface{}{ "status": 2, // 已终止 "finished_at": now, }).Error; err != nil { tx.Rollback() return err } // 更新业务状态 if err := e.updateAppointmentStatus(tx, &instance, 2); err != nil { tx.Rollback() return err } return tx.Commit().Error } // checkNodeCompletion 检查节点所有任务是否完成 func (e *Engine) checkNodeCompletion(instanceID uint, nodeID string) (*Node, bool, error) { // 检查是否还有待处理的任务 var pendingCount int64 e.db.Model(&model.Task{}). Where("instance_id = ? AND node_id = ? AND status = 0", instanceID, nodeID). Count(&pendingCount) return nil, pendingCount == 0, nil } // updateAppointmentStatus 工作流完成后更新业务表状态 func (e *Engine) updateAppointmentStatus(tx *gorm.DB, instance *model.ProcessInstance, apptStatus int) error { if instance.BusinessType != "appointment" { return nil // 只处理预约业务 } return tx.Model(&model.Appointment{}).Where("id = ?", instance.BusinessID).Updates(map[string]interface{}{ "status": apptStatus, }).Error }