Skip to main content

opc_da_client/backend/
opc_da.rs

1use 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
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
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}