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 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}