Skip to main content

opcda_bridge/
client.rs

1//! The connected gRPC client: [`Client`] and its typed methods.
2
3use 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/// A connected client for an opcda-bridge gateway's gRPC API.
11///
12/// Every method takes `&mut self`, matching the generated `BridgeClient`'s
13/// own requirement (it buffers per-call codec state); the underlying
14/// `tonic` channel itself is a cheap-to-reuse, multiplexed HTTP/2
15/// connection, so a single `Client` is meant to be held and reused across
16/// many calls rather than reconnected per request (unlike
17/// `opcda-bridge-client`'s CLI, which is a fresh process per invocation and
18/// so never notices the difference).
19#[derive(Debug)]
20pub struct Client {
21    inner: BridgeClient<Channel>,
22}
23
24impl Client {
25    /// Connect to a gateway at `host` (e.g. `"localhost:7600"`).
26    ///
27    /// Matches `opcda-bridge-client`'s long-standing `http://{host}` scheme
28    /// assumption: the gateway only ever serves plaintext HTTP/2 (no TLS).
29    pub async fn connect(host: &str) -> Result<Self> {
30        let inner = BridgeClient::connect(format!("http://{host}")).await?;
31        Ok(Self { inner })
32    }
33
34    /// List the OPC DA servers registered on the gateway's host.
35    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    /// Browse one level of `server`'s tag tree rooted at `path` (empty for
46    /// the top level), fully materialized into a `Vec` rather than a raw
47    /// stream. Every current caller (the CLI, and any bhtune-style
48    /// consumer) wants a complete result before doing anything else, so
49    /// this drains the stream internally rather than exposing it, sparing
50    /// callers a dependency on `tokio-stream`/`futures` just to consume it.
51    ///
52    /// `flat` and `max_tags` are forwarded to the gateway unchanged; the
53    /// gateway alone decides how they affect what's returned (e.g. `flat`
54    /// yielding every descendant tag instead of one level). `path` selects
55    /// which branch of the tree to browse.
56    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    /// Read one or more tag values from `server`.
85    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    /// Write `value` to `tag` on `server`.
107    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        // `opcda-bridge-client`'s commands convert this crate's `Error`
180        // into `anyhow::Error` via a bare `?`; this must render identically
181        // to today's direct `tonic::transport::Error` -> `anyhow::Error`
182        // conversion (the connect helper's pre-this-crate implementation),
183        // or the CLI's printed error text would silently change.
184        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        // Drop the service response directly instead of relying on gRPC
315        // buffering to propagate a remote client disconnect.
316        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}