use crate::crazyflie::LINK_PORT;
use crate::{Error, Result};
use crazyflie_link::Packet;
use flume as channel;
use futures::lock::Mutex;
use std::sync::Arc;
use std::time::{Duration, Instant};
const ECHO_CHANNEL: u8 = 0;
const SOURCE_CHANNEL: u8 = 1;
const SINK_CHANNEL: u8 = 2;
const MAX_DATA_SIZE: usize = 30;
const FILL_PATTERN: u8 = 0xAA;
#[derive(Debug, Clone)]
pub struct Statistics {
pub link_quality: Option<f32>,
pub uplink_rate: Option<f32>,
pub downlink_rate: Option<f32>,
pub radio_send_rate: Option<f32>,
pub avg_retries: Option<f32>,
pub power_detector_rate: Option<f32>,
pub rssi: Option<f32>,
}
#[derive(Debug, Clone)]
pub struct BandwidthResult {
pub uplink_bytes_per_sec: f64,
pub downlink_bytes_per_sec: f64,
pub packets_per_sec: f64,
}
pub struct LinkService {
uplink: channel::Sender<Packet>,
echo_downlink: Mutex<channel::Receiver<Packet>>,
source_downlink: Mutex<channel::Receiver<Packet>>,
link: Arc<crazyflie_link::Connection>,
}
impl LinkService {
pub(crate) fn new(
uplink: channel::Sender<Packet>,
downlink: channel::Receiver<Packet>,
link: Arc<crazyflie_link::Connection>,
) -> Self {
let (echo_downlink, source_downlink, _, _) =
crate::crtp_utils::crtp_channel_dispatcher(downlink);
Self {
uplink,
echo_downlink: Mutex::new(echo_downlink),
source_downlink: Mutex::new(source_downlink),
link,
}
}
pub async fn ping(&self) -> Result<f64> {
let echo_downlink = self.echo_downlink.lock().await;
const PING_PAYLOAD: [u8; 1] = [0x01];
let start = Instant::now();
let pk = Packet::new(LINK_PORT, ECHO_CHANNEL, PING_PAYLOAD.to_vec());
self.uplink
.send_async(pk)
.await
.map_err(|_| Error::Disconnected)?;
let answer = tokio::time::timeout(Duration::from_secs(1), echo_downlink.recv_async())
.await
.map_err(|_| Error::Timeout)?
.map_err(|_| Error::Disconnected)?;
if answer.get_data() != &PING_PAYLOAD {
return Err(Error::ProtocolError("Ping got wrong echo back".to_string()));
}
Ok(start.elapsed().as_secs_f64() * 1000.0)
}
pub async fn test_uplink_bandwidth(&self, n_packets: u64) -> Result<f64> {
let data = vec![FILL_PATTERN; MAX_DATA_SIZE];
let start = Instant::now();
let mut total_bytes: u64 = 0;
for _ in 0..n_packets {
let pk = Packet::new(LINK_PORT, SINK_CHANNEL, data.clone());
self.uplink
.send_async(pk)
.await
.map_err(|_| Error::Disconnected)?;
total_bytes += MAX_DATA_SIZE as u64;
}
const ECHO_PAYLOAD: [u8; 1] = [0x00];
let echo = Packet::new(LINK_PORT, ECHO_CHANNEL, ECHO_PAYLOAD.to_vec());
let echo_downlink = self.echo_downlink.lock().await;
self.uplink
.send_async(echo)
.await
.map_err(|_| Error::Disconnected)?;
let answer = tokio::time::timeout(Duration::from_secs(10), echo_downlink.recv_async())
.await
.map_err(|_| Error::Timeout)?
.map_err(|_| Error::Disconnected)?;
if answer.get_data() != &ECHO_PAYLOAD {
return Err(Error::ProtocolError(
"Echo got wrong payload back".to_string(),
));
}
let elapsed = start.elapsed().as_secs_f64();
Ok(total_bytes as f64 / elapsed)
}
pub async fn test_downlink_bandwidth(&self, n_packets: u64) -> Result<f64> {
let source_downlink = self.source_downlink.lock().await;
let start = Instant::now();
let mut total_bytes: u64 = 0;
for _ in 0..n_packets {
let pk = Packet::new(LINK_PORT, SOURCE_CHANNEL, vec![0x00]);
self.uplink
.send_async(pk)
.await
.map_err(|_| Error::Disconnected)?;
}
for _ in 0..n_packets {
let response =
tokio::time::timeout(Duration::from_secs(1), source_downlink.recv_async())
.await
.map_err(|_| Error::Timeout)?
.map_err(|_| Error::Disconnected)?;
total_bytes += response.get_data().len() as u64;
}
let elapsed = start.elapsed().as_secs_f64();
Ok(total_bytes as f64 / elapsed)
}
pub async fn test_echo_bandwidth(&self, n_packets: u64) -> Result<BandwidthResult> {
let echo_downlink = self.echo_downlink.lock().await;
let data = vec![FILL_PATTERN; MAX_DATA_SIZE];
let start = Instant::now();
let mut packets: u64 = 0;
for _ in 0..n_packets {
let pk = Packet::new(LINK_PORT, ECHO_CHANNEL, data.clone());
self.uplink
.send_async(pk)
.await
.map_err(|_| Error::Disconnected)?;
}
for _ in 0..n_packets {
let answer = tokio::time::timeout(Duration::from_secs(1), echo_downlink.recv_async())
.await
.map_err(|_| Error::Timeout)?
.map_err(|_| Error::Disconnected)?;
if answer.get_data() != data.as_slice() {
return Err(Error::ProtocolError(
"Echo got wrong payload back".to_string(),
));
}
packets += 1;
}
let elapsed = start.elapsed().as_secs_f64();
let bytes = packets as f64 * MAX_DATA_SIZE as f64;
Ok(BandwidthResult {
uplink_bytes_per_sec: bytes / elapsed,
downlink_bytes_per_sec: bytes / elapsed,
packets_per_sec: packets as f64 / elapsed,
})
}
pub async fn get_statistics(&self) -> Statistics {
let radio_stats = self.link.link_statistics().await;
Statistics {
link_quality: radio_stats.as_ref().map(|s| s.link_quality),
uplink_rate: radio_stats.as_ref().map(|s| s.uplink_rate),
downlink_rate: radio_stats.as_ref().map(|s| s.downlink_rate),
radio_send_rate: radio_stats.as_ref().map(|s| s.radio_send_rate),
avg_retries: radio_stats.as_ref().map(|s| s.avg_retries),
power_detector_rate: radio_stats.as_ref().map(|s| s.power_detector_rate),
rssi: radio_stats.and_then(|s| s.rssi),
}
}
}