1use 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#[derive(Debug)]
26pub struct Client {
27 inner: BridgeClient<Channel>,
28}
29
30#[derive(Debug)]
35pub struct SearchStream {
36 inner: Streaming<opcda_bridge_proto::bridge::SearchEvent>,
37}
38
39impl SearchStream {
40 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 pub async fn connect(host: &str) -> Result<Self> {
53 let inner = BridgeClient::connect(format!("http://{host}")).await?;
54 Ok(Self { inner })
55 }
56
57 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}