FlowRunner.cs 15 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393
  1. using System;
  2. using System.Collections.Generic;
  3. using System.Diagnostics;
  4. using System.Linq;
  5. using System.Threading;
  6. using System.Threading.Tasks;
  7. using TeamAAS.Communication;
  8. using TeamAAS.Communication.Config;
  9. using TeamAAS.Communication.Interfaces;
  10. using TeamAAS.FlowEditor;
  11. using TeamAAS.FlowEditor.Execution;
  12. using TeamAAS.FlowEditor.Models;
  13. namespace TeamAAS.FlowEngine.Execution
  14. {
  15. /// <summary>
  16. /// 产品运行引擎:把"运行产品"与"心跳触发"真正接上线。
  17. /// - LoadRunProduct:加载产品的独立运行态流程(与编辑态互不影响)
  18. /// - RunOnceAsync:执行所有"主流程触发"的流程一次,并更新 RunStatistics
  19. /// - StartHeartbeat:订阅通讯设备 DataReceivedEvent,按心跳间隔节流触发一次运行
  20. /// 注意:主程序需先设置 <see cref="Communication"/>(CoreManager.Communication)。
  21. /// </summary>
  22. public class FlowRunner
  23. {
  24. private static FlowRunner _instance;
  25. public static FlowRunner Instance => _instance ?? (_instance = new FlowRunner());
  26. private readonly object _sync = new object();
  27. private readonly Dictionary<string, DateTime> _lastTrigger = new Dictionary<string, DateTime>();
  28. private readonly Dictionary<Guid, List<Action<object, string>>> _subscriptions =
  29. new Dictionary<Guid, List<Action<object, string>>>();
  30. private FlowRunner() { }
  31. /// <summary>通讯管理器(由主程序注入 CoreManager.Communication)</summary>
  32. public CommunicationManager Communication { get; set; }
  33. /// <summary>运行态流程(从产品文件加载的独立副本)</summary>
  34. public List<FlowTabItem> RunTabs { get; private set; } = new List<FlowTabItem>();
  35. /// <summary>当前运行产品名</summary>
  36. public string RunProductName { get; private set; } = "";
  37. public bool IsRunning
  38. {
  39. get { lock (_sync) return _isRunning; }
  40. private set
  41. {
  42. bool changed;
  43. lock (_sync)
  44. {
  45. changed = _isRunning != value;
  46. _isRunning = value;
  47. }
  48. if (changed)
  49. RunStateChanged?.Invoke();
  50. }
  51. }
  52. private bool _isRunning;
  53. private CancellationTokenSource _runCts;
  54. // 停止后后台循环仍在"排空"本轮(硬件节点无法中途打断),用该标记防止重复启动
  55. private volatile bool _draining;
  56. // 运行态流程的磁盘指纹:仅当流程文件变化时才重载,避免高频触发每次读盘反序列化
  57. private string _runTabsStamp;
  58. /// <summary>运行状态变化(开始/停止)——供首页按钮、编辑器锁定监听</summary>
  59. public event Action RunStateChanged;
  60. /// <summary>产品加载完成(参数:产品名)</summary>
  61. public event Action<string> ProductLoaded;
  62. /// <summary>一次运行结束</summary>
  63. public event Action RunCompleted;
  64. /// <summary>运行态 RunTabs 被重新加载(RefreshRunTabs 完成)——监控端订阅以自动同步显示</summary>
  65. public event Action RunTabsRefreshed;
  66. /// <summary>
  67. /// 加载产品为运行产品。
  68. /// </summary>
  69. public bool LoadRunProduct(string productName)
  70. {
  71. if (string.IsNullOrWhiteSpace(productName)) return false;
  72. // 切换运行产品前必须先停掉当前运行(首页"加载"也依赖此逻辑兜底)
  73. Stop();
  74. // 释放旧运行产品的所有插件资源(防止切换后旧插件句柄/连接残留)
  75. var oldTabs = RunTabs;
  76. RunTabs = new List<FlowTabItem>();
  77. ProductManager.DisposeFlowTabs(oldTabs);
  78. ProductManager.Instance.LoadProductToRun(productName);
  79. RunProductName = productName;
  80. RefreshRunTabs();
  81. RunStatistics.Instance.Log($"运行产品已加载: {productName}({RunTabs.Count} 个流程)");
  82. ProductLoaded?.Invoke(productName);
  83. return true;
  84. }
  85. /// <summary>
  86. /// 重新解析运行态流程:始终从磁盘加载独立副本,编辑态与运行态彻底分离。
  87. /// 监控端(FlowEditorShellViewModel.IsMonitorMode=true)通过订阅 RunTabsRefreshed
  88. /// 主动挂接 RunTabs 显示运行状态,不再隐式共享编辑器 FlowTabs。
  89. /// </summary>
  90. private void RefreshRunTabs()
  91. {
  92. // 重载运行态流程前,先完全释放 Main(含上次产品的各流程与子流程残留),
  93. // 保证任何加载入口(首页/PLC 序号切换/心跳重载)都从干净的 Main 开始
  94. ResultRegistry.Main.ClearAllFlows();
  95. // 运行态流程绑定主页面专用注册表(与产品管理 Debug 注册表隔离)
  96. var oldTabs = RunTabs;
  97. RunTabs = ProductManager.Instance.LoadAllFlows(RunProductName, ResultRegistry.Main);
  98. // 释放上一轮运行态 Tab(插件 Dispose + Main 注册表清理),
  99. // 防止每次开始/停止累积一份带相机图像等大对象的旧插件实例
  100. if (oldTabs != null && oldTabs.Count > 0)
  101. ProductManager.DisposeFlowTabs(oldTabs);
  102. RunTabsRefreshed?.Invoke();
  103. _runTabsStamp = ProductManager.Instance.GetProductFlowStamp(RunProductName);
  104. }
  105. /// <summary>
  106. /// 仅当产品流程文件发生变化(或尚未加载)时才从磁盘重载运行态流程;否则复用内存中的 RunTabs。
  107. /// 用于高频触发(心跳/事件)路径,避免每次触发都读盘+反序列化+重建插件实例。
  108. /// </summary>
  109. private void RefreshRunTabsIfChanged()
  110. {
  111. var stamp = ProductManager.Instance.GetProductFlowStamp(RunProductName);
  112. if (RunTabs != null && RunTabs.Count > 0 && stamp == _runTabsStamp)
  113. return;
  114. RefreshRunTabs();
  115. }
  116. /// <summary>
  117. /// 执行一次:并行执行所有"主流程触发"的流程,全部结束后统计结果。
  118. /// </summary>
  119. public async Task<bool> RunOnceAsync()
  120. {
  121. if (IsRunning || _draining) return false;
  122. RefreshRunTabsIfChanged();
  123. var mains = RunTabs.Where(t => t.TriggerType == TriggerType.MainFlow).ToList();
  124. if (mains.Count == 0)
  125. {
  126. RunStatistics.Instance.Log("未加载运行产品或产品中没有主流程");
  127. return false;
  128. }
  129. IsRunning = true;
  130. _draining = false;
  131. _runCts = new CancellationTokenSource();
  132. bool success = true;
  133. var sw = Stopwatch.StartNew();
  134. try
  135. {
  136. RunStatistics.Instance.Log("开始运行...");
  137. var tasks = mains.Select(t => ExecuteTabAsync(t, _runCts.Token)).ToArray();
  138. var results = await Task.WhenAll(tasks);
  139. success = results.All(r => r);
  140. }
  141. finally
  142. {
  143. sw.Stop();
  144. _runCts?.Dispose();
  145. _runCts = null;
  146. _draining = false;
  147. IsRunning = false;
  148. RunStatistics.Instance.AddRun(success, sw.ElapsedMilliseconds);
  149. RunCompleted?.Invoke();
  150. }
  151. return success;
  152. }
  153. /// <summary>
  154. /// 循环运行:每个"主流程触发"的流程在独立的循环 Task 中运行,互不等待。
  155. /// 流程 A 慢、流程 B 快时,B 不会被 A 拖住等待;任何一个流程本轮完成立即进入下一轮。
  156. /// 直到 Stop() 被调用(取消 token)所有循环退出。
  157. /// </summary>
  158. public async Task RunLoopAsync()
  159. {
  160. if (IsRunning || _draining) return;
  161. RefreshRunTabs();
  162. var mains = RunTabs.Where(t => t.TriggerType == TriggerType.MainFlow).ToList();
  163. if (mains.Count == 0)
  164. {
  165. RunStatistics.Instance.Log("未加载运行产品或产品中没有主流程");
  166. return;
  167. }
  168. IsRunning = true;
  169. _draining = false;
  170. _runCts = new CancellationTokenSource();
  171. var token = _runCts.Token;
  172. RunStatistics.Instance.Log("开始循环运行(点击停止按钮结束)...");
  173. try
  174. {
  175. // 每个流程一个独立循环,互不等待
  176. var loopTasks = mains.Select(tab => RunSingleTabLoopAsync(tab, token)).ToList();
  177. await Task.WhenAll(loopTasks);
  178. }
  179. finally
  180. {
  181. _runCts?.Dispose();
  182. _runCts = null;
  183. _draining = false;
  184. IsRunning = false;
  185. RunStatistics.Instance.Log("运行已停止");
  186. RunCompleted?.Invoke();
  187. }
  188. }
  189. /// <summary>
  190. /// 单个流程的独立循环:完成一轮立即进入下一轮,不等其他流程。
  191. /// </summary>
  192. private async Task RunSingleTabLoopAsync(FlowTabItem tab, CancellationToken token)
  193. {
  194. while (!token.IsCancellationRequested)
  195. {
  196. var sw = Stopwatch.StartNew();
  197. bool ok = await ExecuteTabAsync(tab, token);
  198. sw.Stop();
  199. if (token.IsCancellationRequested) break;
  200. RunStatistics.Instance.AddRun(ok, sw.ElapsedMilliseconds);
  201. // 循环间隔:单轮结束后等待用户设置的毫秒数再开始下一轮(最小1ms,防止空流程卡CPU)
  202. try { await Task.Delay(tab.LoopIntervalMs, token); }
  203. catch (OperationCanceledException) { break; }
  204. }
  205. }
  206. /// <summary>
  207. /// 停止循环运行 / 单次运行。
  208. /// 立即解除"运行中"状态(编辑器解锁、按钮复位);后台循环在完成本轮后自行结束。
  209. /// </summary>
  210. public void Stop()
  211. {
  212. _runCts?.Cancel();
  213. IsRunning = false;
  214. // 完全释放运行态结果缓存(含各流程与其子流程条目),释放挂住的大对象引用
  215. try { ResultRegistry.Main.ClearAllFlows(); }
  216. catch { }
  217. // 强制全代 GC,确保 LOH 大对象(图像、CogRecord 等)被回收
  218. GC.Collect(2, GCCollectionMode.Forced, true);
  219. GC.WaitForPendingFinalizers();
  220. GC.Collect(2, GCCollectionMode.Forced, true);
  221. }
  222. /// <summary>
  223. /// 按名称执行一个运行态流程(供事件触发等外部场景调用)。
  224. /// </summary>
  225. public async Task<bool> RunFlowByNameAsync(string flowName)
  226. {
  227. var tab = RunTabs.FirstOrDefault(t => t.Name == flowName);
  228. if (tab == null)
  229. {
  230. AppLogger.Warning($"运行流程不存在: {flowName}", nameof(FlowRunner));
  231. return false;
  232. }
  233. return await ExecuteTabAsync(tab, CancellationToken.None);
  234. }
  235. private async Task<bool> ExecuteTabAsync(FlowTabItem tab, CancellationToken token)
  236. {
  237. if (tab?.Graph == null) return false;
  238. var editorVm = tab.EditorVm;
  239. SetEditorRunning(editorVm, true);
  240. try
  241. {
  242. var executor = new FlowExecutor(tab.Graph);
  243. var tcs = new TaskCompletionSource<bool>();
  244. executor.ExecutionCompleted += ok => tcs.TrySetResult(ok);
  245. await executor.ExecuteAsync(token);
  246. return await tcs.Task;
  247. }
  248. catch (Exception ex)
  249. {
  250. AppLogger.Error($"流程执行异常: {tab.Name}", ex, nameof(FlowRunner));
  251. return false;
  252. }
  253. }
  254. /// <summary>
  255. /// 在 UI 线程同步编辑器 ViewModel 的 IsRunning,触发 RunFlowCommand 等按钮状态刷新。
  256. /// </summary>
  257. private static void SetEditorRunning(FlowEditorViewModel editorVm, bool running)
  258. {
  259. if (editorVm == null) return;
  260. System.Windows.Application.Current?.Dispatcher.Invoke(() =>
  261. {
  262. editorVm.IsRunning = running;
  263. });
  264. }
  265. #region 心跳触发
  266. /// <summary>
  267. /// 根据全局事件配置订阅通讯设备的数据到达事件,作为产品运行触发器。
  268. /// 同一设备按 IntervalMs 节流,防止高频数据把执行引擎打爆。
  269. /// </summary>
  270. public void StartHeartbeat()
  271. {
  272. StopHeartbeat();
  273. var cfg = GlobalVariableManager.LoadGlobalEventConfig();
  274. foreach (var hb in cfg.Heartbeats)
  275. {
  276. if (hb == null || !hb.Enabled) continue;
  277. var device = FindDevice(hb);
  278. if (device == null)
  279. {
  280. AppLogger.Warning($"心跳设备未找到: {hb.DeviceName}", nameof(FlowRunner));
  281. continue;
  282. }
  283. Action<object, string> handler = (s, data) => OnHeartbeatData(hb);
  284. device.DataReceivedEvent += handler;
  285. lock (_sync)
  286. {
  287. if (!_subscriptions.TryGetValue(device.Id, out var list))
  288. {
  289. list = new List<Action<object, string>>();
  290. _subscriptions[device.Id] = list;
  291. }
  292. list.Add(handler);
  293. }
  294. }
  295. AppLogger.Info($"心跳监听已启动,订阅 {_subscriptions.Count} 个设备", nameof(FlowRunner));
  296. }
  297. public void StopHeartbeat()
  298. {
  299. lock (_sync)
  300. {
  301. foreach (var kvp in _subscriptions)
  302. {
  303. var device = Communication?.Get(kvp.Key);
  304. if (device == null) continue;
  305. foreach (var handler in kvp.Value)
  306. device.DataReceivedEvent -= handler;
  307. }
  308. _subscriptions.Clear();
  309. _lastTrigger.Clear();
  310. }
  311. }
  312. private async void OnHeartbeatData(HeartbeatConfig hb)
  313. {
  314. if (hb == null || !hb.Enabled) return;
  315. var now = DateTime.Now;
  316. lock (_sync)
  317. {
  318. string key = hb.DeviceId.ToString();
  319. if (_lastTrigger.TryGetValue(key, out var last)
  320. && (now - last).TotalMilliseconds < hb.IntervalMs)
  321. return;
  322. _lastTrigger[key] = now;
  323. }
  324. await RunOnceAsync();
  325. }
  326. private ICommunication FindDevice(HeartbeatConfig hb)
  327. {
  328. if (Communication == null) return null;
  329. if (hb.DeviceId != Guid.Empty)
  330. {
  331. var byId = Communication.Get(hb.DeviceId);
  332. if (byId != null) return byId;
  333. }
  334. if (!string.IsNullOrWhiteSpace(hb.DeviceName))
  335. return Communication.GetByName(hb.DeviceName);
  336. return null;
  337. }
  338. #endregion
  339. }
  340. }