use std::{
collections::{HashMap, VecDeque},
io,
sync::{Arc, Mutex},
task::{Context, Poll, Waker},
};
use iroh::{
endpoint::transports::{CustomEndpoint, CustomSender, CustomTransport, RecvInfo, Transmit},
EndpointAddr, EndpointId, TransportAddr,
};
use iroh_base::CustomAddr;
use crate::iroh_carrier::{DEFAULT_PACKET_QUEUE_CAPACITY, MAX_INNER_PACKET_BYTES};
use crate::iroh_carrier_kind::IrohCarrierKind;
#[derive(Debug)]
struct InboundPacket {
remote: CustomAddr,
local: CustomAddr,
bytes: Vec<u8>,
}
#[derive(Debug)]
struct PacketCarrierState {
inbound: VecDeque<InboundPacket>,
inbound_capacity: usize,
recv_waker: Option<Waker>,
outbound: HashMap<CustomAddr, OutboundSession>,
outbound_space_wakers: HashMap<CustomAddr, (u64, Waker)>,
next_session_id: u64,
inbound_space: Arc<tokio::sync::Notify>,
closed: bool,
}
#[derive(Debug, Clone)]
struct OutboundSession {
id: u64,
sender: async_channel::Sender<Vec<u8>>,
}
#[derive(Debug)]
pub struct PacketCarrierSession {
session_id: u64,
remote_addr: CustomAddr,
outbound: async_channel::Receiver<Vec<u8>>,
state: Arc<Mutex<PacketCarrierState>>,
local_addr: CustomAddr,
}
impl PacketCarrierSession {
pub fn remote_addr(&self) -> &CustomAddr {
&self.remote_addr
}
pub async fn recv_outbound(&self) -> Result<Vec<u8>, PacketCarrierQueueError> {
let packet = self
.outbound
.recv()
.await
.map_err(|_| PacketCarrierQueueError::Closed)?;
self.wake_outbound_sender();
Ok(packet)
}
pub fn try_recv_outbound(&self) -> Result<Vec<u8>, PacketCarrierQueueError> {
let packet = self.outbound.try_recv().map_err(|error| match error {
async_channel::TryRecvError::Empty => PacketCarrierQueueError::Empty,
async_channel::TryRecvError::Closed => PacketCarrierQueueError::Closed,
})?;
self.wake_outbound_sender();
Ok(packet)
}
pub fn try_deliver_inbound(&self, bytes: Vec<u8>) -> Result<(), PacketCarrierQueueError> {
if bytes.is_empty() || bytes.len() > MAX_INNER_PACKET_BYTES {
return Err(PacketCarrierQueueError::PacketSize(bytes.len()));
}
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
if state.closed {
return Err(PacketCarrierQueueError::Closed);
}
if state.inbound.len() >= state.inbound_capacity {
return Err(PacketCarrierQueueError::Backpressure);
}
state.inbound.push_back(InboundPacket {
remote: self.remote_addr.clone(),
local: self.local_addr.clone(),
bytes,
});
if let Some(waker) = state.recv_waker.take() {
waker.wake();
}
Ok(())
}
pub async fn deliver_inbound(&self, bytes: Vec<u8>) -> Result<(), PacketCarrierQueueError> {
loop {
let ready = {
let state = self.state.lock().expect("packet carrier mutex poisoned");
state.inbound_space.clone()
};
let notified = ready.notified();
match self.try_deliver_inbound(bytes.clone()) {
Err(PacketCarrierQueueError::Backpressure) => notified.await,
result => return result,
}
}
}
fn wake_outbound_sender(&self) {
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
let waker = state
.outbound_space_wakers
.get(&self.remote_addr)
.filter(|(session_id, _)| *session_id == self.session_id)
.map(|(_, waker)| waker.clone());
if waker.is_some() {
state.outbound_space_wakers.remove(&self.remote_addr);
}
drop(state);
if let Some(waker) = waker {
waker.wake();
}
}
pub fn close(&self) {
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
let is_current = state
.outbound
.get(&self.remote_addr)
.is_some_and(|session| session.id == self.session_id);
if !is_current {
return;
}
state.outbound.remove(&self.remote_addr);
let waker = state
.outbound_space_wakers
.remove(&self.remote_addr)
.map(|(_, waker)| waker);
drop(state);
if let Some(waker) = waker {
waker.wake();
}
}
}
impl Drop for PacketCarrierSession {
fn drop(&mut self) {
self.close();
}
}
#[derive(Debug, Clone)]
pub struct PacketCarrierTransport {
transport_id: u64,
local_addr: CustomAddr,
local_addrs: n0_watcher::Watchable<Vec<CustomAddr>>,
state: Arc<Mutex<PacketCarrierState>>,
queue_capacity: usize,
}
impl PacketCarrierTransport {
pub fn new(transport_id: u64, endpoint_id: EndpointId) -> io::Result<Self> {
Self::with_capacity(transport_id, endpoint_id, DEFAULT_PACKET_QUEUE_CAPACITY)
}
pub fn with_capacity(
transport_id: u64,
endpoint_id: EndpointId,
queue_capacity: usize,
) -> io::Result<Self> {
if transport_id == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"custom transport id must be non-zero",
));
}
if queue_capacity == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"packet carrier queue capacity must be non-zero",
));
}
let local_addr = CustomAddr::from_parts(transport_id, endpoint_id.as_bytes());
Ok(Self {
transport_id,
local_addrs: n0_watcher::Watchable::new(vec![local_addr.clone()]),
local_addr,
state: Arc::new(Mutex::new(PacketCarrierState {
inbound: VecDeque::new(),
inbound_capacity: queue_capacity,
recv_waker: None,
outbound: HashMap::new(),
outbound_space_wakers: HashMap::new(),
next_session_id: 0,
inbound_space: Arc::new(tokio::sync::Notify::new()),
closed: false,
})),
queue_capacity,
})
}
pub fn transport_id(&self) -> u64 {
self.transport_id
}
pub fn local_addr(&self) -> &CustomAddr {
&self.local_addr
}
pub fn session_addr(&self, carrier_session_id: [u8; 16]) -> CustomAddr {
CustomAddr::from_parts(self.transport_id, &carrier_session_id)
}
pub fn peer_addr(&self, endpoint_id: EndpointId) -> CustomAddr {
CustomAddr::from_parts(self.transport_id, endpoint_id.as_bytes())
}
pub fn attach_session(
&self,
remote_addr: CustomAddr,
) -> Result<PacketCarrierSession, PacketCarrierQueueError> {
if remote_addr.id() != self.transport_id {
return Err(PacketCarrierQueueError::WrongTransport {
expected: self.transport_id,
actual: remote_addr.id(),
});
}
let (sender, receiver) = async_channel::bounded(self.queue_capacity);
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
if state.closed {
return Err(PacketCarrierQueueError::Closed);
}
if state.outbound.contains_key(&remote_addr) {
return Err(PacketCarrierQueueError::DuplicateSession);
}
state.next_session_id = state.next_session_id.wrapping_add(1).max(1);
let session_id = state.next_session_id;
state.outbound.insert(
remote_addr.clone(),
OutboundSession {
id: session_id,
sender,
},
);
drop(state);
Ok(PacketCarrierSession {
session_id,
remote_addr,
outbound: receiver,
state: self.state.clone(),
local_addr: self.local_addr.clone(),
})
}
pub fn has_session(&self, remote_addr: &CustomAddr) -> bool {
self.state
.lock()
.expect("packet carrier mutex poisoned")
.outbound
.contains_key(remote_addr)
}
pub fn close(&self) {
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
state.closed = true;
state.outbound.clear();
for (_, (_, waker)) in state.outbound_space_wakers.drain() {
waker.wake();
}
state.inbound.clear();
state.inbound_space.notify_waiters();
if let Some(waker) = state.recv_waker.take() {
waker.wake();
}
}
}
#[derive(Debug)]
pub struct PacketCarrierAddressProvider {
carrier_kind: IrohCarrierKind,
transport: Arc<PacketCarrierTransport>,
active_peers: Mutex<std::collections::HashSet<EndpointId>>,
}
impl PacketCarrierAddressProvider {
pub fn new(carrier_kind: IrohCarrierKind, transport: Arc<PacketCarrierTransport>) -> Self {
Self {
carrier_kind,
transport,
active_peers: Mutex::new(std::collections::HashSet::new()),
}
}
pub fn carrier_kind(&self) -> IrohCarrierKind {
self.carrier_kind
}
pub fn activate(
&self,
endpoint_id: EndpointId,
) -> Result<PacketCarrierSession, PacketCarrierQueueError> {
let remote_addr = self.transport.peer_addr(endpoint_id);
let session = self.transport.attach_session(remote_addr)?;
self.active_peers
.lock()
.expect("packet carrier provider mutex poisoned")
.insert(endpoint_id);
Ok(session)
}
pub fn deactivate(&self, endpoint_id: EndpointId) {
self.active_peers
.lock()
.expect("packet carrier provider mutex poisoned")
.remove(&endpoint_id);
}
pub fn prepared_endpoint_addr(&self, endpoint_id: EndpointId) -> anyhow::Result<EndpointAddr> {
anyhow::ensure!(
self.active_peers
.lock()
.expect("packet carrier provider mutex poisoned")
.contains(&endpoint_id),
"packet carrier session is not active for {endpoint_id}"
);
Ok(EndpointAddr::new(endpoint_id)
.with_addrs([TransportAddr::Custom(self.transport.peer_addr(endpoint_id))]))
}
}
#[cfg(not(target_arch = "wasm32"))]
#[async_trait::async_trait]
impl crate::client::UpgradeProvider for PacketCarrierAddressProvider {
fn kind(&self) -> crate::client::IrohPathKind {
match self.carrier_kind {
IrohCarrierKind::WebRtc => crate::client::IrohPathKind::IrohWebRtc,
IrohCarrierKind::Moq => crate::client::IrohPathKind::IrohMoq,
}
}
fn transport_id(&self) -> u64 {
self.transport.transport_id()
}
async fn prepare_endpoint_addr(&self, endpoint_id: EndpointId) -> anyhow::Result<EndpointAddr> {
self.prepared_endpoint_addr(endpoint_id)
}
}
impl CustomTransport for PacketCarrierTransport {
fn bind(&self) -> io::Result<Box<dyn CustomEndpoint>> {
Ok(Box::new(PacketCarrierEndpoint {
local_addrs: self.local_addrs.clone(),
state: self.state.clone(),
}))
}
}
#[derive(Debug)]
struct PacketCarrierEndpoint {
local_addrs: n0_watcher::Watchable<Vec<CustomAddr>>,
state: Arc<Mutex<PacketCarrierState>>,
}
impl CustomEndpoint for PacketCarrierEndpoint {
fn watch_local_addrs(&self) -> n0_watcher::Direct<Vec<CustomAddr>> {
self.local_addrs.watch()
}
fn create_sender(&self) -> Arc<dyn CustomSender> {
Arc::new(PacketCarrierSender {
state: self.state.clone(),
})
}
fn poll_recv(
&mut self,
cx: &mut Context,
bufs: &mut [io::IoSliceMut<'_>],
metas: &mut [noq_udp::RecvMeta],
recv_infos: &mut [RecvInfo],
) -> Poll<io::Result<usize>> {
debug_assert_eq!(bufs.len(), metas.len());
debug_assert_eq!(bufs.len(), recv_infos.len());
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
let mut count = 0;
while count < bufs.len() {
let Some(packet) = state.inbound.pop_front() else {
break;
};
if packet.bytes.len() > bufs[count].len() {
state.inbound.push_front(packet);
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::InvalidData,
"Iroh receive buffer is smaller than carrier packet",
)));
}
bufs[count][..packet.bytes.len()].copy_from_slice(&packet.bytes);
let mut meta = noq_udp::RecvMeta::default();
meta.len = packet.bytes.len();
meta.stride = packet.bytes.len();
metas[count] = meta;
recv_infos[count] = RecvInfo::new(packet.remote, Some(packet.local));
count += 1;
}
if count > 0 {
for _ in 0..count {
state.inbound_space.notify_one();
}
Poll::Ready(Ok(count))
} else if state.closed {
Poll::Ready(Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"packet carrier closed",
)))
} else {
state.recv_waker = Some(cx.waker().clone());
Poll::Pending
}
}
}
#[derive(Debug)]
struct PacketCarrierSender {
state: Arc<Mutex<PacketCarrierState>>,
}
impl PacketCarrierSender {
fn poll_send_packet(
&self,
cx: &mut Context<'_>,
dst: &CustomAddr,
contents: &[u8],
) -> Poll<io::Result<()>> {
let mut state = self.state.lock().expect("packet carrier mutex poisoned");
if state.closed {
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"packet carrier closed",
)));
}
let Some(session) = state.outbound.get(dst).cloned() else {
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::NotConnected,
"carrier session is not attached",
)));
};
match session.sender.try_send(contents.to_vec()) {
Ok(()) => Poll::Ready(Ok(())),
Err(async_channel::TrySendError::Full(_)) => {
state
.outbound_space_wakers
.insert(dst.clone(), (session.id, cx.waker().clone()));
Poll::Pending
}
Err(async_channel::TrySendError::Closed(_)) => Poll::Ready(Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"carrier session closed",
))),
}
}
}
impl CustomSender for PacketCarrierSender {
fn is_valid_send_addr(&self, addr: &CustomAddr) -> bool {
self.state
.lock()
.map(|state| state.outbound.contains_key(addr))
.unwrap_or(false)
}
fn poll_send(
&self,
cx: &mut Context,
dst: &CustomAddr,
_src: Option<&CustomAddr>,
transmit: &Transmit<'_>,
) -> Poll<io::Result<()>> {
if transmit.segment_size.is_some() {
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::InvalidInput,
"packet carrier does not accept GSO batches",
)));
}
if transmit.contents.is_empty() || transmit.contents.len() > MAX_INNER_PACKET_BYTES {
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::InvalidData,
"invalid inner Iroh packet size",
)));
}
self.poll_send_packet(cx, dst, transmit.contents)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PacketCarrierQueueError {
Empty,
Closed,
Backpressure,
PacketSize(usize),
WrongTransport { expected: u64, actual: u64 },
DuplicateSession,
}
impl std::fmt::Display for PacketCarrierQueueError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(formatter, "{self:?}")
}
}
impl std::error::Error for PacketCarrierQueueError {}
#[cfg(test)]
mod tests {
use super::*;
use iroh::endpoint::transports::CustomTransport;
use iroh::Watcher;
#[derive(Debug, Default)]
struct WakeFlag(std::sync::atomic::AtomicBool);
impl std::task::Wake for WakeFlag {
fn wake(self: Arc<Self>) {
self.0.store(true, std::sync::atomic::Ordering::SeqCst);
}
fn wake_by_ref(self: &Arc<Self>) {
self.0.store(true, std::sync::atomic::Ordering::SeqCst);
}
}
#[test]
fn attach_is_transport_scoped_and_duplicate_safe() {
let transport =
PacketCarrierTransport::new(0x57_52_54_43, iroh::SecretKey::generate().public())
.unwrap();
let remote = transport.session_addr([9; 16]);
let session = transport.attach_session(remote.clone()).unwrap();
assert_eq!(session.remote_addr(), &remote);
assert_eq!(
transport.attach_session(remote).unwrap_err(),
PacketCarrierQueueError::DuplicateSession
);
assert_eq!(
session.try_recv_outbound(),
Err(PacketCarrierQueueError::Empty),
"a rejected duplicate must not replace or close the live session"
);
}
#[test]
fn explicit_retirement_allows_replacement_and_stale_drop_cannot_remove_it() {
let transport =
PacketCarrierTransport::new(0x4d_4f_51, iroh::SecretKey::generate().public()).unwrap();
let remote = transport.session_addr([7; 16]);
let stale = transport.attach_session(remote.clone()).unwrap();
stale.close();
let replacement = transport.attach_session(remote.clone()).unwrap();
drop(stale);
let sender = PacketCarrierSender {
state: transport.state.clone(),
};
let waker = futures::task::noop_waker();
let mut context = Context::from_waker(&waker);
assert!(matches!(
sender.poll_send_packet(&mut context, &remote, b"replacement"),
Poll::Ready(Ok(()))
));
assert_eq!(replacement.try_recv_outbound().unwrap(), b"replacement");
}
#[test]
fn inbound_queue_applies_bounded_backpressure() {
let transport = PacketCarrierTransport::with_capacity(
0x4d_4f_51,
iroh::SecretKey::generate().public(),
1,
)
.unwrap();
let session = transport
.attach_session(transport.session_addr([3; 16]))
.unwrap();
session.try_deliver_inbound(vec![1]).unwrap();
assert_eq!(
session.try_deliver_inbound(vec![2]),
Err(PacketCarrierQueueError::Backpressure)
);
}
#[test]
fn outbound_queue_wakes_iroh_after_bounded_backpressure() {
let transport = PacketCarrierTransport::with_capacity(
0x4d_4f_51,
iroh::SecretKey::generate().public(),
1,
)
.unwrap();
let remote = transport.session_addr([5; 16]);
let session = transport.attach_session(remote.clone()).unwrap();
let sender = PacketCarrierSender {
state: transport.state.clone(),
};
let wake_flag = Arc::new(WakeFlag::default());
let waker = std::task::Waker::from(wake_flag.clone());
let mut context = Context::from_waker(&waker);
assert!(matches!(
sender.poll_send_packet(&mut context, &remote, &[1]),
Poll::Ready(Ok(()))
));
assert!(matches!(
sender.poll_send_packet(&mut context, &remote, &[2]),
Poll::Pending
));
assert!(!wake_flag.0.load(std::sync::atomic::Ordering::SeqCst));
assert_eq!(session.try_recv_outbound().unwrap(), vec![1]);
assert!(wake_flag.0.load(std::sync::atomic::Ordering::SeqCst));
assert!(matches!(
sender.poll_send_packet(&mut context, &remote, &[2]),
Poll::Ready(Ok(()))
));
}
#[tokio::test]
async fn async_inbound_delivery_waits_for_bounded_queue_space() {
let transport = PacketCarrierTransport::with_capacity(
0x4d_4f_51,
iroh::SecretKey::generate().public(),
1,
)
.unwrap();
let session = Arc::new(
transport
.attach_session(transport.session_addr([4; 16]))
.unwrap(),
);
session.try_deliver_inbound(vec![1]).unwrap();
let pending = {
let session = session.clone();
tokio::spawn(async move { session.deliver_inbound(vec![2]).await })
};
tokio::task::yield_now().await;
assert!(!pending.is_finished());
let mut endpoint = transport.bind().unwrap();
let mut bytes = [0_u8; MAX_INNER_PACKET_BYTES];
let mut buffers = [io::IoSliceMut::new(&mut bytes)];
let mut metas = [noq_udp::RecvMeta::default()];
let mut infos = [RecvInfo::default()];
let waker = futures::task::noop_waker();
let mut context = Context::from_waker(&waker);
assert!(matches!(
endpoint.poll_recv(&mut context, &mut buffers, &mut metas, &mut infos),
Poll::Ready(Ok(1))
));
pending.await.unwrap().unwrap();
}
#[test]
fn custom_transport_binds_without_owning_an_iroh_endpoint() {
let transport =
PacketCarrierTransport::new(0x57_52_54_43, iroh::SecretKey::generate().public())
.unwrap();
let endpoint = transport.bind().unwrap();
assert_eq!(
endpoint.watch_local_addrs().get(),
vec![transport.local_addr().clone()]
);
}
#[tokio::test]
async fn prepared_address_is_custom_only_and_requires_active_session() {
let transport_id = 0x57_52_54_43;
let local = iroh::SecretKey::generate().public();
let remote = iroh::SecretKey::generate().public();
let transport = Arc::new(PacketCarrierTransport::new(transport_id, local).unwrap());
let provider = PacketCarrierAddressProvider::new(IrohCarrierKind::WebRtc, transport);
assert!(provider.prepared_endpoint_addr(remote).is_err());
let _session = provider.activate(remote).unwrap();
let addr = provider.prepared_endpoint_addr(remote).unwrap();
assert_eq!(addr.addrs.len(), 1);
assert!(addr.addrs.iter().all(TransportAddr::is_custom));
}
#[tokio::test]
async fn normal_iroh_alpn_and_bi_stream_cross_packet_carrier() {
use iroh::{endpoint::presets, Endpoint, RelayMode};
const TEST_ALPN: &[u8] = b"openrtc/test/packet-carrier/1";
const TEST_TRANSPORT_ID: u64 = 0x4f_52_54_54;
let secret_a = iroh::SecretKey::generate();
let secret_b = iroh::SecretKey::generate();
let endpoint_id_a = secret_a.public();
let endpoint_id_b = secret_b.public();
let transport_a =
Arc::new(PacketCarrierTransport::new(TEST_TRANSPORT_ID, endpoint_id_a).unwrap());
let transport_b =
Arc::new(PacketCarrierTransport::new(TEST_TRANSPORT_ID, endpoint_id_b).unwrap());
let session_a = Arc::new(
transport_a
.attach_session(transport_a.peer_addr(endpoint_id_b))
.unwrap(),
);
let session_b = Arc::new(
transport_b
.attach_session(transport_b.peer_addr(endpoint_id_a))
.unwrap(),
);
let endpoint_a = Endpoint::builder(presets::N0)
.secret_key(secret_a)
.relay_mode(RelayMode::Disabled)
.clear_ip_transports()
.alpns(vec![TEST_ALPN.to_vec()])
.add_custom_transport(transport_a)
.bind()
.await
.unwrap();
let endpoint_b = Endpoint::builder(presets::N0)
.secret_key(secret_b)
.relay_mode(RelayMode::Disabled)
.clear_ip_transports()
.alpns(vec![TEST_ALPN.to_vec()])
.add_custom_transport(transport_b)
.bind()
.await
.unwrap();
let pump_a = {
let source = session_a.clone();
let destination = session_b.clone();
tokio::spawn(async move {
while let Ok(packet) = source.recv_outbound().await {
destination.try_deliver_inbound(packet).unwrap();
}
})
};
let pump_b = {
let source = session_b.clone();
let destination = session_a.clone();
tokio::spawn(async move {
while let Ok(packet) = source.recv_outbound().await {
destination.try_deliver_inbound(packet).unwrap();
}
})
};
let server = {
let endpoint_b = endpoint_b.clone();
tokio::spawn(async move {
let incoming = endpoint_b.accept().await.expect("incoming connection");
let connection = incoming.await.expect("accepted packet-carrier connection");
assert_eq!(connection.alpn(), TEST_ALPN);
let (mut send, mut recv) = connection.accept_bi().await.unwrap();
let request = recv.read_to_end(1024).await.unwrap();
assert_eq!(request, b"normal-iroh-stream");
send.write_all(b"protected-response").await.unwrap();
send.finish().unwrap();
let _ = connection.closed().await;
})
};
let address_b = EndpointAddr::new(endpoint_id_b)
.with_addrs([TransportAddr::Custom(session_a.remote_addr().clone())]);
let connection = tokio::time::timeout(
std::time::Duration::from_secs(10),
endpoint_a.connect(address_b, TEST_ALPN),
)
.await
.expect("packet-carrier connect timed out")
.expect("packet-carrier connect failed");
let selected_custom = connection.paths().iter().any(|path| {
path.is_selected()
&& matches!(
path.remote_addr(),
TransportAddr::Custom(address) if address.id() == TEST_TRANSPORT_ID
)
});
assert!(selected_custom, "custom carrier path was not selected");
let (mut send, mut recv) = connection.open_bi().await.unwrap();
send.write_all(b"normal-iroh-stream").await.unwrap();
send.finish().unwrap();
let response = recv.read_to_end(1024).await.unwrap();
assert_eq!(response, b"protected-response");
connection.close(0u8.into(), b"test-complete");
server.await.unwrap();
endpoint_a.close().await;
endpoint_b.close().await;
pump_a.abort();
pump_b.abort();
}
}