diff --git a/mobile/lib/features/channels/channels_provider.dart b/mobile/lib/features/channels/channels_provider.dart index e44464585..9d49f6936 100644 --- a/mobile/lib/features/channels/channels_provider.dart +++ b/mobile/lib/features/channels/channels_provider.dart @@ -47,6 +47,7 @@ class ChannelsNotifier extends AsyncNotifier> { String? _subscriptionRelayBaseUrl; Timer? _backstopTimer; Timer? _unreadCatchUpTimer; + int _unreadCatchUpGeneration = 0; final Map _latestObservedByChannel = {}; final Map> _observedUnreadEventsByChannel = {}; @@ -525,8 +526,12 @@ class ChannelsNotifier extends AsyncNotifier> { if (ref.read(relayConfigProvider).baseUrl != relayBaseUrl) return; final channelIds = _desiredLiveChannelIds; + var didEstablishSubscription = false; if (_subscriptionRelayBaseUrl != relayBaseUrl || !_sameStringSet(_subscribedLiveChannelIds, channelIds)) { + _unreadCatchUpGeneration++; + _unreadCatchUpTimer?.cancel(); + _unreadCatchUpTimer = null; _liveUnsubscribe?.call(); _liveUnsubscribe = null; _subscribedLiveChannelIds = const {}; @@ -552,6 +557,7 @@ class ChannelsNotifier extends AsyncNotifier> { } _liveUnsubscribe = unsubscribe; _subscribedLiveChannelIds = Set.unmodifiable(channelIds); + didEstablishSubscription = true; } catch (error) { debugPrint('[ChannelsNotifier] live subscription failed: $error'); } @@ -563,14 +569,18 @@ class ChannelsNotifier extends AsyncNotifier> { return; } - _unreadCatchUpTimer?.cancel(); - _unreadCatchUpTimer = Timer(_unreadCatchUpDelay, () { - if (subscriptionVersion != _subscriptionVersion || - ref.read(relaySessionProvider).status != SessionStatus.connected) { - return; - } - unawaited(_catchUpUnreadEvents(channels)); - }); + if (didEstablishSubscription) { + final catchUpGeneration = ++_unreadCatchUpGeneration; + _unreadCatchUpTimer?.cancel(); + _unreadCatchUpTimer = Timer(_unreadCatchUpDelay, () { + if (catchUpGeneration != _unreadCatchUpGeneration || + subscriptionVersion != _subscriptionVersion || + ref.read(relaySessionProvider).status != SessionStatus.connected) { + return; + } + unawaited(_catchUpUnreadEvents(channels, catchUpGeneration)); + }); + } _backstopTimer?.cancel(); _backstopTimer = Timer.periodic( @@ -579,7 +589,10 @@ class ChannelsNotifier extends AsyncNotifier> { ); } - Future _catchUpUnreadEvents(List channels) async { + Future _catchUpUnreadEvents( + List channels, + int catchUpGeneration, + ) async { final myPk = ref.read(myPubkeyProvider); if (myPk == null) return; @@ -597,7 +610,6 @@ class ChannelsNotifier extends AsyncNotifier> { for (final channel in channels) { if (!channel.isMember || channel.isArchived) continue; final readAt = readState.effectiveTimestamp(channel.id); - if (readAt == null) continue; tasks.add( () => _catchUpUnreadEventsForChannel( session, @@ -610,7 +622,8 @@ class ChannelsNotifier extends AsyncNotifier> { } for (final task in tasks) { - if (ref.read(relaySessionProvider).status != SessionStatus.connected) { + if (catchUpGeneration != _unreadCatchUpGeneration || + ref.read(relaySessionProvider).status != SessionStatus.connected) { return; } await task(); @@ -850,6 +863,7 @@ class ChannelsNotifier extends AsyncNotifier> { _backstopTimer = null; _unreadCatchUpTimer?.cancel(); _unreadCatchUpTimer = null; + _unreadCatchUpGeneration++; } } diff --git a/mobile/test/features/channels/channels_provider_test.dart b/mobile/test/features/channels/channels_provider_test.dart index 83f9a7863..f3a07f421 100644 --- a/mobile/test/features/channels/channels_provider_test.dart +++ b/mobile/test/features/channels/channels_provider_test.dart @@ -548,6 +548,51 @@ void main() { }, ); + test('unchanged refresh does not restart unread catch-up', () async { + final session = _FakeRelaySession( + memberships: [_membership(_channelA, myPk)], + metadata: [_meta(id: _channelA, name: 'general')], + ); + final container = _buildContainer(session: session); + addTearDown(container.dispose); + + await container.read(channelsProvider.future); + await Future.delayed(const Duration(milliseconds: 2100)); + expect(session.unreadCatchUpQueryCount, 1); + + await container.read(channelsProvider.notifier).refresh(); + await Future.delayed(const Duration(milliseconds: 2100)); + + expect(session.unreadCatchUpQueryCount, 1); + }); + + test('changed channel set starts one replacement unread catch-up', () async { + final session = _FakeRelaySession( + memberships: [_membership(_channelA, myPk)], + metadata: [_meta(id: _channelA, name: 'general')], + ); + final container = _buildContainer(session: session); + addTearDown(container.dispose); + + await container.read(channelsProvider.future); + await Future.delayed(const Duration(milliseconds: 2100)); + expect(session.unreadCatchUpQueryCount, 1); + session.memberships = [ + _membership(_channelA, myPk), + _membership(_channelB, myPk), + ]; + session.metadata = [ + _meta(id: _channelA, name: 'general'), + _meta(id: _channelB, name: 'random'), + ]; + + await container.read(channelsProvider.notifier).refresh(); + await Future.delayed(const Duration(milliseconds: 2400)); + + expect(session.unreadCatchUpQueryCount, 3); + expect(session.totalSubscribeCount, 2); + }); + test('initial fetch issues membership + metadata queries', () async { final session = _FakeRelaySession( memberships: [_membership(_channelA, myPk)], @@ -677,6 +722,15 @@ class _FakeRelaySession extends RelaySessionNotifier { int unsubscribeCount = 0; int totalSubscribeCount = 0; + int get unreadCatchUpQueryCount => historyFilters + .where( + (filter) => + filter.tags['#h']?.length == 1 && + filter.kinds.length == EventKind.channelMessageEventKinds.length && + filter.kinds.every(EventKind.channelMessageEventKinds.contains), + ) + .length; + Set get activeChannels => { for (final (filter, _) in _subscriptions.values) ...?filter.tags['#h'], }; @@ -744,11 +798,7 @@ class _FakeRelaySession extends RelaySessionNotifier { List filters, { Duration timeout = const Duration(seconds: 8), }) async { - final events = []; - for (final filter in filters) { - events.addAll(await fetchHistory(filter, timeout: timeout)); - } - return events; + return const []; } @override