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        let server = server.into();
212        self.inner
213            .refresh_search_index(RefreshSearchIndexRequest {
214                server: server.clone(),
215                force,
216            })
217            .await
218            .map_err(|status| index_operation_error("indexed-search refresh", server, status))?
219            .into_inner()
220            .try_into()
221    }
222
223    /// Pause, resume, or cancel a persistent namespace-index build.
224    pub async fn control_search_index(
225        &mut self,
226        server: impl Into<String>,
227        action: SearchIndexControlAction,
228    ) -> Result<SearchIndexStatus> {
229        let server = server.into();
230        self.inner
231            .control_search_index(ControlSearchIndexRequest {
232                server: server.clone(),
233                action: opcda_bridge_proto::bridge::SearchIndexControlAction::from(action) as i32,
234            })
235            .await
236            .map_err(|status| index_operation_error("indexed-search control", server, status))?
237            .into_inner()
238            .try_into()
239    }
240
241    /// Enable or disable completion-time-based scheduled refreshes for an
242    /// enrolled namespace index. Disabling preserves searchable data.
243    pub async fn set_search_index_auto_refresh(
244        &mut self,
245        server: impl Into<String>,
246        enabled: bool,
247    ) -> Result<SearchIndexStatus> {
248        self.control_search_index(
249            server,
250            if enabled {
251                SearchIndexControlAction::EnableAutoRefresh
252            } else {
253                SearchIndexControlAction::DisableAutoRefresh
254            },
255        )
256        .await
257    }
258
259    /// Cancel any active build and permanently remove an enrolled namespace
260    /// index, including its generations and scheduler metadata.
261    pub async fn delete_search_index(
262        &mut self,
263        server: impl Into<String>,
264    ) -> Result<SearchIndexStatus> {
265        self.control_search_index(server, SearchIndexControlAction::Delete)
266            .await
267    }
268
269    /// Search the gateway-owned persistent namespace index.
270    ///
271    /// This never falls back to live namespace traversal.
272    pub async fn search_index(
273        &mut self,
274        request: SearchIndexRequest,
275    ) -> Result<SearchIndexResponse> {
276        self.inner
277            .search_index(opcda_bridge_proto::bridge::SearchIndexRequest::from(
278                request,
279            ))
280            .await
281            .map_err(|status| feature_error("indexed search", status))?
282            .into_inner()
283            .try_into()
284    }
285
286    /// Read one or more exact OPC DA ItemIDs from `server`.
287    pub async fn read(&mut self, server: String, tags: Vec<String>) -> Result<Vec<TagValue>> {
288        let response = self
289            .inner
290            .read(ReadRequest {
291                server,
292                tag_ids: tags,
293            })
294            .await?;
295        Ok(response
296            .into_inner()
297            .values
298            .into_iter()
299            .map(|v| TagValue {
300                tag_id: v.tag_id,
301                value: v.value,
302                quality: v.quality,
303                timestamp: v.timestamp,
304            })
305            .collect())
306    }
307
308    /// Write `value` to one exact OPC DA ItemID on `server`.
309    pub async fn write(
310        &mut self,
311        server: String,
312        tag: String,
313        value: Value,
314    ) -> Result<WriteResult> {
315        let typed_value = match value {
316            Value::String(s) => TypedValue::StringValue(s),
317            Value::Int(i) => TypedValue::IntValue(i),
318            Value::Float(f) => TypedValue::FloatValue(f),
319            Value::Bool(b) => TypedValue::BoolValue(b),
320        };
321        let response = self
322            .inner
323            .write(WriteRequest {
324                server,
325                tag_id: tag,
326                typed_value: Some(typed_value),
327            })
328            .await?;
329        let result = response.into_inner();
330        Ok(WriteResult {
331            tag_id: result.tag_id,
332            success: result.success,
333            error: result.error,
334        })
335    }
336}
337
338fn feature_error(operation: &'static str, status: tonic::Status) -> Error {
339    if status.code() == Code::Unimplemented {
340        Error::IncompatibleGateway { operation }
341    } else {
342        Error::Rpc(status)
343    }
344}
345
346fn index_operation_error(operation: &'static str, server: String, status: tonic::Status) -> Error {
347    match status.code() {
348        Code::Unimplemented => Error::IncompatibleGateway { operation },
349        Code::InvalidArgument => Error::UnknownIndexServer { server },
350        Code::NotFound => Error::IndexNotEnrolled { server },
351        Code::FailedPrecondition
352            if status.message()
353                == format!("namespace index for OPC DA server {server:?} is being deleted") =>
354        {
355            Error::IndexDeleting { server }
356        }
357        _ => Error::Rpc(status),
358    }
359}
360
361#[cfg(test)]
362mod tests {
363    use super::*;
364    use crate::error::Error;
365    use crate::test_support::{MockBridgeService, start_mock_server};
366    use crate::{
367        BrowseNodeKind, BrowseSource, NamespaceOrganization, SearchIndexControlAction,
368        SearchIndexRequest, SearchIndexState, SearchMatchMode,
369    };
370    use opcda_bridge_proto::bridge::search_event;
371    use opcda_bridge_proto::bridge::{
372        BrowseNode as ProtoBrowseNode, BrowsePage as ProtoBrowsePage,
373        BrowseSource as ProtoBrowseSource, GetCapabilitiesResponse, GetGatewayInfoResponse,
374        IndexedSearchMatch, IndexedSearchProgress, ListServersResponse,
375        NamespaceOrganization as ProtoOrganization, ProtocolFeature, ProtocolFeatureKind,
376        ReadResponse, SearchCompleted, SearchEvent as ProtoSearchEvent, SearchIndexResponse,
377        SearchIndexState as ProtoSearchIndexState, SearchIndexStatus, SearchProgress,
378        TagValue as ProtoTagValue, WriteResponse,
379    };
380    use std::sync::Arc;
381    use std::time::Duration;
382    use tonic::Status;
383
384    fn item_node() -> ProtoBrowseNode {
385        ProtoBrowseNode {
386            node_key: "node".into(),
387            display_name: "PV".into(),
388            kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
389            item_id: Some("FCS!TAG.PV".into()),
390        }
391    }
392
393    #[test]
394    fn index_operation_error_maps_unclassified_status_to_rpc() {
395        let error = index_operation_error(
396            "indexed-search refresh",
397            "S".into(),
398            tonic::Status::permission_denied("denied"),
399        );
400        assert!(matches!(
401            error,
402            Error::Rpc(status) if status.code() == tonic::Code::PermissionDenied
403        ));
404    }
405
406    #[test]
407    fn index_operation_error_only_classifies_the_gateway_deletion_message() {
408        let error = index_operation_error(
409            "indexed-search refresh",
410            "S".into(),
411            tonic::Status::failed_precondition(
412                r#"namespace index for OPC DA server "S" is being deleted"#,
413            ),
414        );
415        assert!(matches!(error, Error::IndexDeleting { server } if server == "S"));
416
417        let error = index_operation_error(
418            "indexed-search refresh",
419            "S".into(),
420            tonic::Status::failed_precondition("some other precondition failed"),
421        );
422        assert!(matches!(
423            error,
424            Error::Rpc(status) if status.code() == tonic::Code::FailedPrecondition
425        ));
426    }
427
428    #[tokio::test]
429    async fn connect_success_and_failure_are_typed() {
430        let host = start_mock_server(MockBridgeService::default()).await;
431        Client::connect(&host).await.unwrap();
432        assert!(matches!(
433            Client::connect("127.0.0.1:1").await.unwrap_err(),
434            Error::Connect(_)
435        ));
436    }
437
438    #[tokio::test]
439    async fn mock_server_shutdown_completes() {
440        let service = MockBridgeService::default();
441        let shutdown = Arc::clone(&service.server_shutdown);
442        let stopped = Arc::clone(&service.server_stopped);
443        let _host = start_mock_server(service).await;
444        shutdown.notify_one();
445        tokio::time::timeout(Duration::from_secs(1), stopped.notified())
446            .await
447            .unwrap();
448    }
449
450    #[tokio::test]
451    async fn capabilities_maps_fields_and_request() {
452        let service = MockBridgeService {
453            capabilities_response: GetCapabilitiesResponse {
454                application_version: "0.3.0".into(),
455                protocol_version: "0.3".into(),
456                max_page_size: 1000,
457                supports_browse_sessions: true,
458                supports_search: true,
459                organization: ProtoOrganization::Hierarchical as i32,
460                source: ProtoBrowseSource::Da2 as i32,
461                supports_indexed_search: true,
462                indexed_search_protocol_version: "1".into(),
463                max_indexed_search_results: 50,
464                search_index_state: ProtoSearchIndexState::Ready as i32,
465                search_index_promoting: true,
466            },
467            ..Default::default()
468        };
469        let requests = Arc::clone(&service.capabilities_requests);
470        let host = start_mock_server(service).await;
471        let mut client = Client::connect(&host).await.unwrap();
472        let capabilities = client.capabilities("S").await.unwrap();
473        assert_eq!(
474            capabilities.organization,
475            NamespaceOrganization::Hierarchical
476        );
477        assert_eq!(capabilities.source, BrowseSource::Da2);
478        assert!(capabilities.supports_indexed_search);
479        assert_eq!(capabilities.indexed_search_protocol_version, "1");
480        assert_eq!(capabilities.max_indexed_search_results, 50);
481        assert_eq!(capabilities.search_index_state, SearchIndexState::Ready);
482        assert!(capabilities.search_index_promoting);
483        assert_eq!(requests.lock().unwrap()[0].server, "S");
484    }
485
486    #[tokio::test]
487    async fn gateway_info_and_compatibility_report_are_typed() {
488        let service = MockBridgeService {
489            gateway_info_response: GetGatewayInfoResponse {
490                application_version: "0.5.0".into(),
491                compatibility_schema_version: 1,
492                features: vec![
493                    ProtocolFeature {
494                        kind: ProtocolFeatureKind::Core as i32,
495                        min_version: 1,
496                        max_version: 1,
497                    },
498                    ProtocolFeature {
499                        kind: ProtocolFeatureKind::Namespace as i32,
500                        min_version: 2,
501                        max_version: 2,
502                    },
503                    ProtocolFeature {
504                        kind: ProtocolFeatureKind::IndexedSearch as i32,
505                        min_version: 2,
506                        max_version: 2,
507                    },
508                ],
509            },
510            ..Default::default()
511        };
512        let requests = Arc::clone(&service.gateway_info_requests);
513        let host = start_mock_server(service).await;
514        let mut client = Client::connect(&host).await.unwrap();
515        let info = client.gateway_info().await.unwrap();
516        assert_eq!(info.application_version, "0.5.0");
517        let report = client
518            .compatibility_with_client_version(None, "0.5.0")
519            .await
520            .unwrap();
521        assert_eq!(report.status, crate::CompatibilityStatus::Full);
522        assert_eq!(report.library_version, env!("CARGO_PKG_VERSION"));
523        assert_eq!(requests.lock().unwrap().len(), 2);
524    }
525
526    #[tokio::test]
527    async fn compatibility_wrapper_and_gateway_info_errors_are_typed() {
528        let host = start_mock_server(MockBridgeService {
529            gateway_info_response: GetGatewayInfoResponse {
530                application_version: "0.5.0".into(),
531                compatibility_schema_version: 1,
532                features: vec![
533                    ProtocolFeature {
534                        kind: ProtocolFeatureKind::Core as i32,
535                        min_version: 1,
536                        max_version: 1,
537                    },
538                    ProtocolFeature {
539                        kind: ProtocolFeatureKind::Namespace as i32,
540                        min_version: 2,
541                        max_version: 2,
542                    },
543                    ProtocolFeature {
544                        kind: ProtocolFeatureKind::IndexedSearch as i32,
545                        min_version: 2,
546                        max_version: 2,
547                    },
548                ],
549            },
550            ..Default::default()
551        })
552        .await;
553        let mut client = Client::connect(&host).await.unwrap();
554        assert_eq!(
555            client.compatibility(None).await.unwrap().status,
556            crate::CompatibilityStatus::Full
557        );
558
559        let host = start_mock_server(MockBridgeService {
560            gateway_info_error: Some(Status::internal("gateway unavailable")),
561            ..Default::default()
562        })
563        .await;
564        let mut client = Client::connect(&host).await.unwrap();
565        assert!(matches!(
566            client.compatibility(None).await.unwrap_err(),
567            Error::Rpc(_)
568        ));
569    }
570
571    #[tokio::test]
572    async fn compatibility_falls_back_to_legacy_or_reports_unknown() {
573        let service = MockBridgeService {
574            gateway_info_error: Some(Status::unimplemented("old gateway")),
575            capabilities_response: GetCapabilitiesResponse {
576                application_version: "0.3.2".into(),
577                protocol_version: "2".into(),
578                supports_indexed_search: false,
579                ..Default::default()
580            },
581            ..Default::default()
582        };
583        let host = start_mock_server(service).await;
584        let mut client = Client::connect(&host).await.unwrap();
585        let report = client
586            .compatibility_with_client_version(Some("S"), "0.4.3")
587            .await
588            .unwrap();
589        assert_eq!(
590            report.source,
591            crate::CompatibilitySource::LegacyCapabilities
592        );
593        assert_eq!(report.status, crate::CompatibilityStatus::Partial);
594
595        let host = start_mock_server(MockBridgeService {
596            gateway_info_error: Some(Status::unimplemented("old gateway")),
597            ..Default::default()
598        })
599        .await;
600        let mut client = Client::connect(&host).await.unwrap();
601        let report = client
602            .compatibility_with_client_version(None, "0.4.3")
603            .await
604            .unwrap();
605        assert_eq!(report.status, crate::CompatibilityStatus::Unknown);
606    }
607
608    #[tokio::test]
609    async fn capabilities_rpc_error_is_typed() {
610        let host = start_mock_server(MockBridgeService {
611            capabilities_error: Some(Status::unimplemented("old gateway")),
612            ..Default::default()
613        })
614        .await;
615        let mut client = Client::connect(&host).await.unwrap();
616        assert!(matches!(
617            client.capabilities("S").await.unwrap_err(),
618            Error::IncompatibleGateway { .. }
619        ));
620    }
621
622    #[tokio::test]
623    async fn list_servers_maps_data_and_errors() {
624        let host = start_mock_server(MockBridgeService {
625            list_servers_response: ListServersResponse {
626                servers: vec!["S1".into(), "S2".into()],
627            },
628            ..Default::default()
629        })
630        .await;
631        let mut client = Client::connect(&host).await.unwrap();
632        assert_eq!(client.list_servers().await.unwrap(), ["S1", "S2"]);
633
634        let host = start_mock_server(MockBridgeService {
635            list_servers_error: Some(Status::internal("boom")),
636            ..Default::default()
637        })
638        .await;
639        let mut client = Client::connect(&host).await.unwrap();
640        assert!(matches!(
641            client.list_servers().await.unwrap_err(),
642            Error::Rpc(_)
643        ));
644    }
645
646    #[tokio::test]
647    async fn browse_returns_one_typed_page_without_draining() {
648        let service = MockBridgeService {
649            browse_response: ProtoBrowsePage {
650                session_id: "session".into(),
651                nodes: vec![item_node()],
652                next_page_token: Some("next".into()),
653                complete: false,
654                organization: ProtoOrganization::Hierarchical as i32,
655                source: ProtoBrowseSource::Da3 as i32,
656                warning: None,
657            },
658            ..Default::default()
659        };
660        let requests = Arc::clone(&service.browse_requests);
661        let host = start_mock_server(service).await;
662        let mut client = Client::connect(&host).await.unwrap();
663        let page = client.browse("S", 25).await.unwrap();
664        assert_eq!(page.nodes[0].kind, BrowseNodeKind::Item);
665        assert_eq!(page.next_page_token.as_deref(), Some("next"));
666        let request = &requests.lock().unwrap()[0];
667        assert_eq!(request.page_size, 25);
668        assert!(request.session_id.is_none());
669    }
670
671    #[tokio::test]
672    async fn browse_page_forwards_session_parent_token_and_refresh() {
673        let service = MockBridgeService::default();
674        let requests = Arc::clone(&service.browse_requests);
675        let host = start_mock_server(service).await;
676        let mut client = Client::connect(&host).await.unwrap();
677        client
678            .browse_page(
679                BrowsePageRequest::next("S", "session", Some("parent".into()), "token", 10)
680                    .with_refresh(true),
681            )
682            .await
683            .unwrap();
684        let request = &requests.lock().unwrap()[0];
685        assert_eq!(request.session_id.as_deref(), Some("session"));
686        assert_eq!(request.parent_node_key.as_deref(), Some("parent"));
687        assert_eq!(request.page_token.as_deref(), Some("token"));
688        assert!(request.refresh);
689    }
690
691    #[tokio::test]
692    async fn browse_rpc_and_protocol_errors_are_typed() {
693        let host = start_mock_server(MockBridgeService {
694            browse_error: Some(Status::failed_precondition("expired")),
695            ..Default::default()
696        })
697        .await;
698        let mut client = Client::connect(&host).await.unwrap();
699        assert!(matches!(
700            client.open_browse("S", 20).await.unwrap_err(),
701            Error::Rpc(_)
702        ));
703
704        let host = start_mock_server(MockBridgeService {
705            browse_response: ProtoBrowsePage::default(),
706            ..Default::default()
707        })
708        .await;
709        let mut client = Client::connect(&host).await.unwrap();
710        assert!(matches!(
711            client.open_browse("S", 20).await.unwrap_err(),
712            Error::Protocol(_)
713        ));
714
715        let host = start_mock_server(MockBridgeService {
716            browse_error: Some(Status::unimplemented("old gateway")),
717            ..Default::default()
718        })
719        .await;
720        let mut client = Client::connect(&host).await.unwrap();
721        assert!(matches!(
722            client.open_browse("S", 20).await.unwrap_err(),
723            Error::IncompatibleGateway { .. }
724        ));
725    }
726
727    #[tokio::test]
728    async fn close_browse_session_forwards_id_and_error() {
729        let service = MockBridgeService::default();
730        let requests = Arc::clone(&service.close_requests);
731        let host = start_mock_server(service).await;
732        let mut client = Client::connect(&host).await.unwrap();
733        client.close_browse_session("session").await.unwrap();
734        assert_eq!(requests.lock().unwrap()[0].session_id, "session");
735
736        let host = start_mock_server(MockBridgeService {
737            close_error: Some(Status::not_found("missing")),
738            ..Default::default()
739        })
740        .await;
741        let mut client = Client::connect(&host).await.unwrap();
742        assert!(matches!(
743            client.close_browse_session("missing").await.unwrap_err(),
744            Error::Rpc(_)
745        ));
746    }
747
748    #[tokio::test]
749    async fn search_stream_and_collect_map_events_and_request() {
750        let events = vec![
751            ProtoSearchEvent {
752                event: Some(search_event::Event::Progress(SearchProgress {
753                    visited_nodes: 5,
754                    matches: 0,
755                    partial: true,
756                })),
757            },
758            ProtoSearchEvent {
759                event: Some(search_event::Event::Completed(SearchCompleted {
760                    complete: true,
761                    cancelled: false,
762                    truncated: false,
763                    warning: None,
764                })),
765            },
766        ];
767        let service = MockBridgeService {
768            search_events: events,
769            ..Default::default()
770        };
771        let requests = Arc::clone(&service.search_requests);
772        let host = start_mock_server(service).await;
773        let mut client = Client::connect(&host).await.unwrap();
774        let mut request = SearchRequest::new("S", "PV", SearchMatchMode::Prefix);
775        request.session_id = Some("session".into());
776        request.scope_node_key = Some("scope".into());
777        request.max_results = 50;
778        request.include_branches = true;
779        request.refresh = true;
780        let found = client.search(request).await.unwrap();
781        assert_eq!(found.len(), 2);
782        let request = &requests.lock().unwrap()[0];
783        assert_eq!(request.query, "PV");
784        assert_eq!(
785            request.match_mode,
786            opcda_bridge_proto::bridge::SearchMatchMode::Prefix as i32
787        );
788        assert!(request.include_branches);
789        assert!(request.refresh);
790    }
791
792    #[tokio::test]
793    async fn search_initial_stream_and_protocol_errors_are_typed() {
794        let host = start_mock_server(MockBridgeService {
795            search_initial_error: Some(Status::unavailable("down")),
796            ..Default::default()
797        })
798        .await;
799        let mut client = Client::connect(&host).await.unwrap();
800        assert!(matches!(
801            client
802                .search_stream(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
803                .await
804                .unwrap_err(),
805            Error::Rpc(_)
806        ));
807
808        let host = start_mock_server(MockBridgeService {
809            search_stream_error: Some(Status::deadline_exceeded("slow")),
810            ..Default::default()
811        })
812        .await;
813        let mut client = Client::connect(&host).await.unwrap();
814        assert!(matches!(
815            client
816                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
817                .await
818                .unwrap_err(),
819            Error::Rpc(_)
820        ));
821
822        let host = start_mock_server(MockBridgeService {
823            search_events: vec![ProtoSearchEvent::default()],
824            ..Default::default()
825        })
826        .await;
827        let mut client = Client::connect(&host).await.unwrap();
828        assert!(matches!(
829            client
830                .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
831                .await
832                .unwrap_err(),
833            Error::Protocol(_)
834        ));
835
836        assert!(matches!(
837            feature_error("test", Status::unimplemented("old")),
838            Error::IncompatibleGateway { .. }
839        ));
840    }
841
842    fn index_status(state: ProtoSearchIndexState) -> SearchIndexStatus {
843        SearchIndexStatus {
844            server: "S".into(),
845            state: state as i32,
846            configured: true,
847            active_generation: 4,
848            entry_count: 100,
849            unique_item_count: 99,
850            started_at: Some("start".into()),
851            completed_at: Some("complete".into()),
852            last_error: None,
853            database_bytes: 1024,
854            organization: ProtoOrganization::Hierarchical as i32,
855            source: ProtoBrowseSource::Da3 as i32,
856            progress: Some(IndexedSearchProgress {
857                branches_visited: 2,
858                entries_seen: 3,
859                unique_items: 3,
860                active_time_ms: 4,
861                paused_time_ms: 5,
862                items_per_second: 6.0,
863                estimated_remaining_ms: Some(7),
864            }),
865            ..Default::default()
866        }
867    }
868
869    #[tokio::test]
870    async fn indexed_search_methods_map_requests_responses_and_errors() {
871        let service = MockBridgeService {
872            search_index_status_response: index_status(ProtoSearchIndexState::Ready),
873            refresh_search_index_response: index_status(ProtoSearchIndexState::Refreshing),
874            control_search_index_response: index_status(ProtoSearchIndexState::Partial),
875            search_index_response: SearchIndexResponse {
876                matches: vec![IndexedSearchMatch {
877                    item_id: "Exact.ItemID".into(),
878                    display_name: "PV".into(),
879                    kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
880                    breadcrumbs: vec!["Area".into()],
881                }],
882                has_more: true,
883                status: Some(index_status(ProtoSearchIndexState::Stale)),
884            },
885            ..Default::default()
886        };
887        let status_requests = Arc::clone(&service.search_index_status_requests);
888        let refresh_requests = Arc::clone(&service.refresh_search_index_requests);
889        let control_requests = Arc::clone(&service.control_search_index_requests);
890        let search_requests = Arc::clone(&service.search_index_requests);
891        let host = start_mock_server(service).await;
892        let mut client = Client::connect(&host).await.unwrap();
893
894        assert_eq!(
895            client.search_index_status("S").await.unwrap().state,
896            SearchIndexState::Ready
897        );
898        assert_eq!(
899            client.refresh_search_index("S", true).await.unwrap().state,
900            SearchIndexState::Refreshing
901        );
902        assert_eq!(
903            client
904                .control_search_index("S", SearchIndexControlAction::Pause)
905                .await
906                .unwrap()
907                .state,
908            SearchIndexState::Partial
909        );
910        client
911            .set_search_index_auto_refresh("S", false)
912            .await
913            .unwrap();
914        client
915            .set_search_index_auto_refresh("S", true)
916            .await
917            .unwrap();
918        client.delete_search_index("S").await.unwrap();
919        let mut request = SearchIndexRequest::new("S", "PV", SearchMatchMode::Contains);
920        request.max_results = 25;
921        let response = client.search_index(request).await.unwrap();
922        assert_eq!(response.matches[0].item_id, "Exact.ItemID");
923        assert!(response.has_more);
924
925        assert_eq!(status_requests.lock().unwrap()[0].server, "S");
926        assert!(refresh_requests.lock().unwrap()[0].force);
927        assert_eq!(
928            control_requests.lock().unwrap()[0].action,
929            opcda_bridge_proto::bridge::SearchIndexControlAction::Pause as i32
930        );
931        assert_eq!(
932            control_requests.lock().unwrap()[1].action,
933            opcda_bridge_proto::bridge::SearchIndexControlAction::DisableAutoRefresh as i32
934        );
935        assert_eq!(
936            control_requests.lock().unwrap()[2].action,
937            opcda_bridge_proto::bridge::SearchIndexControlAction::EnableAutoRefresh as i32
938        );
939        assert_eq!(
940            control_requests.lock().unwrap()[3].action,
941            opcda_bridge_proto::bridge::SearchIndexControlAction::Delete as i32
942        );
943        assert_eq!(search_requests.lock().unwrap()[0].max_results, 25);
944
945        for (field, operation) in [
946            ("search_index_status_error", "status"),
947            ("refresh_search_index_error", "refresh"),
948            ("control_search_index_error", "control"),
949            ("search_index_error", "search"),
950        ] {
951            let mut service = MockBridgeService::default();
952            let error = Some(Status::unimplemented("old gateway"));
953            match field {
954                "search_index_status_error" => service.search_index_status_error = error,
955                "refresh_search_index_error" => service.refresh_search_index_error = error,
956                "control_search_index_error" => service.control_search_index_error = error,
957                "search_index_error" => service.search_index_error = error,
958                _ => unreachable!(),
959            }
960            let host = start_mock_server(service).await;
961            let mut client = Client::connect(&host).await.unwrap();
962            let error = match operation {
963                "status" => client.search_index_status("S").await.unwrap_err(),
964                "refresh" => client.refresh_search_index("S", false).await.unwrap_err(),
965                "control" => client
966                    .control_search_index("S", SearchIndexControlAction::Cancel)
967                    .await
968                    .unwrap_err(),
969                "search" => client
970                    .search_index(SearchIndexRequest::new(
971                        "S",
972                        "PV",
973                        SearchMatchMode::Contains,
974                    ))
975                    .await
976                    .unwrap_err(),
977                _ => unreachable!(),
978            };
979            assert!(matches!(error, Error::IncompatibleGateway { .. }));
980        }
981    }
982
983    #[tokio::test]
984    async fn index_enrollment_errors_are_typed() {
985        let host = start_mock_server(MockBridgeService {
986            refresh_search_index_error: Some(Status::invalid_argument("unknown server")),
987            ..Default::default()
988        })
989        .await;
990        let mut client = Client::connect(&host).await.unwrap();
991        assert!(matches!(
992            client.refresh_search_index("Typo.Server", false).await,
993            Err(Error::UnknownIndexServer { server }) if server == "Typo.Server"
994        ));
995
996        let host = start_mock_server(MockBridgeService {
997            control_search_index_error: Some(Status::not_found("not enrolled")),
998            ..Default::default()
999        })
1000        .await;
1001        let mut client = Client::connect(&host).await.unwrap();
1002        assert!(matches!(
1003            client
1004                .set_search_index_auto_refresh("Known.Server", false)
1005                .await,
1006            Err(Error::IndexNotEnrolled { server }) if server == "Known.Server"
1007        ));
1008
1009        let host = start_mock_server(MockBridgeService {
1010            search_index_error: Some(Status::permission_denied("denied")),
1011            ..Default::default()
1012        })
1013        .await;
1014        let mut client = Client::connect(&host).await.unwrap();
1015        assert!(matches!(
1016            client
1017                .search_index(SearchIndexRequest::new(
1018                    "Known.Server",
1019                    "PV",
1020                    SearchMatchMode::Contains,
1021                ))
1022                .await,
1023            Err(Error::Rpc(status)) if status.code() == Code::PermissionDenied
1024        ));
1025    }
1026
1027    #[tokio::test]
1028    async fn read_maps_data_and_errors() {
1029        let host = start_mock_server(MockBridgeService {
1030            read_response: ReadResponse {
1031                values: ["AUT", "", "A\"B", "\"AUT\""]
1032                    .into_iter()
1033                    .enumerate()
1034                    .map(|(index, value)| ProtoTagValue {
1035                        tag_id: format!("t{index}"),
1036                        value: value.into(),
1037                        quality: "Good".into(),
1038                        timestamp: "now".into(),
1039                    })
1040                    .collect(),
1041            },
1042            ..Default::default()
1043        })
1044        .await;
1045        let mut client = Client::connect(&host).await.unwrap();
1046        let values = client.read("S".into(), vec![]).await.unwrap();
1047        assert_eq!(
1048            values
1049                .iter()
1050                .map(|value| value.value.as_str())
1051                .collect::<Vec<_>>(),
1052            vec!["AUT", "", "A\"B", "\"AUT\""]
1053        );
1054
1055        let host = start_mock_server(MockBridgeService {
1056            read_error: Some(Status::internal("boom")),
1057            ..Default::default()
1058        })
1059        .await;
1060        let mut client = Client::connect(&host).await.unwrap();
1061        assert!(matches!(
1062            client.read("S".into(), vec![]).await.unwrap_err(),
1063            Error::Rpc(_)
1064        ));
1065    }
1066
1067    #[tokio::test]
1068    async fn write_maps_every_value_and_result_or_error() {
1069        for value in [
1070            Value::Bool(true),
1071            Value::Int(42),
1072            Value::Float(3.5),
1073            Value::String("text".into()),
1074        ] {
1075            let host = start_mock_server(MockBridgeService {
1076                write_response: WriteResponse {
1077                    tag_id: "t".into(),
1078                    success: true,
1079                    error: None,
1080                },
1081                ..Default::default()
1082            })
1083            .await;
1084            let mut client = Client::connect(&host).await.unwrap();
1085            assert!(
1086                client
1087                    .write("S".into(), "t".into(), value)
1088                    .await
1089                    .unwrap()
1090                    .success
1091            );
1092        }
1093
1094        let host = start_mock_server(MockBridgeService {
1095            write_error: Some(Status::internal("boom")),
1096            ..Default::default()
1097        })
1098        .await;
1099        let mut client = Client::connect(&host).await.unwrap();
1100        assert!(matches!(
1101            client
1102                .write("S".into(), "t".into(), Value::Int(1))
1103                .await
1104                .unwrap_err(),
1105            Error::Rpc(_)
1106        ));
1107    }
1108}