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
13pub 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
30impl 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 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}