Forward stream queries in DatabaseConnection constructor

This commit is contained in:
Simon Binder
2025-10-07 14:45:03 +02:00
parent 434bde3e87
commit f01899e049
2 changed files with 27 additions and 2 deletions
+5 -2
View File
@@ -34,8 +34,11 @@ class DatabaseConnection implements QueryExecutor {
this.connectionData,
bool closeStreamsSynchronously = false,
}) : streamQueries = streamQueries ??
StreamQueryStore(
closeStreamsSynchronously: closeStreamsSynchronously);
switch (executor) {
DatabaseConnection() => executor.streamQueries,
_ => StreamQueryStore(
closeStreamsSynchronously: closeStreamsSynchronously)
};
/// Constructs a [DatabaseConnection] from the [QueryExecutor] by using the
/// default type system and a new [StreamQueryStore].
@@ -2,9 +2,11 @@
@TestOn('vm')
library;
import 'package:async/async.dart';
import 'package:drift/drift.dart';
import 'package:drift/native.dart';
import 'package:drift/src/runtime/cancellation_zone.dart';
import 'package:drift/src/runtime/executor/stream_queries.dart';
import 'package:sqlite3/sqlite3.dart';
import 'package:test/test.dart';
@@ -47,6 +49,26 @@ void main() {
driftDb.select(driftDb.categories).get(), completion(hasLength(1)));
});
test('DatabaseConnection constructor can wrap inner', () async {
final raw = NativeDatabase.memory();
final streams = StreamQueryStore();
final db = TodoDb(
DatabaseConnection(DatabaseConnection(raw, streamQueries: streams)));
await db
.into(db.categories)
.insert(CategoriesCompanion.insert(description: 'description'));
final query = StreamQueue(db.categories.all().watch());
await expectLater(query, emits(hasLength(1)));
await raw.runCustom('DELETE FROM categories');
streams.handleTableUpdates({TableUpdate('categories')});
await expectLater(query, emits(isEmpty));
await query.cancel();
await db.close();
});
group('nested transactions', () {
test(
'outer transaction does not see inner writes after rollback',