|
@@ -15,6 +15,9 @@ public class FlowEngineService : ITransient
|
|
|
private readonly SqlSugarRepository<ApprovalFlowLog> _logRep;
|
|
private readonly SqlSugarRepository<ApprovalFlowLog> _logRep;
|
|
|
private readonly SqlSugarRepository<SysUserRole> _userRoleRep;
|
|
private readonly SqlSugarRepository<SysUserRole> _userRoleRep;
|
|
|
private readonly SqlSugarRepository<SysUser> _userRep;
|
|
private readonly SqlSugarRepository<SysUser> _userRep;
|
|
|
|
|
+ private readonly SqlSugarRepository<SysOrg> _orgRep;
|
|
|
|
|
+ private readonly SqlSugarRepository<ApprovalFlowDelegate> _delegateRep;
|
|
|
|
|
+ private readonly SqlSugarRepository<ApprovalFlowCompletedNode> _completedNodeRep;
|
|
|
private readonly UserManager _userManager;
|
|
private readonly UserManager _userManager;
|
|
|
private readonly SysOrgService _sysOrgService;
|
|
private readonly SysOrgService _sysOrgService;
|
|
|
private readonly FlowNotifyService _notifyService;
|
|
private readonly FlowNotifyService _notifyService;
|
|
@@ -26,6 +29,9 @@ public class FlowEngineService : ITransient
|
|
|
SqlSugarRepository<ApprovalFlowLog> logRep,
|
|
SqlSugarRepository<ApprovalFlowLog> logRep,
|
|
|
SqlSugarRepository<SysUserRole> userRoleRep,
|
|
SqlSugarRepository<SysUserRole> userRoleRep,
|
|
|
SqlSugarRepository<SysUser> userRep,
|
|
SqlSugarRepository<SysUser> userRep,
|
|
|
|
|
+ SqlSugarRepository<SysOrg> orgRep,
|
|
|
|
|
+ SqlSugarRepository<ApprovalFlowDelegate> delegateRep,
|
|
|
|
|
+ SqlSugarRepository<ApprovalFlowCompletedNode> completedNodeRep,
|
|
|
UserManager userManager,
|
|
UserManager userManager,
|
|
|
SysOrgService sysOrgService,
|
|
SysOrgService sysOrgService,
|
|
|
FlowNotifyService notifyService)
|
|
FlowNotifyService notifyService)
|
|
@@ -36,6 +42,9 @@ public class FlowEngineService : ITransient
|
|
|
_logRep = logRep;
|
|
_logRep = logRep;
|
|
|
_userRoleRep = userRoleRep;
|
|
_userRoleRep = userRoleRep;
|
|
|
_userRep = userRep;
|
|
_userRep = userRep;
|
|
|
|
|
+ _orgRep = orgRep;
|
|
|
|
|
+ _delegateRep = delegateRep;
|
|
|
|
|
+ _completedNodeRep = completedNodeRep;
|
|
|
_userManager = userManager;
|
|
_userManager = userManager;
|
|
|
_sysOrgService = sysOrgService;
|
|
_sysOrgService = sysOrgService;
|
|
|
_notifyService = notifyService;
|
|
_notifyService = notifyService;
|
|
@@ -93,14 +102,19 @@ public class FlowEngineService : ITransient
|
|
|
n.Type is "bpmn:startEvent" or "start-node")
|
|
n.Type is "bpmn:startEvent" or "start-node")
|
|
|
?? throw Oops.Oh("流程图中未找到开始节点");
|
|
?? throw Oops.Oh("流程图中未找到开始节点");
|
|
|
|
|
|
|
|
- var firstTaskNodeId = FindNextNodeId(flowData, startNode.Id);
|
|
|
|
|
- instance.CurrentNodeId = firstTaskNodeId;
|
|
|
|
|
-
|
|
|
|
|
await _instanceRep.InsertAsync(instance);
|
|
await _instanceRep.InsertAsync(instance);
|
|
|
|
|
|
|
|
await WriteLog(instance.Id, null, startNode.Id, FlowLogActionEnum.Submit, input.Comment);
|
|
await WriteLog(instance.Id, null, startNode.Id, FlowLogActionEnum.Submit, input.Comment);
|
|
|
|
|
|
|
|
- await CreateTasksForNode(instance, flowData, firstTaskNodeId);
|
|
|
|
|
|
|
+ // 开始节点登记为已完成,随后沿出边推进(兼容首节点即为并行网关 Fork 的拓扑)
|
|
|
|
|
+ await MarkNodeCompleted(instance.Id, startNode);
|
|
|
|
|
+ var firstOutgoing = flowData.Edges.Where(e => e.SourceNodeId == startNode.Id).Select(e => e.TargetNodeId).ToList();
|
|
|
|
|
+ if (firstOutgoing.Count == 0)
|
|
|
|
|
+ throw Oops.Oh("流程图开始节点未连接任何后继节点");
|
|
|
|
|
+ foreach (var target in firstOutgoing)
|
|
|
|
|
+ {
|
|
|
|
|
+ await ProcessNextNode(instance, flowData, target);
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
await InvokeHandler(input.BizType, h => h.OnFlowStarted(input.BizId, instance.Id));
|
|
await InvokeHandler(input.BizType, h => h.OnFlowStarted(input.BizId, instance.Id));
|
|
|
|
|
|
|
@@ -118,6 +132,9 @@ public class FlowEngineService : ITransient
|
|
|
task.ActionTime = DateTime.Now;
|
|
task.ActionTime = DateTime.Now;
|
|
|
await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
|
|
|
|
|
+ // P1-6 审批代理:同步取消配对的本人/代理任务
|
|
|
|
|
+ await CancelPairedDelegateTask(task);
|
|
|
|
|
+
|
|
|
await WriteLog(task.InstanceId, taskId, task.NodeId, FlowLogActionEnum.Approve, comment);
|
|
await WriteLog(task.InstanceId, taskId, task.NodeId, FlowLogActionEnum.Approve, comment);
|
|
|
|
|
|
|
|
var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
|
|
var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
|
|
@@ -144,7 +161,8 @@ public class FlowEngineService : ITransient
|
|
|
task.ActionTime = DateTime.Now;
|
|
task.ActionTime = DateTime.Now;
|
|
|
await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
|
|
|
- await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
|
|
|
|
|
|
|
+ // Reject 导致流程终止:取消整个实例所有剩余 Pending 任务(含并行分支、配对代理任务等)
|
|
|
|
|
+ await CancelAllPendingTasks(task.InstanceId, task.Id);
|
|
|
|
|
|
|
|
var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
|
|
var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
|
|
|
?? throw Oops.Oh("流程实例不存在");
|
|
?? throw Oops.Oh("流程实例不存在");
|
|
@@ -302,6 +320,65 @@ public class FlowEngineService : ITransient
|
|
|
await _notifyService.NotifyAddSign(targetUserId, task.InstanceId, instance?.Title ?? "", _userManager.RealName);
|
|
await _notifyService.NotifyAddSign(targetUserId, task.InstanceId, instance?.Title ?? "", _userManager.RealName);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 手动升级 — 当前审批人主动将任务升级到更高层级
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ public async Task Escalate(long taskId, string? comment)
|
|
|
|
|
+ {
|
|
|
|
|
+ var task = await GetPendingTask(taskId);
|
|
|
|
|
+ var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
|
|
|
|
|
+ ?? throw Oops.Oh("流程实例不存在");
|
|
|
|
|
+
|
|
|
|
|
+ var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
|
|
|
|
|
+ var node = flowData.Nodes.FirstOrDefault(n => n.Id == task.NodeId)
|
|
|
|
|
+ ?? throw Oops.Oh("节点不存在");
|
|
|
|
|
+
|
|
|
|
|
+ var props = node.Properties;
|
|
|
|
|
+ if (props?.EnableManualEscalation != true
|
|
|
|
|
+ || string.IsNullOrWhiteSpace(props.EscalationApproverType)
|
|
|
|
|
+ || string.IsNullOrWhiteSpace(props.EscalationApproverIds))
|
|
|
|
|
+ throw Oops.Oh("该节点未配置升级目标,无法升级");
|
|
|
|
|
+
|
|
|
|
|
+ task.Status = FlowTaskStatusEnum.Escalated;
|
|
|
|
|
+ task.Comment = comment;
|
|
|
|
|
+ task.ActionTime = DateTime.Now;
|
|
|
|
|
+ await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
|
|
|
|
|
+
|
|
|
|
|
+ var escalationApprovers = await ResolveApprovers(
|
|
|
|
|
+ new FlowProperties
|
|
|
|
|
+ {
|
|
|
|
|
+ ApproverType = props.EscalationApproverType,
|
|
|
|
|
+ ApproverIds = props.EscalationApproverIds,
|
|
|
|
|
+ ApproverNames = props.EscalationApproverNames,
|
|
|
|
|
+ },
|
|
|
|
|
+ instance.InitiatorId);
|
|
|
|
|
+
|
|
|
|
|
+ if (escalationApprovers.Count == 0)
|
|
|
|
|
+ throw Oops.Oh("升级目标审批人列表为空");
|
|
|
|
|
+
|
|
|
|
|
+ var newTasks = escalationApprovers.Select(a => new ApprovalFlowTask
|
|
|
|
|
+ {
|
|
|
|
|
+ InstanceId = instance.Id,
|
|
|
|
|
+ NodeId = task.NodeId,
|
|
|
|
|
+ NodeName = task.NodeName,
|
|
|
|
|
+ AssigneeId = a.userId,
|
|
|
|
|
+ AssigneeName = a.userName,
|
|
|
|
|
+ Status = FlowTaskStatusEnum.Pending,
|
|
|
|
|
+ OrgId = instance.OrgId,
|
|
|
|
|
+ }).ToList();
|
|
|
|
|
+ await _taskRep.AsInsertable(newTasks).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ var targetNames = string.Join(", ", escalationApprovers.Select(a => a.userName));
|
|
|
|
|
+ await WriteLog(instance.Id, taskId, task.NodeId, FlowLogActionEnum.Escalate,
|
|
|
|
|
+ $"{comment} → 升级给 {targetNames}");
|
|
|
|
|
+
|
|
|
|
|
+ var targetUserIds = escalationApprovers.Select(a => a.userId).Distinct().ToList();
|
|
|
|
|
+ await _notifyService.NotifyEscalated(targetUserIds, instance.Id, instance.Title,
|
|
|
|
|
+ _userManager.RealName, task.NodeName);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
/// <summary>
|
|
/// <summary>
|
|
|
/// 催办
|
|
/// 催办
|
|
|
/// </summary>
|
|
/// </summary>
|
|
@@ -321,19 +398,188 @@ public class FlowEngineService : ITransient
|
|
|
await _notifyService.NotifyUrge(userIds, instanceId, instance.Title);
|
|
await _notifyService.NotifyUrge(userIds, instanceId, instance.Title);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // ═══════════════════════════════════════════
|
|
|
|
|
+ // 超时自动处理(由 FlowTimeoutJob 调用,无 UserManager 上下文)
|
|
|
|
|
+ // ═══════════════════════════════════════════
|
|
|
|
|
+
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 处理单个超时任务(由定时任务调用)
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ public async Task HandleTimeoutTask(long taskId)
|
|
|
|
|
+ {
|
|
|
|
|
+ var task = await _taskRep.GetByIdAsync(taskId);
|
|
|
|
|
+ if (task == null || task.Status != FlowTaskStatusEnum.Pending) return;
|
|
|
|
|
+
|
|
|
|
|
+ var instance = await _instanceRep.GetByIdAsync(task.InstanceId);
|
|
|
|
|
+ if (instance == null || instance.Status != FlowInstanceStatusEnum.Running) return;
|
|
|
|
|
+
|
|
|
|
|
+ var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
|
|
|
|
|
+ var node = flowData.Nodes.FirstOrDefault(n => n.Id == task.NodeId);
|
|
|
|
|
+ var props = node?.Properties;
|
|
|
|
|
+ if (props == null) return;
|
|
|
|
|
+
|
|
|
|
|
+ switch (props.TimeoutAction)
|
|
|
|
|
+ {
|
|
|
|
|
+ case "Notify":
|
|
|
|
|
+ await _notifyService.NotifyTimeout(
|
|
|
|
|
+ new List<long> { task.AssigneeId }, instance.Id, instance.Title);
|
|
|
|
|
+ await WriteSystemLog(instance.Id, task.Id, task.NodeId,
|
|
|
|
|
+ FlowLogActionEnum.AutoTimeout, "审批超时,已发送提醒通知");
|
|
|
|
|
+ break;
|
|
|
|
|
+ case "AutoApprove":
|
|
|
|
|
+ await AutoApproveTask(task, instance);
|
|
|
|
|
+ break;
|
|
|
|
|
+ case "AutoReject":
|
|
|
|
|
+ await AutoRejectTask(task, instance);
|
|
|
|
|
+ break;
|
|
|
|
|
+ case "AutoEscalate":
|
|
|
|
|
+ await AutoEscalateTask(task, props, instance);
|
|
|
|
|
+ break;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private async Task AutoApproveTask(ApprovalFlowTask task, ApprovalFlowInstance instance)
|
|
|
|
|
+ {
|
|
|
|
|
+ task.Status = FlowTaskStatusEnum.Approved;
|
|
|
|
|
+ task.Comment = "系统自动通过(超时)";
|
|
|
|
|
+ task.ActionTime = DateTime.Now;
|
|
|
|
|
+ await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ await WriteSystemLog(instance.Id, task.Id, task.NodeId,
|
|
|
|
|
+ FlowLogActionEnum.AutoTimeout, "审批超时,系统自动通过");
|
|
|
|
|
+
|
|
|
|
|
+ if (await IsNodeCompleted(instance, task.NodeId))
|
|
|
|
|
+ {
|
|
|
|
|
+ await InvokeHandler(instance.BizType,
|
|
|
|
|
+ h => h.OnNodeCompleted(instance.BizId, task.NodeId, task.NodeName ?? ""));
|
|
|
|
|
+ var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
|
|
|
|
|
+ await AdvanceToNext(instance, flowData, task.NodeId);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private async Task AutoRejectTask(ApprovalFlowTask task, ApprovalFlowInstance instance)
|
|
|
|
|
+ {
|
|
|
|
|
+ task.Status = FlowTaskStatusEnum.Rejected;
|
|
|
|
|
+ task.Comment = "系统自动拒绝(超时)";
|
|
|
|
|
+ task.ActionTime = DateTime.Now;
|
|
|
|
|
+ await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
|
|
|
|
|
+
|
|
|
|
|
+ instance.Status = FlowInstanceStatusEnum.Rejected;
|
|
|
|
|
+ instance.EndTime = DateTime.Now;
|
|
|
|
|
+ await _instanceRep.AsUpdateable(instance).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ await WriteSystemLog(instance.Id, task.Id, task.NodeId,
|
|
|
|
|
+ FlowLogActionEnum.AutoTimeout, "审批超时,系统自动拒绝");
|
|
|
|
|
+
|
|
|
|
|
+ await InvokeHandler(instance.BizType,
|
|
|
|
|
+ h => h.OnFlowCompleted(instance.BizId, FlowInstanceStatusEnum.Rejected));
|
|
|
|
|
+ await _notifyService.NotifyFlowCompleted(instance.InitiatorId, instance.Id,
|
|
|
|
|
+ instance.Title, FlowInstanceStatusEnum.Rejected);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private async Task AutoEscalateTask(ApprovalFlowTask task, FlowProperties nodeProps, ApprovalFlowInstance instance)
|
|
|
|
|
+ {
|
|
|
|
|
+ if (string.IsNullOrWhiteSpace(nodeProps.EscalationApproverType)
|
|
|
|
|
+ || string.IsNullOrWhiteSpace(nodeProps.EscalationApproverIds))
|
|
|
|
|
+ return;
|
|
|
|
|
+
|
|
|
|
|
+ task.Status = FlowTaskStatusEnum.Escalated;
|
|
|
|
|
+ task.Comment = "系统自动升级(超时)";
|
|
|
|
|
+ task.ActionTime = DateTime.Now;
|
|
|
|
|
+ await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
|
|
|
|
|
+
|
|
|
|
|
+ var approvers = await ResolveApprovers(
|
|
|
|
|
+ new FlowProperties
|
|
|
|
|
+ {
|
|
|
|
|
+ ApproverType = nodeProps.EscalationApproverType,
|
|
|
|
|
+ ApproverIds = nodeProps.EscalationApproverIds,
|
|
|
|
|
+ },
|
|
|
|
|
+ instance.InitiatorId);
|
|
|
|
|
+
|
|
|
|
|
+ if (approvers.Count == 0) return;
|
|
|
|
|
+
|
|
|
|
|
+ var newTasks = approvers.Select(a => new ApprovalFlowTask
|
|
|
|
|
+ {
|
|
|
|
|
+ InstanceId = instance.Id,
|
|
|
|
|
+ NodeId = task.NodeId,
|
|
|
|
|
+ NodeName = task.NodeName,
|
|
|
|
|
+ AssigneeId = a.userId,
|
|
|
|
|
+ AssigneeName = a.userName,
|
|
|
|
|
+ Status = FlowTaskStatusEnum.Pending,
|
|
|
|
|
+ OrgId = instance.OrgId,
|
|
|
|
|
+ }).ToList();
|
|
|
|
|
+ await _taskRep.AsInsertable(newTasks).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ var targetNames = string.Join(", ", approvers.Select(a => a.userName));
|
|
|
|
|
+ await WriteSystemLog(instance.Id, task.Id, task.NodeId,
|
|
|
|
|
+ FlowLogActionEnum.AutoTimeout, $"审批超时,自动升级给 {targetNames}");
|
|
|
|
|
+
|
|
|
|
|
+ var targetUserIds = approvers.Select(a => a.userId).Distinct().ToList();
|
|
|
|
|
+ await _notifyService.NotifyEscalated(targetUserIds, instance.Id, instance.Title,
|
|
|
|
|
+ "系统", task.NodeName);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ private async Task WriteSystemLog(long instanceId, long? taskId, string? nodeId,
|
|
|
|
|
+ FlowLogActionEnum action, string? comment)
|
|
|
|
|
+ {
|
|
|
|
|
+ await _logRep.InsertAsync(new ApprovalFlowLog
|
|
|
|
|
+ {
|
|
|
|
|
+ InstanceId = instanceId,
|
|
|
|
|
+ TaskId = taskId,
|
|
|
|
|
+ NodeId = nodeId,
|
|
|
|
|
+ Action = action,
|
|
|
|
|
+ OperatorId = 0,
|
|
|
|
|
+ OperatorName = "系统",
|
|
|
|
|
+ Comment = comment,
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// ═══════════════════════════════════════════
|
|
// ═══════════════════════════════════════════
|
|
|
// 内部引擎方法
|
|
// 内部引擎方法
|
|
|
// ═══════════════════════════════════════════
|
|
// ═══════════════════════════════════════════
|
|
|
|
|
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 推进到下一节点。支持并行网关(Fork / Join)。
|
|
|
|
|
+ /// 行为约定:
|
|
|
|
|
+ /// - 进入本方法前,调用方(如 <see cref="Approve"/>)应已将 <paramref name="currentNodeId"/>(userTask)标记为完成节点;
|
|
|
|
|
+ /// 本方法会将途经的网关节点也写入 <see cref="ApprovalFlowCompletedNode"/>。
|
|
|
|
|
+ /// - 并行网关 Fork(出边>=2):沿每条出边递归推进;
|
|
|
|
|
+ /// 并行网关 Join(入边>=2):校验所有前驱节点是否都在"已完成"集合中,
|
|
|
|
|
+ /// 任一尚未完成则**静默等待**(不报错、不推进),由后续分支完成后再次触发 Join 校验。
|
|
|
|
|
+ /// </summary>
|
|
|
private async Task AdvanceToNext(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string currentNodeId)
|
|
private async Task AdvanceToNext(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string currentNodeId)
|
|
|
{
|
|
{
|
|
|
- var nextNodeId = FindNextNodeId(flowData, currentNodeId);
|
|
|
|
|
- if (nextNodeId == null)
|
|
|
|
|
|
|
+ // 当前节点可能是 userTask(完成记录由 Approve 写入)或网关(入口处已写入);此处统一确保幂等入库
|
|
|
|
|
+ var currentNode = flowData.Nodes.FirstOrDefault(n => n.Id == currentNodeId);
|
|
|
|
|
+ await MarkNodeCompleted(instance.Id, currentNode);
|
|
|
|
|
+
|
|
|
|
|
+ var outgoingEdges = flowData.Edges.Where(e => e.SourceNodeId == currentNodeId).ToList();
|
|
|
|
|
+ if (outgoingEdges.Count == 0)
|
|
|
{
|
|
{
|
|
|
await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
|
|
await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ // 非并行网关场景:当前节点一般只有 1 条出边
|
|
|
|
|
+ foreach (var edge in outgoingEdges)
|
|
|
|
|
+ {
|
|
|
|
|
+ await ProcessNextNode(instance, flowData, edge.TargetNodeId);
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 处理某个"下一节点"。根据节点类型分发:
|
|
|
|
|
+ /// - endEvent:所有分支任务都结束时触发实例完成
|
|
|
|
|
+ /// - exclusiveGateway:按条件选择分支
|
|
|
|
|
+ /// - parallelGateway:Fork 并行分发;Join 等待所有前驱完成
|
|
|
|
|
+ /// - userTask / 其他:创建任务
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task ProcessNextNode(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string nextNodeId)
|
|
|
|
|
+ {
|
|
|
var nextNode = flowData.Nodes.FirstOrDefault(n => n.Id == nextNodeId);
|
|
var nextNode = flowData.Nodes.FirstOrDefault(n => n.Id == nextNodeId);
|
|
|
if (nextNode == null)
|
|
if (nextNode == null)
|
|
|
{
|
|
{
|
|
@@ -343,27 +589,98 @@ public class FlowEngineService : ITransient
|
|
|
|
|
|
|
|
if (nextNode.Type is "bpmn:endEvent" or "end-node")
|
|
if (nextNode.Type is "bpmn:endEvent" or "end-node")
|
|
|
{
|
|
{
|
|
|
|
|
+ // 所有并行分支都已结束(无其它 Pending 任务)才真正完成实例
|
|
|
|
|
+ var hasOtherPending = await _taskRep.AsQueryable()
|
|
|
|
|
+ .AnyAsync(t => t.InstanceId == instance.Id && t.Status == FlowTaskStatusEnum.Pending);
|
|
|
|
|
+ if (hasOtherPending) return;
|
|
|
|
|
+
|
|
|
await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
|
|
await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
if (nextNode.Type is "bpmn:exclusiveGateway")
|
|
if (nextNode.Type is "bpmn:exclusiveGateway")
|
|
|
{
|
|
{
|
|
|
|
|
+ await MarkNodeCompleted(instance.Id, nextNode);
|
|
|
var bizData = await GetBizData(instance.BizType, instance.BizId);
|
|
var bizData = await GetBizData(instance.BizType, instance.BizId);
|
|
|
var targetNodeId = EvaluateGateway(nextNode.Properties?.Conditions, flowData, nextNode.Id, bizData);
|
|
var targetNodeId = EvaluateGateway(nextNode.Properties?.Conditions, flowData, nextNode.Id, bizData);
|
|
|
instance.CurrentNodeId = targetNodeId;
|
|
instance.CurrentNodeId = targetNodeId;
|
|
|
await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
|
|
await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
|
|
|
- await AdvanceToNext(instance, flowData, nextNode.Id);
|
|
|
|
|
|
|
+ await ProcessNextNode(instance, flowData, targetNodeId);
|
|
|
return;
|
|
return;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ if (nextNode.Type is "bpmn:parallelGateway")
|
|
|
|
|
+ {
|
|
|
|
|
+ var incoming = flowData.Edges.Where(e => e.TargetNodeId == nextNode.Id).Select(e => e.SourceNodeId).ToList();
|
|
|
|
|
+ var outgoing = flowData.Edges.Where(e => e.SourceNodeId == nextNode.Id).Select(e => e.TargetNodeId).ToList();
|
|
|
|
|
+
|
|
|
|
|
+ // Join 语义:入边 >= 2,需等所有前驱都已完成
|
|
|
|
|
+ if (incoming.Count >= 2)
|
|
|
|
|
+ {
|
|
|
|
|
+ var completedSet = await GetCompletedNodeIdSet(instance.Id);
|
|
|
|
|
+ if (!incoming.All(p => completedSet.Contains(p)))
|
|
|
|
|
+ {
|
|
|
|
|
+ // 未汇合,静默等待后续分支抵达
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // Fork 或 Join 通过:标记网关完成,沿所有出边推进
|
|
|
|
|
+ await MarkNodeCompleted(instance.Id, nextNode);
|
|
|
|
|
+ foreach (var target in outgoing)
|
|
|
|
|
+ {
|
|
|
|
|
+ await ProcessNextNode(instance, flowData, target);
|
|
|
|
|
+ }
|
|
|
|
|
+ return;
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ // userTask 或其他:创建任务
|
|
|
instance.CurrentNodeId = nextNodeId;
|
|
instance.CurrentNodeId = nextNodeId;
|
|
|
await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
|
|
await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
|
|
|
await CreateTasksForNode(instance, flowData, nextNodeId);
|
|
await CreateTasksForNode(instance, flowData, nextNodeId);
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 标记节点已完成(幂等:重复写入被唯一索引拦截后忽略)
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task MarkNodeCompleted(long instanceId, ApprovalFlowNodeItem? node)
|
|
|
|
|
+ {
|
|
|
|
|
+ if (node == null) return;
|
|
|
|
|
+ try
|
|
|
|
|
+ {
|
|
|
|
|
+ await _completedNodeRep.InsertAsync(new ApprovalFlowCompletedNode
|
|
|
|
|
+ {
|
|
|
|
|
+ InstanceId = instanceId,
|
|
|
|
|
+ NodeId = node.Id,
|
|
|
|
|
+ NodeName = node.Properties?.NodeName ?? node.Text?.Value,
|
|
|
|
|
+ NodeType = node.Type,
|
|
|
|
|
+ CompletedTime = DateTime.Now,
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+ catch
|
|
|
|
|
+ {
|
|
|
|
|
+ // 并发场景下可能触发唯一索引冲突,忽略(已存在即可)
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 查询实例已完成节点 Id 集合
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task<HashSet<string>> GetCompletedNodeIdSet(long instanceId)
|
|
|
|
|
+ {
|
|
|
|
|
+ var ids = await _completedNodeRep.AsQueryable()
|
|
|
|
|
+ .Where(c => c.InstanceId == instanceId)
|
|
|
|
|
+ .Select(c => c.NodeId)
|
|
|
|
|
+ .ToListAsync();
|
|
|
|
|
+ return new HashSet<string>(ids);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
private async Task CompleteInstance(ApprovalFlowInstance instance, FlowInstanceStatusEnum status)
|
|
private async Task CompleteInstance(ApprovalFlowInstance instance, FlowInstanceStatusEnum status)
|
|
|
{
|
|
{
|
|
|
|
|
+ // 幂等:并行分支同时到达 end 时避免重复完成
|
|
|
|
|
+ var latest = await _instanceRep.GetByIdAsync(instance.Id);
|
|
|
|
|
+ if (latest == null || latest.Status != FlowInstanceStatusEnum.Running) return;
|
|
|
|
|
+
|
|
|
instance.Status = status;
|
|
instance.Status = status;
|
|
|
instance.EndTime = DateTime.Now;
|
|
instance.EndTime = DateTime.Now;
|
|
|
await _instanceRep.AsUpdateable(instance)
|
|
await _instanceRep.AsUpdateable(instance)
|
|
@@ -383,21 +700,101 @@ public class FlowEngineService : ITransient
|
|
|
if (approvers.Count == 0)
|
|
if (approvers.Count == 0)
|
|
|
throw Oops.Oh($"节点 [{node.Properties?.NodeName ?? nodeId}] 未配置审批人或审批人列表为空");
|
|
throw Oops.Oh($"节点 [{node.Properties?.NodeName ?? nodeId}] 未配置审批人或审批人列表为空");
|
|
|
|
|
|
|
|
|
|
+ var nodeName = node.Properties?.NodeName ?? node.Text?.Value;
|
|
|
var tasks = approvers.Select(a => new ApprovalFlowTask
|
|
var tasks = approvers.Select(a => new ApprovalFlowTask
|
|
|
{
|
|
{
|
|
|
InstanceId = instance.Id,
|
|
InstanceId = instance.Id,
|
|
|
NodeId = nodeId,
|
|
NodeId = nodeId,
|
|
|
- NodeName = node.Properties?.NodeName ?? node.Text?.Value,
|
|
|
|
|
|
|
+ NodeName = nodeName,
|
|
|
AssigneeId = a.userId,
|
|
AssigneeId = a.userId,
|
|
|
AssigneeName = a.userName,
|
|
AssigneeName = a.userName,
|
|
|
Status = FlowTaskStatusEnum.Pending,
|
|
Status = FlowTaskStatusEnum.Pending,
|
|
|
OrgId = instance.OrgId,
|
|
OrgId = instance.OrgId,
|
|
|
}).ToList();
|
|
}).ToList();
|
|
|
|
|
|
|
|
- await _taskRep.AsInsertable(tasks).ExecuteCommandAsync();
|
|
|
|
|
|
|
+ // P1-6 审批代理:为每个原审批人检查是否存在有效代理,是则并行创建一条代理任务
|
|
|
|
|
+ var delegateTasks = new List<ApprovalFlowTask>();
|
|
|
|
|
+ foreach (var a in approvers)
|
|
|
|
|
+ {
|
|
|
|
|
+ var del = await FindEffectiveDelegate(a.userId, instance.BizType);
|
|
|
|
|
+ if (del == null) continue;
|
|
|
|
|
+ // 代理人和原审批人不能重复,代理人也不能是原审批人列表中其他人(避免同一人两条任务)
|
|
|
|
|
+ if (approvers.Any(x => x.userId == del.DelegateUserId)) continue;
|
|
|
|
|
+
|
|
|
|
|
+ delegateTasks.Add(new ApprovalFlowTask
|
|
|
|
|
+ {
|
|
|
|
|
+ InstanceId = instance.Id,
|
|
|
|
|
+ NodeId = nodeId,
|
|
|
|
|
+ NodeName = nodeName,
|
|
|
|
|
+ AssigneeId = del.DelegateUserId,
|
|
|
|
|
+ AssigneeName = del.DelegateUserName,
|
|
|
|
|
+ Status = FlowTaskStatusEnum.Pending,
|
|
|
|
|
+ OrgId = instance.OrgId,
|
|
|
|
|
+ IsDelegate = true,
|
|
|
|
|
+ DelegateForUserId = a.userId,
|
|
|
|
|
+ DelegateForUserName = a.userName,
|
|
|
|
|
+ });
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ var allTasks = tasks.Concat(delegateTasks).ToList();
|
|
|
|
|
+ await _taskRep.AsInsertable(allTasks).ExecuteCommandAsync();
|
|
|
|
|
+
|
|
|
|
|
+ var assigneeIds = allTasks.Select(t => t.AssigneeId).Distinct().ToList();
|
|
|
|
|
+ await _notifyService.NotifyNewTask(assigneeIds, instance.Id, instance.Title, nodeName);
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 查找指定用户当前生效的审批代理(时间窗口内 + 已启用 + BizType 匹配或全局)
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task<ApprovalFlowDelegate?> FindEffectiveDelegate(long userId, string? bizType)
|
|
|
|
|
+ {
|
|
|
|
|
+ var now = DateTime.Now;
|
|
|
|
|
+ return await _delegateRep.AsQueryable()
|
|
|
|
|
+ .Where(d => d.UserId == userId
|
|
|
|
|
+ && d.IsEnabled
|
|
|
|
|
+ && d.StartTime <= now
|
|
|
|
|
+ && d.EndTime >= now
|
|
|
|
|
+ && (string.IsNullOrEmpty(d.BizType) || d.BizType == bizType))
|
|
|
|
|
+ .OrderBy(d => d.BizType == null ? 1 : 0) // 优先匹配指定 BizType 的代理
|
|
|
|
|
+ .OrderByDescending(d => d.CreateTime)
|
|
|
|
|
+ .FirstAsync();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 取消与已完成任务配对的代理任务(P1-6)
|
|
|
|
|
+ /// - 本人任务完成:取消对应的代理任务
|
|
|
|
|
+ /// - 代理任务完成:取消对应的本人任务
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task CancelPairedDelegateTask(ApprovalFlowTask completedTask)
|
|
|
|
|
+ {
|
|
|
|
|
+ ApprovalFlowTask? paired;
|
|
|
|
|
+ if (completedTask.IsDelegate)
|
|
|
|
|
+ {
|
|
|
|
|
+ var originalUserId = completedTask.DelegateForUserId ?? 0;
|
|
|
|
|
+ if (originalUserId == 0) return;
|
|
|
|
|
+ paired = await _taskRep.AsQueryable()
|
|
|
|
|
+ .Where(t => t.InstanceId == completedTask.InstanceId
|
|
|
|
|
+ && t.NodeId == completedTask.NodeId
|
|
|
|
|
+ && t.Status == FlowTaskStatusEnum.Pending
|
|
|
|
|
+ && t.AssigneeId == originalUserId
|
|
|
|
|
+ && !t.IsDelegate)
|
|
|
|
|
+ .FirstAsync();
|
|
|
|
|
+ }
|
|
|
|
|
+ else
|
|
|
|
|
+ {
|
|
|
|
|
+ paired = await _taskRep.AsQueryable()
|
|
|
|
|
+ .Where(t => t.InstanceId == completedTask.InstanceId
|
|
|
|
|
+ && t.NodeId == completedTask.NodeId
|
|
|
|
|
+ && t.Status == FlowTaskStatusEnum.Pending
|
|
|
|
|
+ && t.IsDelegate
|
|
|
|
|
+ && t.DelegateForUserId == completedTask.AssigneeId)
|
|
|
|
|
+ .FirstAsync();
|
|
|
|
|
+ }
|
|
|
|
|
|
|
|
- var assigneeIds = tasks.Select(t => t.AssigneeId).Distinct().ToList();
|
|
|
|
|
- await _notifyService.NotifyNewTask(assigneeIds, instance.Id, instance.Title, node.Properties?.NodeName ?? node.Text?.Value);
|
|
|
|
|
|
|
+ if (paired == null) return;
|
|
|
|
|
+ paired.Status = FlowTaskStatusEnum.Cancelled;
|
|
|
|
|
+ paired.ActionTime = DateTime.Now;
|
|
|
|
|
+ await _taskRep.AsUpdateable(paired).ExecuteCommandAsync();
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
private async Task<List<(long userId, string userName)>> ResolveApprovers(FlowProperties? props, long initiatorId)
|
|
private async Task<List<(long userId, string userName)>> ResolveApprovers(FlowProperties? props, long initiatorId)
|
|
@@ -446,6 +843,30 @@ public class FlowEngineService : ITransient
|
|
|
return users.Select(u => (u.Id, u.RealName ?? "")).ToList();
|
|
return users.Select(u => (u.Id, u.RealName ?? "")).ToList();
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ if (approverType == nameof(ApproverTypeEnum.DepartmentLeader))
|
|
|
|
|
+ {
|
|
|
|
|
+ var initiator = await _userRep.GetByIdAsync(initiatorId);
|
|
|
|
|
+ if (initiator == null || initiator.OrgId <= 0)
|
|
|
|
|
+ return new List<(long, string)>();
|
|
|
|
|
+
|
|
|
|
|
+ var org = await _orgRep.GetByIdAsync(initiator.OrgId);
|
|
|
|
|
+ if (org?.DirectorId != null && org.DirectorId > 0)
|
|
|
|
|
+ {
|
|
|
|
|
+ var director = await _userRep.GetByIdAsync(org.DirectorId.Value);
|
|
|
|
|
+ if (director != null)
|
|
|
|
|
+ return new List<(long, string)> { (director.Id, director.RealName ?? "") };
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ if (initiator.ManagerUserId != null && initiator.ManagerUserId > 0)
|
|
|
|
|
+ {
|
|
|
|
|
+ var manager = await _userRep.GetByIdAsync(initiator.ManagerUserId.Value);
|
|
|
|
|
+ if (manager != null)
|
|
|
|
|
+ return new List<(long, string)> { (manager.Id, manager.RealName ?? "") };
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ return new List<(long, string)>();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
return new List<(long, string)>();
|
|
return new List<(long, string)>();
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -596,6 +1017,24 @@ public class FlowEngineService : ITransient
|
|
|
await _taskRep.AsUpdateable(tasks).ExecuteCommandAsync();
|
|
await _taskRep.AsUpdateable(tasks).ExecuteCommandAsync();
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ /// <summary>
|
|
|
|
|
+ /// 取消整个实例下所有剩余 Pending 任务(Reject 时跨并行分支使用)
|
|
|
|
|
+ /// </summary>
|
|
|
|
|
+ private async Task CancelAllPendingTasks(long instanceId, long? excludeTaskId = null)
|
|
|
|
|
+ {
|
|
|
|
|
+ var tasks = await _taskRep.AsQueryable()
|
|
|
|
|
+ .Where(t => t.InstanceId == instanceId && t.Status == FlowTaskStatusEnum.Pending)
|
|
|
|
|
+ .WhereIF(excludeTaskId.HasValue, t => t.Id != excludeTaskId!.Value)
|
|
|
|
|
+ .ToListAsync();
|
|
|
|
|
+ foreach (var t in tasks)
|
|
|
|
|
+ {
|
|
|
|
|
+ t.Status = FlowTaskStatusEnum.Cancelled;
|
|
|
|
|
+ t.ActionTime = DateTime.Now;
|
|
|
|
|
+ }
|
|
|
|
|
+ if (tasks.Count > 0)
|
|
|
|
|
+ await _taskRep.AsUpdateable(tasks).ExecuteCommandAsync();
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
private async Task<ApprovalFlowTask> GetPendingTask(long taskId)
|
|
private async Task<ApprovalFlowTask> GetPendingTask(long taskId)
|
|
|
{
|
|
{
|
|
|
var task = await _taskRep.GetByIdAsync(taskId)
|
|
var task = await _taskRep.GetByIdAsync(taskId)
|