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::sync::Arc;
10use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
11use tokio::sync::mpsc;
12
13/// Concrete [`OpcProvider`] implementation for Windows OPC DA.
14///
15/// Uses native `windows-rs` COM interop via the internal `opc_da` module.
16pub 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
30/// Returns the default `OpcDaClient` using native COM settings.
31///
32/// # Panics
33///
34/// Panics if the background COM worker thread cannot be started or COM
35/// Multi-Threaded Apartment (MTA) initialization fails on the worker thread.
36/// Use [`OpcDaClient::new`] for fallible construction.
37impl 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    /// Creates a new `OpcDaClient` with the given connector.
55    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    async fn read_tag_values_with_presentation(
68        &self,
69        server: &str,
70        tag_ids: Vec<String>,
71        presentation: ReadPresentation,
72    ) -> OpcResult<Vec<TagValue>> {
73        let server_owned = server.to_string();
74        self.worker
75            .send_request(|reply| ComRequest::ReadTagValues {
76                server: server_owned,
77                tag_ids,
78                presentation,
79                reply,
80            })
81            .await
82    }
83}
84
85#[allow(clippy::too_many_lines)]
86#[async_trait]
87impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
88    async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
89        let host_owned = host.to_string();
90        self.worker
91            .send_request(|reply| ComRequest::ListServers {
92                host: host_owned,
93                reply,
94            })
95            .await
96    }
97
98    async fn browse_tags(
99        &self,
100        server: &str,
101        max_tags: usize,
102        progress: Arc<AtomicUsize>,
103        tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
104    ) -> OpcResult<Vec<String>> {
105        let server_owned = server.to_string();
106        self.worker
107            .send_request(|reply| ComRequest::BrowseTags {
108                server: server_owned,
109                max_tags,
110                progress,
111                tags_sink,
112                reply,
113            })
114            .await
115    }
116
117    async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
118        let server_owned = server.to_string();
119        self.worker
120            .send_request(|reply| ComRequest::BrowseCapabilities {
121                server: server_owned,
122                reply,
123            })
124            .await
125    }
126
127    async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
128        let server_owned = server.to_string();
129        self.worker
130            .send_request(|reply| ComRequest::OpenBrowseSession {
131                server: server_owned,
132                reply,
133            })
134            .await
135    }
136
137    async fn browse_page(
138        &self,
139        session: &BrowseSessionToken,
140        request: BrowsePageRequest,
141    ) -> OpcResult<BrowsePage> {
142        let session = *session;
143        self.worker
144            .send_request(|reply| ComRequest::BrowsePage {
145                session,
146                request,
147                reply,
148            })
149            .await
150    }
151
152    async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
153        let session = *session;
154        self.worker
155            .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
156            .await
157    }
158
159    async fn start_inventory(
160        &self,
161        server: &str,
162        options: InventoryOptions,
163    ) -> OpcResult<InventoryStream> {
164        if options.batch_size == 0 || options.batch_size > 1_000 {
165            return Err(crate::opc_da::errors::OpcError::InvalidState(
166                "Inventory batch size must be between 1 and 1000".to_string(),
167            ));
168        }
169        if self.inventory_active.swap(true, Ordering::AcqRel) {
170            return Err(crate::opc_da::errors::OpcError::InvalidState(
171                "An OPC namespace inventory is already running".to_string(),
172            ));
173        }
174
175        let (sender, receiver) = mpsc::channel(64);
176        let control = InventoryControl::new();
177        let worker_control = control.clone();
178        let active = Arc::clone(&self.inventory_active);
179        let connector = Arc::clone(&self.connector);
180        let server = server.to_string();
181        let spawn_result = std::thread::Builder::new()
182            .name("opc-da-inventory".to_string())
183            .spawn(move || {
184                let _active_guard = InventoryActiveGuard(active);
185                let result = (|| {
186                    let _guard = crate::ComGuard::new().map_err(|error| {
187                        crate::opc_da::errors::OpcError::Internal(error.to_string())
188                    })?;
189                    crate::inventory::run_inventory(
190                        &*connector,
191                        &server,
192                        options,
193                        &worker_control,
194                        &sender,
195                    )
196                })();
197                if let Err(error) = result {
198                    let _ = sender.blocking_send(Err(error));
199                }
200            });
201
202        let worker = match spawn_result {
203            Ok(worker) => worker,
204            Err(error) => {
205                self.inventory_active.store(false, Ordering::Release);
206                return Err(crate::opc_da::errors::OpcError::Internal(format!(
207                    "failed to start OPC inventory worker: {error}"
208                )));
209            }
210        };
211
212        Ok(InventoryStream::new(receiver, control, worker))
213    }
214
215    async fn read_tag_values(
216        &self,
217        server: &str,
218        tag_ids: Vec<String>,
219    ) -> OpcResult<Vec<TagValue>> {
220        self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Semantic)
221            .await
222    }
223
224    async fn read_tag_values_for_display(
225        &self,
226        server: &str,
227        tag_ids: Vec<String>,
228    ) -> OpcResult<Vec<TagValue>> {
229        self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Display)
230            .await
231    }
232
233    async fn write_tag_value(
234        &self,
235        server: &str,
236        tag_id: &str,
237        value: OpcValue,
238    ) -> OpcResult<WriteResult> {
239        let server_owned = server.to_string();
240        let tag_id_owned = tag_id.to_string();
241        self.worker
242            .send_request(|reply| ComRequest::WriteTagValue {
243                server: server_owned,
244                tag_id: tag_id_owned,
245                value,
246                reply,
247            })
248            .await
249    }
250}