FlowEngineService.cs 53 KB

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