use crate::{
constants::{
INBOUND_RPC_TIMEOUT_MS, MAX_CONCURRENT_INBOUND_RPCS, MAX_CONCURRENT_OUTBOUND_RPCS,
MAX_FRAME_SIZE, NETWORK_CHANNEL_SIZE,
},
peer::{DisconnectReason, Peer, PeerNotification, PeerRequest},
peer_manager::TransportNotification,
protocols::{
direct_send::Message,
rpc::{error::RpcError, InboundRpcRequest, OutboundRpcRequest},
wire::{
handshake::v1::{MessagingProtocolVersion, ProtocolIdSet},
messaging::v1::{
DirectSendMsg, NetworkMessage, NetworkMessageSink, NetworkMessageStream,
RpcRequest, RpcResponse,
},
},
},
transport::{Connection, ConnectionId, ConnectionMetadata},
ProtocolId,
};
use aptos_config::{config::PeerRole, network_id::NetworkContext};
use aptos_time_service::{MockTimeService, TimeService};
use aptos_types::{network_address::NetworkAddress, PeerId};
use bytes::Bytes;
use channel::{self, aptos_channel, message_queues::QueueStyle};
use futures::{
channel::oneshot,
future::{self, FutureExt},
io::{AsyncRead, AsyncWrite, AsyncWriteExt},
stream::{StreamExt, TryStreamExt},
SinkExt,
};
use memsocket::MemorySocket;
use netcore::transport::ConnectionOrigin;
use std::{collections::HashSet, str::FromStr, time::Duration};
use tokio::runtime::{Handle, Runtime};
use tokio_util::compat::{
FuturesAsyncReadCompatExt, TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt,
};
static PROTOCOL: ProtocolId = ProtocolId::MempoolDirectSend;
fn build_test_peer(
executor: Handle,
time_service: TimeService,
origin: ConnectionOrigin,
) -> (
Peer<MemorySocket>,
PeerHandle,
MemorySocket,
channel::Receiver<TransportNotification<MemorySocket>>,
aptos_channel::Receiver<ProtocolId, PeerNotification>,
) {
let (a, b) = MemorySocket::new_pair();
let peer_id = PeerId::random();
let connection = Connection {
metadata: ConnectionMetadata::new(
peer_id,
ConnectionId::default(),
NetworkAddress::from_str("/ip4/127.0.0.1/tcp/8081").unwrap(),
origin,
MessagingProtocolVersion::V1,
ProtocolIdSet::empty(),
PeerRole::Unknown,
),
socket: a,
};
let (connection_notifs_tx, connection_notifs_rx) = channel::new_test(1);
let (peer_reqs_tx, peer_reqs_rx) =
aptos_channel::new(QueueStyle::FIFO, NETWORK_CHANNEL_SIZE, None);
let (peer_notifs_tx, peer_notifs_rx) =
aptos_channel::new(QueueStyle::FIFO, NETWORK_CHANNEL_SIZE, None);
let peer = Peer::new(
NetworkContext::mock(),
executor,
time_service,
connection,
connection_notifs_tx,
peer_reqs_rx,
peer_notifs_tx,
Duration::from_millis(INBOUND_RPC_TIMEOUT_MS),
MAX_CONCURRENT_INBOUND_RPCS,
MAX_CONCURRENT_OUTBOUND_RPCS,
MAX_FRAME_SIZE,
None,
None,
);
let peer_handle = PeerHandle(peer_reqs_tx);
(peer, peer_handle, b, connection_notifs_rx, peer_notifs_rx)
}
fn build_test_connected_peers(
executor: Handle,
time_service: TimeService,
) -> (
(
Peer<MemorySocket>,
PeerHandle,
channel::Receiver<TransportNotification<MemorySocket>>,
aptos_channel::Receiver<ProtocolId, PeerNotification>,
),
(
Peer<MemorySocket>,
PeerHandle,
channel::Receiver<TransportNotification<MemorySocket>>,
aptos_channel::Receiver<ProtocolId, PeerNotification>,
),
) {
let (peer_a, peer_handle_a, connection_a, connection_notifs_rx_a, peer_notifs_rx_a) =
build_test_peer(
executor.clone(),
time_service.clone(),
ConnectionOrigin::Inbound,
);
let (mut peer_b, peer_handle_b, _connection_b, connection_notifs_rx_b, peer_notifs_rx_b) =
build_test_peer(executor, time_service, ConnectionOrigin::Outbound);
peer_b.connection = Some(connection_a);
(
(
peer_a,
peer_handle_a,
connection_notifs_rx_a,
peer_notifs_rx_a,
),
(
peer_b,
peer_handle_b,
connection_notifs_rx_b,
peer_notifs_rx_b,
),
)
}
fn build_network_sink_stream(
connection: &mut MemorySocket,
) -> (
NetworkMessageSink<impl AsyncWrite + '_>,
NetworkMessageStream<impl AsyncRead + '_>,
) {
let (read_half, write_half) = tokio::io::split(connection.compat());
let sink = NetworkMessageSink::new(write_half.compat_write(), MAX_FRAME_SIZE, None);
let stream = NetworkMessageStream::new(read_half.compat(), MAX_FRAME_SIZE, None);
(sink, stream)
}
async fn assert_disconnected_event(
peer_id: PeerId,
reason: DisconnectReason,
connection_notifs_rx: &mut channel::Receiver<TransportNotification<MemorySocket>>,
) {
match connection_notifs_rx.next().await {
Some(TransportNotification::Disconnected(metadata, actual_reason)) => {
assert_eq!(metadata.remote_peer_id, peer_id);
assert_eq!(actual_reason, reason);
}
event => panic!("Expected a Disconnected, received: {:?}", event),
}
}
#[derive(Clone)]
struct PeerHandle(aptos_channel::Sender<ProtocolId, PeerRequest>);
impl PeerHandle {
fn send_direct_send(&mut self, message: Message) {
self.0
.push(message.protocol_id, PeerRequest::SendDirectSend(message))
.unwrap()
}
async fn send_rpc_request(
&mut self,
protocol_id: ProtocolId,
data: Bytes,
timeout: Duration,
) -> Result<Bytes, RpcError> {
let (res_tx, res_rx) = oneshot::channel();
let request = OutboundRpcRequest {
protocol_id,
data,
res_tx,
timeout,
};
self.0.push(protocol_id, PeerRequest::SendRpc(request))?;
let response_data = res_rx.await??;
Ok(response_data)
}
}
#[test]
fn peer_send_message() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, mut peer_handle, mut connection, _connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut client_sink, mut client_stream) = build_network_sink_stream(&mut connection);
let send_msg = Message {
protocol_id: PROTOCOL,
mdata: Bytes::from("hello world"),
};
let recv_msg = NetworkMessage::DirectSendMsg(DirectSendMsg {
protocol_id: PROTOCOL,
priority: 0,
raw_msg: Vec::from("hello world"),
});
let client = async {
for _ in 0..30 {
let msg = client_stream.next().await.unwrap().unwrap();
assert_eq!(msg, recv_msg);
}
client_sink.close().await.unwrap();
};
let server = async {
for _ in 0..30 {
peer_handle.send_direct_send(send_msg.clone());
}
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peer_recv_message() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, _peer_handle, connection, _connection_notifs_rx, mut peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let send_msg = NetworkMessage::DirectSendMsg(DirectSendMsg {
protocol_id: PROTOCOL,
priority: 0,
raw_msg: Vec::from("hello world"),
});
let recv_msg = PeerNotification::RecvMessage(Message {
protocol_id: PROTOCOL,
mdata: Bytes::from("hello world"),
});
let client = async move {
let mut connection = NetworkMessageSink::new(connection, MAX_FRAME_SIZE, None);
for _ in 0..30 {
connection.send(&send_msg).await.unwrap();
}
connection.close().await.unwrap();
};
let server = async move {
for _ in 0..30 {
let received = peer_notifs_rx.next().await.unwrap();
assert_eq!(recv_msg, received);
}
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peers_send_message_concurrent() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (
(peer_a, mut peer_handle_a, mut connection_notifs_rx_a, mut peer_notifs_rx_a),
(peer_b, mut peer_handle_b, mut connection_notifs_rx_b, mut peer_notifs_rx_b),
) = build_test_connected_peers(rt.handle().clone(), TimeService::mock());
let remote_peer_id_a = peer_a.remote_peer_id();
let remote_peer_id_b = peer_b.remote_peer_id();
let test = async move {
let msg_a = Message {
protocol_id: PROTOCOL,
mdata: Bytes::from("hello world"),
};
let msg_b = Message {
protocol_id: PROTOCOL,
mdata: Bytes::from("namaste"),
};
peer_handle_a.send_direct_send(msg_a.clone());
peer_handle_b.send_direct_send(msg_b.clone());
let notif_a = peer_notifs_rx_a.next().await;
let notif_b = peer_notifs_rx_b.next().await;
assert_eq!(notif_a, Some(PeerNotification::RecvMessage(msg_b)));
assert_eq!(notif_b, Some(PeerNotification::RecvMessage(msg_a)));
drop(peer_handle_a);
assert_disconnected_event(
remote_peer_id_a,
DisconnectReason::Requested,
&mut connection_notifs_rx_a,
)
.await;
assert_disconnected_event(
remote_peer_id_b,
DisconnectReason::ConnectionLost,
&mut connection_notifs_rx_b,
)
.await;
};
rt.block_on(future::join3(peer_a.start(), peer_b.start(), test));
}
#[test]
fn peer_recv_rpc() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, _peer_handle, mut connection, _connection_notifs_rx, mut peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut client_sink, mut client_stream) = build_network_sink_stream(&mut connection);
let send_msg = NetworkMessage::RpcRequest(RpcRequest {
request_id: 123,
protocol_id: PROTOCOL,
priority: 0,
raw_request: Vec::from("hello world"),
});
let recv_msg = PeerNotification::RecvRpc(InboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from("hello world"),
res_tx: oneshot::channel().0,
});
let resp_msg = NetworkMessage::RpcResponse(RpcResponse {
request_id: 123,
priority: 0,
raw_response: Vec::from("goodbye world"),
});
let client = async move {
for _ in 0..30 {
client_sink.send(&send_msg).await.unwrap();
let received = client_stream.next().await.unwrap().unwrap();
assert_eq!(received, resp_msg);
}
client_sink.close().await.unwrap();
};
let server = async move {
for _ in 0..30 {
let received = peer_notifs_rx.next().await.unwrap();
assert_eq!(recv_msg, received);
match received {
PeerNotification::RecvRpc(req) => {
let response = Ok(Bytes::from("goodbye world"));
req.res_tx.send(response).unwrap()
}
_ => panic!("Unexpected PeerNotification: {:?}", received),
}
}
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peer_recv_rpc_concurrent() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, _peer_handle, mut connection, _connection_notifs_rx, mut peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut client_sink, mut client_stream) = build_network_sink_stream(&mut connection);
let send_msg = NetworkMessage::RpcRequest(RpcRequest {
request_id: 123,
protocol_id: PROTOCOL,
priority: 0,
raw_request: Vec::from("hello world"),
});
let recv_msg = PeerNotification::RecvRpc(InboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from("hello world"),
res_tx: oneshot::channel().0,
});
let resp_msg = NetworkMessage::RpcResponse(RpcResponse {
request_id: 123,
priority: 0,
raw_response: Vec::from("goodbye world"),
});
let client = async move {
for _ in 0..30 {
client_sink.send(&send_msg).await.unwrap();
}
for _ in 0..30 {
let received = client_stream.next().await.unwrap().unwrap();
assert_eq!(received, resp_msg);
}
client_sink.close().await.unwrap();
};
let server = async move {
let mut res_txs = vec![];
for _ in 0..30 {
let received = peer_notifs_rx.next().await.unwrap();
assert_eq!(recv_msg, received);
match received {
PeerNotification::RecvRpc(req) => res_txs.push(req.res_tx),
_ => panic!("Unexpected PeerNotification: {:?}", received),
};
}
for res_tx in res_txs.into_iter() {
let response = Bytes::from("goodbye world");
res_tx.send(Ok(response)).unwrap();
}
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peer_recv_rpc_timeout() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let mock_time = MockTimeService::new();
let (peer, _peer_handle, mut connection, _connection_notifs_rx, mut peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
mock_time.clone().into(),
ConnectionOrigin::Inbound,
);
let (mut client_sink, client_stream) = build_network_sink_stream(&mut connection);
let send_msg = NetworkMessage::RpcRequest(RpcRequest {
request_id: 123,
protocol_id: PROTOCOL,
priority: 0,
raw_request: Vec::from("hello world"),
});
let recv_msg = PeerNotification::RecvRpc(InboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from("hello world"),
res_tx: oneshot::channel().0,
});
let test = async move {
client_sink.send(&send_msg).await.unwrap();
let received = peer_notifs_rx.next().await.unwrap();
assert_eq!(received, recv_msg);
let mut res_tx = match received {
PeerNotification::RecvRpc(req) => req.res_tx,
_ => panic!("Unexpected PeerNotification: {:?}", received),
};
assert!(!res_tx.is_canceled());
mock_time.advance_ms_async(INBOUND_RPC_TIMEOUT_MS).await;
assert!(res_tx.is_canceled());
res_tx.cancellation().await;
client_sink.close().await.unwrap();
let messages = client_stream.try_collect::<Vec<_>>().await.unwrap();
assert_eq!(messages, vec![]);
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_recv_rpc_cancel() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, _peer_handle, mut connection, _connection_notifs_rx, mut peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut client_sink, client_stream) = build_network_sink_stream(&mut connection);
let send_msg = NetworkMessage::RpcRequest(RpcRequest {
request_id: 123,
protocol_id: PROTOCOL,
priority: 0,
raw_request: Vec::from("hello world"),
});
let recv_msg = PeerNotification::RecvRpc(InboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from("hello world"),
res_tx: oneshot::channel().0,
});
let test = async move {
client_sink.send(&send_msg).await.unwrap();
let received = peer_notifs_rx.next().await.unwrap();
assert_eq!(received, recv_msg);
let res_tx = match received {
PeerNotification::RecvRpc(req) => req.res_tx,
_ => panic!("Unexpected PeerNotification: {:?}", received),
};
assert!(!res_tx.is_canceled());
drop(res_tx);
client_sink.close().await.unwrap();
let messages = client_stream.try_collect::<Vec<_>>().await.unwrap();
assert_eq!(messages, vec![]);
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_send_rpc() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, mut peer_handle, mut connection, _connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut server_sink, mut server_stream) = build_network_sink_stream(&mut connection);
let timeout = Duration::from_millis(10_000);
let mut request_ids = HashSet::new();
let client = async move {
for _ in 0..30 {
let response = peer_handle
.send_rpc_request(PROTOCOL, Bytes::from(&b"hello world"[..]), timeout)
.await
.unwrap();
assert_eq!(response, Bytes::from(&b"goodbye world"[..]));
}
};
let server = async move {
for _ in 0..30 {
let received = server_stream.next().await.unwrap().unwrap();
let received = match received {
NetworkMessage::RpcRequest(request) => request,
_ => panic!("Expected RpcRequest; unexpected: {:?}", received),
};
assert_eq!(received.protocol_id, PROTOCOL);
assert_eq!(received.priority, 0);
assert_eq!(received.raw_request, b"hello world");
assert!(
request_ids.insert(received.request_id),
"should not receive requests with duplicate request ids: {}",
received.request_id,
);
let response = NetworkMessage::RpcResponse(RpcResponse {
request_id: received.request_id,
priority: 0,
raw_response: Vec::from(&b"goodbye world"[..]),
});
server_sink.send(&response).await.unwrap();
}
assert!(matches!(server_stream.next().await, None));
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peer_send_rpc_concurrent() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, peer_handle, mut connection, _connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut server_sink, mut server_stream) = build_network_sink_stream(&mut connection);
let timeout = Duration::from_millis(10_000);
let mut request_ids = HashSet::new();
let client = async move {
let mut send_recv_futures = Vec::new();
for _ in 0..30 {
let mut peer_handle = peer_handle.clone();
let send_recv = async move {
let response = peer_handle
.send_rpc_request(PROTOCOL, Bytes::from(&b"hello world"[..]), timeout)
.await
.unwrap();
assert_eq!(response, Bytes::from(&b"goodbye world"[..]));
};
send_recv_futures.push(send_recv.boxed());
}
future::join_all(send_recv_futures).await;
};
let server = async move {
for _ in 0..30 {
let received = server_stream.next().await.unwrap().unwrap();
let received = match received {
NetworkMessage::RpcRequest(request) => request,
_ => panic!("Expected RpcRequest; unexpected: {:?}", received),
};
assert_eq!(received.protocol_id, PROTOCOL);
assert_eq!(received.priority, 0);
assert_eq!(received.raw_request, b"hello world");
assert!(
request_ids.insert(received.request_id),
"should not receive requests with duplicate request ids: {}",
received.request_id,
);
let response = NetworkMessage::RpcResponse(RpcResponse {
request_id: received.request_id,
priority: 0,
raw_response: Vec::from(&b"goodbye world"[..]),
});
server_sink.send(&response).await.unwrap();
}
assert!(matches!(server_stream.next().await, None));
};
rt.block_on(future::join3(peer.start(), server, client));
}
#[test]
fn peer_send_rpc_cancel() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, peer_handle, mut connection, _connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let (mut server_sink, mut server_stream) = build_network_sink_stream(&mut connection);
let timeout = Duration::from_millis(10_000);
let test = async move {
let (response_tx, mut response_rx) = oneshot::channel();
let request = PeerRequest::SendRpc(OutboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from(&b"hello world"[..]),
res_tx: response_tx,
timeout,
});
peer_handle.0.push(PROTOCOL, request).unwrap();
let received = server_stream.next().await.unwrap().unwrap();
let received = match received {
NetworkMessage::RpcRequest(request) => request,
_ => panic!("Expected RpcRequest; unexpected: {:?}", received),
};
assert_eq!(received.protocol_id, PROTOCOL);
assert_eq!(received.priority, 0);
assert_eq!(received.raw_request, b"hello world");
assert!(matches!(response_rx.try_recv(), Ok(None)));
drop(response_rx);
let response = NetworkMessage::RpcResponse(RpcResponse {
request_id: received.request_id,
priority: 0,
raw_response: Vec::from(&b"goodbye world"[..]),
});
server_sink.send(&response).await.unwrap();
tokio::task::yield_now().await;
drop(peer_handle);
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_send_rpc_timeout() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let mock_time = MockTimeService::new();
let (peer, peer_handle, mut connection, _connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
mock_time.clone().into(),
ConnectionOrigin::Inbound,
);
let (mut server_sink, mut server_stream) = build_network_sink_stream(&mut connection);
let timeout = Duration::from_millis(10_000);
let test = async move {
let (response_tx, mut response_rx) = oneshot::channel();
let request = PeerRequest::SendRpc(OutboundRpcRequest {
protocol_id: PROTOCOL,
data: Bytes::from(&b"hello world"[..]),
res_tx: response_tx,
timeout,
});
peer_handle.0.push(PROTOCOL, request).unwrap();
let received = server_stream.next().await.unwrap().unwrap();
let received = match received {
NetworkMessage::RpcRequest(request) => request,
_ => panic!("Expected RpcRequest; unexpected: {:?}", received),
};
assert_eq!(received.protocol_id, PROTOCOL);
assert_eq!(received.priority, 0);
assert_eq!(received.raw_request, b"hello world");
assert!(matches!(response_rx.try_recv(), Ok(None)));
mock_time.advance_async(timeout).await;
assert!(matches!(response_rx.await, Ok(Err(RpcError::TimedOut))));
let response = NetworkMessage::RpcResponse(RpcResponse {
request_id: received.request_id,
priority: 0,
raw_response: Vec::from(&b"goodbye world"[..]),
});
server_sink.send(&response).await.unwrap();
tokio::task::yield_now().await;
drop(peer_handle);
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_disconnect_request() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, peer_handle, _connection, mut connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let remote_peer_id = peer.remote_peer_id();
let test = async move {
drop(peer_handle);
assert_disconnected_event(
remote_peer_id,
DisconnectReason::Requested,
&mut connection_notifs_rx,
)
.await;
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_disconnect_connection_lost() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, _peer_handle, mut connection, mut connection_notifs_rx, _peer_notifs_rx) =
build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let remote_peer_id = peer.remote_peer_id();
let test = async move {
connection.close().await.unwrap();
assert_disconnected_event(
remote_peer_id,
DisconnectReason::ConnectionLost,
&mut connection_notifs_rx,
)
.await;
};
rt.block_on(future::join(peer.start(), test));
}
#[test]
fn peer_terminates_when_request_tx_has_dropped() {
::aptos_logger::Logger::init_for_testing();
let rt = Runtime::new().unwrap();
let (peer, peer_handle, _connection, _connection_notifs_rx, _peer_notifs_rx) = build_test_peer(
rt.handle().clone(),
TimeService::mock(),
ConnectionOrigin::Inbound,
);
let drop = async move {
drop(peer_handle);
};
rt.block_on(future::join(peer.start(), drop));
}