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