diff --git a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_document_reference.dart b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_document_reference.dart index e957d4c63fbb..4eeec6a5ee4f 100644 --- a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_document_reference.dart +++ b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_document_reference.dart @@ -118,8 +118,10 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform { controller; // ignore: close_sinks StreamSubscription? snapshotStreamSubscription; + var listenGeneration = 0; controller = StreamController.broadcast( onListen: () async { + final generation = ++listenGeneration; final observerId = await MethodChannelFirebaseFirestore.pigeonChannel .documentReferenceSnapshot( pigeonApp, @@ -130,6 +132,7 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform { includeMetadataChanges, listenSource, ); + if (generation != listenGeneration) return; snapshotStreamSubscription = MethodChannelFirebaseFirestore.documentSnapshotChannel(observerId) .receiveGuardedBroadcastStream( @@ -154,7 +157,9 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform { ); }, onCancel: () { + listenGeneration++; snapshotStreamSubscription?.cancel(); + snapshotStreamSubscription = null; }, ); diff --git a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_firestore.dart b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_firestore.dart index a29d7bdb725b..3d8edf9b1bce 100644 --- a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_firestore.dart +++ b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_firestore.dart @@ -187,10 +187,13 @@ class MethodChannelFirebaseFirestore extends FirebaseFirestorePlatform { Stream snapshotsInSync() { StreamSubscription? snapshotStreamSubscription; late StreamController controller; // ignore: close_sinks + var listenGeneration = 0; controller = StreamController.broadcast( onListen: () async { + final generation = ++listenGeneration; final observerId = await pigeonChannel.snapshotsInSyncSetup(pigeonApp); + if (generation != listenGeneration) return; snapshotStreamSubscription = MethodChannelFirebaseFirestore.snapshotsInSyncChannel(observerId) @@ -203,7 +206,9 @@ class MethodChannelFirebaseFirestore extends FirebaseFirestorePlatform { ); }, onCancel: () { + listenGeneration++; snapshotStreamSubscription?.cancel(); + snapshotStreamSubscription = null; }, ); diff --git a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_query.dart b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_query.dart index 12604b472d7c..33bd72799b0d 100644 --- a/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_query.dart +++ b/packages/cloud_firestore/cloud_firestore_platform_interface/lib/src/method_channel/method_channel_query.dart @@ -161,9 +161,11 @@ class MethodChannelQuery extends QueryPlatform { controller; // ignore: close_sinks StreamSubscription? snapshotStreamSubscription; + var listenGeneration = 0; controller = StreamController.broadcast( onListen: () async { + final generation = ++listenGeneration; final observerId = await MethodChannelFirebaseFirestore.pigeonChannel.querySnapshot( pigeonApp, @@ -177,6 +179,7 @@ class MethodChannelQuery extends QueryPlatform { includeMetadataChanges, listenSource, ); + if (generation != listenGeneration) return; snapshotStreamSubscription = MethodChannelFirebaseFirestore.querySnapshotChannel(observerId) @@ -195,7 +198,9 @@ class MethodChannelQuery extends QueryPlatform { ); }, onCancel: () { + listenGeneration++; snapshotStreamSubscription?.cancel(); + snapshotStreamSubscription = null; }, ); diff --git a/packages/cloud_firestore/cloud_firestore_platform_interface/test/method_channel_firestore_test.dart b/packages/cloud_firestore/cloud_firestore_platform_interface/test/method_channel_firestore_test.dart index 2b58d1d4f8d0..d8954f69d2a3 100644 --- a/packages/cloud_firestore/cloud_firestore_platform_interface/test/method_channel_firestore_test.dart +++ b/packages/cloud_firestore/cloud_firestore_platform_interface/test/method_channel_firestore_test.dart @@ -18,6 +18,47 @@ class _MockFirebaseFirestoreHostApi extends Mock implements TestFirebaseFirestoreHostApi { final Completer storeResultCalled = Completer(); final Completer releaseStoreResult = Completer(); + final Completer snapshotObserverId = Completer(); + final Completer snapshotsInSyncObserverId = Completer(); + + /// When non-empty, each `documentReferenceSnapshot`/`querySnapshot` call + /// consumes the next completer from the front of the queue instead of the + /// shared [snapshotObserverId]. Used to hand distinct observer ids to + /// successive (re-)listen attempts. + final List> queuedSnapshotObserverIds = + >[]; + + Future _nextSnapshotObserverId() { + if (queuedSnapshotObserverIds.isNotEmpty) { + return queuedSnapshotObserverIds.removeAt(0).future; + } + return snapshotObserverId.future; + } + + @override + Future documentReferenceSnapshot( + FirestorePigeonFirebaseApp app, + DocumentReferenceRequest parameters, + bool includeMetadataChanges, + ListenSource source, + ) => + _nextSnapshotObserverId(); + + @override + Future snapshotsInSyncSetup(FirestorePigeonFirebaseApp app) => + snapshotsInSyncObserverId.future; + + @override + Future querySnapshot( + FirestorePigeonFirebaseApp app, + String path, + bool isCollectionGroup, + InternalQueryParameters parameters, + InternalGetOptions options, + bool includeMetadataChanges, + ListenSource source, + ) => + _nextSnapshotObserverId(); @override Future transactionCreate( @@ -135,4 +176,161 @@ void main() { await Future.delayed(Duration.zero); }, ); + + group('snapshot listener cancelled while it registers', () { + const observerId = 'observer-id'; + const channels = [ + 'plugins.flutter.io/firebase_firestore/document/$observerId', + 'plugins.flutter.io/firebase_firestore/query/$observerId', + ]; + late List eventChannelCalls; + + setUp(() { + eventChannelCalls = []; + for (final channel in channels) { + messenger.setMockMessageHandler(channel, (ByteData? message) async { + eventChannelCalls.add(codec.decodeMethodCall(message).method); + return codec.encodeSuccessEnvelope(null); + }); + } + }); + + tearDown(() { + for (final channel in channels) { + messenger.setMockMessageHandler(channel, null); + } + }); + + Future expectNoNativeListen(Stream stream) async { + await stream.listen((_) {}).cancel(); + hostApi.snapshotObserverId.complete(observerId); + await pumpEventQueue(); + expect(eventChannelCalls, isEmpty); + } + + late MethodChannelFirebaseFirestore firestore; + + setUp(() { + firestore = MethodChannelFirebaseFirestore( + app: app, + databaseId: '(default)', + ); + }); + + test('DocumentReference.snapshots() does not attach a native listener', + () async { + await expectNoNativeListen( + firestore + .doc('foo/bar') + .snapshots(listenSource: ListenSource.defaultSource), + ); + }); + + test('Query.snapshots() does not attach a native listener', () async { + await expectNoNativeListen( + firestore + .collection('foo') + .snapshots(listenSource: ListenSource.defaultSource), + ); + }); + + test('a listener that is not cancelled still attaches', () async { + final subscription = firestore + .doc('foo/bar') + .snapshots(listenSource: ListenSource.defaultSource) + .listen((_) {}); + hostApi.snapshotObserverId.complete(observerId); + await pumpEventQueue(); + expect(eventChannelCalls, ['listen']); + await subscription.cancel(); + await pumpEventQueue(); + expect(eventChannelCalls, ['listen', 'cancel']); + }); + + test( + 're-listen while the first pigeon call is still pending attaches only ' + 'the second listener', + () async { + const firstObserverId = 'observer-id-1'; + const secondObserverId = 'observer-id-2'; + + final firstObserver = Completer(); + final secondObserver = Completer(); + hostApi.queuedSnapshotObserverIds + ..add(firstObserver) + ..add(secondObserver); + + final relistenChannels = >{ + 'plugins.flutter.io/firebase_firestore/document/$firstObserverId': + [], + 'plugins.flutter.io/firebase_firestore/document/$secondObserverId': + [], + }; + relistenChannels.forEach((channel, calls) { + messenger.setMockMessageHandler(channel, (ByteData? message) async { + calls.add(codec.decodeMethodCall(message).method); + return codec.encodeSuccessEnvelope(null); + }); + }); + addTearDown(() { + for (final channel in relistenChannels.keys) { + messenger.setMockMessageHandler(channel, null); + } + }); + + final stream = firestore + .doc('foo/bar') + .snapshots(listenSource: ListenSource.defaultSource); + + // First listen: the pigeon setup future stays pending. + await stream.listen((_) {}).cancel(); + // Re-listen while the first setup future is still pending. + final secondSubscription = stream.listen((_) {}); + + // Complete the first (stale/cancelled) observer id, then the second. + firstObserver.complete(firstObserverId); + await pumpEventQueue(); + secondObserver.complete(secondObserverId); + await pumpEventQueue(); + + // The stale first completion must not attach; only the second channel + // receives the native `listen`. + expect( + relistenChannels[ + 'plugins.flutter.io/firebase_firestore/document/$firstObserverId'], + isEmpty, + ); + expect( + relistenChannels[ + 'plugins.flutter.io/firebase_firestore/document/$secondObserverId'], + ['listen'], + ); + + await secondSubscription.cancel(); + }, + ); + + test( + 'snapshotsInSync() does not attach a native listener when cancelled ' + 'before setup completes', + () async { + const syncObserverId = 'sync-observer-id'; + const syncChannel = + 'plugins.flutter.io/firebase_firestore/snapshotsInSync/$syncObserverId'; + final syncChannelCalls = []; + messenger.setMockMessageHandler(syncChannel, (ByteData? message) async { + syncChannelCalls.add(codec.decodeMethodCall(message).method); + return codec.encodeSuccessEnvelope(null); + }); + addTearDown(() => messenger.setMockMessageHandler(syncChannel, null)); + + // Cancel before the snapshotsInSyncSetup future completes. + await firestore.snapshotsInSync().listen((_) {}).cancel(); + hostApi.snapshotsInSyncObserverId.complete(syncObserverId); + await pumpEventQueue(); + + expect(syncChannelCalls, isEmpty); + }, + ); + }); }