Files
bioz-host-rs/src/communication.rs

289 lines
14 KiB
Rust

use std::time::SystemTime;
use log::{error, info};
use tokio::select;
use tokio::sync::{watch, mpsc::{Receiver, Sender}};
use std::sync::atomic::{AtomicU32, Ordering};
use atomic_float::AtomicF32;
use std::sync::{Arc, Mutex};
use bioz_icd_rs::{SweepPoints, MultiplexerCapability};
use crate::state::{HardwareState, MeasurementDataState};
use crate::icd;
use crate::client::WorkbookClient;
use crate::logging::LoggingStates;
use crate::plot::{TimeSeriesPlot, BodePlot};
use crate::signals::{LoggingSignal, StartStopSignal};
use crate::state::HardwareConnected;
pub async fn communicate_with_hardware(
mut run_impedancemeter_rx: Receiver<StartStopSignal>,
run_impedancemeter_tx: Sender<StartStopSignal>,
measurement_data: Arc<MeasurementDataState>,
hardware_state_tx: watch::Sender<HardwareState>,
gui_logging_state: Arc<Mutex<LoggingStates>>,
log_tx: Sender<LoggingSignal>,
) {
let data_counter = Arc::new(AtomicU32::new(0));
let data_counter_clone = data_counter.clone();
let sampling_rate_clone = measurement_data.sampling_rate.clone();
tokio::spawn(async move {
loop {
tokio::time::sleep(tokio::time::Duration::from_secs(1)).await;
sampling_rate_clone.store(data_counter.load(Ordering::Relaxed) as f32, Ordering::Relaxed);
data_counter.store(0, Ordering::Relaxed);
}
});
#[derive(Default, Clone, Copy)]
struct Settings {
mode: Option<StartStopSignal>,
}
let settings = Arc::new(Mutex::new(Settings::default()));
loop {
let mut hardware_state = HardwareState::default();
let workbook_client = match WorkbookClient::new() {
Ok(client) => {
info!("Connected to hardware successfully.");
if let Some(mode) = settings.lock().unwrap().mode {
run_impedancemeter_tx.send(mode).await.unwrap();
}
match client.get_device_info().await.unwrap() {
MultiplexerCapability::Absent => {
hardware_state.connected = HardwareConnected::WithoutMultiplexer;
hardware_state_tx.send(hardware_state.clone()).unwrap();
info!("Connected device: Without Multiplexer");
},
MultiplexerCapability::Present => {
hardware_state.connected = HardwareConnected::WithMultiplexer;
hardware_state_tx.send(hardware_state.clone()).unwrap();
info!("Connected device: With Multiplexer");
},
}
client
},
Err(e) => {
error!("Failed to connect to hardware: {:?}", e);
tokio::time::sleep(tokio::time::Duration::from_secs(1)).await;
continue;
}
};
// Subscribe to SingleImpedanceOutputTopic
let mut single_impedance_sub = workbook_client
.client
.subscribe_multi::<icd::SingleImpedanceOutputTopic>(8)
.await
.unwrap();
let data = (measurement_data.magnitude_series.clone(), measurement_data.phase_series.clone(), measurement_data.magnitude.clone(), measurement_data.phase.clone());
let data_counter_clone_single = data_counter_clone.clone();
// Clone log_tx for the task
let gui_logging_state_clone = gui_logging_state.clone();
let log_tx_clone = log_tx.clone();
let settings_clone = settings.clone();
tokio::spawn(async move {
while let Ok(val) = single_impedance_sub.recv().await {
{
let mut mag_plot = data.0.lock().unwrap();
let mut phase_plot = data.1.lock().unwrap();
let mut mag_val = data.2.lock().unwrap();
let mut phase_val = data.3.lock().unwrap();
*mag_val = val.magnitude;
*phase_val = val.phase;
mag_plot.add(val.magnitude as f64);
phase_plot.add(val.phase as f64);
data_counter_clone_single.fetch_add(1, Ordering::Relaxed);
}
// Send logging signal
if *gui_logging_state_clone.lock().unwrap() == LoggingStates::Logging {
let settings = *settings_clone.lock().unwrap();
match settings.mode {
Some(StartStopSignal::StartSingle(freq, _, _, _)) => {
if let Err(e) = log_tx_clone.try_send(LoggingSignal::SingleImpedance(SystemTime::now(), freq, val.magnitude, val.phase)) {
error!("Failed to send logging signal: {:?}", e);
}
},
_ => {
error!("Frequency not set for single impedance logging.");
},
}
// if let Err(e) = log_tx_clone.send(LoggingSignal::SingleImpedance(SystemTime::now(), settings.frequency, val.magnitude, val.phase)).await {
// error!("Failed to send logging signal: {:?}", e);
// }
}
}
info!("SingleImpedanceOutputTopic subscription ended.");
});
// Subscribe to SweepImpedanceOutputTopic8
let mut sweep_impedance_sub = workbook_client
.client
.subscribe_multi::<icd::SweepImpedanceOutputTopic>(8)
.await
.unwrap();
let data = measurement_data.bode_plot.clone();
let data_counter_clone_sweep = data_counter_clone.clone();
// Clone log_tx for the task
let gui_logging_state_clone = gui_logging_state.clone();
let log_tx_clone = log_tx.clone();
tokio::spawn(async move {
while let Ok(val) = sweep_impedance_sub.recv().await {
match val.points {
SweepPoints::Eight => {
let magnitudes: Vec<f32> = val.magnitudes_8.into_iter().collect();
let phases: Vec<f32> = val.phases_8.into_iter().collect();
{
let mut bode_plot = data.lock().unwrap();
bode_plot.update_magnitudes(SweepPoints::Eight, magnitudes.clone());
bode_plot.update_phases(SweepPoints::Eight, phases.clone());
}
if *gui_logging_state_clone.lock().unwrap() == LoggingStates::Logging {
if let Err(e) = log_tx_clone.send(LoggingSignal::SweepImpedance(SystemTime::now(), SweepPoints::Eight.values().to_vec(), magnitudes.clone(), phases.clone())).await {
error!("Failed to send logging signal: {:?}", e);
}
}
},
SweepPoints::Eighteen => {
let magnitudes: Vec<f32> = val.magnitudes_18.into_iter().collect();
let phases: Vec<f32> = val.phases_18.into_iter().collect();
{
let mut bode_plot = data.lock().unwrap();
bode_plot.update_magnitudes(SweepPoints::Eighteen, magnitudes.clone());
bode_plot.update_phases(SweepPoints::Eighteen, phases.clone());
}
if *gui_logging_state_clone.lock().unwrap() == LoggingStates::Logging {
if let Err(e) = log_tx_clone.send(LoggingSignal::SweepImpedance(SystemTime::now(), SweepPoints::Eighteen.values().to_vec(), magnitudes.clone(), phases.clone())).await {
error!("Failed to send logging signal: {:?}", e);
}
}
},
}
data_counter_clone_sweep.fetch_add(1, Ordering::Relaxed);
}
});
let log_tx_clone = log_tx.clone();
loop {
select! {
Some(frequency) = run_impedancemeter_rx.recv() => {
match frequency {
StartStopSignal::StartSingle(freq, lead_mode, electrode_config, dft_num) => {
match workbook_client.start_impedancemeter_single(freq, lead_mode, electrode_config, dft_num).await {
Ok(Ok(periods)) => {
info!("Impedance meter started at frequency: {} with periods per DFT: {}", freq, periods);
settings.lock().unwrap().mode = Some(StartStopSignal::StartSingle(freq, lead_mode, electrode_config, dft_num));
hardware_state.periods_per_dft = Some(periods);
hardware_state_tx.send(hardware_state.clone()).unwrap();
// When logging add electrode configuration to logging file
if *gui_logging_state.lock().unwrap() == LoggingStates::Logging {
if let Err(e) = log_tx_clone.send(LoggingSignal::ElectrodeCongiguration(electrode_config)).await {
error!("Failed to send logging signal: {:?}", e);
}
}
},
Ok(Err(e)) => {
error!("Failed to init on hardware: {:?}", e);
hardware_state.periods_per_dft = None;
hardware_state_tx.send(hardware_state.clone()).unwrap();
},
Err(e) => {
error!("Communication error when starting impedancemeter: {:?}", e);
hardware_state.periods_per_dft = None;
hardware_state_tx.send(hardware_state.clone()).unwrap();
}
}
},
StartStopSignal::StartSweep(lead_mode, electrode_config, num_points) => {
match workbook_client.start_impedancemeter_sweep(lead_mode, electrode_config, num_points).await {
Ok(Ok(periods)) => {
settings.lock().unwrap().mode = Some(StartStopSignal::StartSweep(lead_mode, electrode_config, num_points));
info!("Sweep Impedancemeter started.");
match num_points {
SweepPoints::Eight => {
hardware_state.periods_per_dft_sweep = (num_points.values().iter().copied().collect(), Some(periods.periods_per_dft_8.into_iter().collect()));
hardware_state_tx.send(hardware_state.clone()).unwrap();
},
SweepPoints::Eighteen => {
hardware_state.periods_per_dft_sweep = (num_points.values().iter().copied().collect(), Some(periods.periods_per_dft_18.into_iter().collect()));
hardware_state_tx.send(hardware_state.clone()).unwrap();
},
}
// When logging add electrode configuration to logging file
if *gui_logging_state.lock().unwrap() == LoggingStates::Logging {
if let Err(e) = log_tx_clone.send(LoggingSignal::ElectrodeCongiguration(electrode_config)).await {
error!("Failed to send logging signal: {:?}", e);
}
}
},
Ok(Err(e)) => {
error!("Failed to sweep-init on hardware: {:?}", e);
hardware_state.periods_per_dft_sweep = (num_points.values().iter().copied().collect(), None);
hardware_state_tx.send(hardware_state.clone()).unwrap();
},
Err(e) => {
error!("Communication error when starting impedancemeter: {:?}", e);
hardware_state.periods_per_dft_sweep = (num_points.values().iter().copied().collect(), None);
hardware_state_tx.send(hardware_state.clone()).unwrap();
}
}
},
StartStopSignal::Stop => {
if let Err(e) = workbook_client.stop_impedancemeter().await {
error!("Failed to stop impedancemeter: {:?}", e);
} else {
settings.lock().unwrap().mode = Some(StartStopSignal::Stop);
hardware_state.periods_per_dft = None;
let (freq, _) = hardware_state.periods_per_dft_sweep.clone();
hardware_state.periods_per_dft_sweep = (freq, None);
hardware_state_tx.send(hardware_state.clone()).unwrap();
info!("Impedancemeter stopped.");
}
},
}
}
_ = workbook_client.wait_closed() => {
// Handle client closure
info!("Client connection closed.");
break;
}
else => {
// All channels closed
break;
}
}
}
info!("Communication with hardware ended.");
hardware_state.connected = HardwareConnected::None;
hardware_state_tx.send(hardware_state).unwrap();
tokio::time::sleep(tokio::time::Duration::from_secs(1)).await;
}
}