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        match PeerNext::root(self.peer_layers.clone())
390            .admit(request)
391            .await
392        {
393            Ok(admitted) => match admitted.verified() {
394                Some(verified) => {
395                    self.verified_peers
396                        .lock()
397                        .await
398                        .insert(session.clone(), verified);
399                    if let Some(observation) =
400                        self.candidate_identities.lock().await.remove(&session)
401                    {
402                        observation.send_replace(Some(remote.clone()));
403                    }
404                    PeerAdmission::Admitted(remote)
405                }
406                None => PeerAdmission::Rejected("peer admission produced no VerifiedPeer".into()),
407            },
408            Err(error) => PeerAdmission::Rejected(error.message),
409        }
410    }
411
412    async fn invoke_application(
413        self: &Arc<Self>,
414        effect: EffectId,
415        invocation: ApplicationInvocation,
416        handle: ProtocolCoreHandle,
417    ) -> Option<CoreInput> {
418        let abort = self.cancellation.child_token();
419        self.dispatching.lock().await.insert(effect, abort.clone());
420        let permit = self
421            .dispatch_permits
422            .lock()
423            .await
424            .remove(&invocation.reservation);
425        let Some(_permit) = permit else {
426            self.dispatching.lock().await.remove(&effect);
427            return Some(dispatch_failure(
428                effect,
429                ErrorCode::Busy,
430                "capacity reservation expired",
431            ));
432        };
433        if abort.is_cancelled() {
434            self.dispatching.lock().await.remove(&effect);
435            return Some(dispatch_failure(
436                effect,
437                ErrorCode::Cancelled,
438                "request cancelled",
439            ));
440        }
441        let origin = match invocation.origin {
442            ApplicationOrigin::Client { session } => Origin::Client {
443                session: session.to_string(),
444            },
445            ApplicationOrigin::Peer { session, peer } => Origin::Peer {
446                peer: self
447                    .verified_peers
448                    .lock()
449                    .await
450                    .get(&session)
451                    .cloned()
452                    .unwrap_or_else(|| VerifiedPeer::from_identity(&peer)),
453                session: session.to_string(),
454            },
455        };
456        let mut envelope = invocation.frame.clone().into_envelope();
457        let streaming_body = if let Some(body) = &invocation.frame.body {
458            let Some(body) = handle.claim_body(&invocation.stream.session, body.as_str()) else {
459                self.dispatching.lock().await.remove(&effect);
460                return Some(dispatch_failure(
461                    effect,
462                    ErrorCode::Protocol,
463                    "application body unavailable",
464                ));
465            };
466            match body {
467                unb_runtime::WireBody::Bytes(payload) => {
468                    envelope.payload = payload;
469                    None
470                }
471                unb_runtime::WireBody::Stream(stream) => Some(stream),
472            }
473        } else {
474            None
475        };
476        let snapshot = self.snapshot.load_full();
477        let mut request = match Self::inbound_request(&envelope) {
478            Ok(request) => request,
479            Err(error) => {
480                self.dispatching.lock().await.remove(&effect);
481                return Some(dispatch_failure(effect, error.code, error.message));
482            }
483        };
484        if let Some(stream) = streaming_body {
485            request
486                .extensions_mut()
487                .insert(crate::service::StreamingBody(std::sync::Arc::new(
488                    std::sync::Mutex::new(Some(stream)),
489                )));
490        }
491        let outcome = tokio::select! {
492            biased;
493            () = abort.cancelled() => {
494                self.dispatching.lock().await.remove(&effect);
495                return Some(dispatch_failure(effect, ErrorCode::Cancelled, "request cancelled"));
496            }
497            outcome = self.run_service(snapshot.clone(), request, origin) => outcome,
498        };
499        self.dispatching.lock().await.remove(&effect);
500        let outcome = match outcome {
501            Some(Ok(outcome)) => outcome,
502            Some(Err(error)) => return Some(dispatch_failure(effect, error.code, error.message)),
503            None => {
504                let error = Self::teach_unknown_subject(&snapshot, &invocation.frame.head.subject);
505                return Some(dispatch_failure(effect, error.code, error.message));
506            }
507        };
508        let (parts, body) = outcome.into_parts();
509        match body {
510            ServiceBody::Unary(payload) => Some(CoreInput::DispatchCompleted {
511                effect,
512                result: unary_result(parts, payload, &handle, &invocation.stream.session),
513            }),
514            ServiceBody::Stream(mut stream) => {
515                let key = (
516                    invocation.stream.session.clone(),
517                    invocation.stream.corr.as_str().to_string(),
518                );
519                let cancel = self.cancellation.child_token();
520                self.active.lock().await.insert(key.clone(), cancel.clone());
521                let active = self.active.clone();
522                let response_handle = handle.clone();
523                let response_session = invocation.stream.session.clone();
524                unb_runtime::RuntimeHandle::current().spawn(async move {
525                    'pump: loop {
526                        let item = tokio::select! {
527                            biased;
528                            () = cancel.cancelled() => break,
529                            item = stream.next() => item,
530                        };
531                        let mut result =
532                            stream_result(item, &parts, &response_handle, &response_session);
533                        let mut batch = Vec::new();
534                        let terminal = loop {
535                            let terminal = !matches!(result, Ok(ApplicationResult::Event(_)));
536                            batch.push(CoreInput::DispatchCompleted { effect, result });
537                            if terminal || batch.len() >= STREAM_BATCH {
538                                break terminal;
539                            }
540                            match stream.next().now_or_never() {
541                                Some(item) => {
542                                    result = stream_result(
543                                        item,
544                                        &parts,
545                                        &response_handle,
546                                        &response_session,
547                                    )
548                                }
549                                None => break false,
550                            }
551                        };
552                        if handle.submit_batch(batch).await.is_err() || terminal {
553                            break 'pump;
554                        }
555                    }
556                    active.lock().await.remove(&key);
557                });
558                None
559            }
560        }
561    }
562
563    async fn session_established(
564        self: &Arc<Self>,
565        session: SessionId,
566        peer: unb_core::NodeIdentity,
567    ) {
568        let wire = loop {
569            if let Some(wire) = self.session(session.as_str()).await {
570                break wire;
571            }
572            if self.cancellation.is_cancelled() {
573                return;
574            }
575            tokio::task::yield_now().await;
576        };
577        let outbound = self.outbound_sessions.lock().await.contains(&session);
578        if !outbound
579            && self
580                .connection(&peer.node_id)
581                .is_some_and(|connection| connection.is_terminal())
582        {
583            wire.shutdown();
584            return;
585        }
586        self.session_peers
587            .write()
588            .await
589            .insert(session.to_string(), peer.node_id.clone());
590        let replaced = self.peers.write().await.insert(
591            peer.node_id.clone(),
592            PeerLink {
593                session_id: session.to_string(),
594                wire: wire.clone(),
595                instance_id: peer.instance_id.clone(),
596                outbound,
597            },
598        );
599        if let Some(old) = replaced {
600            if old.session_id != session.as_str() {
601                old.wire.shutdown();
602            }
603        }
604        self.publish_route_change();
605        {
606            let mut connections = self
607                .connections
608                .write()
609                .unwrap_or_else(|poisoned| poisoned.into_inner());
610            if let Some(connection) = connections.get(&peer.node_id).cloned() {
611                if outbound {
612                    return;
613                }
614                if !connection.bind(peer, session.to_string(), wire.clone()) {
615                    wire.shutdown();
616                }
617            } else {
618                connections.insert(
619                    peer.node_id.clone(),
620                    crate::PeerConnection::passive(
621                        Arc::downgrade(self),
622                        peer,
623                        session.to_string(),
624                        wire,
625                    ),
626                );
627            }
628        }
629    }
630
631    async fn session_retired(&self, session: SessionId, reason: RetirementReason) {
632        self.verified_peers.lock().await.remove(&session);
633        self.candidate_identities.lock().await.remove(&session);
634        self.outbound_sessions.lock().await.remove(&session);
635        let peer = self.session_peers.write().await.remove(session.as_str());
636        if let Some(connection) = peer.and_then(|peer| self.connection(&peer)) {
637            connection.retire(session.as_str(), reason);
638        }
639        if let Some(wire) = self.sessions.write().await.remove(session.as_str()) {
640            wire.shutdown();
641        }
642        self.cleanup_session(session.as_str()).await;
643    }
644
645    #[cfg(feature = "hosting")]
646    pub fn serve_ws_upgrade(
647        self: &Arc<Self>,
648        upgrade: axum::extract::ws::WebSocketUpgrade,
649    ) -> axum::response::Response {
650        let node = self.clone();
651        upgrade
652            .max_message_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
653            .max_frame_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
654            .on_upgrade(move |socket| async move {
655                let (pipe, initiator) = unb_transport::ws::accept(socket);
656                let _ = node.attach(Pipe::Piped { pipe, initiator }, None).await;
657            })
658    }
659
660    #[cfg(feature = "hosting")]
661    pub async fn serve_webtransport(
662        self: &Arc<Self>,
663        connection: unb_transport::webtransport::wtransport::Connection,
664    ) -> Result<Arc<Wire>, WsError> {
665        let (pipe, initiator, bodies) = unb_transport::webtransport::accept(connection).await?;
666        Ok(self
667            .attach(Pipe::piped_with_streams(pipe, initiator, bodies), None)
668            .await
669            .0)
670    }
671
672    pub async fn serve_transport(self: &Arc<Self>, transport: Pipe) -> Arc<Wire> {
673        self.attach(transport, None).await.0
674    }
675
676    pub async fn connect_transport(
677        self: &Arc<Self>,
678        peer: &str,
679        transport: Pipe,
680    ) -> Result<Arc<Wire>, WsError> {
681        let candidate = self.establish(transport, Some(peer.to_string())).await;
682        match candidate.outcome(peer).await {
683            Ok(CandidateOutcome::Promoted(_)) => {
684                let _ = n0_future::time::timeout(ROUTE_SYNC_TIMEOUT, candidate.wire.routes_acked())
685                    .await;
686                Ok(candidate.wire)
687            }
688            Ok(CandidateOutcome::Duplicate(_)) => Ok(candidate.wire),
689            Err(error) => {
690                candidate.wire.shutdown();
691                Err(error)
692            }
693        }
694    }
695
696    pub(crate) async fn establish(
697        self: &Arc<Self>,
698        transport: Pipe,
699        expected_peer: Option<String>,
700    ) -> CandidateSession {
701        let (wire, _session, cleaned, identity) = self.attach_inner(transport, expected_peer).await;
702        CandidateSession {
703            wire,
704            cleaned,
705            identity,
706        }
707    }
708
709    pub(crate) async fn attach(
710        self: &Arc<Self>,
711        transport: Pipe,
712        expected_peer: Option<String>,
713    ) -> (Arc<Wire>, SessionId, tokio::sync::oneshot::Receiver<()>) {
714        let session = self.next_session_id();
715        let (wire, cleaned) = self
716            .attach_session(session.clone(), transport, expected_peer)
717            .await;
718        (wire, session, cleaned)
719    }
720
721    async fn attach_inner(
722        self: &Arc<Self>,
723        transport: Pipe,
724        expected_peer: Option<String>,
725    ) -> (
726        Arc<Wire>,
727        SessionId,
728        tokio::sync::oneshot::Receiver<()>,
729        tokio::sync::watch::Receiver<Option<unb_core::NodeIdentity>>,
730    ) {
731        let session = self.next_session_id();
732        self.outbound_sessions.lock().await.insert(session.clone());
733        let (identity_tx, identity) = tokio::sync::watch::channel(None);
734        self.candidate_identities
735            .lock()
736            .await
737            .insert(session.clone(), identity_tx);
738        let (wire, cleaned) = self
739            .attach_session(session.clone(), transport, expected_peer)
740            .await;
741        (wire, session, cleaned, identity)
742    }
743
744    fn next_session_id(&self) -> SessionId {
745        SessionId::from(format!(
746            "sess-{}",
747            self.next_session
748                .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
749                + 1
750        ))
751    }
752
753    async fn attach_session(
754        self: &Arc<Self>,
755        session: SessionId,
756        transport: Pipe,
757        expected_peer: Option<String>,
758    ) -> (Arc<Wire>, tokio::sync::oneshot::Receiver<()>) {
759        let bridge = SessionBridge {
760            node: Arc::downgrade(self),
761            session: session.clone(),
762        };
763        let wire = self
764            .protocol
765            .attach_with_ceiling(
766                session.clone(),
767                transport,
768                expected_peer,
769                bridge,
770                self.ws_collect_ceiling,
771            )
772            .await
773            .expect("protocol core actor unavailable");
774        self.sessions
775            .write()
776            .await
777            .insert(session.to_string(), wire.clone());
778        let (cleaned_tx, cleaned_rx) = tokio::sync::oneshot::channel();
779        let cancellation = self.cancellation.child_token();
780        let closed = wire.clone();
781        unb_runtime::RuntimeHandle::current().spawn(async move {
782            tokio::select! {
783                biased;
784                () = cancellation.cancelled() => {}
785                () = closed.closed() => {}
786            }
787            let _ = cleaned_tx.send(());
788        });
789        (wire, cleaned_rx)
790    }
791
792    async fn cleanup_session(&self, session: &str) {
793        self.peers
794            .write()
795            .await
796            .retain(|_, link| link.session_id != session);
797        self.publish_route_change();
798        self.active.lock().await.retain(|(owner, _), cancel| {
799            let keep = owner.as_str() != session;
800            if !keep {
801                cancel.cancel();
802            }
803            keep
804        });
805    }
806}
807
808fn unary_result(
809    parts: http::response::Parts,
810    payload: Bytes,
811    handle: &ProtocolCoreHandle,
812    session: &SessionId,
813) -> Result<ApplicationResult, ApplicationFailure> {
814    if parts.status.is_client_error() || parts.status.is_server_error() {
815        let value: Value = serde_json::from_slice(&payload).unwrap_or(Value::Null);
816        let code = value
817            .get("code")
818            .and_then(|code| serde_json::from_value(code.clone()).ok())
819            .unwrap_or_else(|| ErrorCode::from_status(parts.status));
820        let message = value
821            .get("message")
822            .and_then(Value::as_str)
823            .map(str::to_string)
824            .unwrap_or_else(|| String::from_utf8_lossy(&payload).into_owned());
825        return Err(ApplicationFailure { code, message });
826    }
827    Ok(ApplicationResult::Response(application_response(
828        parts.status,
829        &parts.headers,
830        register_response_body(handle, session, payload)?,
831    )))
832}
833
834fn stream_result(
835    item: Option<Result<Bytes, crate::handler::HandlerError>>,
836    parts: &http::response::Parts,
837    handle: &ProtocolCoreHandle,
838    session: &SessionId,
839) -> Result<ApplicationResult, ApplicationFailure> {
840    match item {
841        Some(Ok(payload)) => Ok(ApplicationResult::Event(application_response(
842            parts.status,
843            &parts.headers,
844            register_response_body(handle, session, payload)?,
845        ))),
846        Some(Err(error)) => Err(ApplicationFailure {
847            code: error.code,
848            message: error.message,
849        }),
850        None => Ok(ApplicationResult::Finished(application_response(
851            parts.status,
852            &parts.headers,
853            None,
854        ))),
855    }
856}
857
858fn register_response_body(
859    handle: &ProtocolCoreHandle,
860    session: &SessionId,
861    payload: Bytes,
862) -> Result<Option<unb_core::BodyId>, ApplicationFailure> {
863    if payload.is_empty() {
864        return Ok(None);
865    }
866    handle
867        .register_body(session, unb_runtime::WireBody::Bytes(payload))
868        .map(Some)
869        .map_err(|error| ApplicationFailure {
870            code: ErrorCode::Busy,
871            message: error.to_string(),
872        })
873}
874
875fn application_response(
876    status: http::StatusCode,
877    headers: &http::HeaderMap,
878    body: Option<unb_core::BodyId>,
879) -> ApplicationResponse {
880    let mut head = http::Response::new(());
881    *head.status_mut() = status;
882    *head.headers_mut() = headers.clone();
883    ApplicationResponse { head, body }
884}
885
886fn dispatch_failure(effect: EffectId, code: ErrorCode, message: impl Into<String>) -> CoreInput {
887    CoreInput::DispatchCompleted {
888        effect,
889        result: Err(ApplicationFailure {
890            code,
891            message: message.into(),
892        }),
893    }
894}
895
896pub(crate) fn retirement_error(peer: &str, reason: RetirementReason) -> WsError {
897    WsError::Connect(format!(
898        "connection to {peer:?} retired during establishment: {reason:?}"
899    ))
900}