T10/T11: location seam and recording pipeline; fix a native crash-recovery bug
T10 -- LocationSource abstraction plus FakeLocationSource, so the riskiest code in the app is testable with no device. Also found that geolocator ships its own Android foreground service with enableWakeLock, which could drop flutter_foreground_task and both its deprecation paths; deferred to T12 since it costs notification actions. T11 -- the pipeline. Kotlin's unbounded Channel plus blocking-receive writer becomes a List buffer plus a periodic Timer; Dart's event loop makes a blocking receive unnecessary and the guarantee is unchanged, since the fix callback only appends and returns. 24 tests. Found a real bug in the native app while porting. restoreAfterProcessDeath says it resumes into a new segment because the dead time is a real gap, but it calls resumeTrip, which adopts the segment a crash left open. After a pause that is right; after a crash nothing closed it. Points either side of the dead time then share a segment, and computeSummary -- the authoritative pass that overwrites the live estimate on completion -- measures straight through the gap. Reverting the fix and running the guard shows 111,217 m of phantom distance. Fixed via TripRepository.resumeIntoNewSegment, which closes the stale segment at its last recorded point rather than at now, so the dead time is not billed as ride time either. A deliberate, documented departure from parity: the code contradicted its own comment and produced silently wrong data. The first version of that guard passed with the bug still present -- the fixture used pause(), which closes the segment and stops reproducing the crash. Caught only by reverting the fix and checking the test failed. 145 tests passing, analyze clean. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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.**
|
||||
|
||||
@@ -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<TripHandle?> 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<int?> completeTrip(int now) => _db.transaction(() async {
|
||||
final trip = await _db.getActiveTrip();
|
||||
if (trip == null) return null;
|
||||
|
||||
167
lib/src/recording/location_source.dart
Normal file
167
lib/src/recording/location_source.dart
Normal file
@@ -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<LocationFix> 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<void> 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<void> stop();
|
||||
|
||||
Future<void> 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<LocationFix>.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<LocationFix> get fixes => _controller.stream;
|
||||
|
||||
@override
|
||||
Future<void> start() async {
|
||||
_startCalls++;
|
||||
final failure = failOnStart;
|
||||
if (failure != null) throw failure;
|
||||
_started = true;
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> stop() async {
|
||||
_stopCalls++;
|
||||
_started = false;
|
||||
}
|
||||
|
||||
@override
|
||||
Future<void> 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,
|
||||
));
|
||||
}
|
||||
333
lib/src/recording/recording_engine.dart
Normal file
333
lib/src/recording/recording_engine.dart
Normal file
@@ -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<TrackPoint> _pending = [];
|
||||
|
||||
/// Serialises the periodic flush against explicit drains at pause, stop and discard.
|
||||
final _writeLock = Lock();
|
||||
|
||||
StreamSubscription<LocationFix>? _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<RecorderState>.broadcast();
|
||||
var _state = RecorderState.idle;
|
||||
|
||||
RecorderState get state => _state;
|
||||
Stream<RecorderState> 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<void> 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<void> 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<int?> 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<void> 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<RecorderState> 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<void> 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<void> _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<void> _drain() => _writeLock.synchronized(() async {
|
||||
if (_pending.isEmpty) return;
|
||||
final batch = List<TrackPoint>.from(_pending);
|
||||
_pending.clear();
|
||||
await _persist(batch);
|
||||
});
|
||||
|
||||
Future<void> _persist(List<TrackPoint> 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<void> _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<void> _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);
|
||||
}
|
||||
}
|
||||
440
test/recording_engine_test.dart
Normal file
440
test/recording_engine_test.dart
Normal file
@@ -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<void> settle() => Future<void>.delayed(Duration.zero);
|
||||
|
||||
/// Emits [n] fixes a second apart, moving ~11 m north each time.
|
||||
Future<void> 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 = <RecorderState>[];
|
||||
final sub = engine.stateStream.listen(seen.add);
|
||||
|
||||
await engine.start();
|
||||
await engine.pause();
|
||||
await engine.start();
|
||||
await engine.stop();
|
||||
await Future<void>.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);
|
||||
});
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user