1use std::future::Future;
2use std::sync::Arc;
3use std::time::Duration;
4
5use bytes::Bytes;
6use futures_util::{stream, StreamExt};
7use n0_future::time::Instant;
8use serde_json::Value;
9use unb_core::{Envelope, ErrorCode, Kind, Resolution, TargetPath, DEFAULT_HOPS};
10use unb_runtime::WireBody;
11
12use crate::handler::HandlerError;
13use crate::layer::Origin;
14use crate::node::{Node, NodeSnapshot};
15
16pub(crate) const CALL_TIMEOUT: Duration = Duration::from_secs(30);
17
18pub trait IntoBody: Send {
19 fn into_body(self) -> Bytes;
20}
21
22impl IntoBody for Bytes {
23 fn into_body(self) -> Bytes {
24 self
25 }
26}
27
28impl IntoBody for Vec<u8> {
29 fn into_body(self) -> Bytes {
30 self.into()
31 }
32}
33
34impl IntoBody for String {
35 fn into_body(self) -> Bytes {
36 self.into()
37 }
38}
39
40impl IntoBody for &str {
41 fn into_body(self) -> Bytes {
42 Bytes::copy_from_slice(self.as_bytes())
43 }
44}
45
46impl IntoBody for Value {
47 fn into_body(self) -> Bytes {
48 Envelope::encode_payload(&self)
49 }
50}
51
52impl IntoBody for () {
53 fn into_body(self) -> Bytes {
54 Bytes::new()
55 }
56}
57
58pub trait Destination: Send + Sync {
59 fn send(
60 &self,
61 request: http::Request<Bytes>,
62 ) -> impl Future<Output = Result<http::Response<Bytes>, HandlerError>> + Send;
63}
64
65impl Destination for Arc<Node> {
66 async fn send(
67 &self,
68 request: http::Request<Bytes>,
69 ) -> Result<http::Response<Bytes>, HandlerError> {
70 let response = self.fetch(request).await?;
71 let (parts, body) = response.into_parts();
72 match body {
73 crate::layer::ServiceBody::Unary(payload) => {
74 Ok(http::Response::from_parts(parts, payload))
75 }
76 crate::layer::ServiceBody::Stream(_) => Err(HandlerError::new(
77 ErrorCode::InvalidInput,
78 "subscribe is not available over send; use Node::subscribe",
79 )),
80 }
81 }
82}
83
84impl<D: Destination + ?Sized> Destination for &D {
85 async fn send(
86 &self,
87 request: http::Request<Bytes>,
88 ) -> Result<http::Response<Bytes>, HandlerError> {
89 D::send(self, request).await
90 }
91}
92
93#[cfg(feature = "hosting")]
94impl Destination for &str {
95 async fn send(
96 &self,
97 request: http::Request<Bytes>,
98 ) -> Result<http::Response<Bytes>, HandlerError> {
99 n0_future::time::timeout(CALL_TIMEOUT, one_shot_http(self, request))
100 .await
101 .map_err(|_| {
102 HandlerError::new(
103 ErrorCode::PeerUnreachable,
104 format!("{self:?} did not answer within the call timeout"),
105 )
106 })?
107 }
108}
109
110#[cfg(feature = "hosting")]
111impl Destination for String {
112 async fn send(
113 &self,
114 request: http::Request<Bytes>,
115 ) -> Result<http::Response<Bytes>, HandlerError> {
116 self.as_str().send(request).await
117 }
118}
119
120pub trait SendExt<T> {
121 fn send<D: Destination>(
122 self,
123 destination: D,
124 ) -> impl Future<Output = Result<http::Response<Bytes>, HandlerError>> + Send;
125}
126
127impl<T: IntoBody> SendExt<T> for http::Request<T> {
128 async fn send<D: Destination>(
129 self,
130 destination: D,
131 ) -> Result<http::Response<Bytes>, HandlerError> {
132 let (mut parts, body) = self.into_parts();
133 if parts.method == http::Method::GET {
134 parts.method = http::Method::POST;
135 }
136 destination
137 .send(http::Request::from_parts(parts, body.into_body()))
138 .await
139 }
140}
141
142impl<T: IntoBody> SendExt<T> for Result<http::Request<T>, http::Error> {
143 async fn send<D: Destination>(
144 self,
145 destination: D,
146 ) -> Result<http::Response<Bytes>, HandlerError> {
147 match self {
148 Ok(request) => request.send(destination).await,
149 Err(error) => Err(HandlerError::new(
150 ErrorCode::InvalidInput,
151 error.to_string(),
152 )),
153 }
154 }
155}
156
157#[cfg(feature = "hosting")]
158async fn one_shot_http(
159 address: &str,
160 request: http::Request<Bytes>,
161) -> Result<http::Response<Bytes>, HandlerError> {
162 let (tls, remainder) = if let Some(rest) = address.strip_prefix("https://") {
163 (true, rest)
164 } else if let Some(rest) = address.strip_prefix("wss://") {
165 (true, rest)
166 } else if let Some(rest) = address.strip_prefix("http://") {
167 (false, rest)
168 } else if let Some(rest) = address.strip_prefix("ws://") {
169 (false, rest)
170 } else {
171 (false, address)
172 };
173 let authority = remainder
174 .split('/')
175 .next()
176 .filter(|authority| !authority.is_empty())
177 .ok_or_else(|| {
178 HandlerError::new(
179 ErrorCode::InvalidInput,
180 format!("{address:?} names no host to send to"),
181 )
182 })?;
183 let has_port = match authority.rfind(']') {
184 Some(bracket) => authority[bracket + 1..].contains(':'),
185 None => authority.contains(':'),
186 };
187 let authority = if has_port {
188 authority.to_string()
189 } else {
190 format!("{authority}:{}", if tls { 443 } else { 80 })
191 };
192 let unreachable = |error: String| HandlerError::new(ErrorCode::PeerUnreachable, error);
193 let stream = tokio::net::TcpStream::connect(&authority)
194 .await
195 .map_err(|error| unreachable(error.to_string()))?;
196 if tls {
197 let host = authority
198 .rsplit_once(':')
199 .map(|(host, _)| host)
200 .unwrap_or(&authority)
201 .trim_start_matches('[')
202 .trim_end_matches(']');
203 let server_name = rustls_pki_types::ServerName::try_from(host.to_string())
204 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
205 let config = unb_transport::ws::tls_client_config()
206 .map_err(|error| unreachable(error.to_string()))?;
207 let stream = tokio_rustls::TlsConnector::from(config)
208 .connect(server_name, stream)
209 .await
210 .map_err(|error| unreachable(error.to_string()))?;
211 exchange_http1(stream, &authority, request).await
212 } else {
213 exchange_http1(stream, &authority, request).await
214 }
215}
216
217#[cfg(feature = "hosting")]
218async fn exchange_http1<T>(
219 stream: T,
220 authority: &str,
221 request: http::Request<Bytes>,
222) -> Result<http::Response<Bytes>, HandlerError>
223where
224 T: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static,
225{
226 let target = request
227 .uri()
228 .path_and_query()
229 .map(|target| target.as_str())
230 .filter(|target| !target.is_empty())
231 .unwrap_or("/");
232 let unreachable = |error: String| HandlerError::new(ErrorCode::PeerUnreachable, error);
233 let (mut sender, connection) =
234 hyper::client::conn::http1::handshake(hyper_util::rt::TokioIo::new(stream))
235 .await
236 .map_err(|error| unreachable(error.to_string()))?;
237 tokio::spawn(async move {
238 let _ = connection.await;
239 });
240 let mut outbound = http::Request::builder()
241 .method(http::Method::POST)
242 .uri(target)
243 .header(http::header::HOST, authority);
244 for (name, value) in request.headers() {
245 if matches!(
246 *name,
247 http::header::HOST
248 | http::header::CONTENT_LENGTH
249 | http::header::TRANSFER_ENCODING
250 | http::header::CONNECTION
251 ) {
252 continue;
253 }
254 outbound = outbound.header(name, value);
255 }
256 let outbound = outbound
257 .body(http_body_util::Full::new(request.into_body()))
258 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
259 let response = sender
260 .send_request(outbound)
261 .await
262 .map_err(|error| unreachable(error.to_string()))?;
263 let (parts, body) = response.into_parts();
264 let body = http_body_util::Limited::new(body, unb_transport::DEFAULT_MAX_FRAME_SIZE);
265 let body = http_body_util::BodyExt::collect(body)
266 .await
267 .map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))?
268 .to_bytes();
269 let mut projected = http::Response::builder().status(parts.status);
270 for (name, value) in &parts.headers {
271 if matches!(
272 *name,
273 http::header::CONNECTION
274 | http::header::CONTENT_LENGTH
275 | http::header::TRANSFER_ENCODING
276 | http::header::DATE
277 ) {
278 continue;
279 }
280 projected = projected.header(name, value);
281 }
282 projected
283 .body(body)
284 .map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))
285}
286
287impl Node {
288 pub async fn fetch_body(
289 self: &Arc<Self>,
290 request: http::Request<WireBody>,
291 ) -> Result<http::Response<WireBody>, HandlerError> {
292 self.fetch_body_until(request, Instant::now() + CALL_TIMEOUT)
293 .await
294 }
295
296 pub async fn fetch_body_with_timeout(
297 self: &Arc<Self>,
298 request: http::Request<WireBody>,
299 timeout: Duration,
300 ) -> Result<http::Response<WireBody>, HandlerError> {
301 self.fetch_body_until(request, Instant::now() + timeout)
302 .await
303 }
304
305 pub(crate) async fn fetch_body_until(
306 self: &Arc<Self>,
307 request: http::Request<WireBody>,
308 deadline: Instant,
309 ) -> Result<http::Response<WireBody>, HandlerError> {
310 if let Some(name) = request
311 .headers()
312 .keys()
313 .find(|name| name.as_str().starts_with("unb-"))
314 {
315 return Err(HandlerError::new(
316 ErrorCode::InvalidInput,
317 format!("{name}: unb-* headers are reserved for framing metadata"),
318 ));
319 }
320 let (parts, body) = request.into_parts();
321 if parts.uri.query().is_some() {
322 return Err(HandlerError::new(
323 ErrorCode::InvalidInput,
324 "application request targets must not carry a query string",
325 ));
326 }
327 let target_path = TargetPath::parse_application(parts.uri.path())
328 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
329 let target = target_path.target().to_owned();
330 let subject = target_path.subject().to_owned();
331 let (snapshot, resolution) = self.resolve_unary_until(&target, deadline).await?;
332 match resolution {
333 Resolution::Local => {
334 let (payload, streaming_body) = match body {
335 WireBody::Bytes(payload) => (payload, None),
336 WireBody::Stream(body) => (Bytes::new(), Some(body)),
337 };
338 let envelope = Envelope::from_request(http::Request::from_parts(parts, payload))
339 .map_err(|error| {
340 HandlerError::new(ErrorCode::InvalidInput, error.to_string())
341 })?;
342 let mut request = Self::inbound_request(&envelope)?;
343 if let Some(body) = streaming_body {
344 request
345 .extensions_mut()
346 .insert(crate::service::StreamingBody(Arc::new(
347 std::sync::Mutex::new(Some(body)),
348 )));
349 }
350 match self
351 .run_service(snapshot.clone(), request, Origin::Local)
352 .await
353 {
354 Some(outcome) => {
355 let response = outcome?;
356 let (parts, body) = response.into_parts();
357 let body = match body {
358 crate::layer::ServiceBody::Unary(payload) => WireBody::Bytes(payload),
359 crate::layer::ServiceBody::Stream(body) => {
360 WireBody::Stream(Box::pin(body.map(|item| {
361 item.map_err(|error| {
362 unb_core::CoreError::Malformed(error.to_string())
363 })
364 })))
365 }
366 };
367 Ok(http::Response::from_parts(parts, body))
368 }
369 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
370 }
371 }
372 Resolution::Route(peer_name) => {
373 let link = self.route_link(&peer_name).await?;
374 let remaining = Self::remaining_unary_time(deadline)?;
375 let mut response = link
376 .wire
377 .client_session()
378 .fetch_body(http::Request::from_parts(parts, body), remaining)
379 .await
380 .map_err(Self::client_error)?;
381 let reserved = response
382 .headers()
383 .keys()
384 .filter(|name| {
385 name.as_str().starts_with("unb-") && name.as_str() != unb_core::UNB_CODE
386 })
387 .cloned()
388 .collect::<Vec<_>>();
389 for name in reserved {
390 response.headers_mut().remove(name);
391 }
392 Ok(response)
393 }
394 Resolution::Conflicted { owners } => Err(HandlerError::new(
395 ErrorCode::PeerUnreachable,
396 format!(
397 "destination node {target:?} has multiple live incarnations: {}",
398 owners.join(", ")
399 ),
400 )),
401 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
402 }
403 }
404
405 pub async fn fetch(
406 self: &Arc<Self>,
407 request: http::Request<bytes::Bytes>,
408 ) -> Result<http::Response<crate::layer::ServiceBody>, HandlerError> {
409 self.fetch_until(request, Instant::now() + CALL_TIMEOUT)
410 .await
411 }
412
413 pub(crate) async fn fetch_until(
414 self: &Arc<Self>,
415 request: http::Request<bytes::Bytes>,
416 deadline: Instant,
417 ) -> Result<http::Response<crate::layer::ServiceBody>, HandlerError> {
418 let (parts, body) = request.into_parts();
419 let response = self
420 .fetch_body_until(
421 http::Request::from_parts(parts, WireBody::Bytes(body)),
422 deadline,
423 )
424 .await?;
425 let (parts, body) = response.into_parts();
426 let body = match body {
427 WireBody::Bytes(payload) => crate::layer::ServiceBody::Unary(payload),
428 WireBody::Stream(body) => {
429 crate::layer::ServiceBody::Stream(Box::pin(body.map(|item| {
430 item.map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))
431 })))
432 }
433 };
434 Ok(http::Response::from_parts(parts, body))
435 }
436
437 pub async fn subscribe(
438 self: &Arc<Self>,
439 target_path: &str,
440 payload: Value,
441 ) -> Result<crate::EventStream, HandlerError> {
442 self.subscribe_with(target_path, payload, serde_json::Map::new())
443 .await
444 }
445
446 pub async fn subscribe_with(
447 self: &Arc<Self>,
448 target_path: &str,
449 payload: Value,
450 headers: serde_json::Map<String, Value>,
451 ) -> Result<crate::EventStream, HandlerError> {
452 self.subscribe_bytes(target_path, Envelope::encode_payload(&payload), headers)
453 .await
454 }
455
456 pub async fn subscribe_bytes(
457 self: &Arc<Self>,
458 target_path: &str,
459 payload: Bytes,
460 headers: serde_json::Map<String, Value>,
461 ) -> Result<crate::EventStream, HandlerError> {
462 let target_path = TargetPath::parse_application(target_path)
463 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
464 let target = target_path.target().to_owned();
465 let subject = target_path.subject().to_owned();
466 let target_path = target_path.to_string();
467 let (snapshot, resolution) = self
468 .resolve_unary_until(&target, Instant::now() + CALL_TIMEOUT)
469 .await?;
470 match resolution {
471 Resolution::Local => {
472 let request =
473 self.local_request(Kind::Subscribe, &subject, payload, headers.clone())?;
474 match self
475 .run_service(snapshot.clone(), request, Origin::Local)
476 .await
477 {
478 Some(Ok(response)) => match response.into_body() {
479 crate::layer::ServiceBody::Stream(stream) => Ok(stream),
480 crate::layer::ServiceBody::Unary(_) => Err(HandlerError::new(
481 ErrorCode::Internal,
482 "a streaming operation produced a unary response",
483 )),
484 },
485 Some(Err(error)) => Err(error),
486 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
487 }
488 }
489 Resolution::Route(peer_name) => {
490 let link = self.peer(&peer_name).await.ok_or_else(|| {
491 HandlerError::new(
492 ErrorCode::PeerUnreachable,
493 format!("no live connection to peer {peer_name:?}"),
494 )
495 })?;
496 let stream = link
497 .wire
498 .client_session()
499 .start(
500 &target_path,
501 Kind::Subscribe,
502 payload,
503 Some(DEFAULT_HOPS),
504 headers,
505 )
506 .await
507 .map_err(|error| {
508 HandlerError::new(ErrorCode::PeerUnreachable, error.to_string())
509 })?;
510 Ok(Box::pin(stream::unfold(stream, |mut stream| async move {
511 let item = match stream.next().await {
512 Ok(Some(envelope)) if envelope.kind == Kind::Event => {
513 Some(Ok(envelope.payload))
514 }
515 Ok(Some(envelope)) if envelope.kind == Kind::Response => None,
516 Ok(Some(_)) => Some(Err(HandlerError::new(
517 ErrorCode::Protocol,
518 "unexpected frame in subscription",
519 ))),
520 Ok(None) => None,
521 Err(error) => Some(Err(Node::client_error(error))),
522 };
523 item.map(|item| (item, stream))
524 })))
525 }
526 Resolution::Conflicted { owners } => Err(HandlerError::new(
527 ErrorCode::PeerUnreachable,
528 format!(
529 "destination node {target:?} has multiple live incarnations: {}",
530 owners.join(", ")
531 ),
532 )),
533 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
534 }
535 }
536
537 pub(crate) async fn call_nested(
538 self: &Arc<Self>,
539 target_path: &str,
540 payload: Value,
541 headers: serde_json::Map<String, Value>,
542 ) -> Result<Value, HandlerError> {
543 self.call_with_origin(target_path, payload, headers, Origin::Nested)
544 .await
545 }
546
547 async fn call_with_origin(
548 self: &Arc<Self>,
549 target_path: &str,
550 payload: Value,
551 headers: serde_json::Map<String, Value>,
552 origin: Origin,
553 ) -> Result<Value, HandlerError> {
554 let target_path = TargetPath::parse_application(target_path)
555 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
556 let target = target_path.target().to_owned();
557 let subject = target_path.subject().to_owned();
558 let target_path = target_path.to_string();
559 let deadline = Instant::now() + CALL_TIMEOUT;
560 let (snapshot, resolution) = self.resolve_unary_until(&target, deadline).await?;
561 match resolution {
562 Resolution::Local => {
563 let request = self.local_request(
564 Kind::Request,
565 &subject,
566 Envelope::encode_payload(&payload),
567 headers,
568 )?;
569 let outcome = self.run_service(snapshot.clone(), request, origin).await;
570 match outcome {
571 Some(outcome) => match outcome?.into_body() {
572 crate::layer::ServiceBody::Unary(payload) => {
573 Self::json_profile_payload(&payload)
574 }
575 crate::layer::ServiceBody::Stream(_) => Err(HandlerError::new(
576 ErrorCode::Internal,
577 "a unary operation produced a stream",
578 )),
579 },
580 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
581 }
582 }
583 Resolution::Route(peer_name) => {
584 let link = self.route_link(&peer_name).await?;
585 self.call_peer(link, &target_path, payload, headers, DEFAULT_HOPS, deadline)
586 .await
587 }
588 Resolution::Conflicted { owners } => Err(HandlerError::new(
589 ErrorCode::PeerUnreachable,
590 format!(
591 "destination node {target:?} has multiple live incarnations: {}",
592 owners.join(", ")
593 ),
594 )),
595 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
596 }
597 }
598
599 pub(crate) async fn route_link(
600 &self,
601 peer_name: &str,
602 ) -> Result<crate::node::PeerLink, HandlerError> {
603 self.peer(peer_name).await.ok_or_else(|| {
604 HandlerError::new(
605 ErrorCode::PeerUnreachable,
606 format!("no live connection to peer {peer_name:?}"),
607 )
608 })
609 }
610
611 async fn call_peer(
612 &self,
613 link: crate::node::PeerLink,
614 target_path: &str,
615 payload: Value,
616 headers: serde_json::Map<String, Value>,
617 hops: u8,
618 deadline: Instant,
619 ) -> Result<Value, HandlerError> {
620 let reply = self
621 .call_peer_envelope(
622 link,
623 target_path,
624 Envelope::encode_payload(&payload),
625 headers,
626 hops,
627 deadline,
628 )
629 .await?;
630 Self::json_profile_payload(&reply.payload)
631 }
632
633 async fn call_peer_envelope(
634 &self,
635 link: crate::node::PeerLink,
636 target_path: &str,
637 payload: bytes::Bytes,
638 headers: serde_json::Map<String, Value>,
639 hops: u8,
640 deadline: Instant,
641 ) -> Result<Envelope, HandlerError> {
642 let remaining = Self::remaining_unary_time(deadline)?;
643 let operation = async move {
644 let mut stream = link
645 .wire
646 .client_session()
647 .start(target_path, Kind::Request, payload, Some(hops), headers)
648 .await
649 .map_err(|error| {
650 HandlerError::new(ErrorCode::PeerUnreachable, error.to_string())
651 })?;
652 match stream.next().await {
653 Ok(Some(envelope)) if envelope.kind == Kind::Response => Ok(envelope),
654 Err(error) => Err(Self::client_error(error)),
655 Ok(Some(_)) => Err(HandlerError::new(
656 ErrorCode::Protocol,
657 "downstream call returned an unexpected frame",
658 )),
659 Ok(None) => Err(HandlerError::new(
660 ErrorCode::PeerUnreachable,
661 "downstream call did not complete",
662 )),
663 }
664 };
665 match n0_future::time::timeout(remaining, operation).await {
666 Ok(result) => result,
667 Err(_) => Err(HandlerError::new(
668 ErrorCode::PeerUnreachable,
669 "downstream call did not complete before its deadline",
670 )),
671 }
672 }
673
674 pub(crate) async fn resolve_unary_until(
675 &self,
676 target: &str,
677 deadline: Instant,
678 ) -> Result<(Arc<NodeSnapshot>, Resolution), HandlerError> {
679 let mut route_changes = self.route_changes();
680 let readiness_waits = self.readiness_waits_for_destination(target);
681 let snapshot = self.snapshot.load_full();
682 let resolution = snapshot.node_core.resolve(target);
683 if !matches!(resolution, Resolution::Unknown) {
684 return Ok((snapshot, resolution));
685 }
686 if readiness_waits.is_empty() {
687 return Ok((snapshot, Resolution::Unknown));
688 }
689 let mut waiting = stream::FuturesUnordered::new();
690 for readiness in readiness_waits {
691 waiting.push(readiness.wait());
692 }
693 let mut restored = false;
694 while !waiting.is_empty() {
695 let remaining = Self::remaining_unary_time(deadline)?;
696 let wake = n0_future::time::timeout(remaining, async {
697 tokio::select! {
698 result = waiting.next() => (result, false),
699 changed = route_changes.changed() => {
700 let _ = changed;
701 (None, true)
702 }
703 }
704 })
705 .await
706 .map_err(|_| {
707 HandlerError::new(
708 ErrorCode::PeerUnreachable,
709 format!(
710 "recovery did not restore target {target:?} before the request deadline"
711 ),
712 )
713 })?;
714 if !wake.1 && wake.0.is_some_and(|result| result.is_ok()) {
715 restored = true;
716 }
717 let snapshot = self.snapshot.load_full();
718 let resolution = snapshot.node_core.resolve(target);
719 if !matches!(resolution, Resolution::Unknown) {
720 return Ok((snapshot, resolution));
721 }
722 if waiting.is_empty() {
723 if restored {
724 return Ok((snapshot, Resolution::Unknown));
725 }
726 return Err(HandlerError::new(
727 ErrorCode::PeerUnreachable,
728 format!("recovery did not restore target {target:?}"),
729 ));
730 }
731 }
732 Ok((snapshot, Resolution::Unknown))
733 }
734
735 pub(crate) async fn await_target_readiness(
736 &self,
737 target: &str,
738 cancellation: &unb_runtime::CancellationToken,
739 ) -> Result<(), String> {
740 let mut route_changes = self.route_changes();
741 let readiness_waits = self.readiness_waits_for_destination(target);
742 if !matches!(
743 self.snapshot.load().node_core.resolve(target),
744 Resolution::Unknown
745 ) {
746 return Ok(());
747 }
748 if readiness_waits.is_empty() {
749 return Err(format!(
750 "no recovering connection carried target {target:?}"
751 ));
752 }
753 let mut waiting = stream::FuturesUnordered::new();
754 for readiness in readiness_waits {
755 waiting.push(readiness.wait());
756 }
757 loop {
758 if !matches!(
759 self.snapshot.load().node_core.resolve(target),
760 Resolution::Unknown
761 ) {
762 return Ok(());
763 }
764 if waiting.is_empty() {
765 return Err(format!("recovery did not restore target {target:?}"));
766 }
767 tokio::select! {
768 biased;
769 () = cancellation.cancelled() => {
770 return Err(format!("source stream ended while waiting for target {target:?}"));
771 }
772 changed = route_changes.changed() => {
773 if changed.is_err() {
774 return Err("route readiness notifications closed".into());
775 }
776 }
777 result = waiting.next() => {
778 if result.is_some_and(|result| result.is_err()) && waiting.is_empty() {
779 return Err(format!("recovery did not restore target {target:?}"));
780 }
781 }
782 }
783 }
784 }
785
786 fn remaining_unary_time(deadline: Instant) -> Result<Duration, HandlerError> {
787 let remaining = deadline.saturating_duration_since(Instant::now());
788 if remaining.is_zero() {
789 Err(HandlerError::new(
790 ErrorCode::PeerUnreachable,
791 "unary request deadline elapsed",
792 ))
793 } else {
794 Ok(remaining)
795 }
796 }
797
798 fn json_profile_payload(payload: &Bytes) -> Result<Value, HandlerError> {
799 if payload.is_empty() {
800 return Ok(Value::Null);
801 }
802 serde_json::from_slice(payload)
803 .map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))
804 }
805
806 fn client_error(error: unb_runtime::ClientError) -> HandlerError {
807 match error {
808 unb_runtime::ClientError::Protocol { code, message, .. } => {
809 HandlerError::new(code, message)
810 }
811 unb_runtime::ClientError::Cancelled(_) => {
812 HandlerError::new(ErrorCode::Cancelled, error.to_string())
813 }
814 unb_runtime::ClientError::Invalid(message) => {
815 HandlerError::new(ErrorCode::InvalidInput, message)
816 }
817 _ => HandlerError::new(ErrorCode::PeerUnreachable, error.to_string()),
818 }
819 }
820}