use crate::command_message::*;
use crate::connection_config::*;
use crate::connection_status::*;
use crate::connections::*;
use crate::data_messages::*;
use crate::dispatcher::*;
use crate::ping_response::*;
use crate::receive_error::*;
use crate::statistics::*;
use std::ops::Drop;
use std::sync::{Arc, Mutex, MutexGuard};
pub const DEFAULT_RETRIES: u32 = 2;
pub const DEFAULT_TIMEOUT: u32 = 500;
pub(crate) type InternalConnection = Arc<Mutex<Box<dyn GenericConnection + Send>>>;
pub struct Connection {
dropped: Arc<Mutex<bool>>,
pub(crate) internal: InternalConnection,
}
impl Connection {
pub fn new(config: &ConnectionConfig) -> Self {
let connection = Self {
dropped: Arc::new(Mutex::new(false)),
internal: Arc::new(Mutex::new(match config {
ConnectionConfig::UsbConnectionConfig(config) => Box::new(UsbConnection::new(config)),
ConnectionConfig::SerialConnectionConfig(config) => Box::new(SerialConnection::new(config)),
ConnectionConfig::TcpConnectionConfig(config) => Box::new(TcpConnection::new(config)),
ConnectionConfig::UdpConnectionConfig(config) => Box::new(UdpConnection::new(config)),
ConnectionConfig::BluetoothConnectionConfig(config) => Box::new(BluetoothConnection::new(config)),
ConnectionConfig::FileConnectionConfig(config) => Box::new(FileConnection::new(config)),
ConnectionConfig::MuxConnectionConfig(config) => Box::new(MuxConnection::new(config)),
})),
};
let dropped = connection.dropped.clone();
let receiver = connection.internal.lock().unwrap().get_receiver();
let mut previous_instant = std::time::Instant::now();
let mut previous_statistics: Statistics = Default::default();
std::thread::spawn(move || loop {
std::thread::sleep(std::time::Duration::from_secs(1));
if let Ok(mut receiver) = receiver.lock() {
let delta_time = previous_instant.elapsed().as_secs_f32();
let delta_data = receiver.statistics.data_total - previous_statistics.data_total;
let delta_message = receiver.statistics.message_total - previous_statistics.message_total;
let delta_error = receiver.statistics.error_total - previous_statistics.error_total;
receiver.statistics.data_rate = (delta_data as f32 / delta_time).round() as u32;
receiver.statistics.message_rate = (delta_message as f32 / delta_time).round() as u32;
receiver.statistics.error_rate = (delta_error as f32 / delta_time).round() as u32;
previous_instant = std::time::Instant::now();
previous_statistics = receiver.statistics;
receiver.dispatcher.sender.send(DispatcherData::Statistics(receiver.statistics)).ok();
}
if *dropped.lock().unwrap() {
return;
}
});
connection
}
pub fn open(&self) -> std::io::Result<()> {
self.internal.lock().unwrap().open()
}
pub fn open_async(&self, closure: Box<dyn FnOnce(std::io::Result<()>) + Send>) {
let internal = self.internal.clone();
let dropped = self.dropped.clone();
std::thread::spawn(move || {
let result = internal.lock().unwrap().open();
if let Ok(dropped) = dropped.lock() {
if *dropped {
return;
}
closure(result);
}
});
}
pub fn close(&self) {
self.internal.lock().unwrap().close();
}
pub fn ping(&self) -> Option<PingResponse> {
Self::ping_internal(&self.internal)
}
pub fn ping_async(&self, closure: Box<dyn FnOnce(Option<PingResponse>) + Send>) {
let internal = self.internal.clone();
let dropped = self.dropped.clone();
std::thread::spawn(move || {
let response = Self::ping_internal(&internal);
if let Ok(dropped) = dropped.lock() {
if *dropped {
return;
}
closure(response);
}
});
}
pub(crate) fn ping_internal(internal: &InternalConnection) -> Option<PingResponse> {
let responses = Self::send_commands_internal(internal, vec!["{\"ping\":null}".into()], 4, 200);
PingResponse::parse(&responses.first()?.as_ref()?.json)
}
pub fn send_command(&self, command: Vec<u8>, retries: u32, timeout: u32) -> Option<CommandMessage> {
self.send_commands(vec![command], retries, timeout).first()?.clone()
}
pub fn send_commands(&self, commands: Vec<Vec<u8>>, retries: u32, timeout: u32) -> Vec<Option<CommandMessage>> {
Self::send_commands_internal(&self.internal, commands, retries, timeout)
}
pub fn send_command_async(&self, command: Vec<u8>, retries: u32, timeout: u32, closure: Box<dyn FnOnce(Option<CommandMessage>) + Send>) {
self.send_commands_async(vec![command], retries, timeout, Box::new(move |responses: Vec<Option<CommandMessage>>| closure(responses.first().cloned().flatten())));
}
pub fn send_commands_async(&self, commands: Vec<Vec<u8>>, retries: u32, timeout: u32, closure: Box<dyn FnOnce(Vec<Option<CommandMessage>>) + Send>) {
let internal = self.internal.clone();
let dropped = self.dropped.clone();
std::thread::spawn(move || {
let responses = Self::send_commands_internal(&internal, commands, retries, timeout);
if let Ok(dropped) = dropped.lock() {
if *dropped {
return;
}
closure(responses);
}
});
}
fn send_commands_internal(internal: &InternalConnection, commands: Vec<Vec<u8>>, retries: u32, timeout: u32) -> Vec<Option<CommandMessage>> {
let receiver = internal.lock().unwrap().get_receiver();
let Some(write_sender) = internal.lock().unwrap().get_write_sender() else {
return vec![None; commands.len()];
};
struct Transaction {
command: Option<CommandMessage>,
response: Option<CommandMessage>,
}
let mut transactions: Vec<Transaction> = commands
.iter()
.map(|command| Transaction {
command: CommandMessage::parse(command),
response: None,
})
.collect();
let (response_sender, response_receiver) = crossbeam::channel::unbounded();
let closure_id = receiver.lock().unwrap().dispatcher.add_command_closure(Box::new(move |command| {
response_sender.send(command).ok();
}));
'outer: for _ in 0..(1 + retries) {
for command in transactions.iter().filter_map(|transaction| transaction.command.clone()) {
write_sender.send([command.json.as_slice(), &[b'\n']].concat()).ok();
}
let start_time = std::time::Instant::now();
while start_time.elapsed() < std::time::Duration::from_millis(timeout as u64) {
if let Ok(response) = response_receiver.try_recv() {
for transaction in transactions.iter_mut() {
if let Some(command) = transaction.command.as_ref() {
if response.key == command.key {
*transaction = Transaction {
command: None,
response: Some(response.clone()),
};
}
}
}
}
if transactions.iter().all(|transaction| transaction.command.as_ref().is_none()) {
break 'outer;
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
}
receiver.lock().unwrap().dispatcher.remove_closure(closure_id);
transactions.iter().map(|transaction| transaction.response.clone()).collect()
}
pub fn get_config(&self) -> ConnectionConfig {
self.internal.lock().unwrap().get_config()
}
pub fn get_status(&self) -> ConnectionStatus {
self.internal.lock().unwrap().get_status()
}
pub fn get_statistics(&self) -> Statistics {
self.internal.lock().unwrap().get_receiver().lock().unwrap().statistics
}
pub fn get_inertial_message(&self, consume: bool) -> Option<InertialMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.inertial_message.lock().unwrap(), consume)
}
pub fn get_magnetometer_message(&self, consume: bool) -> Option<MagnetometerMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.magnetometer_message.lock().unwrap(), consume)
}
pub fn get_high_g_accelerometer_message(&self, consume: bool) -> Option<HighGAccelerometerMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.high_g_accelerometer_message.lock().unwrap(), consume)
}
pub fn get_quaternion_message(&self, consume: bool) -> Option<QuaternionMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.quaternion_message.lock().unwrap(), consume)
}
pub fn get_rotation_matrix_message(&self, consume: bool) -> Option<RotationMatrixMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.rotation_matrix_message.lock().unwrap(), consume)
}
pub fn get_euler_angles_message(&self, consume: bool) -> Option<EulerAnglesMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.euler_angles_message.lock().unwrap(), consume)
}
pub fn get_linear_acceleration_message(&self, consume: bool) -> Option<LinearAccelerationMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.linear_acceleration_message.lock().unwrap(), consume)
}
pub fn get_earth_acceleration_message(&self, consume: bool) -> Option<EarthAccelerationMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.earth_acceleration_message.lock().unwrap(), consume)
}
pub fn get_ahrs_status_message(&self, consume: bool) -> Option<AhrsStatusMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.ahrs_status_message.lock().unwrap(), consume)
}
pub fn get_serial_accessory_message(&self, consume: bool) -> Option<SerialAccessoryMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.serial_accessory_message.lock().unwrap(), consume)
}
pub fn get_sync_message(&self, consume: bool) -> Option<SyncMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.sync_message.lock().unwrap(), consume)
}
pub fn get_ltc_message(&self, consume: bool) -> Option<LtcMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.ltc_message.lock().unwrap(), consume)
}
pub fn get_temperature_message(&self, consume: bool) -> Option<TemperatureMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.temperature_message.lock().unwrap(), consume)
}
pub fn get_battery_message(&self, consume: bool) -> Option<BatteryMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.battery_message.lock().unwrap(), consume)
}
pub fn get_rssi_message(&self, consume: bool) -> Option<RssiMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.rssi_message.lock().unwrap(), consume)
}
pub fn get_button_message(&self, consume: bool) -> Option<ButtonMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.button_message.lock().unwrap(), consume)
}
pub fn get_notification_message(&self, consume: bool) -> Option<NotificationMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.notification_message.lock().unwrap(), consume)
}
pub fn get_error_message(&self, consume: bool) -> Option<ErrorMessage> {
Self::take_or_clone(self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.error_message.lock().unwrap(), consume)
}
fn take_or_clone<T: Clone>(mut guard: MutexGuard<Option<T>>, take: bool) -> Option<T> {
if take {
guard.take()
} else {
guard.clone()
}
}
pub fn add_receive_error_closure(&self, closure: Box<dyn Fn(ReceiveError) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_receive_error_closure(closure)
}
pub fn add_status_closure(&self, closure: Box<dyn Fn(ConnectionStatus) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_status_closure(closure)
}
pub fn add_statistics_closure(&self, closure: Box<dyn Fn(Statistics) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_statistics_closure(closure)
}
pub fn add_inertial_closure(&self, closure: Box<dyn Fn(InertialMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_inertial_closure(closure)
}
pub fn add_magnetometer_closure(&self, closure: Box<dyn Fn(MagnetometerMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_magnetometer_closure(closure)
}
pub fn add_high_g_accelerometer_closure(&self, closure: Box<dyn Fn(HighGAccelerometerMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_high_g_accelerometer_closure(closure)
}
pub fn add_quaternion_closure(&self, closure: Box<dyn Fn(QuaternionMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_quaternion_closure(closure)
}
pub fn add_rotation_matrix_closure(&self, closure: Box<dyn Fn(RotationMatrixMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_rotation_matrix_closure(closure)
}
pub fn add_euler_angles_closure(&self, closure: Box<dyn Fn(EulerAnglesMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_euler_angles_closure(closure)
}
pub fn add_linear_acceleration_closure(&self, closure: Box<dyn Fn(LinearAccelerationMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_linear_acceleration_closure(closure)
}
pub fn add_earth_acceleration_closure(&self, closure: Box<dyn Fn(EarthAccelerationMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_earth_acceleration_closure(closure)
}
pub fn add_ahrs_status_closure(&self, closure: Box<dyn Fn(AhrsStatusMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_ahrs_status_closure(closure)
}
pub fn add_serial_accessory_closure(&self, closure: Box<dyn Fn(SerialAccessoryMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_serial_accessory_closure(closure)
}
pub fn add_sync_closure(&self, closure: Box<dyn Fn(SyncMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_sync_closure(closure)
}
pub fn add_ltc_closure(&self, closure: Box<dyn Fn(LtcMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_ltc_closure(closure)
}
pub fn add_temperature_closure(&self, closure: Box<dyn Fn(TemperatureMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_temperature_closure(closure)
}
pub fn add_battery_closure(&self, closure: Box<dyn Fn(BatteryMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_battery_closure(closure)
}
pub fn add_rssi_closure(&self, closure: Box<dyn Fn(RssiMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_rssi_closure(closure)
}
pub fn add_button_closure(&self, closure: Box<dyn Fn(ButtonMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_button_closure(closure)
}
pub fn add_notification_closure(&self, closure: Box<dyn Fn(NotificationMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_notification_closure(closure)
}
pub fn add_error_closure(&self, closure: Box<dyn Fn(ErrorMessage) + Send>) -> u64 {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.add_error_closure(closure)
}
pub fn remove_closure(&self, id: u64) {
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.remove_closure(id);
}
}
impl Drop for Connection {
fn drop(&mut self) {
*self.dropped.lock().unwrap() = true;
self.internal.lock().unwrap().get_receiver().lock().unwrap().dispatcher.remove_all_closures();
self.close(); }
}