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::panic::{AssertUnwindSafe, catch_unwind};
10use std::sync::Arc;
11use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
12use tokio::sync::mpsc;
13
14pub struct OpcDaClient<C: ServerConnector + 'static = ComConnector> {
18 pub worker: ComWorker<C>,
19 connector: Arc<C>,
20 inventory_active: Arc<AtomicBool>,
21}
22
23struct InventoryActiveGuard(Arc<AtomicBool>);
24
25impl Drop for InventoryActiveGuard {
26 fn drop(&mut self) {
27 self.0.store(false, Ordering::Release);
28 }
29}
30
31impl Default for OpcDaClient<ComConnector> {
39 fn default() -> Self {
40 match Self::new(ComConnector) {
41 Ok(client) => client,
42 Err(err) => {
43 tracing::error!(error = ?err, "Failed to initialize default OpcDaClient");
44 Self {
45 worker: ComWorker::closed(),
46 connector: Arc::new(ComConnector),
47 inventory_active: Arc::new(AtomicBool::new(false)),
48 }
49 }
50 }
51 }
52}
53
54impl<C: ServerConnector + 'static> OpcDaClient<C> {
55 pub fn new(connector: C) -> OpcResult<Self> {
57 tracing::info!("Initializing OpcDaClient...");
58 let connector = Arc::new(connector);
59 let worker = ComWorker::start(Arc::clone(&connector))?;
60 tracing::info!("OpcDaClient initialized successfully");
61 Ok(Self {
62 worker,
63 connector,
64 inventory_active: Arc::new(AtomicBool::new(false)),
65 })
66 }
67
68 async fn read_tag_values_with_presentation(
69 &self,
70 server: &str,
71 tag_ids: Vec<String>,
72 presentation: ReadPresentation,
73 ) -> OpcResult<Vec<TagValue>> {
74 let server_owned = server.to_string();
75 self.worker
76 .send_request(|reply| ComRequest::ReadTagValues {
77 server: server_owned,
78 tag_ids,
79 presentation,
80 reply,
81 })
82 .await
83 }
84}
85
86#[allow(clippy::too_many_lines)]
87#[async_trait]
88impl<C: ServerConnector + 'static> OpcProvider for OpcDaClient<C> {
89 async fn list_servers(&self, host: &str) -> OpcResult<Vec<String>> {
90 let host_owned = host.to_string();
91 self.worker
92 .send_request(|reply| ComRequest::ListServers {
93 host: host_owned,
94 reply,
95 })
96 .await
97 }
98
99 async fn browse_tags(
100 &self,
101 server: &str,
102 max_tags: usize,
103 progress: Arc<AtomicUsize>,
104 tags_sink: Arc<std::sync::Mutex<Vec<String>>>,
105 ) -> OpcResult<Vec<String>> {
106 let server_owned = server.to_string();
107 self.worker
108 .send_request(|reply| ComRequest::BrowseTags {
109 server: server_owned,
110 max_tags,
111 progress,
112 tags_sink,
113 reply,
114 })
115 .await
116 }
117
118 async fn browse_capabilities(&self, server: &str) -> OpcResult<BrowseCapabilities> {
119 let server_owned = server.to_string();
120 self.worker
121 .send_request(|reply| ComRequest::BrowseCapabilities {
122 server: server_owned,
123 reply,
124 })
125 .await
126 }
127
128 async fn open_browse_session(&self, server: &str) -> OpcResult<BrowseSessionToken> {
129 let server_owned = server.to_string();
130 self.worker
131 .send_request(|reply| ComRequest::OpenBrowseSession {
132 server: server_owned,
133 reply,
134 })
135 .await
136 }
137
138 async fn browse_page(
139 &self,
140 session: &BrowseSessionToken,
141 request: BrowsePageRequest,
142 ) -> OpcResult<BrowsePage> {
143 let session = *session;
144 self.worker
145 .send_request(|reply| ComRequest::BrowsePage {
146 session,
147 request,
148 reply,
149 })
150 .await
151 }
152
153 async fn close_browse_session(&self, session: &BrowseSessionToken) -> OpcResult<()> {
154 let session = *session;
155 self.worker
156 .send_request(|reply| ComRequest::CloseBrowseSession { session, reply })
157 .await
158 }
159
160 async fn start_inventory(
161 &self,
162 server: &str,
163 options: InventoryOptions,
164 ) -> OpcResult<InventoryStream> {
165 if options.batch_size == 0 || options.batch_size > 1_000 {
166 return Err(crate::opc_da::errors::OpcError::InvalidState(
167 "Inventory batch size must be between 1 and 1000".to_string(),
168 ));
169 }
170 if self.inventory_active.swap(true, Ordering::AcqRel) {
171 return Err(crate::opc_da::errors::OpcError::InvalidState(
172 "An OPC namespace inventory is already running".to_string(),
173 ));
174 }
175
176 let (sender, receiver) = mpsc::channel(64);
177 let control = InventoryControl::new();
178 let worker_control = control.clone();
179 let active = Arc::clone(&self.inventory_active);
180 let connector = Arc::clone(&self.connector);
181 let server = server.to_string();
182 let spawn_result = std::thread::Builder::new()
183 .name("opc-da-inventory".to_string())
184 .spawn(move || {
185 let _active_guard = InventoryActiveGuard(active);
186 let result = catch_unwind(AssertUnwindSafe(|| {
187 let _guard = crate::ComGuard::new().map_err(|error| {
188 crate::opc_da::errors::OpcError::Internal(error.to_string())
189 })?;
190 crate::inventory::run_inventory(
191 &*connector,
192 &server,
193 options,
194 &worker_control,
195 &sender,
196 )
197 }));
198 match result {
199 Ok(Ok(())) => {}
200 Ok(Err(error)) => {
201 tracing::error!(server = %server, error = %error, "OPC namespace inventory failed");
202 if sender.blocking_send(Err(error)).is_err() {
203 tracing::debug!(
204 server = %server,
205 "OPC namespace inventory receiver was closed before failure could be delivered"
206 );
207 }
208 }
209 Err(payload) => {
210 let payload_type = panic_payload_type(&*payload);
211 tracing::error!(
212 server = %server,
213 payload_type,
214 "OPC namespace inventory worker panicked"
215 );
216 let error = crate::opc_da::errors::OpcError::Internal(
217 "OPC namespace inventory worker panicked".to_string(),
218 );
219 if sender.blocking_send(Err(error)).is_err() {
220 tracing::debug!(
221 server = %server,
222 "OPC namespace inventory receiver was closed before panic could be delivered"
223 );
224 }
225 }
226 }
227 });
228
229 let worker = match spawn_result {
230 Ok(worker) => worker,
231 Err(error) => {
232 self.inventory_active.store(false, Ordering::Release);
233 return Err(crate::opc_da::errors::OpcError::Internal(format!(
234 "failed to start OPC inventory worker: {error}"
235 )));
236 }
237 };
238
239 Ok(InventoryStream::new(receiver, control, worker))
240 }
241
242 async fn read_tag_values(
243 &self,
244 server: &str,
245 tag_ids: Vec<String>,
246 ) -> OpcResult<Vec<TagValue>> {
247 self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Semantic)
248 .await
249 }
250
251 async fn read_tag_values_for_display(
252 &self,
253 server: &str,
254 tag_ids: Vec<String>,
255 ) -> OpcResult<Vec<TagValue>> {
256 self.read_tag_values_with_presentation(server, tag_ids, ReadPresentation::Display)
257 .await
258 }
259
260 async fn write_tag_value(
261 &self,
262 server: &str,
263 tag_id: &str,
264 value: OpcValue,
265 ) -> OpcResult<WriteResult> {
266 let server_owned = server.to_string();
267 let tag_id_owned = tag_id.to_string();
268 self.worker
269 .send_request(|reply| ComRequest::WriteTagValue {
270 server: server_owned,
271 tag_id: tag_id_owned,
272 value,
273 reply,
274 })
275 .await
276 }
277}
278
279fn panic_payload_type(payload: &(dyn std::any::Any + Send)) -> &'static str {
280 if payload.is::<&'static str>() {
281 "&str"
282 } else if payload.is::<String>() {
283 "String"
284 } else {
285 "other"
286 }
287}