OPCuaClientPLC.cs 27 KB

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