FlowEngineService.cs 53 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234
  1. using System.Text.Json;
  2. using Admin.NET.Core.Service;
  3. using Microsoft.Extensions.Logging;
  4. namespace Admin.NET.Plugin.ApprovalFlow.Service;
  5. /// <summary>
  6. /// 流程推进引擎 — 核心状态机
  7. /// 不暴露为 API,由其他 Service 内部调用
  8. /// </summary>
  9. public class FlowEngineService : ITransient
  10. {
  11. private readonly SqlSugarRepository<ApprovalFlow> _flowRep;
  12. private readonly SqlSugarRepository<ApprovalFlowInstance> _instanceRep;
  13. private readonly SqlSugarRepository<ApprovalFlowTask> _taskRep;
  14. private readonly SqlSugarRepository<ApprovalFlowLog> _logRep;
  15. private readonly SqlSugarRepository<SysUserRole> _userRoleRep;
  16. private readonly SqlSugarRepository<SysUser> _userRep;
  17. private readonly SqlSugarRepository<SysOrg> _orgRep;
  18. private readonly SqlSugarRepository<ApprovalFlowDelegate> _delegateRep;
  19. private readonly SqlSugarRepository<ApprovalFlowCompletedNode> _completedNodeRep;
  20. private readonly UserManager _userManager;
  21. private readonly SysOrgService _sysOrgService;
  22. private readonly FlowNotifyService _notifyService;
  23. private readonly ILogger<FlowEngineService> _logger;
  24. public FlowEngineService(
  25. SqlSugarRepository<ApprovalFlow> flowRep,
  26. SqlSugarRepository<ApprovalFlowInstance> instanceRep,
  27. SqlSugarRepository<ApprovalFlowTask> taskRep,
  28. SqlSugarRepository<ApprovalFlowLog> logRep,
  29. SqlSugarRepository<SysUserRole> userRoleRep,
  30. SqlSugarRepository<SysUser> userRep,
  31. SqlSugarRepository<SysOrg> orgRep,
  32. SqlSugarRepository<ApprovalFlowDelegate> delegateRep,
  33. SqlSugarRepository<ApprovalFlowCompletedNode> completedNodeRep,
  34. UserManager userManager,
  35. SysOrgService sysOrgService,
  36. FlowNotifyService notifyService,
  37. ILogger<FlowEngineService> logger)
  38. {
  39. _flowRep = flowRep;
  40. _instanceRep = instanceRep;
  41. _taskRep = taskRep;
  42. _logRep = logRep;
  43. _userRoleRep = userRoleRep;
  44. _userRep = userRep;
  45. _orgRep = orgRep;
  46. _delegateRep = delegateRep;
  47. _completedNodeRep = completedNodeRep;
  48. _userManager = userManager;
  49. _sysOrgService = sysOrgService;
  50. _notifyService = notifyService;
  51. _logger = logger;
  52. }
  53. /// <summary>
  54. /// 机构数据权限依赖实体的 OrgId;插入时必须与当前用户机构一致,否则非超管在「我发起的/待办」中查不到(OrgId=0 会被过滤)。
  55. /// </summary>
  56. private async Task<long> ResolveOrgIdForNewFlowEntityAsync()
  57. {
  58. var oid = _userManager.OrgId;
  59. if (oid > 0) return oid;
  60. var list = await _sysOrgService.GetUserOrgIdList();
  61. return list is { Count: > 0 } ? list[0] : 0;
  62. }
  63. // ═══════════════════════════════════════════
  64. // 核心生命周期
  65. // ═══════════════════════════════════════════
  66. /// <summary>
  67. /// 发起流程
  68. /// </summary>
  69. public async Task<long> StartFlow(StartFlowInput input)
  70. {
  71. // S8-S1-EXCEPTION-FLOW-SYNC-FIX-1:流程定义是「全局配置」,不应受 SqlSugarFilter 的数据范围(DataScope)隔离。
  72. // ApprovalFlow 继承 EntityBaseOrgDel,会被数据范围过滤命中:
  73. // - 角色 DataScope=Self(仅本人) → 加「CreateUserId == 当前用户」过滤(按 entityType 注册,非 IOrgIdFilter);
  74. // - 角色 DataScope=Dept/DeptChild → 加 IOrgIdFilter「OrgId ∈ 用户机构集」过滤。
  75. // 既有定义(如 TB001/EXCEPTION_REPORT)的 CreateUserId 为流程创建者(非发起人)、OrgId=0,非超管用户在任一受限数据范围下都会被过滤掉
  76. // → 此查询返回 null → 抛「未找到已发布流程定义」→ StartFlow 失败 → 上游静默吞、建单不起流。
  77. // ClearFilter() 清除该查询全部全局过滤(Org/Self 数据范围 + 软删);软删由 WHERE 显式 !IsDelete 补回;
  78. // ApprovalFlow 无租户过滤(EntityBase 未实现 ITenantIdFilter),不存在绕过租户隔离风险。
  79. var tenantId = _userManager.TenantId;
  80. var flow = await _flowRep.AsQueryable()
  81. .ClearFilter()
  82. .Where(u => u.BizType == input.BizType && u.IsPublished && !u.IsDelete)
  83. .OrderBy(u => new dynamic[] {
  84. u.TenantId == tenantId ? 1 : 0,
  85. u.Version
  86. }, OrderByType.Desc)
  87. .FirstAsync() ?? throw Oops.Oh($"未找到业务类型 [{input.BizType}] 的已发布流程定义");
  88. if (string.IsNullOrWhiteSpace(flow.FlowJson))
  89. throw Oops.Oh("流程定义的 FlowJson 为空,请先设计流程图");
  90. var flowData = JsonSerializer.Deserialize<ApprovalFlowItem>(flow.FlowJson)
  91. ?? throw Oops.Oh("FlowJson 反序列化失败");
  92. var orgId = await ResolveOrgIdForNewFlowEntityAsync();
  93. var instance = new ApprovalFlowInstance
  94. {
  95. FlowId = flow.Id,
  96. FlowVersion = flow.Version,
  97. BizType = input.BizType,
  98. BizId = input.BizId,
  99. BizNo = input.BizNo,
  100. Title = input.Title ?? $"{flow.Name}-{input.BizNo}",
  101. InitiatorId = _userManager.UserId,
  102. InitiatorName = _userManager.RealName,
  103. Status = FlowInstanceStatusEnum.Running,
  104. FlowJsonSnapshot = flow.FlowJson,
  105. StartTime = DateTime.Now,
  106. OrgId = orgId,
  107. };
  108. var startNode = flowData.Nodes.FirstOrDefault(n =>
  109. n.Type is "bpmn:startEvent" or "start-node")
  110. ?? throw Oops.Oh("流程图中未找到开始节点");
  111. await _instanceRep.InsertAsync(instance);
  112. await WriteLog(instance.Id, null, startNode.Id, FlowLogActionEnum.Submit, input.Comment);
  113. // 开始节点登记为已完成,随后沿出边推进(兼容首节点即为并行网关 Fork 的拓扑)
  114. await MarkNodeCompleted(instance.Id, startNode);
  115. var firstOutgoing = flowData.Edges.Where(e => e.SourceNodeId == startNode.Id).Select(e => e.TargetNodeId).ToList();
  116. if (firstOutgoing.Count == 0)
  117. throw Oops.Oh("流程图开始节点未连接任何后继节点");
  118. foreach (var target in firstOutgoing)
  119. {
  120. await ProcessNextNode(instance, flowData, target);
  121. }
  122. await InvokeHandler(input.BizType, instance.Id, h => h.OnFlowStarted(input.BizId, instance.Id));
  123. return instance.Id;
  124. }
  125. /// <summary>
  126. /// 同意
  127. /// </summary>
  128. public async Task Approve(long taskId, string? comment)
  129. {
  130. var task = await GetPendingTask(taskId);
  131. task.Status = FlowTaskStatusEnum.Approved;
  132. task.Comment = comment;
  133. task.ActionTime = DateTime.Now;
  134. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  135. // P1-6 审批代理:同步取消配对的本人/代理任务
  136. await CancelPairedDelegateTask(task);
  137. await WriteLog(task.InstanceId, taskId, task.NodeId, FlowLogActionEnum.Approve, comment);
  138. var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
  139. ?? throw Oops.Oh("流程实例不存在");
  140. if (await IsNodeCompleted(instance, task.NodeId))
  141. {
  142. await InvokeHandler(instance.BizType, instance.Id,
  143. h => h.OnNodeCompleted(instance.BizId, instance.Id, task.NodeId, task.NodeName ?? "", task.AssigneeId));
  144. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  145. await AdvanceToNext(instance, flowData, task.NodeId);
  146. }
  147. }
  148. /// <summary>
  149. /// 拒绝
  150. /// </summary>
  151. public async Task Reject(long taskId, string? comment)
  152. {
  153. var task = await GetPendingTask(taskId);
  154. task.Status = FlowTaskStatusEnum.Rejected;
  155. task.Comment = comment;
  156. task.ActionTime = DateTime.Now;
  157. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  158. // Reject 导致流程终止:取消整个实例所有剩余 Pending 任务(含并行分支、配对代理任务等)
  159. await CancelAllPendingTasks(task.InstanceId, task.Id);
  160. var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
  161. ?? throw Oops.Oh("流程实例不存在");
  162. instance.Status = FlowInstanceStatusEnum.Rejected;
  163. instance.EndTime = DateTime.Now;
  164. await _instanceRep.AsUpdateable(instance).ExecuteCommandAsync();
  165. await WriteLog(instance.Id, taskId, task.NodeId, FlowLogActionEnum.Reject, comment);
  166. // P4-17: 驳回人 = 当前被指派的审批人
  167. await InvokeHandler(instance.BizType, instance.Id,
  168. h => h.OnFlowCompleted(instance.BizId, instance.Id, FlowInstanceStatusEnum.Rejected, task.AssigneeId));
  169. await _notifyService.NotifyFlowCompleted(instance.InitiatorId, instance.Id, instance.Title, FlowInstanceStatusEnum.Rejected, instance.BizType);
  170. }
  171. // ═══════════════════════════════════════════
  172. // 扩展操作
  173. // ═══════════════════════════════════════════
  174. /// <summary>
  175. /// 转办
  176. /// </summary>
  177. public async Task Transfer(long taskId, long targetUserId, string? comment)
  178. {
  179. var task = await GetPendingTask(taskId);
  180. task.Status = FlowTaskStatusEnum.Transferred;
  181. task.Comment = comment;
  182. task.ActionTime = DateTime.Now;
  183. task.TransferToId = targetUserId;
  184. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  185. var targetUser = await _userRep.GetByIdAsync(targetUserId);
  186. var instForOrg = await _instanceRep.GetByIdAsync(task.InstanceId);
  187. var newTask = new ApprovalFlowTask
  188. {
  189. InstanceId = task.InstanceId,
  190. NodeId = task.NodeId,
  191. NodeName = task.NodeName,
  192. AssigneeId = targetUserId,
  193. AssigneeName = targetUser?.RealName,
  194. Status = FlowTaskStatusEnum.Pending,
  195. OrgId = instForOrg?.OrgId ?? 0,
  196. };
  197. await _taskRep.InsertAsync(newTask);
  198. await WriteLog(task.InstanceId, taskId, task.NodeId, FlowLogActionEnum.Transfer,
  199. $"{comment} → 转办给 {targetUser?.RealName}");
  200. var instance = await _instanceRep.GetByIdAsync(task.InstanceId);
  201. await _notifyService.NotifyTransferred(targetUserId, task.InstanceId, instance?.Title ?? "", _userManager.RealName, instance?.BizType ?? "");
  202. }
  203. /// <summary>
  204. /// 撤回(发起人撤回)
  205. /// </summary>
  206. public async Task Withdraw(long instanceId)
  207. {
  208. var instance = await _instanceRep.GetByIdAsync(instanceId)
  209. ?? throw Oops.Oh("流程实例不存在");
  210. if (instance.InitiatorId != _userManager.UserId)
  211. throw Oops.Oh("只有发起人可以撤回");
  212. if (instance.Status != FlowInstanceStatusEnum.Running)
  213. throw Oops.Oh("当前流程状态不允许撤回");
  214. var pendingTasks = await _taskRep.AsQueryable()
  215. .Where(t => t.InstanceId == instanceId && t.Status == FlowTaskStatusEnum.Pending)
  216. .ToListAsync();
  217. var doneTasks = await _taskRep.AsQueryable()
  218. .Where(t => t.InstanceId == instanceId &&
  219. t.Status != FlowTaskStatusEnum.Pending &&
  220. t.Status != FlowTaskStatusEnum.Cancelled)
  221. .CountAsync();
  222. if (doneTasks > 0)
  223. throw Oops.Oh("已有人审批过,不可撤回");
  224. var cancelledUserIds = pendingTasks.Select(t => t.AssigneeId).Distinct().ToList();
  225. foreach (var t in pendingTasks)
  226. {
  227. t.Status = FlowTaskStatusEnum.Cancelled;
  228. t.ActionTime = DateTime.Now;
  229. }
  230. await _taskRep.AsUpdateable(pendingTasks).ExecuteCommandAsync();
  231. instance.Status = FlowInstanceStatusEnum.Cancelled;
  232. instance.EndTime = DateTime.Now;
  233. await _instanceRep.AsUpdateable(instance).ExecuteCommandAsync();
  234. await WriteLog(instanceId, null, instance.CurrentNodeId, FlowLogActionEnum.Withdraw, null);
  235. // P4-17: 撤回场景无"审批人",lastApproverId = null
  236. await InvokeHandler(instance.BizType, instanceId,
  237. h => h.OnFlowCompleted(instance.BizId, instanceId, FlowInstanceStatusEnum.Cancelled, null));
  238. await _notifyService.NotifyWithdrawn(cancelledUserIds, instanceId, instance.Title, instance.InitiatorName, instance.BizType);
  239. }
  240. /// <summary>
  241. /// 退回上一步
  242. /// </summary>
  243. public async Task ReturnToPrev(long taskId, string? comment)
  244. {
  245. var task = await GetPendingTask(taskId);
  246. var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
  247. ?? throw Oops.Oh("流程实例不存在");
  248. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  249. var prevNodeId = FindPrevUserTaskNodeId(flowData, task.NodeId);
  250. if (prevNodeId == null)
  251. throw Oops.Oh("已是第一个审批节点,无法退回");
  252. await CancelPendingTasks(task.InstanceId, task.NodeId);
  253. task.Status = FlowTaskStatusEnum.Returned;
  254. task.Comment = comment;
  255. task.ActionTime = DateTime.Now;
  256. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  257. instance.CurrentNodeId = prevNodeId;
  258. await _instanceRep.AsUpdateable(instance).ExecuteCommandAsync();
  259. await CreateTasksForNode(instance, flowData, prevNodeId);
  260. await WriteLog(instance.Id, taskId, task.NodeId, FlowLogActionEnum.Return, comment);
  261. var returnedTasks = await _taskRep.AsQueryable()
  262. .Where(t => t.InstanceId == instance.Id && t.NodeId == prevNodeId && t.Status == FlowTaskStatusEnum.Pending)
  263. .ToListAsync();
  264. var returnedUserIds = returnedTasks.Select(t => t.AssigneeId).Distinct().ToList();
  265. await _notifyService.NotifyReturned(returnedUserIds, instance.Id, instance.Title, _userManager.RealName, instance.BizType);
  266. }
  267. /// <summary>
  268. /// 加签
  269. /// </summary>
  270. public async Task AddSign(long taskId, long targetUserId, string? comment)
  271. {
  272. var task = await GetPendingTask(taskId);
  273. var targetUser = await _userRep.GetByIdAsync(targetUserId);
  274. var instForAddSign = await _instanceRep.GetByIdAsync(task.InstanceId);
  275. var newTask = new ApprovalFlowTask
  276. {
  277. InstanceId = task.InstanceId,
  278. NodeId = task.NodeId,
  279. NodeName = task.NodeName,
  280. AssigneeId = targetUserId,
  281. AssigneeName = targetUser?.RealName,
  282. Status = FlowTaskStatusEnum.Pending,
  283. IsAddSign = true,
  284. AddSignById = _userManager.UserId,
  285. OrgId = instForAddSign?.OrgId ?? 0,
  286. };
  287. await _taskRep.InsertAsync(newTask);
  288. await WriteLog(task.InstanceId, taskId, task.NodeId, FlowLogActionEnum.AddSign,
  289. $"{comment} → 加签给 {targetUser?.RealName}");
  290. var instance = await _instanceRep.GetByIdAsync(task.InstanceId);
  291. await _notifyService.NotifyAddSign(targetUserId, task.InstanceId, instance?.Title ?? "", _userManager.RealName, instance?.BizType ?? "");
  292. }
  293. /// <summary>
  294. /// 手动升级 — 当前审批人主动将任务升级到更高层级
  295. /// </summary>
  296. public async Task Escalate(long taskId, string? comment)
  297. {
  298. var task = await GetPendingTask(taskId);
  299. var instance = await _instanceRep.GetByIdAsync(task.InstanceId)
  300. ?? throw Oops.Oh("流程实例不存在");
  301. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  302. var node = flowData.Nodes.FirstOrDefault(n => n.Id == task.NodeId)
  303. ?? throw Oops.Oh("节点不存在");
  304. var props = node.Properties;
  305. if (props?.EnableManualEscalation != true
  306. || string.IsNullOrWhiteSpace(props.EscalationApproverType)
  307. || string.IsNullOrWhiteSpace(props.EscalationApproverIds))
  308. throw Oops.Oh("该节点未配置升级目标,无法升级");
  309. task.Status = FlowTaskStatusEnum.Escalated;
  310. task.Comment = comment;
  311. task.ActionTime = DateTime.Now;
  312. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  313. await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
  314. var escalationApprovers = await ResolveApprovers(
  315. new FlowProperties
  316. {
  317. ApproverType = props.EscalationApproverType,
  318. ApproverIds = props.EscalationApproverIds,
  319. ApproverNames = props.EscalationApproverNames,
  320. },
  321. instance.InitiatorId);
  322. if (escalationApprovers.Count == 0)
  323. throw Oops.Oh("升级目标审批人列表为空");
  324. var newTasks = escalationApprovers.Select(a => new ApprovalFlowTask
  325. {
  326. InstanceId = instance.Id,
  327. NodeId = task.NodeId,
  328. NodeName = task.NodeName,
  329. AssigneeId = a.userId,
  330. AssigneeName = a.userName,
  331. Status = FlowTaskStatusEnum.Pending,
  332. OrgId = instance.OrgId,
  333. }).ToList();
  334. await _taskRep.AsInsertable(newTasks).ExecuteCommandAsync();
  335. var targetNames = string.Join(", ", escalationApprovers.Select(a => a.userName));
  336. await WriteLog(instance.Id, taskId, task.NodeId, FlowLogActionEnum.Escalate,
  337. $"{comment} → 升级给 {targetNames}");
  338. var targetUserIds = escalationApprovers.Select(a => a.userId).Distinct().ToList();
  339. await _notifyService.NotifyEscalated(targetUserIds, instance.Id, instance.Title,
  340. _userManager.RealName, task.NodeName, instance.BizType);
  341. }
  342. /// <summary>
  343. /// 催办
  344. /// </summary>
  345. public async Task Urge(long instanceId)
  346. {
  347. var instance = await _instanceRep.GetByIdAsync(instanceId)
  348. ?? throw Oops.Oh("流程实例不存在");
  349. if (instance.Status != FlowInstanceStatusEnum.Running)
  350. throw Oops.Oh("当前流程不在审批中");
  351. await WriteLog(instanceId, null, instance.CurrentNodeId, FlowLogActionEnum.Urge, "催办");
  352. var pendingTasks = await _taskRep.AsQueryable()
  353. .Where(t => t.InstanceId == instanceId && t.Status == FlowTaskStatusEnum.Pending)
  354. .ToListAsync();
  355. var userIds = pendingTasks.Select(t => t.AssigneeId).Distinct().ToList();
  356. await _notifyService.NotifyUrge(userIds, instanceId, instance.Title, instance.BizType);
  357. }
  358. // ═══════════════════════════════════════════
  359. // 超时自动处理(由 FlowTimeoutJob 调用,无 UserManager 上下文)
  360. // ═══════════════════════════════════════════
  361. /// <summary>
  362. /// 扫描所有 Pending 且已超过 timeoutHours 的任务,按节点 timeoutAction 执行对应动作。
  363. /// 与 FlowTimeoutJob 扫描等价,抽出公共方法以便管理员接口立即触发(运维/E2E 测试用)。
  364. /// 返回本次处理的任务条数。
  365. /// </summary>
  366. public async Task<int> ScanTimeoutTasks(CancellationToken stoppingToken = default)
  367. {
  368. var pendingTasks = await _taskRep.Context.Queryable<ApprovalFlowTask>()
  369. .InnerJoin<ApprovalFlowInstance>((t, i) => t.InstanceId == i.Id)
  370. .Where((t, i) => t.Status == FlowTaskStatusEnum.Pending
  371. && i.Status == FlowInstanceStatusEnum.Running)
  372. .Select((t, i) => new
  373. {
  374. TaskId = t.Id,
  375. TaskCreatedAt = t.CreateTime,
  376. NodeId = t.NodeId,
  377. FlowJsonSnapshot = i.FlowJsonSnapshot,
  378. })
  379. .ToListAsync();
  380. var now = DateTime.Now;
  381. var processed = 0;
  382. foreach (var item in pendingTasks)
  383. {
  384. if (stoppingToken.IsCancellationRequested) break;
  385. if (string.IsNullOrWhiteSpace(item.FlowJsonSnapshot)) continue;
  386. ApprovalFlowItem? flowData;
  387. try { flowData = JsonSerializer.Deserialize<ApprovalFlowItem>(item.FlowJsonSnapshot); }
  388. catch { continue; }
  389. var node = flowData?.Nodes?.FirstOrDefault(n => n.Id == item.NodeId);
  390. var props = node?.Properties;
  391. if (props?.TimeoutHours == null || props.TimeoutHours <= 0) continue;
  392. if (string.IsNullOrWhiteSpace(props.TimeoutAction)) continue;
  393. var deadline = item.TaskCreatedAt.AddHours(props.TimeoutHours.Value);
  394. if (now < deadline) continue;
  395. // Notify 动作幂等:已记录过 AutoTimeout 日志的跳过,避免每轮重复催办
  396. if (props.TimeoutAction == "Notify")
  397. {
  398. var alreadyNotified = await _logRep.AsQueryable()
  399. .AnyAsync(log => log.TaskId == item.TaskId
  400. && log.Action == FlowLogActionEnum.AutoTimeout);
  401. if (alreadyNotified) continue;
  402. }
  403. try
  404. {
  405. await HandleTimeoutTask(item.TaskId);
  406. processed++;
  407. }
  408. catch
  409. {
  410. // 单任务异常隔离:吞掉继续下一个,由 Job/Admin API 层统一记日志
  411. }
  412. }
  413. return processed;
  414. }
  415. /// <summary>
  416. /// 处理单个超时任务(由定时任务调用)
  417. /// </summary>
  418. public async Task HandleTimeoutTask(long taskId)
  419. {
  420. var task = await _taskRep.GetByIdAsync(taskId);
  421. if (task == null || task.Status != FlowTaskStatusEnum.Pending) return;
  422. var instance = await _instanceRep.GetByIdAsync(task.InstanceId);
  423. if (instance == null || instance.Status != FlowInstanceStatusEnum.Running) return;
  424. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  425. var node = flowData.Nodes.FirstOrDefault(n => n.Id == task.NodeId);
  426. var props = node?.Properties;
  427. if (props == null) return;
  428. switch (props.TimeoutAction)
  429. {
  430. case "Notify":
  431. await _notifyService.NotifyTimeout(
  432. new List<long> { task.AssigneeId }, instance.Id, instance.Title, instance.BizType);
  433. await WriteSystemLog(instance.Id, task.Id, task.NodeId,
  434. FlowLogActionEnum.AutoTimeout, "审批超时,已发送提醒通知");
  435. break;
  436. case "AutoApprove":
  437. await AutoApproveTask(task, instance);
  438. break;
  439. case "AutoReject":
  440. await AutoRejectTask(task, instance);
  441. break;
  442. case "AutoEscalate":
  443. await AutoEscalateTask(task, props, instance);
  444. break;
  445. }
  446. }
  447. private async Task AutoApproveTask(ApprovalFlowTask task, ApprovalFlowInstance instance)
  448. {
  449. task.Status = FlowTaskStatusEnum.Approved;
  450. task.Comment = "系统自动通过(超时)";
  451. task.ActionTime = DateTime.Now;
  452. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  453. await WriteSystemLog(instance.Id, task.Id, task.NodeId,
  454. FlowLogActionEnum.AutoTimeout, "审批超时,系统自动通过");
  455. if (await IsNodeCompleted(instance, task.NodeId))
  456. {
  457. // P4-17: 系统自动通过,无人工审批人,approverUserId = null
  458. await InvokeHandler(instance.BizType, instance.Id,
  459. h => h.OnNodeCompleted(instance.BizId, instance.Id, task.NodeId, task.NodeName ?? "", null));
  460. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  461. await AdvanceToNext(instance, flowData, task.NodeId);
  462. }
  463. }
  464. private async Task AutoRejectTask(ApprovalFlowTask task, ApprovalFlowInstance instance)
  465. {
  466. task.Status = FlowTaskStatusEnum.Rejected;
  467. task.Comment = "系统自动拒绝(超时)";
  468. task.ActionTime = DateTime.Now;
  469. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  470. await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
  471. instance.Status = FlowInstanceStatusEnum.Rejected;
  472. instance.EndTime = DateTime.Now;
  473. await _instanceRep.AsUpdateable(instance).ExecuteCommandAsync();
  474. await WriteSystemLog(instance.Id, task.Id, task.NodeId,
  475. FlowLogActionEnum.AutoTimeout, "审批超时,系统自动拒绝");
  476. // P4-17: 系统自动拒绝,无人工审批人,lastApproverId = null
  477. await InvokeHandler(instance.BizType, instance.Id,
  478. h => h.OnFlowCompleted(instance.BizId, instance.Id, FlowInstanceStatusEnum.Rejected, null));
  479. await _notifyService.NotifyFlowCompleted(instance.InitiatorId, instance.Id,
  480. instance.Title, FlowInstanceStatusEnum.Rejected, instance.BizType);
  481. }
  482. private async Task AutoEscalateTask(ApprovalFlowTask task, FlowProperties nodeProps, ApprovalFlowInstance instance)
  483. {
  484. if (string.IsNullOrWhiteSpace(nodeProps.EscalationApproverType)
  485. || string.IsNullOrWhiteSpace(nodeProps.EscalationApproverIds))
  486. return;
  487. task.Status = FlowTaskStatusEnum.Escalated;
  488. task.Comment = "系统自动升级(超时)";
  489. task.ActionTime = DateTime.Now;
  490. await _taskRep.AsUpdateable(task).ExecuteCommandAsync();
  491. await CancelPendingTasks(task.InstanceId, task.NodeId, task.Id);
  492. var approvers = await ResolveApprovers(
  493. new FlowProperties
  494. {
  495. ApproverType = nodeProps.EscalationApproverType,
  496. ApproverIds = nodeProps.EscalationApproverIds,
  497. },
  498. instance.InitiatorId);
  499. if (approvers.Count == 0) return;
  500. var newTasks = approvers.Select(a => new ApprovalFlowTask
  501. {
  502. InstanceId = instance.Id,
  503. NodeId = task.NodeId,
  504. NodeName = task.NodeName,
  505. AssigneeId = a.userId,
  506. AssigneeName = a.userName,
  507. Status = FlowTaskStatusEnum.Pending,
  508. OrgId = instance.OrgId,
  509. }).ToList();
  510. await _taskRep.AsInsertable(newTasks).ExecuteCommandAsync();
  511. var targetNames = string.Join(", ", approvers.Select(a => a.userName));
  512. await WriteSystemLog(instance.Id, task.Id, task.NodeId,
  513. FlowLogActionEnum.AutoTimeout, $"审批超时,自动升级给 {targetNames}");
  514. var targetUserIds = approvers.Select(a => a.userId).Distinct().ToList();
  515. await _notifyService.NotifyEscalated(targetUserIds, instance.Id, instance.Title,
  516. "系统", task.NodeName, instance.BizType);
  517. }
  518. private async Task WriteSystemLog(long instanceId, long? taskId, string? nodeId,
  519. FlowLogActionEnum action, string? comment)
  520. {
  521. await _logRep.InsertAsync(new ApprovalFlowLog
  522. {
  523. InstanceId = instanceId,
  524. TaskId = taskId,
  525. NodeId = nodeId,
  526. Action = action,
  527. OperatorId = 0,
  528. OperatorName = "系统",
  529. Comment = comment,
  530. });
  531. }
  532. // ═══════════════════════════════════════════
  533. // 内部引擎方法
  534. // ═══════════════════════════════════════════
  535. /// <summary>
  536. /// 推进到下一节点。支持并行网关(Fork / Join)。
  537. /// 行为约定:
  538. /// - 进入本方法前,调用方(如 <see cref="Approve"/>)应已将 <paramref name="currentNodeId"/>(userTask)标记为完成节点;
  539. /// 本方法会将途经的网关节点也写入 <see cref="ApprovalFlowCompletedNode"/>。
  540. /// - 并行网关 Fork(出边&gt;=2):沿每条出边递归推进;
  541. /// 并行网关 Join(入边&gt;=2):校验所有前驱节点是否都在"已完成"集合中,
  542. /// 任一尚未完成则**静默等待**(不报错、不推进),由后续分支完成后再次触发 Join 校验。
  543. /// </summary>
  544. private async Task AdvanceToNext(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string currentNodeId)
  545. {
  546. // 当前节点可能是 userTask(完成记录由 Approve 写入)或网关(入口处已写入);此处统一确保幂等入库
  547. var currentNode = flowData.Nodes.FirstOrDefault(n => n.Id == currentNodeId);
  548. await MarkNodeCompleted(instance.Id, currentNode);
  549. var outgoingEdges = flowData.Edges.Where(e => e.SourceNodeId == currentNodeId).ToList();
  550. if (outgoingEdges.Count == 0)
  551. {
  552. await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
  553. return;
  554. }
  555. // 非并行网关场景:当前节点一般只有 1 条出边
  556. foreach (var edge in outgoingEdges)
  557. {
  558. await ProcessNextNode(instance, flowData, edge.TargetNodeId);
  559. }
  560. }
  561. /// <summary>
  562. /// 处理某个"下一节点"。根据节点类型分发:
  563. /// - endEvent:所有分支任务都结束时触发实例完成
  564. /// - exclusiveGateway:按条件选择分支
  565. /// - parallelGateway:Fork 并行分发;Join 等待所有前驱完成
  566. /// - userTask / 其他:创建任务
  567. /// </summary>
  568. private async Task ProcessNextNode(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string nextNodeId)
  569. {
  570. var nextNode = flowData.Nodes.FirstOrDefault(n => n.Id == nextNodeId);
  571. if (nextNode == null)
  572. {
  573. await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
  574. return;
  575. }
  576. if (nextNode.Type is "bpmn:endEvent" or "end-node")
  577. {
  578. // 所有并行分支都已结束(无其它 Pending 任务)才真正完成实例
  579. var hasOtherPending = await _taskRep.AsQueryable()
  580. .AnyAsync(t => t.InstanceId == instance.Id && t.Status == FlowTaskStatusEnum.Pending);
  581. if (hasOtherPending) return;
  582. await CompleteInstance(instance, FlowInstanceStatusEnum.Approved);
  583. return;
  584. }
  585. if (nextNode.Type is "bpmn:exclusiveGateway")
  586. {
  587. await MarkNodeCompleted(instance.Id, nextNode);
  588. var bizData = await GetBizData(instance.BizType, instance.BizId);
  589. var targetNodeId = EvaluateGateway(nextNode.Properties?.Conditions, flowData, nextNode.Id, bizData);
  590. instance.CurrentNodeId = targetNodeId;
  591. await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
  592. await ProcessNextNode(instance, flowData, targetNodeId);
  593. return;
  594. }
  595. if (nextNode.Type is "bpmn:parallelGateway")
  596. {
  597. var incoming = flowData.Edges.Where(e => e.TargetNodeId == nextNode.Id).Select(e => e.SourceNodeId).ToList();
  598. var outgoing = flowData.Edges.Where(e => e.SourceNodeId == nextNode.Id).Select(e => e.TargetNodeId).ToList();
  599. // Join 语义:入边 >= 2,需等所有前驱都已完成
  600. if (incoming.Count >= 2)
  601. {
  602. var completedSet = await GetCompletedNodeIdSet(instance.Id);
  603. if (!incoming.All(p => completedSet.Contains(p)))
  604. {
  605. // 未汇合,静默等待后续分支抵达
  606. return;
  607. }
  608. }
  609. // Fork 或 Join 通过:标记网关完成,沿所有出边推进
  610. await MarkNodeCompleted(instance.Id, nextNode);
  611. foreach (var target in outgoing)
  612. {
  613. await ProcessNextNode(instance, flowData, target);
  614. }
  615. return;
  616. }
  617. // userTask 或其他:创建任务
  618. instance.CurrentNodeId = nextNodeId;
  619. await _instanceRep.AsUpdateable(instance).UpdateColumns(i => new { i.CurrentNodeId }).ExecuteCommandAsync();
  620. await CreateTasksForNode(instance, flowData, nextNodeId);
  621. }
  622. /// <summary>
  623. /// 标记节点已完成(幂等:重复写入被唯一索引拦截后忽略)
  624. /// </summary>
  625. private async Task MarkNodeCompleted(long instanceId, ApprovalFlowNodeItem? node)
  626. {
  627. if (node == null) return;
  628. try
  629. {
  630. await _completedNodeRep.InsertAsync(new ApprovalFlowCompletedNode
  631. {
  632. InstanceId = instanceId,
  633. NodeId = node.Id,
  634. NodeName = node.Properties?.NodeName ?? node.Text?.Value,
  635. NodeType = node.Type,
  636. CompletedTime = DateTime.Now,
  637. });
  638. }
  639. catch
  640. {
  641. // 并发场景下可能触发唯一索引冲突,忽略(已存在即可)
  642. }
  643. }
  644. /// <summary>
  645. /// 查询实例已完成节点 Id 集合
  646. /// </summary>
  647. private async Task<HashSet<string>> GetCompletedNodeIdSet(long instanceId)
  648. {
  649. var ids = await _completedNodeRep.AsQueryable()
  650. .Where(c => c.InstanceId == instanceId)
  651. .Select(c => c.NodeId)
  652. .ToListAsync();
  653. return new HashSet<string>(ids);
  654. }
  655. private async Task CompleteInstance(ApprovalFlowInstance instance, FlowInstanceStatusEnum status)
  656. {
  657. // 幂等:并行分支同时到达 end 时避免重复完成
  658. var latest = await _instanceRep.GetByIdAsync(instance.Id);
  659. if (latest == null || latest.Status != FlowInstanceStatusEnum.Running) return;
  660. instance.Status = status;
  661. instance.EndTime = DateTime.Now;
  662. await _instanceRep.AsUpdateable(instance)
  663. .UpdateColumns(i => new { i.Status, i.EndTime })
  664. .ExecuteCommandAsync();
  665. // P4-17: 正常完成路径的 lastApproverId = 最后一条人工 Approve 日志的操作人(系统自动通过不算)
  666. var lastApproverId = status == FlowInstanceStatusEnum.Approved
  667. ? await GetLastHumanApproverIdAsync(instance.Id)
  668. : null;
  669. await InvokeHandler(instance.BizType, instance.Id,
  670. h => h.OnFlowCompleted(instance.BizId, instance.Id, status, lastApproverId));
  671. await _notifyService.NotifyFlowCompleted(instance.InitiatorId, instance.Id, instance.Title, status, instance.BizType);
  672. }
  673. /// <summary>
  674. /// P4-17: 查询该实例最后一条人工 Approve 日志的操作人 UserId
  675. /// (系统超时自动通过写入的是 AutoTimeout,不会被此查询命中)
  676. /// </summary>
  677. private async Task<long?> GetLastHumanApproverIdAsync(long instanceId)
  678. {
  679. var lastLog = await _logRep.AsQueryable()
  680. .Where(l => l.InstanceId == instanceId && l.Action == FlowLogActionEnum.Approve)
  681. .OrderByDescending(l => l.CreateTime)
  682. .FirstAsync();
  683. return lastLog?.OperatorId > 0 ? lastLog.OperatorId : null;
  684. }
  685. private async Task CreateTasksForNode(ApprovalFlowInstance instance, ApprovalFlowItem flowData, string nodeId)
  686. {
  687. var node = flowData.Nodes.FirstOrDefault(n => n.Id == nodeId)
  688. ?? throw Oops.Oh($"FlowJson 中未找到节点 [{nodeId}]");
  689. var approvers = await ResolveApprovers(node.Properties, instance.InitiatorId);
  690. if (approvers.Count == 0)
  691. throw Oops.Oh($"节点 [{node.Properties?.NodeName ?? nodeId}] 未配置审批人或审批人列表为空");
  692. var nodeName = node.Properties?.NodeName ?? node.Text?.Value;
  693. var tasks = approvers.Select(a => new ApprovalFlowTask
  694. {
  695. InstanceId = instance.Id,
  696. NodeId = nodeId,
  697. NodeName = nodeName,
  698. AssigneeId = a.userId,
  699. AssigneeName = a.userName,
  700. Status = FlowTaskStatusEnum.Pending,
  701. OrgId = instance.OrgId,
  702. }).ToList();
  703. // P1-6 审批代理:为每个原审批人检查是否存在有效代理,是则并行创建一条代理任务
  704. var delegateTasks = new List<ApprovalFlowTask>();
  705. foreach (var a in approvers)
  706. {
  707. var del = await FindEffectiveDelegate(a.userId, instance.BizType);
  708. if (del == null) continue;
  709. // 代理人和原审批人不能重复,代理人也不能是原审批人列表中其他人(避免同一人两条任务)
  710. if (approvers.Any(x => x.userId == del.DelegateUserId)) continue;
  711. delegateTasks.Add(new ApprovalFlowTask
  712. {
  713. InstanceId = instance.Id,
  714. NodeId = nodeId,
  715. NodeName = nodeName,
  716. AssigneeId = del.DelegateUserId,
  717. AssigneeName = del.DelegateUserName,
  718. Status = FlowTaskStatusEnum.Pending,
  719. OrgId = instance.OrgId,
  720. IsDelegate = true,
  721. DelegateForUserId = a.userId,
  722. DelegateForUserName = a.userName,
  723. });
  724. }
  725. var allTasks = tasks.Concat(delegateTasks).ToList();
  726. await _taskRep.AsInsertable(allTasks).ExecuteCommandAsync();
  727. var assigneeIds = allTasks.Select(t => t.AssigneeId).Distinct().ToList();
  728. await _notifyService.NotifyNewTask(assigneeIds, instance.Id, instance.Title, nodeName, instance.BizType);
  729. }
  730. /// <summary>
  731. /// 查找指定用户当前生效的审批代理(时间窗口内 + 已启用 + BizType 匹配或全局)
  732. /// </summary>
  733. private async Task<ApprovalFlowDelegate?> FindEffectiveDelegate(long userId, string? bizType)
  734. {
  735. var now = DateTime.Now;
  736. return await _delegateRep.AsQueryable()
  737. .Where(d => d.UserId == userId
  738. && d.IsEnabled
  739. && d.StartTime <= now
  740. && d.EndTime >= now
  741. && (string.IsNullOrEmpty(d.BizType) || d.BizType == bizType))
  742. .OrderBy(d => d.BizType == null ? 1 : 0) // 优先匹配指定 BizType 的代理
  743. .OrderByDescending(d => d.CreateTime)
  744. .FirstAsync();
  745. }
  746. /// <summary>
  747. /// 取消与已完成任务配对的代理任务(P1-6)
  748. /// - 本人任务完成:取消对应的代理任务
  749. /// - 代理任务完成:取消对应的本人任务
  750. /// </summary>
  751. private async Task CancelPairedDelegateTask(ApprovalFlowTask completedTask)
  752. {
  753. ApprovalFlowTask? paired;
  754. if (completedTask.IsDelegate)
  755. {
  756. var originalUserId = completedTask.DelegateForUserId ?? 0;
  757. if (originalUserId == 0) return;
  758. paired = await _taskRep.AsQueryable()
  759. .Where(t => t.InstanceId == completedTask.InstanceId
  760. && t.NodeId == completedTask.NodeId
  761. && t.Status == FlowTaskStatusEnum.Pending
  762. && t.AssigneeId == originalUserId
  763. && !t.IsDelegate)
  764. .FirstAsync();
  765. }
  766. else
  767. {
  768. paired = await _taskRep.AsQueryable()
  769. .Where(t => t.InstanceId == completedTask.InstanceId
  770. && t.NodeId == completedTask.NodeId
  771. && t.Status == FlowTaskStatusEnum.Pending
  772. && t.IsDelegate
  773. && t.DelegateForUserId == completedTask.AssigneeId)
  774. .FirstAsync();
  775. }
  776. if (paired == null) return;
  777. paired.Status = FlowTaskStatusEnum.Cancelled;
  778. paired.ActionTime = DateTime.Now;
  779. await _taskRep.AsUpdateable(paired).ExecuteCommandAsync();
  780. }
  781. private async Task<List<(long userId, string userName)>> ResolveApprovers(FlowProperties? props, long initiatorId)
  782. {
  783. if (props == null || string.IsNullOrWhiteSpace(props.ApproverType))
  784. return new List<(long, string)>();
  785. var approverType = props.ApproverType;
  786. if (approverType == nameof(ApproverTypeEnum.Initiator))
  787. {
  788. var initiator = await _userRep.GetByIdAsync(initiatorId);
  789. return initiator != null
  790. ? new List<(long, string)> { (initiator.Id, initiator.RealName ?? "") }
  791. : new List<(long, string)>();
  792. }
  793. if (string.IsNullOrWhiteSpace(props.ApproverIds))
  794. return new List<(long, string)>();
  795. var tokens = props.ApproverIds.Split(',', StringSplitOptions.RemoveEmptyEntries)
  796. .Select(s => s.Trim()).Where(s => s.Length > 0).ToList();
  797. var ids = tokens.Select(s => long.TryParse(s, out var v) ? v : 0).Where(id => id > 0).ToList();
  798. var codes = tokens.Where(s => !long.TryParse(s, out _)).ToList();
  799. if (approverType == nameof(ApproverTypeEnum.SpecificUser))
  800. {
  801. var users = await _userRep.AsQueryable()
  802. .Where(u => ids.Contains(u.Id)).ToListAsync();
  803. return users.Select(u => (u.Id, u.RealName ?? "")).ToList();
  804. }
  805. if (approverType == nameof(ApproverTypeEnum.Role))
  806. {
  807. if (codes.Count > 0)
  808. {
  809. var codeRoleIds = await _userRep.Context.Queryable<SysRole>()
  810. .ClearFilter()
  811. .Where(r => r.TenantId == _userManager.TenantId
  812. && r.Status == StatusEnum.Enable
  813. && codes.Contains(r.Code))
  814. .Select(r => r.Id)
  815. .ToListAsync();
  816. ids = ids.Concat(codeRoleIds).Distinct().ToList();
  817. }
  818. if (ids.Count == 0)
  819. return new List<(long, string)>();
  820. var userIds = await _userRoleRep.AsQueryable()
  821. .Where(ur => ids.Contains(ur.RoleId))
  822. .Select(ur => ur.UserId)
  823. .ToListAsync();
  824. // S8-S1-EXCEPTION-FLOW-SYNC-FIX-1:审批人解析是「流程配置」,不应受发起人数据范围(DataScope)隔离。
  825. // SysUser 继承 EntityBaseTenantOrg(→EntityBaseOrg),发起人 DataScope=Self 时被「CreateUserId==发起人」过滤、
  826. // Dept/DeptChild 时被 OrgId 过滤,会把角色成员(甚至发起人自己,CreateUserId 可能为 NULL)过滤掉
  827. // → 审批人列表为空 → ProcessNextNode 抛错、不建任务、实例悬挂。
  828. // ClearFilter() 清除数据范围过滤;显式补回租户隔离(TenantId==当前登录租户),等价于原全局租户过滤,无跨租户泄漏。
  829. var users = await _userRep.AsQueryable()
  830. .ClearFilter()
  831. .Where(u => userIds.Contains(u.Id) && u.TenantId == _userManager.TenantId)
  832. .ToListAsync();
  833. return users.Select(u => (u.Id, u.RealName ?? "")).ToList();
  834. }
  835. if (approverType == nameof(ApproverTypeEnum.Department))
  836. {
  837. var users = await _userRep.AsQueryable()
  838. .Where(u => ids.Contains(u.OrgId)).ToListAsync();
  839. return users.Select(u => (u.Id, u.RealName ?? "")).ToList();
  840. }
  841. if (approverType == nameof(ApproverTypeEnum.DepartmentLeader))
  842. {
  843. var initiator = await _userRep.GetByIdAsync(initiatorId);
  844. if (initiator == null || initiator.OrgId <= 0)
  845. return new List<(long, string)>();
  846. var org = await _orgRep.GetByIdAsync(initiator.OrgId);
  847. if (org?.DirectorId != null && org.DirectorId > 0)
  848. {
  849. var director = await _userRep.GetByIdAsync(org.DirectorId.Value);
  850. if (director != null)
  851. return new List<(long, string)> { (director.Id, director.RealName ?? "") };
  852. }
  853. if (initiator.ManagerUserId != null && initiator.ManagerUserId > 0)
  854. {
  855. var manager = await _userRep.GetByIdAsync(initiator.ManagerUserId.Value);
  856. if (manager != null)
  857. return new List<(long, string)> { (manager.Id, manager.RealName ?? "") };
  858. }
  859. return new List<(long, string)>();
  860. }
  861. return new List<(long, string)>();
  862. }
  863. /// <summary>
  864. /// 评估排他网关 — 依次尝试各非默认分支的条件表达式,首个匹配的获胜;
  865. /// 全不匹配则走默认分支;无默认则走第一条出边
  866. /// 支持简单比较表达式:variable op value(op: ==,!=,>,>=,&lt;,&lt;=)
  867. /// </summary>
  868. private string EvaluateGateway(List<GatewayCondition>? conditions, ApprovalFlowItem flowData, string gatewayNodeId, Dictionary<string, object>? bizData)
  869. {
  870. if (conditions != null && conditions.Count > 0 && bizData != null && bizData.Count > 0)
  871. {
  872. foreach (var cond in conditions.Where(c => !c.IsDefault))
  873. {
  874. if (!string.IsNullOrWhiteSpace(cond.Expression) && EvalSimpleExpression(cond.Expression, bizData))
  875. return cond.TargetNodeId;
  876. }
  877. var defaultBranch = conditions.FirstOrDefault(c => c.IsDefault);
  878. if (defaultBranch != null)
  879. return defaultBranch.TargetNodeId;
  880. }
  881. else if (conditions != null && conditions.Count > 0)
  882. {
  883. var defaultBranch = conditions.FirstOrDefault(c => c.IsDefault);
  884. if (defaultBranch != null) return defaultBranch.TargetNodeId;
  885. return conditions.First().TargetNodeId;
  886. }
  887. var edge = flowData.Edges.FirstOrDefault(e => e.SourceNodeId == gatewayNodeId);
  888. return edge?.TargetNodeId ?? throw Oops.Oh("排他网关没有出边");
  889. }
  890. /// <summary>
  891. /// 简单表达式求值:支持 "field op value" 格式(如 "urgent == 1", "customLevel >= 3", "amount > 10000")
  892. /// 多条件用 &amp;&amp; 连接
  893. /// </summary>
  894. private static bool EvalSimpleExpression(string expression, Dictionary<string, object> bizData)
  895. {
  896. var parts = expression.Split("&&", StringSplitOptions.TrimEntries);
  897. foreach (var part in parts)
  898. {
  899. if (!EvalSingleComparison(part.Trim(), bizData))
  900. return false;
  901. }
  902. return true;
  903. }
  904. private static bool EvalSingleComparison(string expr, Dictionary<string, object> bizData)
  905. {
  906. string[] ops = { ">=", "<=", "!=", "==", ">", "<" };
  907. foreach (var op in ops)
  908. {
  909. var idx = expr.IndexOf(op, StringComparison.Ordinal);
  910. if (idx < 0) continue;
  911. var fieldName = expr[..idx].Trim();
  912. var valueStr = expr[(idx + op.Length)..].Trim().Trim('"', '\'');
  913. if (!bizData.TryGetValue(fieldName, out var fieldValue))
  914. return false;
  915. if (decimal.TryParse(fieldValue?.ToString(), out var numLeft) && decimal.TryParse(valueStr, out var numRight))
  916. {
  917. return op switch
  918. {
  919. "==" => numLeft == numRight,
  920. "!=" => numLeft != numRight,
  921. ">" => numLeft > numRight,
  922. ">=" => numLeft >= numRight,
  923. "<" => numLeft < numRight,
  924. "<=" => numLeft <= numRight,
  925. _ => false,
  926. };
  927. }
  928. var strLeft = fieldValue?.ToString() ?? "";
  929. return op switch
  930. {
  931. "==" => strLeft.Equals(valueStr, StringComparison.OrdinalIgnoreCase),
  932. "!=" => !strLeft.Equals(valueStr, StringComparison.OrdinalIgnoreCase),
  933. _ => false,
  934. };
  935. }
  936. return false;
  937. }
  938. private async Task<Dictionary<string, object>?> GetBizData(string bizType, long bizId)
  939. {
  940. var handlers = App.GetServices<IFlowBizHandler>();
  941. var handler = handlers?.FirstOrDefault(h => h.BizType == bizType);
  942. if (handler == null) return null;
  943. return await handler.GetBizData(bizId);
  944. }
  945. private string? FindNextNodeId(ApprovalFlowItem flowData, string currentNodeId)
  946. {
  947. var edge = flowData.Edges.FirstOrDefault(e => e.SourceNodeId == currentNodeId);
  948. return edge?.TargetNodeId;
  949. }
  950. private string? FindPrevUserTaskNodeId(ApprovalFlowItem flowData, string currentNodeId)
  951. {
  952. var inEdge = flowData.Edges.FirstOrDefault(e => e.TargetNodeId == currentNodeId);
  953. if (inEdge == null) return null;
  954. var prevNode = flowData.Nodes.FirstOrDefault(n => n.Id == inEdge.SourceNodeId);
  955. if (prevNode == null) return null;
  956. if (prevNode.Type is "bpmn:userTask" or "user-node" or "task-node")
  957. return prevNode.Id;
  958. // 递归跳过网关等非用户任务节点
  959. return FindPrevUserTaskNodeId(flowData, prevNode.Id);
  960. }
  961. private async Task<bool> IsNodeCompleted(ApprovalFlowInstance instance, string nodeId)
  962. {
  963. var flowData = DeserializeFlowJson(instance.FlowJsonSnapshot);
  964. var node = flowData.Nodes.FirstOrDefault(n => n.Id == nodeId);
  965. var mode = node?.Properties?.MultiApproveMode;
  966. if (mode == nameof(MultiApproveModeEnum.All))
  967. {
  968. var pendingCount = await _taskRep.AsQueryable()
  969. .Where(t => t.InstanceId == instance.Id && t.NodeId == nodeId && t.Status == FlowTaskStatusEnum.Pending)
  970. .CountAsync();
  971. return pendingCount == 0;
  972. }
  973. // 默认或签(Any):一人通过即完成,取消其他 Pending
  974. await CancelPendingTasks(instance.Id, nodeId);
  975. return true;
  976. }
  977. private async Task CancelPendingTasks(long instanceId, string nodeId, long? excludeTaskId = null)
  978. {
  979. var tasks = await _taskRep.AsQueryable()
  980. .Where(t => t.InstanceId == instanceId && t.NodeId == nodeId && t.Status == FlowTaskStatusEnum.Pending)
  981. .WhereIF(excludeTaskId.HasValue, t => t.Id != excludeTaskId!.Value)
  982. .ToListAsync();
  983. foreach (var t in tasks)
  984. {
  985. t.Status = FlowTaskStatusEnum.Cancelled;
  986. t.ActionTime = DateTime.Now;
  987. }
  988. if (tasks.Count > 0)
  989. await _taskRep.AsUpdateable(tasks).ExecuteCommandAsync();
  990. }
  991. /// <summary>
  992. /// 取消整个实例下所有剩余 Pending 任务(Reject 时跨并行分支使用)
  993. /// </summary>
  994. private async Task CancelAllPendingTasks(long instanceId, long? excludeTaskId = null)
  995. {
  996. var tasks = await _taskRep.AsQueryable()
  997. .Where(t => t.InstanceId == instanceId && t.Status == FlowTaskStatusEnum.Pending)
  998. .WhereIF(excludeTaskId.HasValue, t => t.Id != excludeTaskId!.Value)
  999. .ToListAsync();
  1000. foreach (var t in tasks)
  1001. {
  1002. t.Status = FlowTaskStatusEnum.Cancelled;
  1003. t.ActionTime = DateTime.Now;
  1004. }
  1005. if (tasks.Count > 0)
  1006. await _taskRep.AsUpdateable(tasks).ExecuteCommandAsync();
  1007. }
  1008. private async Task<ApprovalFlowTask> GetPendingTask(long taskId)
  1009. {
  1010. var task = await _taskRep.GetByIdAsync(taskId)
  1011. ?? throw Oops.Oh("审批任务不存在");
  1012. if (task.Status != FlowTaskStatusEnum.Pending)
  1013. throw Oops.Oh("该任务已处理");
  1014. if (task.AssigneeId != _userManager.UserId)
  1015. throw Oops.Oh("当前用户不是该任务的审批人");
  1016. return task;
  1017. }
  1018. private async Task WriteLog(long instanceId, long? taskId, string? nodeId, FlowLogActionEnum action, string? comment)
  1019. {
  1020. await _logRep.InsertAsync(new ApprovalFlowLog
  1021. {
  1022. InstanceId = instanceId,
  1023. TaskId = taskId,
  1024. NodeId = nodeId,
  1025. Action = action,
  1026. OperatorId = _userManager.UserId,
  1027. OperatorName = _userManager.RealName,
  1028. Comment = comment,
  1029. });
  1030. }
  1031. private static ApprovalFlowItem DeserializeFlowJson(string? json)
  1032. {
  1033. if (string.IsNullOrWhiteSpace(json))
  1034. throw Oops.Oh("FlowJson 快照为空");
  1035. return JsonSerializer.Deserialize<ApprovalFlowItem>(json)
  1036. ?? throw Oops.Oh("FlowJson 反序列化失败");
  1037. }
  1038. private async Task InvokeHandler(string bizType, long instanceId, Func<IFlowBizHandler, Task> action)
  1039. {
  1040. var handlers = App.GetServices<IFlowBizHandler>();
  1041. var handler = handlers?.FirstOrDefault(h => h.BizType == bizType);
  1042. if (handler == null)
  1043. return;
  1044. try
  1045. {
  1046. await action(handler);
  1047. }
  1048. catch (Exception ex)
  1049. {
  1050. _logger.LogError(ex, "FlowBizHandler 执行失败: BizType={BizType}, InstanceId={InstanceId}, ErrorType={ErrorType}, ErrorMessage={ErrorMessage}",
  1051. bizType, instanceId, ex.GetType().FullName, ex.Message);
  1052. var notifiers = App.GetServices<IFlowHandlerFailureNotifier>();
  1053. if (notifiers != null)
  1054. {
  1055. foreach (var notifier in notifiers)
  1056. {
  1057. try
  1058. {
  1059. await notifier.NotifyAsync(bizType, instanceId, ex);
  1060. }
  1061. catch (Exception notifyEx)
  1062. {
  1063. _logger.LogError(notifyEx, "FlowHandlerFailureNotifier 执行失败: NotifierType={NotifierType}, BizType={BizType}, InstanceId={InstanceId}",
  1064. notifier.GetType().FullName, bizType, instanceId);
  1065. }
  1066. }
  1067. }
  1068. throw;
  1069. }
  1070. }
  1071. }