opc_da_client/backend/
opc_da.rs1use crate::backend::connector::{ComConnector, ServerConnector};
2use crate::com_worker::{ComRequest, ComWorker};
3use crate::opc_da::errors::OpcResult;
4use crate::provider::{
5 BrowseCapabilities, BrowsePage, BrowsePageRequest, BrowseSessionToken, OpcProvider, OpcValue,
6 TagValue, WriteResult,
7};
8use async_trait::async_trait;
9use std::sync::Arc;
10use std::sync::atomic::AtomicUsize;
11
12pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
16 pub worker: ComWorker<C>,
17}
18
19impl Default for OpcDaClient<ComConnector> {
27 fn default() -> Self {
28 match Self::new(ComConnector) {
29 Ok(client) => client,
30 Err(err) => {
31 tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
32 Self {
33 worker: ComWorker::closed(),
34 }
35 }
36 }
37 }
38}
39
40impl<C: ServerConnector + 'static> OpcDaClient<C> {
41 pub fn new(connector: C) -> OpcResult<Self> {
43 tracing::info!("Initializing OpcDaClient...");
44 let worker = ComWorker::start(Arc::new(connector))?;
45 tracing::info!("OpcDaClient initialized successfully");
46 Ok(Self { worker })
47 }
48}
49
50#[allow(clippy::too_many_lines)]
51#[async_trait]
52impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
53 async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
54 let host_owned = host.to_string();
55 self.worker
56 .send_request(|reply| ComRequest::ListServers {
57 host: host_owned,
58 reply,
59 })
60 .await
61 }
62
63 async fn browse_tags(
64 &self,
65 server: &str,
66 max_tags: usize,
67 progress: Arc<AtomicUsize>,
68 tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
69 ) -> OpcResult<Vec<String>> {
70 let server_owned = server.to_string();
71 self.worker
72 .send_request(|reply| ComRequest::BrowseTags {
73 server: server_owned,
74 max_tags,
75 progress,
76 tags_sink,
77 reply,
78 })
79 .await
80 }
81
82 async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
83 let server_owned = server.to_string();
84 self.worker
85 .send_request(|reply| ComRequest::BrowseCapabilities {
86 server: server_owned,
87 reply,
88 })
89 .await
90 }
91
92 async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
93 let server_owned = server.to_string();
94 self.worker
95 .send_request(|reply| ComRequest::OpenBrowseSession {
96 server: server_owned,
97 reply,
98 })
99 .await
100 }
101
102 async fn browse_page(
103 &self,
104 session: &BrowseSessionToken,
105 request: BrowsePageRequest,
106 ) -> OpcResult<BrowsePage> {
107 let session = *session;
108 self.worker
109 .send_request(|reply| ComRequest::BrowsePage {
110 session,
111 request,
112 reply,
113 })
114 .await
115 }
116
117 async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
118 let session = *session;
119 self.worker
120 .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
121 .await
122 }
123
124 async fn read_tag_values(
125 &self,
126 server: &str,
127 tag_ids: Vec<String>,
128 ) -> OpcResult<Vec<TagValue>> {
129 let server_owned = server.to_string();
130 self.worker
131 .send_request(|reply| ComRequest::ReadTagValues {
132 server: server_owned,
133 tag_ids,
134 reply,
135 })
136 .await
137 }
138
139 async fn write_tag_value(
140 &self,
141 server: &str,
142 tag_id: &str,
143 value: OpcValue,
144 ) -> OpcResult<WriteResult> {
145 let server_owned = server.to_string();
146 let tag_id_owned = tag_id.to_string();
147 self.worker
148 .send_request(|reply| ComRequest::WriteTagValue {
149 server: server_owned,
150 tag_id: tag_id_owned,
151 value,
152 reply,
153 })
154 .await
155 }
156}