use crate::command_message::*;
use crate::data_messages::*;
use crate::dispatcher::*;
use crate::mux_message::*;
use crate::receive_error::*;
use crate::statistics::*;
const BUFFER_SIZE: usize = 4096;
pub struct Receiver {
buffer: [u8; BUFFER_SIZE],
index: usize,
pub statistics: Statistics,
pub dispatcher: Dispatcher,
}
impl Receiver {
pub fn new() -> Receiver {
Self {
buffer: [0; BUFFER_SIZE],
index: 0,
statistics: Default::default(),
dispatcher: Dispatcher::new(),
}
}
pub fn receive_bytes(&mut self, bytes: &[u8]) {
self.statistics.data_total += bytes.len() as u64;
for byte in bytes {
self.buffer[self.index] = *byte;
self.index += 1;
if self.index >= self.buffer.len() {
self.statistics.error_total += 1;
self.dispatcher.sender.send(DispatcherData::ReceiveError(ReceiveError::BufferOverrun)).ok();
self.index = 0;
continue;
}
if *byte == b'\n' {
match self.receive_message() {
Ok(_) => self.statistics.message_total += 1,
Err(receive_error) => {
self.statistics.error_total += 1;
self.dispatcher.sender.send(DispatcherData::ReceiveError(receive_error)).ok();
}
}
self.index = 0;
}
}
}
fn receive_message(&mut self) -> Result<(), ReceiveError> {
match self.buffer[0] {
b'{' => self.receive_command_message(),
b'^' => self.receive_mux_message(),
_ => self.receive_data_message(),
}
}
fn receive_command_message(&self) -> Result<(), ReceiveError> {
let response = CommandMessage::parse(&self.buffer[..self.index]).ok_or(ReceiveError::InvalidJson)?;
self.dispatcher.sender.send(DispatcherData::Command(response)).ok();
Ok(())
}
fn receive_mux_message(&self) -> Result<(), ReceiveError> {
let message = MuxMessage::parse(&self.buffer[..self.index])?;
self.dispatcher.sender.send(DispatcherData::Mux(message)).ok();
Ok(())
}
fn receive_data_message(&mut self) -> Result<(), ReceiveError> {
let message = Self::undo_byte_stuffing(&mut self.buffer[..self.index])?;
macro_rules! parse {
($data_message:ident, $dispatcher_data:ident) => {{
match $data_message::parse(message) {
Ok(message) => {
self.dispatcher.sender.send(DispatcherData::$dispatcher_data(message)).ok();
return Ok(());
}
Err(ReceiveError::InvalidMessageIdentifier) => {}
Err(error) => return Err(error),
}
}};
}
parse!(InertialMessage, Inertial);
parse!(MagnetometerMessage, Magnetometer);
parse!(HighGAccelerometerMessage, HighGAccelerometer);
parse!(QuaternionMessage, Quaternion);
parse!(RotationMatrixMessage, RotationMatrix);
parse!(EulerAnglesMessage, EulerAngles);
parse!(LinearAccelerationMessage, LinearAcceleration);
parse!(EarthAccelerationMessage, EarthAcceleration);
parse!(AhrsStatusMessage, AhrsStatus);
parse!(SerialAccessoryMessage, SerialAccessory);
parse!(SyncMessage, Sync);
parse!(LtcMessage, Ltc);
parse!(TemperatureMessage, Temperature);
parse!(BatteryMessage, Battery);
parse!(RssiMessage, Rssi);
parse!(ButtonMessage, Button);
parse!(NotificationMessage, Notification);
parse!(ErrorMessage, Error);
Err(ReceiveError::InvalidMessageIdentifier)
}
fn undo_byte_stuffing(message: &mut [u8]) -> Result<&[u8], ReceiveError> {
const BYTE_STUFFING_END: u8 = 0x0A;
const BYTE_STUFFING_ESC: u8 = 0xDB;
const BYTE_STUFFING_ESC_END: u8 = 0xDC;
const BYTE_STUFFING_ESC_ESC: u8 = 0xDD;
let mut source_index = 0;
let mut destination_index = 0;
while source_index < message.len() {
if message[source_index] == BYTE_STUFFING_ESC {
source_index += 1;
if source_index >= message.len() {
return Err(ReceiveError::InvalidEscapeSequence);
}
match message[source_index] {
BYTE_STUFFING_ESC_END => message[destination_index] = BYTE_STUFFING_END,
BYTE_STUFFING_ESC_ESC => message[destination_index] = BYTE_STUFFING_ESC,
_ => return Err(ReceiveError::InvalidEscapeSequence),
}
} else {
message[destination_index] = message[source_index];
}
source_index += 1;
destination_index += 1;
}
Ok(&message[..destination_index])
}
}