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, 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
15/// Concrete [`OpcProvider`] implementation for Windows OPC DA.
16///
17/// Uses native `windows-rs` COM interop via the internal `opc_da` module.
18pub 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
32/// Returns the default `OpcDaClient` using native COM settings.
33///
34/// # Panics
35///
36/// Panics if the background COM worker thread cannot be started or COM
37/// Multi-Threaded Apartment (MTA) initialization fails on the worker thread.
38/// Use [`OpcDaClient::new`] for fallible construction.
39impl 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    /// Creates a new `OpcDaClient` with the given connector.
57    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}