use crate::twamp_control::Accept;
use crate::twamp_control::AcceptSession;
use crate::twamp_control::RequestTwSession;
use crate::twamp_control::SecurityMode;
use crate::twamp_control::ServerGreeting;
use crate::twamp_control::ServerStart;
use crate::twamp_control::SetUpResponse;
use crate::twamp_control::StartAck;
use crate::twamp_control::StartSessions;
use crate::twamp_control::StopSessions;
use anyhow::{Result, anyhow};
use deku::prelude::*;
use std::mem::size_of;
use std::net::IpAddr;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::TcpStream;
use tokio::sync::oneshot;
use tracing::*;
#[derive(Debug)]
pub struct ControlClient {
pub stream: Option<TcpStream>,
}
impl ControlClient {
pub fn new() -> Self {
Self { stream: None }
}
pub async fn do_twamp_control(
&mut self,
twamp_control: TcpStream,
start_session_tx: oneshot::Sender<()>,
reflector_port_tx: oneshot::Sender<u16>,
responder_reflect_port: u16,
controller_port: u16,
reflector_timeout: u64,
twamp_test_complete_rx: oneshot::Receiver<()>,
) -> Result<()> {
self.stream = Some(twamp_control);
self.read_server_greeting().await?;
self.send_set_up_response().await?;
self.read_server_start().await?;
self.send_request_tw_session(responder_reflect_port, controller_port, reflector_timeout)
.await?;
let accept_session = self.read_accept_session().await?;
if accept_session.accept != Accept::Ok {
return Err(anyhow!("Did not receive Ok in Accept-Session"));
};
debug!("Responder provided port: {}", accept_session.port);
reflector_port_tx.send(accept_session.port).unwrap();
self.send_start_sessions().await?;
let start_ack = self.read_start_ack().await?;
if start_ack.accept != Accept::Ok {
return Err(anyhow!("Start-Ack should be zero"));
}
start_session_tx.send(()).unwrap();
debug!(
"Waiting for Session-Sender to complete, Control-Client will then send Stop-Sessions."
);
let _ = twamp_test_complete_rx.await;
debug!("Received confirmation that TWAMP-Test is complete. Sending Stop-Sessions");
self.send_stop_sessions().await?;
Ok(())
}
pub async fn read_server_greeting(&mut self) -> Result<ServerGreeting> {
let mut buf = [0; size_of::<ServerGreeting>()];
info!("Reading ServerGreeting");
self.stream.as_mut().unwrap().read_exact(&mut buf).await?;
let (_rest, server_greeting) = ServerGreeting::from_bytes((&buf, 0)).expect("should have received valid ServerGreeting");
debug!("Server greeting: {:?}", server_greeting);
info!("Done reading ServerGreeting");
Ok(server_greeting)
}
pub async fn send_set_up_response(&mut self) -> Result<()> {
info!("Preparing to send Set-Up-Response");
let set_up_response = SetUpResponse::new(SecurityMode::Unauthenticated);
debug!("Set-Up-Response: {:?}", set_up_response);
let encoded = set_up_response.unwrap().to_bytes().unwrap();
self.stream
.as_mut()
.unwrap()
.write_all(&encoded[..])
.await?;
info!("Set-Up-Response sent");
Ok(())
}
pub async fn read_server_start(&mut self) -> Result<ServerStart> {
let mut buf = [0; size_of::<ServerStart>()];
info!("Reading Server-Start");
self.stream.as_mut().unwrap().read_exact(&mut buf).await?;
let (_rest, server_start) = ServerStart::from_bytes((&buf, 0)).unwrap();
debug!("Server-Start: {:?}", server_start);
info!("Done reading Server-Start");
Ok(server_start)
}
pub async fn send_request_tw_session(
&mut self,
session_reflector_port: u16,
controller_port: u16,
timeout: u64,
) -> Result<RequestTwSession> {
info!("Preparing to send Request-TW-Session");
let stream = self.stream.as_ref().unwrap();
let sender_address = match stream.local_addr().unwrap().ip() {
IpAddr::V4(ip) => ip,
IpAddr::V6(ip) => panic!("da hail did v6 come from: {ip}"),
};
let receiver_address = match stream.peer_addr().unwrap().ip() {
IpAddr::V4(ip) => ip,
IpAddr::V6(ip) => panic!("da hail did v6 come from: {ip}"),
};
debug!(
"Request-TW-Session reflector port: {}",
session_reflector_port
);
let request_tw_session = RequestTwSession::new(
sender_address,
controller_port,
receiver_address,
session_reflector_port,
None,
timeout,
);
debug!("request-tw-session: {:?}", request_tw_session);
let encoded = request_tw_session.to_bytes().unwrap();
self.stream
.as_mut()
.unwrap()
.write_all(&encoded[..])
.await?;
info!("Request-TW-Session sent");
Ok(request_tw_session)
}
pub async fn read_accept_session(&mut self) -> Result<AcceptSession> {
let mut buf = [0; size_of::<AcceptSession>()];
info!("Reading Accept-Session");
self.stream.as_mut().unwrap().read_exact(&mut buf).await?;
let (_rest, accept_session) = AcceptSession::from_bytes((&buf, 0)).unwrap();
debug!("Accept-Session: {:?}", accept_session);
info!("Read Accept-Session");
Ok(accept_session)
}
pub async fn send_start_sessions(&mut self) -> Result<()> {
info!("Preparing to send Start-Sessions");
let start_sessions = StartSessions::new();
debug!("Start-Sessions: {:?}", start_sessions);
let encoded = start_sessions.to_bytes().unwrap();
self.stream
.as_mut()
.unwrap()
.write_all(&encoded[..])
.await?;
info!("Start-Sessions sent");
Ok(())
}
pub async fn read_start_ack(&mut self) -> Result<StartAck> {
let mut buf = [0; size_of::<StartAck>()];
info!("Reading Start-Ack");
self.stream.as_mut().unwrap().read_exact(&mut buf).await?;
let (_rest, start_ack) = StartAck::from_bytes((&buf, 0)).unwrap();
debug!("Start-Ack: {:?}", start_ack);
info!("Done reading Start-Ack");
Ok(start_ack)
}
pub async fn send_stop_sessions(&mut self) -> Result<()> {
info!("Preparing to send Stop-Sessions");
let stop_sessions = StopSessions::new(Accept::Ok);
debug!("Stop-Sessions: {:?}", stop_sessions);
let encoded = stop_sessions.to_bytes().unwrap();
self.stream
.as_mut()
.unwrap()
.write_all(&encoded[..])
.await?;
info!("Stop-Sessions sent");
Ok(())
}
}
impl Default for ControlClient {
fn default() -> Self {
ControlClient { stream: None }
}
}