xref: /webrtc/sctp/examples/throughput.rs (revision a4acb6a4)
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 {}..", format!("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 mut buf = Vec::with_capacity(LEN);
133                 unsafe {
134                     buf.set_len(LEN);
135                 }
136 
137                 let mut now = tokio::time::Instant::now();
138                 let mut pkt_num = 0;
139                 while stream.write(&buf.clone().into()).await.is_ok() {
140                     pkt_num += 1;
141                     if now.elapsed().as_secs() == 1 {
142                         println!("Send {} pkts", pkt_num);
143                         now = tokio::time::Instant::now();
144                         pkt_num = 0;
145                     }
146                 }
147                 Result::<(), Error>::Ok(())
148             })
149     });
150     loop {}
151 }
152