T08: Drift schema mirroring the Room baseline

Three tables with CASCADE foreign keys, indices, WAL, and a real MigrationStrategy
from schema version 1. No destructive fallback, ever.

Two things needed deliberate handling. SQLite defaults foreign_keys to OFF and
Drift, unlike Room, does not enable it -- without the pragma every CASCADE is
decorative, so there is now a test that reads the pragma back. And Drift both
snake_cases columns and names row classes after tables; build.yaml sets
case_from_dart_to_sql: preserve so the schema stays column-for-column identical
to Room's, and @DataClassName keeps row types from colliding with the domain
models.

SchemaTest was instrumented and needed a device; the Drift version is a plain
unit test that runs in under a second with no emulator. 16 tests.

95 tests passing, analyze clean.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-15 01:13:45 -05:00
parent 398157b1f5
commit d7d9854dc9
5 changed files with 3749 additions and 0 deletions

429
lib/src/data/database.dart Normal file
View File

@@ -0,0 +1,429 @@
/// Drift schema, mirroring the Room schema at
/// `app/schemas/com.rippr.data.AppDatabase/2.json` in the native repo.
///
/// Ported from `com.rippr.data.AppDatabase` plus the three DAOs.
///
/// ## Two invariants carried over verbatim
///
/// **No destructive migration, ever.** v2 removed `fallbackToDestructiveMigration()`
/// because rides are real data. This starts at Dart schema version 1 with a real
/// [MigrationStrategy] from day one; reinstating a destructive fallback would silently
/// delete every stored ride on the next version bump.
///
/// **Foreign keys must be switched on explicitly.** SQLite defaults `foreign_keys` to
/// OFF. Room turned it on for us; Drift does not. Without the pragma in [beforeOpen] the
/// `CASCADE` deletes below are decorative, and deleting a trip would silently orphan
/// every one of its points. There is a test for exactly this.
library;
import 'package:drift/drift.dart';
import '../domain/models.dart' as domain;
part 'database.g.dart';
/// One ride, from pressing Start to pressing Stop.
///
/// The aggregate columns are denormalised on purpose — accumulated as points arrive and
/// recomputed authoritatively on completion, so the trips list never touches the point
/// table.
@DataClassName('TripRow')
@TableIndex(name: 'idx_trips_ended', columns: {#endedAt})
class Trips extends Table {
@override
String get tableName => 'trips';
IntColumn get id => integer().autoIncrement()();
IntColumn get startedAt => integer()();
/// Null while the ride is still active. This column *is* the fact of an in-progress
/// ride, which is why recording state survives process death.
IntColumn get endedAt => integer().nullable()();
/// Null means the UI derives a label from [startedAt]. Never store an empty string.
TextColumn get name => text().nullable()();
TextColumn get state => textEnum<domain.TripState>()();
RealColumn get distanceM => real().withDefault(const Constant(0))();
IntColumn get movingMillis => integer().withDefault(const Constant(0))();
RealColumn get maxSpeedKmh => real().withDefault(const Constant(0))();
RealColumn get elevationGainM => real().withDefault(const Constant(0))();
IntColumn get pointCount => integer().withDefault(const Constant(0))();
}
/// One pause-free stretch of recording within a trip.
///
/// This layer is what makes pause correct rather than cosmetic.
@DataClassName('SegmentRow')
@TableIndex(name: 'idx_segments_trip', columns: {#tripId})
class Segments extends Table {
@override
String get tableName => 'segments';
IntColumn get id => integer().autoIncrement()();
IntColumn get tripId =>
integer().references(Trips, #id, onDelete: KeyAction.cascade)();
IntColumn get startedAt => integer()();
/// Null while this segment is still being recorded into.
IntColumn get endedAt => integer().nullable()();
}
/// A single GPS fix.
@DataClassName('TrackPointRow')
@TableIndex(name: 'idx_points_trip', columns: {#tripId})
@TableIndex(name: 'idx_points_segment', columns: {#segmentId})
@TableIndex(name: 'idx_points_synced', columns: {#synced})
class TrackPoints extends Table {
@override
String get tableName => 'track_points';
IntColumn get id => integer().autoIncrement()();
IntColumn get tripId =>
integer().references(Trips, #id, onDelete: KeyAction.cascade)();
IntColumn get segmentId =>
integer().references(Segments, #id, onDelete: KeyAction.cascade)();
IntColumn get timestamp => integer()();
RealColumn get latitude => real()();
RealColumn get longitude => real()();
RealColumn get speedKmh => real()();
RealColumn get altitudeM => real()();
RealColumn get accuracyM => real().withDefault(const Constant(0))();
RealColumn get bearingDeg => real().withDefault(const Constant(0))();
/// Set once the point has been accepted by the remote endpoint.
BoolColumn get synced => boolean().withDefault(const Constant(false))();
}
@DriftDatabase(tables: [Trips, Segments, TrackPoints])
class AppDatabase extends _$AppDatabase {
AppDatabase(super.e);
@override
int get schemaVersion => 1;
@override
MigrationStrategy get migration => MigrationStrategy(
onCreate: (m) => m.createAll(),
beforeOpen: (details) async {
// Non-negotiable: without this the CASCADE relationships above do nothing.
await customStatement('PRAGMA foreign_keys = ON');
// A ride is unrecoverable if a write is lost to a crash mid-flush, but full
// sync on every insert at 2 Hz burns battery. WAL with NORMAL sync is the
// standard compromise and survives app crashes; only an OS-level crash can
// lose the last few points.
await customStatement('PRAGMA journal_mode = WAL');
await customStatement('PRAGMA synchronous = NORMAL');
},
);
// --- Trips ---------------------------------------------------------------
/// The in-progress ride, or null. This is the source of truth for "are we recording" —
/// it survives process death, which an in-memory flag cannot.
Stream<domain.Trip?> watchActiveTrip() => (select(trips)
..where((t) => t.endedAt.isNull())
..orderBy([(t) => OrderingTerm.desc(t.id)])
..limit(1))
.watchSingleOrNull()
.map((r) => r == null ? null : _toTrip(r));
Future<domain.Trip?> getActiveTrip() async {
final row = await (select(trips)
..where((t) => t.endedAt.isNull())
..orderBy([(t) => OrderingTerm.desc(t.id)])
..limit(1))
.getSingleOrNull();
return row == null ? null : _toTrip(row);
}
Stream<List<domain.Trip>> watchCompletedTrips() => (select(trips)
..where((t) => t.endedAt.isNotNull())
..orderBy([(t) => OrderingTerm.desc(t.startedAt)]))
.watch()
.map((rows) => rows.map(_toTrip).toList());
Stream<domain.Trip?> watchTrip(int id) =>
(select(trips)..where((t) => t.id.equals(id)))
.watchSingleOrNull()
.map((r) => r == null ? null : _toTrip(r));
Future<domain.Trip?> getTrip(int id) async {
final row =
await (select(trips)..where((t) => t.id.equals(id))).getSingleOrNull();
return row == null ? null : _toTrip(row);
}
Future<int> insertTrip(domain.Trip trip) => into(trips).insert(
TripsCompanion.insert(
startedAt: trip.startedAt,
endedAt: Value(trip.endedAt),
name: Value(trip.name),
state: trip.state,
distanceM: Value(trip.distanceM),
movingMillis: Value(trip.movingMillis),
maxSpeedKmh: Value(trip.maxSpeedKmh),
elevationGainM: Value(trip.elevationGainM),
pointCount: Value(trip.pointCount),
),
);
Future<void> renameTrip(int id, String? name) =>
(update(trips)..where((t) => t.id.equals(id)))
.write(TripsCompanion(name: Value(name)));
Future<void> setTripState(int id, domain.TripState state) =>
(update(trips)..where((t) => t.id.equals(id)))
.write(TripsCompanion(state: Value(state)));
Future<void> closeTrip(int id, int endedAt,
{domain.TripState state = domain.TripState.completed}) =>
(update(trips)..where((t) => t.id.equals(id))).write(
TripsCompanion(endedAt: Value(endedAt), state: Value(state)));
/// Persists the running totals. Called once per writer flush (~every 2 s), so it stays
/// a narrow targeted update rather than a full row rewrite.
Future<void> updateAggregates({
required int id,
required double distanceM,
required int movingMillis,
required double maxSpeedKmh,
required double elevationGainM,
required int pointCount,
}) =>
(update(trips)..where((t) => t.id.equals(id))).write(TripsCompanion(
distanceM: Value(distanceM),
movingMillis: Value(movingMillis),
maxSpeedKmh: Value(maxSpeedKmh),
elevationGainM: Value(elevationGainM),
pointCount: Value(pointCount),
));
/// Segments and points go with it via CASCADE.
Future<void> deleteTrip(int id) =>
(delete(trips)..where((t) => t.id.equals(id))).go();
Future<int> countTrips() async =>
(await select(trips).get()).length;
/// Test/maintenance helper. Segments and points follow via CASCADE.
Future<void> deleteAllTrips() => delete(trips).go();
// --- Segments ------------------------------------------------------------
Future<int> insertSegment(int tripId, int startedAt) =>
into(segments).insert(
SegmentsCompanion.insert(tripId: tripId, startedAt: startedAt),
);
Future<List<domain.Segment>> segmentsForTrip(int tripId) async {
final rows = await (select(segments)
..where((s) => s.tripId.equals(tripId))
..orderBy([(s) => OrderingTerm.asc(s.id)]))
.get();
return rows.map(_toSegment).toList();
}
/// The segment currently being recorded into, if any.
Future<domain.Segment?> openSegment(int tripId) async {
final row = await (select(segments)
..where((s) => s.tripId.equals(tripId) & s.endedAt.isNull())
..orderBy([(s) => OrderingTerm.desc(s.id)])
..limit(1))
.getSingleOrNull();
return row == null ? null : _toSegment(row);
}
Future<void> closeSegment(int id, int endedAt) =>
(update(segments)..where((s) => s.id.equals(id)))
.write(SegmentsCompanion(endedAt: Value(endedAt)));
/// Used by merge: re-parents a trip's segments onto the surviving trip.
Future<void> reparentSegments(int oldTripId, int newTripId) =>
(update(segments)..where((s) => s.tripId.equals(oldTripId)))
.write(SegmentsCompanion(tripId: Value(newTripId)));
Future<int> countSegmentsForTrip(int tripId) async =>
(await (select(segments)..where((s) => s.tripId.equals(tripId))).get())
.length;
// --- Track points --------------------------------------------------------
/// Batched insert — the recorder buffers points and flushes them in groups.
Future<void> insertPoints(List<domain.TrackPoint> points) async {
if (points.isEmpty) return;
await batch((b) {
b.insertAll(
trackPoints,
points.map((p) => TrackPointsCompanion.insert(
tripId: p.tripId,
segmentId: p.segmentId,
timestamp: p.timestamp,
latitude: p.latitude,
longitude: p.longitude,
speedKmh: p.speedKmh,
altitudeM: p.altitudeM,
accuracyM: Value(p.accuracyM),
bearingDeg: Value(p.bearingDeg),
synced: Value(p.synced),
)),
);
});
}
Future<int> insertPoint(domain.TrackPoint p) =>
into(trackPoints).insert(TrackPointsCompanion.insert(
tripId: p.tripId,
segmentId: p.segmentId,
timestamp: p.timestamp,
latitude: p.latitude,
longitude: p.longitude,
speedKmh: p.speedKmh,
altitudeM: p.altitudeM,
accuracyM: Value(p.accuracyM),
bearingDeg: Value(p.bearingDeg),
synced: Value(p.synced),
));
/// Ordered by segment then id so consumers walk the ride in recording order with pause
/// boundaries intact. **Not** ordered by timestamp: that value is GPS-derived and can
/// jump, whereas id is monotonic in write order.
Future<List<domain.TrackPoint>> pointsForTrip(int tripId) async {
final rows = await (select(trackPoints)
..where((p) => p.tripId.equals(tripId))
..orderBy([
(p) => OrderingTerm.asc(p.segmentId),
(p) => OrderingTerm.asc(p.id),
]))
.get();
return rows.map(_toPoint).toList();
}
Future<List<domain.TrackPoint>> pointsForSegment(int segmentId) async {
final rows = await (select(trackPoints)
..where((p) => p.segmentId.equals(segmentId))
..orderBy([(p) => OrderingTerm.asc(p.id)]))
.get();
return rows.map(_toPoint).toList();
}
/// The anchor a restarted recorder needs to continue accumulating distance.
Future<domain.TrackPoint?> lastInSegment(int segmentId) async {
final row = await (select(trackPoints)
..where((p) => p.segmentId.equals(segmentId))
..orderBy([(p) => OrderingTerm.desc(p.id)])
..limit(1))
.getSingleOrNull();
return row == null ? null : _toPoint(row);
}
Future<int> countPointsForTrip(int tripId) async =>
(await (select(trackPoints)..where((p) => p.tripId.equals(tripId))).get())
.length;
/// Live stats for one trip.
///
/// The COALESCE guards in the Kotlin original existed because Room throws on null for
/// non-null fields; they are kept because the semantics are what matter — a trip with
/// no points yet must read as zeros, not as an error.
Stream<domain.RideStats> watchTripStats(int tripId) {
final q = customSelect(
'''
SELECT COUNT(*) AS pointCount,
COALESCE(MAX(speedKmh), 0.0) AS maxSpeedKmh,
COALESCE(AVG(speedKmh), 0.0) AS avgSpeedKmh,
COALESCE(MIN(timestamp), 0) AS firstTimestamp,
COALESCE(MAX(timestamp), 0) AS lastTimestamp,
COALESCE(SUM(CASE WHEN synced = 0 THEN 1 ELSE 0 END), 0) AS pendingUpload
FROM track_points
WHERE tripId = ?
''',
variables: [Variable.withInt(tripId)],
readsFrom: {trackPoints},
);
return q.watchSingle().map((row) => domain.RideStats(
pointCount: row.read<int>('pointCount'),
maxSpeedKmh: row.read<double>('maxSpeedKmh'),
avgSpeedKmh: row.read<double>('avgSpeedKmh'),
firstTimestamp: row.read<int>('firstTimestamp'),
lastTimestamp: row.read<int>('lastTimestamp'),
pendingUpload: row.read<int>('pendingUpload'),
));
}
// --- Upload backlog ------------------------------------------------------
// Batches are drawn by id and can straddle a segment or trip boundary, which is why
// the payload carries trip/segment identity per point rather than per batch.
Future<List<domain.TrackPoint>> unsyncedPoints(int limit) async {
final rows = await (select(trackPoints)
..where((p) => p.synced.equals(false))
..orderBy([(p) => OrderingTerm.asc(p.id)])
..limit(limit))
.get();
return rows.map(_toPoint).toList();
}
Future<void> markSynced(List<int> ids) async {
if (ids.isEmpty) return;
await (update(trackPoints)..where((p) => p.id.isIn(ids)))
.write(const TrackPointsCompanion(synced: Value(true)));
}
Future<int> countUnsynced() async =>
(await (select(trackPoints)..where((p) => p.synced.equals(false))).get())
.length;
// --- Maintenance ---------------------------------------------------------
/// Used by merge: points carry tripId directly, so they re-parent alongside segments.
Future<void> reparentPoints(int oldTripId, int newTripId) =>
(update(trackPoints)..where((p) => p.tripId.equals(oldTripId)))
.write(TrackPointsCompanion(tripId: Value(newTripId)));
Future<List<domain.TrackPoint>> allPoints() async {
final rows = await (select(trackPoints)
..orderBy([(p) => OrderingTerm.asc(p.id)]))
.get();
return rows.map(_toPoint).toList();
}
}
// --- Row → domain mapping ---------------------------------------------------
// Kept as free functions rather than extension getters so the domain layer stays
// entirely unaware that Drift exists.
domain.Trip _toTrip(TripRow r) => domain.Trip(
id: r.id,
startedAt: r.startedAt,
endedAt: r.endedAt,
name: r.name,
state: r.state,
distanceM: r.distanceM,
movingMillis: r.movingMillis,
maxSpeedKmh: r.maxSpeedKmh,
elevationGainM: r.elevationGainM,
pointCount: r.pointCount,
);
domain.Segment _toSegment(SegmentRow r) => domain.Segment(
id: r.id,
tripId: r.tripId,
startedAt: r.startedAt,
endedAt: r.endedAt,
);
domain.TrackPoint _toPoint(TrackPointRow r) => domain.TrackPoint(
id: r.id,
tripId: r.tripId,
segmentId: r.segmentId,
timestamp: r.timestamp,
latitude: r.latitude,
longitude: r.longitude,
speedKmh: r.speedKmh,
altitudeM: r.altitudeM,
accuracyM: r.accuracyM,
bearingDeg: r.bearingDeg,
synced: r.synced,
);