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
{
///
/// 产品运行引擎:把"运行产品"与"心跳触发"真正接上线。
/// - LoadRunProduct:加载产品的独立运行态流程(与编辑态互不影响)
/// - RunOnceAsync:执行所有"主流程触发"的流程一次,并更新 RunStatistics
/// - StartHeartbeat:订阅通讯设备 DataReceivedEvent,按心跳间隔节流触发一次运行
/// 注意:主程序需先设置 (CoreManager.Communication)。
///
public class FlowRunner
{
private static FlowRunner _instance;
public static FlowRunner Instance => _instance ?? (_instance = new FlowRunner());
private readonly object _sync = new object();
private readonly Dictionary _lastTrigger = new Dictionary();
private readonly Dictionary>> _subscriptions =
new Dictionary>>();
private FlowRunner() { }
/// 通讯管理器(由主程序注入 CoreManager.Communication)
public CommunicationManager Communication { get; set; }
/// 运行态流程(从产品文件加载的独立副本)
public List RunTabs { get; private set; } = new List();
/// 当前运行产品名
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;
/// 运行状态变化(开始/停止)——供首页按钮、编辑器锁定监听
public event Action RunStateChanged;
/// 产品加载完成(参数:产品名)
public event Action ProductLoaded;
/// 一次运行结束
public event Action RunCompleted;
/// 运行态 RunTabs 被重新加载(RefreshRunTabs 完成)——监控端订阅以自动同步显示
public event Action RunTabsRefreshed;
///
/// 加载产品为运行产品。
///
public bool LoadRunProduct(string productName)
{
if (string.IsNullOrWhiteSpace(productName)) return false;
// 切换运行产品前必须先停掉当前运行(首页"加载"也依赖此逻辑兜底)
Stop();
// 释放旧运行产品的所有插件资源(防止切换后旧插件句柄/连接残留)
var oldTabs = RunTabs;
RunTabs = new List();
ProductManager.DisposeFlowTabs(oldTabs);
ProductManager.Instance.LoadProductToRun(productName);
RunProductName = productName;
RefreshRunTabs();
RunStatistics.Instance.Log($"运行产品已加载: {productName}({RunTabs.Count} 个流程)");
ProductLoaded?.Invoke(productName);
return true;
}
///
/// 重新解析运行态流程:始终从磁盘加载独立副本,编辑态与运行态彻底分离。
/// 监控端(FlowEditorShellViewModel.IsMonitorMode=true)通过订阅 RunTabsRefreshed
/// 主动挂接 RunTabs 显示运行状态,不再隐式共享编辑器 FlowTabs。
///
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);
}
///
/// 仅当产品流程文件发生变化(或尚未加载)时才从磁盘重载运行态流程;否则复用内存中的 RunTabs。
/// 用于高频触发(心跳/事件)路径,避免每次触发都读盘+反序列化+重建插件实例。
///
private void RefreshRunTabsIfChanged()
{
var stamp = ProductManager.Instance.GetProductFlowStamp(RunProductName);
if (RunTabs != null && RunTabs.Count > 0 && stamp == _runTabsStamp)
return;
RefreshRunTabs();
}
///
/// 执行一次:并行执行所有"主流程触发"的流程,全部结束后统计结果。
///
public async Task 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;
}
///
/// 循环运行:每个"主流程触发"的流程在独立的循环 Task 中运行,互不等待。
/// 流程 A 慢、流程 B 快时,B 不会被 A 拖住等待;任何一个流程本轮完成立即进入下一轮。
/// 直到 Stop() 被调用(取消 token)所有循环退出。
///
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();
}
}
///
/// 单个流程的独立循环:完成一轮立即进入下一轮,不等其他流程。
///
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; }
}
}
///
/// 停止循环运行 / 单次运行。
/// 立即解除"运行中"状态(编辑器解锁、按钮复位);后台循环在完成本轮后自行结束。
///
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);
}
///
/// 按名称执行一个运行态流程(供事件触发等外部场景调用)。
///
public async Task 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 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();
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;
}
}
///
/// 在 UI 线程同步编辑器 ViewModel 的 IsRunning,触发 RunFlowCommand 等按钮状态刷新。
///
private static void SetEditorRunning(FlowEditorViewModel editorVm, bool running)
{
if (editorVm == null) return;
System.Windows.Application.Current?.Dispatcher.Invoke(() =>
{
editorVm.IsRunning = running;
});
}
#region 心跳触发
///
/// 根据全局事件配置订阅通讯设备的数据到达事件,作为产品运行触发器。
/// 同一设备按 IntervalMs 节流,防止高频数据把执行引擎打爆。
///
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