fix(mobile): deduplicate unread catch-up lifecycle

Only schedule unread reconciliation after establishing a new aggregate
subscription, cancel superseded passes, and preserve historical unread for
channels without a read marker.

Co-authored-by: Carl <c7ebe626f000404285d3686e1dc74cc07cc60a9754a150041ba132e14bd3e2ec@buzz.block.builderlab.xyz>
Signed-off-by: Wes <wesbillman@users.noreply.github.com>
This commit is contained in:
Wes
2026-08-13 16:14:02 -06:00
co-authored by Carl
parent a67440977d
commit 413d8f4d86
2 changed files with 80 additions and 16 deletions
@@ -47,6 +47,7 @@ class ChannelsNotifier extends AsyncNotifier<List<Channel>> {
String? _subscriptionRelayBaseUrl;
Timer? _backstopTimer;
Timer? _unreadCatchUpTimer;
int _unreadCatchUpGeneration = 0;
final Map<String, int> _latestObservedByChannel = {};
final Map<String, Map<String, ObservedUnreadEvent>>
_observedUnreadEventsByChannel = {};
@@ -525,8 +526,12 @@ class ChannelsNotifier extends AsyncNotifier<List<Channel>> {
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<List<Channel>> {
}
_liveUnsubscribe = unsubscribe;
_subscribedLiveChannelIds = Set.unmodifiable(channelIds);
didEstablishSubscription = true;
} catch (error) {
debugPrint('[ChannelsNotifier] live subscription failed: $error');
}
@@ -563,14 +569,18 @@ class ChannelsNotifier extends AsyncNotifier<List<Channel>> {
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<List<Channel>> {
);
}
Future<void> _catchUpUnreadEvents(List<Channel> channels) async {
Future<void> _catchUpUnreadEvents(
List<Channel> channels,
int catchUpGeneration,
) async {
final myPk = ref.read(myPubkeyProvider);
if (myPk == null) return;
@@ -597,7 +610,6 @@ class ChannelsNotifier extends AsyncNotifier<List<Channel>> {
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<List<Channel>> {
}
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<List<Channel>> {
_backstopTimer = null;
_unreadCatchUpTimer?.cancel();
_unreadCatchUpTimer = null;
_unreadCatchUpGeneration++;
}
}
@@ -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<void>.delayed(const Duration(milliseconds: 2100));
expect(session.unreadCatchUpQueryCount, 1);
await container.read(channelsProvider.notifier).refresh();
await Future<void>.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<void>.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<void>.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<String> get activeChannels => {
for (final (filter, _) in _subscriptions.values) ...?filter.tags['#h'],
};
@@ -744,11 +798,7 @@ class _FakeRelaySession extends RelaySessionNotifier {
List<NostrFilter> filters, {
Duration timeout = const Duration(seconds: 8),
}) async {
final events = <NostrEvent>[];
for (final filter in filters) {
events.addAll(await fetchHistory(filter, timeout: timeout));
}
return events;
return const [];
}
@override