Skip to main content

unb_server/
session.rs

1use std::sync::{Arc, Weak};
2
3use bytes::Bytes;
4use futures_util::{FutureExt, StreamExt};
5use serde_json::Value;
6use unb_core::{
7    ApplicationFailure, ApplicationInvocation, ApplicationOrigin, ApplicationResponse,
8    ApplicationResult, CapacityResult, CoreEffect, CoreInput, EffectId, Envelope, ErrorCode,
9    PeerAdmission, RelayOpenResult, RetirementReason, SessionId, TargetPath, TargetReadinessResult,
10};
11use unb_runtime::{
12    EffectExecutor, EffectFuture, Pipe, ProtocolCoreHandle, SessionHandler, SessionOutcome, Wire,
13    WsError,
14};
15
16use crate::layer::{Origin, ServiceBody};
17use crate::node::{Node, PeerLink};
18use crate::peer::{PeerNext, PeerRequest, VerifiedPeer};
19
20const STREAM_BATCH: usize = 8;
21pub(crate) const ROUTE_SYNC_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
22
23pub(crate) struct CandidateSession {
24    pub wire: Arc<Wire>,
25    pub cleaned: tokio::sync::oneshot::Receiver<()>,
26    identity: tokio::sync::watch::Receiver<Option<unb_core::NodeIdentity>>,
27}
28
29impl CandidateSession {
30    pub async fn outcome(&self, peer: &str) -> Result<CandidateOutcome, WsError> {
31        self.observed_outcome()
32            .await
33            .map_err(|failure| match failure {
34                CandidateFailure::Session(error) => error,
35                CandidateFailure::Retired { reason, .. } => retirement_error(peer, reason),
36                CandidateFailure::MissingIdentity => WsError::Connect(format!(
37                    "connection to {peer:?} completed without an admitted identity"
38                )),
39            })
40    }
41
42    pub(crate) async fn observed_outcome(&self) -> Result<CandidateOutcome, CandidateFailure> {
43        match self.wire.session_outcome().await? {
44            SessionOutcome::Established => {
45                Ok(CandidateOutcome::Promoted(self.observed_identity().await?))
46            }
47            SessionOutcome::Retired(
48                RetirementReason::DuplicateSession | RetirementReason::DuplicateSessionReplaced,
49            ) => Ok(CandidateOutcome::Duplicate(self.observed_identity().await?)),
50            SessionOutcome::Retired(reason) => Err(CandidateFailure::Retired {
51                reason,
52                identity: self.identity.borrow().clone(),
53            }),
54        }
55    }
56
57    async fn observed_identity(&self) -> Result<unb_core::NodeIdentity, CandidateFailure> {
58        let mut identity = self.identity.clone();
59        loop {
60            if let Some(identity) = identity.borrow().clone() {
61                return Ok(identity);
62            }
63            if identity.changed().await.is_err() {
64                return Err(CandidateFailure::MissingIdentity);
65            }
66        }
67    }
68}
69
70pub(crate) enum CandidateFailure {
71    Session(WsError),
72    Retired {
73        reason: RetirementReason,
74        identity: Option<unb_core::NodeIdentity>,
75    },
76    MissingIdentity,
77}
78
79impl From<WsError> for CandidateFailure {
80    fn from(error: WsError) -> Self {
81        CandidateFailure::Session(error)
82    }
83}
84
85#[derive(Debug)]
86pub(crate) enum CandidateOutcome {
87    Promoted(unb_core::NodeIdentity),
88    Duplicate(unb_core::NodeIdentity),
89}
90
91pub(crate) struct ServerEffectExecutor {
92    pub(crate) node: Weak<Node>,
93}
94
95impl EffectExecutor for ServerEffectExecutor {
96    fn execute(&self, effect: CoreEffect, handle: ProtocolCoreHandle) -> EffectFuture {
97        let node = self.node.clone();
98        Box::pin(async move {
99            let node = node.upgrade()?;
100            match effect {
101                CoreEffect::RequestPeerAdmission {
102                    effect,
103                    session,
104                    remote,
105                } => Some(CoreInput::PeerAdmissionCompleted {
106                    effect,
107                    result: node.admit_peer(session, remote).await,
108                }),
109                CoreEffect::CheckDispatchCapacity { effect, .. } => {
110                    let result = match node.dispatch_slots.clone().try_acquire_owned() {
111                        Ok(permit) => {
112                            node.dispatch_permits.lock().await.insert(effect, permit);
113                            CapacityResult::Available
114                        }
115                        Err(_) => {
116                            #[cfg(feature = "observability")]
117                            metrics::counter!("unb_dispatch_busy").increment(1);
118                            CapacityResult::Busy
119                        }
120                    };
121                    Some(CoreInput::CapacityChecked { effect, result })
122                }
123                CoreEffect::InvokeApplication { effect, invocation } => {
124                    node.invoke_application(effect, invocation, handle).await
125                }
126                CoreEffect::OpenRelay {
127                    effect,
128                    source,
129                    peer: _,
130                    frame,
131                } => {
132                    let cancellation = node.cancellation.child_token();
133                    node.dispatching
134                        .lock()
135                        .await
136                        .insert(effect, cancellation.clone());
137                    let target = frame.head.target.clone();
138                    let link = async {
139                        let mut route_changes = node.route_changes();
140                        loop {
141                            match node.snapshot.load().node_core.resolve(&target) {
142                                unb_core::Resolution::Route(peer) => {
143                                    if let Some(link) = node.peer(&peer).await {
144                                        break Ok(link);
145                                    }
146                                }
147                                unb_core::Resolution::Unknown => {
148                                    node.await_target_readiness(&target, &cancellation).await?;
149                                    continue;
150                                }
151                                unb_core::Resolution::Conflicted { owners } => {
152                                    break Err(format!(
153                                        "destination node {target:?} has multiple live incarnations: {}",
154                                        owners.join(", ")
155                                    ));
156                                }
157                                unb_core::Resolution::Local => {
158                                    break Err(format!(
159                                        "relay target {target:?} resolved to the local node"
160                                    ));
161                                }
162                            }
163                            tokio::select! {
164                                biased;
165                                () = cancellation.cancelled() => {
166                                    break Err(format!(
167                                        "source stream ended while waiting for target {target:?}"
168                                    ));
169                                }
170                                changed = route_changes.changed() => {
171                                    if changed.is_err() {
172                                        break Err("route readiness notifications closed".into());
173                                    }
174                                }
175                            }
176                        }
177                    }
178                    .await;
179                    let result = match link {
180                        Ok(link) if !cancellation.is_cancelled() => {
181                            let expects_body = frame.body.is_some();
182                            let claimed = frame
183                                .body
184                                .as_ref()
185                                .and_then(|body| handle.claim_body(&source.session, body.as_str()));
186                            if expects_body && claimed.is_none() {
187                                RelayOpenResult::Failed(ApplicationFailure {
188                                    code: ErrorCode::Cancelled,
189                                    message: "relay body was released before downstream admission"
190                                        .into(),
191                                })
192                            } else {
193                                let (payload, body) = match claimed {
194                                    Some(unb_runtime::WireBody::Bytes(payload)) => (payload, None),
195                                    Some(unb_runtime::WireBody::Stream(body)) => {
196                                        (Bytes::new(), Some(body))
197                                    }
198                                    None => (Bytes::new(), None),
199                                };
200                                let forwarded = frame.into_envelope();
201                                let target_path =
202                                    TargetPath::application(&forwarded.target, &forwarded.subject);
203                                let opened = async {
204                                    let target_path =
205                                        target_path.map_err(|error| error.to_string())?;
206                                    link.wire
207                                        .open_forward_with(
208                                            &target_path.to_string(),
209                                            forwarded.kind,
210                                            payload,
211                                            forwarded.hops,
212                                            forwarded.headers,
213                                            body,
214                                            |_| async {},
215                                        )
216                                        .await
217                                        .map_err(|error| error.to_string())
218                                };
219                                match tokio::select! {
220                                    biased;
221                                    () = cancellation.cancelled() => Err(format!(
222                                        "source stream ended while opening target {target:?}"
223                                    )),
224                                    opened = opened => opened,
225                                } {
226                                    Ok(corr) => RelayOpenResult::Opened(unb_core::StreamKey {
227                                        session: link.session_id.into(),
228                                        corr: corr.into(),
229                                    }),
230                                    Err(error) => RelayOpenResult::Failed(ApplicationFailure {
231                                        code: ErrorCode::PeerUnreachable,
232                                        message: error,
233                                    }),
234                                }
235                            }
236                        }
237                        Ok(_) => RelayOpenResult::Failed(ApplicationFailure {
238                            code: ErrorCode::PeerUnreachable,
239                            message: format!(
240                                "source stream ended while waiting for target {target:?}"
241                            ),
242                        }),
243                        Err(message) => RelayOpenResult::Failed(ApplicationFailure {
244                            code: ErrorCode::PeerUnreachable,
245                            message,
246                        }),
247                    };
248                    node.dispatching.lock().await.remove(&effect);
249                    Some(CoreInput::RelayOpenCompleted { effect, result })
250                }
251                CoreEffect::AwaitTargetReadiness {
252                    effect,
253                    stream: _,
254                    target,
255                } => {
256                    let cancellation = node.cancellation.child_token();
257                    node.dispatching
258                        .lock()
259                        .await
260                        .insert(effect, cancellation.clone());
261                    let result = node
262                        .await_target_readiness(&target, &cancellation)
263                        .await
264                        .map_or_else(
265                            |message| TargetReadinessResult::Unavailable { message },
266                            |_| TargetReadinessResult::Ready,
267                        );
268                    node.dispatching.lock().await.remove(&effect);
269                    Some(CoreInput::TargetReadinessCompleted { effect, result })
270                }
271                CoreEffect::QueryDiscoveryNeighbor { stream, peer, plan } => {
272                    node.query_discovery_neighbor(stream, peer, plan, handle);
273                    None
274                }
275                CoreEffect::QueryDiscoveryTarget {
276                    stream,
277                    peer,
278                    target_path,
279                    plan,
280                } => {
281                    node.query_discovery_target(stream, peer, target_path, plan, handle);
282                    None
283                }
284                CoreEffect::SessionEstablished { session, peer } => {
285                    node.session_established(session.clone(), peer).await;
286                    None
287                }
288                CoreEffect::RouteSnapshotApplied {
289                    session,
290                    peer,
291                    snapshot,
292                    ..
293                } => {
294                    if let Some(connection) = node.connection(&peer) {
295                        connection.replace_destinations(&snapshot);
296                    }
297                    node.snapshot.rcu(|snapshot_state| {
298                        let mut next = (**snapshot_state).clone();
299                        let _ = next
300                            .node_core
301                            .apply_snapshot(session.as_str(), &peer, &snapshot);
302                        next
303                    });
304                    node.publish_route_change();
305                    None
306                }
307                CoreEffect::RouteDeltaApplied {
308                    session,
309                    peer,
310                    delta,
311                    ..
312                } => {
313                    if let Some(connection) = node.connection(&peer) {
314                        connection.apply_destination_delta(&delta);
315                    }
316                    node.snapshot.rcu(|snapshot_state| {
317                        let mut next = (**snapshot_state).clone();
318                        let _ = next.node_core.apply_delta(session.as_str(), &peer, &delta);
319                        next
320                    });
321                    node.publish_route_change();
322                    None
323                }
324                CoreEffect::RouteSessionWithdrawn { session, .. } => {
325                    node.snapshot.rcu(|snapshot_state| {
326                        let mut next = (**snapshot_state).clone();
327                        next.node_core.leave(session.as_str());
328                        next
329                    });
330                    node.publish_route_change();
331                    None
332                }
333                CoreEffect::AbortDispatch { effect, .. } => {
334                    if let Some(abort) = node.dispatching.lock().await.remove(&effect) {
335                        abort.cancel();
336                    }
337                    None
338                }
339                CoreEffect::SessionRetired { session, reason } => {
340                    node.session_retired(session, reason).await;
341                    None
342                }
343                _ => None,
344            }
345        })
346    }
347}
348
349struct SessionBridge {
350    node: Weak<Node>,
351    session: SessionId,
352}
353
354impl SessionHandler for SessionBridge {
355    async fn deliver(&mut self, _envelope: Envelope) {}
356
357    async fn stream_closed(&mut self, operation: unb_core::ClientOperationId) {
358        let Some(node) = self.node.upgrade() else {
359            return;
360        };
361        if let Some(cancel) = node
362            .active
363            .lock()
364            .await
365            .remove(&(self.session.clone(), operation.as_str().to_owned()))
366        {
367            cancel.cancel();
368        };
369    }
370}
371
372impl Node {
373    async fn admit_peer(
374        &self,
375        session: SessionId,
376        remote: unb_core::NodeIdentity,
377    ) -> PeerAdmission {
378        let outbound = self.outbound_sessions.lock().await.contains(&session);
379        if !outbound
380            && self
381                .connection(&remote.node_id)
382                .is_some_and(|connection| connection.is_terminal())
383        {
384            return PeerAdmission::Rejected(
385                "the local peer connection was explicitly disconnected".into(),
386            );
387        }
388        let request = PeerRequest::new(self.identity.clone(), remote.clone());
389        let cancellation = self.cancellation.child_token();
390        self.peer_admissions
391            .lock()
392            .await
393            .insert(session.clone(), cancellation.clone());
394        let admission = PeerNext::root(self.peer_layers.clone()).admit(request);
395        tokio::pin!(admission);
396        let result = tokio::select! {
397            biased;
398            () = cancellation.cancelled() => Err(crate::HandlerError::new(
399                ErrorCode::Cancelled,
400                "peer admission session retired",
401            )),
402            result = &mut admission => result,
403        };
404        self.peer_admissions.lock().await.remove(&session);
405        match result {
406            Ok(admitted) => match admitted.verified() {
407                Some(verified) => {
408                    self.verified_peers
409                        .lock()
410                        .await
411                        .insert(session.clone(), verified);
412                    if let Some(observation) =
413                        self.candidate_identities.lock().await.remove(&session)
414                    {
415                        observation.send_replace(Some(remote.clone()));
416                    }
417                    PeerAdmission::Admitted(remote)
418                }
419                None => PeerAdmission::Rejected("peer admission produced no VerifiedPeer".into()),
420            },
421            Err(error) => PeerAdmission::Rejected(error.message),
422        }
423    }
424
425    async fn invoke_application(
426        self: &Arc<Self>,
427        effect: EffectId,
428        invocation: ApplicationInvocation,
429        handle: ProtocolCoreHandle,
430    ) -> Option<CoreInput> {
431        let abort = self.cancellation.child_token();
432        self.dispatching.lock().await.insert(effect, abort.clone());
433        let permit = self
434            .dispatch_permits
435            .lock()
436            .await
437            .remove(&invocation.reservation);
438        let Some(_permit) = permit else {
439            self.dispatching.lock().await.remove(&effect);
440            return Some(dispatch_failure(
441                effect,
442                ErrorCode::Busy,
443                "capacity reservation expired",
444            ));
445        };
446        if abort.is_cancelled() {
447            self.dispatching.lock().await.remove(&effect);
448            return Some(dispatch_failure(
449                effect,
450                ErrorCode::Cancelled,
451                "request cancelled",
452            ));
453        }
454        let origin = match invocation.origin {
455            ApplicationOrigin::Client { session } => Origin::Client {
456                session: session.to_string(),
457            },
458            ApplicationOrigin::Peer { session, peer } => Origin::Peer {
459                peer: self
460                    .verified_peers
461                    .lock()
462                    .await
463                    .get(&session)
464                    .cloned()
465                    .unwrap_or_else(|| VerifiedPeer::from_identity(&peer)),
466                session: session.to_string(),
467            },
468        };
469        let mut envelope = invocation.frame.clone().into_envelope();
470        let streaming_body = if let Some(body) = &invocation.frame.body {
471            let Some(body) = handle.claim_body(&invocation.stream.session, body.as_str()) else {
472                self.dispatching.lock().await.remove(&effect);
473                return Some(dispatch_failure(
474                    effect,
475                    ErrorCode::Protocol,
476                    "application body unavailable",
477                ));
478            };
479            match body {
480                unb_runtime::WireBody::Bytes(payload) => {
481                    envelope.payload = payload;
482                    None
483                }
484                unb_runtime::WireBody::Stream(stream) => Some(stream),
485            }
486        } else {
487            None
488        };
489        let snapshot = self.snapshot.load_full();
490        let mut request = match Self::inbound_request(&envelope) {
491            Ok(request) => request,
492            Err(error) => {
493                self.dispatching.lock().await.remove(&effect);
494                return Some(dispatch_failure(effect, error.code, error.message));
495            }
496        };
497        if let Some(stream) = streaming_body {
498            request
499                .extensions_mut()
500                .insert(crate::service::StreamingBody(std::sync::Arc::new(
501                    std::sync::Mutex::new(Some(stream)),
502                )));
503        }
504        let outcome = tokio::select! {
505            biased;
506            () = abort.cancelled() => {
507                self.dispatching.lock().await.remove(&effect);
508                return Some(dispatch_failure(effect, ErrorCode::Cancelled, "request cancelled"));
509            }
510            outcome = self.run_service(snapshot.clone(), request, origin) => outcome,
511        };
512        self.dispatching.lock().await.remove(&effect);
513        let outcome = match outcome {
514            Some(Ok(outcome)) => outcome,
515            Some(Err(error)) => return Some(dispatch_failure(effect, error.code, error.message)),
516            None => {
517                let error = Self::teach_unknown_subject(&snapshot, &invocation.frame.head.subject);
518                return Some(dispatch_failure(effect, error.code, error.message));
519            }
520        };
521        let (parts, body) = outcome.into_parts();
522        match body {
523            ServiceBody::Unary(payload) => Some(CoreInput::DispatchCompleted {
524                effect,
525                result: unary_result(parts, payload, &handle, &invocation.stream.session),
526            }),
527            ServiceBody::Stream(mut stream) => {
528                let key = (
529                    invocation.stream.session.clone(),
530                    invocation.stream.corr.as_str().to_string(),
531                );
532                let cancel = self.cancellation.child_token();
533                self.active.lock().await.insert(key.clone(), cancel.clone());
534                let active = self.active.clone();
535                let response_handle = handle.clone();
536                let response_session = invocation.stream.session.clone();
537                unb_runtime::RuntimeHandle::current().spawn(async move {
538                    'pump: loop {
539                        let item = tokio::select! {
540                            biased;
541                            () = cancel.cancelled() => break,
542                            item = stream.next() => item,
543                        };
544                        let mut result =
545                            stream_result(item, &parts, &response_handle, &response_session);
546                        let mut batch = Vec::new();
547                        let terminal = loop {
548                            let terminal = !matches!(result, Ok(ApplicationResult::Event(_)));
549                            batch.push(CoreInput::DispatchCompleted { effect, result });
550                            if terminal || batch.len() >= STREAM_BATCH {
551                                break terminal;
552                            }
553                            match stream.next().now_or_never() {
554                                Some(item) => {
555                                    result = stream_result(
556                                        item,
557                                        &parts,
558                                        &response_handle,
559                                        &response_session,
560                                    )
561                                }
562                                None => break false,
563                            }
564                        };
565                        if handle.submit_batch(batch).await.is_err() || terminal {
566                            break 'pump;
567                        }
568                    }
569                    active.lock().await.remove(&key);
570                });
571                None
572            }
573        }
574    }
575
576    async fn session_established(
577        self: &Arc<Self>,
578        session: SessionId,
579        peer: unb_core::NodeIdentity,
580    ) {
581        let wire = loop {
582            if let Some(wire) = self.session(session.as_str()).await {
583                break wire;
584            }
585            if self.cancellation.is_cancelled() {
586                return;
587            }
588            tokio::task::yield_now().await;
589        };
590        let outbound = self.outbound_sessions.lock().await.contains(&session);
591        if !outbound
592            && self
593                .connection(&peer.node_id)
594                .is_some_and(|connection| connection.is_terminal())
595        {
596            wire.shutdown();
597            return;
598        }
599        self.session_peers
600            .write()
601            .await
602            .insert(session.to_string(), peer.node_id.clone());
603        let replaced = self.peers.write().await.insert(
604            peer.node_id.clone(),
605            PeerLink {
606                session_id: session.to_string(),
607                wire: wire.clone(),
608                instance_id: peer.instance_id.clone(),
609                outbound,
610            },
611        );
612        if let Some(old) = replaced {
613            if old.session_id != session.as_str() {
614                old.wire.shutdown();
615            }
616        }
617        self.publish_route_change();
618        {
619            let mut connections = self
620                .connections
621                .write()
622                .unwrap_or_else(|poisoned| poisoned.into_inner());
623            if let Some(connection) = connections.get(&peer.node_id).cloned() {
624                if outbound {
625                    return;
626                }
627                if !connection.bind(peer, session.to_string(), wire.clone()) {
628                    wire.shutdown();
629                }
630            } else {
631                connections.insert(
632                    peer.node_id.clone(),
633                    crate::PeerConnection::passive(
634                        Arc::downgrade(self),
635                        peer,
636                        session.to_string(),
637                        wire,
638                    ),
639                );
640            }
641        }
642    }
643
644    async fn session_retired(&self, session: SessionId, reason: RetirementReason) {
645        if let Some(admission) = self.peer_admissions.lock().await.remove(&session) {
646            admission.cancel();
647        }
648        self.verified_peers.lock().await.remove(&session);
649        self.candidate_identities.lock().await.remove(&session);
650        self.outbound_sessions.lock().await.remove(&session);
651        let peer = self.session_peers.write().await.remove(session.as_str());
652        if let Some(connection) = peer.and_then(|peer| self.connection(&peer)) {
653            connection.retire(session.as_str(), reason);
654        }
655        if let Some(wire) = self.sessions.write().await.remove(session.as_str()) {
656            wire.shutdown();
657        }
658        self.cleanup_session(session.as_str()).await;
659    }
660
661    #[cfg(feature = "hosting")]
662    pub fn serve_ws_upgrade(
663        self: &Arc<Self>,
664        upgrade: axum::extract::ws::WebSocketUpgrade,
665    ) -> axum::response::Response {
666        let node = self.clone();
667        upgrade
668            .max_message_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
669            .max_frame_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
670            .on_upgrade(move |socket| async move {
671                let (pipe, initiator) = unb_transport::ws::accept(socket);
672                let _ = node.attach(Pipe::Piped { pipe, initiator }, None).await;
673            })
674    }
675
676    #[cfg(feature = "hosting")]
677    pub async fn serve_webtransport(
678        self: &Arc<Self>,
679        connection: unb_transport::webtransport::wtransport::Connection,
680    ) -> Result<Arc<Wire>, WsError> {
681        let (pipe, initiator, bodies) = unb_transport::webtransport::accept(connection).await?;
682        Ok(self
683            .attach(Pipe::piped_with_streams(pipe, initiator, bodies), None)
684            .await
685            .0)
686    }
687
688    pub async fn serve_transport(self: &Arc<Self>, transport: Pipe) -> Arc<Wire> {
689        self.attach(transport, None).await.0
690    }
691
692    pub async fn connect_transport(
693        self: &Arc<Self>,
694        peer: &str,
695        transport: Pipe,
696    ) -> Result<Arc<Wire>, WsError> {
697        let candidate = self.establish(transport, Some(peer.to_string())).await;
698        match candidate.outcome(peer).await {
699            Ok(CandidateOutcome::Promoted(_)) => {
700                let synchronized =
701                    n0_future::time::timeout(ROUTE_SYNC_TIMEOUT, candidate.wire.routes_acked())
702                        .await
703                        .map_err(|_| {
704                            WsError::Connect(format!(
705                        "connected peer {peer:?} did not acknowledge its synchronized routes"
706                    ))
707                        })
708                        .and_then(|result| result);
709                if let Err(error) = synchronized {
710                    candidate.wire.shutdown();
711                    let _ = candidate.cleaned.await;
712                    return Err(error);
713                }
714                Ok(candidate.wire)
715            }
716            Ok(CandidateOutcome::Duplicate(_)) => Ok(candidate.wire),
717            Err(error) => {
718                candidate.wire.shutdown();
719                let _ = candidate.cleaned.await;
720                Err(error)
721            }
722        }
723    }
724
725    pub(crate) async fn establish(
726        self: &Arc<Self>,
727        transport: Pipe,
728        expected_peer: Option<String>,
729    ) -> CandidateSession {
730        let (wire, _session, cleaned, identity) = self.attach_inner(transport, expected_peer).await;
731        CandidateSession {
732            wire,
733            cleaned,
734            identity,
735        }
736    }
737
738    pub(crate) async fn attach(
739        self: &Arc<Self>,
740        transport: Pipe,
741        expected_peer: Option<String>,
742    ) -> (Arc<Wire>, SessionId, tokio::sync::oneshot::Receiver<()>) {
743        let session = self.next_session_id();
744        let (wire, cleaned) = self
745            .attach_session(session.clone(), transport, expected_peer)
746            .await;
747        (wire, session, cleaned)
748    }
749
750    async fn attach_inner(
751        self: &Arc<Self>,
752        transport: Pipe,
753        expected_peer: Option<String>,
754    ) -> (
755        Arc<Wire>,
756        SessionId,
757        tokio::sync::oneshot::Receiver<()>,
758        tokio::sync::watch::Receiver<Option<unb_core::NodeIdentity>>,
759    ) {
760        let session = self.next_session_id();
761        self.outbound_sessions.lock().await.insert(session.clone());
762        let (identity_tx, identity) = tokio::sync::watch::channel(None);
763        self.candidate_identities
764            .lock()
765            .await
766            .insert(session.clone(), identity_tx);
767        let (wire, cleaned) = self
768            .attach_session(session.clone(), transport, expected_peer)
769            .await;
770        (wire, session, cleaned, identity)
771    }
772
773    fn next_session_id(&self) -> SessionId {
774        SessionId::from(format!(
775            "sess-{}",
776            self.next_session
777                .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
778                + 1
779        ))
780    }
781
782    async fn attach_session(
783        self: &Arc<Self>,
784        session: SessionId,
785        transport: Pipe,
786        expected_peer: Option<String>,
787    ) -> (Arc<Wire>, tokio::sync::oneshot::Receiver<()>) {
788        let bridge = SessionBridge {
789            node: Arc::downgrade(self),
790            session: session.clone(),
791        };
792        let wire = self
793            .protocol
794            .attach_with_ceiling(
795                session.clone(),
796                transport,
797                expected_peer,
798                bridge,
799                self.ws_collect_ceiling,
800            )
801            .await
802            .expect("protocol core actor unavailable");
803        self.sessions
804            .write()
805            .await
806            .insert(session.to_string(), wire.clone());
807        let (cleaned_tx, cleaned_rx) = tokio::sync::oneshot::channel();
808        let cancellation = self.cancellation.child_token();
809        let closed = wire.clone();
810        unb_runtime::RuntimeHandle::current().spawn(async move {
811            tokio::select! {
812                biased;
813                () = cancellation.cancelled() => {}
814                () = closed.closed() => {}
815            }
816            let _ = cleaned_tx.send(());
817        });
818        (wire, cleaned_rx)
819    }
820
821    async fn cleanup_session(&self, session: &str) {
822        self.peers
823            .write()
824            .await
825            .retain(|_, link| link.session_id != session);
826        self.publish_route_change();
827        self.active.lock().await.retain(|(owner, _), cancel| {
828            let keep = owner.as_str() != session;
829            if !keep {
830                cancel.cancel();
831            }
832            keep
833        });
834    }
835}
836
837fn unary_result(
838    parts: http::response::Parts,
839    payload: Bytes,
840    handle: &ProtocolCoreHandle,
841    session: &SessionId,
842) -> Result<ApplicationResult, ApplicationFailure> {
843    if parts.status.is_client_error() || parts.status.is_server_error() {
844        let value: Value = serde_json::from_slice(&payload).unwrap_or(Value::Null);
845        let code = value
846            .get("code")
847            .and_then(|code| serde_json::from_value(code.clone()).ok())
848            .unwrap_or_else(|| ErrorCode::from_status(parts.status));
849        let message = value
850            .get("message")
851            .and_then(Value::as_str)
852            .map(str::to_string)
853            .unwrap_or_else(|| String::from_utf8_lossy(&payload).into_owned());
854        return Err(ApplicationFailure { code, message });
855    }
856    Ok(ApplicationResult::Response(application_response(
857        parts.status,
858        &parts.headers,
859        register_response_body(handle, session, payload)?,
860    )))
861}
862
863fn stream_result(
864    item: Option<Result<Bytes, crate::handler::HandlerError>>,
865    parts: &http::response::Parts,
866    handle: &ProtocolCoreHandle,
867    session: &SessionId,
868) -> Result<ApplicationResult, ApplicationFailure> {
869    match item {
870        Some(Ok(payload)) => Ok(ApplicationResult::Event(application_response(
871            parts.status,
872            &parts.headers,
873            register_response_body(handle, session, payload)?,
874        ))),
875        Some(Err(error)) => Err(ApplicationFailure {
876            code: error.code,
877            message: error.message,
878        }),
879        None => Ok(ApplicationResult::Finished(application_response(
880            parts.status,
881            &parts.headers,
882            None,
883        ))),
884    }
885}
886
887fn register_response_body(
888    handle: &ProtocolCoreHandle,
889    session: &SessionId,
890    payload: Bytes,
891) -> Result<Option<unb_core::BodyId>, ApplicationFailure> {
892    if payload.is_empty() {
893        return Ok(None);
894    }
895    handle
896        .register_body(session, unb_runtime::WireBody::Bytes(payload))
897        .map(Some)
898        .map_err(|error| ApplicationFailure {
899            code: ErrorCode::Busy,
900            message: error.to_string(),
901        })
902}
903
904fn application_response(
905    status: http::StatusCode,
906    headers: &http::HeaderMap,
907    body: Option<unb_core::BodyId>,
908) -> ApplicationResponse {
909    let mut head = http::Response::new(());
910    *head.status_mut() = status;
911    *head.headers_mut() = headers.clone();
912    ApplicationResponse { head, body }
913}
914
915fn dispatch_failure(effect: EffectId, code: ErrorCode, message: impl Into<String>) -> CoreInput {
916    CoreInput::DispatchCompleted {
917        effect,
918        result: Err(ApplicationFailure {
919            code,
920            message: message.into(),
921        }),
922    }
923}
924
925pub(crate) fn retirement_error(peer: &str, reason: RetirementReason) -> WsError {
926    WsError::Connect(format!(
927        "connection to {peer:?} retired during establishment: {reason:?}"
928    ))
929}