FlowEngineService.cs 61 KB

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