Core ride logic, FTMS client, FIT encoder and probe CLI

Adds backing state for Resistance and Erg control modes, which had no
value to hold and so could never satisfy FR-4.3/FR-4.6.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-05 13:34:27 +02:00
co-authored by Claude Opus 5
parent 3e106de2c5
commit 7c17ca6158
61 changed files with 20933 additions and 55 deletions
+634
View File
@@ -0,0 +1,634 @@
//! [`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<usize>,
/// 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<u64>,
/// 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,
}
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,
}
}
}
/// 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<RideSnapshot>) -> Result<(), Box<dyn std::error::Error>> {
/// 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<u64>,
last_elapsed_ms: u64,
/// Elapsed time at which an open (unterminated) gap began.
open_gap_at: Option<u64>,
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<Path>, opts: RecorderOptions) -> Result<Self, FitError> {
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<Path>,
opts: RecorderOptions,
started_at: DateTime<Utc>,
utc_offset_secs: i32,
) -> Result<Self, FitError> {
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,
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<bool, FitError> {
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<bool, FitError> {
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<String>) -> 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<Path>) -> Result<FitSummary, FitError> {
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<u8>, 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> {
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,
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);
}
}