1use crate::error::Result;
4use crate::types::{BrowseNode, TagValue, Value, WriteResult};
5use opcda_bridge_proto::bridge::bridge_client::BridgeClient;
6use opcda_bridge_proto::bridge::write_request::TypedValue;
7use opcda_bridge_proto::bridge::{BrowseRequest, ListServersRequest, ReadRequest, WriteRequest};
8use tonic::transport::Channel;
9
10#[derive(Debug)]
20pub struct Client {
21 inner: BridgeClient<Channel>,
22}
23
24impl Client {
25 pub async fn connect(host: &str) -> Result<Self> {
30 let inner = BridgeClient::connect(format!("http://{host}")).await?;
31 Ok(Self { inner })
32 }
33
34 pub async fn list_servers(&mut self) -> Result<Vec<String>> {
36 let response = self
37 .inner
38 .list_servers(ListServersRequest {
39 host: "localhost".to_string(),
40 })
41 .await?;
42 Ok(response.into_inner().servers)
43 }
44
45 pub async fn browse(
57 &mut self,
58 server: String,
59 flat: bool,
60 path: String,
61 max_tags: u32,
62 ) -> Result<Vec<BrowseNode>> {
63 let mut stream = self
64 .inner
65 .browse(BrowseRequest {
66 server,
67 flat,
68 path,
69 max_tags,
70 })
71 .await?
72 .into_inner();
73
74 let mut nodes = Vec::new();
75 while let Some(response) = stream.message().await? {
76 nodes.push(BrowseNode {
77 tag_id: response.tag_id,
78 node_type: response.node_type,
79 });
80 }
81 Ok(nodes)
82 }
83
84 pub async fn read(&mut self, server: String, tags: Vec<String>) -> Result<Vec<TagValue>> {
86 let response = self
87 .inner
88 .read(ReadRequest {
89 server,
90 tag_ids: tags,
91 })
92 .await?;
93 Ok(response
94 .into_inner()
95 .values
96 .into_iter()
97 .map(|v| TagValue {
98 tag_id: v.tag_id,
99 value: v.value,
100 quality: v.quality,
101 timestamp: v.timestamp,
102 })
103 .collect())
104 }
105
106 pub async fn write(
108 &mut self,
109 server: String,
110 tag: String,
111 value: Value,
112 ) -> Result<WriteResult> {
113 let typed_value = match value {
114 Value::String(s) => TypedValue::StringValue(s),
115 Value::Int(i) => TypedValue::IntValue(i),
116 Value::Float(f) => TypedValue::FloatValue(f),
117 Value::Bool(b) => TypedValue::BoolValue(b),
118 };
119 let response = self
120 .inner
121 .write(WriteRequest {
122 server,
123 tag_id: tag,
124 typed_value: Some(typed_value),
125 })
126 .await?;
127 let r = response.into_inner();
128 Ok(WriteResult {
129 tag_id: r.tag_id,
130 success: r.success,
131 error: r.error,
132 })
133 }
134}
135
136#[cfg(test)]
137mod tests {
138 use super::*;
139 use crate::error::Error;
140 use crate::test_support::{MockBridgeService, start_mock_server};
141 use opcda_bridge_proto::bridge::bridge_server::Bridge;
142 use opcda_bridge_proto::bridge::{
143 BrowseResponse, bridge_client::BridgeClient as ProtoBridgeClient,
144 };
145 use opcda_bridge_proto::bridge::{
146 ListServersResponse, ReadResponse, TagValue as ProtoTagValue, WriteResponse,
147 };
148 use std::sync::Arc;
149 use std::time::Duration;
150 use tonic::{Request, Status};
151
152 #[tokio::test]
153 async fn test_connect_success() {
154 let host = start_mock_server(MockBridgeService::default()).await;
155 Client::connect(&host).await.unwrap();
156 }
157
158 #[tokio::test]
159 async fn test_mock_server_shutdown_completes_background_task() {
160 let service = MockBridgeService::default();
161 let server_shutdown = Arc::clone(&service.server_shutdown);
162 let server_stopped = Arc::clone(&service.server_stopped);
163 let _host = start_mock_server(service).await;
164
165 server_shutdown.notify_one();
166 tokio::time::timeout(Duration::from_secs(1), server_stopped.notified())
167 .await
168 .expect("mock server did not stop after shutdown");
169 }
170
171 #[tokio::test]
172 async fn test_connect_failure_is_connect_variant() {
173 let err = Client::connect("127.0.0.1:1").await.unwrap_err();
174 assert!(matches!(err, Error::Connect(_)));
175 }
176
177 #[tokio::test]
178 async fn test_connect_failure_anyhow_debug_matches_bare_transport_error() {
179 let bare_err = ProtoBridgeClient::connect("http://127.0.0.1:1".to_string())
185 .await
186 .unwrap_err();
187 let bare = anyhow::Error::from(bare_err);
188
189 let wrapped_err = Client::connect("127.0.0.1:1").await.unwrap_err();
190 let wrapped = anyhow::Error::from(wrapped_err);
191
192 assert_eq!(format!("{bare:?}"), format!("{wrapped:?}"));
193 assert_eq!(bare.to_string(), wrapped.to_string());
194 }
195
196 #[tokio::test]
197 async fn test_list_servers_empty() {
198 let host = start_mock_server(MockBridgeService::default()).await;
199 let mut client = Client::connect(&host).await.unwrap();
200 assert_eq!(client.list_servers().await.unwrap(), Vec::<String>::new());
201 }
202
203 #[tokio::test]
204 async fn test_list_servers_with_data() {
205 let svc = MockBridgeService {
206 list_servers_response: ListServersResponse {
207 servers: vec!["Server1".into(), "Server2".into()],
208 },
209 ..Default::default()
210 };
211 let host = start_mock_server(svc).await;
212 let mut client = Client::connect(&host).await.unwrap();
213 assert_eq!(
214 client.list_servers().await.unwrap(),
215 vec!["Server1".to_string(), "Server2".to_string()]
216 );
217 }
218
219 #[tokio::test]
220 async fn test_list_servers_rpc_error() {
221 let svc = MockBridgeService {
222 list_servers_error: Some(Status::internal("boom")),
223 ..Default::default()
224 };
225 let host = start_mock_server(svc).await;
226 let mut client = Client::connect(&host).await.unwrap();
227 let err = client.list_servers().await.unwrap_err();
228 assert!(matches!(err, Error::Rpc(_)));
229 }
230
231 #[tokio::test]
232 async fn test_browse_empty() {
233 let host = start_mock_server(MockBridgeService::default()).await;
234 let mut client = Client::connect(&host).await.unwrap();
235 let nodes = client
236 .browse("S".into(), false, String::new(), 1000)
237 .await
238 .unwrap();
239 assert!(nodes.is_empty());
240 }
241
242 #[tokio::test]
243 async fn test_browse_with_data_maps_fields() {
244 let svc = MockBridgeService {
245 browse_responses: vec![
246 BrowseResponse {
247 tag_id: "tag1".into(),
248 node_type: "Leaf".into(),
249 },
250 BrowseResponse {
251 tag_id: "tag2".into(),
252 node_type: "Branch".into(),
253 },
254 ],
255 ..Default::default()
256 };
257 let host = start_mock_server(svc).await;
258 let mut client = Client::connect(&host).await.unwrap();
259 let nodes = client
260 .browse("S".into(), true, String::new(), 1000)
261 .await
262 .unwrap();
263 assert_eq!(
264 nodes,
265 vec![
266 BrowseNode {
267 tag_id: "tag1".into(),
268 node_type: "Leaf".into(),
269 },
270 BrowseNode {
271 tag_id: "tag2".into(),
272 node_type: "Branch".into(),
273 },
274 ]
275 );
276 }
277
278 #[tokio::test]
279 async fn test_browse_initial_rpc_error() {
280 let svc = MockBridgeService {
281 browse_initial_error: Some(Status::unavailable("gateway down")),
282 ..Default::default()
283 };
284 let host = start_mock_server(svc).await;
285 let mut client = Client::connect(&host).await.unwrap();
286 let err = client
287 .browse("S".into(), false, String::new(), 1000)
288 .await
289 .unwrap_err();
290 assert!(matches!(err, Error::Rpc(_)));
291 }
292
293 #[tokio::test]
294 async fn test_browse_stream_error_after_items() {
295 let svc = MockBridgeService {
296 browse_responses: vec![BrowseResponse {
297 tag_id: "tag1".into(),
298 node_type: "Leaf".into(),
299 }],
300 browse_stream_error: Some(Status::internal("stream broke")),
301 ..Default::default()
302 };
303 let host = start_mock_server(svc).await;
304 let mut client = Client::connect(&host).await.unwrap();
305 let err = client
306 .browse("S".into(), false, String::new(), 1000)
307 .await
308 .unwrap_err();
309 assert!(matches!(err, Error::Rpc(_)));
310 }
311
312 #[tokio::test]
313 async fn test_browse_drop_stops_server_send_loop() {
314 let svc = MockBridgeService {
317 browse_responses: (0..300)
318 .map(|i| BrowseResponse {
319 tag_id: format!("tag{i}"),
320 node_type: "Leaf".into(),
321 })
322 .collect(),
323 ..Default::default()
324 };
325 let browse_send_failure = Arc::clone(&svc.browse_send_failure);
326 let response = svc
327 .browse(Request::new(BrowseRequest {
328 server: "S".into(),
329 flat: false,
330 path: String::new(),
331 max_tags: 1000,
332 }))
333 .await
334 .unwrap();
335 drop(response);
336 tokio::time::timeout(Duration::from_secs(1), browse_send_failure.notified())
337 .await
338 .expect("mock sender did not observe the dropped browse stream");
339 }
340
341 #[tokio::test]
342 async fn test_read_empty() {
343 let host = start_mock_server(MockBridgeService::default()).await;
344 let mut client = Client::connect(&host).await.unwrap();
345 let values = client.read("S".into(), vec![]).await.unwrap();
346 assert!(values.is_empty());
347 }
348
349 #[tokio::test]
350 async fn test_read_with_data_maps_fields() {
351 let svc = MockBridgeService {
352 read_response: ReadResponse {
353 values: vec![ProtoTagValue {
354 tag_id: "t1".into(),
355 value: "42".into(),
356 quality: "Good".into(),
357 timestamp: "now".into(),
358 }],
359 },
360 ..Default::default()
361 };
362 let host = start_mock_server(svc).await;
363 let mut client = Client::connect(&host).await.unwrap();
364 let values = client.read("S".into(), vec!["t1".into()]).await.unwrap();
365 assert_eq!(
366 values,
367 vec![TagValue {
368 tag_id: "t1".into(),
369 value: "42".into(),
370 quality: "Good".into(),
371 timestamp: "now".into(),
372 }]
373 );
374 }
375
376 #[tokio::test]
377 async fn test_read_rpc_error() {
378 let svc = MockBridgeService {
379 read_error: Some(Status::internal("boom")),
380 ..Default::default()
381 };
382 let host = start_mock_server(svc).await;
383 let mut client = Client::connect(&host).await.unwrap();
384 let err = client.read("S".into(), vec![]).await.unwrap_err();
385 assert!(matches!(err, Error::Rpc(_)));
386 }
387
388 #[tokio::test]
389 async fn test_write_bool_value() {
390 let host = start_mock_server(MockBridgeService::default()).await;
391 let mut client = Client::connect(&host).await.unwrap();
392 client
393 .write("S".into(), "tag1".into(), Value::Bool(true))
394 .await
395 .unwrap();
396 }
397
398 #[tokio::test]
399 async fn test_write_int_value() {
400 let host = start_mock_server(MockBridgeService::default()).await;
401 let mut client = Client::connect(&host).await.unwrap();
402 client
403 .write("S".into(), "tag1".into(), Value::Int(42))
404 .await
405 .unwrap();
406 }
407
408 #[tokio::test]
409 async fn test_write_float_value() {
410 let host = start_mock_server(MockBridgeService::default()).await;
411 let mut client = Client::connect(&host).await.unwrap();
412 client
413 .write("S".into(), "tag1".into(), Value::Float(9.5))
414 .await
415 .unwrap();
416 }
417
418 #[tokio::test]
419 async fn test_write_string_value() {
420 let host = start_mock_server(MockBridgeService::default()).await;
421 let mut client = Client::connect(&host).await.unwrap();
422 client
423 .write(
424 "S".into(),
425 "tag1".into(),
426 Value::String("hello world".into()),
427 )
428 .await
429 .unwrap();
430 }
431
432 #[tokio::test]
433 async fn test_write_maps_success_result() {
434 let svc = MockBridgeService {
435 write_response: WriteResponse {
436 tag_id: "t1".into(),
437 success: true,
438 error: None,
439 },
440 ..Default::default()
441 };
442 let host = start_mock_server(svc).await;
443 let mut client = Client::connect(&host).await.unwrap();
444 let result = client
445 .write("S".into(), "t1".into(), Value::Int(1))
446 .await
447 .unwrap();
448 assert_eq!(
449 result,
450 WriteResult {
451 tag_id: "t1".into(),
452 success: true,
453 error: None,
454 }
455 );
456 }
457
458 #[tokio::test]
459 async fn test_write_maps_failure_result_with_error() {
460 let svc = MockBridgeService {
461 write_response: WriteResponse {
462 tag_id: "bad".into(),
463 success: false,
464 error: Some("access denied".into()),
465 },
466 ..Default::default()
467 };
468 let host = start_mock_server(svc).await;
469 let mut client = Client::connect(&host).await.unwrap();
470 let result = client
471 .write("S".into(), "bad".into(), Value::Int(0))
472 .await
473 .unwrap();
474 assert_eq!(
475 result,
476 WriteResult {
477 tag_id: "bad".into(),
478 success: false,
479 error: Some("access denied".into()),
480 }
481 );
482 }
483
484 #[tokio::test]
485 async fn test_write_rpc_error() {
486 let svc = MockBridgeService {
487 write_error: Some(Status::internal("boom")),
488 ..Default::default()
489 };
490 let host = start_mock_server(svc).await;
491 let mut client = Client::connect(&host).await.unwrap();
492 let err = client
493 .write("S".into(), "t1".into(), Value::Int(1))
494 .await
495 .unwrap_err();
496 assert!(matches!(err, Error::Rpc(_)));
497 }
498}