Skip to main content

opc_da_client/backend/
opc_da.rs

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