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