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(crate) async fn fetch_body_until(
297 self: &Arc<Self>,
298 request: http::Request<WireBody>,
299 deadline: Instant,
300 ) -> Result<http::Response<WireBody>, HandlerError> {
301 if let Some(name) = request
302 .headers()
303 .keys()
304 .find(|name| name.as_str().starts_with("unb-"))
305 {
306 return Err(HandlerError::new(
307 ErrorCode::InvalidInput,
308 format!("{name}: unb-* headers are reserved for framing metadata"),
309 ));
310 }
311 let (parts, body) = request.into_parts();
312 if parts.uri.query().is_some() {
313 return Err(HandlerError::new(
314 ErrorCode::InvalidInput,
315 "application request targets must not carry a query string",
316 ));
317 }
318 let target_path = TargetPath::parse_application(parts.uri.path())
319 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
320 let target = target_path.target().to_owned();
321 let subject = target_path.subject().to_owned();
322 let (snapshot, resolution) = self.resolve_unary_until(&target, deadline).await?;
323 match resolution {
324 Resolution::Local => {
325 let (payload, streaming_body) = match body {
326 WireBody::Bytes(payload) => (payload, None),
327 WireBody::Stream(body) => (Bytes::new(), Some(body)),
328 };
329 let envelope = Envelope::from_request(http::Request::from_parts(parts, payload))
330 .map_err(|error| {
331 HandlerError::new(ErrorCode::InvalidInput, error.to_string())
332 })?;
333 let mut request = Self::inbound_request(&envelope)?;
334 if let Some(body) = streaming_body {
335 request
336 .extensions_mut()
337 .insert(crate::service::StreamingBody(Arc::new(
338 std::sync::Mutex::new(Some(body)),
339 )));
340 }
341 match self
342 .run_service(snapshot.clone(), request, Origin::Local)
343 .await
344 {
345 Some(outcome) => {
346 let response = outcome?;
347 let (parts, body) = response.into_parts();
348 let body = match body {
349 crate::layer::ServiceBody::Unary(payload) => WireBody::Bytes(payload),
350 crate::layer::ServiceBody::Stream(body) => {
351 WireBody::Stream(Box::pin(body.map(|item| {
352 item.map_err(|error| {
353 unb_core::CoreError::Malformed(error.to_string())
354 })
355 })))
356 }
357 };
358 Ok(http::Response::from_parts(parts, body))
359 }
360 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
361 }
362 }
363 Resolution::Route(peer_name) => {
364 let link = self.route_link(&peer_name).await?;
365 let remaining = Self::remaining_unary_time(deadline)?;
366 let mut response = link
367 .wire
368 .client_session()
369 .fetch_body(http::Request::from_parts(parts, body), remaining)
370 .await
371 .map_err(Self::client_error)?;
372 let reserved = response
373 .headers()
374 .keys()
375 .filter(|name| {
376 name.as_str().starts_with("unb-") && name.as_str() != unb_core::UNB_CODE
377 })
378 .cloned()
379 .collect::<Vec<_>>();
380 for name in reserved {
381 response.headers_mut().remove(name);
382 }
383 Ok(response)
384 }
385 Resolution::Conflicted { owners } => Err(HandlerError::new(
386 ErrorCode::PeerUnreachable,
387 format!(
388 "destination node {target:?} has multiple live incarnations: {}",
389 owners.join(", ")
390 ),
391 )),
392 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
393 }
394 }
395
396 pub async fn fetch(
397 self: &Arc<Self>,
398 request: http::Request<bytes::Bytes>,
399 ) -> Result<http::Response<crate::layer::ServiceBody>, HandlerError> {
400 self.fetch_until(request, Instant::now() + CALL_TIMEOUT)
401 .await
402 }
403
404 pub(crate) async fn fetch_until(
405 self: &Arc<Self>,
406 request: http::Request<bytes::Bytes>,
407 deadline: Instant,
408 ) -> Result<http::Response<crate::layer::ServiceBody>, HandlerError> {
409 let (parts, body) = request.into_parts();
410 let response = self
411 .fetch_body_until(
412 http::Request::from_parts(parts, WireBody::Bytes(body)),
413 deadline,
414 )
415 .await?;
416 let (parts, body) = response.into_parts();
417 let body = match body {
418 WireBody::Bytes(payload) => crate::layer::ServiceBody::Unary(payload),
419 WireBody::Stream(body) => {
420 crate::layer::ServiceBody::Stream(Box::pin(body.map(|item| {
421 item.map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))
422 })))
423 }
424 };
425 Ok(http::Response::from_parts(parts, body))
426 }
427
428 pub async fn subscribe(
429 self: &Arc<Self>,
430 target_path: &str,
431 payload: Value,
432 ) -> Result<crate::EventStream, HandlerError> {
433 self.subscribe_with(target_path, payload, serde_json::Map::new())
434 .await
435 }
436
437 pub async fn subscribe_with(
438 self: &Arc<Self>,
439 target_path: &str,
440 payload: Value,
441 headers: serde_json::Map<String, Value>,
442 ) -> Result<crate::EventStream, HandlerError> {
443 self.subscribe_bytes(target_path, Envelope::encode_payload(&payload), headers)
444 .await
445 }
446
447 pub async fn subscribe_bytes(
448 self: &Arc<Self>,
449 target_path: &str,
450 payload: Bytes,
451 headers: serde_json::Map<String, Value>,
452 ) -> Result<crate::EventStream, HandlerError> {
453 let target_path = TargetPath::parse_application(target_path)
454 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
455 let target = target_path.target().to_owned();
456 let subject = target_path.subject().to_owned();
457 let target_path = target_path.to_string();
458 let (snapshot, resolution) = self
459 .resolve_unary_until(&target, Instant::now() + CALL_TIMEOUT)
460 .await?;
461 match resolution {
462 Resolution::Local => {
463 let request =
464 self.local_request(Kind::Subscribe, &subject, payload, headers.clone())?;
465 match self
466 .run_service(snapshot.clone(), request, Origin::Local)
467 .await
468 {
469 Some(Ok(response)) => match response.into_body() {
470 crate::layer::ServiceBody::Stream(stream) => Ok(stream),
471 crate::layer::ServiceBody::Unary(_) => Err(HandlerError::new(
472 ErrorCode::Internal,
473 "a streaming operation produced a unary response",
474 )),
475 },
476 Some(Err(error)) => Err(error),
477 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
478 }
479 }
480 Resolution::Route(peer_name) => {
481 let link = self.peer(&peer_name).await.ok_or_else(|| {
482 HandlerError::new(
483 ErrorCode::PeerUnreachable,
484 format!("no live connection to peer {peer_name:?}"),
485 )
486 })?;
487 let stream = link
488 .wire
489 .client_session()
490 .start(
491 &target_path,
492 Kind::Subscribe,
493 payload,
494 Some(DEFAULT_HOPS),
495 headers,
496 )
497 .await
498 .map_err(|error| {
499 HandlerError::new(ErrorCode::PeerUnreachable, error.to_string())
500 })?;
501 Ok(Box::pin(stream::unfold(stream, |mut stream| async move {
502 let item = match stream.next().await {
503 Ok(Some(envelope)) if envelope.kind == Kind::Event => {
504 Some(Ok(envelope.payload))
505 }
506 Ok(Some(envelope)) if envelope.kind == Kind::Response => None,
507 Ok(Some(_)) => Some(Err(HandlerError::new(
508 ErrorCode::Protocol,
509 "unexpected frame in subscription",
510 ))),
511 Ok(None) => None,
512 Err(error) => Some(Err(Node::client_error(error))),
513 };
514 item.map(|item| (item, stream))
515 })))
516 }
517 Resolution::Conflicted { owners } => Err(HandlerError::new(
518 ErrorCode::PeerUnreachable,
519 format!(
520 "destination node {target:?} has multiple live incarnations: {}",
521 owners.join(", ")
522 ),
523 )),
524 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
525 }
526 }
527
528 pub(crate) async fn call_nested(
529 self: &Arc<Self>,
530 target_path: &str,
531 payload: Value,
532 headers: serde_json::Map<String, Value>,
533 ) -> Result<Value, HandlerError> {
534 self.call_with_origin(target_path, payload, headers, Origin::Nested)
535 .await
536 }
537
538 async fn call_with_origin(
539 self: &Arc<Self>,
540 target_path: &str,
541 payload: Value,
542 headers: serde_json::Map<String, Value>,
543 origin: Origin,
544 ) -> Result<Value, HandlerError> {
545 let target_path = TargetPath::parse_application(target_path)
546 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
547 let target = target_path.target().to_owned();
548 let subject = target_path.subject().to_owned();
549 let target_path = target_path.to_string();
550 let deadline = Instant::now() + CALL_TIMEOUT;
551 let (snapshot, resolution) = self.resolve_unary_until(&target, deadline).await?;
552 match resolution {
553 Resolution::Local => {
554 let request = self.local_request(
555 Kind::Request,
556 &subject,
557 Envelope::encode_payload(&payload),
558 headers,
559 )?;
560 let outcome = self.run_service(snapshot.clone(), request, origin).await;
561 match outcome {
562 Some(outcome) => match outcome?.into_body() {
563 crate::layer::ServiceBody::Unary(payload) => {
564 Self::json_profile_payload(&payload)
565 }
566 crate::layer::ServiceBody::Stream(_) => Err(HandlerError::new(
567 ErrorCode::Internal,
568 "a unary operation produced a stream",
569 )),
570 },
571 None => Err(Self::teach_unknown_subject(&snapshot, &subject)),
572 }
573 }
574 Resolution::Route(peer_name) => {
575 let link = self.route_link(&peer_name).await?;
576 self.call_peer(link, &target_path, payload, headers, DEFAULT_HOPS, deadline)
577 .await
578 }
579 Resolution::Conflicted { owners } => Err(HandlerError::new(
580 ErrorCode::PeerUnreachable,
581 format!(
582 "destination node {target:?} has multiple live incarnations: {}",
583 owners.join(", ")
584 ),
585 )),
586 Resolution::Unknown => Err(Self::teach_unknown_target(&snapshot, &target)),
587 }
588 }
589
590 pub(crate) async fn route_link(
591 &self,
592 peer_name: &str,
593 ) -> Result<crate::node::PeerLink, HandlerError> {
594 self.peer(peer_name).await.ok_or_else(|| {
595 HandlerError::new(
596 ErrorCode::PeerUnreachable,
597 format!("no live connection to peer {peer_name:?}"),
598 )
599 })
600 }
601
602 async fn call_peer(
603 &self,
604 link: crate::node::PeerLink,
605 target_path: &str,
606 payload: Value,
607 headers: serde_json::Map<String, Value>,
608 hops: u8,
609 deadline: Instant,
610 ) -> Result<Value, HandlerError> {
611 let reply = self
612 .call_peer_envelope(
613 link,
614 target_path,
615 Envelope::encode_payload(&payload),
616 headers,
617 hops,
618 deadline,
619 )
620 .await?;
621 Self::json_profile_payload(&reply.payload)
622 }
623
624 async fn call_peer_envelope(
625 &self,
626 link: crate::node::PeerLink,
627 target_path: &str,
628 payload: bytes::Bytes,
629 headers: serde_json::Map<String, Value>,
630 hops: u8,
631 deadline: Instant,
632 ) -> Result<Envelope, HandlerError> {
633 let remaining = Self::remaining_unary_time(deadline)?;
634 let operation = async move {
635 let mut stream = link
636 .wire
637 .client_session()
638 .start(target_path, Kind::Request, payload, Some(hops), headers)
639 .await
640 .map_err(|error| {
641 HandlerError::new(ErrorCode::PeerUnreachable, error.to_string())
642 })?;
643 match stream.next().await {
644 Ok(Some(envelope)) if envelope.kind == Kind::Response => Ok(envelope),
645 Err(error) => Err(Self::client_error(error)),
646 Ok(Some(_)) => Err(HandlerError::new(
647 ErrorCode::Protocol,
648 "downstream call returned an unexpected frame",
649 )),
650 Ok(None) => Err(HandlerError::new(
651 ErrorCode::PeerUnreachable,
652 "downstream call did not complete",
653 )),
654 }
655 };
656 match n0_future::time::timeout(remaining, operation).await {
657 Ok(result) => result,
658 Err(_) => Err(HandlerError::new(
659 ErrorCode::PeerUnreachable,
660 "downstream call did not complete before its deadline",
661 )),
662 }
663 }
664
665 pub(crate) async fn resolve_unary_until(
666 &self,
667 target: &str,
668 deadline: Instant,
669 ) -> Result<(Arc<NodeSnapshot>, Resolution), HandlerError> {
670 let mut route_changes = self.route_changes();
671 let readiness_waits = self.readiness_waits_for_destination(target);
672 let snapshot = self.snapshot.load_full();
673 let resolution = snapshot.node_core.resolve(target);
674 if !matches!(resolution, Resolution::Unknown) {
675 return Ok((snapshot, resolution));
676 }
677 if readiness_waits.is_empty() {
678 return Ok((snapshot, Resolution::Unknown));
679 }
680 let mut waiting = stream::FuturesUnordered::new();
681 for readiness in readiness_waits {
682 waiting.push(readiness.wait());
683 }
684 let mut restored = false;
685 while !waiting.is_empty() {
686 let remaining = Self::remaining_unary_time(deadline)?;
687 let wake = n0_future::time::timeout(remaining, async {
688 tokio::select! {
689 result = waiting.next() => (result, false),
690 changed = route_changes.changed() => {
691 let _ = changed;
692 (None, true)
693 }
694 }
695 })
696 .await
697 .map_err(|_| {
698 HandlerError::new(
699 ErrorCode::PeerUnreachable,
700 format!(
701 "recovery did not restore target {target:?} before the request deadline"
702 ),
703 )
704 })?;
705 if !wake.1 && wake.0.is_some_and(|result| result.is_ok()) {
706 restored = true;
707 }
708 let snapshot = self.snapshot.load_full();
709 let resolution = snapshot.node_core.resolve(target);
710 if !matches!(resolution, Resolution::Unknown) {
711 return Ok((snapshot, resolution));
712 }
713 if waiting.is_empty() {
714 if restored {
715 return Ok((snapshot, Resolution::Unknown));
716 }
717 return Err(HandlerError::new(
718 ErrorCode::PeerUnreachable,
719 format!("recovery did not restore target {target:?}"),
720 ));
721 }
722 }
723 Ok((snapshot, Resolution::Unknown))
724 }
725
726 pub(crate) async fn await_target_readiness(
727 &self,
728 target: &str,
729 cancellation: &unb_runtime::CancellationToken,
730 ) -> Result<(), String> {
731 let mut route_changes = self.route_changes();
732 let readiness_waits = self.readiness_waits_for_destination(target);
733 if !matches!(
734 self.snapshot.load().node_core.resolve(target),
735 Resolution::Unknown
736 ) {
737 return Ok(());
738 }
739 if readiness_waits.is_empty() {
740 return Err(format!(
741 "no recovering connection carried target {target:?}"
742 ));
743 }
744 let mut waiting = stream::FuturesUnordered::new();
745 for readiness in readiness_waits {
746 waiting.push(readiness.wait());
747 }
748 loop {
749 if !matches!(
750 self.snapshot.load().node_core.resolve(target),
751 Resolution::Unknown
752 ) {
753 return Ok(());
754 }
755 if waiting.is_empty() {
756 return Err(format!("recovery did not restore target {target:?}"));
757 }
758 tokio::select! {
759 biased;
760 () = cancellation.cancelled() => {
761 return Err(format!("source stream ended while waiting for target {target:?}"));
762 }
763 changed = route_changes.changed() => {
764 if changed.is_err() {
765 return Err("route readiness notifications closed".into());
766 }
767 }
768 result = waiting.next() => {
769 if result.is_some_and(|result| result.is_err()) && waiting.is_empty() {
770 return Err(format!("recovery did not restore target {target:?}"));
771 }
772 }
773 }
774 }
775 }
776
777 fn remaining_unary_time(deadline: Instant) -> Result<Duration, HandlerError> {
778 let remaining = deadline.saturating_duration_since(Instant::now());
779 if remaining.is_zero() {
780 Err(HandlerError::new(
781 ErrorCode::PeerUnreachable,
782 "unary request deadline elapsed",
783 ))
784 } else {
785 Ok(remaining)
786 }
787 }
788
789 fn json_profile_payload(payload: &Bytes) -> Result<Value, HandlerError> {
790 if payload.is_empty() {
791 return Ok(Value::Null);
792 }
793 serde_json::from_slice(payload)
794 .map_err(|error| HandlerError::new(ErrorCode::Protocol, error.to_string()))
795 }
796
797 fn client_error(error: unb_runtime::ClientError) -> HandlerError {
798 match error {
799 unb_runtime::ClientError::Protocol { code, message, .. } => {
800 HandlerError::new(code, message)
801 }
802 unb_runtime::ClientError::Cancelled(_) => {
803 HandlerError::new(ErrorCode::Cancelled, error.to_string())
804 }
805 unb_runtime::ClientError::Invalid(message) => {
806 HandlerError::new(ErrorCode::InvalidInput, message)
807 }
808 _ => HandlerError::new(ErrorCode::PeerUnreachable, error.to_string()),
809 }
810 }
811}