1 use clap::{App, AppSettings, Arg}; 2 use std::io::Write; 3 use std::sync::Arc; 4 use tokio::net::UdpSocket; 5 use util::{conn::conn_disconnected_packet::DisconnectedPacketConn, Conn}; 6 use webrtc_sctp::association::*; 7 use webrtc_sctp::chunk::chunk_payload_data::PayloadProtocolIdentifier; 8 use webrtc_sctp::stream::*; 9 use webrtc_sctp::Error; 10 11 fn main() -> Result<(), Error> { 12 env_logger::Builder::new() 13 .format(|buf, record| { 14 writeln!( 15 buf, 16 "{}:{} [{}] {} - {}", 17 record.file().unwrap_or("unknown"), 18 record.line().unwrap_or(0), 19 record.level(), 20 chrono::Local::now().format("%H:%M:%S.%6f"), 21 record.args() 22 ) 23 }) 24 .filter(None, log::LevelFilter::Warn) 25 .init(); 26 27 let mut app = App::new("SCTP Throughput") 28 .version("0.1.0") 29 .about("An example of SCTP Server") 30 .setting(AppSettings::DeriveDisplayOrder) 31 .setting(AppSettings::SubcommandsNegateReqs) 32 .arg( 33 Arg::with_name("FULLHELP") 34 .help("Prints more detailed help information") 35 .long("fullhelp"), 36 ) 37 .arg( 38 Arg::with_name("port") 39 .required_unless("FULLHELP") 40 .takes_value(true) 41 .long("port") 42 .help("use port ."), 43 ); 44 45 let matches = app.clone().get_matches(); 46 47 if matches.is_present("FULLHELP") { 48 app.print_long_help().unwrap(); 49 std::process::exit(0); 50 } 51 52 let port1 = matches.value_of("port").unwrap().to_owned(); 53 let port2 = port1.clone(); 54 55 std::thread::spawn(|| { 56 tokio::runtime::Runtime::new() 57 .unwrap() 58 .block_on(async move { 59 let conn = DisconnectedPacketConn::new(Arc::new( 60 UdpSocket::bind(format!("127.0.0.1:{}", port1)) 61 .await 62 .unwrap(), 63 )); 64 println!("listening {}...", conn.local_addr().unwrap()); 65 66 let config = Config { 67 net_conn: Arc::new(conn), 68 max_receive_buffer_size: 0, 69 max_message_size: 0, 70 name: "recver".to_owned(), 71 }; 72 let a = Association::server(config).await?; 73 println!("created a server"); 74 75 let stream = a.accept_stream().await.unwrap(); 76 println!("accepted a stream"); 77 78 // set unordered = true and 10ms treshold for dropping packets 79 stream.set_reliability_params(true, ReliabilityType::Rexmit, 0); 80 81 let mut buff = [0u8; 65535]; 82 let mut recv = 0; 83 let mut pkt_num = 0; 84 let mut loop_num = 0; 85 let mut now = tokio::time::Instant::now(); 86 while let Ok(n) = stream.read(&mut buff).await { 87 recv += n; 88 if n != 0 { 89 pkt_num += 1; 90 } 91 loop_num += 1; 92 if now.elapsed().as_secs() == 1 { 93 println!( 94 "Throughput: {} Bytes/s, {} pkts, {} loops", 95 recv, pkt_num, loop_num 96 ); 97 now = tokio::time::Instant::now(); 98 recv = 0; 99 loop_num = 0; 100 pkt_num = 0; 101 } 102 } 103 Result::<(), Error>::Ok(()) 104 }) 105 }); 106 107 std::thread::spawn(|| { 108 tokio::runtime::Runtime::new() 109 .unwrap() 110 .block_on(async move { 111 let conn = Arc::new(UdpSocket::bind("0.0.0.0:0").await.unwrap()); 112 conn.connect(format!("127.0.0.1:{}", port2)).await.unwrap(); 113 println!("connecting 127.0.0.1:{}..", port2); 114 115 let config = Config { 116 net_conn: conn, 117 max_receive_buffer_size: 0, 118 max_message_size: 0, 119 name: "sender".to_owned(), 120 }; 121 let a = Association::client(config).await.unwrap(); 122 println!("created a client"); 123 124 let stream = a 125 .open_stream(0, PayloadProtocolIdentifier::Binary) 126 .await 127 .unwrap(); 128 println!("opened a stream"); 129 130 //const LEN: usize = 1200; 131 const LEN: usize = 65535; 132 let buf = vec![0; LEN]; 133 let bytes = bytes::Bytes::from(buf); 134 135 let mut now = tokio::time::Instant::now(); 136 let mut pkt_num = 0; 137 while stream.write(&bytes).await.is_ok() { 138 pkt_num += 1; 139 if now.elapsed().as_secs() == 1 { 140 println!("Send {} pkts", pkt_num); 141 now = tokio::time::Instant::now(); 142 pkt_num = 0; 143 } 144 } 145 Result::<(), Error>::Ok(()) 146 }) 147 }); 148 loop {} 149 } 150