use crate::{
constants,
peer::Peer,
protocols::wire::{
handshake::v1::{MessagingProtocolVersion, ProtocolIdSet},
messaging::v1::{NetworkMessage, NetworkMessageSink},
},
testutils::fake_socket::ReadOnlyTestSocketVec,
transport::{Connection, ConnectionId, ConnectionMetadata},
};
use aptos_config::{config::PeerRole, network_id::NetworkContext};
use aptos_proptest_helpers::ValueGenerator;
use aptos_time_service::TimeService;
use aptos_types::{network_address::NetworkAddress, PeerId};
use channel::{aptos_channel, message_queues::QueueStyle};
use futures::{executor::block_on, future, io::AsyncReadExt, sink::SinkExt, stream::StreamExt};
use memsocket::MemorySocket;
use netcore::transport::ConnectionOrigin;
use proptest::{arbitrary::any, collection::vec};
use std::time::Duration;
pub fn generate_corpus(gen: &mut ValueGenerator) -> Vec<u8> {
let network_msgs = gen.generate(vec(any::<NetworkMessage>(), 1..20));
let (write_socket, mut read_socket) = MemorySocket::new_pair();
let mut writer = NetworkMessageSink::new(write_socket, constants::MAX_FRAME_SIZE, None);
let f_send = async move {
for network_msg in &network_msgs {
writer.send(network_msg).await.unwrap();
}
};
let f_recv = async move {
let mut buf = Vec::new();
read_socket.read_to_end(&mut buf).await.unwrap();
buf
};
let (_, buf) = block_on(future::join(f_send, f_recv));
buf
}
pub fn fuzz(data: &[u8]) {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
let executor = rt.handle().clone();
let peer_id = PeerId::ZERO;
let remote_peer_id = PeerId::random();
let network_context = NetworkContext::mock_with_peer_id(peer_id);
let socket = ReadOnlyTestSocketVec::new(data.to_vec());
let metadata = ConnectionMetadata::new(
remote_peer_id,
ConnectionId::from(123),
NetworkAddress::mock(),
ConnectionOrigin::Inbound,
MessagingProtocolVersion::V1,
ProtocolIdSet::all_known(),
PeerRole::Unknown,
);
let connection = Connection { socket, metadata };
let (connection_notifs_tx, connection_notifs_rx) = channel::new_test(8);
let channel_size = 8;
let (peer_reqs_tx, peer_reqs_rx) = aptos_channel::new(QueueStyle::FIFO, channel_size, None);
let (peer_notifs_tx, peer_notifs_rx) = aptos_channel::new(QueueStyle::FIFO, channel_size, None);
let peer = Peer::new(
network_context,
executor.clone(),
TimeService::mock(),
connection,
connection_notifs_tx,
peer_reqs_rx,
peer_notifs_tx,
Duration::from_millis(constants::INBOUND_RPC_TIMEOUT_MS),
constants::MAX_CONCURRENT_INBOUND_RPCS,
constants::MAX_CONCURRENT_OUTBOUND_RPCS,
constants::MAX_FRAME_SIZE,
None,
None,
);
executor.spawn(peer.start());
rt.block_on(async move {
connection_notifs_rx.collect::<Vec<_>>().await;
drop(peer_reqs_tx);
peer_notifs_rx.collect::<Vec<_>>().await;
});
}
#[test]
fn test_peer_fuzzers() {
let mut value_gen = ValueGenerator::deterministic();
for _ in 0..50 {
let corpus = generate_corpus(&mut value_gen);
fuzz(&corpus);
}
}