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 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 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 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 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 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}