Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -118,8 +118,10 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform {
controller; // ignore: close_sinks

StreamSubscription<dynamic>? snapshotStreamSubscription;
var listenGeneration = 0;
controller = StreamController<DocumentSnapshotPlatform>.broadcast(
onListen: () async {
final generation = ++listenGeneration;
final observerId = await MethodChannelFirebaseFirestore.pigeonChannel
.documentReferenceSnapshot(
pigeonApp,
Expand All @@ -130,6 +132,7 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform {
includeMetadataChanges,
listenSource,
);
if (generation != listenGeneration) return;
snapshotStreamSubscription =
MethodChannelFirebaseFirestore.documentSnapshotChannel(observerId)
.receiveGuardedBroadcastStream(
Expand All @@ -154,7 +157,9 @@ class MethodChannelDocumentReference extends DocumentReferencePlatform {
);
},
onCancel: () {
listenGeneration++;
snapshotStreamSubscription?.cancel();
snapshotStreamSubscription = null;
},
);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -187,10 +187,13 @@ class MethodChannelFirebaseFirestore extends FirebaseFirestorePlatform {
Stream<void> snapshotsInSync() {
StreamSubscription<dynamic>? snapshotStreamSubscription;
late StreamController<void> controller; // ignore: close_sinks
var listenGeneration = 0;

controller = StreamController<void>.broadcast(
onListen: () async {
final generation = ++listenGeneration;
final observerId = await pigeonChannel.snapshotsInSyncSetup(pigeonApp);
if (generation != listenGeneration) return;

snapshotStreamSubscription =
MethodChannelFirebaseFirestore.snapshotsInSyncChannel(observerId)
Expand All @@ -203,7 +206,9 @@ class MethodChannelFirebaseFirestore extends FirebaseFirestorePlatform {
);
},
onCancel: () {
listenGeneration++;
snapshotStreamSubscription?.cancel();
snapshotStreamSubscription = null;
},
);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -161,9 +161,11 @@ class MethodChannelQuery extends QueryPlatform {
controller; // ignore: close_sinks

StreamSubscription<dynamic>? snapshotStreamSubscription;
var listenGeneration = 0;

controller = StreamController<QuerySnapshotPlatform>.broadcast(
onListen: () async {
final generation = ++listenGeneration;
final observerId =
await MethodChannelFirebaseFirestore.pigeonChannel.querySnapshot(
pigeonApp,
Expand All @@ -177,6 +179,7 @@ class MethodChannelQuery extends QueryPlatform {
includeMetadataChanges,
listenSource,
);
if (generation != listenGeneration) return;

snapshotStreamSubscription =
MethodChannelFirebaseFirestore.querySnapshotChannel(observerId)
Expand All @@ -195,7 +198,9 @@ class MethodChannelQuery extends QueryPlatform {
);
},
onCancel: () {
listenGeneration++;
snapshotStreamSubscription?.cancel();
snapshotStreamSubscription = null;
},
);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,47 @@ class _MockFirebaseFirestoreHostApi extends Mock
implements TestFirebaseFirestoreHostApi {
final Completer<void> storeResultCalled = Completer<void>();
final Completer<void> releaseStoreResult = Completer<void>();
final Completer<String> snapshotObserverId = Completer<String>();
final Completer<String> snapshotsInSyncObserverId = Completer<String>();

/// 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<Completer<String>> queuedSnapshotObserverIds =
<Completer<String>>[];

Future<String> _nextSnapshotObserverId() {
if (queuedSnapshotObserverIds.isNotEmpty) {
return queuedSnapshotObserverIds.removeAt(0).future;
}
return snapshotObserverId.future;
}

@override
Future<String> documentReferenceSnapshot(
FirestorePigeonFirebaseApp app,
DocumentReferenceRequest parameters,
bool includeMetadataChanges,
ListenSource source,
) =>
_nextSnapshotObserverId();

@override
Future<String> snapshotsInSyncSetup(FirestorePigeonFirebaseApp app) =>
snapshotsInSyncObserverId.future;

@override
Future<String> querySnapshot(
FirestorePigeonFirebaseApp app,
String path,
bool isCollectionGroup,
InternalQueryParameters parameters,
InternalGetOptions options,
bool includeMetadataChanges,
ListenSource source,
) =>
_nextSnapshotObserverId();

@override
Future<String> transactionCreate(
Expand Down Expand Up @@ -135,4 +176,161 @@ void main() {
await Future<void>.delayed(Duration.zero);
},
);

group('snapshot listener cancelled while it registers', () {
const observerId = 'observer-id';
const channels = <String>[
'plugins.flutter.io/firebase_firestore/document/$observerId',
'plugins.flutter.io/firebase_firestore/query/$observerId',
];
late List<String> eventChannelCalls;

setUp(() {
eventChannelCalls = <String>[];
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<void> expectNoNativeListen(Stream<Object?> 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, <String>['listen']);
await subscription.cancel();
await pumpEventQueue();
expect(eventChannelCalls, <String>['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<String>();
final secondObserver = Completer<String>();
hostApi.queuedSnapshotObserverIds
..add(firstObserver)
..add(secondObserver);

final relistenChannels = <String, List<String>>{
'plugins.flutter.io/firebase_firestore/document/$firstObserverId':
<String>[],
'plugins.flutter.io/firebase_firestore/document/$secondObserverId':
<String>[],
};
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'],
<String>['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 = <String>[];
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);
},
);
});
}
Loading