use crazyflie_link::Packet;
use flume::{Receiver, Sender};
use async_broadcast::{broadcast, Receiver as BroadcastReceiver};
use futures::Stream;
use half::f16;
use crate::{Error, Result};
use crate::crazyflie::{LOCALIZATION_PORT, SUPERVISOR_PORT};
use crate::subsystems::supervisor::{
CMD_EMERGENCY_STOP, CMD_EMERGENCY_STOP_WATCHDOG, SUPERVISOR_CH_COMMAND,
};
const POSITION_CHANNEL: u8 = 0;
const GENERIC_CHANNEL: u8 = 1;
const _RANGE_STREAM_REPORT: u8 = 0;
const _RANGE_STREAM_REPORT_FP16: u8 = 1;
const LPS_SHORT_LPP_PACKET: u8 = 2;
const _EMERGENCY_STOP: u8 = 3;
const _EMERGENCY_STOP_WATCHDOG: u8 = 4;
const _COMM_GNSS_NMEA: u8 = 6;
const _COMM_GNSS_PROPRIETARY: u8 = 7;
const EXT_POSE: u8 = 8;
const _EXT_POSE_PACKED: u8 = 9;
const LH_ANGLE_STREAM: u8 = 10;
const LH_PERSIST_DATA: u8 = 11;
#[derive(Debug, Clone)]
pub struct LighthouseAngleData {
pub base_station: u8,
pub angles: LighthouseAngles,
}
#[derive(Debug, Clone)]
pub struct LighthouseAngles {
pub x: [f32; 4],
pub y: [f32; 4],
}
pub struct Localization{
pub emergency: EmergencyControl,
pub external_pose: ExternalPose,
pub lighthouse: Lighthouse,
pub loco_positioning: LocoPositioning,
}
impl Localization {
pub(crate) fn new(uplink: Sender<Packet>, downlink: Receiver<Packet>) -> Self {
let emergency = EmergencyControl { uplink: uplink.clone() };
let external_pose = ExternalPose { uplink: uplink.clone() };
let (mut angle_broadcast, angle_receiver) = broadcast(100);
let (mut persist_broadcast, persist_receiver) = broadcast(10);
angle_broadcast.set_overflow(true);
persist_broadcast.set_overflow(true);
tokio::spawn(async move {
while let Ok(pk) = downlink.recv_async().await {
if pk.get_channel() != GENERIC_CHANNEL || pk.get_data().is_empty() {
continue;
}
let packet_type = pk.get_data()[0];
let data = &pk.get_data()[1..];
match packet_type {
LH_ANGLE_STREAM => {
if let Ok(angle_data) = decode_lh_angle(data) {
let _ = angle_broadcast.broadcast(angle_data).await;
}
}
LH_PERSIST_DATA if !data.is_empty() => {
let success = data[0] != 0;
let _ = persist_broadcast.broadcast(success).await;
}
_ => {} }
}
});
let lighthouse = Lighthouse {
uplink: uplink.clone(),
angle_stream_receiver: angle_receiver,
persist_receiver,
};
let loco_positioning = LocoPositioning { uplink: uplink.clone() };
Self { emergency, external_pose, lighthouse, loco_positioning }
}
}
fn decode_lh_angle(data: &[u8]) -> Result<LighthouseAngleData> {
if data.len() < 21 {
return Err(Error::ProtocolError("LH_ANGLE_STREAM packet too short".to_owned()));
}
let base_station = data[0];
let x0 = f32::from_le_bytes([data[1], data[2], data[3], data[4]]);
let x1_diff_i16 = i16::from_le_bytes([data[5], data[6]]);
let x2_diff_i16 = i16::from_le_bytes([data[7], data[8]]);
let x3_diff_i16 = i16::from_le_bytes([data[9], data[10]]);
let x1 = x0 - f16::from_bits(x1_diff_i16 as u16).to_f32();
let x2 = x0 - f16::from_bits(x2_diff_i16 as u16).to_f32();
let x3 = x0 - f16::from_bits(x3_diff_i16 as u16).to_f32();
let y0 = f32::from_le_bytes([data[11], data[12], data[13], data[14]]);
let y1_diff_i16 = i16::from_le_bytes([data[15], data[16]]);
let y2_diff_i16 = i16::from_le_bytes([data[17], data[18]]);
let y3_diff_i16 = i16::from_le_bytes([data[19], data[20]]);
let y1 = y0 - f16::from_bits(y1_diff_i16 as u16).to_f32();
let y2 = y0 - f16::from_bits(y2_diff_i16 as u16).to_f32();
let y3 = y0 - f16::from_bits(y3_diff_i16 as u16).to_f32();
Ok(LighthouseAngleData {
base_station,
angles: LighthouseAngles {
x: [x0, x1, x2, x3],
y: [y0, y1, y2, y3],
},
})
}
pub struct EmergencyControl {
uplink: Sender<Packet>,
}
impl EmergencyControl {
#[deprecated(since = "0.8.1", note = "Use [`Supervisor::send_emergency_stop`](crate::subsystems::supervisor::Supervisor::send_emergency_stop) instead")]
pub async fn send_emergency_stop(&self) -> Result<()> {
let pk = Packet::new(SUPERVISOR_PORT, SUPERVISOR_CH_COMMAND, vec![CMD_EMERGENCY_STOP]);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
#[deprecated(since = "0.8.1", note = "Use [`Supervisor::send_emergency_stop_watchdog`](crate::subsystems::supervisor::Supervisor::send_emergency_stop_watchdog) instead")]
pub async fn send_emergency_stop_watchdog(&self) -> Result<()> {
let pk = Packet::new(
SUPERVISOR_PORT,
SUPERVISOR_CH_COMMAND,
vec![CMD_EMERGENCY_STOP_WATCHDOG],
);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
}
pub struct ExternalPose {
uplink: Sender<Packet>,
}
impl ExternalPose {
pub async fn send_external_position(&self, pos: [f32; 3]) -> Result<()> {
let mut payload = Vec::with_capacity(3 * 4);
payload.extend_from_slice(&pos[0].to_le_bytes());
payload.extend_from_slice(&pos[1].to_le_bytes());
payload.extend_from_slice(&pos[2].to_le_bytes());
let pk = Packet::new(LOCALIZATION_PORT, POSITION_CHANNEL, payload);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
pub async fn send_external_pose(&self, pos: [f32; 3], quat: [f32; 4]) -> Result<()> {
let mut payload = Vec::with_capacity(1 + 7 * 4);
payload.push(EXT_POSE);
payload.extend_from_slice(&pos[0].to_le_bytes());
payload.extend_from_slice(&pos[1].to_le_bytes());
payload.extend_from_slice(&pos[2].to_le_bytes());
payload.extend_from_slice(&quat[0].to_le_bytes());
payload.extend_from_slice(&quat[1].to_le_bytes());
payload.extend_from_slice(&quat[2].to_le_bytes());
payload.extend_from_slice(&quat[3].to_le_bytes());
let pk = Packet::new(LOCALIZATION_PORT, GENERIC_CHANNEL, payload);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
}
pub struct LocoPositioning {
uplink: Sender<Packet>,
}
impl LocoPositioning {
pub async fn send_short_lpp_packet(&self, dest_id: u8, data: &[u8]) -> Result<()> {
let mut payload = Vec::with_capacity(2 + data.len());
payload.push(LPS_SHORT_LPP_PACKET);
payload.push(dest_id);
payload.extend_from_slice(data);
let pk = Packet::new(LOCALIZATION_PORT, GENERIC_CHANNEL, payload);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
}
pub struct Lighthouse {
uplink: Sender<Packet>,
angle_stream_receiver: BroadcastReceiver<LighthouseAngleData>,
persist_receiver: BroadcastReceiver<bool>,
}
impl Lighthouse {
pub async fn angle_stream(&self) -> impl Stream<Item = LighthouseAngleData> + use<> {
self.angle_stream_receiver.clone()
}
pub async fn persist_lighthouse_data(&self, geo_list: &[u8], calib_list: &[u8]) -> Result<bool> {
self.send_lh_persist_data_packet(geo_list, calib_list).await?;
self.wait_persist_confirmation().await
}
async fn wait_persist_confirmation(&self) -> Result<bool> {
let mut receiver = self.persist_receiver.clone();
match tokio::time::timeout(
std::time::Duration::from_secs(5),
receiver.recv()
).await {
Ok(Ok(success)) => Ok(success),
Ok(Err(_)) => Err(Error::Disconnected),
Err(_) => Err(Error::Timeout),
}
}
async fn send_lh_persist_data_packet(&self, geo_list: &[u8], calib_list: &[u8]) -> Result<()> {
const MAX_BS_NR: u8 = 15;
for &bs in geo_list {
if bs > MAX_BS_NR {
return Err(Error::ProtocolError(format!(
"Invalid geometry base station ID: {} (max: {})", bs, MAX_BS_NR
)));
}
}
for &bs in calib_list {
if bs > MAX_BS_NR {
return Err(Error::ProtocolError(format!(
"Invalid calibration base station ID: {} (max: {})", bs, MAX_BS_NR
)));
}
}
let mut mask_geo: u16 = 0;
let mut mask_calib: u16 = 0;
for &bs in geo_list {
mask_geo |= 1 << bs;
}
for &bs in calib_list {
mask_calib |= 1 << bs;
}
let mut payload = Vec::with_capacity(5);
payload.push(LH_PERSIST_DATA);
payload.extend_from_slice(&mask_geo.to_le_bytes());
payload.extend_from_slice(&mask_calib.to_le_bytes());
let pk = Packet::new(LOCALIZATION_PORT, GENERIC_CHANNEL, payload);
self.uplink.send_async(pk).await.map_err(|_| Error::Disconnected)?;
Ok(())
}
}