use std::{
net::{Ipv4Addr, SocketAddr},
time::Duration,
};
use chia_protocol::Bytes;
use dig_peer_protocol::{DigLink, DigMessage, LinkOptions, HOLDINGS_ANNOUNCE};
use tokio::io::DuplexStream;
use tokio_tungstenite::{tungstenite::protocol::Role, WebSocketStream};
const PATIENCE: Duration = Duration::from_secs(5);
async fn linked_pair(
options: LinkOptions,
) -> (
DigLink,
tokio::sync::mpsc::Receiver<DigMessage>,
DigLink,
tokio::sync::mpsc::Receiver<DigMessage>,
) {
let (left, right) = tokio::io::duplex(8 * 1024 * 1024);
let addr = SocketAddr::from((Ipv4Addr::LOCALHOST, 8444));
let client: WebSocketStream<DuplexStream> =
WebSocketStream::from_raw_socket(left, Role::Client, None).await;
let server: WebSocketStream<DuplexStream> =
WebSocketStream::from_raw_socket(right, Role::Server, None).await;
let (a, a_rx) = DigLink::from_server_websocket(client, addr, options);
let (b, b_rx) = DigLink::from_server_websocket(server, addr, options);
(a, a_rx, b, b_rx)
}
#[tokio::test]
async fn an_unsendably_large_message_is_refused_rather_than_retried_forever() {
let (link, _rx, _peer, _peer_rx) = linked_pair(LinkOptions::default()).await;
let over_cap = Bytes::new(vec![0u8; 1024 * 1024 + 1]);
let refusal = tokio::time::timeout(PATIENCE, link.send_dig(HOLDINGS_ANNOUNCE, over_cap))
.await
.expect("send spun instead of refusing a message that can never fit");
assert!(
refusal.is_err(),
"an oversized message reported success without being sendable"
);
let at_cap = Bytes::new(vec![0u8; 1024 * 1024]);
tokio::time::timeout(PATIENCE, link.send_dig(HOLDINGS_ANNOUNCE, at_cap))
.await
.expect("the at-cap control send spun")
.expect("the at-cap control send was refused, so the refusal above proves nothing");
}
#[tokio::test]
async fn a_flood_of_unmatched_ids_cannot_wedge_the_correlated_reply_path() {
let mut options = LinkOptions::default();
options.rate_limit_factor = 1.0;
let (peer, mut peer_rx, requester, _requester_rx) = linked_pair(options).await;
let peer_task = tokio::spawn(async move {
let request = peer_rx.recv().await.expect("the request arrives");
for offset in 0..64u16 {
peer.send_message(DigMessage::new(
HOLDINGS_ANNOUNCE,
Some(u16::MAX - offset),
Bytes::new(b"unmatched".to_vec()),
))
.await
.expect("send an unmatched frame");
}
peer.send_message(DigMessage::new(
HOLDINGS_ANNOUNCE,
request.id,
Bytes::new(b"pong".to_vec()),
))
.await
.expect("send the correlated reply");
});
let response = tokio::time::timeout(
PATIENCE,
requester.request_dig(HOLDINGS_ANNOUNCE, Bytes::new(b"ping".to_vec())),
)
.await
.expect("the flood wedged the reader: no correlated reply was ever routed")
.expect("the request failed");
assert_eq!(response.data.as_ref(), b"pong");
peer_task.await.expect("the peer finished");
}
#[tokio::test]
async fn an_unanswered_request_errors_on_its_deadline() {
let mut options = LinkOptions::default();
options.request_timeout = Duration::from_millis(300);
let (_peer, _peer_rx, requester, _requester_rx) = linked_pair(options).await;
let outcome = tokio::time::timeout(
PATIENCE,
requester.request_dig(HOLDINGS_ANNOUNCE, Bytes::new(b"ping".to_vec())),
)
.await
.expect("the request hung past its deadline");
assert!(
outcome.is_err(),
"an unanswered request resolved successfully"
);
}
#[tokio::test]
async fn a_late_reply_to_a_timed_out_request_is_not_misrouted_to_the_next_waiter() {
let mut options = LinkOptions::default();
options.request_timeout = Duration::from_millis(200);
let (peer, mut peer_rx, requester, _requester_rx) = linked_pair(options).await;
let a_id = {
let _ = tokio::time::timeout(
PATIENCE,
requester.request_dig(HOLDINGS_ANNOUNCE, Bytes::new(b"question-A".to_vec())),
)
.await
.expect("request A should have timed out, not hung");
peer_rx.recv().await.expect("peer receives question-A").id
};
let peer_task = tokio::spawn(async move {
let b_msg = peer_rx.recv().await.expect("peer receives question-B");
peer.send_message(DigMessage::new(
HOLDINGS_ANNOUNCE,
a_id,
Bytes::new(b"answer-to-A".to_vec()),
))
.await
.expect("peer sends late answer-to-A");
peer.send_message(DigMessage::new(
HOLDINGS_ANNOUNCE,
b_msg.id,
Bytes::new(b"answer-to-B".to_vec()),
))
.await
.expect("peer sends answer-to-B");
});
let b_response = tokio::time::timeout(
PATIENCE,
requester.request_dig(HOLDINGS_ANNOUNCE, Bytes::new(b"question-B".to_vec())),
)
.await
.expect("request B hung")
.expect("request B errored");
assert_eq!(
b_response.data.as_ref(),
b"answer-to-B",
"MISROUTED: question-B received a reply intended for question-A"
);
peer_task.await.expect("peer task finished");
}