//! [`Recorder`] — the ride-time half of the crate. //! //! The recorder owns the raw journal. It is fed snapshots at whatever rate the //! ride engine ticks, throttles them to 1 Hz (FR-8.1), and appends each one as //! a complete line that is flushed immediately (FR-8.6). Nothing about the FIT //! file is decided until [`Recorder::finish`]. use std::fs::{File, OpenOptions}; use std::io::Write; use std::path::{Path, PathBuf}; use bikecontrol_core::RideSnapshot; use chrono::{DateTime, Local, Utc}; use crate::builder::{encode_activity, FitSummary}; use crate::rawlog::{entry_to_line, read_log, LogEntry, Sample, SessionStart, LOG_FORMAT_VERSION}; use crate::FitError; /// Tuning for [`Recorder`]. #[derive(Debug, Clone, PartialEq)] pub struct RecorderOptions { /// Minimum spacing between recorded samples, ms. Snapshots arriving sooner /// are dropped. The FIT `record` timestamp has one-second resolution, so /// there is nothing to gain from a faster journal. pub sample_interval_ms: u64, /// Force the journal to stable storage every N samples. `None` relies on /// the OS page cache, which is fast but loses the tail on a hard power cut. /// The default trades roughly ten seconds of exposure for one `fsync` per /// ten samples. pub fsync_every: Option, /// A silence longer than this is recorded as a BLE dropout (FR-8.5). /// `None` disables automatic gap detection; gaps can still be marked /// explicitly with [`Recorder::mark_gap`]. pub auto_gap_after_ms: Option, /// FIT `sub_sport`. `virtual_activity` makes Strava file the ride as a /// Virtual Ride; `indoor_cycling` is the alternative for a plain /// trainer session with no simulated course. pub sub_sport: u8, /// Written to `file_id.product_name` and `device_info.product_name`. pub product_name: String, /// Application version scaled by 100 — 1.20 is `120`. pub software_version: u16, /// Device serial number. Zero means unset. pub serial_number: u32, /// Rider mass in kilograms, recorded so the calorie estimate in the /// finished activity can include the resting term. Zero means unknown. pub rider_kg: f32, } impl Default for RecorderOptions { fn default() -> Self { Self { sample_interval_ms: 1000, fsync_every: Some(10), auto_gap_after_ms: Some(5_000), sub_sport: crate::profile::enums::SUB_SPORT_VIRTUAL_ACTIVITY, product_name: "BikeControl".to_string(), software_version: 100, serial_number: 0, rider_kg: 0.0, } } } /// Records a ride to a crash-safe journal and finalises it to a FIT activity. /// /// ```no_run /// # use bikecontrol_fit::{Recorder, RecorderOptions}; /// # use bikecontrol_core::RideSnapshot; /// # fn demo(snapshots: Vec) -> Result<(), Box> { /// let mut rec = Recorder::create("/tmp/ride.jsonl", RecorderOptions::default())?; /// for snap in &snapshots { /// rec.record(snap)?; // throttled to 1 Hz internally /// } /// rec.mark_lap(60_000, true)?; // controller pressed lap /// let summary = rec.finish("/tmp/ride.fit")?; /// println!("{} records, {:.1} km", summary.records, summary.total_distance_m / 1000.0); /// # Ok(()) /// # } /// ``` #[derive(Debug)] pub struct Recorder { file: File, log_path: PathBuf, opts: RecorderOptions, start: SessionStart, samples_written: usize, since_sync: usize, last_sample_ms: Option, last_elapsed_ms: u64, /// Elapsed time at which an open (unterminated) gap began. open_gap_at: Option, finished: bool, } impl Recorder { /// Start recording, creating the journal at `log_path`. /// /// The ride's wall-clock start is taken as "now", and the local UTC offset /// is captured with it so the activity shows the right time of day. pub fn create(log_path: impl AsRef, opts: RecorderOptions) -> Result { Self::create_at(log_path, opts, Utc::now(), local_utc_offset_secs()) } /// Start recording with an explicit wall-clock start and UTC offset. /// Used by tests, and by anything that needs a reproducible file. pub fn create_at( log_path: impl AsRef, opts: RecorderOptions, started_at: DateTime, utc_offset_secs: i32, ) -> Result { let log_path = log_path.as_ref().to_path_buf(); if let Some(parent) = log_path.parent() { if !parent.as_os_str().is_empty() { std::fs::create_dir_all(parent).map_err(|source| FitError::Io { path: parent.to_path_buf(), source, })?; } } // Truncate rather than append: a journal holds exactly one ride, and // silently concatenating two would produce a nonsense activity. let file = OpenOptions::new() .create(true) .write(true) .truncate(true) .open(&log_path) .map_err(|source| FitError::Io { path: log_path.clone(), source, })?; let start = SessionStart { start_unix_ms: started_at.timestamp_millis(), utc_offset_secs, sub_sport: opts.sub_sport, product_name: opts.product_name.clone(), software_version: opts.software_version, serial_number: opts.serial_number, rider_kg: opts.rider_kg, log_format: LOG_FORMAT_VERSION, }; let mut rec = Self { file, log_path, opts, start: start.clone(), samples_written: 0, since_sync: 0, last_sample_ms: None, last_elapsed_ms: 0, open_gap_at: None, finished: false, }; rec.append(&LogEntry::Start(start))?; rec.sync()?; Ok(rec) } /// Record a snapshot, subject to the 1 Hz throttle. /// /// Returns `true` if the sample was written, `false` if it was throttled /// away. Safe to call on every engine tick. pub fn record(&mut self, snapshot: &RideSnapshot) -> Result { self.record_sample(Sample::from_snapshot(snapshot)) } /// Record a fully-formed sample, subject to the same throttle. Use this /// when there is more to record than the snapshot carries — a virtual gear, /// or an absolute altitude from a loaded route. pub fn record_sample(&mut self, sample: Sample) -> Result { let t = sample.elapsed_ms; self.last_elapsed_ms = self.last_elapsed_ms.max(t); if let Some(prev) = self.last_sample_ms { if t < prev.saturating_add(self.opts.sample_interval_ms) { return Ok(false); } // Telemetry has been silent long enough to call it a dropout. if let Some(threshold) = self.opts.auto_gap_after_ms { if t.saturating_sub(prev) >= threshold && self.open_gap_at.is_none() { self.append(&LogEntry::Gap { at_ms: prev, until_ms: Some(t), reason: "no telemetry".to_string(), })?; } } } // An explicitly opened gap closes as soon as telemetry returns. if let Some(at_ms) = self.open_gap_at.take() { self.append(&LogEntry::Gap { at_ms, until_ms: Some(t), reason: "telemetry resumed".to_string(), })?; } self.last_sample_ms = Some(t); self.append(&LogEntry::Sample(sample))?; self.samples_written += 1; self.since_sync += 1; if self .opts .fsync_every .is_some_and(|n| n > 0 && self.since_sync >= n) { self.sync()?; } Ok(true) } /// Mark the start of a BLE dropout (FR-8.5). Recording continues; the gap /// is closed automatically by the next sample, or at the end of the ride. /// /// Calling this is optional — `auto_gap_after_ms` catches dropouts on its /// own — but a caller that *knows* the peripheral disconnected can record a /// reason and the exact moment. pub fn mark_gap(&mut self, at_ms: u64, reason: impl Into) -> Result<(), FitError> { if self.open_gap_at.is_some() { return Ok(()); // already inside a dropout } self.open_gap_at = Some(at_ms); self.last_elapsed_ms = self.last_elapsed_ms.max(at_ms); // Written now, unterminated, so it survives a crash during the dropout. self.append(&LogEntry::Gap { at_ms, until_ms: None, reason: reason.into(), })?; self.sync() } /// Mark a lap boundary (FR-8.7). `from_controller` distinguishes a Click /// button press from an on-screen tap. pub fn mark_lap(&mut self, at_ms: u64, from_controller: bool) -> Result<(), FitError> { self.last_elapsed_ms = self.last_elapsed_ms.max(at_ms); self.append(&LogEntry::Lap { at_ms, from_controller, })?; self.sync() } /// Pause the ride timer. Time until [`Recorder::resume`] counts towards /// elapsed time but not timer time. pub fn pause(&mut self, at_ms: u64) -> Result<(), FitError> { self.last_elapsed_ms = self.last_elapsed_ms.max(at_ms); self.append(&LogEntry::Pause { at_ms })?; self.sync() } /// Resume the ride timer. pub fn resume(&mut self, at_ms: u64) -> Result<(), FitError> { self.last_elapsed_ms = self.last_elapsed_ms.max(at_ms); self.append(&LogEntry::Resume { at_ms })?; self.sync() } /// Close the journal and write the FIT activity to `fit_path`. /// /// The FIT is built from the journal on disk, by the same code path crash /// recovery uses — so the file a rider gets after a clean ride and the file /// they get after a crash are produced identically. pub fn finish(mut self, fit_path: impl AsRef) -> Result { let end_ms = self.last_elapsed_ms; // Close any dropout that was still open when the ride ended. if let Some(at_ms) = self.open_gap_at.take() { self.append(&LogEntry::Gap { at_ms, until_ms: Some(end_ms), reason: "ride ended during dropout".to_string(), })?; } self.append(&LogEntry::End { at_ms: end_ms })?; self.sync()?; self.finished = true; let log_path = self.log_path.clone(); drop(self); crate::build_fit_from_log(&log_path, fit_path) } /// Close the journal without producing a FIT file. The journal remains on /// disk and can be turned into an activity later. pub fn abandon(mut self) -> PathBuf { self.finished = true; self.log_path.clone() } /// Where the journal is being written. pub fn log_path(&self) -> &Path { &self.log_path } /// How many samples have been committed to the journal. pub fn samples_written(&self) -> usize { self.samples_written } /// The session header written at the top of the journal. pub fn session_start(&self) -> &SessionStart { &self.start } /// Build a FIT from the journal *as it currently stands*, without ending /// the ride. Useful for a mid-ride preview or export, and the cheapest way /// to convince yourself the recording is sound before the ride ends. pub fn snapshot_fit(&mut self) -> Result<(Vec, FitSummary), FitError> { self.sync()?; let log = read_log(&self.log_path)?; encode_activity(&log) } fn append(&mut self, entry: &LogEntry) -> Result<(), FitError> { let line = entry_to_line(entry)?; // One `write_all` per entry: a torn write can only ever damage the // final line, which the reader drops. self.file .write_all(line.as_bytes()) .map_err(|source| FitError::Io { path: self.log_path.clone(), source, }) } fn sync(&mut self) -> Result<(), FitError> { self.since_sync = 0; self.file.flush().map_err(|source| FitError::Io { path: self.log_path.clone(), source, })?; self.file.sync_data().map_err(|source| FitError::Io { path: self.log_path.clone(), source, }) } } impl Drop for Recorder { fn drop(&mut self) { if !self.finished { // Best effort: get whatever is buffered onto disk. A ride // interrupted by a panic is still recoverable from the journal. let _ = self.file.flush(); let _ = self.file.sync_data(); } } } /// The machine's current UTC offset in seconds. fn local_utc_offset_secs() -> i32 { Local::now().offset().local_minus_utc() } #[cfg(test)] mod tests { use super::*; use crate::rawlog::parse_log; use bikecontrol_core::{ControlMode, Telemetry}; use chrono::TimeZone; fn tmpdir(name: &str) -> PathBuf { let dir = std::env::temp_dir().join(format!("bikecontrol-fit-{name}-{}", std::process::id())); std::fs::create_dir_all(&dir).unwrap(); dir } fn started_at() -> DateTime { Utc.with_ymd_and_hms(2026, 8, 5, 9, 0, 0).unwrap() } fn snapshot(elapsed_ms: u64) -> RideSnapshot { RideSnapshot { elapsed_ms, telemetry: Telemetry { elapsed_ms, power_w: Some(210), cadence_rpm: Some(88.0), ..Default::default() }, virtual_speed_kph: 32.4, virtual_distance_m: elapsed_ms as f64 * 0.009, gradient_pct: 1.5, elevation_gain_m: elapsed_ms as f32 * 0.000_135, gear: 6, gear_count: 12, development_m: 5.7, target_cadence_rpm: 94.7, speed_source: bikecontrol_core::types::SpeedSource::Drivetrain, pedal_force_n: 120.0, mode: ControlMode::ManualGrade, target: None, profile_progress: None, } } fn recorder(name: &str, opts: RecorderOptions) -> (Recorder, PathBuf) { let dir = tmpdir(name); let log = dir.join("ride.jsonl"); let rec = Recorder::create_at(&log, opts, started_at(), 7200).unwrap(); (rec, dir) } #[test] fn samples_are_throttled_to_one_hertz() { let (mut rec, dir) = recorder("throttle", RecorderOptions::default()); // 10 Hz input for 3 seconds. let mut accepted = 0; for i in 0..30u64 { if rec.record(&snapshot(i * 100)).unwrap() { accepted += 1; } } assert_eq!(accepted, 3, "0 ms, 1000 ms, 2000 ms"); assert_eq!(rec.samples_written(), 3); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn the_throttle_can_be_turned_off() { let (mut rec, dir) = recorder("nothrottle", RecorderOptions { sample_interval_ms: 0, auto_gap_after_ms: None, ..Default::default() }); for i in 0..10u64 { assert!(rec.record(&snapshot(i * 100)).unwrap()); } let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn every_sample_is_on_disk_before_the_call_returns() { // The crash-safety claim, tested directly: read the journal back with // the recorder still open and still holding the file. let (mut rec, dir) = recorder("durable", RecorderOptions::default()); for i in 0..5u64 { rec.record(&snapshot(i * 1000)).unwrap(); let text = std::fs::read_to_string(rec.log_path()).unwrap(); let log = parse_log(&text).unwrap(); assert_eq!( log.samples().count(), (i + 1) as usize, "sample {i} was not durable when record() returned" ); } let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn an_abandoned_journal_still_makes_a_fit() { // Simulates a crash: the process dies, nothing calls finish(), and the // journal is later handed to build_fit_from_log. let (mut rec, dir) = recorder("crash", RecorderOptions::default()); for i in 0..20u64 { rec.record(&snapshot(i * 1000)).unwrap(); } let log_path = rec.abandon(); let fit_path = dir.join("recovered.fit"); let summary = crate::build_fit_from_log(&log_path, &fit_path).unwrap(); assert_eq!(summary.records, 20); assert!(summary.recovered_from_crash); assert!(crate::encode::verify(&std::fs::read(&fit_path).unwrap()).is_ok()); let _ = std::fs::remove_dir_all(dir); } #[test] fn finish_writes_a_verifiable_fit_and_a_clean_journal() { let (mut rec, dir) = recorder("finish", RecorderOptions::default()); for i in 0..30u64 { rec.record(&snapshot(i * 1000)).unwrap(); } rec.mark_lap(10_000, true).unwrap(); let log_path = rec.log_path().to_path_buf(); let fit_path = dir.join("ride.fit"); let summary = rec.finish(&fit_path).unwrap(); assert_eq!(summary.records, 30); assert_eq!(summary.laps, 2); assert!(!summary.recovered_from_crash); assert_eq!(summary.skipped_log_lines, 0); let bytes = std::fs::read(&fit_path).unwrap(); assert_eq!(bytes.len(), summary.bytes); assert!(crate::encode::verify(&bytes).is_ok()); // The journal is still there and still describes the same ride. let log = crate::read_log(&log_path).unwrap(); assert!(log.clean_shutdown); assert_eq!(log.samples().count(), 30); let _ = std::fs::remove_dir_all(dir); } #[test] fn a_recovered_file_is_identical_to_the_clean_one() { // The strongest form of FR-8.4: recovery is not a degraded path, it is // the same path. let (mut rec, dir) = recorder("identical", RecorderOptions::default()); for i in 0..15u64 { rec.record(&snapshot(i * 1000)).unwrap(); } let log_path = rec.log_path().to_path_buf(); let clean = dir.join("clean.fit"); rec.finish(&clean).unwrap(); let rebuilt = dir.join("rebuilt.fit"); crate::build_fit_from_log(&log_path, &rebuilt).unwrap(); assert_eq!( std::fs::read(&clean).unwrap(), std::fs::read(&rebuilt).unwrap() ); let _ = std::fs::remove_dir_all(dir); } #[test] fn a_dropout_is_detected_automatically() { let (mut rec, dir) = recorder("autogap", RecorderOptions::default()); rec.record(&snapshot(0)).unwrap(); rec.record(&snapshot(1000)).unwrap(); // Ten seconds of silence, then telemetry returns. rec.record(&snapshot(11_000)).unwrap(); let log = parse_log(&std::fs::read_to_string(rec.log_path()).unwrap()).unwrap(); assert_eq!(log.gaps(11_000), vec![(1000, 11_000)]); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn an_explicit_gap_is_closed_when_telemetry_returns() { let (mut rec, dir) = recorder("explicitgap", RecorderOptions { auto_gap_after_ms: None, ..Default::default() }); rec.record(&snapshot(0)).unwrap(); rec.mark_gap(1000, "peripheral disconnected").unwrap(); // A second mark while already in a dropout is a no-op. rec.mark_gap(2000, "still gone").unwrap(); rec.record(&snapshot(9000)).unwrap(); let log = parse_log(&std::fs::read_to_string(rec.log_path()).unwrap()).unwrap(); let gaps = log.gaps(9000); // The unterminated marker written at 1000 ms, plus its closure. assert!(gaps.contains(&(1000, 9000))); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn a_dropout_open_at_the_end_of_the_ride_is_closed_by_finish() { let (mut rec, dir) = recorder("opengap", RecorderOptions { auto_gap_after_ms: None, ..Default::default() }); rec.record(&snapshot(0)).unwrap(); rec.record(&snapshot(5000)).unwrap(); rec.mark_gap(6000, "trainer lost").unwrap(); let fit = dir.join("ride.fit"); let summary = rec.finish(&fit).unwrap(); assert!(summary.gaps >= 1); assert!(crate::encode::verify(&std::fs::read(&fit).unwrap()).is_ok()); let _ = std::fs::remove_dir_all(dir); } #[test] fn pause_and_resume_are_journalled() { let (mut rec, dir) = recorder("pause", RecorderOptions::default()); for i in 0..5u64 { rec.record(&snapshot(i * 1000)).unwrap(); } rec.pause(5000).unwrap(); rec.resume(20_000).unwrap(); for i in 20..25u64 { rec.record(&snapshot(i * 1000)).unwrap(); } let fit = dir.join("ride.fit"); let summary = rec.finish(&fit).unwrap(); assert_eq!(summary.total_elapsed_s, 24.0); assert_eq!(summary.total_timer_s, 9.0, "24 s elapsed less 15 s paused"); let _ = std::fs::remove_dir_all(dir); } #[test] fn a_mid_ride_snapshot_is_a_valid_fit() { let (mut rec, dir) = recorder("midride", RecorderOptions::default()); for i in 0..8u64 { rec.record(&snapshot(i * 1000)).unwrap(); } let (bytes, summary) = rec.snapshot_fit().unwrap(); assert!(crate::encode::verify(&bytes).is_ok()); assert_eq!(summary.records, 8); assert!(summary.recovered_from_crash, "no end marker yet"); // Recording continues afterwards. assert!(rec.record(&snapshot(8000)).unwrap()); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn creating_a_recorder_makes_missing_directories() { let dir = tmpdir("mkdir").join("a").join("b"); let log = dir.join("ride.jsonl"); let rec = Recorder::create_at(&log, RecorderOptions::default(), started_at(), 0).unwrap(); assert!(log.exists()); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(tmpdir("mkdir")); } #[test] fn the_session_header_captures_the_start_and_offset() { let (rec, dir) = recorder("header", RecorderOptions::default()); let start = rec.session_start(); assert_eq!(start.start_unix_ms, started_at().timestamp_millis()); assert_eq!(start.utc_offset_secs, 7200); assert_eq!(start.log_format, LOG_FORMAT_VERSION); assert_eq!( start.sub_sport, crate::profile::enums::SUB_SPORT_VIRTUAL_ACTIVITY ); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } #[test] fn a_gear_can_be_recorded_alongside_the_snapshot() { let (mut rec, dir) = recorder("gear", RecorderOptions::default()); rec.record_sample(Sample::from_snapshot(&snapshot(0)).with_gear(11)) .unwrap(); let log = parse_log(&std::fs::read_to_string(rec.log_path()).unwrap()).unwrap(); let s = log.samples().next().unwrap(); assert_eq!(s.gear, Some(11)); assert_eq!(s.mode, Some(ControlMode::ManualGrade)); let _ = rec.abandon(); let _ = std::fs::remove_dir_all(dir); } }