1use crate::backend::connector::{ComConnector, ServerConnector};
2use crate::com_worker::{ComRequest, ComWorker, ReadPresentation};
3use crate::opc_da::errors::OpcResult;
4use crate::provider::{
5 BrowseCapabilities, BrowsePage, BrowsePageRequest, BrowseSessionToken, InventoryControl,
6 InventoryOptions, InventoryStream, MAX_INVENTORY_BATCH_SIZE, OpcProvider, OpcValue, TagValue,
7 WriteResult,
8};
9use async_trait::async_trait;
10use std::panic::{AssertUnwindSafe, catch_unwind};
11use std::sync::Arc;
12use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
13use tokio::sync::mpsc;
14
15pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
19 pub worker: ComWorker<C>,
20 connector: Arc<C>,
21 inventory_active: Arc<AtomicBool>,
22}
23
24struct InventoryActiveGuard(Arc<AtomicBool>);
25
26impl Drop for InventoryActiveGuard {
27 fn drop(&mut self) {
28 self.0.store(false, Ordering::Release);
29 }
30}
31
32impl Default for OpcDaClient<ComConnector> {
40 fn default() -> Self {
41 match Self::new(ComConnector) {
42 Ok(client) => client,
43 Err(err) => {
44 tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
45 Self {
46 worker: ComWorker::closed(),
47 connector: Arc::new(ComConnector),
48 inventory_active: Arc::new(AtomicBool::new(false)),
49 }
50 }
51 }
52 }
53}
54
55impl<C: ServerConnector + 'static> OpcDaClient<C> {
56 pub fn new(connector: C) -> OpcResult<Self> {
58 tracing::info!("Initializing OpcDaClient...");
59 let connector = Arc::new(connector);
60 let worker = ComWorker::start(Arc::clone(&connector))?;
61 tracing::info!("OpcDaClient initialized successfully");
62 Ok(Self {
63 worker,
64 connector,
65 inventory_active: Arc::new(AtomicBool::new(false)),
66 })
67 }
68
69 async fn read_tag_values_with_presentation(
70 &self,
71 server: &str,
72 tag_ids: Vec<String>,
73 presentation: ReadPresentation,
74 ) -> OpcResult<Vec<TagValue>> {
75 let server_owned = server.to_string();
76 self.worker
77 .send_request(|reply| ComRequest::ReadTagValues {
78 server: server_owned,
79 tag_ids,
80 presentation,
81 reply,
82 })
83 .await
84 }
85}
86
87#[allow(clippy::too_many_lines)]
88#[async_trait]
89impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
90 async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
91 let host_owned = host.to_string();
92 self.worker
93 .send_request(|reply| ComRequest::ListServers {
94 host: host_owned,
95 reply,
96 })
97 .await
98 }
99
100 async fn browse_tags(
101 &self,
102 server: &str,
103 max_tags: usize,
104 progress: Arc<AtomicUsize>,
105 tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
106 ) -> OpcResult<Vec<String>> {
107 let server_owned = server.to_string();
108 self.worker
109 .send_request(|reply| ComRequest::BrowseTags {
110 server: server_owned,
111 max_tags,
112 progress,
113 tags_sink,
114 reply,
115 })
116 .await
117 }
118
119 async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
120 let server_owned = server.to_string();
121 self.worker
122 .send_request(|reply| ComRequest::BrowseCapabilities {
123 server: server_owned,
124 reply,
125 })
126 .await
127 }
128
129 async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
130 let server_owned = server.to_string();
131 self.worker
132 .send_request(|reply| ComRequest::OpenBrowseSession {
133 server: server_owned,
134 reply,
135 })
136 .await
137 }
138
139 async fn browse_page(
140 &self,
141 session: &BrowseSessionToken,
142 request: BrowsePageRequest,
143 ) -> OpcResult<BrowsePage> {
144 let session = *session;
145 self.worker
146 .send_request(|reply| ComRequest::BrowsePage {
147 session,
148 request,
149 reply,
150 })
151 .await
152 }
153
154 async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
155 let session = *session;
156 self.worker
157 .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
158 .await
159 }
160
161 async fn start_inventory(
162 &self,
163 server: &str,
164 options: InventoryOptions,
165 ) -> OpcResult<InventoryStream> {
166 if options.batch_size == 0 || options.batch_size > MAX_INVENTORY_BATCH_SIZE {
167 return Err(crate::opc_da::errors::OpcError::InvalidState(format!(
168 "Inventory batch size must be between 1 and {MAX_INVENTORY_BATCH_SIZE}"
169 )));
170 }
171 if self.inventory_active.swap(true, Ordering::AcqRel) {
172 return Err(crate::opc_da::errors::OpcError::InvalidState(
173 "An OPC namespace inventory is already running".to_string(),
174 ));
175 }
176
177 let (sender, receiver) = mpsc::channel(64);
178 let control = InventoryControl::new_with_batch_size(options.batch_size);
179 let worker_control = control.clone();
180 let active = Arc::clone(&self.inventory_active);
181 let connector = Arc::clone(&self.connector);
182 let server = server.to_string();
183 let spawn_result = std::thread::Builder::new()
184 .name("opc-da-inventory".to_string())
185 .spawn(move || {
186 let _active_guard = InventoryActiveGuard(active);
187 let result = catch_unwind(AssertUnwindSafe(|| {
188 let _guard = crate::ComGuard::new().map_err(|error| {
189 crate::opc_da::errors::OpcError::Internal(error.to_string())
190 })?;
191 crate::inventory::run_inventory(
192 &*connector,
193 &server,
194 options,
195 &worker_control,
196 &sender,
197 )
198 }));
199 match result {
200 Ok(Ok(())) => {}
201 Ok(Err(error)) => {
202 tracing::error!(server = %server, error = %error, "OPC namespace inventory failed");
203 if sender.blocking_send(Err(error)).is_err() {
204 tracing::debug!(
205 server = %server,
206 "OPC namespace inventory receiver was closed before failure could be delivered"
207 );
208 }
209 }
210 Err(payload) => {
211 let payload_type = panic_payload_type(&*payload);
212 tracing::error!(
213 server = %server,
214 payload_type,
215 "OPC namespace inventory worker panicked"
216 );
217 let error = crate::opc_da::errors::OpcError::Internal(
218 "OPC namespace inventory worker panicked".to_string(),
219 );
220 if sender.blocking_send(Err(error)).is_err() {
221 tracing::debug!(
222 server = %server,
223 "OPC namespace inventory receiver was closed before panic could be delivered"
224 );
225 }
226 }
227 }
228 });
229
230 let worker = match spawn_result {
231 Ok(worker) => worker,
232 Err(error) => {
233 self.inventory_active.store(false, Ordering::Release);
234 return Err(crate::opc_da::errors::OpcError::Internal(format!(
235 "failed to start OPC inventory worker: {error}"
236 )));
237 }
238 };
239
240 Ok(InventoryStream::new(receiver, control, worker))
241 }
242
243 async fn read_tag_values(
244 &self,
245 server: &str,
246 tag_ids: Vec<String>,
247 ) -> OpcResult<Vec<TagValue>> {
248 self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Semantic)
249 .await
250 }
251
252 async fn read_tag_values_for_display(
253 &self,
254 server: &str,
255 tag_ids: Vec<String>,
256 ) -> OpcResult<Vec<TagValue>> {
257 self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Display)
258 .await
259 }
260
261 async fn write_tag_value(
262 &self,
263 server: &str,
264 tag_id: &str,
265 value: OpcValue,
266 ) -> OpcResult<WriteResult> {
267 let server_owned = server.to_string();
268 let tag_id_owned = tag_id.to_string();
269 self.worker
270 .send_request(|reply| ComRequest::WriteTagValue {
271 server: server_owned,
272 tag_id: tag_id_owned,
273 value,
274 reply,
275 })
276 .await
277 }
278}
279
280fn panic_payload_type(payload: &(dyn std::any::Any + Send)) -> &'static str {
281 if payload.is::<&'static str>() {
282 "&str"
283 } else if payload.is::<String>() {
284 "String"
285 } else {
286 "other"
287 }
288}