Skip to content

Commit 95c5a02

Browse files
committed
feat: remove task spawn in session-sender
controller already creates a task for session-sender's recv functionality removed extra spawn inside session-sender's recv function
1 parent 19251a5 commit 95c5a02

3 files changed

Lines changed: 88 additions & 38 deletions

File tree

examples/controller/controller.rs

Lines changed: 8 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ impl Controller {
5757
let (start_session_tx, start_session_rx) = oneshot::channel::<()>();
5858
let (twamp_test_complete_tx, twamp_test_complete_rx) = oneshot::channel::<()>();
5959
let (reflector_port_tx, reflector_port_rx) = oneshot::channel::<u16>();
60-
let control_client_handle = spawn(async move {
60+
let task_control_client = spawn(async move {
6161
self.control_client
6262
.do_twamp_control(
6363
twamp_control,
@@ -71,10 +71,9 @@ impl Controller {
7171
.await
7272
.unwrap();
7373
});
74-
let reflected_pkts_vec: Arc<Mutex<Vec<(TwampTestPacketUnauthReflected, TimeStamp)>>> =
75-
Arc::new(Mutex::new(Vec::new()));
74+
let reflected_pkts_vec = Arc::new(Mutex::new(Vec::new()));
7675
let reflected_pkts_vec_cloned = Arc::clone(&reflected_pkts_vec);
77-
let session_sender_handle = spawn(async move {
76+
let task_session_sender = spawn(async move {
7877
// Wait until we get the Accept-Session's port.
7978
let final_port = reflector_port_rx.await.unwrap();
8079
debug!("Received reflector port: {}", final_port);
@@ -94,31 +93,31 @@ impl Controller {
9493
));
9594
let session_sender_send = Arc::clone(self.session_sender.as_ref().unwrap());
9695
let session_sender_recv = Arc::clone(self.session_sender.as_ref().unwrap());
97-
let send_task = spawn(async move {
96+
let task_session_sender_send = spawn(async move {
9897
let _ = session_sender_send.send_it(number_of_test_packets).await;
9998
info!("Sent all test packets");
10099
});
101-
let recv_task = spawn(async move {
100+
let task_session_sender_recv = spawn(async move {
102101
let _ = session_sender_recv
103102
.recv(number_of_test_packets, reflected_pkts_vec_cloned)
104103
.await;
105104
info!("Got back all test packets");
106105
});
107106
// wait for all test pkts to be sent.
108-
send_task.await.unwrap();
107+
task_session_sender_send.await.unwrap();
109108

110109
select! {
111110
// If stop-session-sleep duration finishes before all pkts are received, drop
112111
// recv task and conclude.
113112
_ = sleep(Duration::from_secs(stop_session_sleep)) => (),
114113
// Ignore stop-session-sleep duration if session-sender got all test pkts before
115114
// duration.
116-
_ = recv_task => ()
115+
_ = task_session_sender_recv => ()
117116
}
118117
// Inform Control-Client to send Stop-Sessions
119118
twamp_test_complete_tx.send(()).unwrap();
120119
});
121-
try_join!(control_client_handle, session_sender_handle).unwrap();
120+
try_join!(task_control_client, task_session_sender).unwrap();
122121
debug!("Control-Client & Session-Sender tasks completed.");
123122
let acquired_vec = reflected_pkts_vec.lock().await;
124123
debug!("Reflected pkts len: {}", acquired_vec.len());

src/session_reflector/mod.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ use thiserror::Error;
99
use tokio::{net::UdpSocket, time::timeout};
1010
use tracing::*;
1111

12+
/// Errors that can be raised by [SessionReflector].
1213
#[derive(Error, Debug)]
1314
pub enum SessionReflectorError {
1415
/// Raised when "connect" call to a Session-Sender fails.

src/session_sender/mod.rs

Lines changed: 79 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,53 @@
11
use crate::timestamp::TimeStamp;
22
use crate::twamp_test::{TwampTestPacketUnauth, TwampTestPacketUnauthReflected};
3-
use anyhow::Result;
43
use deku::prelude::*;
54
use std::{
5+
io,
66
net::{SocketAddr, SocketAddrV4},
77
sync::Arc,
88
};
9-
use tokio::{net::UdpSocket, spawn, sync::Mutex};
9+
use thiserror::Error;
10+
use tokio::{net::UdpSocket, sync::Mutex};
1011
use tracing::*;
1112

13+
/// Errors that can be raised by [SessionSender].
14+
#[derive(Error, Debug)]
15+
pub enum SessionSenderError {
16+
/// Raised when "connect" call to a Session-Reflector fails.
17+
///
18+
/// - First field is the socket address of the Session-Reflector it was trying to connect to.
19+
#[error("failed to connect to Session-Reflector: {0}")]
20+
SessionReflectorConnectError(SocketAddrV4, #[source] io::Error),
21+
22+
/// Indicates read error on TWAMP-Test **after** it was initiated.
23+
///
24+
/// - First field is the number of packets that were reflected back to Session-Sender before
25+
/// this error.
26+
#[error("failed to read UDP datagram from socket after {0} packets were reflected")]
27+
SessionReflectorReadError(u32, #[source] io::Error),
28+
29+
#[error("failed to send UDP datagram to Session-Reflector after {0} packets were sent")]
30+
SessionReflectorWriteError(u32, #[source] io::Error),
31+
32+
/// Indicates Session-Sender not being bound to a socket.
33+
///
34+
/// This usually indicates an OS level error, since UdpSocket is made from
35+
/// "bind", and this error indicates something unbinded the UdpSocket after it's creation.
36+
#[error("Session-Sender was not not bound to a socket")]
37+
NotBound(#[source] io::Error),
38+
39+
/// Indicates Session-Sender was not "connected" to Session-Reflector's UdpSocket.
40+
///
41+
/// TODO: This can also indicate an OS level issue, since "connection" is established in
42+
/// [SessionSender::new].
43+
#[error("Session-Reflector was not connected to socket")]
44+
SessionReflectorNotConnected(#[source] io::Error),
45+
46+
/// Indicates underlying (de)serialization failure.
47+
#[error("failed to convert between rust and wire format")]
48+
WireConversionError(#[source] deku::error::DekuError),
49+
}
50+
1251
#[derive(Debug)]
1352
pub struct SessionSender {
1453
pub socket: Arc<UdpSocket>,
@@ -23,16 +62,28 @@ impl SessionSender {
2362
}
2463
}
2564

26-
pub async fn send_it(&self, number_of_packets: u32) -> Result<()> {
65+
pub async fn send_it(&self, number_of_packets: u32) -> Result<(), SessionSenderError> {
2766
info!("Sending Twamp-Test packets to {}", self.dest);
2867
for i in 0..number_of_packets {
2968
let twamp_test = TwampTestPacketUnauth::new(i, 0, true);
3069
trace!("Twamp-Test: {:?}", twamp_test);
31-
let encoded = twamp_test.to_bytes().unwrap();
32-
let l = self.socket.local_addr().unwrap();
33-
let p = self.socket.peer_addr().unwrap();
34-
trace!("Sending pkt from {} to {}", l, p);
35-
let len = self.socket.send(&encoded[..]).await?;
70+
let encoded = twamp_test
71+
.to_bytes()
72+
.map_err(|err| SessionSenderError::WireConversionError(err))?;
73+
let local_addr = self
74+
.socket
75+
.local_addr()
76+
.map_err(|err| SessionSenderError::NotBound(err))?;
77+
let peer_addr = self
78+
.socket
79+
.peer_addr()
80+
.map_err(|err| SessionSenderError::SessionReflectorNotConnected(err))?;
81+
trace!("Sending pkt from {} to {}", local_addr, peer_addr);
82+
let len = self
83+
.socket
84+
.send(&encoded[..])
85+
.await
86+
.map_err(|err| SessionSenderError::SessionReflectorWriteError(i, err))?;
3687
trace!("Twamp-Test sent of bytes: {}", len);
3788
}
3889
Ok(())
@@ -42,28 +93,27 @@ impl SessionSender {
4293
&self,
4394
number_of_packets: u32,
4495
reflected_pkts_shared: Arc<Mutex<Vec<(TwampTestPacketUnauthReflected, TimeStamp)>>>,
45-
) {
46-
let sock_clone = Arc::clone(&self.socket);
47-
let reflect_task = spawn(async move {
48-
let mut count: u32 = 1;
49-
loop {
50-
let mut buf = [0u8; 1024]; // Buffer to hold incoming packets
51-
let bytes_read = sock_clone.recv(&mut buf).await.unwrap();
52-
trace!("Bytes read: {}", bytes_read);
53-
let (_rest, reflected_pkt) =
54-
TwampTestPacketUnauthReflected::from_bytes((&buf, 0)).unwrap();
55-
trace!("Received reflected pkt: {:?}", reflected_pkt);
56-
//debug!("Adding reflector pkt to vec");
57-
let mut acquired_vec = reflected_pkts_shared.lock().await;
58-
//debug!("Added reflector pkt to vec");
59-
acquired_vec.push((reflected_pkt, TimeStamp::default()));
60-
if count == number_of_packets {
61-
break;
62-
}
63-
count += 1;
96+
) -> Result<(), SessionSenderError> {
97+
let mut count: u32 = 0;
98+
loop {
99+
let mut buf = [0u8; 1024]; // Buffer to hold incoming packets
100+
let bytes_read = self
101+
.socket
102+
.recv(&mut buf)
103+
.await
104+
.map_err(|err| SessionSenderError::SessionReflectorReadError(count, err))?;
105+
trace!("Bytes read: {}", bytes_read);
106+
let (_rest, reflected_pkt) = TwampTestPacketUnauthReflected::from_bytes((&buf, 0))
107+
.map_err(|err| SessionSenderError::WireConversionError(err))?;
108+
trace!("Received reflected pkt: {:?}", reflected_pkt);
109+
let mut acquired_vec = reflected_pkts_shared.lock().await;
110+
acquired_vec.push((reflected_pkt, TimeStamp::default()));
111+
count += 1;
112+
if count == number_of_packets {
113+
break;
64114
}
65-
});
66-
reflect_task.await.unwrap()
115+
}
116+
Ok(())
67117
}
68118
}
69119

0 commit comments

Comments
 (0)