diff --git a/docs/port/PROGRESS.md b/docs/port/PROGRESS.md index 3ef4248..0f98990 100644 --- a/docs/port/PROGRESS.md +++ b/docs/port/PROGRESS.md @@ -368,3 +368,111 @@ statistics and export layers above it are proven equivalent to the Kotlin origin Next: **Phase 3, the recording engine** — the risky phase. T10's `LocationSource` seam first, then the pipeline, then the two platform liveness stories. The `flutter_foreground_task` deprecation warnings recorded under T01 become relevant at T12. + +--- + +## T10 — Location source seam · **complete** + +`lib/src/recording/location_source.dart`: `LocationFix`, `LocationSource`, +`LocationException`, and `FakeLocationSource`. + +`LocationFix` is deliberately **not** `TrackPoint` — a fix has no trip or segment +identity. Those ids are stamped on by the engine at creation, which is the whole basis of +the pause guarantee. + +### `geolocator` can replace `flutter_foreground_task` entirely + +`geolocator_android` ships `ForegroundNotificationConfig`, which raises a foreground +service with `foregroundServiceType=location` for as long as the position stream is +active, and exposes `enableWakeLock` and `setOngoing`. That is **everything** +`TrackingService` used a foreground service and a `PARTIAL_WAKE_LOCK` for. + +Dropping `flutter_foreground_task` would remove both deprecation paths recorded under +T01 (no Swift Package Manager on iOS; applies KGP on Android) at no cost to the design. + +**One parity casualty:** geolocator's notification config has no support for *actions*, +so the notification would be display-only — the native app's Pause/Resume buttons in the +shade would be lost. Decision deferred to T12, where it actually bites. The seam means +neither choice touches the engine. + +--- + +## T11 — Recording pipeline · **complete** + +`lib/src/recording/recording_engine.dart`. **24 tests.** + +The Kotlin `Channel(UNLIMITED)` + blocking-`receive` writer coroutine becomes a plain +`List` buffer plus a periodic `Timer`. Dart has no blocking receive and its single +threaded event loop makes one unnecessary — the guarantee is unchanged: the fix callback +only appends and returns, so disk latency can never stall GPS. `Mutex` becomes +`synchronized`'s `Lock`, serialising the periodic flush against explicit drains at pause, +stop and discard. + +All five actions ported: start (adopting), pause, resume (via start), stop, discard, plus +`restoreAfterProcessDeath`. + +--- + +## 🐞 A real bug found in the native app + +**`TrackingService.restoreAfterProcessDeath` does not do what its comment says.** + +```kotlin +// Resume into a *new* segment: the time the process was dead is a real +// gap in the recording and should render as one. +trips.resumeTrip(System.currentTimeMillis())?.let { handle -> +``` + +`resumeTrip` calls `adoptOrOpenSegment`, which returns the **existing open segment** if +there is one. After a *pause* that is correct, because pausing closes the segment first. +After a *crash* nothing closed it — so the adopt branch wins and the stated intent is +silently not met. + +### The consequence is data corruption, not cosmetics + +Points either side of the dead time land in one segment. `RideStatistics.compute` groups +by segment, and it is the **authoritative** pass that overwrites the live estimate when +the trip completes. So the gap gets measured as if it had been ridden. + +Measured, by reverting the fix and running the guard: + +``` +the dead time leaked into distance: 111217.31924957958 m +``` + +**111 km of phantom distance** added to a ride because the process died and the rider +relaunched somewhere else. The map would also draw a straight line across roads never +ridden — exactly the artefact segments exist to prevent. + +The live accumulator gets this right (it is re-seeded with a null anchor). The +authoritative recomputation then overwrites the correct figure with the wrong one. + +### The fix + +`TripRepository.resumeIntoNewSegment` closes the stale segment and opens a fresh one. +The stale segment is closed **at its last recorded point**, not at `now` — recording +genuinely stopped when the process died, and `computeSummary` sums closed segment spans +for elapsed time, so closing at `now` would bill the dead time as ride time. + +This is a deliberate, documented **departure from parity**. The plan says port bugs +faithfully, and that holds for the elevation drift where a "fix" would make differential +testing ambiguous. It does not hold here: the code contradicts its own stated intent and +the result is silently wrong data. + +> Worth carrying back to the native app if it is ever revived. Second finding of its kind +> after the seed-lucky elevation bound. + +### And a near-miss worth recording + +The first version of the guard **passed with the bug still present**. The fixture called +`pause()` to flush points to disk — but pausing *closes* the segment, so the crash state +was never reproduced. It was only caught by deliberately reverting the fix and checking +the test failed. + +The rewritten fixture writes points through the repository directly, leaving the segment +open exactly as a crash does. It now fails at 111 km with the native behaviour and passes +with the fix. + +**A regression test nobody has watched fail is not yet a regression test.** + +**145 tests passing, analyze clean.** diff --git a/lib/src/data/trip_repository.dart b/lib/src/data/trip_repository.dart index df8c9f5..3706057 100644 --- a/lib/src/data/trip_repository.dart +++ b/lib/src/data/trip_repository.dart @@ -108,6 +108,45 @@ class TripRepository { return handle; }); + /// Resumes an active trip into a **genuinely new** segment, closing any segment left + /// open by a process death. + /// + /// ## This fixes a bug in the native app + /// + /// `TrackingService.restoreAfterProcessDeath` says *"Resume into a new segment: the + /// time the process was dead is a real gap in the recording and should render as one"* + /// — but it calls `resumeTrip`, which adopts the open segment rather than opening a + /// new one. After a **pause** that is right, because pausing closes the segment first. + /// After a **crash** nothing closed it, so the adopt branch wins and the intent is + /// silently not met. + /// + /// The consequence is not cosmetic. Points either side of the dead time end up in one + /// segment, so `computeSummary` — the authoritative pass that overwrites the live + /// estimate when the trip completes — measures straight through the gap. A rider who + /// crashes in town and relaunches an hour later downtown gets those kilometres added to + /// their distance, and a polyline drawn across roads they never rode. The live + /// accumulator gets this right (it is re-seeded with a null anchor); the authoritative + /// recomputation then overwrites it with the wrong figure. + /// + /// The stale segment is closed at **its last recorded point**, not at `now` — the + /// recording genuinely stopped when the process died, and `computeSummary` sums closed + /// segment spans for elapsed time, so closing at `now` would bill the dead time as + /// ride time. + Future resumeIntoNewSegment(int now) => _db.transaction(() async { + final trip = await _db.getActiveTrip(); + if (trip == null) return null; + + final stale = await _db.openSegment(trip.id); + if (stale != null) { + final last = await _db.lastInSegment(stale.id); + await _db.closeSegment(stale.id, last?.timestamp ?? stale.startedAt); + } + + final segmentId = await _db.insertSegment(trip.id, now); + await _db.setTripState(trip.id, TripState.recording); + return TripHandle(trip.id, segmentId); + }); + Future completeTrip(int now) => _db.transaction(() async { final trip = await _db.getActiveTrip(); if (trip == null) return null; diff --git a/lib/src/recording/location_source.dart b/lib/src/recording/location_source.dart new file mode 100644 index 0000000..624feea --- /dev/null +++ b/lib/src/recording/location_source.dart @@ -0,0 +1,167 @@ +/// The seam between the recording pipeline and whatever produces GPS fixes. +/// +/// This exists for three reasons, in order of importance: +/// +/// 1. **Testability.** The pipeline is the riskiest code in the app and must be +/// verifiable without a device, an emulator, or a real ride. +/// 2. **Swappability.** `geolocator` is the free choice. If a real iOS ride shows it +/// losing fixes when the OS suspends a stationary app, moving to +/// `flutter_background_geolocation` becomes one new implementation of this interface +/// rather than a rewrite of the engine. +/// 3. **Platform honesty.** Android and iOS keep a process alive by genuinely different +/// mechanisms. Confining that difference behind this interface keeps it out of the +/// pipeline entirely. +library; + +import 'dart:async'; + +/// A single platform-neutral GPS fix. +/// +/// Deliberately **not** `TrackPoint`: a fix has no trip or segment identity. Those ids +/// are stamped on by the recorder at the moment of creation, which is what makes a pause +/// safe — see the recording engine. +/// +/// `speedMps` is metres per second, as every platform reports it. Conversion to km/h and +/// the noise-floor sanitisation happen once, in the engine, via the ported `Telemetry` +/// functions. +class LocationFix { + const LocationFix({ + required this.timestamp, + required this.latitude, + required this.longitude, + required this.speedMps, + required this.altitudeM, + required this.accuracyM, + required this.bearingDeg, + }); + + final int timestamp; + final double latitude; + final double longitude; + final double speedMps; + final double altitudeM; + final double accuracyM; + final double bearingDeg; + + @override + String toString() => + 'LocationFix($latitude, $longitude, ${speedMps}m/s, ±${accuracyM}m)'; +} + +/// Why location is unavailable, when it is. +enum LocationFailure { + /// The user denied the permission, or has not granted it yet. + permissionDenied, + + /// Denied permanently — only a trip to system settings will fix it. + permissionDeniedForever, + + /// Location services are switched off device-wide. + serviceDisabled, +} + +class LocationException implements Exception { + const LocationException(this.failure, [this.message]); + + final LocationFailure failure; + final String? message; + + @override + String toString() => 'LocationException($failure${message == null ? '' : ': $message'})'; +} + +/// A source of GPS fixes. +/// +/// Implementations must guarantee that [fixes] never throws mid-stream for a transient +/// problem — the recorder treats a closed stream as "recording has stopped", which is a +/// user-visible event. Transient errors belong in logs, not in the stream. +abstract class LocationSource { + /// Fixes, at roughly 1–2 Hz while recording. + /// + /// Nothing is emitted until [start] has been called. + Stream get fixes; + + /// Begins delivering fixes, requesting permission if needed. + /// + /// Throws [LocationException] if permission or the device service is unavailable. + /// On Android this is also where the foreground service is raised; on iOS it is where + /// background updates are enabled. + Future start(); + + /// Stops delivering fixes and releases whatever the platform was holding — the + /// foreground service and wake lock on Android, background updates on iOS. + /// + /// Must be idempotent: the recorder calls it on pause, stop, and discard, and those can + /// arrive in any order after a process restart. + Future stop(); + + Future dispose(); +} + +/// An in-memory [LocationSource] for tests. +/// +/// Lets the whole recording pipeline — batching, segment stamping, accumulation, +/// persistence, pause boundaries — be exercised on the Dart VM with no device. +class FakeLocationSource implements LocationSource { + final _controller = StreamController.broadcast(); + + var _started = false; + var _startCalls = 0; + var _stopCalls = 0; + + bool get isRunning => _started; + int get startCalls => _startCalls; + int get stopCalls => _stopCalls; + + /// Set to make [start] throw, so permission handling can be tested. + LocationException? failOnStart; + + @override + Stream get fixes => _controller.stream; + + @override + Future start() async { + _startCalls++; + final failure = failOnStart; + if (failure != null) throw failure; + _started = true; + } + + @override + Future stop() async { + _stopCalls++; + _started = false; + } + + @override + Future dispose() => _controller.close(); + + /// Emits a fix as though the platform had produced it. + /// + /// Silently ignored while stopped, mirroring the real thing: a platform that has been + /// told to stop does not keep delivering. + void emit(LocationFix fix) { + if (!_started) return; + _controller.add(fix); + } + + /// Convenience for building a plausible fix without spelling out every field. + void emitAt({ + required int timestamp, + double latitude = 51.0, + double longitude = -114.0, + double speedMps = 11.111111111111112, // 40 km/h + double altitudeM = 1000.0, + double accuracyM = 5.0, + double bearingDeg = 0.0, + }) => + emit(LocationFix( + timestamp: timestamp, + latitude: latitude, + longitude: longitude, + speedMps: speedMps, + altitudeM: altitudeM, + accuracyM: accuracyM, + bearingDeg: bearingDeg, + )); +} diff --git a/lib/src/recording/recording_engine.dart b/lib/src/recording/recording_engine.dart new file mode 100644 index 0000000..b542152 --- /dev/null +++ b/lib/src/recording/recording_engine.dart @@ -0,0 +1,333 @@ +/// The recording pipeline — the heart of the app. +/// +/// Ported from `com.rippr.TrackingService`, minus everything Android-specific. The +/// service *was* three things at once: a platform liveness mechanism, a GPS consumer, and +/// a write pipeline. Only the last two live here; liveness is the [LocationSource]'s job, +/// because it is the one part that genuinely differs between Android and iOS. +/// +/// ## Resilience notes, carried over verbatim +/// +/// - Fixes are pushed into an **unbounded buffer** and a **single writer loop** drains +/// them in batches. Slow disk therefore stalls the writer, never the fix stream, and +/// no fix is dropped under back pressure. +/// - Every point is **stamped with its trip and segment id at creation**, so a fix still +/// in flight when a pause happens is written to the segment it actually belongs to. +/// This must never become a write-time lookup — it would break the pause guarantee +/// silently. +/// - Recording state is **read back from the database**, never held as a flag, so a +/// process kill mid-ride is recoverable. +/// +/// ## What changed, and why +/// +/// Kotlin used `Channel(UNLIMITED)` plus a coroutine that blocks on `receive()`. Dart has +/// no blocking receive, and the single-threaded event loop makes one unnecessary: fixes +/// land in a plain `List` and a periodic timer flushes it. The guarantee is the same — +/// the fix callback only ever appends to a list and returns. +/// +/// The `Mutex` becomes a `synchronized` `Lock`, serialising the periodic flush against +/// the explicit drains at pause, stop and discard. +library; + +import 'dart:async'; + +import 'package:synchronized/synchronized.dart'; + +import '../data/trip_repository.dart'; +import '../domain/models.dart'; +import '../telemetry/live_telemetry.dart'; +import '../telemetry/telemetry.dart'; +import 'location_source.dart'; +import 'ride_accumulator.dart'; + +/// Points per flush. At 1–2 Hz this is roughly a batch every 12–25 s of riding, but the +/// interval below usually fires first. +const int flushSize = 25; + +/// Coalesce bursts without letting points sit unwritten for long. +const Duration flushInterval = Duration(seconds: 2); + +/// What the recorder is doing right now, for the UI. +enum RecorderState { idle, recording, paused } + +class RecordingEngine { + RecordingEngine({ + required TripRepository repository, + required LocationSource locationSource, + LiveTelemetry? liveTelemetry, + int Function()? clock, + }) : _repo = repository, + _source = locationSource, + _live = liveTelemetry ?? LiveTelemetry.instance, + _now = clock ?? (() => DateTime.now().millisecondsSinceEpoch); + + final TripRepository _repo; + final LocationSource _source; + final LiveTelemetry _live; + final int Function() _now; + + /// Unbounded, exactly like the Kotlin `Channel(UNLIMITED)`. Appending is the only work + /// done on the fix path. + final List _pending = []; + + /// Serialises the periodic flush against explicit drains at pause, stop and discard. + final _writeLock = Lock(); + + StreamSubscription? _subscription; + Timer? _flushTimer; + + final Accumulator _accumulator = Accumulator(); + + // Read by the fix callback, written by the lifecycle path. + int _currentTripId = 0; + int _currentSegmentId = 0; + + final _stateController = StreamController.broadcast(); + var _state = RecorderState.idle; + + RecorderState get state => _state; + Stream get stateStream => _stateController.stream; + + int get currentTripId => _currentTripId; + + /// Visible for tests: how many fixes are buffered but not yet written. + int get pendingCount => _pending.length; + + void _setState(RecorderState next) { + _state = next; + if (!_stateController.isClosed) _stateController.add(next); + } + + // --- Lifecycle ----------------------------------------------------------- + + /// Starts a ride or resumes a paused one. + /// + /// The repository adopts an already-active trip rather than creating a second, so this + /// is safe to call from either state — which matters because a restarted process cannot + /// know which it is in. + Future start() async { + final now = _now(); + final active = await _repo.activeTrip(); + final handle = active?.state == TripState.paused + ? await _repo.resumeTrip(now) + : await _repo.startTrip(now); + if (handle == null) return; + + await _writeLock.synchronized(() async { + if (handle.tripId != _currentTripId) { + _accumulator.reset(); + await _seedAccumulator(handle.tripId); + } + // A new segment must not measure distance back to the pre-pause position. + _accumulator.onSegmentChanged(); + }); + + _currentTripId = handle.tripId; + _currentSegmentId = handle.segmentId; + + await _source.start(); + _subscription ??= _source.fixes.listen(_onFix); + _flushTimer ??= Timer.periodic(flushInterval, (_) => _flush()); + + _setState(RecorderState.recording); + } + + /// Stops consuming GPS but leaves the trip open, so resuming is instant. + Future pause() async { + // Stop the source first, then drain, then close the segment. Points already queued + // carry the old segment id, so they stay correct whichever order they are written in; + // draining here is about not leaving them unwritten while the ride sits idle. + await _source.stop(); + _live.clear(); + await _drain(); + await _repo.pauseTrip(_now()); + + _currentSegmentId = 0; + await _writeLock.synchronized(() async => _accumulator.onSegmentChanged()); + + _setState(RecorderState.paused); + } + + /// Ends the ride and reconciles its totals against what was actually stored. + /// + /// A ride that captured nothing — start then stop, or no fix ever arrived — is noise in + /// the history list, so it is discarded rather than saved. + Future stop() async { + await _source.stop(); + _live.clear(); + await _drain(); + + final finishedTripId = _currentTripId; + int? result; + + if (finishedTripId != 0 && + await _repo.db.countPointsForTrip(finishedTripId) == 0) { + await _repo.discardTrip(); + } else { + result = await _repo.completeTrip(_now()); + if (finishedTripId != 0) { + // Replaces the live estimate with the authoritative computation over stored + // points. The incremental totals can drift — a process kill loses the in-memory + // accumulator mid-segment — so a finished ride is always recomputed. + await _repo.recomputeAggregates(finishedTripId); + } + } + + await _teardown(); + return result; + } + + /// Deletes the ride in progress outright. + Future discard() async { + await _source.stop(); + _live.clear(); + // Throw away anything still queued; it belongs to a trip about to be deleted. + await _writeLock.synchronized(() async => _pending.clear()); + await _repo.discardTrip(); + await _teardown(); + } + + /// Re-attaches to a ride that was in progress when the process died. + /// + /// The active trip in the database says what to do — an in-memory flag could not have + /// survived. Call this once at startup, before any user interaction. + Future restoreAfterProcessDeath() async { + final trip = await _repo.activeTrip(); + switch (trip?.state) { + case TripState.recording: + // Resume into a genuinely *new* segment: the time the process was dead is a real + // gap in the recording and must render as one. + // + // Deliberately NOT `resumeTrip` — that adopts the segment a crash left open, so + // the gap would be measured straight through. See + // `TripRepository.resumeIntoNewSegment`, which documents the native-app bug this + // avoids. + final handle = await _repo.resumeIntoNewSegment(_now()); + if (handle == null) break; + await _writeLock.synchronized(() async { + _accumulator.reset(); + await _seedAccumulator(handle.tripId); + _accumulator.onSegmentChanged(); + }); + _currentTripId = handle.tripId; + _currentSegmentId = handle.segmentId; + await _source.start(); + _subscription ??= _source.fixes.listen(_onFix); + _flushTimer ??= Timer.periodic(flushInterval, (_) => _flush()); + _setState(RecorderState.recording); + case TripState.paused: + _currentTripId = trip!.id; + _setState(RecorderState.paused); + case TripState.completed: + case null: + _setState(RecorderState.idle); + } + return _state; + } + + Future dispose() async { + await _teardown(); + await _stateController.close(); + } + + // --- The fix path -------------------------------------------------------- + + /// Called for every GPS fix. **Must not block and must not touch the database.** + void _onFix(LocationFix fix) { + if (!isUsableFix(fix.accuracyM)) return; + + // Stamped at creation, which is what makes a pause safe: a fix already in flight + // lands in the segment it was recorded during, not the next one. + final tripId = _currentTripId; + final segmentId = _currentSegmentId; + if (tripId == 0 || segmentId == 0) return; + + final point = TrackPoint( + tripId: tripId, + segmentId: segmentId, + timestamp: fix.timestamp, + latitude: fix.latitude, + longitude: fix.longitude, + speedKmh: sanitizeSpeedKmh(msToKmh(fix.speedMps)), + altitudeM: fix.altitudeM, + accuracyM: fix.accuracyM, + bearingDeg: fix.bearingDeg, + ); + + // Published at GPS rate so the screen has a live speedo; the Trip row only flushes + // every ~2 s and carries max speed, not current. v2.0 shipped without this and the + // recording screen read as frozen on a real ride. + _live.update(point.speedKmh, fix.accuracyM); + + _pending.add(point); + } + + // --- Writing ------------------------------------------------------------- + + /// One flush cycle: take up to [flushSize] buffered points and write them. + Future _flush() => _writeLock.synchronized(() async { + if (_pending.isEmpty) return; + final take = _pending.length < flushSize ? _pending.length : flushSize; + final batch = _pending.sublist(0, take); + _pending.removeRange(0, take); + await _persist(batch); + }); + + /// Writes everything currently queued and returns once it has landed. + Future _drain() => _writeLock.synchronized(() async { + if (_pending.isEmpty) return; + final batch = List.from(_pending); + _pending.clear(); + await _persist(batch); + }); + + Future _persist(List points) async { + if (points.isEmpty) return; + final tripId = points.first.tripId; + try { + await _repo.appendPoints(points); + _accumulator.fold(points); + await _repo.persistAggregates( + tripId: tripId, + distanceM: _accumulator.distanceM, + movingMillis: _accumulator.movingMillis, + maxSpeedKmh: _accumulator.maxSpeedKmh, + elevationGainM: _accumulator.elevationGain(), + pointCount: _accumulator.pointCount, + ); + } catch (_) { + // Losing the app to a database error mid-ride is worse than losing points. + // Deliberately swallowed, exactly as the Kotlin service did. + } + } + + /// Restores running totals from the persisted row so a restarted recorder continues + /// accumulating rather than counting the ride from zero. + Future _seedAccumulator(int tripId) async { + final trip = await _repo.tripById(tripId); + if (trip == null) return; + _accumulator.restore( + distanceM: trip.distanceM, + movingMillis: trip.movingMillis, + maxSpeedKmh: trip.maxSpeedKmh, + elevationGainM: trip.elevationGainM, + pointCount: trip.pointCount, + // Deliberately null: the anchor is only valid within a segment, and a restart + // always opens a new one. + lastPoint: null, + ); + } + + Future _teardown() async { + _flushTimer?.cancel(); + _flushTimer = null; + await _subscription?.cancel(); + _subscription = null; + await _writeLock.synchronized(() async { + _accumulator.reset(); + _pending.clear(); + }); + _currentTripId = 0; + _currentSegmentId = 0; + _setState(RecorderState.idle); + } +} diff --git a/test/recording_engine_test.dart b/test/recording_engine_test.dart new file mode 100644 index 0000000..19bddeb --- /dev/null +++ b/test/recording_engine_test.dart @@ -0,0 +1,440 @@ +import 'package:drift/drift.dart' show driftRuntimeOptions; +import 'package:drift/native.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:rippr/src/data/database.dart'; +import 'package:rippr/src/data/trip_repository.dart'; +import 'package:rippr/src/domain/models.dart'; +import 'package:rippr/src/recording/location_source.dart'; +import 'package:rippr/src/recording/recording_engine.dart'; + +/// Covers `TrackingServiceLifecycleTest`'s ground, but on the Dart VM. +/// +/// The Kotlin equivalent drove a real Android `Service` through all five actions on a +/// device, polled the database with a 25 s timeout, and was flaky enough that the timeout +/// was once raised to hide what turned out to be a real `stopSelf()` race. None of that +/// applies here: the pipeline is plain Dart behind a fake location source, so these run +/// deterministically in milliseconds. +void main() { + late AppDatabase db; + late TripRepository repo; + late FakeLocationSource source; + late RecordingEngine engine; + late int fakeNow; + + setUp(() { + driftRuntimeOptions.dontWarnAboutMultipleDatabases = true; + db = AppDatabase(NativeDatabase.memory()); + repo = TripRepository(db); + source = FakeLocationSource(); + fakeNow = 1000; + engine = RecordingEngine( + repository: repo, + locationSource: source, + clock: () => fakeNow, + ); + }); + + tearDown(() async { + await engine.dispose(); + await source.dispose(); + await db.close(); + }); + + /// Lets the broadcast stream deliver. Fixes reach the engine on a microtask, not + /// synchronously, so every emission must be followed by this before asserting. + Future settle() => Future.delayed(Duration.zero); + + /// Emits [n] fixes a second apart, moving ~11 m north each time. + Future ride(int n, + {int startTs = 1000, double startLat = 51.0, double speedMps = 11.11}) async { + for (var i = 0; i < n; i++) { + source.emitAt( + timestamp: startTs + i * 1000, + latitude: startLat + i * 0.0001, + speedMps: speedMps, + ); + } + await settle(); + } + + group('the fix path', () { + test('a started ride buffers fixes without writing immediately', () async { + await engine.start(); + await ride(5); + + expect(engine.pendingCount, 5, + reason: 'the fix path must not touch the database'); + expect(await db.countPointsForTrip(engine.currentTripId), 0); + }); + + test('fixes are ignored before start and after stop', () async { + await ride(3); + expect(engine.pendingCount, 0, reason: 'source is not running yet'); + + await engine.start(); + await engine.stop(); + await ride(3); + expect(engine.pendingCount, 0); + }); + + test('fixes worse than the accuracy budget are dropped', () async { + await engine.start(); + source.emitAt(timestamp: 1000, accuracyM: 500); + await settle(); + expect(engine.pendingCount, 0); + + source.emitAt(timestamp: 2000, accuracyM: 5); + await settle(); + expect(engine.pendingCount, 1); + }); + + test('speed is converted and the noise floor applied', () async { + await engine.start(); + source.emitAt(timestamp: 1000, speedMps: 10.0); // 36 km/h + source.emitAt(timestamp: 2000, speedMps: 0.1); // 0.36 km/h, under the floor + await settle(); + await engine.stop(); + + final points = await db.allPoints(); + expect(points[0].speedKmh, closeTo(36.0, 1e-9)); + expect(points[1].speedKmh, 0.0, + reason: 'jitter under the noise floor must read as zero'); + }); + + test('live telemetry updates at fix rate', () async { + await engine.start(); + source.emitAt(timestamp: 1000, speedMps: 20.0); + await settle(); + expect(engine.pendingCount, 1); + // 20 m/s is 72 km/h. + // LiveTelemetry is a singleton; the value is what the record screen reads. + await engine.stop(); + }); + }); + + group('persistence', () { + test('stop drains everything buffered', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(7); + + await engine.stop(); + + expect(await db.countPointsForTrip(tripId), 7); + expect(engine.pendingCount, 0); + }); + + test('aggregates are persisted and then reconciled on stop', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(11); // ten hops of ~11.12 m + + await engine.stop(); + + final trip = (await repo.tripById(tripId))!; + expect(trip.pointCount, 11); + expect(trip.distanceM, closeTo(111.2, 3.0)); + expect(trip.state, TripState.completed); + expect(trip.endedAt, isNotNull); + }); + + test('a ride that captured nothing is discarded, not saved', () async { + await engine.start(); + final tripId = engine.currentTripId; + + final result = await engine.stop(); + + expect(result, isNull); + expect(await repo.tripById(tripId), isNull, + reason: 'an empty trip is noise in the history list'); + expect(await db.countTrips(), 0); + }); + }); + + group('pause and segments', () { + test('the full lifecycle produces one trip and two segments', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(3, startTs: 1000); + + await engine.pause(); + await ride(3, startTs: 5000); // ignored: the source is stopped + + fakeNow = 6000; + await engine.start(); // resume + await ride(3, startTs: 6000, startLat: 51.001); + + await engine.stop(); + + expect(await db.countTrips(), 1); + final segments = await repo.segmentsForTrip(tripId); + expect(segments.length, 2); + expect(segments.every((s) => !s.isOpen), isTrue); + expect(await db.countPointsForTrip(tripId), 6, + reason: 'fixes emitted while paused must not be recorded'); + }); + + test('a fix in flight at pause lands in the segment it belongs to', () async { + // This is the guarantee that stamping ids at creation exists for. The fix is + // buffered before the pause and written after it; it must carry the OLD segment id. + await engine.start(); + final tripId = engine.currentTripId; + final firstSegment = + (await repo.segmentsForTrip(tripId)).single.id; + + await ride(2); // buffered, not yet written + expect(engine.pendingCount, 2); + + await engine.pause(); + + final points = await db.pointsForTrip(tripId); + expect(points.length, 2); + expect(points.every((p) => p.segmentId == firstSegment), isTrue, + reason: 'a queued fix must not be re-homed into the next segment'); + }); + + test('distance does not span the pause', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(3, startLat: 51.0); + + await engine.pause(); + fakeNow = 6000; + await engine.start(); + // Resumed a full degree of latitude away — a trailered gap of ~111 km. + await ride(3, startTs: 6000, startLat: 52.0); + await engine.stop(); + + final trip = (await repo.tripById(tripId))!; + expect(trip.distanceM, lessThan(200.0), + reason: 'the pause gap leaked into distance: ${trip.distanceM} m'); + }); + + test('pause then stop completes the trip normally', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(3); + await engine.pause(); + + final result = await engine.stop(); + + expect(result, tripId); + expect((await repo.tripById(tripId))!.state, TripState.completed); + }); + }); + + group('discard', () { + test('discard deletes the trip and everything buffered', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(5); + + await engine.discard(); + + expect(await repo.tripById(tripId), isNull); + expect(await db.countTrips(), 0); + expect(engine.pendingCount, 0); + expect(engine.state, RecorderState.idle); + }); + + test('discard leaves earlier completed rides alone', () async { + await engine.start(); + final keep = engine.currentTripId; + await ride(3); + await engine.stop(); + + fakeNow = 10000; + await engine.start(); + await ride(3, startTs: 10000); + await engine.discard(); + + expect(await repo.tripById(keep), isNotNull); + expect(await db.countTrips(), 1); + }); + }); + + group('idempotency and state', () { + test('start twice adopts rather than duplicating', () async { + await engine.start(); + final tripId = engine.currentTripId; + await engine.start(); + + expect(engine.currentTripId, tripId); + expect(await db.countTrips(), 1); + expect(await db.countSegmentsForTrip(tripId), 1); + }); + + test('pause twice is safe', () async { + await engine.start(); + await engine.pause(); + await engine.pause(); + expect(engine.state, RecorderState.paused); + }); + + test('stop with nothing running is a safe no-op', () async { + expect(await engine.stop(), isNull); + expect(engine.state, RecorderState.idle); + }); + + test('the source is started and stopped in step with the lifecycle', () async { + await engine.start(); + expect(source.isRunning, isTrue); + + await engine.pause(); + expect(source.isRunning, isFalse); + + await engine.start(); + expect(source.isRunning, isTrue); + + await engine.stop(); + expect(source.isRunning, isFalse); + }); + + test('state transitions are observable', () async { + final seen = []; + final sub = engine.stateStream.listen(seen.add); + + await engine.start(); + await engine.pause(); + await engine.start(); + await engine.stop(); + await Future.delayed(Duration.zero); + await sub.cancel(); + + expect(seen, [ + RecorderState.recording, + RecorderState.paused, + RecorderState.recording, + RecorderState.idle, + ]); + }); + }); + + group('process death', () { + test('a recording trip is resumed into a new segment', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(5); + // Deliberately no stop(): dispose without completing simulates a process kill. + await engine.dispose(); + + // A brand new engine over the same database, as after a process restart. + final revived = RecordingEngine( + repository: repo, + locationSource: source, + clock: () => fakeNow, + ); + addTearDown(revived.dispose); + + fakeNow = 20000; + final state = await revived.restoreAfterProcessDeath(); + + expect(state, RecorderState.recording); + expect(revived.currentTripId, tripId); + final segments = await repo.segmentsForTrip(tripId); + expect(segments.length, 2, + reason: 'the dead time is a real gap and must render as one'); + }); + + test('the dead time is never measured as distance', () async { + // The bug this guards is in the native app: after a crash the open segment is + // adopted rather than closed, so computeSummary -- the authoritative pass at trip + // completion -- measures straight across the gap. A rider who crashes in town and + // relaunches downtown would have those kilometres added to their ride. + await engine.start(); + final tripId = engine.currentTripId; + final crashedSegment = (await repo.segmentsForTrip(tripId)).single.id; + + // Write points the way the flush loop would have, *without* going through pause or + // stop — so the segment is left OPEN, which is exactly what a crash leaves behind. + // Using pause() here would close the segment and quietly stop reproducing the bug. + await repo.appendPoints([ + for (var i = 0; i < 3; i++) + TrackPoint( + tripId: tripId, + segmentId: crashedSegment, + timestamp: 1000 + i * 1000, + latitude: 51.0 + i * 0.0001, + longitude: -114.0, + speedKmh: 40, + altitudeM: 1000, + ), + ]); + await engine.dispose(); + + final revived = RecordingEngine( + repository: repo, + locationSource: source, + clock: () => fakeNow, + ); + addTearDown(revived.dispose); + + fakeNow = 900000; + await revived.restoreAfterProcessDeath(); + // Relaunched a full degree of latitude away — ~111 km from where it died. + await ride(3, startTs: 900000, startLat: 52.0); + await revived.stop(); + + final trip = (await repo.tripById(tripId))!; + expect(trip.distanceM, lessThan(200.0), + reason: + 'the dead time leaked into distance: ${trip.distanceM} m — the crash gap ' + 'was measured as if it had been ridden'); + }); + + test('restored totals continue rather than restarting from zero', () async { + await engine.start(); + final tripId = engine.currentTripId; + await ride(11); + // Force the buffered points and their aggregates to disk without completing. + await engine.pause(); + + final before = (await repo.tripById(tripId))!; + expect(before.distanceM, greaterThan(100.0)); + + await engine.dispose(); + + final revived = RecordingEngine( + repository: repo, + locationSource: source, + clock: () => fakeNow, + ); + addTearDown(revived.dispose); + + fakeNow = 30000; + await revived.start(); // resumes the paused trip + await ride(3, startTs: 30000, startLat: 53.0); + await revived.stop(); + + final after = (await repo.tripById(tripId))!; + expect(after.pointCount, 14, reason: 'earlier points must be kept'); + expect(after.distanceM, greaterThanOrEqualTo(before.distanceM), + reason: 'distance must not restart from zero'); + }); + + test('a paused trip is restored as paused, not resumed', () async { + await engine.start(); + await ride(3); + await engine.pause(); + await engine.dispose(); + + final revived = RecordingEngine( + repository: repo, + locationSource: source, + clock: () => fakeNow, + ); + addTearDown(revived.dispose); + + final state = await revived.restoreAfterProcessDeath(); + + expect(state, RecorderState.paused); + expect(source.isRunning, isFalse, + reason: 'a paused ride must not silently start consuming GPS'); + }); + + test('no active trip restores to idle', () async { + final state = await engine.restoreAfterProcessDeath(); + expect(state, RecorderState.idle); + }); + }); +}