| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393 |
- using System;
- using System.Collections.Generic;
- using System.Diagnostics;
- using System.Linq;
- using System.Threading;
- using System.Threading.Tasks;
- using TeamAAS.Communication;
- using TeamAAS.Communication.Config;
- using TeamAAS.Communication.Interfaces;
- using TeamAAS.FlowEditor;
- using TeamAAS.FlowEditor.Execution;
- using TeamAAS.FlowEditor.Models;
- namespace TeamAAS.FlowEngine.Execution
- {
- /// <summary>
- /// 产品运行引擎:把"运行产品"与"心跳触发"真正接上线。
- /// - LoadRunProduct:加载产品的独立运行态流程(与编辑态互不影响)
- /// - RunOnceAsync:执行所有"主流程触发"的流程一次,并更新 RunStatistics
- /// - StartHeartbeat:订阅通讯设备 DataReceivedEvent,按心跳间隔节流触发一次运行
- /// 注意:主程序需先设置 <see cref="Communication"/>(CoreManager.Communication)。
- /// </summary>
- public class FlowRunner
- {
- private static FlowRunner _instance;
- public static FlowRunner Instance => _instance ?? (_instance = new FlowRunner());
- private readonly object _sync = new object();
- private readonly Dictionary<string, DateTime> _lastTrigger = new Dictionary<string, DateTime>();
- private readonly Dictionary<Guid, List<Action<object, string>>> _subscriptions =
- new Dictionary<Guid, List<Action<object, string>>>();
- private FlowRunner() { }
- /// <summary>通讯管理器(由主程序注入 CoreManager.Communication)</summary>
- public CommunicationManager Communication { get; set; }
- /// <summary>运行态流程(从产品文件加载的独立副本)</summary>
- public List<FlowTabItem> RunTabs { get; private set; } = new List<FlowTabItem>();
- /// <summary>当前运行产品名</summary>
- public string RunProductName { get; private set; } = "";
- public bool IsRunning
- {
- get { lock (_sync) return _isRunning; }
- private set
- {
- bool changed;
- lock (_sync)
- {
- changed = _isRunning != value;
- _isRunning = value;
- }
- if (changed)
- RunStateChanged?.Invoke();
- }
- }
- private bool _isRunning;
- private CancellationTokenSource _runCts;
- // 停止后后台循环仍在"排空"本轮(硬件节点无法中途打断),用该标记防止重复启动
- private volatile bool _draining;
- // 运行态流程的磁盘指纹:仅当流程文件变化时才重载,避免高频触发每次读盘反序列化
- private string _runTabsStamp;
- /// <summary>运行状态变化(开始/停止)——供首页按钮、编辑器锁定监听</summary>
- public event Action RunStateChanged;
- /// <summary>产品加载完成(参数:产品名)</summary>
- public event Action<string> ProductLoaded;
- /// <summary>一次运行结束</summary>
- public event Action RunCompleted;
- /// <summary>运行态 RunTabs 被重新加载(RefreshRunTabs 完成)——监控端订阅以自动同步显示</summary>
- public event Action RunTabsRefreshed;
- /// <summary>
- /// 加载产品为运行产品。
- /// </summary>
- public bool LoadRunProduct(string productName)
- {
- if (string.IsNullOrWhiteSpace(productName)) return false;
- // 切换运行产品前必须先停掉当前运行(首页"加载"也依赖此逻辑兜底)
- Stop();
- // 释放旧运行产品的所有插件资源(防止切换后旧插件句柄/连接残留)
- var oldTabs = RunTabs;
- RunTabs = new List<FlowTabItem>();
- ProductManager.DisposeFlowTabs(oldTabs);
- ProductManager.Instance.LoadProductToRun(productName);
- RunProductName = productName;
- RefreshRunTabs();
- RunStatistics.Instance.Log($"运行产品已加载: {productName}({RunTabs.Count} 个流程)");
- ProductLoaded?.Invoke(productName);
- return true;
- }
- /// <summary>
- /// 重新解析运行态流程:始终从磁盘加载独立副本,编辑态与运行态彻底分离。
- /// 监控端(FlowEditorShellViewModel.IsMonitorMode=true)通过订阅 RunTabsRefreshed
- /// 主动挂接 RunTabs 显示运行状态,不再隐式共享编辑器 FlowTabs。
- /// </summary>
- private void RefreshRunTabs()
- {
- // 重载运行态流程前,先完全释放 Main(含上次产品的各流程与子流程残留),
- // 保证任何加载入口(首页/PLC 序号切换/心跳重载)都从干净的 Main 开始
- ResultRegistry.Main.ClearAllFlows();
- // 运行态流程绑定主页面专用注册表(与产品管理 Debug 注册表隔离)
- var oldTabs = RunTabs;
- RunTabs = ProductManager.Instance.LoadAllFlows(RunProductName, ResultRegistry.Main);
- // 释放上一轮运行态 Tab(插件 Dispose + Main 注册表清理),
- // 防止每次开始/停止累积一份带相机图像等大对象的旧插件实例
- if (oldTabs != null && oldTabs.Count > 0)
- ProductManager.DisposeFlowTabs(oldTabs);
- RunTabsRefreshed?.Invoke();
- _runTabsStamp = ProductManager.Instance.GetProductFlowStamp(RunProductName);
- }
- /// <summary>
- /// 仅当产品流程文件发生变化(或尚未加载)时才从磁盘重载运行态流程;否则复用内存中的 RunTabs。
- /// 用于高频触发(心跳/事件)路径,避免每次触发都读盘+反序列化+重建插件实例。
- /// </summary>
- private void RefreshRunTabsIfChanged()
- {
- var stamp = ProductManager.Instance.GetProductFlowStamp(RunProductName);
- if (RunTabs != null && RunTabs.Count > 0 && stamp == _runTabsStamp)
- return;
- RefreshRunTabs();
- }
- /// <summary>
- /// 执行一次:并行执行所有"主流程触发"的流程,全部结束后统计结果。
- /// </summary>
- public async Task<bool> RunOnceAsync()
- {
- if (IsRunning || _draining) return false;
- RefreshRunTabsIfChanged();
- var mains = RunTabs.Where(t => t.TriggerType == TriggerType.MainFlow).ToList();
- if (mains.Count == 0)
- {
- RunStatistics.Instance.Log("未加载运行产品或产品中没有主流程");
- return false;
- }
- IsRunning = true;
- _draining = false;
- _runCts = new CancellationTokenSource();
- bool success = true;
- var sw = Stopwatch.StartNew();
- try
- {
- RunStatistics.Instance.Log("开始运行...");
- var tasks = mains.Select(t => ExecuteTabAsync(t, _runCts.Token)).ToArray();
- var results = await Task.WhenAll(tasks);
- success = results.All(r => r);
- }
- finally
- {
- sw.Stop();
- _runCts?.Dispose();
- _runCts = null;
- _draining = false;
- IsRunning = false;
- RunStatistics.Instance.AddRun(success, sw.ElapsedMilliseconds);
- RunCompleted?.Invoke();
- }
- return success;
- }
- /// <summary>
- /// 循环运行:每个"主流程触发"的流程在独立的循环 Task 中运行,互不等待。
- /// 流程 A 慢、流程 B 快时,B 不会被 A 拖住等待;任何一个流程本轮完成立即进入下一轮。
- /// 直到 Stop() 被调用(取消 token)所有循环退出。
- /// </summary>
- public async Task RunLoopAsync()
- {
- if (IsRunning || _draining) return;
- RefreshRunTabs();
- var mains = RunTabs.Where(t => t.TriggerType == TriggerType.MainFlow).ToList();
- if (mains.Count == 0)
- {
- RunStatistics.Instance.Log("未加载运行产品或产品中没有主流程");
- return;
- }
- IsRunning = true;
- _draining = false;
- _runCts = new CancellationTokenSource();
- var token = _runCts.Token;
- RunStatistics.Instance.Log("开始循环运行(点击停止按钮结束)...");
- try
- {
- // 每个流程一个独立循环,互不等待
- var loopTasks = mains.Select(tab => RunSingleTabLoopAsync(tab, token)).ToList();
- await Task.WhenAll(loopTasks);
- }
- finally
- {
- _runCts?.Dispose();
- _runCts = null;
- _draining = false;
- IsRunning = false;
- RunStatistics.Instance.Log("运行已停止");
- RunCompleted?.Invoke();
- }
- }
- /// <summary>
- /// 单个流程的独立循环:完成一轮立即进入下一轮,不等其他流程。
- /// </summary>
- private async Task RunSingleTabLoopAsync(FlowTabItem tab, CancellationToken token)
- {
- while (!token.IsCancellationRequested)
- {
- var sw = Stopwatch.StartNew();
- bool ok = await ExecuteTabAsync(tab, token);
- sw.Stop();
- if (token.IsCancellationRequested) break;
- RunStatistics.Instance.AddRun(ok, sw.ElapsedMilliseconds);
- // 循环间隔:单轮结束后等待用户设置的毫秒数再开始下一轮(最小1ms,防止空流程卡CPU)
- try { await Task.Delay(tab.LoopIntervalMs, token); }
- catch (OperationCanceledException) { break; }
- }
- }
- /// <summary>
- /// 停止循环运行 / 单次运行。
- /// 立即解除"运行中"状态(编辑器解锁、按钮复位);后台循环在完成本轮后自行结束。
- /// </summary>
- public void Stop()
- {
- _runCts?.Cancel();
- IsRunning = false;
- // 完全释放运行态结果缓存(含各流程与其子流程条目),释放挂住的大对象引用
- try { ResultRegistry.Main.ClearAllFlows(); }
- catch { }
- // 强制全代 GC,确保 LOH 大对象(图像、CogRecord 等)被回收
- GC.Collect(2, GCCollectionMode.Forced, true);
- GC.WaitForPendingFinalizers();
- GC.Collect(2, GCCollectionMode.Forced, true);
- }
- /// <summary>
- /// 按名称执行一个运行态流程(供事件触发等外部场景调用)。
- /// </summary>
- public async Task<bool> RunFlowByNameAsync(string flowName)
- {
- var tab = RunTabs.FirstOrDefault(t => t.Name == flowName);
- if (tab == null)
- {
- AppLogger.Warning($"运行流程不存在: {flowName}", nameof(FlowRunner));
- return false;
- }
- return await ExecuteTabAsync(tab, CancellationToken.None);
- }
- private async Task<bool> ExecuteTabAsync(FlowTabItem tab, CancellationToken token)
- {
- if (tab?.Graph == null) return false;
- var editorVm = tab.EditorVm;
- SetEditorRunning(editorVm, true);
- try
- {
- var executor = new FlowExecutor(tab.Graph);
- var tcs = new TaskCompletionSource<bool>();
- executor.ExecutionCompleted += ok => tcs.TrySetResult(ok);
- await executor.ExecuteAsync(token);
- return await tcs.Task;
- }
- catch (Exception ex)
- {
- AppLogger.Error($"流程执行异常: {tab.Name}", ex, nameof(FlowRunner));
- return false;
- }
- }
- /// <summary>
- /// 在 UI 线程同步编辑器 ViewModel 的 IsRunning,触发 RunFlowCommand 等按钮状态刷新。
- /// </summary>
- private static void SetEditorRunning(FlowEditorViewModel editorVm, bool running)
- {
- if (editorVm == null) return;
- System.Windows.Application.Current?.Dispatcher.Invoke(() =>
- {
- editorVm.IsRunning = running;
- });
- }
- #region 心跳触发
- /// <summary>
- /// 根据全局事件配置订阅通讯设备的数据到达事件,作为产品运行触发器。
- /// 同一设备按 IntervalMs 节流,防止高频数据把执行引擎打爆。
- /// </summary>
- public void StartHeartbeat()
- {
- StopHeartbeat();
- var cfg = GlobalVariableManager.LoadGlobalEventConfig();
- foreach (var hb in cfg.Heartbeats)
- {
- if (hb == null || !hb.Enabled) continue;
- var device = FindDevice(hb);
- if (device == null)
- {
- AppLogger.Warning($"心跳设备未找到: {hb.DeviceName}", nameof(FlowRunner));
- continue;
- }
- Action<object, string> handler = (s, data) => OnHeartbeatData(hb);
- device.DataReceivedEvent += handler;
- lock (_sync)
- {
- if (!_subscriptions.TryGetValue(device.Id, out var list))
- {
- list = new List<Action<object, string>>();
- _subscriptions[device.Id] = list;
- }
- list.Add(handler);
- }
- }
- AppLogger.Info($"心跳监听已启动,订阅 {_subscriptions.Count} 个设备", nameof(FlowRunner));
- }
- public void StopHeartbeat()
- {
- lock (_sync)
- {
- foreach (var kvp in _subscriptions)
- {
- var device = Communication?.Get(kvp.Key);
- if (device == null) continue;
- foreach (var handler in kvp.Value)
- device.DataReceivedEvent -= handler;
- }
- _subscriptions.Clear();
- _lastTrigger.Clear();
- }
- }
- private async void OnHeartbeatData(HeartbeatConfig hb)
- {
- if (hb == null || !hb.Enabled) return;
- var now = DateTime.Now;
- lock (_sync)
- {
- string key = hb.DeviceId.ToString();
- if (_lastTrigger.TryGetValue(key, out var last)
- && (now - last).TotalMilliseconds < hb.IntervalMs)
- return;
- _lastTrigger[key] = now;
- }
- await RunOnceAsync();
- }
- private ICommunication FindDevice(HeartbeatConfig hb)
- {
- if (Communication == null) return null;
- if (hb.DeviceId != Guid.Empty)
- {
- var byId = Communication.Get(hb.DeviceId);
- if (byId != null) return byId;
- }
- if (!string.IsNullOrWhiteSpace(hb.DeviceName))
- return Communication.GetByName(hb.DeviceName);
- return null;
- }
- #endregion
- }
- }
|