|
| 1 | +use discv5::Discv5Event; |
| 2 | +use ntest::timeout; |
| 3 | +use ssz::Encode; |
| 4 | +use std::{ |
| 5 | + net::{IpAddr, Ipv4Addr, SocketAddr}, |
| 6 | + str::FromStr, |
| 7 | + sync::Arc, |
| 8 | +}; |
| 9 | +use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender}; |
| 10 | +use trin_core::{ |
| 11 | + portalnet::{ |
| 12 | + discovery::Discovery, |
| 13 | + types::messages::{PortalnetConfig, ProtocolId}, |
| 14 | + Enr, |
| 15 | + }, |
| 16 | + utp::{ |
| 17 | + stream::{UtpListener, UtpListenerEvent, UtpListenerRequest, UtpStream, BUF_SIZE}, |
| 18 | + trin_helpers::{ |
| 19 | + UtpAccept, UtpMessage, |
| 20 | + UtpStreamId::{AcceptStream, OfferStream}, |
| 21 | + }, |
| 22 | + }, |
| 23 | +}; |
| 24 | + |
| 25 | +fn next_test_port() -> u16 { |
| 26 | + use std::sync::atomic::{AtomicUsize, Ordering}; |
| 27 | + static NEXT_OFFSET: AtomicUsize = AtomicUsize::new(0); |
| 28 | + const BASE_PORT: u16 = 11600; |
| 29 | + BASE_PORT + NEXT_OFFSET.fetch_add(1, Ordering::Relaxed) as u16 |
| 30 | +} |
| 31 | + |
| 32 | +/// Spawn uTP listener instance and start discv5 event handler |
| 33 | +async fn spawn_utp_listener() -> ( |
| 34 | + Enr, |
| 35 | + UnboundedSender<UtpListenerRequest>, |
| 36 | + UnboundedReceiver<UtpListenerEvent>, |
| 37 | +) { |
| 38 | + let ip_addr = IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)); |
| 39 | + let port = next_test_port(); |
| 40 | + let config = PortalnetConfig { |
| 41 | + listen_port: port, |
| 42 | + external_addr: Some(SocketAddr::new(ip_addr, port)), |
| 43 | + ..Default::default() |
| 44 | + }; |
| 45 | + let mut discv5 = Discovery::new(config).unwrap(); |
| 46 | + let enr = discv5.discv5.local_enr(); |
| 47 | + discv5.start().await.unwrap(); |
| 48 | + |
| 49 | + let discv5 = Arc::new(discv5); |
| 50 | + |
| 51 | + let (utp_event_tx, utp_listener_tx, utp_listener_rx, mut utp_listener) = |
| 52 | + UtpListener::new(Arc::clone(&discv5)); |
| 53 | + |
| 54 | + tokio::spawn(async move { |
| 55 | + let mut receiver = discv5.discv5.event_stream().await.unwrap(); |
| 56 | + while let Some(event) = receiver.recv().await { |
| 57 | + match event { |
| 58 | + Discv5Event::TalkRequest(request) => { |
| 59 | + let protocol_id = |
| 60 | + ProtocolId::from_str(&hex::encode_upper(request.protocol())).unwrap(); |
| 61 | + |
| 62 | + match protocol_id { |
| 63 | + ProtocolId::Utp => utp_event_tx.send(request).unwrap(), |
| 64 | + _ => continue, |
| 65 | + } |
| 66 | + } |
| 67 | + _ => continue, |
| 68 | + } |
| 69 | + } |
| 70 | + }); |
| 71 | + tokio::spawn(async move { utp_listener.start().await }); |
| 72 | + |
| 73 | + (enr, utp_listener_tx, utp_listener_rx) |
| 74 | +} |
| 75 | + |
| 76 | +#[tokio::test] |
| 77 | +#[timeout(100)] |
| 78 | +/// Simulate simple OFFER -> ACCEPT uTP payload transfer |
| 79 | +async fn utp_listener_events() { |
| 80 | + let protocol_id = ProtocolId::History; |
| 81 | + |
| 82 | + // Initialize offer uTP listener |
| 83 | + let (enr_offer, listener_tx_offer, mut listener_rx_offer) = spawn_utp_listener().await; |
| 84 | + // Initialize acceptor uTP listener |
| 85 | + let (enr_accept, listener_tx_accept, mut listener_rx_accept) = spawn_utp_listener().await; |
| 86 | + |
| 87 | + // Prepare to receive uTP stream from the offer node |
| 88 | + let (requested_content_key, requested_content_value) = (vec![1], vec![1, 1, 1, 1]); |
| 89 | + let stream_id = AcceptStream(vec![requested_content_key.clone()]); |
| 90 | + let conn_id = 1234; |
| 91 | + let request = UtpListenerRequest::InitiateConnection( |
| 92 | + enr_offer.clone(), |
| 93 | + protocol_id.clone(), |
| 94 | + stream_id, |
| 95 | + conn_id, |
| 96 | + ); |
| 97 | + listener_tx_accept.send(request).unwrap(); |
| 98 | + |
| 99 | + // Initialise an OFFER stream and send handshake uTP packet to the acceptor node |
| 100 | + let stream_id = OfferStream; |
| 101 | + let (tx, rx) = tokio::sync::oneshot::channel::<UtpStream>(); |
| 102 | + let offer_request = UtpListenerRequest::Connect( |
| 103 | + conn_id, |
| 104 | + enr_accept.clone(), |
| 105 | + protocol_id.clone(), |
| 106 | + stream_id, |
| 107 | + tx, |
| 108 | + ); |
| 109 | + listener_tx_offer.send(offer_request).unwrap(); |
| 110 | + |
| 111 | + // Handle STATE packet for SYN handshake in the offer node |
| 112 | + let mut conn = rx.await.unwrap(); |
| 113 | + assert_eq!(conn.connected_to, enr_accept); |
| 114 | + |
| 115 | + let mut buf = [0; BUF_SIZE]; |
| 116 | + conn.recv(&mut buf).await.unwrap(); |
| 117 | + |
| 118 | + // Send content key with content value to the acceptor node |
| 119 | + let content_items = vec![( |
| 120 | + requested_content_key.clone(), |
| 121 | + requested_content_value.clone(), |
| 122 | + )]; |
| 123 | + |
| 124 | + let content_message = UtpAccept { |
| 125 | + message: content_items, |
| 126 | + }; |
| 127 | + |
| 128 | + let utp_payload = UtpMessage::new(content_message.as_ssz_bytes()).encode(); |
| 129 | + let expected_utp_payload = utp_payload.clone(); |
| 130 | + |
| 131 | + tokio::spawn(async move { |
| 132 | + // Send the content to the acceptor over a uTP stream |
| 133 | + conn.send_to(&utp_payload).await.unwrap(); |
| 134 | + // Close uTP connection |
| 135 | + conn.close().await.unwrap(); |
| 136 | + }); |
| 137 | + |
| 138 | + // Check if the expected uTP listener events match the events in offer and accept nodes |
| 139 | + let offer_event = listener_rx_offer.recv().await.unwrap(); |
| 140 | + let expected_offer_event = |
| 141 | + UtpListenerEvent::ClosedStream(vec![], protocol_id.clone(), OfferStream); |
| 142 | + assert_eq!(offer_event, expected_offer_event); |
| 143 | + |
| 144 | + let accept_event = listener_rx_accept.recv().await.unwrap(); |
| 145 | + let expected_accept_event = UtpListenerEvent::ClosedStream( |
| 146 | + expected_utp_payload, |
| 147 | + protocol_id.clone(), |
| 148 | + AcceptStream(vec![requested_content_key]), |
| 149 | + ); |
| 150 | + assert_eq!(accept_event, expected_accept_event); |
| 151 | +} |
0 commit comments