diff --git a/mobile/lib/features/activity/activity_page.dart b/mobile/lib/features/activity/activity_page.dart index d38ac3e1952..31362b30b0a 100644 --- a/mobile/lib/features/activity/activity_page.dart +++ b/mobile/lib/features/activity/activity_page.dart @@ -253,8 +253,9 @@ class ActivityPage extends HookConsumerWidget { // Deep-link to the represented message: oldest unread in the group, // falling back to the latest event. - final readAt = resolveInboxItemReadAt(item, markerOf: markerOf); - final target = item.deepLinkTarget(readAt); + final target = item.deepLinkTarget( + (event) => isInboxEventRead(event, markerOf: markerOf), + ); final thread = threadReferenceOf(target.tags); final threadRootId = isBroadcastReply(target.tags) ? null diff --git a/mobile/lib/features/activity/inbox_item.dart b/mobile/lib/features/activity/inbox_item.dart index 3fe618fc54e..4910c73d083 100644 --- a/mobile/lib/features/activity/inbox_item.dart +++ b/mobile/lib/features/activity/inbox_item.dart @@ -113,20 +113,18 @@ class InboxItem { return null; } - /// The event the row should deep-link to: the oldest event in the group - /// newer than [readAt] (oldest unread), falling back to the latest event. - FeedItem deepLinkTarget(int? readAt) { - if (readAt != null) { - FeedItem? oldestUnread; - for (final candidate in groupItems) { - if (candidate.createdAt <= readAt) continue; - if (oldestUnread == null || - candidate.createdAt < oldestUnread.createdAt) { - oldestUnread = candidate; - } + /// The event the row should deep-link to: the oldest grouped event that + /// [isRead] rejects (oldest unread), falling back to the latest event. + FeedItem deepLinkTarget(bool Function(FeedItem event) isRead) { + FeedItem? oldestUnread; + for (final candidate in groupItems) { + if (isRead(candidate)) continue; + if (oldestUnread == null || + candidate.createdAt < oldestUnread.createdAt) { + oldestUnread = candidate; } - if (oldestUnread != null) return oldestUnread; } + if (oldestUnread != null) return oldestUnread; return item; } } diff --git a/mobile/lib/features/activity/inbox_read_state.dart b/mobile/lib/features/activity/inbox_read_state.dart index 8626b574f55..b92e55fdb50 100644 --- a/mobile/lib/features/activity/inbox_read_state.dart +++ b/mobile/lib/features/activity/inbox_read_state.dart @@ -1,30 +1,48 @@ import '../../shared/read_state/read_state_format.dart'; +import 'feed_item.dart'; import 'inbox_item.dart'; -/// Resolves the effective NIP-RS read marker for one inbox row, mirroring -/// desktop's `resolveInboxItemReadAt`: -/// - thread rows use `max(thread:, msg:)` -/// - channel rows use the channel context marker -/// - rows with no channel have no marker (caller falls back to local state) -int? resolveInboxItemReadAt( - InboxItem item, { +/// The newest read marker that reads one grouped [event], following the +/// unread badge's `observedUnreadEventReadAt`: the channel mark, the +/// event's own `msg:` mark, and for a thread reply `thread:` and +/// `thread-activity:`. Returns null for an event with no channel. +/// +/// `activity:` never counts here. Activity rows hold only mentions, +/// needs-action and agent events addressed to the reader, and DM messages, +/// and that catch-up mark reads none of those. +int? inboxEventReadAt( + FeedItem event, { required int? Function(String contextId) markerOf, }) { - final channelId = item.item.channelId; - final threadRootId = item.threadRootId; - if (threadRootId != null) { - return maxReadAt([ - markerOf(threadContextKey(threadRootId)), - markerOf(msgContextKey(item.item.id)), - ]); - } - return channelId == null ? null : markerOf(channelId); + final channelId = event.channelId; + if (channelId == null) return null; + final rootId = isThreadReply(event.tags) + ? threadReferenceOf(event.tags).rootId + : null; + return maxReadAt([ + markerOf(channelId), + markerOf(msgContextKey(event.id)), + if (rootId != null) ...[ + markerOf(threadContextKey(rootId)), + markerOf(threadActivityContextKey(rootId)), + ], + ]); +} + +/// Whether the marks read one grouped [event]. +bool isInboxEventRead( + FeedItem event, { + required int? Function(String contextId) markerOf, +}) { + final readAt = inboxEventReadAt(event, markerOf: markerOf); + return readAt != null && event.createdAt <= readAt; } -/// Whether the row is read ("done"), mirroring desktop's -/// `useHomeInboxReadState` projection: a local unread override always wins; -/// otherwise the row is done when no grouped activity is newer than the -/// shared read marker; channel-less rows fall back to the local done set. +/// Whether the row is read ("done"): a local unread override always wins; +/// otherwise the row is done when the marks read every grouped event. Each +/// event is checked on its own, because reading a channel marks mentions +/// with per-message marks, so reading the newest says nothing about the +/// rest. Channel-less rows fall back to the local done set. bool isInboxItemDone( InboxItem item, { required int? Function(String contextId) markerOf, @@ -34,13 +52,16 @@ bool isInboxItemDone( final ids = groupedInboxItemIds(item); if (ids.any(localUnreadOverrides.contains)) return false; - final readAt = resolveInboxItemReadAt(item, markerOf: markerOf); - if (readAt != null) return item.latestActivityAt <= readAt; - - if (item.threadRootId != null || item.item.channelId != null) return false; - return localDoneSet.contains(item.id); + if (item.item.channelId == null) return localDoneSet.contains(item.id); + return _groupedEvents( + item, + ).every((event) => isInboxEventRead(event, markerOf: markerOf)); } +Iterable _groupedEvents(InboxItem item) => { + for (final event in [item.item, ...item.groupItems]) event.id: event, +}.values; + /// All event ids identified with the row — desktop's `getGroupedInboxItemIds`. List groupedInboxItemIds(InboxItem item) { return {item.id, item.item.id, ...item.groupItems.map((i) => i.id)}.toList(); diff --git a/mobile/lib/features/channels/channel_detail_page.dart b/mobile/lib/features/channels/channel_detail_page.dart index 5549d6af635..fbd68a707fe 100644 --- a/mobile/lib/features/channels/channel_detail_page.dart +++ b/mobile/lib/features/channels/channel_detail_page.dart @@ -78,6 +78,7 @@ import 'small_avatar.dart'; import 'sticky_date_header.dart'; import 'thread_detail_page.dart'; import 'thread_replies_provider.dart'; +import 'reading_marks.dart'; import 'timeline_message.dart'; part 'channel_detail_page/message_list.dart'; @@ -206,29 +207,21 @@ Future _subscribeToDmIdentityUpdates( } } -int? _channelReadTimestamp({ +/// The time that opening [channel] reads it through, or null when opening it +/// reads nothing. Forums read through their last activity. DMs read through +/// the newest loaded message, including replies, or the last activity before +/// messages load. +int? _openReadTimestamp({ required Channel channel, required AsyncValue> messagesState, }) { - if (channel.isForum) { - return dateTimeToUnixSeconds(channel.lastMessageAt); + if (channel.isForum) return dateTimeToUnixSeconds(channel.lastMessageAt); + if (!channel.isDm) return null; + var latest = 0; + for (final event in messagesState.value ?? const []) { + if (event.createdAt > latest) latest = event.createdAt; } - - final events = messagesState.value; - if (events != null && events.isNotEmpty) { - var latest = 0; - for (final event in events) { - if (event.threadReference.parentId != null) continue; - if (event.createdAt > latest) { - latest = event.createdAt; - } - } - if (latest > 0) { - return latest; - } - } - - return dateTimeToUnixSeconds(channel.lastMessageAt); + return latest > 0 ? latest : dateTimeToUnixSeconds(channel.lastMessageAt); } bool _isOneToOneAgentDm(Channel channel, Set agentPubkeys) { @@ -512,7 +505,11 @@ class ChannelDetailPage extends HookConsumerWidget { !messagesNotifier.hasLoadedMessages; final appBarTitleContentHeight = _twoLineAppBarTitleContentHeight(context); - final readTimestamp = _channelReadTimestamp( + // Opening a forum or a DM reads the whole channel: a forum has no + // timeline to read row by row, and a DM is all for the reader. Other + // timelines write their read marks while the reader looks at rows; see + // `readingMarks`. + final openReadTimestamp = _openReadTimestamp( channel: resolvedChannel, messagesState: messagesState, ); @@ -545,19 +542,25 @@ class ChannelDetailPage extends HookConsumerWidget { ], ); + // Opening reads only while this page is in front and the app is in use. + // A covered or backgrounded DM must not read a reply that loads meanwhile; + // it reads through it when it is shown again. + final openReadActive = + isAppInUse(useAppLifecycleState()) && + (ModalRoute.of(context)?.isCurrent ?? true); useEffect(() { - if (!readState.isReady || readTimestamp == null) { + if (!readState.isReady || !openReadActive || openReadTimestamp == null) { return null; } return deferReadStateUpdate(context, () { ref .read(readStateProvider.notifier) - .markContextRead(channel.id, readTimestamp); + .markContextRead(channel.id, openReadTimestamp); ref .read(channelsProvider.notifier) - .clearObservedUnreadCoveredByRead(channel.id, readTimestamp); + .clearObservedUnreadCoveredByRead(channel.id, openReadTimestamp); }); - }, [channel.id, readState.isReady, readTimestamp]); + }, [channel.id, readState.isReady, openReadActive, openReadTimestamp]); final dmHeader = resolvedChannel.isDm ? _watchDmHeader(ref, resolvedChannel, currentPubkey) @@ -766,6 +769,7 @@ class ChannelDetailPage extends HookConsumerWidget { initialOldestOrdinaryUnreadMessageId != null), channelId: channel.id, + isDm: resolvedChannel.isDm, currentPubkey: currentPubkey, isMember: resolvedChannel.isMember, isArchived: resolvedChannel.isArchived, diff --git a/mobile/lib/features/channels/channel_detail_page/message_list.dart b/mobile/lib/features/channels/channel_detail_page/message_list.dart index 1a2c3c82115..a5a661bdd3d 100644 --- a/mobile/lib/features/channels/channel_detail_page/message_list.dart +++ b/mobile/lib/features/channels/channel_detail_page/message_list.dart @@ -11,6 +11,7 @@ class _MessageList extends HookConsumerWidget { final Set initialForcedUnreadMessageIds; final bool hasInitialUnread; final String channelId; + final bool isDm; final String? currentPubkey; final bool isMember; final bool isArchived; @@ -30,6 +31,7 @@ class _MessageList extends HookConsumerWidget { required this.initialForcedUnreadMessageIds, required this.hasInitialUnread, required this.channelId, + required this.isDm, required this.currentPubkey, required this.isMember, required this.isArchived, @@ -539,6 +541,83 @@ class _MessageList extends HookConsumerWidget { ], ); + void readVisibleRows() { + final readState = ref.read(readStateProvider); + final viewportHeight = timelineViewportHeight.value; + if (!context.mounted || + !readState.isReady || + isAutoScrolling.value || + viewportHeight <= 0 || + entries.isEmpty) { + return; + } + // Rows hidden under the app bar, the composer or the keyboard are not + // read. The keyboard covers rows even when the list does not follow + // the latest message, so this uses the whole covered height, not the + // list's own bottom inset. + final bottomEdge = + (navigationBottomInset / viewportHeight).clamp(0.0, 1.0) - 0.01; + final topEdge = + 1 - + frostedAppBarHeight( + context, + titleContentHeight: appBarTitleContentHeight, + ) / + viewportHeight + + 0.01; + final visible = []; + for (final position in itemPositionsListener.itemPositions.value) { + if (position.index >= displayEntries.length || + position.itemLeadingEdge < bottomEdge || + position.itemTrailingEdge > topEdge) { + continue; + } + final group = + displayEntries[displayEntries.length - 1 - position.index]; + visible.addAll(group.map((entry) => entry.message)); + } + final latest = entries.last.message; + final atBottom = + latestIsAtBoundary() && + itemPositionsListener.itemPositions.value.any( + (position) => + position.index == 0 && position.itemLeadingEdge >= bottomEdge, + ); + final marks = readingMarks( + readState: readState, + channelId: channelId, + isDm: isDm, + currentPubkey: currentPubkey, + loaded: allMessages, + visible: visible, + bottom: atBottom ? latest : null, + ); + final notifier = ref.read(readStateProvider.notifier); + // Reading to the bottom ends a manual "mark unread" on the channel. + if (atBottom && readState.isForcedUnread(channelId)) { + notifier.clearForcedUnread(channelId); + } + for (final mark in marks.entries) { + notifier.markContextRead(mark.key, mark.value); + } + } + + // Reading starts once the timeline holds still for the dwell while this + // page is in front and the app is in use. + final appInUse = isAppInUse(useAppLifecycleState()); + final readStateReady = ref.watch( + readStateProvider.select((state) => state.isReady), + ); + useReadingDwell( + positions: itemPositionsListener.itemPositions, + active: + readStateReady && + appInUse && + (ModalRoute.of(context)?.isCurrent ?? true), + onDwell: readVisibleRows, + keys: [channelId, readingContentKey(allMessages), navigationBottomInset], + ); + useEffect(() { stickyDateHeaderState.value = StickyDateHeaderState.hidden; stickyDayTimestamp.value = null; diff --git a/mobile/lib/features/channels/reading_marks.dart b/mobile/lib/features/channels/reading_marks.dart new file mode 100644 index 00000000000..1ec4f0624ba --- /dev/null +++ b/mobile/lib/features/channels/reading_marks.dart @@ -0,0 +1,185 @@ +import 'dart:async'; + +import 'package:flutter_hooks/flutter_hooks.dart'; +import 'package:flutter/foundation.dart'; +import 'package:flutter/widgets.dart'; + +import '../../shared/read_state/read_state_format.dart'; +import '../../shared/read_state/read_state_provider.dart'; +import 'timeline_message.dart'; +import 'unread_badge/is_high_priority_event.dart'; + +/// How long a row must stay fully visible before it counts as read. This is +/// the same dwell the web and desktop app (`buzz-app`) uses. +const readingDwell = Duration(milliseconds: 300); + +/// The read marks to write after a reading dwell, following `buzz-app`'s +/// rules (its `docs/unread.md`): +/// +/// - When the newest message is at the bottom, a channel writes +/// `activity:` and a thread writes `thread-activity:`. These +/// catch-up marks read ordinary messages and that thread's replies. They do +/// not read mentions, broadcasts or DMs. +/// - Each fully visible message gets its own `msg:` mark unless its own mark, +/// the channel mark, or (for an ordinary top-level message) `activity:` +/// already reads it. A thread mark never counts here: a reply finds its +/// thread only while the root is loaded, so `buzz-app` keeps reply marks. +/// - Automatic reading never writes the channel mark. Only an explicit +/// Mark read does that. +/// +/// [loaded] is every message loaded for the channel or thread, including +/// replies. [visible] is the messages fully visible for the whole dwell. +/// [bottom] is the newest message when its bottom edge is visible at the +/// bottom of the list; otherwise null. Set [threadRootId] when reading a +/// thread, and [isRootThread] when its head is the thread root rather than a +/// nested reply. [now] is the current Unix time in seconds; it defaults to +/// the clock. The result maps each read-state context to its new time. +Map readingMarks({ + required ReadStateState readState, + required String channelId, + required bool isDm, + required String? currentPubkey, + required Iterable loaded, + required Iterable visible, + TimelineMessage? bottom, + String? threadRootId, + bool isRootThread = true, + int? now, +}) { + final marks = {}; + int? markerOf(String contextId) => + maxReadAt([marks[contextId], readState.effectiveTimestamp(contextId)]); + + // A nested thread view shows only one branch, so reaching its bottom does + // not read the whole thread. + if (bottom != null && (threadRootId == null || isRootThread)) { + final catchUp = _catchUpMark( + channelId: channelId, + loaded: loaded, + bottom: bottom, + threadRootId: threadRootId, + markerOf: markerOf, + now: now ?? DateTime.now().millisecondsSinceEpoch ~/ 1000, + ); + if (catchUp != null) marks[catchUp.key] = catchUp.value; + } + + final self = currentPubkey?.toLowerCase(); + for (final message in visible) { + if (message.isSystem || message.pubkey.toLowerCase() == self) continue; + final contextId = msgContextKey(message.id); + // Automatic reading keeps a message the reader marked unread. + if (readState.isForcedUnread(contextId)) continue; + final ordinary = readByChannelCatchUp( + isDm: isDm, + isReply: message.parentId != null && !_isBroadcast(message), + highPriority: self == null || isHighPriorityEvent(message.tags, self), + ); + final readAt = maxReadAt([ + markerOf(contextId), + markerOf(channelId), + if (ordinary) markerOf(activityContextKey(channelId)), + ]); + if (readAt != null && readAt >= message.createdAt) continue; + marks[contextId] = message.createdAt; + } + return marks; +} + +/// The catch-up mark for [bottom], or null when it would read nothing new. +MapEntry? _catchUpMark({ + required String channelId, + required Iterable loaded, + required TimelineMessage bottom, + required String? threadRootId, + required int? Function(String contextId) markerOf, + required int now, +}) { + // A loaded message newer than the bottom row means the list does not show + // the live bottom, for example after a jump into history. + final showsNewest = !loaded.any( + (message) => + message.createdAt > bottom.createdAt && + (threadRootId != null + ? _threadRootOf(message) == threadRootId + : message.parentId == null || _isBroadcast(message)), + ); + if (!showsNewest) return null; + + // Cut at the bottom row. A newer reply's time would also read a top-level + // message that arrives late with an earlier time. A clock ahead of now, on + // this device or the sender's, must not read messages before they arrive. + final cut = bottom.createdAt < now ? bottom.createdAt : now; + final key = threadRootId != null + ? threadActivityContextKey(threadRootId) + : activityContextKey(channelId); + final covered = maxReadAt([ + markerOf(key), + markerOf(channelId), + if (threadRootId != null) markerOf(threadContextKey(threadRootId)), + ]); + if (covered != null && covered >= cut) return null; + return MapEntry(key, cut); +} + +String? _threadRootOf(TimelineMessage message) => + message.parentId == null ? null : message.rootId ?? message.parentId; + +bool _isBroadcast(TimelineMessage message) => message.tags.any( + (tag) => tag.length >= 2 && tag[0] == 'broadcast' && tag[1] == '1', +); + +/// A stable dwell key for [messages]. Pages rebuild a new message list on +/// every frame that touches them, such as a typing indicator, so the list +/// itself would restart the dwell each time. This changes only when a +/// message arrives, leaves, or moves the newest time. +(int, int, String?) readingContentKey(Iterable messages) { + var count = 0; + TimelineMessage? newest; + for (final message in messages) { + count++; + if (newest == null || message.createdAt > newest.createdAt) { + newest = message; + } + } + return (count, newest?.createdAt ?? 0, newest?.id); +} + +/// Whether the app is in the foreground. Null means the platform has not +/// reported a state yet, which only happens before the first frame and in +/// tests. +bool isAppInUse(AppLifecycleState? state) => + state == null || state == AppLifecycleState.resumed; + +/// Runs [onDwell] once the list's item positions have stayed still for +/// [readingDwell]. Any position change restarts the wait. Changing [keys] +/// restarts it too, for changes that move no rows, such as new content or +/// the route coming back to the front. The wait is cancelled while +/// [active] is false. [onDwell] reads the current state when it runs. +void useReadingDwell({ + required ValueListenable positions, + required bool active, + required VoidCallback onDwell, + List keys = const [], +}) { + final callback = useRef(onDwell); + callback.value = onDwell; + useEffect(() { + if (!active) return null; + Timer? timer; + void restart() { + timer?.cancel(); + timer = Timer(readingDwell, () { + timer = null; + callback.value(); + }); + } + + positions.addListener(restart); + restart(); + return () { + timer?.cancel(); + positions.removeListener(restart); + }; + }, [positions, active, ...keys]); +} diff --git a/mobile/lib/features/channels/thread_detail_page.dart b/mobile/lib/features/channels/thread_detail_page.dart index eaaa8cdac6c..da5cd15d09a 100644 --- a/mobile/lib/features/channels/thread_detail_page.dart +++ b/mobile/lib/features/channels/thread_detail_page.dart @@ -50,11 +50,11 @@ import 'message_action_backdrop_state.dart'; import 'message_long_press_region.dart'; import 'message_content.dart'; import 'reaction_row.dart'; -import '../../shared/read_state/read_state_format.dart'; import '../../shared/read_state/read_state_provider.dart'; import 'send_message_provider.dart'; import 'small_avatar.dart'; import 'sticky_date_header.dart'; +import 'reading_marks.dart'; import 'timeline_message.dart'; part 'thread_detail_page/nested_thread_summary_row.dart'; @@ -65,7 +65,6 @@ part 'thread_detail_helpers.dart'; part 'thread_detail_page/tail_alignment.dart'; part 'thread_detail_page/thread_message.dart'; part 'thread_detail_page/avatar.dart'; -part 'thread_detail_page/read_state.dart'; part 'thread_detail_page/app_bar.dart'; const _landingHighlightDuration = Duration(seconds: 3); @@ -705,7 +704,79 @@ class ThreadDetailPage extends HookConsumerWidget { ); return null; }, [hasFetchedReplies, replies.length, settleGeometry, headHeight.value]); - _useThreadReplyReadState(ref, threadHead.id, replies); + void readVisibleReplies() { + final readState = ref.read(readStateProvider); + if (!context.mounted || + !readState.isReady || + !viewportHeight.isFinite || + viewportHeight <= 0) { + return; + } + // Rows hidden under the app bar, the composer or the keyboard are not + // read. The keyboard covers rows even when the list does not follow + // the tail, so this uses the whole covered height. + final topEdge = topOverlayFraction - 0.01; + final bottomEdge = + 1 - ((Grid.xs + navigationBottomInset) / viewportHeight) + 0.01; + final visible = []; + for (final position in itemPositionsListener.itemPositions.value) { + if (position.itemLeadingEdge < topEdge || + position.itemTrailingEdge > bottomEdge) { + continue; + } + if (position.index == headIndex) { + // The head is read like a reply. A thread opened directly, from + // a link or a notification, may be the only place it is seen. + if (!liveDeletionHidesHead) visible.add(liveHead); + continue; + } + final replyIndex = position.index - indexForReply(0); + if (replyIndex < 0 || replyIndex >= replies.length) continue; + visible.add(replies[replyIndex]); + } + final isDm = + ref + .read(channelsProvider) + .value + ?.where((candidate) => candidate.id == channelId) + .firstOrNull + ?.isDm ?? + false; + final marks = readingMarks( + readState: readState, + channelId: channelId, + isDm: isDm, + currentPubkey: currentPubkey, + loaded: allMsgs, + visible: visible, + bottom: replies.isNotEmpty && threadTailIsVisible() + ? replies.last + : null, + threadRootId: queryRootId, + isRootThread: threadHead.parentId == null, + ); + final notifier = ref.read(readStateProvider.notifier); + for (final mark in marks.entries) { + notifier.markContextRead(mark.key, mark.value); + } + } + + // Reading starts once the thread holds still for the dwell while this + // page is in front and the app is in use. + final appInUse = isAppInUse(useAppLifecycleState()); + final readStateReady = ref.watch( + readStateProvider.select((state) => state.isReady), + ); + useReadingDwell( + positions: itemPositionsListener.itemPositions, + active: + threadViewportVisible && + readStateReady && + appInUse && + (ModalRoute.of(context)?.isCurrent ?? true), + onDwell: readVisibleReplies, + keys: [threadHead.id, readingContentKey(allMsgs), navigationBottomInset], + ); // Thread-scoped typing indicators (exclude self). final allTyping = ref.watch(channelTypingProvider(channelId)); diff --git a/mobile/lib/features/channels/thread_detail_page/read_state.dart b/mobile/lib/features/channels/thread_detail_page/read_state.dart deleted file mode 100644 index 023593769bf..00000000000 --- a/mobile/lib/features/channels/thread_detail_page/read_state.dart +++ /dev/null @@ -1,24 +0,0 @@ -part of '../thread_detail_page.dart'; - -void _useThreadReplyReadState( - WidgetRef ref, - String threadHeadId, - List replies, -) { - final readState = ref.watch(readStateProvider); - final visibleReplyReadKey = replies - .map((reply) => '${reply.id}:${reply.createdAt}') - .join(','); - - useEffect(() { - if (!readState.isReady || replies.isEmpty) return null; - WidgetsBinding.instance.addPostFrameCallback((_) { - for (final reply in replies) { - ref - .read(readStateProvider.notifier) - .markContextRead(msgContextKey(reply.id), reply.createdAt); - } - }); - return null; - }, [threadHeadId, readState.isReady, visibleReplyReadKey]); -} diff --git a/mobile/lib/shared/read_state/read_state_format.dart b/mobile/lib/shared/read_state/read_state_format.dart index 35434bdae7c..449d2cadb1f 100644 --- a/mobile/lib/shared/read_state/read_state_format.dart +++ b/mobile/lib/shared/read_state/read_state_format.dart @@ -30,6 +30,243 @@ bool readByChannelCatchUp({ required bool highPriority, }) => !isDm && !isReply && !highPriority; +/// The plaintext budget for one published read-state slot. The web app +/// (`buzz-app`) uses the same limit, well under NIP-44's 64 KiB maximum. +const readStatePlaintextBytes = 40 * 1024; + +/// The most `msg:` and `thread:` marks this device saves. The desktop app +/// uses the same limit. +const localMaxPrunableContexts = 1000; + +bool _isPrunableContext(String contextId) => + contextId.startsWith(msgContextPrefix) || + contextId.startsWith(threadContextPrefix); + +/// The marks from [contexts] that this device saves, following the desktop +/// app's `pruneStaleContexts`. It drops `msg:` and `thread:` marks older +/// than the fetch horizon, then keeps the newest [localMaxPrunableContexts] +/// of them. Each automatic read writes a `msg:` mark, so without this bound +/// the saved state, and the time to save it, would grow with every read. +/// Channel and catch-up marks are kept: there is one per channel or thread, +/// and losing one would show read messages as unread again. +Map pruneStaleContexts( + Map contexts, { + required int nowUnixSeconds, +}) { + final cutoff = nowUnixSeconds - readStateHorizonSeconds; + final kept = {}; + final prunable = >[]; + for (final entry in contexts.entries) { + if (!_isPrunableContext(entry.key)) { + kept[entry.key] = entry.value; + } else if (entry.value >= cutoff) { + prunable.add(entry); + } + } + prunable.sort((a, b) { + final byTime = b.value.compareTo(a.value); + return byTime != 0 ? byTime : a.key.compareTo(b.key); + }); + for (final entry in prunable.take(localMaxPrunableContexts)) { + kept[entry.key] = entry.value; + } + return kept; +} + +/// Whether [contextId] belongs to the web app's manual-unread overrides +/// (`ov_s:`, `ov_c:`, `ov_b:`) or their escaped frontiers (`esc:`). This app +/// does not read overrides. It does not take them from other slots, but it +/// carries the ones already in its own slot unchanged: NIP-RS forbids +/// dropping them, and an old version of this app copied them there. +bool isOverrideContext(String contextId) => + contextId.startsWith('ov_') || contextId.startsWith('esc:'); + +const _overridePrefixes = ['ov_s:', 'ov_c:', 'ov_b:']; + +/// The raw context ID of an override counter key, or null for other keys. +String? _overrideTarget(String key) { + for (final prefix in _overridePrefixes) { + if (key.startsWith(prefix)) return key.substring(prefix.length); + } + return null; +} + +/// The frontier key of raw context [contextId] on the wire. NIP-RS escapes +/// a raw ID that begins with `ov_` or `esc:` by prepending `esc:`. +String overrideFrontierKey(String contextId) => + contextId.startsWith('ov_') || contextId.startsWith('esc:') + ? 'esc:$contextId' + : contextId; + +/// Whether [keys] form a legal override group for [target]: exactly +/// `ov_s:`, `ov_c:` and `ov_b:`, or `ov_c:` alone as a tombstone. +bool _isOverrideShape(String target, Iterable keys) { + final set = keys.toSet(); + return set.length == _overridePrefixes.length || + (set.length == 1 && set.contains('ov_c:$target')); +} + +/// The override groups in [contexts] that have a legal shape, by target. +Map> _overrideGroups(Map contexts) { + final groups = >{}; + for (final entry in contexts.entries) { + if (_overrideTarget(entry.key) case final target?) { + (groups[target] ??= {})[entry.key] = entry.value; + } + } + return { + for (final MapEntry(key: target, value: group) in groups.entries) + if (_isOverrideShape(target, group.keys)) target: group, + }; +} + +/// The override keys of [contexts] that this app may carry. NIP-RS accepts +/// an override group only whole (see [_isOverrideShape]), and a group must +/// travel with its frontier, found in [contexts] or [frontiers]. Any other +/// group is rejected whole, so none of its keys is carried. Escaped +/// frontiers (`esc:`) are frontier keys, not counters, so each is carried +/// as it is. +Map completeOverrideGroups( + Map contexts, { + Map frontiers = const {}, +}) { + final kept = { + for (final entry in contexts.entries) + if (entry.key.startsWith('esc:')) entry.key: entry.value, + }; + for (final MapEntry(key: target, value: group) in _overrideGroups( + contexts, + ).entries) { + final frontier = overrideFrontierKey(target); + if (contexts.containsKey(frontier) || frontiers.containsKey(frontier)) { + kept.addAll(group); + } + } + return kept; +} + +/// The frontier keys that must travel with the override groups in +/// [contexts] (NIP-RS's co-location rule). +Set overrideGroupFrontierKeys(Map contexts) => { + for (final target in _overrideGroups(contexts).keys) + overrideFrontierKey(target), +}; + +/// Whether this device republishes a mark it merged from another device's +/// slot. Message marks are not republished: each covers one message, and the +/// web app prunes them once a catch-up mark reads the message. Broad marks +/// are republished so they outlive the fetch horizon of the slot that wrote +/// them. +bool republishesMergedContext(String contextId) => + !contextId.startsWith(msgContextPrefix) && !isOverrideContext(contextId); + +/// Keep order for a slot that is over budget, following `buzz-app`'s +/// retention: channel marks first, then thread marks, then catch-up marks, +/// then message marks. +int _retentionScope(String key) { + if (!key.contains(':')) return 0; + if (key.startsWith(threadContextPrefix)) return 1; + if (key.startsWith('activity:') || key.startsWith('thread-activity:')) { + return 2; + } + return 3; +} + +/// The share of the budget that broad marks may fill before recent use +/// decides. The rest always goes to the most recently written marks, so a +/// new read is never left out because old channel or thread marks fill the +/// budget. `buzz-app` uses the same share. +const _scopedShare = 0.75; + +/// The marks from [contexts] that fit one read-state slot for [clientId] +/// within [maxBytes] of plaintext, following `buzz-app`'s retention. +/// +/// Up to three quarters of the budget goes to broad marks first: channel +/// marks, then thread marks, then catch-up marks, then message marks. The +/// rest goes by [recent], the time each mark was last written, newest first, +/// so reading old history still syncs. A slot over NIP-44's limit cannot be +/// encrypted, so without this cap sync would stop. +/// +/// [carried] override groups are kept whole, ahead of every mark, together +/// with each group's frontier from [contexts]: NIP-RS requires a context's +/// frontier and its override keys in the same event. Incomplete groups are +/// left out (see [completeOverrideGroups]). If the groups do not fit, this +/// returns null, and the caller must leave the slot as it is instead of +/// publishing part of a group. +Map? retainReadStateContexts( + Map contexts, { + required String clientId, + Map recent = const {}, + Map carried = const {}, + int maxBytes = readStatePlaintextBytes, +}) { + // JSON-encoded bytes, so escaped characters are counted as published. + int bytesOf(Object? value) => utf8.encode(jsonEncode(value)).length; + final groups = completeOverrideGroups(carried, frontiers: contexts); + final reserved = {...groups}; + for (final key in overrideGroupFrontierKeys(groups)) { + final frontier = contexts[key]; + if (frontier != null && frontier > (reserved[key] ?? -1)) { + reserved[key] = frontier; + } + } + final fixedBytes = bytesOf( + ReadStateBlob(clientId: clientId, contexts: reserved).toJson(), + ); + if (fixedBytes > maxBytes || reserved.length > _maxContexts) return null; + // Marks share what the carried keys leave. + final markBytes = maxBytes - fixedBytes; + final markCount = _maxContexts - reserved.length; + var used = 0; + var count = 0; + bool take(MapEntry entry, double share) { + // A key, its colon, its value and one separating comma. + final cost = bytesOf(entry.key) + 2 + '${entry.value}'.length; + if (used + cost > markBytes * share || count >= markCount * share) { + return false; + } + used += cost; + count++; + return true; + } + + int byUse(MapEntry a, MapEntry b) { + final byRecent = (recent[b.key] ?? 0).compareTo(recent[a.key] ?? 0); + if (byRecent != 0) return byRecent; + final byTime = b.value.compareTo(a.value); + return byTime != 0 ? byTime : a.key.compareTo(b.key); + } + + final retained = {...reserved}; + final scoped = + contexts.entries + .where((entry) => !reserved.containsKey(entry.key)) + .toList() + ..sort((a, b) { + final byScope = _retentionScope( + a.key, + ).compareTo(_retentionScope(b.key)); + return byScope != 0 ? byScope : byUse(a, b); + }); + for (final entry in scoped) { + // Stop at the first broad mark that does not fit, so a narrower mark + // never takes the share ahead of it. + if (!take(entry, _scopedShare)) break; + retained[entry.key] = entry.value; + } + final byRecent = + contexts.entries + .where((entry) => !reserved.containsKey(entry.key)) + .toList() + ..sort(byUse); + for (final entry in byRecent) { + if (!retained.containsKey(entry.key) && take(entry, 1)) { + retained[entry.key] = entry.value; + } + } + return retained; +} + int? maxReadAt(Iterable markers) { int? latest; for (final marker in markers) { @@ -142,16 +379,37 @@ ReadStateBlob? decodeReadStateBlob(String plaintext) { ); } +/// The valid entries of a decoded `contexts` object. Override counters are +/// checked as whole groups first, on the raw values, as NIP-RS requires: one +/// invalid sibling rejects its whole group, so dropping it alone can never +/// turn a live group into a tombstone. Other `ov_` keys are reserved and +/// dropped. Every other entry is checked on its own. Map sanitizeReadStateContexts(Map contexts) { + bool valid(String key, Object? value) => + utf8.encode(key).length <= 256 && + value is int && + value >= 0 && + value <= 4294967295; + final groups = >{}; + for (final entry in contexts.entries) { + if (_overrideTarget(entry.key) case final target?) { + (groups[target] ??= {})[entry.key] = entry.value; + } + } + final rejected = { + for (final MapEntry(key: target, value: group) in groups.entries) + if (!_isOverrideShape(target, group.keys) || + group.entries.any((entry) => !valid(entry.key, entry.value))) + target, + }; final sanitized = {}; for (final entry in contexts.entries) { - if (utf8.encode(entry.key).length > 256) continue; - + if (entry.key.startsWith('ov_')) { + final target = _overrideTarget(entry.key); + if (target == null || rejected.contains(target)) continue; + } final value = entry.value; - if (value is! int) continue; - if (value < 0 || value > 4294967295) continue; - - sanitized[entry.key] = value; + if (valid(entry.key, value)) sanitized[entry.key] = value as int; } return sanitized; } diff --git a/mobile/lib/shared/read_state/read_state_manager.dart b/mobile/lib/shared/read_state/read_state_manager.dart index f8dcb878db6..7f730caa9f4 100644 --- a/mobile/lib/shared/read_state/read_state_manager.dart +++ b/mobile/lib/shared/read_state/read_state_manager.dart @@ -55,6 +55,9 @@ class ReadStateManager { final Map _effectiveState = {}; final Set _publishableContextIds = {}; + // Override keys already in this device's own slot, carried unchanged. + // See [isOverrideContext]. + final Map _carriedOverrides = {}; Map _lastPublishedContexts = {}; Timer? _debounceTimer; @@ -62,6 +65,9 @@ class ReadStateManager { bool _initialized = false; bool _disposed = false; bool _isPublishing = false; + // Set when a publish is asked for while one is running. The running + // publish took its snapshot first, so it publishes again when it ends. + bool _publishAgain = false; Completer? _publishCompleter; bool _remoteUnsupported = false; int _maxFetchedCreatedAt = 0; @@ -116,11 +122,14 @@ class ReadStateManager { } void markContextRead(String contextId, int unixTimestamp) { - _advanceContext(contextId, unixTimestamp, publishable: true); + if (_disposed || isOverrideContext(contextId)) return; + // Set first: retention keeps the most recently written marks, so the + // save that follows must already see this read as the newest. _contextSourceCreatedAt[contextId] = max( currentUnixSeconds(), _maxFetchedCreatedAt + 1, ); + _advanceContext(contextId, unixTimestamp, publishable: true); } void seedContextRead(String contextId, int unixTimestamp) { @@ -243,6 +252,7 @@ class ReadStateManager { } for (final entry in decoded.blob.contexts.entries) { + if (isOverrideContext(entry.key)) continue; final result = _applyRemoteContextTimestamp( contextId: entry.key, timestamp: entry.value, @@ -250,7 +260,9 @@ class ReadStateManager { ); if (result == _ApplyRemoteContextResult.advanced) { _pendingSyncedAdvances.add(entry.key); - _publishableContextIds.add(entry.key); + if (republishesMergedContext(entry.key)) { + _publishableContextIds.add(entry.key); + } } } @@ -262,9 +274,26 @@ class ReadStateManager { } if (ownBlob != null) { + _carryOwnOverrides(ownBlob.contexts); _lastPublishedContexts = Map.from(ownBlob.contexts); - _publishableContextIds.addAll(ownBlob.contexts.keys); + _publishableContextIds.addAll( + ownBlob.contexts.keys.where((key) => !isOverrideContext(key)), + ); + } + } + + /// Merges the complete override groups of this device's own slot into + /// [_carriedOverrides] by `max()`, the NIP-RS merge rule. An incomplete + /// group is rejected whole (see [completeOverrideGroups]). + bool _carryOwnOverrides(Map contexts) { + var changed = false; + for (final entry in completeOverrideGroups(contexts).entries) { + if (entry.value > (_carriedOverrides[entry.key] ?? -1)) { + _carriedOverrides[entry.key] = entry.value; + changed = true; + } } + return changed; } Future _startLiveSubscription() async { @@ -313,8 +342,11 @@ class ReadStateManager { _rotateSlotId(); } - var changed = false; + var changed = + decoded.blob.clientId == _clientId && + _carryOwnOverrides(decoded.blob.contexts); for (final entry in decoded.blob.contexts.entries) { + if (isOverrideContext(entry.key)) continue; final result = _applyRemoteContextTimestamp( contextId: entry.key, timestamp: entry.value, @@ -324,7 +356,9 @@ class ReadStateManager { _pendingSyncedAdvances.add(entry.key); changed = true; } - if (_publishableContextIds.add(entry.key)) { + if ((decoded.blob.clientId == _clientId || + republishesMergedContext(entry.key)) && + _publishableContextIds.add(entry.key)) { changed = true; } } @@ -385,16 +419,45 @@ class ReadStateManager { _signedEventRelay == null) { return; } - if (_isPublishing) return; + if (_isPublishing) { + // The running publish may have taken its snapshot before this change. + _publishAgain = true; + return _publishCompleter?.future; + } final completer = Completer(); _publishCompleter = completer; _isPublishing = true; + try { + do { + _publishAgain = false; + await _publishOnce(_signedEventRelay); + } while (_publishAgain && + (allowDisposed || !_disposed) && + !_remoteUnsupported); + } finally { + _isPublishing = false; + _publishAgain = false; + completer.complete(); + if (_publishCompleter == completer) { + _publishCompleter = null; + } + } + } + + Future _publishOnce(SignedEventRelay signedEventRelay) async { debugPrint('[ReadStateManager] publish starting slotId=$_slotId'); try { await _fetchOwnBlobBeforePublish(); final contexts = _currentContexts(); + if (contexts == null) { + debugPrint( + '[ReadStateManager] publish skipped: carried override keys do not ' + 'fit the slot, so it stays as it is.', + ); + return; + } if (_isIdenticalToLastPublished(contexts)) { return; } @@ -403,7 +466,7 @@ class ReadStateManager { final ciphertext = _crypto.encrypt(jsonEncode(blob.toJson())); final createdAt = max(currentUnixSeconds(), _maxFetchedCreatedAt + 1); - await _signedEventRelay.submit( + await signedEventRelay.submit( kind: EventKind.readState, content: ciphertext, tags: [ @@ -444,12 +507,6 @@ class ReadStateManager { return; } debugPrint('[ReadStateManager] publish failed: $error'); - } finally { - _isPublishing = false; - completer.complete(); - if (_publishCompleter == completer) { - _publishCompleter = null; - } } } @@ -477,7 +534,9 @@ class ReadStateManager { } } - bool _isIdenticalToLastPublished(Map contexts) { + bool _isIdenticalToLastPublished(Map? contexts) { + // Null: the slot cannot be published, so there is nothing to send. + if (contexts == null) return true; if (_lastPublishedContexts.length != contexts.length) { return false; } @@ -495,24 +554,50 @@ class ReadStateManager { return drained; } - Map _currentContexts() { + /// The slot to publish, or null when the carried override keys alone do + /// not fit, so the current slot must stay as it is. + Map? _currentContexts() { final contexts = {}; for (final entry in _effectiveState.entries) { if (_publishableContextIds.contains(entry.key)) { contexts[entry.key] = entry.value; } } - return contexts; + // A carried group's frontier travels with it, even a merged `msg:` mark + // that would not be republished on its own. + for (final key in overrideGroupFrontierKeys(_carriedOverrides)) { + if (_effectiveState[key] case final frontier?) contexts[key] = frontier; + } + return retainReadStateContexts( + contexts, + clientId: _clientId, + recent: _contextSourceCreatedAt, + carried: _carriedOverrides, + ); } void _hydrateFromLocalStorage() { final stored = _storage.read(pubkey); + // Earlier versions merged the web app's override keys and published + // them in this device's slot, so publishable ones are carried. + _carriedOverrides + ..clear() + ..addAll( + completeOverrideGroups({ + for (final entry in stored.contexts.entries) + if (isOverrideContext(entry.key) && + stored.publishableContextIds.contains(entry.key)) + entry.key: entry.value, + }, frontiers: stored.contexts), + ); _effectiveState ..clear() - ..addAll(stored.contexts); + ..addEntries( + stored.contexts.entries.where((entry) => !isOverrideContext(entry.key)), + ); _publishableContextIds ..clear() - ..addAll(stored.publishableContextIds); + ..addAll(stored.publishableContextIds.where(_effectiveState.containsKey)); _contextSourceCreatedAt ..clear() ..addAll(stored.sourceCreatedAt); @@ -520,11 +605,27 @@ class ReadStateManager { } void _persistLocalState() { + // Save a bounded copy. Memory keeps every mark for this session, so a + // message read from old history stays read until the app restarts. + final saved = pruneStaleContexts( + _effectiveState, + nowUnixSeconds: currentUnixSeconds(), + )..addAll(_carriedOverrides); + // A carried group's frontier is never pruned away from it. + for (final key in overrideGroupFrontierKeys(_carriedOverrides)) { + if (_effectiveState[key] case final frontier?) saved[key] = frontier; + } _storage.write( pubkey, - _effectiveState, - _publishableContextIds, - _contextSourceCreatedAt, + saved, + { + ..._publishableContextIds.where(saved.containsKey), + ..._carriedOverrides.keys, + }, + { + for (final entry in _contextSourceCreatedAt.entries) + if (saved.containsKey(entry.key)) entry.key: entry.value, + }, ); } diff --git a/mobile/lib/shared/read_state/read_state_provider.dart b/mobile/lib/shared/read_state/read_state_provider.dart index 22375db959e..b7e2ee645e1 100644 --- a/mobile/lib/shared/read_state/read_state_provider.dart +++ b/mobile/lib/shared/read_state/read_state_provider.dart @@ -180,6 +180,16 @@ class ReadStateNotifier extends Notifier { _refreshForcedState(); } + /// Clear a forced-unread flag without moving any read marker. Reading to + /// the bottom of a channel the reader marked unread uses this: it ends the + /// manual unread but must not read mentions or replies the reader has not + /// seen. + void clearForcedUnread(String contextId) { + if (_forcedUnreadContexts.remove(contextId) != null) { + _refreshForcedState(); + } + } + void _refreshForcedState() { final manager = _manager; if (manager == null) return; diff --git a/mobile/test/features/activity/inbox_item_test.dart b/mobile/test/features/activity/inbox_item_test.dart index 379e85dfc52..ad0f75d1707 100644 --- a/mobile/test/features/activity/inbox_item_test.dart +++ b/mobile/test/features/activity/inbox_item_test.dart @@ -157,20 +157,23 @@ void main() { item(id: 'c', createdAt: 30, tags: replyTags('root1', 'b')), ]); + bool readUpTo(int time, FeedItem event) => event.createdAt <= time; + test('targets the oldest unread event', () { - expect(rows.single.deepLinkTarget(15).id, 'b'); + expect(rows.single.deepLinkTarget((e) => readUpTo(15, e)).id, 'b'); }); test('targets the oldest event when nothing is read', () { - expect(rows.single.deepLinkTarget(5).id, 'a'); + expect(rows.single.deepLinkTarget((_) => false).id, 'a'); }); test('falls back to the latest event when all are read', () { - expect(rows.single.deepLinkTarget(99).id, 'c'); + expect(rows.single.deepLinkTarget((_) => true).id, 'c'); }); - test('falls back to the latest event without a read marker', () { - expect(rows.single.deepLinkTarget(null).id, 'c'); + test('checks each event, not a single time', () { + // Only the middle reply is unread, though the newest is read. + expect(rows.single.deepLinkTarget((e) => e.id != 'b').id, 'b'); }); }); diff --git a/mobile/test/features/activity/inbox_read_state_test.dart b/mobile/test/features/activity/inbox_read_state_test.dart index af8f3b90912..9885dd992e1 100644 --- a/mobile/test/features/activity/inbox_read_state_test.dart +++ b/mobile/test/features/activity/inbox_read_state_test.dart @@ -30,28 +30,38 @@ int? Function(String) markers(Map map) => (contextId) => map[contextId]; void main() { - group('resolveInboxItemReadAt', () { - test('channel rows use the channel marker', () { - final row = buildInboxItems([item(id: 'a')]).single; - expect(resolveInboxItemReadAt(row, markerOf: markers({'ch1': 42})), 42); + group('inboxEventReadAt', () { + test('a top-level event uses the channel and its own mark', () { + final event = item(id: 'a', createdAt: 60); + expect(inboxEventReadAt(event, markerOf: markers({'ch1': 42})), 42); + expect( + inboxEventReadAt(event, markerOf: markers({'ch1': 42, 'msg:a': 60})), + 60, + ); }); - test('thread rows use max of thread and msg markers', () { - final row = buildInboxItems([ - item(id: 'a', tags: replyTags('root1', 'root1')), - ]).single; + test('a reply also uses its thread marks', () { + final event = item(id: 'a', tags: replyTags('root1', 'root1')); expect( - resolveInboxItemReadAt( - row, - markerOf: markers({'thread:root1': 10, 'msg:a': 25}), + inboxEventReadAt( + event, + markerOf: markers({'thread:root1': 10, 'thread-activity:root1': 30}), ), - 25, + 30, ); }); - test('channel-less rows have no marker', () { - final row = buildInboxItems([item(id: 'a', channelId: null)]).single; - expect(resolveInboxItemReadAt(row, markerOf: markers({})), isNull); + test('channel catch-up never reads an Activity event', () { + final event = item(id: 'a', createdAt: 60, category: 'activity'); + expect( + inboxEventReadAt(event, markerOf: markers({'activity:ch1': 99})), + isNull, + ); + }); + + test('an event without a channel has no marker', () { + final event = item(id: 'a', channelId: null); + expect(inboxEventReadAt(event, markerOf: markers({})), isNull); }); }); @@ -78,6 +88,37 @@ void main() { ); }); + test('a channel mention read by its own mark is done', () { + // Reading a channel writes `msg:` for a mention, not the channel mark. + final row = buildInboxItems([item(id: 'm', createdAt: 50)]).single; + expect( + isInboxItemDone( + row, + markerOf: markers({'ch1': 10, 'msg:m': 50}), + localUnreadOverrides: const {}, + localDoneSet: const {}, + ), + isTrue, + ); + }); + + test('a thread mention needs every grouped reply read', () { + final row = buildInboxItems([ + item(id: 'r1', createdAt: 40, tags: replyTags('root1', 'root1')), + item(id: 'r2', createdAt: 50, tags: replyTags('root1', 'r1')), + ]).single; + bool done(Map marks) => isInboxItemDone( + row, + markerOf: markers(marks), + localUnreadOverrides: const {}, + localDoneSet: const {}, + ); + // Seeing only the newest reply leaves the older mention unread. + expect(done({'msg:r2': 50}), isFalse); + expect(done({'msg:r1': 40, 'msg:r2': 50}), isTrue); + expect(done({'thread-activity:root1': 50}), isTrue); + }); + test('a local unread override always wins', () { final row = buildInboxItems([item(id: 'a', createdAt: 50)]).single; expect( diff --git a/mobile/test/features/channels/channel_detail_page_test.dart b/mobile/test/features/channels/channel_detail_page_test.dart index ab6bc9e50da..9eb95eb428c 100644 --- a/mobile/test/features/channels/channel_detail_page_test.dart +++ b/mobile/test/features/channels/channel_detail_page_test.dart @@ -2093,7 +2093,9 @@ void main() { expect(find.byType(SkeletonReveal), findsNothing); }); - testWidgets('defers read-state mark until after build', (tester) async { + testWidgets('reading the bottom writes the catch-up mark after the dwell', ( + tester, + ) async { final readState = _SynchronousReadStateNotifier( const ReadStateState( isReady: true, @@ -2125,11 +2127,620 @@ void main() { expect(tester.takeException(), isNull); await tester.pump(); + // Opening the channel does not read it; holding still at the bottom + // does, without moving the channel mark. + expect(readState.markedContexts, isEmpty); + + await tester.pump(const Duration(milliseconds: 300)); - expect(readState.markedContexts, {_channelId: 1200}); + expect(readState.markedContexts, {'activity:$_channelId': 1200}); expect(tester.takeException(), isNull); }); + testWidgets('reading the bottom ends a manual channel unread only', ( + tester, + ) async { + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + forcedUnreadContexts: {_channelId: _channelId, 'msg:old': _channelId}, + ), + ); + + await tester.pumpWidget( + _buildTestable( + messages: [ + _textMsg( + id: 'msg1', + pubkey: 'alice', + content: 'Latest', + createdAt: 1200, + ), + ], + readStateNotifier: readState, + ), + ); + await tester.pump(); + expect(readState.clearedForcedContexts, isEmpty); + + await tester.pump(const Duration(milliseconds: 300)); + + expect(readState.clearedForcedContexts, [_channelId]); + // The channel mark does not move, and a message the reader marked + // unread stays unread. + expect(readState.markedContexts, {'activity:$_channelId': 1200}); + expect(readState.state.isForcedUnread('msg:old'), isTrue); + }); + + testWidgets('opening a DM reads the whole DM', (tester) async { + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final dmChannel = Channel( + id: _channelId, + name: 'DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Alice'], + participantPubkeys: const ['self', 'alice'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: [ + _textMsg( + id: 'msg1', + pubkey: 'alice', + content: 'First', + createdAt: 1100, + ), + _textMsg( + id: 'reply1', + pubkey: 'alice', + content: 'Newest reply', + createdAt: 1300, + extraTags: const [ + ['e', 'msg1', '', 'reply'], + ], + ), + ], + channel: dmChannel, + relaySessionNotifier: PresenceTestRelay()..emptySnapshots = true, + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + }, + readStateNotifier: readState, + ), + ); + await tester.pump(); + + // No dwell needed: opening the DM reads through its newest message, + // replies included. + expect(readState.markedContexts[_channelId], 1300); + }); + + for (final covered in ['a sheet', 'the app switcher']) { + testWidgets('a DM behind $covered reads a new reply only ' + 'when shown again', (tester) async { + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final first = _textMsg( + id: 'msg1', + pubkey: 'alice', + content: 'First', + createdAt: 1100, + ); + final messages = _FakeMessagesNotifier([first]); + await tester.pumpWidget( + _buildTestable( + messages: const [], + messagesNotifier: messages, + channel: Channel( + id: _channelId, + name: 'DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Alice'], + participantPubkeys: const ['self', 'alice'], + isMember: true, + ), + relaySessionNotifier: PresenceTestRelay()..emptySnapshots = true, + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + }, + readStateNotifier: readState, + ), + ); + await tester.pump(); + expect(readState.markedContexts[_channelId], 1100); + + if (covered == 'a sheet') { + // A sheet keeps the page on screen, so its providers stay live. + unawaited( + showModalBottomSheet( + context: tester.element(find.byType(ChannelDetailPage)), + builder: (_) => const Text('Cover'), + ), + ); + await tester.pumpAndSettle(); + } else { + // Inactive, unlike paused, still draws frames. + tester.binding.handleAppLifecycleStateChanged( + AppLifecycleState.inactive, + ); + await tester.pump(); + } + messages.setMessages([ + first, + _textMsg( + id: 'reply1', + pubkey: 'alice', + content: 'Unseen reply', + createdAt: 1300, + extraTags: const [ + ['e', 'msg1', '', 'reply'], + ], + ), + ]); + await tester.pump(); + await tester.pump(); + + // The reply loaded while the DM was covered, so it stays unread. + expect(readState.markedContexts[_channelId], 1100); + + if (covered == 'a sheet') { + Navigator.of(tester.element(find.text('Cover'))).pop(); + await tester.pumpAndSettle(); + } else { + tester.binding.handleAppLifecycleStateChanged( + AppLifecycleState.resumed, + ); + await tester.pump(); + await tester.pump(); + } + + // Shown again, the DM reads through the reply. + expect(readState.markedContexts[_channelId], 1300); + }); + } + + testWidgets('reading away from the bottom marks only visible rows', ( + tester, + ) async { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final messages = [ + for (var i = 0; i < 40; i++) + _textMsg( + id: 'm$i', + pubkey: 'alice', + content: 'Message $i', + createdAt: 1000 + i, + ), + ]; + + await tester.pumpWidget( + _buildTestable( + messages: messages, + initialMessageId: 'm5', + readStateNotifier: readState, + ), + ); + await tester.pumpAndSettle(); + await tester.pump(const Duration(milliseconds: 300)); + + final marked = readState.markedContexts; + expect(marked['msg:m5'], 1005); + expect(marked.keys.where((key) => !key.startsWith('msg:')), isEmpty); + // Rows below the viewport are not read. + expect(marked, isNot(contains('msg:m30'))); + expect(marked, isNot(contains('msg:m39'))); + }); + + testWidgets('reading a thread tail writes the thread catch-up mark', ( + tester, + ) async { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final root = _textMsg( + id: 'root', + pubkey: 'alice', + content: 'Thread root', + createdAt: 1000, + ); + final replies = [ + for (var i = 0; i < 3; i++) + _textMsg( + id: 'reply$i', + pubkey: 'bob', + content: 'Reply $i', + createdAt: 1100 + i, + extraTags: const [ + ['e', 'root', '', 'reply'], + ], + ), + ]; + + await tester.pumpWidget( + _buildTestable( + messages: [root], + threadReplies: {'root': replies}, + readStateNotifier: readState, + ), + ); + await tester.pumpAndSettle(); + await tester.pump(const Duration(milliseconds: 300)); + expect(readState.markedContexts, {'activity:$_channelId': 1000}); + readState.markedContexts.clear(); + + Navigator.of(tester.element(find.byType(ChannelDetailPage))).push( + MaterialPageRoute( + builder: (_) => ThreadDetailPage( + threadHead: formatTimeline([root]).single, + allMessages: formatTimeline([root, ...replies]), + channelId: _channelId, + currentPubkey: 'self', + isMember: true, + isArchived: false, + ), + ), + ); + await tester.pumpAndSettle(); + await tester.pump(const Duration(milliseconds: 300)); + + // A catch-up mark reads the thread. Each visible reply keeps its own + // mark too, as in `buzz-app`: the thread mark reads a reply only while + // its root is loaded. The covered channel page does not read meanwhile. + expect(readState.markedContexts, { + 'thread-activity:root': 1102, + 'msg:reply0': 1100, + 'msg:reply1': 1101, + 'msg:reply2': 1102, + }); + }); + + testWidgets('a thread opened directly reads its head', (tester) async { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + // A mention: no catch-up mark reads it, so only its own mark does. + final root = _textMsg( + id: 'root', + pubkey: 'alice', + content: 'Thread root for @self', + createdAt: 1000, + extraTags: const [ + ['p', 'self'], + ], + ); + final reply = _textMsg( + id: 'reply0', + pubkey: 'bob', + content: 'Reply 0', + createdAt: 1100, + extraTags: const [ + ['e', 'root', '', 'reply'], + ], + ); + + await tester.pumpWidget( + _buildTestable( + messages: [root, reply], + threadReplies: { + 'root': [reply], + }, + initialThreadRootId: 'root', + readStateNotifier: readState, + ), + ); + await tester.pumpAndSettle(); + await tester.pump(const Duration(milliseconds: 300)); + + expect(find.byType(ThreadDetailPage), findsOneWidget); + expect(readState.markedContexts['msg:root'], 1000); + expect(readState.markedContexts['thread-activity:root'], 1100); + }); + + testWidgets('a thread with no replies reads its head', (tester) async { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final root = _textMsg( + id: 'root', + pubkey: 'alice', + content: 'Thread root for @self', + createdAt: 1000, + extraTags: const [ + ['p', 'self'], + ], + ); + + await tester.pumpWidget( + _buildTestable( + messages: [root], + threadReplies: const {'root': []}, + readStateNotifier: readState, + home: ThreadDetailPage( + threadHead: formatTimeline([root]).single, + allMessages: formatTimeline([root]), + channelId: _channelId, + currentPubkey: 'self', + isMember: true, + isArchived: false, + ), + ), + ); + await tester.pumpAndSettle(); + await tester.pump(const Duration(milliseconds: 300)); + + expect(readState.markedContexts, {'msg:root': 1000}); + }); + + testWidgets('rows under the Android keyboard are not read in history', ( + tester, + ) async { + final previousPlatform = debugDefaultTargetPlatformOverride; + debugDefaultTargetPlatformOverride = TargetPlatform.android; + try { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + tester.view.viewPadding = const FakeViewPadding(bottom: 24); + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final messages = [ + for (var i = 0; i < 40; i++) + _textMsg( + id: 'm$i', + pubkey: 'alice', + content: 'Message $i', + createdAt: 1000 + i, + extraTags: const [ + ['p', 'self'], + ], + ), + ]; + + await tester.pumpWidget( + _buildTestable( + messages: messages, + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + }, + readStateNotifier: readState, + ), + ); + await tester.pumpAndSettle(); + + await tester.tap(find.text('Message #general')); + await tester.pump(); + tester.view.viewInsets = const FakeViewPadding(bottom: 300); + await tester.pump(); + await tester.pump(androidImeMetricsSettleDelay); + await tester.pumpAndSettle(); + + // Leave the latest message, then stop with history rows under the + // keyboard. + final messageList = find.byKey(const ValueKey('channel-message-list')); + final messageListElement = tester.element(messageList); + UserScrollNotification( + metrics: FixedScrollMetrics( + minScrollExtent: 0, + maxScrollExtent: 100, + pixels: 0, + viewportDimension: 100, + axisDirection: AxisDirection.down, + devicePixelRatio: 1, + ), + context: messageListElement, + direction: ScrollDirection.reverse, + ).dispatch(messageListElement); + tester + .widget(messageList) + .itemScrollController! + .jumpTo(index: 20); + await tester.pumpAndSettle(); + readState.markedContexts.clear(); + await tester.pump(const Duration(milliseconds: 300)); + + final textField = tester.widget(find.byType(TextField)); + expect(textField.focusNode?.hasFocus, isTrue); + final coveredTop = tester + .getTopLeft(find.byKey(const ValueKey('channel-composer-dock'))) + .dy; + final marked = readState.markedContexts.keys.toSet(); + var hiddenRows = 0; + for (var i = 0; i < 40; i++) { + final row = find.byKey(ValueKey('channel-message-group-m$i')); + if (row.evaluate().isEmpty) continue; + final hidden = tester.getTopLeft(row).dy >= coveredTop; + if (hidden) { + hiddenRows++; + expect(marked, isNot(contains('msg:m$i')), reason: 'm$i'); + } + if (marked.contains('msg:m$i')) { + // The reading edge allows 1% of the list height. + expect( + tester.getBottomLeft(row).dy, + lessThanOrEqualTo(coveredTop + 10), + reason: 'm$i', + ); + } + } + // Rows sit under the keyboard, and rows above it are read. + expect(hiddenRows, greaterThan(0)); + expect(marked, isNotEmpty); + } finally { + debugDefaultTargetPlatformOverride = previousPlatform; + } + }); + + testWidgets('thread rows under the Android keyboard are not read', ( + tester, + ) async { + final previousPlatform = debugDefaultTargetPlatformOverride; + debugDefaultTargetPlatformOverride = TargetPlatform.android; + try { + tester.view.physicalSize = const Size(400, 800); + tester.view.devicePixelRatio = 1; + tester.view.viewPadding = const FakeViewPadding(bottom: 24); + addTearDown(tester.view.reset); + final readState = _SynchronousReadStateNotifier( + const ReadStateState( + isReady: true, + pubkey: 'self', + contexts: {}, + version: 0, + ), + ); + final root = _textMsg( + id: 'root', + pubkey: 'alice', + content: 'Thread root', + createdAt: 1000, + ); + final replies = [ + for (var i = 0; i < 40; i++) + _textMsg( + id: 'reply$i', + pubkey: 'bob', + content: 'Reply $i', + createdAt: 1100 + i, + extraTags: const [ + ['e', 'root', '', 'reply'], + ], + ), + ]; + + await tester.pumpWidget( + _buildTestable( + messages: [root], + threadReplies: {'root': replies}, + readStateNotifier: readState, + home: ThreadDetailPage( + threadHead: formatTimeline([root]).single, + allMessages: formatTimeline([root, ...replies]), + channelId: _channelId, + currentPubkey: 'self', + isMember: true, + isArchived: false, + ), + ), + ); + await tester.pumpAndSettle(); + + // Scroll into history, open the keyboard, then scroll toward newer + // replies without reaching the tail. An upward drag keeps the + // keyboard open. + final messageList = find.byKey(const ValueKey('thread-message-list')); + await tester.drag(messageList, const Offset(0, 600)); + await tester.pumpAndSettle(); + await tester.tap(find.text('Reply in thread…').hitTestable()); + await tester.pumpAndSettle(); + tester.view.viewInsets = const FakeViewPadding(bottom: 300); + await tester.pump(); + await tester.pump(androidImeMetricsSettleDelay); + await tester.pumpAndSettle(); + await tester.drag(messageList, const Offset(0, -150)); + await tester.pumpAndSettle(); + readState.markedContexts.clear(); + await tester.pump(const Duration(milliseconds: 300)); + + final textField = tester.widget(find.byType(TextField)); + expect(textField.focusNode?.hasFocus, isTrue); + final coveredTop = tester + .getTopLeft(find.byKey(const ValueKey('thread-composer-dock'))) + .dy; + final marked = readState.markedContexts.keys.toSet(); + expect(marked, isNot(contains('thread-activity:root'))); + var hiddenRows = 0; + for (var i = 0; i < 40; i++) { + final row = find.byKey(ValueKey('thread-message-group-reply$i')); + if (row.evaluate().isEmpty) continue; + if (tester.getTopLeft(row).dy >= coveredTop) { + hiddenRows++; + expect(marked, isNot(contains('msg:reply$i')), reason: 'reply$i'); + } + if (marked.contains('msg:reply$i')) { + // The reading edge allows 1% of the list height. + expect( + tester.getBottomLeft(row).dy, + lessThanOrEqualTo(coveredTop + 10), + reason: 'reply$i', + ); + } + } + // Replies sit under the keyboard. Replies above it were read + // before the keyboard opened, so they need no new marks. + expect(hiddenRows, greaterThan(0)); + } finally { + debugDefaultTargetPlatformOverride = previousPlatform; + } + }); + testWidgets('shows forum posts view for forum channels', (tester) async { final forumChannel = Channel( id: _channelId, @@ -15897,6 +16508,23 @@ class _SynchronousReadStateNotifier extends ReadStateNotifier { markedContexts[contextId] = unixTimestamp; state = state.copyWithContext(contextId, unixTimestamp); } + + final List clearedForcedContexts = []; + + @override + void clearForcedUnread(String contextId) { + clearedForcedContexts.add(contextId); + state = ReadStateState( + isReady: state.isReady, + pubkey: state.pubkey, + contexts: state.contexts, + version: state.version + 1, + forcedUnreadContexts: { + for (final entry in state.forcedUnreadContexts.entries) + if (entry.key != contextId) entry.key: entry.value, + }, + ); + } } class _FakeProfileNotifier extends ProfileNotifier { diff --git a/mobile/test/features/channels/read_state/read_state_format_test.dart b/mobile/test/features/channels/read_state/read_state_format_test.dart index 57aa10d4caf..af73d6d0081 100644 --- a/mobile/test/features/channels/read_state/read_state_format_test.dart +++ b/mobile/test/features/channels/read_state/read_state_format_test.dart @@ -144,6 +144,218 @@ void main() { {'channel-a': 10, 'channel-b': 12, 'channel-c': 1}, ); }); + + group('retainReadStateContexts', () { + int publishedBytes(Map contexts) => utf8 + .encode( + jsonEncode( + ReadStateBlob(clientId: 'client-a', contexts: contexts).toJson(), + ), + ) + .length; + + test('counts JSON escaping in the slot budget', () { + // Each quote and backslash doubles when encoded. + final contexts = { + for (var index = 0; index < 400; index++) + 'msg:${'"\\' * 120}$index': index + 1, + }; + + final retained = retainReadStateContexts(contexts, clientId: 'client-a')!; + + expect(retained, isNotEmpty); + expect( + publishedBytes(retained), + lessThanOrEqualTo(readStatePlaintextBytes), + ); + }); + + test('keeps recent reads when old broad marks fill the budget', () { + final contexts = { + for (var index = 0; index < 300; index++) + index.toString().padLeft(36, 'c'): 100, + for (var index = 0; index < 380; index++) + 'thread:${index.toString().padLeft(64, 't')}': 100, + 'activity:fresh-channel': 500, + // An old message, read just now from history. + 'msg:${'a' * 64}': 50, + }; + final recent = { + for (final key in contexts.keys) key: 1000, + 'activity:fresh-channel': 2000, + 'msg:${'a' * 64}': 2000, + }; + expect(publishedBytes(contexts), greaterThan(readStatePlaintextBytes)); + + final retained = retainReadStateContexts( + contexts, + clientId: 'client-a', + recent: recent, + )!; + + expect(retained['activity:fresh-channel'], 500); + expect(retained['msg:${'a' * 64}'], 50); + expect( + publishedBytes(retained), + lessThanOrEqualTo(readStatePlaintextBytes), + ); + // Broad marks still fill most of the slot. + expect(retained.keys.where((key) => !key.contains(':')).length, 300); + }); + + test('keeps carried override keys whole ahead of marks', () { + final carried = { + for (var index = 0; index < 50; index++) ...{ + 'ov_s:channel-$index': 2, + 'ov_c:channel-$index': 1, + 'ov_b:channel-$index': 100, + }, + }; + final contexts = { + for (var index = 0; index < 50; index++) 'channel-$index': 100, + for (var index = 0; index < 1400; index++) + 'msg:${index.toString().padLeft(64, '0')}': index + 1, + }; + + final retained = retainReadStateContexts( + contexts, + clientId: 'client-a', + carried: carried, + )!; + + for (final entry in carried.entries) { + expect(retained[entry.key], entry.value); + } + expect(retained.length, lessThan(contexts.length + carried.length)); + expect( + publishedBytes(retained), + lessThanOrEqualTo(readStatePlaintextBytes), + ); + }); + + test('keeps each carried group with its frontier under pressure', () { + // A legal 240-byte context ID, and a raw ID that NIP-RS escapes. + final long = 'c' * 240; + final carried = { + 'ov_s:$long': 2, + 'ov_c:$long': 1, + 'ov_b:$long': 100, + 'ov_c:ov_x': 7, + 'esc:ov_x': 60, + }; + // 2,000 newer channel marks fill the slot many times over. + final contexts = { + long: 100, + for (var index = 0; index < 2000; index++) + 'channel-${index.toString().padLeft(36, '0')}': 1000 + index, + }; + + final retained = retainReadStateContexts( + contexts, + clientId: 'client-a', + recent: { + for (final key in contexts.keys) + if (key != long) key: 5000, + }, + carried: carried, + )!; + + expect(retained[long], 100); + for (final entry in carried.entries) { + expect(retained[entry.key], entry.value); + } + expect(retained.length, lessThan(contexts.length)); + expect( + publishedBytes(retained), + lessThanOrEqualTo(readStatePlaintextBytes), + ); + }); + + test('leaves out incomplete carried override groups', () { + final retained = retainReadStateContexts( + {'partial': 50, 'live': 60, 'dead': 70}, + clientId: 'client-a', + carried: const { + // No ov_c:, so the whole group is rejected. + 'ov_s:partial': 3, + 'ov_b:partial': 50, + 'ov_s:live': 2, + 'ov_c:live': 1, + 'ov_b:live': 60, + // A tombstone is ov_c: alone. + 'ov_c:dead': 4, + // A tombstone with a baseline is not a legal shape. + 'ov_c:odd': 4, + 'ov_b:odd': 9, + // A complete group must travel with its frontier. + 'ov_s:alone': 2, + 'ov_c:alone': 1, + 'ov_b:alone': 9, + }, + )!; + + expect(retained, { + 'partial': 50, + 'live': 60, + 'dead': 70, + 'ov_s:live': 2, + 'ov_c:live': 1, + 'ov_b:live': 60, + 'ov_c:dead': 4, + }); + }); + + test('returns null instead of splitting carried override keys', () { + final carried = { + for (var index = 0; index < 320; index++) ...{ + 'ov_s:${index.toString().padLeft(36, 'c')}': 2, + 'ov_c:${index.toString().padLeft(36, 'c')}': 1, + 'ov_b:${index.toString().padLeft(36, 'c')}': 100, + }, + }; + + expect( + retainReadStateContexts( + { + 'channel-1': 5, + for (var index = 0; index < 320; index++) + index.toString().padLeft(36, 'c'): 100, + }, + clientId: 'client-a', + carried: carried, + ), + isNull, + ); + }); + }); + + test('pruneStaleContexts bounds message and thread marks only', () { + const now = 10 * readStateHorizonSeconds; + final fresh = now - 60; + final stale = now - readStateHorizonSeconds - 1; + final contexts = { + 'channel-1': stale, + 'activity:channel-1': stale, + 'thread-activity:root': stale, + 'msg:stale': stale, + 'thread:stale': stale, + for (var index = 0; index < localMaxPrunableContexts + 10; index++) + 'msg:$index': fresh + index, + }; + + final kept = pruneStaleContexts(contexts, nowUnixSeconds: now); + + expect(kept['channel-1'], stale); + expect(kept['activity:channel-1'], stale); + expect(kept['thread-activity:root'], stale); + expect(kept, isNot(contains('msg:stale'))); + expect(kept, isNot(contains('thread:stale'))); + final messages = kept.keys.where((key) => key.startsWith('msg:')); + expect(messages.length, localMaxPrunableContexts); + // The oldest are dropped first. + expect(kept, isNot(contains('msg:0'))); + expect(kept, contains('msg:${localMaxPrunableContexts + 9}')); + }); } NostrEvent _event({List>? tags}) { diff --git a/mobile/test/features/channels/read_state/read_state_manager_test.dart b/mobile/test/features/channels/read_state/read_state_manager_test.dart index f453bf56311..ee6eb9cc371 100644 --- a/mobile/test/features/channels/read_state/read_state_manager_test.dart +++ b/mobile/test/features/channels/read_state/read_state_manager_test.dart @@ -1,11 +1,14 @@ import 'dart:async'; import 'dart:convert'; +import 'package:fake_async/fake_async.dart'; import 'package:flutter_test/flutter_test.dart'; import 'package:nostr/nostr.dart' as nostr; import 'package:shared_preferences/shared_preferences.dart'; import 'package:buzz/shared/read_state/read_state_format.dart'; import 'package:buzz/shared/read_state/read_state_manager.dart'; +import 'package:buzz/shared/read_state/read_state_storage.dart'; +import 'package:buzz/shared/read_state/read_state_time.dart'; import 'package:buzz/shared/relay/relay.dart'; void main() { @@ -107,7 +110,7 @@ void main() { }, ); - test('disables remote sync after an oversized local blob', () async { + test('caps a large slot instead of disabling sync', () async { SharedPreferences.setMockInitialValues({}); final prefs = await SharedPreferences.getInstance(); final keychain = nostr.Keys.generate(); @@ -126,20 +129,389 @@ void main() { onChanged: () {}, ); + // About 100 KiB of message marks, more than NIP-44 can encrypt. for (var index = 0; index < 1400; index++) { manager.markContextRead( - 'channel-${index.toString().padLeft(4, '0')}-${'x' * 48}', + 'msg:${index.toString().padLeft(64, '0')}', index + 1, ); } + manager.markContextRead('channel-1', 5); await manager.flush(); - manager.markContextRead('channel-new', 2000); + manager.markContextRead('msg:new', 2000); + await manager.flush(); + + expect(relay.submitCount, 2); + final plaintext = crypto.decrypt(relay.contents.last); + expect(utf8.encode(plaintext).length, lessThan(readStatePlaintextBytes)); + final published = decodeReadStateBlob(plaintext)!.contexts; + // Broad marks come first, then the newest message marks. + expect(published['channel-1'], 5); + expect(published['msg:new'], 2000); + expect(published, isNot(contains('msg:${'0' * 64}'))); + // Marks left out of the slot stay in local state. + expect(manager.getEffectiveTimestamp('msg:${'0' * 64}'), 1); + }); + + test('republishes merged broad marks but not message marks', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + final session = _FakeRelaySession(); + final relay = _FakeSignedEventRelay(); + final manager = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: session, + signedEventRelay: relay, + remoteEnabled: true, + onChanged: () {}, + ); + session.historyEvents = [ + _readStateEvent( + pubkey: keychain.public, + crypto: crypto, + clientId: 'web-client', + slotId: 'web-slot', + contexts: {'channel-1': 100, 'activity:channel-1': 120, 'msg:a': 110}, + createdAt: 100, + ), + ]; + + await manager.initialize(); + manager.markContextRead('msg:mine', 130); + await manager.flush(); + + // The web app may prune `msg:a` once `activity:` reads it, so this + // device must not bring it back. It still reads it locally. + expect(manager.getEffectiveTimestamp('msg:a'), 110); + final published = decodeReadStateBlob( + crypto.decrypt(relay.contents.last), + )!.contexts; + expect(published, { + 'channel-1': 100, + 'activity:channel-1': 120, + 'msg:mine': 130, + }); + }); + + test('publishes again when a debounce fires during a publish', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + + for (final failFirst in [false, true]) { + fakeAsync((async) { + final relay = _ParkedSignedEventRelay(failFirst: failFirst); + final manager = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: null, + signedEventRelay: relay, + remoteEnabled: true, + onChanged: () {}, + ); + + manager.markContextRead('a-$failFirst', 10); + async.elapse(const Duration(seconds: 5)); + // Publish A has its snapshot and waits on the relay. + expect(relay.contents, hasLength(1)); + + manager.markContextRead('b-$failFirst', 20); + async.elapse(const Duration(seconds: 5)); + expect(relay.contents, hasLength(1)); + + relay.release(); + async.flushMicrotasks(); + + expect(relay.contents, hasLength(2), reason: 'failFirst=$failFirst'); + final second = decodeReadStateBlob( + crypto.decrypt(relay.contents.last), + )!.contexts; + expect(second['b-$failFirst'], 20); + expect(second['a-$failFirst'], 10); + manager.dispose(flushPending: false); + }); + } + }); + + test('bounds saved marks across a restart', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + ReadStateManager create() => ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: null, + signedEventRelay: null, + remoteEnabled: false, + onChanged: () {}, + ); + + final now = currentUnixSeconds(); + final first = create(); + first.markContextRead('channel-1', now - 2 * readStateHorizonSeconds); + first.markContextRead('msg:stale', now - 2 * readStateHorizonSeconds); + for (var index = 0; index < localMaxPrunableContexts + 200; index++) { + first.markContextRead('msg:$index', now - 1000 + index % 900); + } + // This session still reads every mark. + expect(first.getEffectiveTimestamp('msg:stale'), isNotNull); + first.dispose(); + + final stored = ReadStateStorage(prefs).read(keychain.public); + final messages = stored.contexts.keys.where((k) => k.startsWith('msg:')); + expect(messages.length, localMaxPrunableContexts); + expect(stored.contexts, isNot(contains('msg:stale'))); + expect(stored.contexts['channel-1'], now - 2 * readStateHorizonSeconds); + // All three saved structures hold the same marks. + expect(stored.publishableContextIds, stored.contexts.keys.toSet()); + expect(stored.sourceCreatedAt.keys.toSet(), stored.contexts.keys.toSet()); + + final restarted = create(); + expect(restarted.getEffectiveTimestamp('msg:stale'), isNull); + expect( + restarted.getEffectiveTimestamp('channel-1'), + now - 2 * readStateHorizonSeconds, + ); + expect( + restarted.effectiveContexts.keys.where((k) => k.startsWith('msg:')), + hasLength(localMaxPrunableContexts), + ); + }); + + test('carries its own override keys and takes none from others', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + final storage = ReadStateStorage(prefs); + final clientId = storage.getOrCreateClientId(keychain.public); + final slotId = storage.getOrCreateSlotId(keychain.public); + const ownGroup = { + 'ov_s:channel-1': 2, + 'ov_c:channel-1': 1, + 'ov_b:channel-1': 90, + }; + final oldMessage = 'msg:${'d' * 64}'; + final oldMessageGroup = {'ov_c:$oldMessage': 96}; + final session = _FakeRelaySession() + ..historyEvents = [ + _readStateEvent( + pubkey: keychain.public, + crypto: crypto, + clientId: clientId, + slotId: slotId, + // `ov_s:channel-9` alone is a partial group, rejected whole. + contexts: { + 'channel-1': 90, + ...ownGroup, + 'ov_s:channel-9': 3, + // Older than the local save horizon, so saving prunes it unless + // its group protects it. + oldMessage: 95, + ...oldMessageGroup, + }, + createdAt: 100, + ), + _readStateEvent( + pubkey: keychain.public, + crypto: crypto, + clientId: 'web-client', + slotId: 'web-slot', + contexts: {'ov_c:channel-2': 4, 'esc:ov_x': 5, 'channel-2': 80}, + createdAt: 100, + ), + ]; + final relay = _FakeSignedEventRelay(); + final manager = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: session, + signedEventRelay: relay, + remoteEnabled: true, + onChanged: () {}, + ); + + await manager.initialize(); + // About 100 KiB of message marks, so retention must leave some out. + for (var index = 0; index < 1400; index++) { + manager.markContextRead('msg:${index.toString().padLeft(64, '0')}', 200); + } + await manager.flush(); + + final published = decodeReadStateBlob( + crypto.decrypt(relay.contents.last), + )!.contexts; + expect(published, containsPair('ov_s:channel-1', 2)); + expect(published, containsPair('ov_c:channel-1', 1)); + expect(published, containsPair('ov_b:channel-1', 90)); + // The group's frontier travels with it. + expect(published, containsPair('channel-1', 90)); + expect(published, isNot(contains('ov_s:channel-9'))); + expect(published, isNot(contains('ov_c:channel-2'))); + expect(published, isNot(contains('esc:ov_x'))); + expect(manager.getEffectiveTimestamp('ov_s:channel-1'), isNull); + manager.dispose(flushPending: false); + + // A restart without the relay still carries the group. + final restartedRelay = _FakeSignedEventRelay(); + final restarted = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: null, + signedEventRelay: restartedRelay, + remoteEnabled: true, + onChanged: () {}, + ); + restarted.markContextRead('channel-3', 300); + await restarted.flush(); + final republished = decodeReadStateBlob( + crypto.decrypt(restartedRelay.contents.last), + )!.contexts; + for (final entry in ownGroup.entries) { + expect(republished, containsPair(entry.key, entry.value)); + } + expect(republished, containsPair('channel-1', 90)); + expect(republished, containsPair(oldMessage, 95)); + expect(republished, containsPair('ov_c:$oldMessage', 96)); + expect(republished, isNot(contains('ov_s:channel-9'))); + expect(republished, isNot(contains('ov_c:channel-2'))); + }); + + test('rejects own-slot override groups that are malformed or have no ' + 'frontier', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + final storage = ReadStateStorage(prefs); + final clientId = storage.getOrCreateClientId(keychain.public); + final slotId = storage.getOrCreateSlotId(keychain.public); + final session = _FakeRelaySession() + ..historyEvents = [ + _readStateEvent( + pubkey: keychain.public, + crypto: crypto, + clientId: clientId, + slotId: slotId, + contexts: const { + // Its ov_b: sibling is invalid, so this is not a tombstone. + 'bad-sibling': 100, + 'ov_c:bad-sibling': 7, + // Complete counters, but no frontier for `no-frontier`. + 'ov_s:no-frontier': 2, + 'ov_c:no-frontier': 1, + 'ov_b:no-frontier': 50, + 'channel-1': 90, + 'ov_s:channel-1': 2, + 'ov_c:channel-1': 1, + 'ov_b:channel-1': 90, + }, + rawContexts: const {'ov_b:bad-sibling': 'invalid'}, + createdAt: 100, + ), + ]; + final relay = _FakeSignedEventRelay(); + final manager = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: session, + signedEventRelay: relay, + remoteEnabled: true, + onChanged: () {}, + ); + + await manager.initialize(); + manager.markContextRead('channel-2', 300); + await manager.flush(); + + final published = decodeReadStateBlob( + crypto.decrypt(relay.contents.last), + )!.contexts; + expect(published, { + // Frontiers stay, including the one whose group was rejected. + 'bad-sibling': 100, + 'channel-1': 90, + 'channel-2': 300, + 'ov_s:channel-1': 2, + 'ov_c:channel-1': 1, + 'ov_b:channel-1': 90, + }); + manager.dispose(flushPending: false); + }); + + test('leaves the slot when carried override keys do not fit', () async { + SharedPreferences.setMockInitialValues({}); + final prefs = await SharedPreferences.getInstance(); + final keychain = nostr.Keys.generate(); + final crypto = ReadStateCrypto.tryCreate( + nsec: keychain.nsec, + pubkey: keychain.public, + )!; + final storage = ReadStateStorage(prefs); + final clientId = storage.getOrCreateClientId(keychain.public); + final slotId = storage.getOrCreateSlotId(keychain.public); + final session = _FakeRelaySession() + ..historyEvents = [ + _readStateEvent( + pubkey: keychain.public, + crypto: crypto, + clientId: clientId, + slotId: slotId, + contexts: { + for (var index = 0; index < 320; index++) ...{ + index.toString().padLeft(36, 'c'): 100, + 'ov_s:${index.toString().padLeft(36, 'c')}': 2, + 'ov_c:${index.toString().padLeft(36, 'c')}': 1, + 'ov_b:${index.toString().padLeft(36, 'c')}': 100, + }, + }, + createdAt: 100, + ), + ]; + final relay = _FakeSignedEventRelay(); + final manager = ReadStateManager( + pubkey: keychain.public, + prefs: prefs, + crypto: crypto, + relaySession: session, + signedEventRelay: relay, + remoteEnabled: true, + onChanged: () {}, + ); + + await manager.initialize(); + manager.markContextRead('channel-1', 300); await manager.flush(); expect(relay.submitCount, 0); - expect(manager.getEffectiveTimestamp('channel-0000-${'x' * 48}'), 1); - expect(manager.getEffectiveTimestamp('channel-new'), 2000); + expect(manager.getEffectiveTimestamp('channel-1'), 300); }); test('remote read-state rollback is ignored', () async { @@ -206,6 +578,7 @@ NostrEvent _stubAckEvent() => const NostrEvent( class _FakeSignedEventRelay implements SignedEventRelay { final Completer<_SubmittedEvent> submitted = Completer<_SubmittedEvent>(); + final List contents = []; int submitCount = 0; @override @@ -220,7 +593,40 @@ class _FakeSignedEventRelay implements SignedEventRelay { void Function(NostrEvent event)? onSigned, }) async { submitCount++; - submitted.complete(_SubmittedEvent(kind: kind, tags: tags)); + contents.add(content); + if (!submitted.isCompleted) { + submitted.complete(_SubmittedEvent(kind: kind, tags: tags)); + } + return _stubAckEvent(); + } +} + +/// Holds the first submit until [release], then fails it if [failFirst]. +class _ParkedSignedEventRelay implements SignedEventRelay { + _ParkedSignedEventRelay({required this.failFirst}); + + final bool failFirst; + final List contents = []; + final Completer _parked = Completer(); + + void release() => _parked.complete(); + + @override + String? get pubkey => null; + + @override + Future submit({ + required int kind, + required String content, + required List> tags, + int? createdAt, + void Function(NostrEvent event)? onSigned, + }) async { + contents.add(content); + if (contents.length == 1) { + await _parked.future; + if (failFirst) throw Exception('relay timeout'); + } return _stubAckEvent(); } } @@ -270,8 +676,11 @@ NostrEvent _readStateEvent({ required String slotId, required Map contexts, required int createdAt, + // Extra raw wire entries, such as values that are not valid timestamps. + Map rawContexts = const {}, }) { - final blob = ReadStateBlob(clientId: clientId, contexts: contexts); + final blob = ReadStateBlob(clientId: clientId, contexts: contexts).toJson(); + blob['contexts'] = {...contexts, ...rawContexts}; return NostrEvent( id: 'event-$clientId-$createdAt', pubkey: pubkey, @@ -281,7 +690,7 @@ NostrEvent _readStateEvent({ ['d', '$readStateDTagPrefix$slotId'], ['t', 'read-state'], ], - content: crypto.encrypt(jsonEncode(blob.toJson())), + content: crypto.encrypt(jsonEncode(blob)), sig: 'sig', ); } diff --git a/mobile/test/features/channels/read_state/read_state_provider_test.dart b/mobile/test/features/channels/read_state/read_state_provider_test.dart index f01aa2e81e4..8a5500ca36f 100644 --- a/mobile/test/features/channels/read_state/read_state_provider_test.dart +++ b/mobile/test/features/channels/read_state/read_state_provider_test.dart @@ -105,6 +105,23 @@ void main() { }, ); + test( + 'clearing a channel force keeps its marker and message forces', + () async { + final notifier = await pumpNotifier(); + + notifier.markContextUnread(channelId, channelId: channelId); + notifier.markContextUnread(msgKey, channelId: channelId); + + // Reading to the bottom ends the channel force without reading anything. + notifier.clearForcedUnread(channelId); + + expect(state().isForcedUnread(channelId), isFalse); + expect(state().isForcedUnread(msgKey), isTrue); + expect(state().effectiveTimestamp(channelId), isNull); + }, + ); + test('automatic channel-open read still clears a channel-level force for ' 'the same channel', () async { final notifier = await pumpNotifier(); diff --git a/mobile/test/features/channels/reading_marks_test.dart b/mobile/test/features/channels/reading_marks_test.dart new file mode 100644 index 00000000000..4b0b510e0f8 --- /dev/null +++ b/mobile/test/features/channels/reading_marks_test.dart @@ -0,0 +1,346 @@ +import 'package:buzz/features/channels/reading_marks.dart'; +import 'package:buzz/features/channels/timeline_message.dart'; +import 'package:buzz/shared/read_state/read_state_provider.dart'; +import 'package:flutter/widgets.dart'; +import 'package:flutter_hooks/flutter_hooks.dart'; +import 'package:flutter_test/flutter_test.dart'; + +const _channel = 'channel-1'; +const _self = 'self'; + +TimelineMessage _msg( + String id, + int createdAt, { + String pubkey = 'alice', + String? parentId, + String? rootId, + List> tags = const [], +}) => TimelineMessage( + id: id, + pubkey: pubkey, + createdAt: createdAt, + content: id, + tags: tags, + parentId: parentId, + rootId: rootId, +); + +ReadStateState _state([ + Map contexts = const {}, + Map forced = const {}, +]) => ReadStateState( + isReady: true, + pubkey: _self, + contexts: contexts, + version: 0, + forcedUnreadContexts: forced, +); + +Map _marks({ + ReadStateState? readState, + bool isDm = false, + required List loaded, + List visible = const [], + TimelineMessage? bottom, + String? threadRootId, + bool isRootThread = true, + int now = 10000, +}) => readingMarks( + readState: readState ?? _state(), + channelId: _channel, + isDm: isDm, + currentPubkey: _self, + loaded: loaded, + visible: visible, + bottom: bottom, + threadRootId: threadRootId, + isRootThread: isRootThread, + now: now, +); + +void main() { + group('channel timeline', () { + test('the bottom writes one catch-up mark instead of message marks', () { + final a = _msg('a', 100); + final b = _msg('b', 200); + expect(_marks(loaded: [a, b], visible: [a, b], bottom: b), { + 'activity:$_channel': 200, + }); + }); + + test('catch-up cuts at the bottom row, not a newer reply', () { + // A top-level message created at 200 that arrives late must stay + // unread, so the cut cannot move up to the reply at 300. + final a = _msg('a', 100); + final reply = _msg('r', 300, parentId: 'a', rootId: 'a'); + expect(_marks(loaded: [a, reply], visible: [a], bottom: a), { + 'activity:$_channel': 100, + }); + }); + + test('catch-up never passes the current time', () { + // A sender clock ahead of this one must not read messages that have + // not arrived yet. + final future = _msg('f', 20000); + expect( + _marks(loaded: [future], visible: [future], bottom: future, now: 10000), + {'activity:$_channel': 10000, 'msg:f': 20000}, + ); + }); + + test('a newer top-level message means the bottom is not live', () { + final a = _msg('a', 100); + final b = _msg('b', 200); + expect(_marks(loaded: [a, b], visible: [a], bottom: a), {'msg:a': 100}); + }); + + test('a newer broadcast reply means the bottom is not live', () { + final a = _msg('a', 100); + final broadcast = _msg( + 'r', + 300, + parentId: 'a', + rootId: 'a', + tags: const [ + ['broadcast', '1'], + ], + ); + expect(_marks(loaded: [a, broadcast], visible: [a], bottom: a), { + 'msg:a': 100, + }); + }); + + test('mentions and broadcasts keep their own marks at the bottom', () { + final mention = _msg( + 'm', + 100, + tags: const [ + ['p', _self], + ], + ); + final b = _msg('b', 200); + expect(_marks(loaded: [mention, b], visible: [mention, b], bottom: b), { + 'activity:$_channel': 200, + 'msg:m': 100, + }); + }); + + test('away from the bottom, visible unread rows get message marks', () { + final a = _msg('a', 100); + final b = _msg('b', 200); + final c = _msg('c', 300); + expect( + _marks( + readState: _state({'activity:$_channel': 150}), + loaded: [a, b, c], + visible: [a, b], + ), + {'msg:b': 200}, + ); + }); + + test('never writes the channel mark', () { + final a = _msg('a', 100); + final marks = _marks(loaded: [a], visible: [a], bottom: a); + expect(marks.containsKey(_channel), isFalse); + }); + + test('writes nothing when existing marks already read everything', () { + final a = _msg('a', 100); + expect( + _marks( + readState: _state({_channel: 100}), + loaded: [a], + visible: [a], + bottom: a, + ), + isEmpty, + ); + expect( + _marks( + readState: _state({'activity:$_channel': 100}), + loaded: [a], + visible: [a], + bottom: a, + ), + isEmpty, + ); + }); + + test('skips own, system and manually unread messages', () { + final own = _msg('own', 100, pubkey: _self); + final system = TimelineMessage( + id: 'sys', + pubkey: 'alice', + createdAt: 110, + content: '', + isSystem: true, + ); + final forced = _msg('forced', 120); + expect( + _marks( + readState: _state(const {}, {'msg:forced': _channel}), + loaded: [own, system, forced, _msg('later', 200)], + visible: [own, system, forced], + ), + isEmpty, + ); + }); + + test('in a DM, catch-up does not read messages', () { + final a = _msg('a', 100); + final b = _msg('b', 200); + expect(_marks(isDm: true, loaded: [a, b], visible: [a, b], bottom: b), { + 'activity:$_channel': 200, + 'msg:a': 100, + 'msg:b': 200, + }); + }); + }); + + group('thread', () { + final root = _msg('root', 50); + TimelineMessage reply(String id, int at, {List>? tags}) => + _msg(id, at, parentId: 'root', rootId: 'root', tags: tags ?? const []); + + test('the tail writes one thread catch-up mark', () { + final r1 = reply('r1', 100); + final mention = reply( + 'r2', + 200, + tags: const [ + ['p', _self], + ], + ); + expect( + _marks( + loaded: [root, r1, mention], + visible: [r1, mention], + bottom: mention, + threadRootId: 'root', + ), + // Thread catch-up reads a reply only while its root is loaded, so + // each visible reply keeps its own mark, as in `buzz-app`. + {'thread-activity:root': 200, 'msg:r1': 100, 'msg:r2': 200}, + ); + }); + + test('thread catch-up uses the newest reply, not newer channel rows', () { + final r1 = reply('r1', 100); + final later = _msg('later', 400); + expect( + _marks( + loaded: [root, r1, later], + visible: [r1], + bottom: r1, + threadRootId: 'root', + ), + {'thread-activity:root': 100, 'msg:r1': 100}, + ); + }); + + test('a newer reply in another branch means the tail is not live', () { + final r1 = reply('r1', 100); + final nested = _msg('n1', 300, parentId: 'other', rootId: 'root'); + expect( + _marks( + loaded: [root, r1, nested], + visible: [r1], + bottom: r1, + threadRootId: 'root', + ), + {'msg:r1': 100}, + ); + }); + + test('away from the tail, visible replies get message marks', () { + final r1 = reply('r1', 100); + final r2 = reply('r2', 200); + expect( + _marks(loaded: [root, r1, r2], visible: [r1], threadRootId: 'root'), + {'msg:r1': 100}, + ); + }); + + test('a thread mark replaces catch-up but not reply marks', () { + final r1 = reply('r1', 100); + expect( + _marks( + readState: _state({'thread:root': 100}), + loaded: [root, r1], + visible: [r1], + bottom: r1, + threadRootId: 'root', + ), + {'msg:r1': 100}, + ); + }); + + test('a nested thread head does not catch up the whole thread', () { + // Opened on a nested reply, the page shows one branch. An older + // mention in another branch must stay unread. + final branch = _msg('b1', 100, parentId: 'root', rootId: 'root'); + final nested = _msg('n1', 200, parentId: 'b1', rootId: 'root'); + expect( + _marks( + loaded: [root, branch, nested], + visible: [nested], + bottom: nested, + threadRootId: 'root', + isRootThread: false, + ), + {'msg:n1': 200}, + ); + }); + }); + + group('readingContentKey', () { + test('ignores a rebuilt list with the same messages', () { + final a = _msg('a', 100); + final b = _msg('b', 200); + expect(readingContentKey([a, b]), readingContentKey([a, b].toList())); + expect(readingContentKey([a, b]), isNot(readingContentKey([a]))); + expect( + readingContentKey([a, b]), + isNot(readingContentKey([a, _msg('c', 300)])), + ); + }); + }); + + testWidgets('the dwell fires while the page rebuilds every 200 ms', ( + tester, + ) async { + // A typing indicator or stream rebuilds the page with a new message + // list. The dwell must still finish while the messages stay the same. + final positions = ValueNotifier(null); + addTearDown(positions.dispose); + final tick = ValueNotifier(0); + addTearDown(tick.dispose); + final a = _msg('a', 100); + var fired = 0; + await tester.pumpWidget( + ValueListenableBuilder( + valueListenable: tick, + builder: (_, _, _) => HookBuilder( + builder: (_) { + final messages = [a]; + useReadingDwell( + positions: positions, + active: true, + onDwell: () => fired++, + keys: [readingContentKey(messages)], + ); + return const SizedBox(); + }, + ), + ), + ); + // Rebuild at 200, 400 and 600 ms. A restarted dwell would not finish. + for (var i = 0; i < 3; i++) { + await tester.pump(const Duration(milliseconds: 200)); + tick.value++; + await tester.pump(); + } + expect(fired, 1); + }); +}