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}