use std::sync::{Arc, Weak};
use bytes::Bytes;
use futures_util::{FutureExt, StreamExt};
use serde_json::Value;
use unb_core::{
ApplicationFailure, ApplicationInvocation, ApplicationOrigin, ApplicationResponse,
ApplicationResult, CapacityResult, CoreEffect, CoreInput, EffectId, Envelope, ErrorCode,
PeerAdmission, RelayOpenResult, RetirementReason, SessionId, TargetPath, TargetReadinessResult,
};
use unb_runtime::{
EffectExecutor, EffectFuture, Pipe, ProtocolCoreHandle, SessionHandler, SessionOutcome, Wire,
WsError,
};
use crate::layer::{Origin, ServiceBody};
use crate::node::{Node, PeerLink};
use crate::peer::{PeerNext, PeerRequest, VerifiedPeer};
const STREAM_BATCH: usize = 8;
pub(crate) const ROUTE_SYNC_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
pub(crate) struct CandidateSession {
pub wire: Arc<Wire>,
pub cleaned: tokio::sync::oneshot::Receiver<()>,
identity: tokio::sync::watch::Receiver<Option<unb_core::NodeIdentity>>,
}
impl CandidateSession {
pub async fn outcome(&self, peer: &str) -> Result<CandidateOutcome, WsError> {
self.observed_outcome()
.await
.map_err(|failure| match failure {
CandidateFailure::Session(error) => error,
CandidateFailure::Retired { reason, .. } => retirement_error(peer, reason),
CandidateFailure::MissingIdentity => WsError::Connect(format!(
"connection to {peer:?} completed without an admitted identity"
)),
})
}
pub(crate) async fn observed_outcome(&self) -> Result<CandidateOutcome, CandidateFailure> {
match self.wire.session_outcome().await? {
SessionOutcome::Established => {
Ok(CandidateOutcome::Promoted(self.observed_identity().await?))
}
SessionOutcome::Retired(
RetirementReason::DuplicateSession | RetirementReason::DuplicateSessionReplaced,
) => Ok(CandidateOutcome::Duplicate(self.observed_identity().await?)),
SessionOutcome::Retired(reason) => Err(CandidateFailure::Retired {
reason,
identity: self.identity.borrow().clone(),
}),
}
}
async fn observed_identity(&self) -> Result<unb_core::NodeIdentity, CandidateFailure> {
let mut identity = self.identity.clone();
loop {
if let Some(identity) = identity.borrow().clone() {
return Ok(identity);
}
if identity.changed().await.is_err() {
return Err(CandidateFailure::MissingIdentity);
}
}
}
}
pub(crate) enum CandidateFailure {
Session(WsError),
Retired {
reason: RetirementReason,
identity: Option<unb_core::NodeIdentity>,
},
MissingIdentity,
}
impl From<WsError> for CandidateFailure {
fn from(error: WsError) -> Self {
CandidateFailure::Session(error)
}
}
#[derive(Debug)]
pub(crate) enum CandidateOutcome {
Promoted(unb_core::NodeIdentity),
Duplicate(unb_core::NodeIdentity),
}
pub(crate) struct ServerEffectExecutor {
pub(crate) node: Weak<Node>,
}
impl EffectExecutor for ServerEffectExecutor {
fn execute(&self, effect: CoreEffect, handle: ProtocolCoreHandle) -> EffectFuture {
let node = self.node.clone();
Box::pin(async move {
let node = node.upgrade()?;
match effect {
CoreEffect::RequestPeerAdmission {
effect,
session,
remote,
} => Some(CoreInput::PeerAdmissionCompleted {
effect,
result: node.admit_peer(session, remote).await,
}),
CoreEffect::CheckDispatchCapacity { effect, .. } => {
let result = match node.dispatch_slots.clone().try_acquire_owned() {
Ok(permit) => {
node.dispatch_permits.lock().await.insert(effect, permit);
CapacityResult::Available
}
Err(_) => {
#[cfg(feature = "observability")]
metrics::counter!("unb_dispatch_busy").increment(1);
CapacityResult::Busy
}
};
Some(CoreInput::CapacityChecked { effect, result })
}
CoreEffect::InvokeApplication { effect, invocation } => {
node.invoke_application(effect, invocation, handle).await
}
CoreEffect::OpenRelay {
effect,
source,
peer: _,
frame,
} => {
let cancellation = node.cancellation.child_token();
node.dispatching
.lock()
.await
.insert(effect, cancellation.clone());
let target = frame.head.target.clone();
let link = async {
let mut route_changes = node.route_changes();
loop {
match node.snapshot.load().node_core.resolve(&target) {
unb_core::Resolution::Route(peer) => {
if let Some(link) = node.peer(&peer).await {
break Ok(link);
}
}
unb_core::Resolution::Unknown => {
node.await_target_readiness(&target, &cancellation).await?;
continue;
}
unb_core::Resolution::Conflicted { owners } => {
break Err(format!(
"destination node {target:?} has multiple live incarnations: {}",
owners.join(", ")
));
}
unb_core::Resolution::Local => {
break Err(format!(
"relay target {target:?} resolved to the local node"
));
}
}
tokio::select! {
biased;
() = cancellation.cancelled() => {
break Err(format!(
"source stream ended while waiting for target {target:?}"
));
}
changed = route_changes.changed() => {
if changed.is_err() {
break Err("route readiness notifications closed".into());
}
}
}
}
}
.await;
let result = match link {
Ok(link) if !cancellation.is_cancelled() => {
let expects_body = frame.body.is_some();
let claimed = frame
.body
.as_ref()
.and_then(|body| handle.claim_body(&source.session, body.as_str()));
if expects_body && claimed.is_none() {
RelayOpenResult::Failed(ApplicationFailure {
code: ErrorCode::Cancelled,
message: "relay body was released before downstream admission"
.into(),
})
} else {
let (payload, body) = match claimed {
Some(unb_runtime::WireBody::Bytes(payload)) => (payload, None),
Some(unb_runtime::WireBody::Stream(body)) => {
(Bytes::new(), Some(body))
}
None => (Bytes::new(), None),
};
let forwarded = frame.into_envelope();
let target_path =
TargetPath::application(&forwarded.target, &forwarded.subject);
let opened = async {
let target_path =
target_path.map_err(|error| error.to_string())?;
link.wire
.open_forward_with(
&target_path.to_string(),
forwarded.kind,
payload,
forwarded.hops,
forwarded.headers,
body,
|_| async {},
)
.await
.map_err(|error| error.to_string())
};
match tokio::select! {
biased;
() = cancellation.cancelled() => Err(format!(
"source stream ended while opening target {target:?}"
)),
opened = opened => opened,
} {
Ok(corr) => RelayOpenResult::Opened(unb_core::StreamKey {
session: link.session_id.into(),
corr: corr.into(),
}),
Err(error) => RelayOpenResult::Failed(ApplicationFailure {
code: ErrorCode::PeerUnreachable,
message: error,
}),
}
}
}
Ok(_) => RelayOpenResult::Failed(ApplicationFailure {
code: ErrorCode::PeerUnreachable,
message: format!(
"source stream ended while waiting for target {target:?}"
),
}),
Err(message) => RelayOpenResult::Failed(ApplicationFailure {
code: ErrorCode::PeerUnreachable,
message,
}),
};
node.dispatching.lock().await.remove(&effect);
Some(CoreInput::RelayOpenCompleted { effect, result })
}
CoreEffect::AwaitTargetReadiness {
effect,
stream: _,
target,
} => {
let cancellation = node.cancellation.child_token();
node.dispatching
.lock()
.await
.insert(effect, cancellation.clone());
let result = node
.await_target_readiness(&target, &cancellation)
.await
.map_or_else(
|message| TargetReadinessResult::Unavailable { message },
|_| TargetReadinessResult::Ready,
);
node.dispatching.lock().await.remove(&effect);
Some(CoreInput::TargetReadinessCompleted { effect, result })
}
CoreEffect::QueryDiscoveryNeighbor { stream, peer, plan } => {
node.query_discovery_neighbor(stream, peer, plan, handle);
None
}
CoreEffect::QueryDiscoveryTarget {
stream,
peer,
target_path,
plan,
} => {
node.query_discovery_target(stream, peer, target_path, plan, handle);
None
}
CoreEffect::SessionEstablished { session, peer } => {
node.session_established(session.clone(), peer).await;
None
}
CoreEffect::RouteSnapshotApplied {
session,
peer,
snapshot,
..
} => {
if let Some(connection) = node.connection(&peer) {
connection.replace_destinations(&snapshot);
}
node.snapshot.rcu(|snapshot_state| {
let mut next = (**snapshot_state).clone();
let _ = next
.node_core
.apply_snapshot(session.as_str(), &peer, &snapshot);
next
});
node.publish_route_change();
None
}
CoreEffect::RouteDeltaApplied {
session,
peer,
delta,
..
} => {
if let Some(connection) = node.connection(&peer) {
connection.apply_destination_delta(&delta);
}
node.snapshot.rcu(|snapshot_state| {
let mut next = (**snapshot_state).clone();
let _ = next.node_core.apply_delta(session.as_str(), &peer, &delta);
next
});
node.publish_route_change();
None
}
CoreEffect::RouteSessionWithdrawn { session, .. } => {
node.snapshot.rcu(|snapshot_state| {
let mut next = (**snapshot_state).clone();
next.node_core.leave(session.as_str());
next
});
node.publish_route_change();
None
}
CoreEffect::AbortDispatch { effect, .. } => {
if let Some(abort) = node.dispatching.lock().await.remove(&effect) {
abort.cancel();
}
None
}
CoreEffect::SessionRetired { session, reason } => {
node.session_retired(session, reason).await;
None
}
_ => None,
}
})
}
}
struct SessionBridge {
node: Weak<Node>,
session: SessionId,
}
impl SessionHandler for SessionBridge {
async fn deliver(&mut self, _envelope: Envelope) {}
async fn stream_closed(&mut self, operation: unb_core::ClientOperationId) {
let Some(node) = self.node.upgrade() else {
return;
};
if let Some(cancel) = node
.active
.lock()
.await
.remove(&(self.session.clone(), operation.as_str().to_owned()))
{
cancel.cancel();
};
}
}
impl Node {
async fn admit_peer(
&self,
session: SessionId,
remote: unb_core::NodeIdentity,
) -> PeerAdmission {
let outbound = self.outbound_sessions.lock().await.contains(&session);
if !outbound
&& self
.connection(&remote.node_id)
.is_some_and(|connection| connection.is_terminal())
{
return PeerAdmission::Rejected(
"the local peer connection was explicitly disconnected".into(),
);
}
let request = PeerRequest::new(self.identity.clone(), remote.clone());
let cancellation = self.cancellation.child_token();
self.peer_admissions
.lock()
.await
.insert(session.clone(), cancellation.clone());
let admission = PeerNext::root(self.peer_layers.clone()).admit(request);
tokio::pin!(admission);
let result = tokio::select! {
biased;
() = cancellation.cancelled() => Err(crate::HandlerError::new(
ErrorCode::Cancelled,
"peer admission session retired",
)),
result = &mut admission => result,
};
self.peer_admissions.lock().await.remove(&session);
match result {
Ok(admitted) => match admitted.verified() {
Some(verified) => {
self.verified_peers
.lock()
.await
.insert(session.clone(), verified);
if let Some(observation) =
self.candidate_identities.lock().await.remove(&session)
{
observation.send_replace(Some(remote.clone()));
}
PeerAdmission::Admitted(remote)
}
None => PeerAdmission::Rejected("peer admission produced no VerifiedPeer".into()),
},
Err(error) => PeerAdmission::Rejected(error.message),
}
}
async fn invoke_application(
self: &Arc<Self>,
effect: EffectId,
invocation: ApplicationInvocation,
handle: ProtocolCoreHandle,
) -> Option<CoreInput> {
let abort = self.cancellation.child_token();
self.dispatching.lock().await.insert(effect, abort.clone());
let permit = self
.dispatch_permits
.lock()
.await
.remove(&invocation.reservation);
let Some(_permit) = permit else {
self.dispatching.lock().await.remove(&effect);
return Some(dispatch_failure(
effect,
ErrorCode::Busy,
"capacity reservation expired",
));
};
if abort.is_cancelled() {
self.dispatching.lock().await.remove(&effect);
return Some(dispatch_failure(
effect,
ErrorCode::Cancelled,
"request cancelled",
));
}
let origin = match invocation.origin {
ApplicationOrigin::Client { session } => Origin::Client {
session: session.to_string(),
},
ApplicationOrigin::Peer { session, peer } => Origin::Peer {
peer: self
.verified_peers
.lock()
.await
.get(&session)
.cloned()
.unwrap_or_else(|| VerifiedPeer::from_identity(&peer)),
session: session.to_string(),
},
};
let mut envelope = invocation.frame.clone().into_envelope();
let streaming_body = if let Some(body) = &invocation.frame.body {
let Some(body) = handle.claim_body(&invocation.stream.session, body.as_str()) else {
self.dispatching.lock().await.remove(&effect);
return Some(dispatch_failure(
effect,
ErrorCode::Protocol,
"application body unavailable",
));
};
match body {
unb_runtime::WireBody::Bytes(payload) => {
envelope.payload = payload;
None
}
unb_runtime::WireBody::Stream(stream) => Some(stream),
}
} else {
None
};
let snapshot = self.snapshot.load_full();
let mut request = match Self::inbound_request(&envelope) {
Ok(request) => request,
Err(error) => {
self.dispatching.lock().await.remove(&effect);
return Some(dispatch_failure(effect, error.code, error.message));
}
};
if let Some(stream) = streaming_body {
request
.extensions_mut()
.insert(crate::service::StreamingBody(std::sync::Arc::new(
std::sync::Mutex::new(Some(stream)),
)));
}
let outcome = tokio::select! {
biased;
() = abort.cancelled() => {
self.dispatching.lock().await.remove(&effect);
return Some(dispatch_failure(effect, ErrorCode::Cancelled, "request cancelled"));
}
outcome = self.run_service(snapshot.clone(), request, origin) => outcome,
};
self.dispatching.lock().await.remove(&effect);
let outcome = match outcome {
Some(Ok(outcome)) => outcome,
Some(Err(error)) => return Some(dispatch_failure(effect, error.code, error.message)),
None => {
let error = Self::teach_unknown_subject(&snapshot, &invocation.frame.head.subject);
return Some(dispatch_failure(effect, error.code, error.message));
}
};
let (parts, body) = outcome.into_parts();
match body {
ServiceBody::Unary(payload) => Some(CoreInput::DispatchCompleted {
effect,
result: unary_result(parts, payload, &handle, &invocation.stream.session),
}),
ServiceBody::Stream(mut stream) => {
let key = (
invocation.stream.session.clone(),
invocation.stream.corr.as_str().to_string(),
);
let cancel = self.cancellation.child_token();
self.active.lock().await.insert(key.clone(), cancel.clone());
let active = self.active.clone();
let response_handle = handle.clone();
let response_session = invocation.stream.session.clone();
unb_runtime::RuntimeHandle::current().spawn(async move {
'pump: loop {
let item = tokio::select! {
biased;
() = cancel.cancelled() => break,
item = stream.next() => item,
};
let mut result =
stream_result(item, &parts, &response_handle, &response_session);
let mut batch = Vec::new();
let terminal = loop {
let terminal = !matches!(result, Ok(ApplicationResult::Event(_)));
batch.push(CoreInput::DispatchCompleted { effect, result });
if terminal || batch.len() >= STREAM_BATCH {
break terminal;
}
match stream.next().now_or_never() {
Some(item) => {
result = stream_result(
item,
&parts,
&response_handle,
&response_session,
)
}
None => break false,
}
};
if handle.submit_batch(batch).await.is_err() || terminal {
break 'pump;
}
}
active.lock().await.remove(&key);
});
None
}
}
}
async fn session_established(
self: &Arc<Self>,
session: SessionId,
peer: unb_core::NodeIdentity,
) {
let wire = loop {
if let Some(wire) = self.session(session.as_str()).await {
break wire;
}
if self.cancellation.is_cancelled() {
return;
}
tokio::task::yield_now().await;
};
let outbound = self.outbound_sessions.lock().await.contains(&session);
if !outbound
&& self
.connection(&peer.node_id)
.is_some_and(|connection| connection.is_terminal())
{
wire.shutdown();
return;
}
self.session_peers
.write()
.await
.insert(session.to_string(), peer.node_id.clone());
let replaced = self.peers.write().await.insert(
peer.node_id.clone(),
PeerLink {
session_id: session.to_string(),
wire: wire.clone(),
instance_id: peer.instance_id.clone(),
outbound,
},
);
if let Some(old) = replaced {
if old.session_id != session.as_str() {
old.wire.shutdown();
}
}
self.publish_route_change();
{
let mut connections = self
.connections
.write()
.unwrap_or_else(|poisoned| poisoned.into_inner());
if let Some(connection) = connections.get(&peer.node_id).cloned() {
if outbound {
return;
}
if !connection.bind(peer, session.to_string(), wire.clone()) {
wire.shutdown();
}
} else {
connections.insert(
peer.node_id.clone(),
crate::PeerConnection::passive(
Arc::downgrade(self),
peer,
session.to_string(),
wire,
),
);
}
}
}
async fn session_retired(&self, session: SessionId, reason: RetirementReason) {
if let Some(admission) = self.peer_admissions.lock().await.remove(&session) {
admission.cancel();
}
self.verified_peers.lock().await.remove(&session);
self.candidate_identities.lock().await.remove(&session);
self.outbound_sessions.lock().await.remove(&session);
let peer = self.session_peers.write().await.remove(session.as_str());
if let Some(connection) = peer.and_then(|peer| self.connection(&peer)) {
connection.retire(session.as_str(), reason);
}
if let Some(wire) = self.sessions.write().await.remove(session.as_str()) {
wire.shutdown();
}
self.cleanup_session(session.as_str()).await;
}
#[cfg(feature = "hosting")]
pub fn serve_ws_upgrade(
self: &Arc<Self>,
upgrade: axum::extract::ws::WebSocketUpgrade,
) -> axum::response::Response {
let node = self.clone();
upgrade
.max_message_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
.max_frame_size(unb_transport::DEFAULT_MAX_FRAME_SIZE)
.on_upgrade(move |socket| async move {
let (pipe, initiator) = unb_transport::ws::accept(socket);
let _ = node.attach(Pipe::Piped { pipe, initiator }, None).await;
})
}
#[cfg(feature = "hosting")]
pub async fn serve_webtransport(
self: &Arc<Self>,
connection: unb_transport::webtransport::wtransport::Connection,
) -> Result<Arc<Wire>, WsError> {
let (pipe, initiator, bodies) = unb_transport::webtransport::accept(connection).await?;
Ok(self
.attach(Pipe::piped_with_streams(pipe, initiator, bodies), None)
.await
.0)
}
pub async fn serve_transport(self: &Arc<Self>, transport: Pipe) -> Arc<Wire> {
self.attach(transport, None).await.0
}
pub async fn connect_transport(
self: &Arc<Self>,
peer: &str,
transport: Pipe,
) -> Result<Arc<Wire>, WsError> {
let candidate = self.establish(transport, Some(peer.to_string())).await;
match candidate.outcome(peer).await {
Ok(CandidateOutcome::Promoted(_)) => {
let synchronized =
n0_future::time::timeout(ROUTE_SYNC_TIMEOUT, candidate.wire.routes_acked())
.await
.map_err(|_| {
WsError::Connect(format!(
"connected peer {peer:?} did not acknowledge its synchronized routes"
))
})
.and_then(|result| result);
if let Err(error) = synchronized {
candidate.wire.shutdown();
let _ = candidate.cleaned.await;
return Err(error);
}
Ok(candidate.wire)
}
Ok(CandidateOutcome::Duplicate(_)) => Ok(candidate.wire),
Err(error) => {
candidate.wire.shutdown();
let _ = candidate.cleaned.await;
Err(error)
}
}
}
pub(crate) async fn establish(
self: &Arc<Self>,
transport: Pipe,
expected_peer: Option<String>,
) -> CandidateSession {
let (wire, _session, cleaned, identity) = self.attach_inner(transport, expected_peer).await;
CandidateSession {
wire,
cleaned,
identity,
}
}
pub(crate) async fn attach(
self: &Arc<Self>,
transport: Pipe,
expected_peer: Option<String>,
) -> (Arc<Wire>, SessionId, tokio::sync::oneshot::Receiver<()>) {
let session = self.next_session_id();
let (wire, cleaned) = self
.attach_session(session.clone(), transport, expected_peer)
.await;
(wire, session, cleaned)
}
async fn attach_inner(
self: &Arc<Self>,
transport: Pipe,
expected_peer: Option<String>,
) -> (
Arc<Wire>,
SessionId,
tokio::sync::oneshot::Receiver<()>,
tokio::sync::watch::Receiver<Option<unb_core::NodeIdentity>>,
) {
let session = self.next_session_id();
self.outbound_sessions.lock().await.insert(session.clone());
let (identity_tx, identity) = tokio::sync::watch::channel(None);
self.candidate_identities
.lock()
.await
.insert(session.clone(), identity_tx);
let (wire, cleaned) = self
.attach_session(session.clone(), transport, expected_peer)
.await;
(wire, session, cleaned, identity)
}
fn next_session_id(&self) -> SessionId {
SessionId::from(format!(
"sess-{}",
self.next_session
.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
+ 1
))
}
async fn attach_session(
self: &Arc<Self>,
session: SessionId,
transport: Pipe,
expected_peer: Option<String>,
) -> (Arc<Wire>, tokio::sync::oneshot::Receiver<()>) {
let bridge = SessionBridge {
node: Arc::downgrade(self),
session: session.clone(),
};
let wire = self
.protocol
.attach_with_ceiling(
session.clone(),
transport,
expected_peer,
bridge,
self.ws_collect_ceiling,
)
.await
.expect("protocol core actor unavailable");
self.sessions
.write()
.await
.insert(session.to_string(), wire.clone());
let (cleaned_tx, cleaned_rx) = tokio::sync::oneshot::channel();
let cancellation = self.cancellation.child_token();
let closed = wire.clone();
unb_runtime::RuntimeHandle::current().spawn(async move {
tokio::select! {
biased;
() = cancellation.cancelled() => {}
() = closed.closed() => {}
}
let _ = cleaned_tx.send(());
});
(wire, cleaned_rx)
}
async fn cleanup_session(&self, session: &str) {
self.peers
.write()
.await
.retain(|_, link| link.session_id != session);
self.publish_route_change();
self.active.lock().await.retain(|(owner, _), cancel| {
let keep = owner.as_str() != session;
if !keep {
cancel.cancel();
}
keep
});
}
}
fn unary_result(
parts: http::response::Parts,
payload: Bytes,
handle: &ProtocolCoreHandle,
session: &SessionId,
) -> Result<ApplicationResult, ApplicationFailure> {
if parts.status.is_client_error() || parts.status.is_server_error() {
let value: Value = serde_json::from_slice(&payload).unwrap_or(Value::Null);
let code = value
.get("code")
.and_then(|code| serde_json::from_value(code.clone()).ok())
.unwrap_or_else(|| ErrorCode::from_status(parts.status));
let message = value
.get("message")
.and_then(Value::as_str)
.map(str::to_string)
.unwrap_or_else(|| String::from_utf8_lossy(&payload).into_owned());
return Err(ApplicationFailure { code, message });
}
Ok(ApplicationResult::Response(application_response(
parts.status,
&parts.headers,
register_response_body(handle, session, payload)?,
)))
}
fn stream_result(
item: Option<Result<Bytes, crate::handler::HandlerError>>,
parts: &http::response::Parts,
handle: &ProtocolCoreHandle,
session: &SessionId,
) -> Result<ApplicationResult, ApplicationFailure> {
match item {
Some(Ok(payload)) => Ok(ApplicationResult::Event(application_response(
parts.status,
&parts.headers,
register_response_body(handle, session, payload)?,
))),
Some(Err(error)) => Err(ApplicationFailure {
code: error.code,
message: error.message,
}),
None => Ok(ApplicationResult::Finished(application_response(
parts.status,
&parts.headers,
None,
))),
}
}
fn register_response_body(
handle: &ProtocolCoreHandle,
session: &SessionId,
payload: Bytes,
) -> Result<Option<unb_core::BodyId>, ApplicationFailure> {
if payload.is_empty() {
return Ok(None);
}
handle
.register_body(session, unb_runtime::WireBody::Bytes(payload))
.map(Some)
.map_err(|error| ApplicationFailure {
code: ErrorCode::Busy,
message: error.to_string(),
})
}
fn application_response(
status: http::StatusCode,
headers: &http::HeaderMap,
body: Option<unb_core::BodyId>,
) -> ApplicationResponse {
let mut head = http::Response::new(());
*head.status_mut() = status;
*head.headers_mut() = headers.clone();
ApplicationResponse { head, body }
}
fn dispatch_failure(effect: EffectId, code: ErrorCode, message: impl Into<String>) -> CoreInput {
CoreInput::DispatchCompleted {
effect,
result: Err(ApplicationFailure {
code,
message: message.into(),
}),
}
}
pub(crate) fn retirement_error(peer: &str, reason: RetirementReason) -> WsError {
WsError::Connect(format!(
"connection to {peer:?} retired during establishment: {reason:?}"
))
}