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 opcda_bridge_proto::bridge::bridge_client::BridgeClient;
10use opcda_bridge_proto::bridge::write_request::TypedValue;
11use opcda_bridge_proto::bridge::{
12 CloseBrowseSessionRequest, ControlSearchIndexRequest, GetCapabilitiesRequest,
13 GetSearchIndexStatusRequest, ListServersRequest, ReadRequest, RefreshSearchIndexRequest,
14 WriteRequest,
15};
16use tonic::Code;
17use tonic::codec::Streaming;
18use tonic::transport::Channel;
19
20#[derive(Debug)]
22pub struct Client {
23 inner: BridgeClient<Channel>,
24}
25
26#[derive(Debug)]
31pub struct SearchStream {
32 inner: Streaming<opcda_bridge_proto::bridge::SearchEvent>,
33}
34
35impl SearchStream {
36 pub async fn message(&mut self) -> Result<Option<SearchEvent>> {
38 self.inner
39 .message()
40 .await?
41 .map(SearchEvent::try_from)
42 .transpose()
43 }
44}
45
46impl Client {
47 pub async fn connect(host: &str) -> Result<Self> {
49 let inner = BridgeClient::connect(format!("http://{host}")).await?;
50 Ok(Self { inner })
51 }
52
53 pub async fn capabilities(&mut self, server: impl Into<String>) -> Result<Capabilities> {
55 self.inner
56 .get_capabilities(GetCapabilitiesRequest {
57 server: server.into(),
58 })
59 .await
60 .map_err(|status| feature_error("capability discovery", status))?
61 .into_inner()
62 .try_into()
63 }
64
65 pub async fn list_servers(&mut self) -> Result<Vec<String>> {
67 let response = self
68 .inner
69 .list_servers(ListServersRequest {
70 host: "localhost".to_string(),
71 })
72 .await?;
73 Ok(response.into_inner().servers)
74 }
75
76 pub async fn browse(
78 &mut self,
79 server: impl Into<String>,
80 page_size: u32,
81 ) -> Result<BrowsePage> {
82 self.open_browse(server, page_size).await
83 }
84
85 pub async fn open_browse(
87 &mut self,
88 server: impl Into<String>,
89 page_size: u32,
90 ) -> Result<BrowsePage> {
91 self.browse_page(BrowsePageRequest::root(server, page_size))
92 .await
93 }
94
95 pub async fn browse_page(&mut self, request: BrowsePageRequest) -> Result<BrowsePage> {
99 self.inner
100 .browse(opcda_bridge_proto::bridge::BrowseRequest::from(request))
101 .await
102 .map_err(|status| feature_error("paged browse", status))?
103 .into_inner()
104 .try_into()
105 }
106
107 pub async fn close_browse_session(&mut self, session_id: impl Into<String>) -> Result<()> {
109 self.inner
110 .close_browse_session(CloseBrowseSessionRequest {
111 session_id: session_id.into(),
112 })
113 .await
114 .map_err(|status| feature_error("browse-session close", status))?;
115 Ok(())
116 }
117
118 pub async fn search_stream(&mut self, request: SearchRequest) -> Result<SearchStream> {
120 let inner = self
121 .inner
122 .search(opcda_bridge_proto::bridge::SearchRequest::from(request))
123 .await
124 .map_err(|status| feature_error("namespace search", status))?
125 .into_inner();
126 Ok(SearchStream { inner })
127 }
128
129 pub async fn search(&mut self, request: SearchRequest) -> Result<Vec<SearchEvent>> {
131 let mut stream = self.search_stream(request).await?;
132 let mut events = Vec::new();
133 while let Some(event) = stream.message().await? {
134 events.push(event);
135 }
136 Ok(events)
137 }
138
139 pub async fn search_index_status(
141 &mut self,
142 server: impl Into<String>,
143 ) -> Result<SearchIndexStatus> {
144 self.inner
145 .get_search_index_status(GetSearchIndexStatusRequest {
146 server: server.into(),
147 })
148 .await
149 .map_err(|status| feature_error("indexed-search status", status))?
150 .into_inner()
151 .try_into()
152 }
153
154 pub async fn refresh_search_index(
156 &mut self,
157 server: impl Into<String>,
158 force: bool,
159 ) -> Result<SearchIndexStatus> {
160 self.inner
161 .refresh_search_index(RefreshSearchIndexRequest {
162 server: server.into(),
163 force,
164 })
165 .await
166 .map_err(|status| feature_error("indexed-search refresh", status))?
167 .into_inner()
168 .try_into()
169 }
170
171 pub async fn control_search_index(
173 &mut self,
174 server: impl Into<String>,
175 action: SearchIndexControlAction,
176 ) -> Result<SearchIndexStatus> {
177 self.inner
178 .control_search_index(ControlSearchIndexRequest {
179 server: server.into(),
180 action: opcda_bridge_proto::bridge::SearchIndexControlAction::from(action) as i32,
181 })
182 .await
183 .map_err(|status| feature_error("indexed-search control", status))?
184 .into_inner()
185 .try_into()
186 }
187
188 pub async fn search_index(
192 &mut self,
193 request: SearchIndexRequest,
194 ) -> Result<SearchIndexResponse> {
195 self.inner
196 .search_index(opcda_bridge_proto::bridge::SearchIndexRequest::from(
197 request,
198 ))
199 .await
200 .map_err(|status| feature_error("indexed search", status))?
201 .into_inner()
202 .try_into()
203 }
204
205 pub async fn read(&mut self, server: String, tags: Vec<String>) -> Result<Vec<TagValue>> {
207 let response = self
208 .inner
209 .read(ReadRequest {
210 server,
211 tag_ids: tags,
212 })
213 .await?;
214 Ok(response
215 .into_inner()
216 .values
217 .into_iter()
218 .map(|v| TagValue {
219 tag_id: v.tag_id,
220 value: v.value,
221 quality: v.quality,
222 timestamp: v.timestamp,
223 })
224 .collect())
225 }
226
227 pub async fn write(
229 &mut self,
230 server: String,
231 tag: String,
232 value: Value,
233 ) -> Result<WriteResult> {
234 let typed_value = match value {
235 Value::String(s) => TypedValue::StringValue(s),
236 Value::Int(i) => TypedValue::IntValue(i),
237 Value::Float(f) => TypedValue::FloatValue(f),
238 Value::Bool(b) => TypedValue::BoolValue(b),
239 };
240 let response = self
241 .inner
242 .write(WriteRequest {
243 server,
244 tag_id: tag,
245 typed_value: Some(typed_value),
246 })
247 .await?;
248 let result = response.into_inner();
249 Ok(WriteResult {
250 tag_id: result.tag_id,
251 success: result.success,
252 error: result.error,
253 })
254 }
255}
256
257fn feature_error(operation: &'static str, status: tonic::Status) -> Error {
258 if status.code() == Code::Unimplemented {
259 Error::IncompatibleGateway { operation }
260 } else {
261 Error::Rpc(status)
262 }
263}
264
265#[cfg(test)]
266mod tests {
267 use super::*;
268 use crate::error::Error;
269 use crate::test_support::{MockBridgeService, start_mock_server};
270 use crate::{
271 BrowseNodeKind, BrowseSource, NamespaceOrganization, SearchIndexControlAction,
272 SearchIndexRequest, SearchIndexState, SearchMatchMode,
273 };
274 use opcda_bridge_proto::bridge::search_event;
275 use opcda_bridge_proto::bridge::{
276 BrowseNode as ProtoBrowseNode, BrowsePage as ProtoBrowsePage,
277 BrowseSource as ProtoBrowseSource, GetCapabilitiesResponse, IndexedSearchMatch,
278 IndexedSearchProgress, ListServersResponse, NamespaceOrganization as ProtoOrganization,
279 ReadResponse, SearchCompleted, SearchEvent as ProtoSearchEvent, SearchIndexResponse,
280 SearchIndexState as ProtoSearchIndexState, SearchIndexStatus, SearchProgress,
281 TagValue as ProtoTagValue, WriteResponse,
282 };
283 use std::sync::Arc;
284 use std::time::Duration;
285 use tonic::Status;
286
287 fn item_node() -> ProtoBrowseNode {
288 ProtoBrowseNode {
289 node_key: "node".into(),
290 display_name: "PV".into(),
291 kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
292 item_id: Some("FCS!TAG.PV".into()),
293 }
294 }
295
296 #[tokio::test]
297 async fn connect_success_and_failure_are_typed() {
298 let host = start_mock_server(MockBridgeService::default()).await;
299 Client::connect(&host).await.unwrap();
300 assert!(matches!(
301 Client::connect("127.0.0.1:1").await.unwrap_err(),
302 Error::Connect(_)
303 ));
304 }
305
306 #[tokio::test]
307 async fn mock_server_shutdown_completes() {
308 let service = MockBridgeService::default();
309 let shutdown = Arc::clone(&service.server_shutdown);
310 let stopped = Arc::clone(&service.server_stopped);
311 let _host = start_mock_server(service).await;
312 shutdown.notify_one();
313 tokio::time::timeout(Duration::from_secs(1), stopped.notified())
314 .await
315 .unwrap();
316 }
317
318 #[tokio::test]
319 async fn capabilities_maps_fields_and_request() {
320 let service = MockBridgeService {
321 capabilities_response: GetCapabilitiesResponse {
322 application_version: "0.3.0".into(),
323 protocol_version: "0.3".into(),
324 max_page_size: 1000,
325 supports_browse_sessions: true,
326 supports_search: true,
327 organization: ProtoOrganization::Hierarchical as i32,
328 source: ProtoBrowseSource::Da2 as i32,
329 supports_indexed_search: true,
330 indexed_search_protocol_version: "1".into(),
331 max_indexed_search_results: 50,
332 search_index_state: ProtoSearchIndexState::Ready as i32,
333 },
334 ..Default::default()
335 };
336 let requests = Arc::clone(&service.capabilities_requests);
337 let host = start_mock_server(service).await;
338 let mut client = Client::connect(&host).await.unwrap();
339 let capabilities = client.capabilities("S").await.unwrap();
340 assert_eq!(
341 capabilities.organization,
342 NamespaceOrganization::Hierarchical
343 );
344 assert_eq!(capabilities.source, BrowseSource::Da2);
345 assert!(capabilities.supports_indexed_search);
346 assert_eq!(capabilities.indexed_search_protocol_version, "1");
347 assert_eq!(capabilities.max_indexed_search_results, 50);
348 assert_eq!(capabilities.search_index_state, SearchIndexState::Ready);
349 assert_eq!(requests.lock().unwrap()[0].server, "S");
350 }
351
352 #[tokio::test]
353 async fn capabilities_rpc_error_is_typed() {
354 let host = start_mock_server(MockBridgeService {
355 capabilities_error: Some(Status::unimplemented("old gateway")),
356 ..Default::default()
357 })
358 .await;
359 let mut client = Client::connect(&host).await.unwrap();
360 assert!(matches!(
361 client.capabilities("S").await.unwrap_err(),
362 Error::IncompatibleGateway { .. }
363 ));
364 }
365
366 #[tokio::test]
367 async fn list_servers_maps_data_and_errors() {
368 let host = start_mock_server(MockBridgeService {
369 list_servers_response: ListServersResponse {
370 servers: vec!["S1".into(), "S2".into()],
371 },
372 ..Default::default()
373 })
374 .await;
375 let mut client = Client::connect(&host).await.unwrap();
376 assert_eq!(client.list_servers().await.unwrap(), ["S1", "S2"]);
377
378 let host = start_mock_server(MockBridgeService {
379 list_servers_error: Some(Status::internal("boom")),
380 ..Default::default()
381 })
382 .await;
383 let mut client = Client::connect(&host).await.unwrap();
384 assert!(matches!(
385 client.list_servers().await.unwrap_err(),
386 Error::Rpc(_)
387 ));
388 }
389
390 #[tokio::test]
391 async fn browse_returns_one_typed_page_without_draining() {
392 let service = MockBridgeService {
393 browse_response: ProtoBrowsePage {
394 session_id: "session".into(),
395 nodes: vec![item_node()],
396 next_page_token: Some("next".into()),
397 complete: false,
398 organization: ProtoOrganization::Hierarchical as i32,
399 source: ProtoBrowseSource::Da3 as i32,
400 warning: None,
401 },
402 ..Default::default()
403 };
404 let requests = Arc::clone(&service.browse_requests);
405 let host = start_mock_server(service).await;
406 let mut client = Client::connect(&host).await.unwrap();
407 let page = client.browse("S", 25).await.unwrap();
408 assert_eq!(page.nodes[0].kind, BrowseNodeKind::Item);
409 assert_eq!(page.next_page_token.as_deref(), Some("next"));
410 let request = &requests.lock().unwrap()[0];
411 assert_eq!(request.page_size, 25);
412 assert!(request.session_id.is_none());
413 }
414
415 #[tokio::test]
416 async fn browse_page_forwards_session_parent_token_and_refresh() {
417 let service = MockBridgeService::default();
418 let requests = Arc::clone(&service.browse_requests);
419 let host = start_mock_server(service).await;
420 let mut client = Client::connect(&host).await.unwrap();
421 client
422 .browse_page(
423 BrowsePageRequest::next("S", "session", Some("parent".into()), "token", 10)
424 .with_refresh(true),
425 )
426 .await
427 .unwrap();
428 let request = &requests.lock().unwrap()[0];
429 assert_eq!(request.session_id.as_deref(), Some("session"));
430 assert_eq!(request.parent_node_key.as_deref(), Some("parent"));
431 assert_eq!(request.page_token.as_deref(), Some("token"));
432 assert!(request.refresh);
433 }
434
435 #[tokio::test]
436 async fn browse_rpc_and_protocol_errors_are_typed() {
437 let host = start_mock_server(MockBridgeService {
438 browse_error: Some(Status::failed_precondition("expired")),
439 ..Default::default()
440 })
441 .await;
442 let mut client = Client::connect(&host).await.unwrap();
443 assert!(matches!(
444 client.open_browse("S", 20).await.unwrap_err(),
445 Error::Rpc(_)
446 ));
447
448 let host = start_mock_server(MockBridgeService {
449 browse_response: ProtoBrowsePage::default(),
450 ..Default::default()
451 })
452 .await;
453 let mut client = Client::connect(&host).await.unwrap();
454 assert!(matches!(
455 client.open_browse("S", 20).await.unwrap_err(),
456 Error::Protocol(_)
457 ));
458
459 let host = start_mock_server(MockBridgeService {
460 browse_error: Some(Status::unimplemented("old gateway")),
461 ..Default::default()
462 })
463 .await;
464 let mut client = Client::connect(&host).await.unwrap();
465 assert!(matches!(
466 client.open_browse("S", 20).await.unwrap_err(),
467 Error::IncompatibleGateway { .. }
468 ));
469 }
470
471 #[tokio::test]
472 async fn close_browse_session_forwards_id_and_error() {
473 let service = MockBridgeService::default();
474 let requests = Arc::clone(&service.close_requests);
475 let host = start_mock_server(service).await;
476 let mut client = Client::connect(&host).await.unwrap();
477 client.close_browse_session("session").await.unwrap();
478 assert_eq!(requests.lock().unwrap()[0].session_id, "session");
479
480 let host = start_mock_server(MockBridgeService {
481 close_error: Some(Status::not_found("missing")),
482 ..Default::default()
483 })
484 .await;
485 let mut client = Client::connect(&host).await.unwrap();
486 assert!(matches!(
487 client.close_browse_session("missing").await.unwrap_err(),
488 Error::Rpc(_)
489 ));
490 }
491
492 #[tokio::test]
493 async fn search_stream_and_collect_map_events_and_request() {
494 let events = vec![
495 ProtoSearchEvent {
496 event: Some(search_event::Event::Progress(SearchProgress {
497 visited_nodes: 5,
498 matches: 0,
499 partial: true,
500 })),
501 },
502 ProtoSearchEvent {
503 event: Some(search_event::Event::Completed(SearchCompleted {
504 complete: true,
505 cancelled: false,
506 truncated: false,
507 warning: None,
508 })),
509 },
510 ];
511 let service = MockBridgeService {
512 search_events: events,
513 ..Default::default()
514 };
515 let requests = Arc::clone(&service.search_requests);
516 let host = start_mock_server(service).await;
517 let mut client = Client::connect(&host).await.unwrap();
518 let mut request = SearchRequest::new("S", "PV", SearchMatchMode::Prefix);
519 request.session_id = Some("session".into());
520 request.scope_node_key = Some("scope".into());
521 request.max_results = 50;
522 request.include_branches = true;
523 request.refresh = true;
524 let found = client.search(request).await.unwrap();
525 assert_eq!(found.len(), 2);
526 let request = &requests.lock().unwrap()[0];
527 assert_eq!(request.query, "PV");
528 assert_eq!(
529 request.match_mode,
530 opcda_bridge_proto::bridge::SearchMatchMode::Prefix as i32
531 );
532 assert!(request.include_branches);
533 assert!(request.refresh);
534 }
535
536 #[tokio::test]
537 async fn search_initial_stream_and_protocol_errors_are_typed() {
538 let host = start_mock_server(MockBridgeService {
539 search_initial_error: Some(Status::unavailable("down")),
540 ..Default::default()
541 })
542 .await;
543 let mut client = Client::connect(&host).await.unwrap();
544 assert!(matches!(
545 client
546 .search_stream(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
547 .await
548 .unwrap_err(),
549 Error::Rpc(_)
550 ));
551
552 let host = start_mock_server(MockBridgeService {
553 search_stream_error: Some(Status::deadline_exceeded("slow")),
554 ..Default::default()
555 })
556 .await;
557 let mut client = Client::connect(&host).await.unwrap();
558 assert!(matches!(
559 client
560 .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
561 .await
562 .unwrap_err(),
563 Error::Rpc(_)
564 ));
565
566 let host = start_mock_server(MockBridgeService {
567 search_events: vec![ProtoSearchEvent::default()],
568 ..Default::default()
569 })
570 .await;
571 let mut client = Client::connect(&host).await.unwrap();
572 assert!(matches!(
573 client
574 .search(SearchRequest::new("S", "PV", SearchMatchMode::Exact))
575 .await
576 .unwrap_err(),
577 Error::Protocol(_)
578 ));
579
580 assert!(matches!(
581 feature_error("test", Status::unimplemented("old")),
582 Error::IncompatibleGateway { .. }
583 ));
584 }
585
586 fn index_status(state: ProtoSearchIndexState) -> SearchIndexStatus {
587 SearchIndexStatus {
588 server: "S".into(),
589 state: state as i32,
590 configured: true,
591 active_generation: 4,
592 entry_count: 100,
593 unique_item_count: 99,
594 started_at: Some("start".into()),
595 completed_at: Some("complete".into()),
596 last_error: None,
597 database_bytes: 1024,
598 organization: ProtoOrganization::Hierarchical as i32,
599 source: ProtoBrowseSource::Da3 as i32,
600 progress: Some(IndexedSearchProgress {
601 branches_visited: 2,
602 entries_seen: 3,
603 unique_items: 3,
604 active_time_ms: 4,
605 paused_time_ms: 5,
606 items_per_second: 6.0,
607 estimated_remaining_ms: Some(7),
608 }),
609 }
610 }
611
612 #[tokio::test]
613 async fn indexed_search_methods_map_requests_responses_and_errors() {
614 let service = MockBridgeService {
615 search_index_status_response: index_status(ProtoSearchIndexState::Ready),
616 refresh_search_index_response: index_status(ProtoSearchIndexState::Refreshing),
617 control_search_index_response: index_status(ProtoSearchIndexState::Partial),
618 search_index_response: SearchIndexResponse {
619 matches: vec![IndexedSearchMatch {
620 item_id: "Exact.ItemID".into(),
621 display_name: "PV".into(),
622 kind: opcda_bridge_proto::bridge::BrowseNodeKind::Item as i32,
623 breadcrumbs: vec!["Area".into()],
624 }],
625 has_more: true,
626 status: Some(index_status(ProtoSearchIndexState::Stale)),
627 },
628 ..Default::default()
629 };
630 let status_requests = Arc::clone(&service.search_index_status_requests);
631 let refresh_requests = Arc::clone(&service.refresh_search_index_requests);
632 let control_requests = Arc::clone(&service.control_search_index_requests);
633 let search_requests = Arc::clone(&service.search_index_requests);
634 let host = start_mock_server(service).await;
635 let mut client = Client::connect(&host).await.unwrap();
636
637 assert_eq!(
638 client.search_index_status("S").await.unwrap().state,
639 SearchIndexState::Ready
640 );
641 assert_eq!(
642 client.refresh_search_index("S", true).await.unwrap().state,
643 SearchIndexState::Refreshing
644 );
645 assert_eq!(
646 client
647 .control_search_index("S", SearchIndexControlAction::Pause)
648 .await
649 .unwrap()
650 .state,
651 SearchIndexState::Partial
652 );
653 let mut request = SearchIndexRequest::new("S", "PV", SearchMatchMode::Contains);
654 request.max_results = 25;
655 let response = client.search_index(request).await.unwrap();
656 assert_eq!(response.matches[0].item_id, "Exact.ItemID");
657 assert!(response.has_more);
658
659 assert_eq!(status_requests.lock().unwrap()[0].server, "S");
660 assert!(refresh_requests.lock().unwrap()[0].force);
661 assert_eq!(
662 control_requests.lock().unwrap()[0].action,
663 opcda_bridge_proto::bridge::SearchIndexControlAction::Pause as i32
664 );
665 assert_eq!(search_requests.lock().unwrap()[0].max_results, 25);
666
667 for (field, operation) in [
668 ("search_index_status_error", "status"),
669 ("refresh_search_index_error", "refresh"),
670 ("control_search_index_error", "control"),
671 ("search_index_error", "search"),
672 ] {
673 let mut service = MockBridgeService::default();
674 let error = Some(Status::unimplemented("old gateway"));
675 match field {
676 "search_index_status_error" => service.search_index_status_error = error,
677 "refresh_search_index_error" => service.refresh_search_index_error = error,
678 "control_search_index_error" => service.control_search_index_error = error,
679 "search_index_error" => service.search_index_error = error,
680 _ => unreachable!(),
681 }
682 let host = start_mock_server(service).await;
683 let mut client = Client::connect(&host).await.unwrap();
684 let error = match operation {
685 "status" => client.search_index_status("S").await.unwrap_err(),
686 "refresh" => client.refresh_search_index("S", false).await.unwrap_err(),
687 "control" => client
688 .control_search_index("S", SearchIndexControlAction::Cancel)
689 .await
690 .unwrap_err(),
691 "search" => client
692 .search_index(SearchIndexRequest::new(
693 "S",
694 "PV",
695 SearchMatchMode::Contains,
696 ))
697 .await
698 .unwrap_err(),
699 _ => unreachable!(),
700 };
701 assert!(matches!(error, Error::IncompatibleGateway { .. }));
702 }
703 }
704
705 #[tokio::test]
706 async fn read_maps_data_and_errors() {
707 let host = start_mock_server(MockBridgeService {
708 read_response: ReadResponse {
709 values: vec![ProtoTagValue {
710 tag_id: "t1".into(),
711 value: "42".into(),
712 quality: "Good".into(),
713 timestamp: "now".into(),
714 }],
715 },
716 ..Default::default()
717 })
718 .await;
719 let mut client = Client::connect(&host).await.unwrap();
720 assert_eq!(
721 client.read("S".into(), vec![]).await.unwrap()[0].value,
722 "42"
723 );
724
725 let host = start_mock_server(MockBridgeService {
726 read_error: Some(Status::internal("boom")),
727 ..Default::default()
728 })
729 .await;
730 let mut client = Client::connect(&host).await.unwrap();
731 assert!(matches!(
732 client.read("S".into(), vec![]).await.unwrap_err(),
733 Error::Rpc(_)
734 ));
735 }
736
737 #[tokio::test]
738 async fn write_maps_every_value_and_result_or_error() {
739 for value in [
740 Value::Bool(true),
741 Value::Int(42),
742 Value::Float(3.5),
743 Value::String("text".into()),
744 ] {
745 let host = start_mock_server(MockBridgeService {
746 write_response: WriteResponse {
747 tag_id: "t".into(),
748 success: true,
749 error: None,
750 },
751 ..Default::default()
752 })
753 .await;
754 let mut client = Client::connect(&host).await.unwrap();
755 assert!(
756 client
757 .write("S".into(), "t".into(), value)
758 .await
759 .unwrap()
760 .success
761 );
762 }
763
764 let host = start_mock_server(MockBridgeService {
765 write_error: Some(Status::internal("boom")),
766 ..Default::default()
767 })
768 .await;
769 let mut client = Client::connect(&host).await.unwrap();
770 assert!(matches!(
771 client
772 .write("S".into(), "t".into(), Value::Int(1))
773 .await
774 .unwrap_err(),
775 Error::Rpc(_)
776 ));
777 }
778}