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                search_index_promoting: true,
386            },
387            ..Default::default()
388        };
389        let requests = Arc::clone(&service.capabilities_requests);
390        let host = start_mock_server(service).await;
391        let mut client = Client::connect(&host).await.unwrap();
392        let capabilities = client.capabilities("S").await.unwrap();
393        assert_eq!(
394            capabilities.organization,
395            NamespaceOrganization::Hierarchical
396        );
397        assert_eq!(capabilities.source, BrowseSource::Da2);
398        assert!(capabilities.supports_indexed_search);
399        assert_eq!(capabilities.indexed_search_protocol_version, "1");
400        assert_eq!(capabilities.max_indexed_search_results, 50);
401        assert_eq!(capabilities.search_index_state, SearchIndexState::Ready);
402        assert!(capabilities.search_index_promoting);
403        assert_eq!(requests.lock().unwrap()[0].server, "S");
404    }
405
406    #[tokio::test]
407    async fn gateway_info_and_compatibility_report_are_typed() {
408        let service = MockBridgeService {
409            gateway_info_response: GetGatewayInfoResponse {
410                application_version: "0.4.3".into(),
411                compatibility_schema_version: 1,
412                features: vec![
413                    ProtocolFeature {
414                        kind: ProtocolFeatureKind::Core as i32,
415                        min_version: 1,
416                        max_version: 1,
417                    },
418                    ProtocolFeature {
419                        kind: ProtocolFeatureKind::Namespace as i32,
420                        min_version: 2,
421                        max_version: 2,
422                    },
423                    ProtocolFeature {
424                        kind: ProtocolFeatureKind::IndexedSearch as i32,
425                        min_version: 1,
426                        max_version: 1,
427                    },
428                ],
429            },
430            ..Default::default()
431        };
432        let requests = Arc::clone(&service.gateway_info_requests);
433        let host = start_mock_server(service).await;
434        let mut client = Client::connect(&host).await.unwrap();
435        let info = client.gateway_info().await.unwrap();
436        assert_eq!(info.application_version, "0.4.3");
437        let report = client
438            .compatibility_with_client_version(None, "0.4.3")
439            .await
440            .unwrap();
441        assert_eq!(report.status, crate::CompatibilityStatus::Full);
442        assert_eq!(report.library_version, env!("CARGO_PKG_VERSION"));
443        assert_eq!(requests.lock().unwrap().len(), 2);
444    }
445
446    #[tokio::test]
447    async fn compatibility_wrapper_and_gateway_info_errors_are_typed() {
448        let host = start_mock_server(MockBridgeService {
449            gateway_info_response: GetGatewayInfoResponse {
450                application_version: "0.4.3".into(),
451                compatibility_schema_version: 1,
452                features: vec![
453                    ProtocolFeature {
454                        kind: ProtocolFeatureKind::Core as i32,
455                        min_version: 1,
456                        max_version: 1,
457                    },
458                    ProtocolFeature {
459                        kind: ProtocolFeatureKind::Namespace as i32,
460                        min_version: 2,
461                        max_version: 2,
462                    },
463                    ProtocolFeature {
464                        kind: ProtocolFeatureKind::IndexedSearch as i32,
465                        min_version: 1,
466                        max_version: 1,
467                    },
468                ],
469            },
470            ..Default::default()
471        })
472        .await;
473        let mut client = Client::connect(&host).await.unwrap();
474        assert_eq!(
475            client.compatibility(None).await.unwrap().status,
476            crate::CompatibilityStatus::Full
477        );
478
479        let host = start_mock_server(MockBridgeService {
480            gateway_info_error: Some(Status::internal("gateway unavailable")),
481            ..Default::default()
482        })
483        .await;
484        let mut client = Client::connect(&host).await.unwrap();
485        assert!(matches!(
486            client.compatibility(None).await.unwrap_err(),
487            Error::Rpc(_)
488        ));
489    }
490
491    #[tokio::test]
492    async fn compatibility_falls_back_to_legacy_or_reports_unknown() {
493        let service = MockBridgeService {
494            gateway_info_error: Some(Status::unimplemented("old gateway")),
495            capabilities_response: GetCapabilitiesResponse {
496                application_version: "0.3.2".into(),
497                protocol_version: "2".into(),
498                supports_indexed_search: false,
499                ..Default::default()
500            },
501            ..Default::default()
502        };
503        let host = start_mock_server(service).await;
504        let mut client = Client::connect(&host).await.unwrap();
505        let report = client
506            .compatibility_with_client_version(Some("S"), "0.4.3")
507            .await
508            .unwrap();
509        assert_eq!(
510            report.source,
511            crate::CompatibilitySource::LegacyCapabilities
512        );
513        assert_eq!(report.status, crate::CompatibilityStatus::Partial);
514
515        let host = start_mock_server(MockBridgeService {
516            gateway_info_error: Some(Status::unimplemented("old gateway")),
517            ..Default::default()
518        })
519        .await;
520        let mut client = Client::connect(&host).await.unwrap();
521        let report = client
522            .compatibility_with_client_version(None, "0.4.3")
523            .await
524            .unwrap();
525        assert_eq!(report.status, crate::CompatibilityStatus::Unknown);
526    }
527
528    #[tokio::test]
529    async fn capabilities_rpc_error_is_typed() {
530        let host = start_mock_server(MockBridgeService {
531            capabilities_error: Some(Status::unimplemented("old gateway")),
532            ..Default::default()
533        })
534        .await;
535        let mut client = Client::connect(&host).await.unwrap();
536        assert!(matches!(
537            client.capabilities("S").await.unwrap_err(),
538            Error::IncompatibleGateway { .. }
539        ));
540    }
541
542    #[tokio::test]
543    async fn list_servers_maps_data_and_errors() {
544        let host = start_mock_server(MockBridgeService {
545            list_servers_response: ListServersResponse {
546                servers: vec!["S1".into(), "S2".into()],
547            },
548            ..Default::default()
549        })
550        .await;
551        let mut client = Client::connect(&host).await.unwrap();
552        assert_eq!(client.list_servers().await.unwrap(), ["S1", "S2"]);
553
554        let host = start_mock_server(MockBridgeService {
555            list_servers_error: Some(Status::internal("boom")),
556            ..Default::default()
557        })
558        .await;
559        let mut client = Client::connect(&host).await.unwrap();
560        assert!(matches!(
561            client.list_servers().await.unwrap_err(),
562            Error::Rpc(_)
563        ));
564    }
565
566    #[tokio::test]
567    async fn browse_returns_one_typed_page_without_draining() {
568        let service = MockBridgeService {
569            browse_response: ProtoBrowsePage {
570                session_id: "session".into(),
571                nodes: vec![item_node()],
572                next_page_token: Some("next".into()),
573                complete: false,
574                organization: ProtoOrganization::Hierarchical as i32,
575                source: ProtoBrowseSource::Da3 as i32,
576                warning: None,
577            },
578            ..Default::default()
579        };
580        let requests = Arc::clone(&service.browse_requests);
581        let host = start_mock_server(service).await;
582        let mut client = Client::connect(&host).await.unwrap();
583        let page = client.browse("S", 25).await.unwrap();
584        assert_eq!(page.nodes[0].kind, BrowseNodeKind::Item);
585        assert_eq!(page.next_page_token.as_deref(), Some("next"));
586        let request = &requests.lock().unwrap()[0];
587        assert_eq!(request.page_size, 25);
588        assert!(request.session_id.is_none());
589    }
590
591    #[tokio::test]
592    async fn browse_page_forwards_session_parent_token_and_refresh() {
593        let service = MockBridgeService::default();
594        let requests = Arc::clone(&service.browse_requests);
595        let host = start_mock_server(service).await;
596        let mut client = Client::connect(&host).await.unwrap();
597        client
598            .browse_page(
599                BrowsePageRequest::next("S", "session", Some("parent".into()), "token", 10)
600                    .with_refresh(true),
601            )
602            .await
603            .unwrap();
604        let request = &requests.lock().unwrap()[0];
605        assert_eq!(request.session_id.as_deref(), Some("session"));
606        assert_eq!(request.parent_node_key.as_deref(), Some("parent"));
607        assert_eq!(request.page_token.as_deref(), Some("token"));
608        assert!(request.refresh);
609    }
610
611    #[tokio::test]
612    async fn browse_rpc_and_protocol_errors_are_typed() {
613        let host = start_mock_server(MockBridgeService {
614            browse_error: Some(Status::failed_precondition("expired")),
615            ..Default::default()
616        })
617        .await;
618        let mut client = Client::connect(&host).await.unwrap();
619        assert!(matches!(
620            client.open_browse("S", 20).await.unwrap_err(),
621            Error::Rpc(_)
622        ));
623
624        let host = start_mock_server(MockBridgeService {
625            browse_response: ProtoBrowsePage::default(),
626            ..Default::default()
627        })
628        .await;
629        let mut client = Client::connect(&host).await.unwrap();
630        assert!(matches!(
631            client.open_browse("S", 20).await.unwrap_err(),
632            Error::Protocol(_)
633        ));
634
635        let host = start_mock_server(MockBridgeService {
636            browse_error: Some(Status::unimplemented("old gateway")),
637            ..Default::default()
638        })
639        .await;
640        let mut client = Client::connect(&host).await.unwrap();
641        assert!(matches!(
642            client.open_browse("S", 20).await.unwrap_err(),
643            Error::IncompatibleGateway { .. }
644        ));
645    }
646
647    #[tokio::test]
648    async fn close_browse_session_forwards_id_and_error() {
649        let service = MockBridgeService::default();
650        let requests = Arc::clone(&service.close_requests);
651        let host = start_mock_server(service).await;
652        let mut client = Client::connect(&host).await.unwrap();
653        client.close_browse_session("session").await.unwrap();
654        assert_eq!(requests.lock().unwrap()[0].session_id, "session");
655
656        let host = start_mock_server(MockBridgeService {
657            close_error: Some(Status::not_found("missing")),
658            ..Default::default()
659        })
660        .await;
661        let mut client = Client::connect(&host).await.unwrap();
662        assert!(matches!(
663            client.close_browse_session("missing").await.unwrap_err(),
664            Error::Rpc(_)
665        ));
666    }
667
668    #[tokio::test]
669    async fn search_stream_and_collect_map_events_and_request() {
670        let events = vec![
671            ProtoSearchEvent {
672                event: Some(search_event::Event::Progress(SearchProgress {
673                    visited_nodes: 5,
674                    matches: 0,
675                    partial: true,
676                })),
677            },
678            ProtoSearchEvent {
679                event: Some(search_event::Event::Completed(SearchCompleted {
680                    complete: true,
681                    cancelled: false,
682                    truncated: false,
683                    warning: None,
684                })),
685            },
686        ];
687        let service = MockBridgeService {
688            search_events: events,
689            ..Default::default()
690        };
691        let requests = Arc::clone(&service.search_requests);
692        let host = start_mock_server(service).await;
693        let mut client = Client::connect(&host).await.unwrap();
694        let mut request = SearchRequest::new("S", "PV", SearchMatchMode::Prefix);
695        request.session_id = Some("session".into());
696        request.scope_node_key = Some("scope".into());
697        request.max_results = 50;
698        request.include_branches = true;
699        request.refresh = true;
700        let found = client.search(request).await.unwrap();
701        assert_eq!(found.len(), 2);
702        let request = &requests.lock().unwrap()[0];
703        assert_eq!(request.query, "PV");
704        assert_eq!(
705            request.match_mode,
706            opcda_bridge_proto::bridge::SearchMatchMode::Prefix as i32
707        );
708        assert!(request.include_branches);
709        assert!(request.refresh);
710    }
711
712    #[tokio::test]
713    async fn search_initial_stream_and_protocol_errors_are_typed() {
714        let host = start_mock_server(MockBridgeService {
715            search_initial_error: Some(Status::unavailable("down")),
716            ..Default::default()
717        })
718        .await;
719        let mut client = Client::connect(&host).await.unwrap();
720        assert!(matches!(
721            client
722                .search_stream(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
723                .await
724                .unwrap_err(),
725            Error::Rpc(_)
726        ));
727
728        let host = start_mock_server(MockBridgeService {
729            search_stream_error: Some(Status::deadline_exceeded("slow")),
730            ..Default::default()
731        })
732        .await;
733        let mut client = Client::connect(&host).await.unwrap();
734        assert!(matches!(
735            client
736                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
737                .await
738                .unwrap_err(),
739            Error::Rpc(_)
740        ));
741
742        let host = start_mock_server(MockBridgeService {
743            search_events: vec![ProtoSearchEvent::default()],
744            ..Default::default()
745        })
746        .await;
747        let mut client = Client::connect(&host).await.unwrap();
748        assert!(matches!(
749            client
750                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
751                .await
752                .unwrap_err(),
753            Error::Protocol(_)
754        ));
755
756        assert!(matches!(
757            feature_error("test", Status::unimplemented("old")),
758            Error::IncompatibleGateway { .. }
759        ));
760    }
761
762    fn index_status(state: ProtoSearchIndexState) -> SearchIndexStatus {
763        SearchIndexStatus {
764            server: "S".into(),
765            state: state as i32,
766            configured: true,
767            active_generation: 4,
768            entry_count: 100,
769            unique_item_count: 99,
770            started_at: Some("start".into()),
771            completed_at: Some("complete".into()),
772            last_error: None,
773            database_bytes: 1024,
774            organization: ProtoOrganization::Hierarchical as i32,
775            source: ProtoBrowseSource::Da3 as i32,
776            progress: Some(IndexedSearchProgress {
777                branches_visited: 2,
778                entries_seen: 3,
779                unique_items: 3,
780                active_time_ms: 4,
781                paused_time_ms: 5,
782                items_per_second: 6.0,
783                estimated_remaining_ms: Some(7),
784            }),
785            ..Default::default()
786        }
787    }
788
789    #[tokio::test]
790    async fn indexed_search_methods_map_requests_responses_and_errors() {
791        let service = MockBridgeService {
792            search_index_status_response: index_status(ProtoSearchIndexState::Ready),
793            refresh_search_index_response: index_status(ProtoSearchIndexState::Refreshing),
794            control_search_index_response: index_status(ProtoSearchIndexState::Partial),
795            search_index_response: SearchIndexResponse {
796                matches: vec![IndexedSearchMatch {
797                    item_id: "Exact.ItemID".into(),
798                    display_name: "PV".into(),
799                    kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
800                    breadcrumbs: vec!["Area".into()],
801                }],
802                has_more: true,
803                status: Some(index_status(ProtoSearchIndexState::Stale)),
804            },
805            ..Default::default()
806        };
807        let status_requests = Arc::clone(&service.search_index_status_requests);
808        let refresh_requests = Arc::clone(&service.refresh_search_index_requests);
809        let control_requests = Arc::clone(&service.control_search_index_requests);
810        let search_requests = Arc::clone(&service.search_index_requests);
811        let host = start_mock_server(service).await;
812        let mut client = Client::connect(&host).await.unwrap();
813
814        assert_eq!(
815            client.search_index_status("S").await.unwrap().state,
816            SearchIndexState::Ready
817        );
818        assert_eq!(
819            client.refresh_search_index("S", true).await.unwrap().state,
820            SearchIndexState::Refreshing
821        );
822        assert_eq!(
823            client
824                .control_search_index("S", SearchIndexControlAction::Pause)
825                .await
826                .unwrap()
827                .state,
828            SearchIndexState::Partial
829        );
830        let mut request = SearchIndexRequest::new("S", "PV", SearchMatchMode::Contains);
831        request.max_results = 25;
832        let response = client.search_index(request).await.unwrap();
833        assert_eq!(response.matches[0].item_id, "Exact.ItemID");
834        assert!(response.has_more);
835
836        assert_eq!(status_requests.lock().unwrap()[0].server, "S");
837        assert!(refresh_requests.lock().unwrap()[0].force);
838        assert_eq!(
839            control_requests.lock().unwrap()[0].action,
840            opcda_bridge_proto::bridge::SearchIndexControlAction::Pause as i32
841        );
842        assert_eq!(search_requests.lock().unwrap()[0].max_results, 25);
843
844        for (field, operation) in [
845            ("search_index_status_error", "status"),
846            ("refresh_search_index_error", "refresh"),
847            ("control_search_index_error", "control"),
848            ("search_index_error", "search"),
849        ] {
850            let mut service = MockBridgeService::default();
851            let error = Some(Status::unimplemented("old gateway"));
852            match field {
853                "search_index_status_error" => service.search_index_status_error = error,
854                "refresh_search_index_error" => service.refresh_search_index_error = error,
855                "control_search_index_error" => service.control_search_index_error = error,
856                "search_index_error" => service.search_index_error = error,
857                _ => unreachable!(),
858            }
859            let host = start_mock_server(service).await;
860            let mut client = Client::connect(&host).await.unwrap();
861            let error = match operation {
862                "status" => client.search_index_status("S").await.unwrap_err(),
863                "refresh" => client.refresh_search_index("S", false).await.unwrap_err(),
864                "control" => client
865                    .control_search_index("S", SearchIndexControlAction::Cancel)
866                    .await
867                    .unwrap_err(),
868                "search" => client
869                    .search_index(SearchIndexRequest::new(
870                        "S",
871                        "PV",
872                        SearchMatchMode::Contains,
873                    ))
874                    .await
875                    .unwrap_err(),
876                _ => unreachable!(),
877            };
878            assert!(matches!(error, Error::IncompatibleGateway { .. }));
879        }
880    }
881
882    #[tokio::test]
883    async fn read_maps_data_and_errors() {
884        let host = start_mock_server(MockBridgeService {
885            read_response: ReadResponse {
886                values: ["AUT", "", "A\"B", "\"AUT\""]
887                    .into_iter()
888                    .enumerate()
889                    .map(|(index, value)| ProtoTagValue {
890                        tag_id: format!("t{index}"),
891                        value: value.into(),
892                        quality: "Good".into(),
893                        timestamp: "now".into(),
894                    })
895                    .collect(),
896            },
897            ..Default::default()
898        })
899        .await;
900        let mut client = Client::connect(&host).await.unwrap();
901        let values = client.read("S".into(), vec![]).await.unwrap();
902        assert_eq!(
903            values
904                .iter()
905                .map(|value| value.value.as_str())
906                .collect::<Vec<_>>(),
907            vec!["AUT", "", "A\"B", "\"AUT\""]
908        );
909
910        let host = start_mock_server(MockBridgeService {
911            read_error: Some(Status::internal("boom")),
912            ..Default::default()
913        })
914        .await;
915        let mut client = Client::connect(&host).await.unwrap();
916        assert!(matches!(
917            client.read("S".into(), vec![]).await.unwrap_err(),
918            Error::Rpc(_)
919        ));
920    }
921
922    #[tokio::test]
923    async fn write_maps_every_value_and_result_or_error() {
924        for value in [
925            Value::Bool(true),
926            Value::Int(42),
927            Value::Float(3.5),
928            Value::String("text".into()),
929        ] {
930            let host = start_mock_server(MockBridgeService {
931                write_response: WriteResponse {
932                    tag_id: "t".into(),
933                    success: true,
934                    error: None,
935                },
936                ..Default::default()
937            })
938            .await;
939            let mut client = Client::connect(&host).await.unwrap();
940            assert!(
941                client
942                    .write("S".into(), "t".into(), value)
943                    .await
944                    .unwrap()
945                    .success
946            );
947        }
948
949        let host = start_mock_server(MockBridgeService {
950            write_error: Some(Status::internal("boom")),
951            ..Default::default()
952        })
953        .await;
954        let mut client = Client::connect(&host).await.unwrap();
955        assert!(matches!(
956            client
957                .write("S".into(), "t".into(), Value::Int(1))
958                .await
959                .unwrap_err(),
960            Error::Rpc(_)
961        ));
962    }
963}