using Opc.Ua; using Opc.Ua.Client; using OpcUaHelper; using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading; using System.Threading.Tasks; using TeamAAS_VP.Enums; using TeamAAS_VP.Interfaces; using TeamAAS_VP.Models; namespace TeamAAS_VP.Core.PLCs { /// /// OPC UA 节点订阅模式。 /// public enum OpcUaSubscriptionMode { /// /// 使用 OPC UA 原生订阅机制,由服务端主动推送节点变更。 /// Native, /// /// 使用客户端轮询方式定时读取节点值,并在检测到变化时触发通知。 /// CustomPolling } /// /// OPC UA 订阅配置选项,用于控制订阅模式、轮询间隔以及首次扫描通知行为。 /// public class OpcUaSubscriptionOptions { /// /// 获取或设置订阅模式。 /// 默认为 。 /// public OpcUaSubscriptionMode Mode { get; set; } = OpcUaSubscriptionMode.Native; /// /// 获取或设置轮询间隔,单位为毫秒。 /// 仅在 时生效。 /// 默认值为 50。 /// public int PollingInterval { get; set; } = 50; /// /// 获取或设置首次扫描时是否立即通知当前值。 /// 仅在自定义轮询订阅模式下生效。 /// public bool NotifyOnFirstScan { get; set; } } public class OpcUaClientPLC : IPlc { private readonly object _pollingSubscriptionLock = new object(); private readonly Dictionary _pollingSubscriptions = new Dictionary(); public Guid Id{ get; private set; } public CommunicationType CommunicationType { get; private set; } public string EndpointUrl { get; private set; } public int Index { get; set; } public string Name { get; private set; } public OpcUaClient OpcUaClient { get; private set; } public bool IsConnected { get { return OpcUaClient.Connected; } } /// /// 节点头 /// public string NodeHeader { get; private set; } /// /// 获取或设置默认的节点订阅模式。 /// 在未显式传入订阅选项时,`SubscribeNodes` 将使用该模式决定采用 OPC UA 原生订阅还是客户端轮询订阅。 /// 默认值为 。 /// public OpcUaSubscriptionMode DefaultSubscriptionMode { get; set; } = OpcUaSubscriptionMode.Native; /// /// 获取或设置默认的订阅轮询间隔,单位为毫秒。 /// 仅当默认订阅模式或实际订阅选项使用自定义轮询时生效。 /// 默认值为 `100`。 /// public int DefaultSubscriptionPollingInterval { get; set; } = 50; /// /// 获取或设置在默认轮询订阅模式下,首次扫描节点时是否立即触发一次数据变更通知。 /// public bool DefaultNotifyOnFirstScan { get; set; } /// /// PLC 连接状态变化事件。 /// 当 OPC UA 客户端完成连接、开始重连或重连完成时触发,第二个参数表示当前连接状态。 /// public event Action ConnectChangedEvent; /// /// 轮询订阅上下文。 /// 用于保存某个自定义轮询订阅任务的节点集合、上次读取值、回调处理器以及取消控制信息。 /// private class PollingSubscriptionContext { /// /// 订阅唯一标识,用于区分不同的订阅任务。 /// public string Key { get; set; } /// /// 当前订阅需要轮询的节点集合。 /// public List NodeIds { get; set; } /// /// 记录每个节点上一次读取到的值,用于比较是否发生变化。 /// public Dictionary LastValues { get; } = new Dictionary(); /// /// 节点值发生变化时执行的回调委托。 /// public Action<(string key, string nodeId, object value)> DataChangeHandler { get; set; } /// /// 用于控制轮询任务取消的令牌源。 /// public CancellationTokenSource CancellationTokenSource { get; set; } /// /// 轮询周期,单位为毫秒。 /// public int PollingInterval { get; set; } /// /// 指示首次扫描时是否立即上报当前节点值。 /// public bool NotifyOnFirstScan { get; set; } } public OpcUaClientPLC(Guid id,int index, string name, string endpointUrl,string nodeHeader, CommunicationType communicationType) { Id = id; Index = index; Name = name; EndpointUrl = endpointUrl; NodeHeader = nodeHeader; CommunicationType = communicationType; OpcUaClient = new OpcUaClient(); // 初始化PLC连接 OpcUaClient.UserIdentity = new UserIdentity(new AnonymousIdentityToken()); // connect to server, this is a sample try { OpcUaClient.ConnectComplete += OpcUaClient_ConnectComplete; OpcUaClient.ReconnectStarting += OpcUaClient_ReconnectStarting; OpcUaClient.ReconnectComplete += OpcUaClient_ReconnectComplete; OpcUaClient.OpcStatusChange += OpcUaClient_OpcStatusChange; } catch (Exception) { return; } } public OpcUaClientPLC(PlcInfo plcInfo) { Id = plcInfo.Id; Index=plcInfo.Index; Name = plcInfo.Name; NodeHeader=plcInfo.NodeHeader; EndpointUrl = plcInfo.EndpointUrl; CommunicationType = plcInfo.CommunicationType; OpcUaClient = new OpcUaClient(); // 初始化PLC连接 OpcUaClient.UserIdentity = new UserIdentity(new AnonymousIdentityToken()); // connect to server, this is a sample try { OpcUaClient.ConnectComplete += OpcUaClient_ConnectComplete; OpcUaClient.ReconnectStarting += OpcUaClient_ReconnectStarting; OpcUaClient.ReconnectComplete += OpcUaClient_ReconnectComplete; OpcUaClient.OpcStatusChange += OpcUaClient_OpcStatusChange; } catch (Exception) { return; } } public void Connect() { OpcUaClient.ConnectServer(EndpointUrl); } public Task ConnectAsync() { return OpcUaClient.ConnectServer(EndpointUrl); } public void Disconnect() { OpcUaClient.Disconnect(); } public void Dispose() { StopAllPollingSubscriptions(); OpcUaClient.Disconnect(); OpcUaClient.ConnectComplete -= OpcUaClient_ConnectComplete; OpcUaClient.ReconnectStarting -= OpcUaClient_ReconnectStarting; OpcUaClient.ReconnectComplete -= OpcUaClient_ReconnectComplete; OpcUaClient.OpcStatusChange -= OpcUaClient_OpcStatusChange; } /// /// 读取多个节点的值 /// /// /// public Dictionary ReadNodes(string[] nodeIds) { var result = new Dictionary(); var readNodeIds = nodeIds.Select(s => { if (s.StartsWith(NodeHeader)) { return s; } else { return NodeHeader + s; } }).ToArray(); List readNodeIdList = new List(); foreach (var readNodeId in readNodeIds) { readNodeIdList.Add(new NodeId(readNodeId)); } var values = OpcUaClient.ReadNodes(readNodeIdList.ToArray()); for (int i = 0; i < nodeIds.Length; i++) { result[nodeIds[i]] = values[i].Value; } return result; } /// /// 异步读取多个节点的值 /// /// /// public async Task> ReadNodesAsync(string[] nodeIds) { var result = new Dictionary(); var readNodeIds = nodeIds.Select(s => { if (s.StartsWith(NodeHeader)) { return s; } else { return NodeHeader + s; } }).ToArray(); List readNodeIdList = new List(); foreach (var readNodeId in readNodeIds) { readNodeIdList.Add(new NodeId(readNodeId)); } var values = await OpcUaClient.ReadNodesAsync(readNodeIdList.ToArray()); for (int i = 0; i < nodeIds.Length; i++) { result[nodeIds[i]] = values[i].Value; } return result; } /// /// 读取单个节点的值 /// /// /// public object ReadNode(string nodeId) { string readNodeId = nodeId; if (!nodeId.StartsWith(NodeHeader)) { readNodeId = NodeHeader + nodeId; } var value = OpcUaClient.ReadNode(new NodeId(readNodeId)); return value.Value; } /// /// 异步读取单个节点的值 /// /// /// public Task ReadNodeAsync(string nodeId) { return Task.Run(() => { return ReadNode(nodeId); }); } /// /// 读取单个节点的值 /// /// /// /// public T ReadNode(string nodeId) { string readNodeId = nodeId; if (!nodeId.StartsWith(NodeHeader)) { readNodeId = NodeHeader + nodeId; } return OpcUaClient.ReadNode(readNodeId); } /// /// 异步读取单个节点的值 /// /// /// /// public async Task ReadNodeAsync(string nodeId) { string readNodeId = nodeId; if (!nodeId.StartsWith(NodeHeader)) { readNodeId = NodeHeader + nodeId; } return await OpcUaClient.ReadNodeAsync(readNodeId); } /// /// 写入单个节点的值 /// /// /// /// /// public bool WriteNode(string nodeId, T value) { string writeNodeId = nodeId; if (!nodeId.StartsWith(NodeHeader)) { writeNodeId = NodeHeader + nodeId; } return OpcUaClient.WriteNode(writeNodeId, value); } /// /// 异步写入单个节点的值 /// /// /// /// /// public async Task WriteNodeAsync(string nodeId, T value) { string writeNodeId = nodeId; if (!nodeId.StartsWith(NodeHeader)) { writeNodeId = NodeHeader + nodeId; } return await OpcUaClient.WriteNodeAsync(writeNodeId, value); } /// /// 写入多个节点 /// /// /// public bool WriteNodes(Dictionary nodeValues) { var writeNodeValues = new Dictionary(); foreach (var kvp in nodeValues) { string writeNodeId = kvp.Key; if (!kvp.Key.StartsWith(NodeHeader)) { writeNodeId = NodeHeader + kvp.Key; } writeNodeValues[writeNodeId] = kvp.Value; } return OpcUaClient.WriteNodes(writeNodeValues.Keys.ToArray(), writeNodeValues.Values.ToArray()); } /// /// 写入多个节点异步 /// /// /// public Task WriteNodesAsync(Dictionary nodeValues) { return Task.Run(() => { return WriteNodes(nodeValues); }); } /// /// 使用当前实例的默认订阅配置订阅多个节点。 /// /// 订阅唯一标识,用于区分和管理同一客户端上的不同订阅任务。 /// 需要订阅的节点标识集合,可传入带或不带 前缀的节点名。 /// 节点值变化时触发的回调,返回订阅键、节点标识和值。 public void SubscribeNodes(string key, List nodeIds, Action<(string key,string nodeId,object value)> dataChangeHandler) { SubscribeNodes(key, nodeIds, dataChangeHandler, new OpcUaSubscriptionOptions { Mode = DefaultSubscriptionMode, PollingInterval = DefaultSubscriptionPollingInterval, NotifyOnFirstScan = DefaultNotifyOnFirstScan }); } /// /// 根据传入的OpcUaSubscriptionMode选项订阅多个节点 /// /// /// /// /// public void SubscribeNodes(string key, List nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionMode mode) { var options = new OpcUaSubscriptionOptions { Mode = mode, PollingInterval = DefaultSubscriptionPollingInterval, NotifyOnFirstScan = DefaultNotifyOnFirstScan }; SubscribeNodes(key, nodeIds, dataChangeHandler, options); } /// /// 根据指定订阅选项订阅多个节点。 /// /// /// 该方法会先停止同一 对应的轮询订阅,再根据 中的模式 /// 选择使用 OPC UA 原生订阅或客户端轮询订阅。节点列表会自动去重并忽略空白项。 /// /// 订阅唯一标识,不能为空或仅包含空白字符。 /// 待订阅的节点集合。 /// 节点值变化通知回调。 /// 订阅配置;为 null 时使用默认配置。 /// 为空或仅包含空白字符时抛出。 /// null 时抛出。 public void SubscribeNodes(string key, List nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionOptions options) { if (string.IsNullOrWhiteSpace(key)) { throw new ArgumentException("订阅Key不能为空。", nameof(key)); } if (nodeIds == null) { throw new ArgumentNullException(nameof(nodeIds)); } if (dataChangeHandler == null) { throw new ArgumentNullException(nameof(dataChangeHandler)); } var subscriptionOptions = options ?? new OpcUaSubscriptionOptions(); var normalizedNodeIds = nodeIds.Where(s => !string.IsNullOrWhiteSpace(s)).Distinct().ToList(); StopPollingSubscription(key); if (subscriptionOptions.Mode == OpcUaSubscriptionMode.CustomPolling) { StartPollingSubscription(key, normalizedNodeIds, dataChangeHandler, subscriptionOptions); return; } SubscribeNodesByOpc(key, normalizedNodeIds, dataChangeHandler); } /// /// 使用 OPC UA 原生订阅机制订阅节点。 /// /// /// 订阅前会自动补齐节点前缀;收到服务端推送后,会移除返回节点中的 前缀, /// 再将变化值通过回调透出给业务层。 /// /// 订阅唯一标识。 /// 节点集合。 /// 节点值变化回调。 private void SubscribeNodesByOpc(string key, List nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler) { OpcUaClient.AddSubscription(key, nodeIds.Select(s => { return NormalizeNodeId(s); }).ToArray(), (key1, monitoredItem, args) => { MonitoredItemNotification notification = args.NotificationValue as MonitoredItemNotification; if (notification == null) { return; } string nodeId = monitoredItem.StartNodeId.ToString(); nodeId = TrimNodeHeader(nodeId); // 触发数据变化事件 dataChangeHandler?.Invoke((key1, nodeId, notification.Value.WrappedValue.Value)); }); } /// /// 启动自定义轮询订阅任务。 /// /// /// 该方法会创建轮询上下文并登记到内部订阅字典中,然后在后台任务中执行轮询循环。 /// 轮询间隔最小限制为 10 毫秒,以避免过于频繁的读取请求。 /// /// 订阅唯一标识。 /// 需要轮询的节点集合。 /// 节点值变化回调。 /// 轮询订阅选项。 private void StartPollingSubscription(string key, List nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionOptions options) { var context = new PollingSubscriptionContext { Key = key, NodeIds = nodeIds, DataChangeHandler = dataChangeHandler, CancellationTokenSource = new CancellationTokenSource(), PollingInterval = Math.Max(10, options.PollingInterval), NotifyOnFirstScan = options.NotifyOnFirstScan }; lock (_pollingSubscriptionLock) { _pollingSubscriptions[key] = context; } Task.Run(() => PollingSubscriptionLoopAsync(context)); } /// /// 执行轮询订阅的后台循环。 /// /// /// 当客户端已连接且存在待轮询节点时,循环会批量读取节点值,并与上次缓存值进行比较: /// 首次读取时根据配置决定是否立即通知,后续仅在值发生变化时触发回调。 /// 当取消令牌触发或发生取消异常时,循环退出。 /// /// 轮询订阅上下文。 /// 表示异步轮询操作的任务。 private async Task PollingSubscriptionLoopAsync(PollingSubscriptionContext context) { while (!context.CancellationTokenSource.IsCancellationRequested) { try { if (IsConnected && context.NodeIds.Count > 0) { var values = await ReadNodesAsync(context.NodeIds.ToArray()).ConfigureAwait(false); foreach (var nodeId in context.NodeIds) { object currentValue; if (!values.TryGetValue(nodeId, out currentValue)) { continue; } object lastValue; if (!context.LastValues.TryGetValue(nodeId, out lastValue)) { context.LastValues[nodeId] = currentValue; if (context.NotifyOnFirstScan) { context.DataChangeHandler?.Invoke((context.Key, nodeId, currentValue)); } continue; } if (!Utils.IsEqual(lastValue, currentValue)) { context.LastValues[nodeId] = currentValue; context.DataChangeHandler?.Invoke((context.Key, nodeId, currentValue)); } } } } catch (OperationCanceledException) { break; } catch (Exception) { } try { await Task.Delay(context.PollingInterval, context.CancellationTokenSource.Token).ConfigureAwait(false); } catch (OperationCanceledException) { break; } } } /// /// 停止指定键对应的轮询订阅。 /// /// 要停止的订阅唯一标识。 private void StopPollingSubscription(string key) { PollingSubscriptionContext context = null; lock (_pollingSubscriptionLock) { if (_pollingSubscriptions.TryGetValue(key, out context)) { _pollingSubscriptions.Remove(key); } } context?.CancellationTokenSource.Cancel(); context?.CancellationTokenSource.Dispose(); } /// /// 停止并清理当前实例中的所有轮询订阅。 /// private void StopAllPollingSubscriptions() { List contexts; lock (_pollingSubscriptionLock) { contexts = _pollingSubscriptions.Values.ToList(); _pollingSubscriptions.Clear(); } foreach (var context in contexts) { context.CancellationTokenSource.Cancel(); context.CancellationTokenSource.Dispose(); } } /// /// 规范化节点标识,确保其包含节点头前缀。 /// /// 原始节点标识。 /// 包含 前缀的完整节点标识。 private string NormalizeNodeId(string nodeId) { if (nodeId.StartsWith(NodeHeader)) { return nodeId; } return NodeHeader + nodeId; } /// /// 去除节点标识中的节点头前缀。 /// /// 完整节点标识。 /// 移除 前缀后的节点标识;若原值不包含此前缀则直接返回原值。 private string TrimNodeHeader(string nodeId) { if (nodeId.StartsWith(NodeHeader)) { return nodeId.Substring(NodeHeader.Length); } return nodeId; } #region OPC事件 private void OpcUaClient_OpcStatusChange(object sender, OpcUaStatusEventArgs e) { //LogHelper.WriteLogInfo($"OPC状态发生改变:Error:{e.Error},Time:{e.Time},Text:{e.Text}"); } private void OpcUaClient_ReconnectComplete(object sender, EventArgs e) { ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected); } private void OpcUaClient_ReconnectStarting(object sender, EventArgs e) { ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected); } private void OpcUaClient_ConnectComplete(object sender, EventArgs e) { ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected); } #endregion } }