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, InventoryControl,
6 InventoryOptions, InventoryStream, OpcProvider, OpcValue, TagValue, WriteResult,
7};
8use async_trait::async_trait;
9use std::sync::Arc;
10use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
11use tokio::sync::mpsc;
12
13pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
17 pub worker: ComWorker<C>,
18 connector: Arc<C>,
19 inventory_active: Arc<AtomicBool>,
20}
21
22struct InventoryActiveGuard(Arc<AtomicBool>);
23
24impl Drop for InventoryActiveGuard {
25 fn drop(&mut self) {
26 self.0.store(false, Ordering::Release);
27 }
28}
29
30impl Default for OpcDaClient<ComConnector> {
38 fn default() -> Self {
39 match Self::new(ComConnector) {
40 Ok(client) => client,
41 Err(err) => {
42 tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
43 Self {
44 worker: ComWorker::closed(),
45 connector: Arc::new(ComConnector),
46 inventory_active: Arc::new(AtomicBool::new(false)),
47 }
48 }
49 }
50 }
51}
52
53impl<C: ServerConnector + 'static> OpcDaClient<C> {
54 pub fn new(connector: C) -> OpcResult<Self> {
56 tracing::info!("Initializing OpcDaClient...");
57 let connector = Arc::new(connector);
58 let worker = ComWorker::start(Arc::clone(&connector))?;
59 tracing::info!("OpcDaClient initialized successfully");
60 Ok(Self {
61 worker,
62 connector,
63 inventory_active: Arc::new(AtomicBool::new(false)),
64 })
65 }
66}
67
68#[allow(clippy::too_many_lines)]
69#[async_trait]
70impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
71 async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
72 let host_owned = host.to_string();
73 self.worker
74 .send_request(|reply| ComRequest::ListServers {
75 host: host_owned,
76 reply,
77 })
78 .await
79 }
80
81 async fn browse_tags(
82 &self,
83 server: &str,
84 max_tags: usize,
85 progress: Arc<AtomicUsize>,
86 tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
87 ) -> OpcResult<Vec<String>> {
88 let server_owned = server.to_string();
89 self.worker
90 .send_request(|reply| ComRequest::BrowseTags {
91 server: server_owned,
92 max_tags,
93 progress,
94 tags_sink,
95 reply,
96 })
97 .await
98 }
99
100 async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
101 let server_owned = server.to_string();
102 self.worker
103 .send_request(|reply| ComRequest::BrowseCapabilities {
104 server: server_owned,
105 reply,
106 })
107 .await
108 }
109
110 async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
111 let server_owned = server.to_string();
112 self.worker
113 .send_request(|reply| ComRequest::OpenBrowseSession {
114 server: server_owned,
115 reply,
116 })
117 .await
118 }
119
120 async fn browse_page(
121 &self,
122 session: &BrowseSessionToken,
123 request: BrowsePageRequest,
124 ) -> OpcResult<BrowsePage> {
125 let session = *session;
126 self.worker
127 .send_request(|reply| ComRequest::BrowsePage {
128 session,
129 request,
130 reply,
131 })
132 .await
133 }
134
135 async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
136 let session = *session;
137 self.worker
138 .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
139 .await
140 }
141
142 async fn start_inventory(
143 &self,
144 server: &str,
145 options: InventoryOptions,
146 ) -> OpcResult<InventoryStream> {
147 if options.batch_size == 0 || options.batch_size > 1_000 {
148 return Err(crate::opc_da::errors::OpcError::InvalidState(
149 "Inventory batch size must be between 1 and 1000".to_string(),
150 ));
151 }
152 if self.inventory_active.swap(true, Ordering::AcqRel) {
153 return Err(crate::opc_da::errors::OpcError::InvalidState(
154 "An OPC namespace inventory is already running".to_string(),
155 ));
156 }
157
158 let (sender, receiver) = mpsc::channel(64);
159 let control = InventoryControl::new();
160 let worker_control = control.clone();
161 let active = Arc::clone(&self.inventory_active);
162 let connector = Arc::clone(&self.connector);
163 let server = server.to_string();
164 let spawn_result = std::thread::Builder::new()
165 .name("opc-da-inventory".to_string())
166 .spawn(move || {
167 let _active_guard = InventoryActiveGuard(active);
168 let result = (|| {
169 let _guard = crate::ComGuard::new().map_err(|error| {
170 crate::opc_da::errors::OpcError::Internal(error.to_string())
171 })?;
172 crate::inventory::run_inventory(
173 &*connector,
174 &server,
175 options,
176 &worker_control,
177 &sender,
178 )
179 })();
180 if let Err(error) = result {
181 let _ = sender.blocking_send(Err(error));
182 }
183 });
184
185 let worker = match spawn_result {
186 Ok(worker) => worker,
187 Err(error) => {
188 self.inventory_active.store(false, Ordering::Release);
189 return Err(crate::opc_da::errors::OpcError::Internal(format!(
190 "failed to start OPC inventory worker: {error}"
191 )));
192 }
193 };
194
195 Ok(InventoryStream::new(receiver, control, worker))
196 }
197
198 async fn read_tag_values(
199 &self,
200 server: &str,
201 tag_ids: Vec<String>,
202 ) -> OpcResult<Vec<TagValue>> {
203 let server_owned = server.to_string();
204 self.worker
205 .send_request(|reply| ComRequest::ReadTagValues {
206 server: server_owned,
207 tag_ids,
208 reply,
209 })
210 .await
211 }
212
213 async fn write_tag_value(
214 &self,
215 server: &str,
216 tag_id: &str,
217 value: OpcValue,
218 ) -> OpcResult<WriteResult> {
219 let server_owned = server.to_string();
220 let tag_id_owned = tag_id.to_string();
221 self.worker
222 .send_request(|reply| ComRequest::WriteTagValue {
223 server: server_owned,
224 tag_id: tag_id_owned,
225 value,
226 reply,
227 })
228 .await
229 }
230}