/// 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; // Dart does not permit a named parameter whose name begins with an underscore, so the // lint's suggested `required this._uploadPending` will not compile here. // ignore_for_file: prefer_initializing_formals import 'dart:async'; import 'package:synchronized/synchronized.dart'; import '../data/trip_repository.dart'; import '../domain/activity_profile.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); /// How often the backlog is offered to the server. /// /// Upload runs on its own timer so a slow or dead endpoint can never interrupt /// recording — the same separation the native service kept between its writer loop and /// its upload loop. const Duration uploadInterval = Duration(seconds: 30); /// 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, Future Function()? uploadPending, void Function(String reason)? onUnexpectedStop, }) : _repo = repository, _source = locationSource, _live = liveTelemetry ?? LiveTelemetry.instance, _now = clock ?? (() => DateTime.now().millisecondsSinceEpoch), _uploadPending = uploadPending, _onUnexpectedStop = onUnexpectedStop; final TripRepository _repo; final LocationSource _source; final LiveTelemetry _live; final int Function() _now; /// Injected rather than constructed here, so the engine has no opinion about HTTP and /// tests need no network. final Future Function()? _uploadPending; Timer? _uploadTimer; /// V3-12: injected rather than importing a crash reporter directly, the same reasoning /// as [_uploadPending] -- the engine stays free of any opinion about where a report /// goes, and tests need no Sentry client. final void Function(String reason)? _onUnexpectedStop; /// 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; /// Rebuilt, not mutated, whenever a new trip session begins — see [start] — so its /// [Accumulator.profile] always matches the trip it is accumulating for. Accumulator _accumulator = Accumulator(); /// Which activity's defaults are currently in force for the fix path — the accuracy /// gate and noise floor in [_onFix] read this directly, since that callback must stay /// synchronous and cannot query the trip's activity from the database per fix. ActivityProfile _activeProfile = ActivityProfile.motorcycle; // 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) { // A different trip session: rebuild the accumulator under that trip's own // activity profile, rather than mutating the previous trip's accumulator in // place — a ride switched from walking to a motorcycle mid-session must not // keep walking's noise floor. final trip = await _repo.tripById(handle.tripId); _activeProfile = ActivityProfile.forActivity( trip?.activity ?? Activity.motorcycle, ); _accumulator = Accumulator(profile: _activeProfile); 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()); _startUploadLoop(); _setState(RecorderState.recording); } void _startUploadLoop() { final upload = _uploadPending; if (upload == null || _uploadTimer != null) return; _uploadTimer = Timer.periodic(uploadInterval, (_) async { // Swallowed on purpose. Nothing about uploading may disturb recording. try { await upload(); } catch (_) {} }); } /// 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: // A trip still marked `recording` at launch means the previous process ended // without ever calling stop() or pause() -- a crash, an OS kill, or a location // permission revoked out from under the app. This is the failure V3-12 exists to // surface: the recording stopped, silently, and nobody chose that. _onUnexpectedStop?.call('process death mid-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; _activeProfile = ActivityProfile.forActivity(trip!.activity); await _writeLock.synchronized(() async { _accumulator = Accumulator(profile: _activeProfile); 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; _activeProfile = ActivityProfile.forActivity(trip.activity); _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, maxAccuracyMeters: _activeProfile.accuracyGateM)) { 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), floorKmh: _activeProfile.noiseFloorKmh, ), 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; _uploadTimer?.cancel(); _uploadTimer = null; await _subscription?.cancel(); _subscription = null; await _writeLock.synchronized(() async { _accumulator.reset(); _pending.clear(); }); _currentTripId = 0; _currentSegmentId = 0; _setState(RecorderState.idle); } }