Close streams before database

This commit is contained in:
Simon Binder
2025-11-18 20:51:16 +05:30
committed by shenlong-tanwen
parent 7295a3a17a
commit f1b58a71c3
5 changed files with 25 additions and 6 deletions
@@ -178,7 +178,7 @@ class StreamQueryStore {
_isShuttingDown = true;
for (final stream in _activeKeyStreams.values) {
stream.close();
await stream.close();
}
// awaiting this is fine - the stream is never exposed to users and we don't
// pause any subscriptions on it.
@@ -209,7 +209,7 @@ class QueryStream<Rows extends Object> {
final queryListener = _QueryStreamListener(listener);
if (_isClosed) {
listener.closeSync();
listener.close();
return;
}
@@ -358,12 +358,13 @@ class QueryStream<Rows extends Object> {
}
}
void close() {
Future<void> close() async {
_isClosed = true;
for (final listener in _listeners) {
listener.controller.close();
}
final listenersDone = Future.wait(
[for (final listener in _listeners) listener.controller.close()]);
_listeners.clear();
await listenersDone;
}
}
+15
View File
@@ -234,6 +234,21 @@ void main() {
subscription.resume();
await subscription.cancel();
}, skip: 'testing out awaited streams');
test('closing database waits for streams', () async {
final stream = db.select(db.users).watch();
final subscription = stream.listen((_) {})..pause();
var closed = false;
db.close().then((_) => closed = true);
await pumpEventQueue();
expect(closed, isFalse);
subscription.resume();
await subscription.cancel();
await pumpEventQueue();
expect(closed, isTrue);
});
group('stream keys', () {
@@ -92,6 +92,7 @@ void main() {
test('can be used in a query stream', () async {
final stream = StreamQueue(db.readView().watch());
addTearDown(stream.cancel);
const entry = Config(
configKey: 'another_key',
configValue: DriftAny('value'),
@@ -12,6 +12,7 @@ void main() {
addTearDown(db.close);
final query = StreamQueue(db.select(db.myView).watch());
addTearDown(query.cancel);
await expectLater(query, emits(isEmpty));
await db.into(db.config).insert(ConfigCompanion.insert(
+1
View File
@@ -243,6 +243,7 @@ void main() {
);
});
stream.cancel();
await db.close();
}