445 lines
15 KiB
Rust
445 lines
15 KiB
Rust
//! Recording and CSV export for pressure sensor data.
|
|
//!
|
|
//! Provides two recording modes (Full and Snapshot), state management
|
|
//! (Idle / Recording / Paused), and CSV import/export for replay.
|
|
|
|
use anyhow::{Context, Result};
|
|
use std::fs::File;
|
|
use std::io::{BufRead, BufReader, BufWriter, Write};
|
|
use std::path::Path;
|
|
use std::sync::{Arc, Mutex};
|
|
use std::time::Instant;
|
|
|
|
// ── Public types ────────────────────────────────────────────────────────
|
|
|
|
/// A single frame of pressure data captured at a point in time.
|
|
#[derive(Debug, Clone)]
|
|
pub struct Frame {
|
|
/// Pressure values, one per channel (same order every frame).
|
|
pub pressures: Vec<u32>,
|
|
/// Milliseconds elapsed since the recording started.
|
|
pub timestamp_ms: u64,
|
|
}
|
|
|
|
/// Whether the recording captures everything from connect (Full) or
|
|
/// only between explicit start/stop calls (Snapshot).
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum RecordingMode {
|
|
Full,
|
|
Snapshot,
|
|
}
|
|
|
|
/// Transient state of a recording.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum RecordingState {
|
|
Idle,
|
|
Recording,
|
|
Paused,
|
|
}
|
|
|
|
// ── Inner (unlocked) recorder ───────────────────────────────────────────
|
|
|
|
struct RecorderInner {
|
|
mode: RecordingMode,
|
|
state: RecordingState,
|
|
frames: Vec<Frame>,
|
|
start: Option<Instant>,
|
|
/// Accumulated wall-clock ms while paused (subtracted from elapsed).
|
|
paused_duration_ms: u64,
|
|
/// Instant when we entered Paused (None when not paused).
|
|
pause_start: Option<Instant>,
|
|
/// Number of channels in the first frame (used for CSV header).
|
|
channel_count: Option<usize>,
|
|
}
|
|
|
|
impl RecorderInner {
|
|
fn new(mode: RecordingMode) -> Self {
|
|
Self {
|
|
mode,
|
|
state: RecordingState::Idle,
|
|
frames: Vec::new(),
|
|
start: None,
|
|
paused_duration_ms: 0,
|
|
pause_start: None,
|
|
channel_count: None,
|
|
}
|
|
}
|
|
|
|
fn elapsed_ms(&self) -> u64 {
|
|
let Some(start) = self.start else { return 0 };
|
|
let raw = start.elapsed().as_millis() as u64;
|
|
let paused = if let Some(ps) = self.pause_start {
|
|
self.paused_duration_ms + ps.elapsed().as_millis() as u64
|
|
} else {
|
|
self.paused_duration_ms
|
|
};
|
|
raw.saturating_sub(paused)
|
|
}
|
|
|
|
fn push_frame(&mut self, pressures: Vec<u32>) {
|
|
if self.channel_count.is_none() {
|
|
self.channel_count = Some(pressures.len());
|
|
}
|
|
let ts = self.elapsed_ms();
|
|
self.frames.push(Frame {
|
|
pressures,
|
|
timestamp_ms: ts,
|
|
});
|
|
}
|
|
}
|
|
|
|
// ── Thread-safe wrapper ─────────────────────────────────────────────────
|
|
|
|
/// Thread-safe recorder. Clone the `Arc` to share across threads.
|
|
#[derive(Clone)]
|
|
pub struct Recorder {
|
|
inner: Arc<Mutex<RecorderInner>>,
|
|
}
|
|
|
|
impl Recorder {
|
|
/// Create a recorder in the given mode (starts in `Idle` state).
|
|
pub fn new(mode: RecordingMode) -> Self {
|
|
Self {
|
|
inner: Arc::new(Mutex::new(RecorderInner::new(mode))),
|
|
}
|
|
}
|
|
|
|
// ── Lifecycle ───────────────────────────────────────────────────
|
|
|
|
/// Start a **Full** recording (records every frame pushed from now on).
|
|
pub fn start_full_recording(&self) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(
|
|
r.state == RecordingState::Idle,
|
|
"can only start from Idle state"
|
|
);
|
|
r.state = RecordingState::Recording;
|
|
r.start = Some(Instant::now());
|
|
r.paused_duration_ms = 0;
|
|
r.pause_start = None;
|
|
r.frames.clear();
|
|
r.channel_count = None;
|
|
Ok(())
|
|
}
|
|
|
|
/// Start a **Snapshot** recording window.
|
|
pub fn start_snapshot_recording(&self) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(
|
|
r.state == RecordingState::Idle,
|
|
"can only start from Idle state"
|
|
);
|
|
r.state = RecordingState::Recording;
|
|
r.start = Some(Instant::now());
|
|
r.paused_duration_ms = 0;
|
|
r.pause_start = None;
|
|
r.frames.clear();
|
|
r.channel_count = None;
|
|
Ok(())
|
|
}
|
|
|
|
/// Stop recording (transitions to `Idle`).
|
|
pub fn stop_recording(&self) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(
|
|
r.state == RecordingState::Recording || r.state == RecordingState::Paused,
|
|
"nothing to stop"
|
|
);
|
|
if r.state == RecordingState::Paused {
|
|
if let Some(ps) = r.pause_start.take() {
|
|
r.paused_duration_ms += ps.elapsed().as_millis() as u64;
|
|
}
|
|
}
|
|
r.state = RecordingState::Idle;
|
|
Ok(())
|
|
}
|
|
|
|
/// Pause an active recording.
|
|
pub fn pause_recording(&self) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(r.state == RecordingState::Recording, "not recording");
|
|
r.pause_start = Some(Instant::now());
|
|
r.state = RecordingState::Paused;
|
|
Ok(())
|
|
}
|
|
|
|
/// Resume a paused recording.
|
|
pub fn resume_recording(&self) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(r.state == RecordingState::Paused, "not paused");
|
|
if let Some(ps) = r.pause_start.take() {
|
|
r.paused_duration_ms += ps.elapsed().as_millis() as u64;
|
|
}
|
|
r.state = RecordingState::Recording;
|
|
Ok(())
|
|
}
|
|
|
|
// ── Data ────────────────────────────────────────────────────────
|
|
|
|
/// Feed one frame of pressure data. Ignored when not recording.
|
|
pub fn add_frame(&self, pressures: &[u32]) {
|
|
let mut r = self.inner.lock().unwrap();
|
|
if r.state == RecordingState::Recording {
|
|
r.push_frame(pressures.to_vec());
|
|
}
|
|
}
|
|
|
|
/// Number of recorded frames.
|
|
pub fn frame_count(&self) -> usize {
|
|
self.inner.lock().unwrap().frames.len()
|
|
}
|
|
|
|
/// Whether the recorder is currently in the `Recording` state.
|
|
pub fn is_recording(&self) -> bool {
|
|
self.inner.lock().unwrap().state == RecordingState::Recording
|
|
}
|
|
|
|
/// Current recording state.
|
|
pub fn state(&self) -> RecordingState {
|
|
self.inner.lock().unwrap().state
|
|
}
|
|
|
|
/// Milliseconds elapsed (excludes paused time).
|
|
pub fn duration_ms(&self) -> u64 {
|
|
self.inner.lock().unwrap().elapsed_ms()
|
|
}
|
|
|
|
/// Snapshot of all recorded frames.
|
|
pub fn recorded_frames(&self) -> Vec<Frame> {
|
|
self.inner.lock().unwrap().frames.clone()
|
|
}
|
|
|
|
// ── CSV export / import ─────────────────────────────────────────
|
|
|
|
/// Export recorded frames to CSV.
|
|
///
|
|
/// Header: `channel1,channel2,...,channelN,timestamp_ms`
|
|
pub fn export_csv<P: AsRef<Path>>(&self, path: P) -> Result<()> {
|
|
let r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(!r.frames.is_empty(), "no frames to export");
|
|
|
|
let file = File::create(path.as_ref())
|
|
.with_context(|| format!("creating {}", path.as_ref().display()))?;
|
|
let mut w = BufWriter::new(file);
|
|
|
|
let n = r.channel_count.unwrap_or(r.frames[0].pressures.len());
|
|
|
|
// Header
|
|
for i in 0..n {
|
|
write!(w, "channel{}", i + 1)?;
|
|
if i < n - 1 {
|
|
write!(w, ",")?;
|
|
}
|
|
}
|
|
writeln!(w, ",timestamp_ms")?;
|
|
|
|
// Rows
|
|
for frame in &r.frames {
|
|
for (i, val) in frame.pressures_iter().enumerate() {
|
|
write!(w, "{}", val)?;
|
|
if i < n - 1 {
|
|
write!(w, ",")?;
|
|
}
|
|
}
|
|
writeln!(w, ",{}", frame.timestamp_ms)?;
|
|
}
|
|
|
|
w.flush()?;
|
|
Ok(())
|
|
}
|
|
|
|
/// Import frames from a CSV file (same format as `export_csv`).
|
|
///
|
|
/// The recorder must be in `Idle` state. After import it stays `Idle`
|
|
/// so you can inspect / re-export; call `start_*` to continue recording.
|
|
pub fn import_csv<P: AsRef<Path>>(&self, path: P) -> Result<()> {
|
|
let mut r = self.inner.lock().unwrap();
|
|
anyhow::ensure!(r.state == RecordingState::Idle, "must be Idle to import");
|
|
|
|
let file = File::open(path.as_ref())
|
|
.with_context(|| format!("opening {}", path.as_ref().display()))?;
|
|
let reader = BufReader::new(file);
|
|
let mut lines = reader.lines();
|
|
|
|
// Parse header to learn channel count
|
|
let header = lines
|
|
.next()
|
|
.context("empty CSV file")?
|
|
.context("reading header")?;
|
|
let cols: Vec<&str> = header.split(',').map(str::trim).collect();
|
|
anyhow::ensure!(cols.len() >= 2, "need at least 1 channel + timestamp_ms");
|
|
let n_channels = cols.len() - 1; // last column is timestamp_ms
|
|
|
|
r.frames.clear();
|
|
r.channel_count = Some(n_channels);
|
|
|
|
let mut first_ts: Option<u64> = None;
|
|
|
|
for (lineno, line) in lines.enumerate() {
|
|
let line = line.with_context(|| format!("reading line {}", lineno + 2))?;
|
|
let line = line.trim();
|
|
if line.is_empty() {
|
|
continue;
|
|
}
|
|
let fields: Vec<&str> = line.split(',').map(str::trim).collect();
|
|
anyhow::ensure!(
|
|
fields.len() == cols.len(),
|
|
"line {}: expected {} columns, got {}",
|
|
lineno + 2,
|
|
cols.len(),
|
|
fields.len()
|
|
);
|
|
|
|
let mut pressures = Vec::with_capacity(n_channels);
|
|
for f in &fields[..n_channels] {
|
|
pressures.push(
|
|
f.parse::<u32>().with_context(|| {
|
|
format!("line {}: bad channel value '{}'", lineno + 2, f)
|
|
})?,
|
|
);
|
|
}
|
|
let raw_ts = fields[n_channels]
|
|
.parse::<u64>()
|
|
.with_context(|| format!("line {}: bad timestamp", lineno + 2))?;
|
|
|
|
// Normalise timestamps so the first frame starts at 0
|
|
let ts = if let Some(first) = first_ts {
|
|
raw_ts.saturating_sub(first)
|
|
} else {
|
|
first_ts = Some(raw_ts);
|
|
0
|
|
};
|
|
|
|
r.frames.push(Frame {
|
|
pressures,
|
|
timestamp_ms: ts,
|
|
});
|
|
}
|
|
|
|
// Set the start instant so duration_ms() reports the imported span
|
|
if !r.frames.is_empty() {
|
|
if let Some(last) = r.frames.last() {
|
|
// Pretend the recording happened `last.timestamp_ms` ago
|
|
// so that elapsed_ms() would return that value.
|
|
// We store a "fake" start by noting the offset.
|
|
r.start =
|
|
Some(Instant::now() - std::time::Duration::from_millis(last.timestamp_ms));
|
|
r.paused_duration_ms = 0;
|
|
}
|
|
}
|
|
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
// ── Convenience constructors ────────────────────────────────────────────
|
|
|
|
impl Recorder {
|
|
/// Shorthand: new Full recorder.
|
|
pub fn full() -> Self {
|
|
Self::new(RecordingMode::Full)
|
|
}
|
|
|
|
/// Shorthand: new Snapshot recorder.
|
|
pub fn snapshot() -> Self {
|
|
Self::new(RecordingMode::Snapshot)
|
|
}
|
|
}
|
|
|
|
// ── Helper extension on Frame ───────────────────────────────────────────
|
|
|
|
impl Frame {
|
|
fn pressures_iter(&self) -> impl Iterator<Item = u32> + '_ {
|
|
self.pressures.iter().copied()
|
|
}
|
|
}
|
|
|
|
// ── Tests ───────────────────────────────────────────────────────────────
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use std::thread;
|
|
|
|
#[test]
|
|
fn full_recording_lifecycle() {
|
|
let rec = Recorder::full();
|
|
assert_eq!(rec.state(), RecordingState::Idle);
|
|
assert!(!rec.is_recording());
|
|
assert_eq!(rec.frame_count(), 0);
|
|
|
|
rec.start_full_recording().unwrap();
|
|
assert!(rec.is_recording());
|
|
|
|
rec.add_frame(&[10, 20, 30]);
|
|
rec.add_frame(&[40, 50, 60]);
|
|
assert_eq!(rec.frame_count(), 2);
|
|
|
|
rec.stop_recording().unwrap();
|
|
assert_eq!(rec.state(), RecordingState::Idle);
|
|
assert!(!rec.is_recording());
|
|
assert!(rec.duration_ms() < 500);
|
|
}
|
|
|
|
#[test]
|
|
fn snapshot_pause_resume() {
|
|
let rec = Recorder::snapshot();
|
|
rec.start_snapshot_recording().unwrap();
|
|
|
|
rec.add_frame(&[1, 2]);
|
|
rec.pause_recording().unwrap();
|
|
assert_eq!(rec.state(), RecordingState::Paused);
|
|
rec.add_frame(&[9, 9]); // should be ignored
|
|
assert_eq!(rec.frame_count(), 1);
|
|
|
|
rec.resume_recording().unwrap();
|
|
rec.add_frame(&[3, 4]);
|
|
assert_eq!(rec.frame_count(), 2);
|
|
|
|
rec.stop_recording().unwrap();
|
|
}
|
|
|
|
#[test]
|
|
fn csv_round_trip() {
|
|
let rec = Recorder::full();
|
|
rec.start_full_recording().unwrap();
|
|
rec.add_frame(&[100, 200, 300]);
|
|
rec.add_frame(&[101, 201, 301]);
|
|
rec.stop_recording().unwrap();
|
|
|
|
let dir = std::env::temp_dir().join("eskin_recording_test");
|
|
std::fs::create_dir_all(&dir).unwrap();
|
|
let csv_path = dir.join("test_roundtrip.csv");
|
|
rec.export_csv(&csv_path).unwrap();
|
|
|
|
let rec2 = Recorder::snapshot();
|
|
rec2.import_csv(&csv_path).unwrap();
|
|
assert_eq!(rec2.frame_count(), 2);
|
|
|
|
let frames = rec2.recorded_frames();
|
|
assert_eq!(frames[0].pressures, vec![100, 200, 300]);
|
|
assert_eq!(frames[1].pressures, vec![101, 201, 301]);
|
|
assert_eq!(frames[0].timestamp_ms, 0);
|
|
|
|
let _ = std::fs::remove_dir_all(&dir);
|
|
}
|
|
|
|
#[test]
|
|
fn thread_safety() {
|
|
let rec = Recorder::full();
|
|
let rec2 = rec.clone();
|
|
|
|
rec.start_full_recording().unwrap();
|
|
|
|
let h = thread::spawn(move || {
|
|
for i in 0..100u32 {
|
|
rec2.add_frame(&[i, i + 1, i + 2]);
|
|
}
|
|
});
|
|
h.join().unwrap();
|
|
|
|
rec.stop_recording().unwrap();
|
|
assert_eq!(rec.frame_count(), 100);
|
|
}
|
|
}
|