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}