Skip to main content

opcda_bridge/
client.rs

1//! The connected gRPC client and typed search stream.
2
3use crate::error::{Error, Result};
4use crate::types::{
5    BrowsePage, BrowsePageRequest, Capabilities, SearchEvent, SearchIndexControlAction,
6    SearchIndexRequest, SearchIndexResponse, SearchIndexStatus, SearchRequest, TagValue, Value,
7    WriteResult,
8};
9use opcda_bridge_proto::bridge::bridge_client::BridgeClient;
10use opcda_bridge_proto::bridge::write_request::TypedValue;
11use opcda_bridge_proto::bridge::{
12    CloseBrowseSessionRequest, ControlSearchIndexRequest, GetCapabilitiesRequest,
13    GetSearchIndexStatusRequest, ListServersRequest, ReadRequest, RefreshSearchIndexRequest,
14    WriteRequest,
15};
16use tonic::Code;
17use tonic::codec::Streaming;
18use tonic::transport::Channel;
19
20/// A connected client for an opcda-bridge gateway's gRPC API.
21#[derive(Debug)]
22pub struct Client {
23    inner: BridgeClient<Channel>,
24}
25
26/// A cancellable stream of typed namespace-search events.
27///
28/// Dropping this value drops the underlying gRPC stream, allowing the gateway
29/// to stop scheduling further search work.
30#[derive(Debug)]
31pub struct SearchStream {
32    inner: Streaming<opcda_bridge_proto::bridge::SearchEvent>,
33}
34
35impl SearchStream {
36    /// Wait for the next event. `None` means the server closed the stream.
37    pub async fn message(&mut self) -> Result<Option<SearchEvent>> {
38        self.inner
39            .message()
40            .await?
41            .map(SearchEvent::try_from)
42            .transpose()
43    }
44}
45
46impl Client {
47    /// Connect to a plaintext gateway at `host` (for example, `localhost:7600`).
48    pub async fn connect(host: &str) -> Result<Self> {
49        let inner = BridgeClient::connect(format!("http://{host}")).await?;
50        Ok(Self { inner })
51    }
52
53    /// Report protocol, paging, browse-session, search, and namespace support.
54    pub async fn capabilities(&mut self, server: impl Into<String>) -> Result<Capabilities> {
55        self.inner
56            .get_capabilities(GetCapabilitiesRequest {
57                server: server.into(),
58            })
59            .await
60            .map_err(|status| feature_error("capability discovery", status))?
61            .into_inner()
62            .try_into()
63    }
64
65    /// List the OPC DA servers registered on the gateway's host.
66    pub async fn list_servers(&mut self) -> Result<Vec<String>> {
67        let response = self
68            .inner
69            .list_servers(ListServersRequest {
70                host: "localhost".to_string(),
71            })
72            .await?;
73        Ok(response.into_inner().servers)
74    }
75
76    /// Open a browse session and return only its first root page.
77    pub async fn browse(
78        &mut self,
79        server: impl Into<String>,
80        page_size: u32,
81    ) -> Result<BrowsePage> {
82        self.open_browse(server, page_size).await
83    }
84
85    /// Open a browse session and return only its first root page.
86    pub async fn open_browse(
87        &mut self,
88        server: impl Into<String>,
89        page_size: u32,
90    ) -> Result<BrowsePage> {
91        self.browse_page(BrowsePageRequest::root(server, page_size))
92            .await
93    }
94
95    /// Request exactly one root, child, or continuation page.
96    ///
97    /// This method never follows `next_page_token` automatically.
98    pub async fn browse_page(&mut self, request: BrowsePageRequest) -> Result<BrowsePage> {
99        self.inner
100            .browse(opcda_bridge_proto::bridge::BrowseRequest::from(request))
101            .await
102            .map_err(|status| feature_error("paged browse", status))?
103            .into_inner()
104            .try_into()
105    }
106
107    /// Explicitly release a gateway browse session.
108    pub async fn close_browse_session(&mut self, session_id: impl Into<String>) -> Result<()> {
109        self.inner
110            .close_browse_session(CloseBrowseSessionRequest {
111                session_id: session_id.into(),
112            })
113            .await
114            .map_err(|status| feature_error("browse-session close", status))?;
115        Ok(())
116    }
117
118    /// Start a bounded search and return its progressive event stream.
119    pub async fn search_stream(&mut self, request: SearchRequest) -> Result<SearchStream> {
120        let inner = self
121            .inner
122            .search(opcda_bridge_proto::bridge::SearchRequest::from(request))
123            .await
124            .map_err(|status| feature_error("namespace search", status))?
125            .into_inner();
126        Ok(SearchStream { inner })
127    }
128
129    /// Explicitly collect a complete search stream into memory.
130    pub async fn search(&mut self, request: SearchRequest) -> Result<Vec<SearchEvent>> {
131        let mut stream = self.search_stream(request).await?;
132        let mut events = Vec::new();
133        while let Some(event) = stream.message().await? {
134            events.push(event);
135        }
136        Ok(events)
137    }
138
139    /// Return the persistent namespace-index status for `server`.
140    pub async fn search_index_status(
141        &mut self,
142        server: impl Into<String>,
143    ) -> Result<SearchIndexStatus> {
144        self.inner
145            .get_search_index_status(GetSearchIndexStatusRequest {
146                server: server.into(),
147            })
148            .await
149            .map_err(|status| feature_error("indexed-search status", status))?
150            .into_inner()
151            .try_into()
152    }
153
154    /// Start or coalesce a persistent namespace-index refresh.
155    pub async fn refresh_search_index(
156        &mut self,
157        server: impl Into<String>,
158        force: bool,
159    ) -> Result<SearchIndexStatus> {
160        self.inner
161            .refresh_search_index(RefreshSearchIndexRequest {
162                server: server.into(),
163                force,
164            })
165            .await
166            .map_err(|status| feature_error("indexed-search refresh", status))?
167            .into_inner()
168            .try_into()
169    }
170
171    /// Pause, resume, or cancel a persistent namespace-index build.
172    pub async fn control_search_index(
173        &mut self,
174        server: impl Into<String>,
175        action: SearchIndexControlAction,
176    ) -> Result<SearchIndexStatus> {
177        self.inner
178            .control_search_index(ControlSearchIndexRequest {
179                server: server.into(),
180                action: opcda_bridge_proto::bridge::SearchIndexControlAction::from(action) as i32,
181            })
182            .await
183            .map_err(|status| feature_error("indexed-search control", status))?
184            .into_inner()
185            .try_into()
186    }
187
188    /// Search the gateway-owned persistent namespace index.
189    ///
190    /// This never falls back to live namespace traversal.
191    pub async fn search_index(
192        &mut self,
193        request: SearchIndexRequest,
194    ) -> Result<SearchIndexResponse> {
195        self.inner
196            .search_index(opcda_bridge_proto::bridge::SearchIndexRequest::from(
197                request,
198            ))
199            .await
200            .map_err(|status| feature_error("indexed search", status))?
201            .into_inner()
202            .try_into()
203    }
204
205    /// Read one or more exact OPC DA ItemIDs from `server`.
206    pub async fn read(&mut self, server: String, tags: Vec<String>) -> Result<Vec<TagValue>> {
207        let response = self
208            .inner
209            .read(ReadRequest {
210                server,
211                tag_ids: tags,
212            })
213            .await?;
214        Ok(response
215            .into_inner()
216            .values
217            .into_iter()
218            .map(|v| TagValue {
219                tag_id: v.tag_id,
220                value: v.value,
221                quality: v.quality,
222                timestamp: v.timestamp,
223            })
224            .collect())
225    }
226
227    /// Write `value` to one exact OPC DA ItemID on `server`.
228    pub async fn write(
229        &mut self,
230        server: String,
231        tag: String,
232        value: Value,
233    ) -> Result<WriteResult> {
234        let typed_value = match value {
235            Value::String(s) => TypedValue::StringValue(s),
236            Value::Int(i) => TypedValue::IntValue(i),
237            Value::Float(f) => TypedValue::FloatValue(f),
238            Value::Bool(b) => TypedValue::BoolValue(b),
239        };
240        let response = self
241            .inner
242            .write(WriteRequest {
243                server,
244                tag_id: tag,
245                typed_value: Some(typed_value),
246            })
247            .await?;
248        let result = response.into_inner();
249        Ok(WriteResult {
250            tag_id: result.tag_id,
251            success: result.success,
252            error: result.error,
253        })
254    }
255}
256
257fn feature_error(operation: &'static str, status: tonic::Status) -> Error {
258    if status.code() == Code::Unimplemented {
259        Error::IncompatibleGateway { operation }
260    } else {
261        Error::Rpc(status)
262    }
263}
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268    use crate::error::Error;
269    use crate::test_support::{MockBridgeService, start_mock_server};
270    use crate::{
271        BrowseNodeKind, BrowseSource, NamespaceOrganization, SearchIndexControlAction,
272        SearchIndexRequest, SearchIndexState, SearchMatchMode,
273    };
274    use opcda_bridge_proto::bridge::search_event;
275    use opcda_bridge_proto::bridge::{
276        BrowseNode as ProtoBrowseNode, BrowsePage as ProtoBrowsePage,
277        BrowseSource as ProtoBrowseSource, GetCapabilitiesResponse, IndexedSearchMatch,
278        IndexedSearchProgress, ListServersResponse, NamespaceOrganization as ProtoOrganization,
279        ReadResponse, SearchCompleted, SearchEvent as ProtoSearchEvent, SearchIndexResponse,
280        SearchIndexState as ProtoSearchIndexState, SearchIndexStatus, SearchProgress,
281        TagValue as ProtoTagValue, WriteResponse,
282    };
283    use std::sync::Arc;
284    use std::time::Duration;
285    use tonic::Status;
286
287    fn item_node() -> ProtoBrowseNode {
288        ProtoBrowseNode {
289            node_key: "node".into(),
290            display_name: "PV".into(),
291            kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
292            item_id: Some("FCS!TAG.PV".into()),
293        }
294    }
295
296    #[tokio::test]
297    async fn connect_success_and_failure_are_typed() {
298        let host = start_mock_server(MockBridgeService::default()).await;
299        Client::connect(&host).await.unwrap();
300        assert!(matches!(
301            Client::connect("127.0.0.1:1").await.unwrap_err(),
302            Error::Connect(_)
303        ));
304    }
305
306    #[tokio::test]
307    async fn mock_server_shutdown_completes() {
308        let service = MockBridgeService::default();
309        let shutdown = Arc::clone(&service.server_shutdown);
310        let stopped = Arc::clone(&service.server_stopped);
311        let _host = start_mock_server(service).await;
312        shutdown.notify_one();
313        tokio::time::timeout(Duration::from_secs(1), stopped.notified())
314            .await
315            .unwrap();
316    }
317
318    #[tokio::test]
319    async fn capabilities_maps_fields_and_request() {
320        let service = MockBridgeService {
321            capabilities_response: GetCapabilitiesResponse {
322                application_version: "0.3.0".into(),
323                protocol_version: "0.3".into(),
324                max_page_size: 1000,
325                supports_browse_sessions: true,
326                supports_search: true,
327                organization: ProtoOrganization::Hierarchical as i32,
328                source: ProtoBrowseSource::Da2 as i32,
329                supports_indexed_search: true,
330                indexed_search_protocol_version: "1".into(),
331                max_indexed_search_results: 50,
332                search_index_state: ProtoSearchIndexState::Ready as i32,
333            },
334            ..Default::default()
335        };
336        let requests = Arc::clone(&service.capabilities_requests);
337        let host = start_mock_server(service).await;
338        let mut client = Client::connect(&host).await.unwrap();
339        let capabilities = client.capabilities("S").await.unwrap();
340        assert_eq!(
341            capabilities.organization,
342            NamespaceOrganization::Hierarchical
343        );
344        assert_eq!(capabilities.source, BrowseSource::Da2);
345        assert!(capabilities.supports_indexed_search);
346        assert_eq!(capabilities.indexed_search_protocol_version, "1");
347        assert_eq!(capabilities.max_indexed_search_results, 50);
348        assert_eq!(capabilities.search_index_state, SearchIndexState::Ready);
349        assert_eq!(requests.lock().unwrap()[0].server, "S");
350    }
351
352    #[tokio::test]
353    async fn capabilities_rpc_error_is_typed() {
354        let host = start_mock_server(MockBridgeService {
355            capabilities_error: Some(Status::unimplemented("old gateway")),
356            ..Default::default()
357        })
358        .await;
359        let mut client = Client::connect(&host).await.unwrap();
360        assert!(matches!(
361            client.capabilities("S").await.unwrap_err(),
362            Error::IncompatibleGateway { .. }
363        ));
364    }
365
366    #[tokio::test]
367    async fn list_servers_maps_data_and_errors() {
368        let host = start_mock_server(MockBridgeService {
369            list_servers_response: ListServersResponse {
370                servers: vec!["S1".into(), "S2".into()],
371            },
372            ..Default::default()
373        })
374        .await;
375        let mut client = Client::connect(&host).await.unwrap();
376        assert_eq!(client.list_servers().await.unwrap(), ["S1", "S2"]);
377
378        let host = start_mock_server(MockBridgeService {
379            list_servers_error: Some(Status::internal("boom")),
380            ..Default::default()
381        })
382        .await;
383        let mut client = Client::connect(&host).await.unwrap();
384        assert!(matches!(
385            client.list_servers().await.unwrap_err(),
386            Error::Rpc(_)
387        ));
388    }
389
390    #[tokio::test]
391    async fn browse_returns_one_typed_page_without_draining() {
392        let service = MockBridgeService {
393            browse_response: ProtoBrowsePage {
394                session_id: "session".into(),
395                nodes: vec![item_node()],
396                next_page_token: Some("next".into()),
397                complete: false,
398                organization: ProtoOrganization::Hierarchical as i32,
399                source: ProtoBrowseSource::Da3 as i32,
400                warning: None,
401            },
402            ..Default::default()
403        };
404        let requests = Arc::clone(&service.browse_requests);
405        let host = start_mock_server(service).await;
406        let mut client = Client::connect(&host).await.unwrap();
407        let page = client.browse("S", 25).await.unwrap();
408        assert_eq!(page.nodes[0].kind, BrowseNodeKind::Item);
409        assert_eq!(page.next_page_token.as_deref(), Some("next"));
410        let request = &requests.lock().unwrap()[0];
411        assert_eq!(request.page_size, 25);
412        assert!(request.session_id.is_none());
413    }
414
415    #[tokio::test]
416    async fn browse_page_forwards_session_parent_token_and_refresh() {
417        let service = MockBridgeService::default();
418        let requests = Arc::clone(&service.browse_requests);
419        let host = start_mock_server(service).await;
420        let mut client = Client::connect(&host).await.unwrap();
421        client
422            .browse_page(
423                BrowsePageRequest::next("S", "session", Some("parent".into()), "token", 10)
424                    .with_refresh(true),
425            )
426            .await
427            .unwrap();
428        let request = &requests.lock().unwrap()[0];
429        assert_eq!(request.session_id.as_deref(), Some("session"));
430        assert_eq!(request.parent_node_key.as_deref(), Some("parent"));
431        assert_eq!(request.page_token.as_deref(), Some("token"));
432        assert!(request.refresh);
433    }
434
435    #[tokio::test]
436    async fn browse_rpc_and_protocol_errors_are_typed() {
437        let host = start_mock_server(MockBridgeService {
438            browse_error: Some(Status::failed_precondition("expired")),
439            ..Default::default()
440        })
441        .await;
442        let mut client = Client::connect(&host).await.unwrap();
443        assert!(matches!(
444            client.open_browse("S", 20).await.unwrap_err(),
445            Error::Rpc(_)
446        ));
447
448        let host = start_mock_server(MockBridgeService {
449            browse_response: ProtoBrowsePage::default(),
450            ..Default::default()
451        })
452        .await;
453        let mut client = Client::connect(&host).await.unwrap();
454        assert!(matches!(
455            client.open_browse("S", 20).await.unwrap_err(),
456            Error::Protocol(_)
457        ));
458
459        let host = start_mock_server(MockBridgeService {
460            browse_error: Some(Status::unimplemented("old gateway")),
461            ..Default::default()
462        })
463        .await;
464        let mut client = Client::connect(&host).await.unwrap();
465        assert!(matches!(
466            client.open_browse("S", 20).await.unwrap_err(),
467            Error::IncompatibleGateway { .. }
468        ));
469    }
470
471    #[tokio::test]
472    async fn close_browse_session_forwards_id_and_error() {
473        let service = MockBridgeService::default();
474        let requests = Arc::clone(&service.close_requests);
475        let host = start_mock_server(service).await;
476        let mut client = Client::connect(&host).await.unwrap();
477        client.close_browse_session("session").await.unwrap();
478        assert_eq!(requests.lock().unwrap()[0].session_id, "session");
479
480        let host = start_mock_server(MockBridgeService {
481            close_error: Some(Status::not_found("missing")),
482            ..Default::default()
483        })
484        .await;
485        let mut client = Client::connect(&host).await.unwrap();
486        assert!(matches!(
487            client.close_browse_session("missing").await.unwrap_err(),
488            Error::Rpc(_)
489        ));
490    }
491
492    #[tokio::test]
493    async fn search_stream_and_collect_map_events_and_request() {
494        let events = vec![
495            ProtoSearchEvent {
496                event: Some(search_event::Event::Progress(SearchProgress {
497                    visited_nodes: 5,
498                    matches: 0,
499                    partial: true,
500                })),
501            },
502            ProtoSearchEvent {
503                event: Some(search_event::Event::Completed(SearchCompleted {
504                    complete: true,
505                    cancelled: false,
506                    truncated: false,
507                    warning: None,
508                })),
509            },
510        ];
511        let service = MockBridgeService {
512            search_events: events,
513            ..Default::default()
514        };
515        let requests = Arc::clone(&service.search_requests);
516        let host = start_mock_server(service).await;
517        let mut client = Client::connect(&host).await.unwrap();
518        let mut request = SearchRequest::new("S", "PV", SearchMatchMode::Prefix);
519        request.session_id = Some("session".into());
520        request.scope_node_key = Some("scope".into());
521        request.max_results = 50;
522        request.include_branches = true;
523        request.refresh = true;
524        let found = client.search(request).await.unwrap();
525        assert_eq!(found.len(), 2);
526        let request = &requests.lock().unwrap()[0];
527        assert_eq!(request.query, "PV");
528        assert_eq!(
529            request.match_mode,
530            opcda_bridge_proto::bridge::SearchMatchMode::Prefix as i32
531        );
532        assert!(request.include_branches);
533        assert!(request.refresh);
534    }
535
536    #[tokio::test]
537    async fn search_initial_stream_and_protocol_errors_are_typed() {
538        let host = start_mock_server(MockBridgeService {
539            search_initial_error: Some(Status::unavailable("down")),
540            ..Default::default()
541        })
542        .await;
543        let mut client = Client::connect(&host).await.unwrap();
544        assert!(matches!(
545            client
546                .search_stream(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
547                .await
548                .unwrap_err(),
549            Error::Rpc(_)
550        ));
551
552        let host = start_mock_server(MockBridgeService {
553            search_stream_error: Some(Status::deadline_exceeded("slow")),
554            ..Default::default()
555        })
556        .await;
557        let mut client = Client::connect(&host).await.unwrap();
558        assert!(matches!(
559            client
560                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
561                .await
562                .unwrap_err(),
563            Error::Rpc(_)
564        ));
565
566        let host = start_mock_server(MockBridgeService {
567            search_events: vec![ProtoSearchEvent::default()],
568            ..Default::default()
569        })
570        .await;
571        let mut client = Client::connect(&host).await.unwrap();
572        assert!(matches!(
573            client
574                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
575                .await
576                .unwrap_err(),
577            Error::Protocol(_)
578        ));
579
580        assert!(matches!(
581            feature_error("test", Status::unimplemented("old")),
582            Error::IncompatibleGateway { .. }
583        ));
584    }
585
586    fn index_status(state: ProtoSearchIndexState) -> SearchIndexStatus {
587        SearchIndexStatus {
588            server: "S".into(),
589            state: state as i32,
590            configured: true,
591            active_generation: 4,
592            entry_count: 100,
593            unique_item_count: 99,
594            started_at: Some("start".into()),
595            completed_at: Some("complete".into()),
596            last_error: None,
597            database_bytes: 1024,
598            organization: ProtoOrganization::Hierarchical as i32,
599            source: ProtoBrowseSource::Da3 as i32,
600            progress: Some(IndexedSearchProgress {
601                branches_visited: 2,
602                entries_seen: 3,
603                unique_items: 3,
604                active_time_ms: 4,
605                paused_time_ms: 5,
606                items_per_second: 6.0,
607                estimated_remaining_ms: Some(7),
608            }),
609        }
610    }
611
612    #[tokio::test]
613    async fn indexed_search_methods_map_requests_responses_and_errors() {
614        let service = MockBridgeService {
615            search_index_status_response: index_status(ProtoSearchIndexState::Ready),
616            refresh_search_index_response: index_status(ProtoSearchIndexState::Refreshing),
617            control_search_index_response: index_status(ProtoSearchIndexState::Partial),
618            search_index_response: SearchIndexResponse {
619                matches: vec![IndexedSearchMatch {
620                    item_id: "Exact.ItemID".into(),
621                    display_name: "PV".into(),
622                    kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
623                    breadcrumbs: vec!["Area".into()],
624                }],
625                has_more: true,
626                status: Some(index_status(ProtoSearchIndexState::Stale)),
627            },
628            ..Default::default()
629        };
630        let status_requests = Arc::clone(&service.search_index_status_requests);
631        let refresh_requests = Arc::clone(&service.refresh_search_index_requests);
632        let control_requests = Arc::clone(&service.control_search_index_requests);
633        let search_requests = Arc::clone(&service.search_index_requests);
634        let host = start_mock_server(service).await;
635        let mut client = Client::connect(&host).await.unwrap();
636
637        assert_eq!(
638            client.search_index_status("S").await.unwrap().state,
639            SearchIndexState::Ready
640        );
641        assert_eq!(
642            client.refresh_search_index("S", true).await.unwrap().state,
643            SearchIndexState::Refreshing
644        );
645        assert_eq!(
646            client
647                .control_search_index("S", SearchIndexControlAction::Pause)
648                .await
649                .unwrap()
650                .state,
651            SearchIndexState::Partial
652        );
653        let mut request = SearchIndexRequest::new("S", "PV", SearchMatchMode::Contains);
654        request.max_results = 25;
655        let response = client.search_index(request).await.unwrap();
656        assert_eq!(response.matches[0].item_id, "Exact.ItemID");
657        assert!(response.has_more);
658
659        assert_eq!(status_requests.lock().unwrap()[0].server, "S");
660        assert!(refresh_requests.lock().unwrap()[0].force);
661        assert_eq!(
662            control_requests.lock().unwrap()[0].action,
663            opcda_bridge_proto::bridge::SearchIndexControlAction::Pause as i32
664        );
665        assert_eq!(search_requests.lock().unwrap()[0].max_results, 25);
666
667        for (field, operation) in [
668            ("search_index_status_error", "status"),
669            ("refresh_search_index_error", "refresh"),
670            ("control_search_index_error", "control"),
671            ("search_index_error", "search"),
672        ] {
673            let mut service = MockBridgeService::default();
674            let error = Some(Status::unimplemented("old gateway"));
675            match field {
676                "search_index_status_error" => service.search_index_status_error = error,
677                "refresh_search_index_error" => service.refresh_search_index_error = error,
678                "control_search_index_error" => service.control_search_index_error = error,
679                "search_index_error" => service.search_index_error = error,
680                _ => unreachable!(),
681            }
682            let host = start_mock_server(service).await;
683            let mut client = Client::connect(&host).await.unwrap();
684            let error = match operation {
685                "status" => client.search_index_status("S").await.unwrap_err(),
686                "refresh" => client.refresh_search_index("S", false).await.unwrap_err(),
687                "control" => client
688                    .control_search_index("S", SearchIndexControlAction::Cancel)
689                    .await
690                    .unwrap_err(),
691                "search" => client
692                    .search_index(SearchIndexRequest::new(
693                        "S",
694                        "PV",
695                        SearchMatchMode::Contains,
696                    ))
697                    .await
698                    .unwrap_err(),
699                _ => unreachable!(),
700            };
701            assert!(matches!(error, Error::IncompatibleGateway { .. }));
702        }
703    }
704
705    #[tokio::test]
706    async fn read_maps_data_and_errors() {
707        let host = start_mock_server(MockBridgeService {
708            read_response: ReadResponse {
709                values: vec![ProtoTagValue {
710                    tag_id: "t1".into(),
711                    value: "42".into(),
712                    quality: "Good".into(),
713                    timestamp: "now".into(),
714                }],
715            },
716            ..Default::default()
717        })
718        .await;
719        let mut client = Client::connect(&host).await.unwrap();
720        assert_eq!(
721            client.read("S".into(), vec![]).await.unwrap()[0].value,
722            "42"
723        );
724
725        let host = start_mock_server(MockBridgeService {
726            read_error: Some(Status::internal("boom")),
727            ..Default::default()
728        })
729        .await;
730        let mut client = Client::connect(&host).await.unwrap();
731        assert!(matches!(
732            client.read("S".into(), vec![]).await.unwrap_err(),
733            Error::Rpc(_)
734        ));
735    }
736
737    #[tokio::test]
738    async fn write_maps_every_value_and_result_or_error() {
739        for value in [
740            Value::Bool(true),
741            Value::Int(42),
742            Value::Float(3.5),
743            Value::String("text".into()),
744        ] {
745            let host = start_mock_server(MockBridgeService {
746                write_response: WriteResponse {
747                    tag_id: "t".into(),
748                    success: true,
749                    error: None,
750                },
751                ..Default::default()
752            })
753            .await;
754            let mut client = Client::connect(&host).await.unwrap();
755            assert!(
756                client
757                    .write("S".into(), "t".into(), value)
758                    .await
759                    .unwrap()
760                    .success
761            );
762        }
763
764        let host = start_mock_server(MockBridgeService {
765            write_error: Some(Status::internal("boom")),
766            ..Default::default()
767        })
768        .await;
769        let mut client = Client::connect(&host).await.unwrap();
770        assert!(matches!(
771            client
772                .write("S".into(), "t".into(), Value::Int(1))
773                .await
774                .unwrap_err(),
775            Error::Rpc(_)
776        ));
777    }
778}