OPCuaClientPLC.cs 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754
  1. using Opc.Ua;
  2. using Opc.Ua.Client;
  3. using OpcUaHelper;
  4. using System;
  5. using System.Collections.Generic;
  6. using System.Diagnostics;
  7. using System.Linq;
  8. using System.Text;
  9. using System.Threading;
  10. using System.Threading.Tasks;
  11. using TeamAAS_VP.Enums;
  12. using TeamAAS_VP.Interfaces;
  13. using TeamAAS_VP.Models;
  14. namespace TeamAAS_VP.Core.PLCs
  15. {
  16. /// <summary>
  17. /// OPC UA 节点订阅模式。
  18. /// </summary>
  19. public enum OpcUaSubscriptionMode
  20. {
  21. /// <summary>
  22. /// 使用 OPC UA 原生订阅机制,由服务端主动推送节点变更。
  23. /// </summary>
  24. Native,
  25. /// <summary>
  26. /// 使用客户端轮询方式定时读取节点值,并在检测到变化时触发通知。
  27. /// </summary>
  28. CustomPolling
  29. }
  30. /// <summary>
  31. /// OPC UA 订阅配置选项,用于控制订阅模式、轮询间隔以及首次扫描通知行为。
  32. /// </summary>
  33. public class OpcUaSubscriptionOptions
  34. {
  35. /// <summary>
  36. /// 获取或设置订阅模式。
  37. /// 默认为 <see cref="OpcUaSubscriptionMode.Native"/>。
  38. /// </summary>
  39. public OpcUaSubscriptionMode Mode { get; set; } = OpcUaSubscriptionMode.Native;
  40. /// <summary>
  41. /// 获取或设置轮询间隔,单位为毫秒。
  42. /// 仅在 <see cref="Mode"/> 为 <see cref="OpcUaSubscriptionMode.CustomPolling"/> 时生效。
  43. /// 默认值为 50。
  44. /// </summary>
  45. public int PollingInterval { get; set; } = 50;
  46. /// <summary>
  47. /// 获取或设置首次扫描时是否立即通知当前值。
  48. /// 仅在自定义轮询订阅模式下生效。
  49. /// </summary>
  50. public bool NotifyOnFirstScan { get; set; }
  51. }
  52. public class OpcUaClientPLC : IPlc
  53. {
  54. private readonly object _pollingSubscriptionLock = new object();
  55. private readonly Dictionary<string, PollingSubscriptionContext> _pollingSubscriptions = new Dictionary<string, PollingSubscriptionContext>();
  56. public Guid Id { get; private set; }
  57. public CommunicationType CommunicationType { get; private set; }
  58. public string EndpointUrl { get; private set; }
  59. public int Index { get; set; }
  60. public string Name { get; private set; }
  61. public OpcUaClient OpcUaClient { get; private set; }
  62. public bool IsConnected { get { return OpcUaClient.Connected; } }
  63. /// <summary>
  64. /// 节点头
  65. /// </summary>
  66. public string NodeHeader { get; private set; }
  67. /// <summary>
  68. /// 获取或设置默认的节点订阅模式。
  69. /// 在未显式传入订阅选项时,`SubscribeNodes` 将使用该模式决定采用 OPC UA 原生订阅还是客户端轮询订阅。
  70. /// 默认值为 <see cref="OpcUaSubscriptionMode.Native"/>。
  71. /// </summary>
  72. public OpcUaSubscriptionMode DefaultSubscriptionMode { get; set; } = OpcUaSubscriptionMode.Native;
  73. /// <summary>
  74. /// 获取或设置默认的订阅轮询间隔,单位为毫秒。
  75. /// 仅当默认订阅模式或实际订阅选项使用自定义轮询时生效。
  76. /// 默认值为 `100`。
  77. /// </summary>
  78. public int DefaultSubscriptionPollingInterval { get; set; } = 100;
  79. /// <summary>
  80. /// 获取或设置在默认轮询订阅模式下,首次扫描节点时是否立即触发一次数据变更通知。
  81. /// </summary>
  82. public bool DefaultNotifyOnFirstScan { get; set; }
  83. /// <summary>
  84. /// PLC 连接状态变化事件。
  85. /// 当 OPC UA 客户端完成连接、开始重连或重连完成时触发,第二个参数表示当前连接状态。
  86. /// </summary>
  87. public event Action<object, bool> ConnectChangedEvent;
  88. /// <summary>
  89. /// 轮询订阅上下文。
  90. /// 用于保存某个自定义轮询订阅任务的节点集合、上次读取值、回调处理器以及取消控制信息。
  91. /// </summary>
  92. private class PollingSubscriptionContext
  93. {
  94. /// <summary>
  95. /// 订阅唯一标识,用于区分不同的订阅任务。
  96. /// </summary>
  97. public string Key { get; set; }
  98. /// <summary>
  99. /// 当前订阅需要轮询的节点集合。
  100. /// </summary>
  101. public List<string> NodeIds { get; set; }
  102. /// <summary>
  103. /// 记录每个节点上一次读取到的值,用于比较是否发生变化。
  104. /// </summary>
  105. public Dictionary<string, object> LastValues { get; } = new Dictionary<string, object>();
  106. /// <summary>
  107. /// 节点值发生变化时执行的回调委托。
  108. /// </summary>
  109. public Action<(string key, string nodeId, object value)> DataChangeHandler { get; set; }
  110. /// <summary>
  111. /// 用于控制轮询任务取消的令牌源。
  112. /// </summary>
  113. public CancellationTokenSource CancellationTokenSource { get; set; }
  114. /// <summary>
  115. /// 轮询周期,单位为毫秒。
  116. /// </summary>
  117. public int PollingInterval { get; set; }
  118. /// <summary>
  119. /// 指示首次扫描时是否立即上报当前节点值。
  120. /// </summary>
  121. public bool NotifyOnFirstScan { get; set; }
  122. }
  123. public OpcUaClientPLC(Guid id, int index, string name, string endpointUrl, string nodeHeader, CommunicationType communicationType)
  124. {
  125. Id = id;
  126. Index = index;
  127. Name = name;
  128. EndpointUrl = endpointUrl;
  129. NodeHeader = nodeHeader;
  130. CommunicationType = communicationType;
  131. OpcUaClient = new OpcUaClient();
  132. // 初始化PLC连接
  133. OpcUaClient.UserIdentity = new UserIdentity(new AnonymousIdentityToken());
  134. // connect to server, this is a sample
  135. try
  136. {
  137. OpcUaClient.ConnectComplete += OpcUaClient_ConnectComplete;
  138. OpcUaClient.ReconnectStarting += OpcUaClient_ReconnectStarting;
  139. OpcUaClient.ReconnectComplete += OpcUaClient_ReconnectComplete;
  140. OpcUaClient.OpcStatusChange += OpcUaClient_OpcStatusChange;
  141. }
  142. catch (Exception)
  143. {
  144. return;
  145. }
  146. }
  147. public OpcUaClientPLC(PlcInfo plcInfo)
  148. {
  149. Id = plcInfo.Id;
  150. Index = plcInfo.Index;
  151. Name = plcInfo.Name;
  152. NodeHeader = plcInfo.NodeHeader;
  153. EndpointUrl = plcInfo.EndpointUrl;
  154. CommunicationType = plcInfo.CommunicationType;
  155. OpcUaClient = new OpcUaClient();
  156. // 初始化PLC连接
  157. OpcUaClient.UserIdentity = new UserIdentity(new AnonymousIdentityToken());
  158. // connect to server, this is a sample
  159. try
  160. {
  161. OpcUaClient.ConnectComplete += OpcUaClient_ConnectComplete;
  162. OpcUaClient.ReconnectStarting += OpcUaClient_ReconnectStarting;
  163. OpcUaClient.ReconnectComplete += OpcUaClient_ReconnectComplete;
  164. OpcUaClient.OpcStatusChange += OpcUaClient_OpcStatusChange;
  165. }
  166. catch (Exception)
  167. {
  168. return;
  169. }
  170. }
  171. public void Connect()
  172. {
  173. OpcUaClient.ConnectServer(EndpointUrl);
  174. }
  175. public Task ConnectAsync()
  176. {
  177. return OpcUaClient.ConnectServer(EndpointUrl);
  178. }
  179. public void Disconnect()
  180. {
  181. OpcUaClient.Disconnect();
  182. }
  183. public void Dispose()
  184. {
  185. StopAllPollingSubscriptions();
  186. OpcUaClient.Disconnect();
  187. OpcUaClient.ConnectComplete -= OpcUaClient_ConnectComplete;
  188. OpcUaClient.ReconnectStarting -= OpcUaClient_ReconnectStarting;
  189. OpcUaClient.ReconnectComplete -= OpcUaClient_ReconnectComplete;
  190. OpcUaClient.OpcStatusChange -= OpcUaClient_OpcStatusChange;
  191. }
  192. /// <summary>
  193. /// 读取多个节点的值
  194. /// </summary>
  195. /// <param name="nodeIds"></param>
  196. /// <returns></returns>
  197. public Dictionary<string, object> ReadNodes(string[] nodeIds)
  198. {
  199. var result = new Dictionary<string, object>();
  200. var readNodeIds = nodeIds.Select(s =>
  201. {
  202. if (s.StartsWith(NodeHeader))
  203. {
  204. return s;
  205. }
  206. else
  207. {
  208. return NodeHeader + s;
  209. }
  210. }).ToArray();
  211. List<NodeId> readNodeIdList = new List<NodeId>();
  212. foreach (var readNodeId in readNodeIds)
  213. {
  214. readNodeIdList.Add(new NodeId(readNodeId));
  215. }
  216. var values = OpcUaClient.ReadNodes(readNodeIdList.ToArray());
  217. for (int i = 0; i < nodeIds.Length; i++)
  218. {
  219. result[nodeIds[i]] = values[i].Value;
  220. }
  221. return result;
  222. }
  223. /// <summary>
  224. /// 异步读取多个节点的值
  225. /// </summary>
  226. /// <param name="nodeIds"></param>
  227. /// <returns></returns>
  228. public async Task<Dictionary<string, object>> ReadNodesAsync(string[] nodeIds)
  229. {
  230. //var readWatch = Stopwatch.StartNew();
  231. var result = new Dictionary<string, object>();
  232. var readNodeIds = nodeIds.Select(s =>
  233. {
  234. if (s.StartsWith(NodeHeader))
  235. {
  236. return s;
  237. }
  238. else
  239. {
  240. return NodeHeader + s;
  241. }
  242. }).ToArray();
  243. List<NodeId> readNodeIdList = new List<NodeId>();
  244. foreach (var readNodeId in readNodeIds)
  245. {
  246. readNodeIdList.Add(new NodeId(readNodeId));
  247. }
  248. var values = await OpcUaClient.ReadNodesAsync(readNodeIdList.ToArray());
  249. // readWatch.Stop();
  250. //if (readWatch.ElapsedMilliseconds > Math.Max(50, DefaultSubscriptionPollingInterval))
  251. //{
  252. // LogHelper.WriteLogInfo($"[PLC-TRACE] ReadNodesAsync slow, plc={Name}, nodeCount={nodeIds.Length}, elapsedMs={readWatch.ElapsedMilliseconds}, thread={Thread.CurrentThread.ManagedThreadId}");
  253. //}
  254. for (int i = 0; i < nodeIds.Length; i++)
  255. {
  256. result[nodeIds[i]] = values[i].Value;
  257. }
  258. return result;
  259. }
  260. /// <summary>
  261. /// 读取单个节点的值
  262. /// </summary>
  263. /// <param name="nodeId"></param>
  264. /// <returns></returns>
  265. public object ReadNode(string nodeId)
  266. {
  267. string readNodeId = nodeId;
  268. if (!nodeId.StartsWith(NodeHeader))
  269. {
  270. readNodeId = NodeHeader + nodeId;
  271. }
  272. var value = OpcUaClient.ReadNode(new NodeId(readNodeId));
  273. return value.Value;
  274. }
  275. /// <summary>
  276. /// 异步读取单个节点的值
  277. /// </summary>
  278. /// <param name="nodeId"></param>
  279. /// <returns></returns>
  280. public Task<object> ReadNodeAsync(string nodeId)
  281. {
  282. return Task.Run(() =>
  283. {
  284. return ReadNode(nodeId);
  285. });
  286. }
  287. /// <summary>
  288. /// 读取单个节点的值
  289. /// </summary>
  290. /// <typeparam name="T"></typeparam>
  291. /// <param name="nodeId"></param>
  292. /// <returns></returns>
  293. public T ReadNode<T>(string nodeId)
  294. {
  295. string readNodeId = nodeId;
  296. if (!nodeId.StartsWith(NodeHeader))
  297. {
  298. readNodeId = NodeHeader + nodeId;
  299. }
  300. return OpcUaClient.ReadNode<T>(readNodeId);
  301. }
  302. /// <summary>
  303. /// 异步读取单个节点的值
  304. /// </summary>
  305. /// <typeparam name="T"></typeparam>
  306. /// <param name="nodeId"></param>
  307. /// <returns></returns>
  308. public async Task<T> ReadNodeAsync<T>(string nodeId)
  309. {
  310. string readNodeId = nodeId;
  311. if (!nodeId.StartsWith(NodeHeader))
  312. {
  313. readNodeId = NodeHeader + nodeId;
  314. }
  315. return await OpcUaClient.ReadNodeAsync<T>(readNodeId);
  316. }
  317. /// <summary>
  318. /// 写入单个节点的值
  319. /// </summary>
  320. /// <typeparam name="T"></typeparam>
  321. /// <param name="nodeId"></param>
  322. /// <param name="value"></param>
  323. /// <returns></returns>
  324. public bool WriteNode<T>(string nodeId, T value)
  325. {
  326. string writeNodeId = nodeId;
  327. if (!nodeId.StartsWith(NodeHeader))
  328. {
  329. writeNodeId = NodeHeader + nodeId;
  330. }
  331. return OpcUaClient.WriteNode<T>(writeNodeId, value);
  332. }
  333. /// <summary>
  334. /// 异步写入单个节点的值
  335. /// </summary>
  336. /// <typeparam name="T"></typeparam>
  337. /// <param name="nodeId"></param>
  338. /// <param name="value"></param>
  339. /// <returns></returns>
  340. public async Task<bool> WriteNodeAsync<T>(string nodeId, T value)
  341. {
  342. //var writeWatch = Stopwatch.StartNew();
  343. string writeNodeId = nodeId;
  344. if (!nodeId.StartsWith(NodeHeader))
  345. {
  346. writeNodeId = NodeHeader + nodeId;
  347. }
  348. // LogHelper.WriteLogInfo($"[PLC-TRACE] WriteNodeAsync start, plc={Name}, node={nodeId}, fullNode={writeNodeId}, value={value}, thread={Thread.CurrentThread.ManagedThreadId}, time={DateTime.Now:HH:mm:ss.fff}");
  349. try
  350. {
  351. bool result = await OpcUaClient.WriteNodeAsync<T>(writeNodeId, value).ConfigureAwait(false);
  352. //writeWatch.Stop();
  353. //LogHelper.WriteLogInfo($"[PLC-TRACE] WriteNodeAsync done, plc={Name}, node={nodeId}, value={value}, result={result}, elapsedMs={writeWatch.ElapsedMilliseconds}, thread={Thread.CurrentThread.ManagedThreadId}, time={DateTime.Now:HH:mm:ss.fff}");
  354. return result;
  355. }
  356. catch (Exception ex)
  357. {
  358. //writeWatch.Stop();
  359. //LogHelper.WriteLogError($"[PLC-TRACE] WriteNodeAsync failed, plc={Name}, node={nodeId}, value={value}, elapsedMs={writeWatch.ElapsedMilliseconds}", ex);
  360. throw;
  361. }
  362. }
  363. /// <summary>
  364. /// 写入多个节点
  365. /// </summary>
  366. /// <param name="nodeValues"></param>
  367. /// <returns></returns>
  368. public bool WriteNodes(Dictionary<string, object> nodeValues)
  369. {
  370. var writeNodeValues = new Dictionary<string, object>();
  371. foreach (var kvp in nodeValues)
  372. {
  373. string writeNodeId = kvp.Key;
  374. if (!kvp.Key.StartsWith(NodeHeader))
  375. {
  376. writeNodeId = NodeHeader + kvp.Key;
  377. }
  378. writeNodeValues[writeNodeId] = kvp.Value;
  379. }
  380. return OpcUaClient.WriteNodes(writeNodeValues.Keys.ToArray(), writeNodeValues.Values.ToArray());
  381. }
  382. /// <summary>
  383. /// 写入多个节点异步
  384. /// </summary>
  385. /// <param name="nodeValues"></param>
  386. /// <returns></returns>
  387. public Task<bool> WriteNodesAsync(Dictionary<string, object> nodeValues)
  388. {
  389. return Task.Run(() =>
  390. {
  391. return WriteNodes(nodeValues);
  392. });
  393. }
  394. /// <summary>
  395. /// 使用当前实例的默认订阅配置订阅多个节点。
  396. /// </summary>
  397. /// <param name="key">订阅唯一标识,用于区分和管理同一客户端上的不同订阅任务。</param>
  398. /// <param name="nodeIds">需要订阅的节点标识集合,可传入带或不带 <see cref="NodeHeader"/> 前缀的节点名。</param>
  399. /// <param name="dataChangeHandler">节点值变化时触发的回调,返回订阅键、节点标识和值。</param>
  400. public void SubscribeNodes(string key, List<string> nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler)
  401. {
  402. SubscribeNodes(key, nodeIds, dataChangeHandler, new OpcUaSubscriptionOptions
  403. {
  404. Mode = DefaultSubscriptionMode,
  405. PollingInterval = DefaultSubscriptionPollingInterval,
  406. NotifyOnFirstScan = DefaultNotifyOnFirstScan
  407. });
  408. }
  409. /// <summary>
  410. /// 根据传入的OpcUaSubscriptionMode选项订阅多个节点
  411. /// </summary>
  412. /// <param name="key"></param>
  413. /// <param name="nodeIds"></param>
  414. /// <param name="dataChangeHandler"></param>
  415. /// <param name="mode"></param>
  416. public void SubscribeNodes(string key, List<string> nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionMode mode)
  417. {
  418. var options = new OpcUaSubscriptionOptions
  419. {
  420. Mode = mode,
  421. PollingInterval = DefaultSubscriptionPollingInterval,
  422. NotifyOnFirstScan = DefaultNotifyOnFirstScan
  423. };
  424. SubscribeNodes(key, nodeIds, dataChangeHandler, options);
  425. }
  426. /// <summary>
  427. /// 根据指定订阅选项订阅多个节点。
  428. /// </summary>
  429. /// <remarks>
  430. /// 该方法会先停止同一 <paramref name="key"/> 对应的轮询订阅,再根据 <paramref name="options"/> 中的模式
  431. /// 选择使用 OPC UA 原生订阅或客户端轮询订阅。节点列表会自动去重并忽略空白项。
  432. /// </remarks>
  433. /// <param name="key">订阅唯一标识,不能为空或仅包含空白字符。</param>
  434. /// <param name="nodeIds">待订阅的节点集合。</param>
  435. /// <param name="dataChangeHandler">节点值变化通知回调。</param>
  436. /// <param name="options">订阅配置;为 <c>null</c> 时使用默认配置。</param>
  437. /// <exception cref="ArgumentException"><paramref name="key"/> 为空或仅包含空白字符时抛出。</exception>
  438. /// <exception cref="ArgumentNullException"><paramref name="nodeIds"/> 或 <paramref name="dataChangeHandler"/> 为 <c>null</c> 时抛出。</exception>
  439. public void SubscribeNodes(string key, List<string> nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionOptions options)
  440. {
  441. if (string.IsNullOrWhiteSpace(key))
  442. {
  443. throw new ArgumentException("订阅Key不能为空。", nameof(key));
  444. }
  445. if (nodeIds == null)
  446. {
  447. throw new ArgumentNullException(nameof(nodeIds));
  448. }
  449. if (dataChangeHandler == null)
  450. {
  451. throw new ArgumentNullException(nameof(dataChangeHandler));
  452. }
  453. var subscriptionOptions = options ?? new OpcUaSubscriptionOptions();
  454. var normalizedNodeIds = nodeIds.Where(s => !string.IsNullOrWhiteSpace(s)).Distinct().ToList();
  455. StopPollingSubscription(key);
  456. if (subscriptionOptions.Mode == OpcUaSubscriptionMode.CustomPolling)
  457. {
  458. StartPollingSubscription(key, normalizedNodeIds, dataChangeHandler, subscriptionOptions);
  459. return;
  460. }
  461. SubscribeNodesByOpc(key, normalizedNodeIds, dataChangeHandler);
  462. }
  463. /// <summary>
  464. /// 使用 OPC UA 原生订阅机制订阅节点。
  465. /// </summary>
  466. /// <remarks>
  467. /// 订阅前会自动补齐节点前缀;收到服务端推送后,会移除返回节点中的 <see cref="NodeHeader"/> 前缀,
  468. /// 再将变化值通过回调透出给业务层。
  469. /// </remarks>
  470. /// <param name="key">订阅唯一标识。</param>
  471. /// <param name="nodeIds">节点集合。</param>
  472. /// <param name="dataChangeHandler">节点值变化回调。</param>
  473. private void SubscribeNodesByOpc(string key, List<string> nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler)
  474. {
  475. OpcUaClient.AddSubscription(key, nodeIds.Select(s =>
  476. {
  477. return NormalizeNodeId(s);
  478. }).ToArray(), (key1, monitoredItem, args) =>
  479. {
  480. MonitoredItemNotification notification = args.NotificationValue as MonitoredItemNotification;
  481. if (notification == null)
  482. {
  483. return;
  484. }
  485. string nodeId = monitoredItem.StartNodeId.ToString();
  486. nodeId = TrimNodeHeader(nodeId);
  487. // 触发数据变化事件
  488. dataChangeHandler?.Invoke((key1, nodeId, notification.Value.WrappedValue.Value));
  489. });
  490. }
  491. /// <summary>
  492. /// 启动自定义轮询订阅任务。
  493. /// </summary>
  494. /// <remarks>
  495. /// 该方法会创建轮询上下文并登记到内部订阅字典中,然后在后台任务中执行轮询循环。
  496. /// 轮询间隔最小限制为 10 毫秒,以避免过于频繁的读取请求。
  497. /// </remarks>
  498. /// <param name="key">订阅唯一标识。</param>
  499. /// <param name="nodeIds">需要轮询的节点集合。</param>
  500. /// <param name="dataChangeHandler">节点值变化回调。</param>
  501. /// <param name="options">轮询订阅选项。</param>
  502. private void StartPollingSubscription(string key, List<string> nodeIds, Action<(string key, string nodeId, object value)> dataChangeHandler, OpcUaSubscriptionOptions options)
  503. {
  504. var context = new PollingSubscriptionContext
  505. {
  506. Key = key,
  507. NodeIds = nodeIds,
  508. DataChangeHandler = dataChangeHandler,
  509. CancellationTokenSource = new CancellationTokenSource(),
  510. PollingInterval = Math.Max(10, options.PollingInterval),
  511. NotifyOnFirstScan = options.NotifyOnFirstScan
  512. };
  513. lock (_pollingSubscriptionLock)
  514. {
  515. _pollingSubscriptions[key] = context;
  516. }
  517. Task.Run(() => PollingSubscriptionLoopAsync(context));
  518. }
  519. /// <summary>
  520. /// 执行轮询订阅的后台循环。
  521. /// </summary>
  522. /// <remarks>
  523. /// 当客户端已连接且存在待轮询节点时,循环会批量读取节点值,并与上次缓存值进行比较:
  524. /// 首次读取时根据配置决定是否立即通知,后续仅在值发生变化时触发回调。
  525. /// 当取消令牌触发或发生取消异常时,循环退出。
  526. /// </remarks>
  527. /// <param name="context">轮询订阅上下文。</param>
  528. /// <returns>表示异步轮询操作的任务。</returns>
  529. private async Task PollingSubscriptionLoopAsync(PollingSubscriptionContext context)
  530. {
  531. while (!context.CancellationTokenSource.IsCancellationRequested)
  532. {
  533. // var cycleWatch = Stopwatch.StartNew();
  534. int notifyCount = 0;
  535. try
  536. {
  537. if (IsConnected && context.NodeIds.Count > 0)
  538. {
  539. //var pollingReadWatch = Stopwatch.StartNew();
  540. var values = await ReadNodesAsync(context.NodeIds.ToArray()).ConfigureAwait(false);
  541. // pollingReadWatch.Stop();
  542. foreach (var nodeId in context.NodeIds)
  543. {
  544. object currentValue;
  545. if (!values.TryGetValue(nodeId, out currentValue))
  546. {
  547. continue;
  548. }
  549. object lastValue;
  550. if (!context.LastValues.TryGetValue(nodeId, out lastValue))
  551. {
  552. context.LastValues[nodeId] = currentValue;
  553. if (context.NotifyOnFirstScan)
  554. {
  555. // var handlerWatch = Stopwatch.StartNew();
  556. context.DataChangeHandler?.Invoke((context.Key, nodeId, currentValue));
  557. //handlerWatch.Stop();
  558. //notifyCount++;
  559. //if (handlerWatch.ElapsedMilliseconds > 20)
  560. //{
  561. // LogHelper.WriteLogInfo($"[PLC-TRACE] Polling handler slow(first), key={context.Key}, node={nodeId}, elapsedMs={handlerWatch.ElapsedMilliseconds}");
  562. //}
  563. }
  564. continue;
  565. }
  566. if (!Utils.IsEqual(lastValue, currentValue))
  567. {
  568. context.LastValues[nodeId] = currentValue;
  569. // var handlerWatch = Stopwatch.StartNew();
  570. context.DataChangeHandler?.Invoke((context.Key, nodeId, currentValue));
  571. //handlerWatch.Stop();
  572. //notifyCount++;
  573. //if (handlerWatch.ElapsedMilliseconds > 20)
  574. //{
  575. // LogHelper.WriteLogInfo($"[PLC-TRACE] Polling handler slow, key={context.Key}, node={nodeId}, elapsedMs={handlerWatch.ElapsedMilliseconds}, value={currentValue}");
  576. //}
  577. }
  578. }
  579. //if (pollingReadWatch.ElapsedMilliseconds > context.PollingInterval)
  580. //{
  581. // LogHelper.WriteLogInfo($"[PLC-TRACE] Polling read slower than interval, key={context.Key}, nodeCount={context.NodeIds.Count}, readMs={pollingReadWatch.ElapsedMilliseconds}, intervalMs={context.PollingInterval}");
  582. //}
  583. }
  584. }
  585. catch (OperationCanceledException)
  586. {
  587. break;
  588. }
  589. catch (Exception ex)
  590. {
  591. LogHelper.WriteLogError($"轮询订阅 [{context.Key}] 发生异常", ex);
  592. }
  593. try
  594. {
  595. //cycleWatch.Stop();
  596. //if (cycleWatch.ElapsedMilliseconds > context.PollingInterval || notifyCount > 0)
  597. //{
  598. // LogHelper.WriteLogInfo($"[PLC-TRACE] Polling cycle, key={context.Key}, nodeCount={context.NodeIds.Count}, notifyCount={notifyCount}, elapsedMs={cycleWatch.ElapsedMilliseconds}, intervalMs={context.PollingInterval}");
  599. //}
  600. await Task.Delay(context.PollingInterval, context.CancellationTokenSource.Token).ConfigureAwait(false);
  601. }
  602. catch (OperationCanceledException)
  603. {
  604. break;
  605. }
  606. }
  607. }
  608. /// <summary>
  609. /// 停止指定键对应的轮询订阅。
  610. /// </summary>
  611. /// <param name="key">要停止的订阅唯一标识。</param>
  612. private void StopPollingSubscription(string key)
  613. {
  614. PollingSubscriptionContext context = null;
  615. lock (_pollingSubscriptionLock)
  616. {
  617. if (_pollingSubscriptions.TryGetValue(key, out context))
  618. {
  619. _pollingSubscriptions.Remove(key);
  620. }
  621. }
  622. context?.CancellationTokenSource.Cancel();
  623. context?.CancellationTokenSource.Dispose();
  624. }
  625. /// <summary>
  626. /// 停止并清理当前实例中的所有轮询订阅。
  627. /// </summary>
  628. private void StopAllPollingSubscriptions()
  629. {
  630. List<PollingSubscriptionContext> contexts;
  631. lock (_pollingSubscriptionLock)
  632. {
  633. contexts = _pollingSubscriptions.Values.ToList();
  634. _pollingSubscriptions.Clear();
  635. }
  636. foreach (var context in contexts)
  637. {
  638. context.CancellationTokenSource.Cancel();
  639. context.CancellationTokenSource.Dispose();
  640. }
  641. }
  642. /// <summary>
  643. /// 规范化节点标识,确保其包含节点头前缀。
  644. /// </summary>
  645. /// <param name="nodeId">原始节点标识。</param>
  646. /// <returns>包含 <see cref="NodeHeader"/> 前缀的完整节点标识。</returns>
  647. private string NormalizeNodeId(string nodeId)
  648. {
  649. if (nodeId.StartsWith(NodeHeader))
  650. {
  651. return nodeId;
  652. }
  653. return NodeHeader + nodeId;
  654. }
  655. /// <summary>
  656. /// 去除节点标识中的节点头前缀。
  657. /// </summary>
  658. /// <param name="nodeId">完整节点标识。</param>
  659. /// <returns>移除 <see cref="NodeHeader"/> 前缀后的节点标识;若原值不包含此前缀则直接返回原值。</returns>
  660. private string TrimNodeHeader(string nodeId)
  661. {
  662. if (nodeId.StartsWith(NodeHeader))
  663. {
  664. return nodeId.Substring(NodeHeader.Length);
  665. }
  666. return nodeId;
  667. }
  668. #region OPC事件
  669. private void OpcUaClient_OpcStatusChange(object sender, OpcUaStatusEventArgs e)
  670. {
  671. //LogHelper.WriteLogInfo($"OPC状态发生改变:Error:{e.Error},Time:{e.Time},Text:{e.Text}");
  672. }
  673. private void OpcUaClient_ReconnectComplete(object sender, EventArgs e)
  674. {
  675. ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected);
  676. }
  677. private void OpcUaClient_ReconnectStarting(object sender, EventArgs e)
  678. {
  679. ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected);
  680. }
  681. private void OpcUaClient_ConnectComplete(object sender, EventArgs e)
  682. {
  683. ConnectChangedEvent?.Invoke(this, OpcUaClient.Connected);
  684. }
  685. #endregion
  686. }
  687. }