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, OpcProvider, OpcValue,
6    TagValue, WriteResult,
7};
8use async_trait::async_trait;
9use std::sync::Arc;
10use std::sync::atomic::AtomicUsize;
11
12/// Concrete [`OpcProvider`] implementation for Windows OPC DA.
13///
14/// Uses native `windows-rs` COM interop via the internal `opc_da` module.
15pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
16    pub worker: ComWorker<C>,
17}
18
19/// Returns the default `OpcDaClient` using native COM settings.
20///
21/// # Panics
22///
23/// Panics if the background COM worker thread cannot be started or COM
24/// Multi-Threaded Apartment (MTA) initialization fails on the worker thread.
25/// Use [`OpcDaClient::new`] for fallible construction.
26impl Default for OpcDaClient<ComConnector> {
27    fn default() -> Self {
28        match Self::new(ComConnector) {
29            Ok(client) => client,
30            Err(err) => {
31                tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
32                Self {
33                    worker: ComWorker::closed(),
34                }
35            }
36        }
37    }
38}
39
40impl<C: ServerConnector + 'static> OpcDaClient<C> {
41    /// Creates a new `OpcDaClient` with the given connector.
42    pub fn new(connector: C) -> OpcResult<Self> {
43        tracing::info!("Initializing OpcDaClient...");
44        let worker = ComWorker::start(Arc::new(connector))?;
45        tracing::info!("OpcDaClient initialized successfully");
46        Ok(Self { worker })
47    }
48}
49
50#[allow(clippy::too_many_lines)]
51#[async_trait]
52impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
53    async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
54        let host_owned = host.to_string();
55        self.worker
56            .send_request(|reply| ComRequest::ListServers {
57                host: host_owned,
58                reply,
59            })
60            .await
61    }
62
63    async fn browse_tags(
64        &self,
65        server: &str,
66        max_tags: usize,
67        progress: Arc<AtomicUsize>,
68        tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
69    ) -> OpcResult<Vec<String>> {
70        let server_owned = server.to_string();
71        self.worker
72            .send_request(|reply| ComRequest::BrowseTags {
73                server: server_owned,
74                max_tags,
75                progress,
76                tags_sink,
77                reply,
78            })
79            .await
80    }
81
82    async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
83        let server_owned = server.to_string();
84        self.worker
85            .send_request(|reply| ComRequest::BrowseCapabilities {
86                server: server_owned,
87                reply,
88            })
89            .await
90    }
91
92    async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
93        let server_owned = server.to_string();
94        self.worker
95            .send_request(|reply| ComRequest::OpenBrowseSession {
96                server: server_owned,
97                reply,
98            })
99            .await
100    }
101
102    async fn browse_page(
103        &self,
104        session: &BrowseSessionToken,
105        request: BrowsePageRequest,
106    ) -> OpcResult<BrowsePage> {
107        let session = *session;
108        self.worker
109            .send_request(|reply| ComRequest::BrowsePage {
110                session,
111                request,
112                reply,
113            })
114            .await
115    }
116
117    async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
118        let session = *session;
119        self.worker
120            .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
121            .await
122    }
123
124    async fn read_tag_values(
125        &self,
126        server: &str,
127        tag_ids: Vec<String>,
128    ) -> OpcResult<Vec<TagValue>> {
129        let server_owned = server.to_string();
130        self.worker
131            .send_request(|reply| ComRequest::ReadTagValues {
132                server: server_owned,
133                tag_ids,
134                reply,
135            })
136            .await
137    }
138
139    async fn write_tag_value(
140        &self,
141        server: &str,
142        tag_id: &str,
143        value: OpcValue,
144    ) -> OpcResult<WriteResult> {
145        let server_owned = server.to_string();
146        let tag_id_owned = tag_id.to_string();
147        self.worker
148            .send_request(|reply| ComRequest::WriteTagValue {
149                server: server_owned,
150                tag_id: tag_id_owned,
151                value,
152                reply,
153            })
154            .await
155    }
156}