#![deny(missing_docs)]
pub(crate) mod address;
pub mod tcp;
pub mod utils;
#[cfg(unix)]
pub mod uds;
#[cfg(feature = "nats-transport")]
pub mod nats;
#[cfg(feature = "grpc")]
pub mod grpc;
#[cfg(feature = "zmq")]
pub mod zmq;
mod transport;
use std::{collections::HashMap, sync::Arc};
use crate::observability::{Direction, TransportRejection, VeloMetrics};
use bytes::Bytes;
use dashmap::DashMap;
use parking_lot::Mutex;
use velo_ext::{InstanceId, PeerInfo, TransportKey, WorkerAddress, WorkerId};
use address::WorkerAddressBuilder;
pub use utils::interfaces::{InterfaceEndpoint, InterfaceFilter};
pub use transport::{
DataStreams, HealthCheckError, InFlightGuard, MessageType, SendBackpressure, SendOutcome,
ShutdownPolicy, ShutdownState, Transport, TransportAdapter, TransportError,
TransportErrorHandler, make_channels, try_send_or_backpressure,
};
#[derive(Debug, thiserror::Error)]
pub enum VeloBackendError {
#[error("No compatible transports found")]
NoCompatibleTransports,
#[error("Transport not found for instance: {0}")]
InstanceNotRegistered(InstanceId),
#[error("Worker not found: {0}")]
WorkerNotRegistered(WorkerId),
#[error("Transport not found: {0}")]
TransportNotFound(TransportKey),
#[error("Invalid transport priority: {0}")]
InvalidTransportPriority(String),
}
pub struct VeloBackend {
instance_id: InstanceId,
address: WorkerAddress,
priorities: Mutex<Vec<TransportKey>>,
transports: HashMap<TransportKey, Arc<dyn Transport>>,
transport_metrics: HashMap<TransportKey, Arc<crate::observability::TransportMetricsHandle>>,
primary_transport: DashMap<InstanceId, Arc<dyn Transport>>,
alternative_transports: DashMap<InstanceId, Vec<TransportKey>>,
workers: DashMap<WorkerId, InstanceId>,
shutdown_state: ShutdownState,
#[allow(dead_code)]
runtime: tokio::runtime::Handle,
}
impl VeloBackend {
pub async fn new(
backend_transports: Vec<Arc<dyn Transport>>,
observability: Option<Arc<VeloMetrics>>,
) -> anyhow::Result<(Self, DataStreams)> {
let instance_id = InstanceId::new_v4();
let mut priorities = Vec::new();
let mut builder = WorkerAddressBuilder::new();
let mut transports = HashMap::new();
let mut transport_metrics = HashMap::new();
let (adapter, data_streams) = transport::make_channels();
let shutdown_state = adapter.shutdown_state.clone();
let runtime = tokio::runtime::Handle::current();
for transport in backend_transports {
let key = transport.key();
if let Some(metrics) = observability.as_ref() {
let handle = Arc::new(metrics.bind_transport(key.as_str()));
transport
.set_observability(handle.clone() as Arc<dyn velo_ext::TransportObservability>);
transport_metrics.insert(key.clone(), handle);
}
transport
.start(instance_id, adapter.clone(), runtime.clone())
.await?;
builder.merge(&transport.address())?;
priorities.push(key.clone());
transports.insert(key, transport);
}
let address = builder.build()?;
Ok((
Self {
instance_id,
address,
transports,
transport_metrics,
priorities: Mutex::new(priorities),
primary_transport: DashMap::new(),
alternative_transports: DashMap::new(),
workers: DashMap::new(),
shutdown_state,
runtime,
},
data_streams,
))
}
pub fn instance_id(&self) -> InstanceId {
self.instance_id
}
pub fn peer_info(&self) -> PeerInfo {
PeerInfo::new(self.instance_id, self.address.clone())
}
pub fn is_registered(&self, instance_id: InstanceId) -> bool {
self.primary_transport.contains_key(&instance_id)
}
pub fn try_translate_worker_id(
&self,
worker_id: WorkerId,
) -> Result<InstanceId, VeloBackendError> {
self.workers
.get(&worker_id)
.map(|entry| *entry)
.ok_or(VeloBackendError::WorkerNotRegistered(worker_id))
}
#[deprecated(since = "0.7.0", note = "Use try_translate_worker_id() instead")]
pub fn translate_worker_id(&self, worker_id: WorkerId) -> Result<InstanceId, VeloBackendError> {
self.try_translate_worker_id(worker_id)
}
pub fn has_instance(&self, instance_id: InstanceId) -> bool {
self.primary_transport.contains_key(&instance_id)
}
pub fn primary_transport_key(&self, target: InstanceId) -> Option<TransportKey> {
self.primary_transport
.get(&target)
.map(|entry| entry.value().key())
}
pub fn alternative_transport_keys(&self, target: InstanceId) -> Option<Vec<TransportKey>> {
self.alternative_transports
.get(&target)
.map(|entry| entry.value().clone())
}
pub fn send_message(
&self,
target: InstanceId,
header: Bytes,
payload: Bytes,
message_type: MessageType,
on_error: Arc<dyn TransportErrorHandler>,
) -> anyhow::Result<SendOutcome> {
let transport = self
.primary_transport
.get(&target)
.ok_or(VeloBackendError::InstanceNotRegistered(target))?;
let transport_key = transport.value().key();
let transport_name = transport_key.to_string();
#[cfg(not(feature = "distributed-tracing"))]
let _ = &transport_name;
let bytes = header.len() + payload.len();
let metrics = self.transport_metrics.get(&transport_key);
let error_handler = instrument_transport_error_handler(metrics.cloned(), on_error);
#[cfg(feature = "distributed-tracing")]
let send_result = {
let span = tracing::info_span!(
"velo.transport.send",
transport = transport_name.as_str(),
message_type = message_type_label(message_type),
bytes
);
let _entered = span.enter();
transport.send_message(target, header, payload, message_type, error_handler)
};
#[cfg(not(feature = "distributed-tracing"))]
let send_result =
transport.send_message(target, header, payload, message_type, error_handler);
Ok(finalize_send_outcome(
send_result,
metrics.cloned(),
message_type,
bytes,
))
}
pub fn send_message_with_transport(
&self,
target: InstanceId,
header: Bytes,
payload: Bytes,
message_type: MessageType,
on_error: Arc<dyn TransportErrorHandler>,
transport_key: TransportKey,
) -> anyhow::Result<SendOutcome> {
let transport = self
.primary_transport
.get(&target)
.ok_or(VeloBackendError::InstanceNotRegistered(target))?;
if transport.value().key() == transport_key {
let _transport_name = transport_key.to_string();
let bytes = header.len() + payload.len();
let metrics = self.transport_metrics.get(&transport_key);
let error_handler = instrument_transport_error_handler(metrics.cloned(), on_error);
let send_result =
transport.send_message(target, header, payload, message_type, error_handler);
return Ok(finalize_send_outcome(
send_result,
metrics.cloned(),
message_type,
bytes,
));
} else {
let alternative_transports = self
.alternative_transports
.get(&target)
.ok_or(VeloBackendError::InstanceNotRegistered(target))?;
for alternative_transport in alternative_transports.iter() {
if *alternative_transport == transport_key
&& let Some(transport) = self.transports.get(alternative_transport)
{
let _transport_name = alternative_transport.to_string();
let bytes = header.len() + payload.len();
let metrics = self.transport_metrics.get(alternative_transport);
let error_handler =
instrument_transport_error_handler(metrics.cloned(), on_error);
let send_result = transport.send_message(
target,
header,
payload,
message_type,
error_handler,
);
return Ok(finalize_send_outcome(
send_result,
metrics.cloned(),
message_type,
bytes,
));
}
}
}
Err(VeloBackendError::NoCompatibleTransports)?
}
pub fn send_message_to_worker(
&self,
worker_id: WorkerId,
header: Bytes,
payload: Bytes,
message_type: MessageType,
on_error: Arc<dyn TransportErrorHandler>,
) -> anyhow::Result<SendOutcome> {
let instance_id = self.try_translate_worker_id(worker_id)?;
self.send_message(instance_id, header, payload, message_type, on_error)
}
pub fn register_peer(&self, peer: PeerInfo) -> Result<(), VeloBackendError> {
let instance_id = peer.instance_id();
let mut compatible_transports = Vec::new();
for (key, transport) in self.transports.iter() {
if transport.register(peer.clone()).is_ok() {
compatible_transports.push(key.clone());
}
}
if compatible_transports.is_empty() {
return Err(VeloBackendError::NoCompatibleTransports);
}
let sorted_transports = self
.priorities
.lock()
.iter()
.filter(|key| compatible_transports.contains(key))
.cloned()
.collect::<Vec<TransportKey>>();
assert!(
!sorted_transports.is_empty(),
"failed to properly sort compatible transports"
);
let primary_transport_key = sorted_transports[0].clone();
let alternative_transport_keys = sorted_transports[1..].to_vec();
let primary_transport = self.transports.get(&primary_transport_key).unwrap();
self.primary_transport
.insert(instance_id, primary_transport.clone());
self.alternative_transports
.insert(instance_id, alternative_transport_keys);
self.workers.insert(instance_id.worker_id(), instance_id);
Ok(())
}
pub fn available_transports(&self) -> Vec<TransportKey> {
self.transports.keys().cloned().collect()
}
pub fn set_transport_priority(
&self,
priorities: Vec<TransportKey>,
) -> Result<(), VeloBackendError> {
let required_transports = self.available_transports();
if required_transports.len() != priorities.len() {
return Err(VeloBackendError::InvalidTransportPriority(format!(
"Required transports: {:?}, provided priorities: {:?}",
required_transports, priorities
)));
}
for priority in &priorities {
if !required_transports.contains(priority) {
return Err(VeloBackendError::InvalidTransportPriority(format!(
"Priority transport not found: {:?}",
priority
)));
}
}
let mut guard = self.priorities.lock();
*guard = priorities;
Ok(())
}
pub fn shutdown_state(&self) -> &ShutdownState {
&self.shutdown_state
}
pub async fn graceful_shutdown(&self, policy: ShutdownPolicy) {
self.shutdown_state.begin_drain();
for transport in self.transports.values() {
transport.begin_drain();
}
match policy {
ShutdownPolicy::WaitForever => {
self.shutdown_state.wait_for_drain().await;
}
ShutdownPolicy::Timeout(duration) => {
let _ = tokio::time::timeout(duration, self.shutdown_state.wait_for_drain()).await;
}
}
self.shutdown_state.teardown_token().cancel();
for transport in self.transports.values() {
transport.shutdown();
}
}
}
pub(crate) fn message_type_label(message_type: MessageType) -> &'static str {
match message_type {
MessageType::Message => "message",
MessageType::Response => "response",
MessageType::Ack => "ack",
MessageType::Event => "event",
MessageType::ShuttingDown => "shutting_down",
}
}
struct InstrumentedTransportErrorHandler {
metrics: Option<Arc<crate::observability::TransportMetricsHandle>>,
inner: Arc<dyn TransportErrorHandler>,
}
impl TransportErrorHandler for InstrumentedTransportErrorHandler {
fn on_error(&self, header: Bytes, payload: Bytes, error: String) {
if let Some(metrics) = self.metrics.as_ref() {
metrics.record_rejection(TransportRejection::SendError);
}
self.inner.on_error(header, payload, error);
}
}
fn instrument_transport_error_handler(
metrics: Option<Arc<crate::observability::TransportMetricsHandle>>,
inner: Arc<dyn TransportErrorHandler>,
) -> Arc<dyn TransportErrorHandler> {
Arc::new(InstrumentedTransportErrorHandler { metrics, inner })
}
fn finalize_send_outcome(
send_result: Result<(), SendBackpressure>,
metrics: Option<Arc<crate::observability::TransportMetricsHandle>>,
message_type: MessageType,
bytes: usize,
) -> SendOutcome {
match send_result {
Ok(()) => {
if let Some(metrics) = metrics {
metrics.record_frame(Direction::Outbound, message_type_label(message_type), bytes);
}
SendOutcome::Enqueued
}
Err(bp) => {
let label = message_type_label(message_type);
SendOutcome::Backpressured(SendBackpressure::new(Box::pin(async move {
bp.await;
if let Some(metrics) = metrics {
metrics.record_frame(Direction::Outbound, label, bytes);
}
})))
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use bytes::Bytes;
use futures::future::BoxFuture;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
struct MockTransport {
key: TransportKey,
address: WorkerAddress,
accept_register: bool,
started: AtomicBool,
drained: AtomicBool,
shut_down: AtomicBool,
send_count: AtomicUsize,
always_backpressure: bool,
}
impl MockTransport {
fn new(key: &str, accept_register: bool) -> Arc<Self> {
let mut builder = WorkerAddressBuilder::new();
builder
.add_entry(key, format!("mock://{}", key).into_bytes())
.unwrap();
let address = builder.build().unwrap();
Arc::new(Self {
key: TransportKey::from(key),
address,
accept_register,
started: AtomicBool::new(false),
drained: AtomicBool::new(false),
shut_down: AtomicBool::new(false),
send_count: AtomicUsize::new(0),
always_backpressure: false,
})
}
fn new_backpressured(key: &str) -> Arc<Self> {
let mut builder = WorkerAddressBuilder::new();
builder
.add_entry(key, format!("mock://{}", key).into_bytes())
.unwrap();
let address = builder.build().unwrap();
Arc::new(Self {
key: TransportKey::from(key),
address,
accept_register: true,
started: AtomicBool::new(false),
drained: AtomicBool::new(false),
shut_down: AtomicBool::new(false),
send_count: AtomicUsize::new(0),
always_backpressure: true,
})
}
}
impl Transport for MockTransport {
fn key(&self) -> TransportKey {
self.key.clone()
}
fn address(&self) -> WorkerAddress {
self.address.clone()
}
fn register(&self, _peer_info: PeerInfo) -> Result<(), TransportError> {
if self.accept_register {
Ok(())
} else {
Err(TransportError::NoEndpoint)
}
}
fn send_message(
&self,
_instance_id: InstanceId,
_header: Bytes,
_payload: Bytes,
_message_type: MessageType,
_on_error: Arc<dyn TransportErrorHandler>,
) -> Result<(), SendBackpressure> {
self.send_count.fetch_add(1, Ordering::Relaxed);
if self.always_backpressure {
Err(SendBackpressure::new(Box::pin(async {})))
} else {
Ok(())
}
}
fn start(
&self,
_instance_id: InstanceId,
_channels: TransportAdapter,
_rt: tokio::runtime::Handle,
) -> BoxFuture<'_, anyhow::Result<()>> {
self.started.store(true, Ordering::Relaxed);
Box::pin(async { Ok(()) })
}
fn shutdown(&self) {
self.shut_down.store(true, Ordering::Relaxed);
}
fn begin_drain(&self) {
self.drained.store(true, Ordering::Relaxed);
}
fn check_health(
&self,
_instance_id: InstanceId,
_timeout: Duration,
) -> std::pin::Pin<
Box<
dyn std::future::Future<Output = Result<(), transport::HealthCheckError>>
+ Send
+ '_,
>,
> {
Box::pin(async { Ok(()) })
}
}
struct NoopErrorHandler;
impl TransportErrorHandler for NoopErrorHandler {
fn on_error(&self, _header: Bytes, _payload: Bytes, _error: String) {}
}
fn make_peer_info(keys: &[&str]) -> PeerInfo {
let instance_id = InstanceId::new_v4();
let mut builder = WorkerAddressBuilder::new();
for key in keys {
builder
.add_entry(*key, format!("mock://{}", key).into_bytes())
.unwrap();
}
let address = builder.build().unwrap();
PeerInfo::new(instance_id, address)
}
#[tokio::test]
async fn test_new_single_transport() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t.clone() as Arc<dyn Transport>], None)
.await
.unwrap();
assert!(t.started.load(Ordering::Relaxed));
assert!(!backend.instance_id().as_bytes().iter().all(|&b| b == 0));
assert_eq!(backend.available_transports().len(), 1);
}
#[tokio::test]
async fn test_new_multiple_transports() {
let t1 = MockTransport::new("tcp", true);
let t2 = MockTransport::new("http", true);
let (backend, _streams) = VeloBackend::new(
vec![
t1.clone() as Arc<dyn Transport>,
t2.clone() as Arc<dyn Transport>,
],
None,
)
.await
.unwrap();
assert!(t1.started.load(Ordering::Relaxed));
assert!(t2.started.load(Ordering::Relaxed));
assert_eq!(backend.available_transports().len(), 2);
}
#[tokio::test]
async fn test_register_peer_selects_primary_by_priority() {
let t1 = MockTransport::new("tcp", true);
let t2 = MockTransport::new("http", true);
let (backend, _streams) = VeloBackend::new(
vec![
t1.clone() as Arc<dyn Transport>,
t2.clone() as Arc<dyn Transport>,
],
None,
)
.await
.unwrap();
let peer = make_peer_info(&["tcp", "http"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
assert!(backend.is_registered(peer_id));
let primary = backend.primary_transport.get(&peer_id).unwrap();
assert_eq!(primary.value().key(), TransportKey::from("tcp"));
}
#[tokio::test]
async fn test_register_peer_no_compatible_transports() {
let t = MockTransport::new("tcp", false);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let result = backend.register_peer(peer);
assert!(matches!(
result,
Err(VeloBackendError::NoCompatibleTransports)
));
}
#[tokio::test]
async fn test_register_peer_stores_worker_mapping() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let peer_id = peer.instance_id();
let worker_id = peer_id.worker_id();
backend.register_peer(peer).unwrap();
let resolved = backend.try_translate_worker_id(worker_id).unwrap();
assert_eq!(resolved, peer_id);
}
#[tokio::test]
async fn test_send_message_routes_to_primary() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t.clone() as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
let outcome = backend
.send_message(
peer_id,
Bytes::from_static(&[1]),
Bytes::from_static(&[2]),
MessageType::Message,
Arc::new(NoopErrorHandler),
)
.unwrap();
assert!(
matches!(outcome, SendOutcome::Enqueued),
"MockTransport returns Ok(()) so backend should report Enqueued"
);
assert_eq!(t.send_count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn test_send_message_backpressured() {
let t = MockTransport::new_backpressured("tcp");
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
let outcome = backend
.send_message(
peer_id,
Bytes::from_static(&[1]),
Bytes::from_static(&[2]),
MessageType::Message,
Arc::new(NoopErrorHandler),
)
.unwrap();
match outcome {
SendOutcome::Backpressured(bp) => {
tokio::time::timeout(Duration::from_secs(1), bp)
.await
.expect("bp should resolve when inner future completes");
}
SendOutcome::Enqueued => {
panic!("backpressured mock should surface SendOutcome::Backpressured")
}
}
}
#[tokio::test]
async fn test_send_message_unregistered_peer() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let result = backend.send_message(
InstanceId::new_v4(),
Bytes::new(),
Bytes::new(),
MessageType::Message,
Arc::new(NoopErrorHandler),
);
assert!(result.is_err());
}
#[tokio::test]
async fn test_send_message_with_transport_primary_match() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t.clone() as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
backend
.send_message_with_transport(
peer_id,
Bytes::from_static(&[1]),
Bytes::from_static(&[2]),
MessageType::Message,
Arc::new(NoopErrorHandler),
TransportKey::from("tcp"),
)
.unwrap();
assert_eq!(t.send_count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn test_send_message_with_transport_alternative() {
let t1 = MockTransport::new("tcp", true);
let t2 = MockTransport::new("http", true);
let (backend, _streams) = VeloBackend::new(
vec![
t1.clone() as Arc<dyn Transport>,
t2.clone() as Arc<dyn Transport>,
],
None,
)
.await
.unwrap();
let peer = make_peer_info(&["tcp", "http"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
backend
.send_message_with_transport(
peer_id,
Bytes::from_static(&[1]),
Bytes::from_static(&[2]),
MessageType::Message,
Arc::new(NoopErrorHandler),
TransportKey::from("http"),
)
.unwrap();
assert_eq!(t2.send_count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn test_send_message_with_transport_not_found() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let peer = make_peer_info(&["tcp"]);
let peer_id = peer.instance_id();
backend.register_peer(peer).unwrap();
let result = backend.send_message_with_transport(
peer_id,
Bytes::new(),
Bytes::new(),
MessageType::Message,
Arc::new(NoopErrorHandler),
TransportKey::from("grpc"),
);
assert!(result.is_err());
}
#[tokio::test]
async fn test_try_translate_worker_id_not_found() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let result = backend.try_translate_worker_id(InstanceId::new_v4().worker_id());
assert!(matches!(
result,
Err(VeloBackendError::WorkerNotRegistered(_))
));
}
#[tokio::test]
async fn test_set_transport_priority_valid() {
let t1 = MockTransport::new("tcp", true);
let t2 = MockTransport::new("http", true);
let (backend, _streams) = VeloBackend::new(
vec![t1 as Arc<dyn Transport>, t2 as Arc<dyn Transport>],
None,
)
.await
.unwrap();
backend
.set_transport_priority(vec![TransportKey::from("http"), TransportKey::from("tcp")])
.unwrap();
}
#[tokio::test]
async fn test_set_transport_priority_wrong_length() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let result = backend
.set_transport_priority(vec![TransportKey::from("tcp"), TransportKey::from("http")]);
assert!(matches!(
result,
Err(VeloBackendError::InvalidTransportPriority(_))
));
}
#[tokio::test]
async fn test_set_transport_priority_unknown_key() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let result = backend.set_transport_priority(vec![TransportKey::from("unknown")]);
assert!(matches!(
result,
Err(VeloBackendError::InvalidTransportPriority(_))
));
}
#[tokio::test]
async fn test_graceful_shutdown_calls_all_transports() {
let t1 = MockTransport::new("tcp", true);
let t2 = MockTransport::new("http", true);
let (backend, _streams) = VeloBackend::new(
vec![
t1.clone() as Arc<dyn Transport>,
t2.clone() as Arc<dyn Transport>,
],
None,
)
.await
.unwrap();
backend
.graceful_shutdown(ShutdownPolicy::Timeout(Duration::from_millis(100)))
.await;
assert!(t1.drained.load(Ordering::Relaxed));
assert!(t2.drained.load(Ordering::Relaxed));
assert!(t1.shut_down.load(Ordering::Relaxed));
assert!(t2.shut_down.load(Ordering::Relaxed));
assert!(backend.shutdown_state().is_draining());
assert!(backend.shutdown_state().teardown_token().is_cancelled());
}
#[tokio::test]
async fn test_peer_info_roundtrip() {
let t = MockTransport::new("tcp", true);
let (backend, _streams) = VeloBackend::new(vec![t as Arc<dyn Transport>], None)
.await
.unwrap();
let info = backend.peer_info();
assert_eq!(info.instance_id(), backend.instance_id());
}
}