Skip to main content

unb_server/
call.rs

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}