diff --git a/mobile/lib/features/channels/channel_detail_page.dart b/mobile/lib/features/channels/channel_detail_page.dart index b195dc9eb81..a5493de11b8 100644 --- a/mobile/lib/features/channels/channel_detail_page.dart +++ b/mobile/lib/features/channels/channel_detail_page.dart @@ -107,17 +107,99 @@ Future _loadDeepLinkEvents( } /// Fetch channel members and preload their profiles into the user cache. -Future _preloadMembers(WidgetRef ref, String channelId) async { +/// One-to-one DMs additionally refresh participant profiles for identity gates. +/// Returns whether identity resolution completed successfully. +Future _preloadMembers( + WidgetRef ref, + String channelId, + List participantPubkeys, { + required bool refreshDmParticipants, +}) async { // Capture references before async gap to avoid using disposed ref. final notifier = ref.read(userCacheProvider.notifier); try { final members = await ref.read(channelMembersProvider(channelId).future); - final pubkeys = members.map((m) => m.pubkey).toList(); - if (pubkeys.isNotEmpty) { - notifier.preload(pubkeys); + await notifier.preload(members.map((member) => member.pubkey).toList()); + if (refreshDmParticipants) { + return notifier.refresh(participantPubkeys); + } + return true; + } catch (_) { + // Identity remains unresolved, so agent-only actions stay hidden. + return false; + } +} + +Future _subscribeToDmIdentityUpdates( + WidgetRef ref, + List participantPubkeys, { + required ValueChanged onReadyChanged, + required ValueChanged> onAgentPubkeysChanged, + required VoidCallback onFailure, +}) async { + final session = ref.read(relaySessionProvider.notifier); + var subscriptionStatus = RelaySubscriptionStatus.retrying; + var directLookupComplete = false; + final agentPubkeys = {}; + + void publishAgentPubkeys() { + onAgentPubkeysChanged(Set.unmodifiable(agentPubkeys)); + } + + void handleEvent(NostrEvent event) { + if (event.kind == 0) { + try { + ref.read(userCacheProvider.notifier).cacheProfileEvent(event); + } catch (error) { + debugPrint('[DmIdentity] invalid live profile: $error'); + onFailure(); + } + } else if (event.kind == 10100) { + agentPubkeys.add(event.pubkey.toLowerCase()); + publishAgentPubkeys(); + ref.invalidate(agentDirectoryProvider); + ref.invalidate(agentOwnersProvider); } + } + + final unsubscribe = await session.subscribeWithStatus( + NostrFilter( + kinds: const [0, 10100], + authors: participantPubkeys, + limit: 100, + ).copyWithSince(DateTime.now().millisecondsSinceEpoch ~/ 1000 - 5), + handleEvent, + onClosed: (_) => onFailure(), + onStatusChanged: (status) { + subscriptionStatus = status; + if (status == RelaySubscriptionStatus.retrying) { + onReadyChanged(false); + } else if (directLookupComplete) { + onReadyChanged(true); + } + }, + ); + + try { + final profiles = await session.fetchHistory( + NostrFilter( + kinds: const [10100], + authors: participantPubkeys, + limit: participantPubkeys.length, + ), + ); + for (final profile in profiles) { + if (profile.kind == 10100) { + agentPubkeys.add(profile.pubkey.toLowerCase()); + } + } + publishAgentPubkeys(); + directLookupComplete = true; + onReadyChanged(subscriptionStatus == RelaySubscriptionStatus.ready); + return unsubscribe; } catch (_) { - // Non-fatal — mentions will just fall back to cache from messages. + unsubscribe(); + rethrow; } } @@ -146,6 +228,16 @@ int? _channelReadTimestamp({ return dateTimeToUnixSeconds(channel.lastMessageAt); } +bool _isOneToOneAgentDm(Channel channel, Set agentPubkeys) { + final participants = channel.participantPubkeys + .map((pubkey) => pubkey.trim().toLowerCase()) + .where((pubkey) => pubkey.isNotEmpty) + .toSet(); + return channel.isDm && + participants.length == 2 && + participants.any(agentPubkeys.contains); +} + /// Controls how a hydrated initial thread is added to the navigation stack. enum InitialThreadRouteBehavior { /// Keep the channel route beneath the thread. @@ -256,10 +348,147 @@ class ChannelDetailPage extends HookConsumerWidget { channel; final resolvedChannel = detailsAsync.whenData(baseChannel.mergeDetails).value ?? baseChannel; + final participantCount = resolvedChannel.participantPubkeys + .map((pubkey) => pubkey.trim().toLowerCase()) + .where((pubkey) => pubkey.isNotEmpty) + .toSet() + .length; + final isOneToOneDm = resolvedChannel.isDm && participantCount == 2; + final memberProfilesPreload = useMemoized( + () => _preloadMembers( + ref, + resolvedChannel.id, + resolvedChannel.participantPubkeys, + refreshDmParticipants: isOneToOneDm, + ), + [ + resolvedChannel.id, + sessionStatus, + isOneToOneDm, + Object.hashAll(resolvedChannel.participantPubkeys), + ], + ); + final memberProfilesPreloadState = useFuture(memberProfilesPreload); final showsComposer = !resolvedChannel.isForum && resolvedChannel.isMember && !resolvedChannel.isArchived; + final profileOwnedAgentPubkeys = []; + for (final participantPubkey in resolvedChannel.participantPubkeys) { + final normalized = participantPubkey.trim().toLowerCase(); + final isProfileOwnedAgent = ref.watch( + userCacheProvider.select( + (cache) => cache[normalized]?.ownerPubkey != null, + ), + ); + if (isProfileOwnedAgent) profileOwnedAgentPubkeys.add(normalized); + } + final agentDirectoryState = ref.watch(agentDirectoryProvider); + final agentOwnersState = ref.watch(agentOwnersProvider); + final channelMembershipUpdateState = isOneToOneDm + ? ref.watch(channelMembershipUpdateProvider(resolvedChannel.id)) + : const ChannelMembershipUpdateState(isReady: true); + final channelBotPubkeysState = ref.watch( + channelBotPubkeysProvider(resolvedChannel.id), + ); + final identitySubscriptionPubkeys = isOneToOneDm + ? (resolvedChannel.participantPubkeys + .map((pubkey) => pubkey.trim().toLowerCase()) + .where((pubkey) => pubkey.isNotEmpty) + .toSet() + .toList() + ..sort()) + : const []; + final identitySubscriptionKey = Object.hashAll(identitySubscriptionPubkeys); + final identitySubscriptionReady = useValueNotifier(false, [ + sessionStatus, + resolvedChannel.id, + identitySubscriptionKey, + ]); + final directlyResolvedAgentPubkeys = useValueNotifier({}, [ + sessionStatus, + resolvedChannel.id, + identitySubscriptionKey, + ]); + final isIdentitySubscriptionReady = useValueListenable( + identitySubscriptionReady, + ); + final directAgentPubkeys = useValueListenable(directlyResolvedAgentPubkeys); + final agentPubkeys = agentPubkeysWithChannelBots( + knownAgentPubkeys: agentPubkeysWithProfileOwners( + knownAgentPubkeys: { + ...ref.watch(knownAgentPubkeysProvider), + ...directAgentPubkeys, + }, + profileOwnedAgentPubkeys: profileOwnedAgentPubkeys, + ), + channelBotPubkeys: + channelBotPubkeysState.asData?.value ?? const {}, + ); + useEffect(() { + if (sessionStatus != SessionStatus.connected || + identitySubscriptionPubkeys.isEmpty) { + return null; + } + var disposed = false; + var subscriptionFailed = false; + void markFailed() { + subscriptionFailed = true; + if (!disposed) identitySubscriptionReady.value = false; + } + + void Function()? unsubscribe; + Future.microtask(() async { + try { + final cleanup = await _subscribeToDmIdentityUpdates( + ref, + identitySubscriptionPubkeys, + onReadyChanged: (isReady) { + if (!disposed && !subscriptionFailed) { + identitySubscriptionReady.value = isReady; + } + }, + onAgentPubkeysChanged: (pubkeys) { + if (!disposed) directlyResolvedAgentPubkeys.value = pubkeys; + }, + onFailure: markFailed, + ); + if (disposed) { + cleanup(); + } else { + unsubscribe = cleanup; + } + } catch (error) { + if (!disposed) { + debugPrint('[DmIdentity] live subscription failed: $error'); + markFailed(); + } + } + }); + return () { + disposed = true; + unsubscribe?.call(); + }; + }, [sessionStatus, resolvedChannel.id, identitySubscriptionKey]); + final isAgentIdentityUnresolved = + isOneToOneDm && + (sessionStatus != SessionStatus.connected || + !isIdentitySubscriptionReady || + agentDirectoryState.isLoading || + agentDirectoryState.hasError || + agentOwnersState.isLoading || + agentOwnersState.hasError || + !channelMembershipUpdateState.isReady || + channelMembershipUpdateState.error != null || + channelBotPubkeysState.isLoading || + channelBotPubkeysState.hasError || + memberProfilesPreloadState.connectionState != + ConnectionState.done || + memberProfilesPreloadState.data != true); + final showsHuddleAction = + showsComposer && + !isAgentIdentityUnresolved && + !_isOneToOneAgentDm(resolvedChannel, agentPubkeys); final messagesNotifier = ref.read( channelMessagesProvider(channel.id).notifier, ); @@ -301,12 +530,6 @@ class ChannelDetailPage extends HookConsumerWidget { return session.registerVisibleChannel(channel.id); }, [channel.id]); - // Preload channel member profiles so @mentions resolve correctly. - useEffect(() { - _preloadMembers(ref, channel.id); - return null; - }, [channel.id]); - useEffect( () { if (channel.isForum) return null; @@ -394,7 +617,7 @@ class ChannelDetailPage extends HookConsumerWidget { ), actions: resolvedChannel.isDm ? [ - if (showsComposer) + if (showsHuddleAction) _HuddleButton( channel: resolvedChannel, events: [ diff --git a/mobile/lib/features/channels/channel_management_provider.dart b/mobile/lib/features/channels/channel_management_provider.dart index 9e021ab5178..119f5dce7b0 100644 --- a/mobile/lib/features/channels/channel_management_provider.dart +++ b/mobile/lib/features/channels/channel_management_provider.dart @@ -488,7 +488,11 @@ final channelDetailsProvider = FutureProvider.family(( /// Channel members from kind:39002 NIP-29 members event. final channelMembersProvider = FutureProvider.autoDispose .family, String>((ref, channelId) async { - ref.watch(channelMembershipUpdateProvider(channelId)); + ref.watch( + channelMembershipUpdateProvider( + channelId, + ).select((update) => update.version), + ); final relayBaseUrl = ref.watch(relayConfigProvider).baseUrl; final pubkey = ref.watch(myPubkeyProvider)?.toLowerCase(); final snapshotCache = ref.read(_channelMembersSnapshotCacheProvider); diff --git a/mobile/lib/shared/mentions/agent_identity_provider.dart b/mobile/lib/shared/mentions/agent_identity_provider.dart index ea6a2ee50fd..35c063d13d9 100644 --- a/mobile/lib/shared/mentions/agent_identity_provider.dart +++ b/mobile/lib/shared/mentions/agent_identity_provider.dart @@ -161,10 +161,24 @@ Map mentionNamesWithDirectoryLabels({ String _agentFallbackLabel(String pubkey) => pubkey.length >= 8 ? pubkey.substring(0, 8) : pubkey; +/// Readiness and refresh state for a channel's live membership subscription. +class ChannelMembershipUpdateState { + final int version; + final bool isReady; + final Object? error; + + const ChannelMembershipUpdateState({ + this.version = 0, + this.isReady = false, + this.error, + }); +} + /// Keeps the role feed alive for consumers that render mentions outside the /// channel timeline, such as search results. A membership change refreshes the /// shared bot-role lookup below, regardless of which surface owns the channel. -class _ChannelBotRoleSubscription extends Notifier { +class _ChannelBotRoleSubscription + extends Notifier { final String channelId; void Function()? _unsubscribe; int _subscriptionVersion = 0; @@ -172,7 +186,7 @@ class _ChannelBotRoleSubscription extends Notifier { _ChannelBotRoleSubscription(this.channelId); @override - int build() { + ChannelMembershipUpdateState build() { final sessionState = ref.watch(relaySessionProvider); final subscriptionVersion = ++_subscriptionVersion; _clearSubscription(); @@ -181,34 +195,64 @@ class _ChannelBotRoleSubscription extends Notifier { _clearSubscription(); }); - if (sessionState.status != SessionStatus.connected) return 0; + if (sessionState.status != SessionStatus.connected) { + return const ChannelMembershipUpdateState(); + } Future.microtask(() => _subscribe(channelId, subscriptionVersion)); - return 0; + return const ChannelMembershipUpdateState(); } Future _subscribe(String channelId, int subscriptionVersion) async { final session = ref.read(relaySessionProvider.notifier); + var subscriptionStatus = RelaySubscriptionStatus.retrying; try { - final unsubscribe = await session.subscribe( + final unsubscribe = await session.subscribeWithStatus( NostrFilter( kinds: const [39002], tags: { - '#h': [channelId], + '#d': [channelId], }, ).copyWithSince(DateTime.now().millisecondsSinceEpoch ~/ 1000), (_) { if (_isCurrent(subscriptionVersion)) { - state++; + state = ChannelMembershipUpdateState( + version: state.version + 1, + isReady: state.isReady, + ); + } + }, + onClosed: (message) { + if (_isCurrent(subscriptionVersion)) { + state = ChannelMembershipUpdateState( + version: state.version, + error: Exception(message), + ); } }, + onStatusChanged: (status) { + subscriptionStatus = status; + if (!_isCurrent(subscriptionVersion)) return; + state = ChannelMembershipUpdateState( + version: state.version, + isReady: status == RelaySubscriptionStatus.ready, + ); + }, ); if (!_isCurrent(subscriptionVersion)) { unsubscribe(); return; } _unsubscribe = unsubscribe; + state = ChannelMembershipUpdateState( + version: state.version, + isReady: subscriptionStatus == RelaySubscriptionStatus.ready, + ); } catch (error) { if (_isCurrent(subscriptionVersion)) { + state = ChannelMembershipUpdateState( + version: state.version, + error: error, + ); debugPrint( '[ChannelBotRoleSubscription] failed for $channelId: $error', ); @@ -229,14 +273,18 @@ class _ChannelBotRoleSubscription extends Notifier { /// changes. Channel-member and agent-role views share this source so remote /// membership updates refresh both snapshots together. final channelMembershipUpdateProvider = NotifierProvider.autoDispose - .family<_ChannelBotRoleSubscription, int, String>( + .family<_ChannelBotRoleSubscription, ChannelMembershipUpdateState, String>( _ChannelBotRoleSubscription.new, ); /// Bot pubkeys currently assigned a channel bot role. final channelBotPubkeysProvider = FutureProvider.autoDispose .family, String>((ref, channelId) async { - ref.watch(channelMembershipUpdateProvider(channelId)); + ref.watch( + channelMembershipUpdateProvider( + channelId, + ).select((update) => update.version), + ); final sessionState = ref.watch(relaySessionProvider); if (sessionState.status != SessionStatus.connected) return const {}; final session = ref.read(relaySessionProvider.notifier); diff --git a/mobile/lib/shared/profile/user_cache_provider.dart b/mobile/lib/shared/profile/user_cache_provider.dart index fd8d3fafcd7..d11db73b1cf 100644 --- a/mobile/lib/shared/profile/user_cache_provider.dart +++ b/mobile/lib/shared/profile/user_cache_provider.dart @@ -12,14 +12,19 @@ import 'user_profile.dart'; /// kind:0 batch query (NIP-01 `authors` filter) every 50ms. class UserCacheNotifier extends Notifier> { final Set _pending = {}; + final Map _profileEventOrders = {}; Timer? _batchTimer; + Completer? _batchCompleter; @override Map build() { ref.watch(relayConfigProvider); + _profileEventOrders.clear(); ref.onDispose(() { _batchTimer?.cancel(); _batchTimer = null; + _batchCompleter?.complete(false); + _batchCompleter = null; }); return {}; } @@ -39,14 +44,52 @@ class UserCacheNotifier extends Notifier> { } /// Preload profiles for a list of pubkeys (e.g. channel members). - void preload(List pubkeys) { - final uncached = pubkeys + /// Returns whether the batch completed successfully. + Future preload(List pubkeys) { + final normalized = pubkeys.map((pk) => pk.toLowerCase()).toSet(); + final alreadyPending = normalized.any(_pending.contains); + final uncached = normalized .map((pk) => pk.toLowerCase()) .where((pk) => !state.containsKey(pk) && !_pending.contains(pk)) .toList(); - if (uncached.isEmpty) return; + if (uncached.isEmpty && !alreadyPending) return Future.value(true); _pending.addAll(uncached); + final completer = _batchCompleter ??= Completer(); _batchTimer ??= Timer(const Duration(milliseconds: 50), _flushPending); + return completer.future; + } + + /// Force-refresh profiles for identity-sensitive gates. + /// + /// Unlike [preload], this fetches cached pubkeys too so stale human profiles + /// cannot be trusted after a verified agent-owner profile was published. + Future refresh(List pubkeys) async { + final normalized = pubkeys + .map((pubkey) => pubkey.toLowerCase()) + .where((pubkey) => pubkey.isNotEmpty) + .toSet() + .toList(); + if (normalized.isEmpty) return true; + try { + final session = ref.read(relaySessionProvider.notifier); + final events = await session.fetchHistory( + NostrFilters.profilesBatch(normalized), + ); + final updated = Map.from(state); + final updatedOrders = Map.from( + _profileEventOrders, + ); + for (final event in events) { + _cacheProfileEvent(event, updated, updatedOrders); + } + _profileEventOrders + ..clear() + ..addAll(updatedOrders); + state = updated; + return true; + } catch (_) { + return false; + } } /// Applies a live kind:0 profile event to the cache. @@ -55,13 +98,14 @@ class UserCacheNotifier extends Notifier> { /// to update names and avatars without discarding the rest of the cache. void cacheProfileEvent(NostrEvent event) { if (event.kind != 0) return; - final profile = _profileFromEvent(event); - state = {...state, profile.pubkey: profile}; + final updated = Map.from(state); + if (_cacheProfileEvent(event, updated)) state = updated; } void _scheduleFetch(String pubkey) { if (state.containsKey(pubkey) || _pending.contains(pubkey)) return; _pending.add(pubkey); + _batchCompleter ??= Completer(); _batchTimer ??= Timer(const Duration(milliseconds: 50), _flushPending); } @@ -71,7 +115,10 @@ class UserCacheNotifier extends Notifier> { final pubkeys = _pending.toList(); _pending.clear(); + final completer = _batchCompleter; + _batchCompleter = null; + var succeeded = false; try { final session = ref.read(relaySessionProvider.notifier); final events = await session.fetchHistory( @@ -79,17 +126,46 @@ class UserCacheNotifier extends Notifier> { ); final updated = Map.from(state); + final updatedOrders = Map.from( + _profileEventOrders, + ); for (final event in events) { - final profile = _profileFromEvent(event); - updated[profile.pubkey] = profile; + _cacheProfileEvent(event, updated, updatedOrders); } + _profileEventOrders + ..clear() + ..addAll(updatedOrders); state = updated; + succeeded = true; } catch (_) { - // Silently fail — we'll just show pubkeys. + // Silently fail — non-gating callers will just show pubkeys. + } finally { + completer?.complete(succeeded); } } + bool _cacheProfileEvent( + NostrEvent event, + Map profiles, [ + Map? orders, + ]) { + if (event.kind != 0) return false; + final eventOrders = orders ?? _profileEventOrders; + final pubkey = event.pubkey.toLowerCase(); + final current = eventOrders[pubkey]; + final isNewer = + current == null || + event.createdAt > current.createdAt || + (event.createdAt == current.createdAt && + event.id.compareTo(current.eventId) < 0); + if (!isNewer) return false; + + profiles[pubkey] = _profileFromEvent(event); + eventOrders[pubkey] = (createdAt: event.createdAt, eventId: event.id); + return true; + } + UserProfile _profileFromEvent(NostrEvent event) { final data = ProfileData.fromEvent(event); final pubkey = data.pubkey.toLowerCase(); diff --git a/mobile/lib/shared/relay/relay_session.dart b/mobile/lib/shared/relay/relay_session.dart index 77e10e66ede..b8d13ee34a4 100644 --- a/mobile/lib/shared/relay/relay_session.dart +++ b/mobile/lib/shared/relay/relay_session.dart @@ -17,17 +17,12 @@ import 'relay_closed_policy.dart'; import 'relay_http_query_client.dart'; import 'relay_provider.dart'; import 'relay_rate_limit_gate.dart'; +import 'relay_session_types.dart'; import 'relay_socket.dart'; -enum SessionStatus { disconnected, connecting, connected, reconnecting } +export 'relay_session_types.dart'; -@immutable -class SessionState { - final SessionStatus status; - final int reconnectAttempt; - - const SessionState({required this.status, this.reconnectAttempt = 0}); -} +part 'relay_session_auth.dart'; class _HistorySubscription { final List events = []; @@ -41,6 +36,7 @@ class _LiveSubscription { final NostrFilter filter; final void Function(NostrEvent) onEvent; final void Function(String message)? onClosed; + final void Function(RelaySubscriptionStatus status)? onStatusChanged; Completer? readyCompleter; int? lastSeenCreatedAt; int closedRetryAttempt = 0; @@ -50,6 +46,7 @@ class _LiveSubscription { required this.filter, required this.onEvent, this.onClosed, + this.onStatusChanged, this.readyCompleter, }); } @@ -57,7 +54,6 @@ class _LiveSubscription { class _ClosedRetry { final _LiveSubscription subscription; final int generation; - _ClosedRetry({required this.subscription, required this.generation}); } @@ -75,16 +71,6 @@ class _BufferedEvent { _BufferedEvent(this.subId, this.event); } -/// Manages websocket subscriptions, batching, reconnection, and pending events. -typedef RelaySocketFactory = - RelaySocket Function({ - required String wsUrl, - required String? nsec, - required void Function(List message) onMessage, - required void Function() onConnected, - required void Function(Object? error) onDisconnected, - }); - class RelaySessionNotifier extends Notifier { RelaySessionNotifier({ http.Client? httpClient, @@ -257,13 +243,38 @@ class RelaySessionNotifier extends Notifier { return completer.future; } - /// Subscribe to live events matching [filter]. Returns an unsubscribe - /// function. Live subscriptions survive reconnects — they are replayed with - /// `since: lastSeenCreatedAt - 5s` on reconnect. Future subscribe( NostrFilter filter, void Function(NostrEvent) onEvent, { void Function(String message)? onClosed, + }) => _subscribe(filter, onEvent, onClosed: onClosed); + + /// Subscribe to a live stream and observe its recovery lifecycle. + /// + /// The returned future completes after initial EOSE (or the existing + /// fallback timeout) and yields a cleanup callback. [onStatusChanged] emits + /// [RelaySubscriptionStatus.ready] at each EOSE after buffered replay events + /// have been delivered, and [RelaySubscriptionStatus.retrying] immediately + /// when a retryable or rate-limited CLOSED begins backoff. Terminal CLOSED + /// invokes [onClosed] and removes the subscription instead of retrying it. + /// Calling the cleanup callback cancels pending retries and sends CLOSE. + Future subscribeWithStatus( + NostrFilter filter, + void Function(NostrEvent) onEvent, { + void Function(String message)? onClosed, + required void Function(RelaySubscriptionStatus status) onStatusChanged, + }) => _subscribe( + filter, + onEvent, + onClosed: onClosed, + onStatusChanged: onStatusChanged, + ); + + Future _subscribe( + NostrFilter filter, + void Function(NostrEvent) onEvent, { + void Function(String message)? onClosed, + void Function(RelaySubscriptionStatus status)? onStatusChanged, }) async { if (_disposed) throw StateError('Relay session is disposed'); final subId = _nextSubId('l'); @@ -273,12 +284,12 @@ class RelaySessionNotifier extends Notifier { filter: filter, onEvent: onEvent, onClosed: onClosed, + onStatusChanged: onStatusChanged, readyCompleter: readyCompleter, ); _sendReq(subId, filter); - // Wait for EOSE or a short fallback timeout. try { await readyCompleter.future.timeout( const Duration(milliseconds: 500), @@ -297,7 +308,6 @@ class RelaySessionNotifier extends Notifier { return () => _unsubscribe(subId); } - /// Publish an event and wait for the relay's OK confirmation. Future publish( NostrEvent event, { Duration timeout = const Duration(seconds: 8), @@ -324,8 +334,6 @@ class RelaySessionNotifier extends Notifier { return completer.future; } - /// Send a raw message over the WebSocket without waiting for acknowledgement. - /// Used for ephemeral events like typing indicators. void sendRaw(List payload) { _socket?.send(payload); } @@ -635,7 +643,8 @@ class RelaySessionNotifier extends Notifier { return; } - // Live subscriptions get batched. + // Live subscriptions get batched. An EVENT proves the stream is active, + // but not that a retry replay is complete; only EOSE is that boundary. final liveSub = _liveSubscriptions[subId]; if (liveSub != null) { _resetClosedRetry(liveSub); @@ -664,19 +673,18 @@ class RelaySessionNotifier extends Notifier { return; } - // Live subscription: signal ready. + // Live subscription: flush replay callbacks before signaling ready. This + // ordering matters for retry replays, whose original ready completer has + // already been released. final liveSub = _liveSubscriptions[subId]; if (liveSub != null) { _resetClosedRetry(liveSub); + _flushBufferedEventsNow(); + liveSub.onStatusChanged?.call(RelaySubscriptionStatus.ready); } if (liveSub != null && liveSub.readyCompleter != null && !liveSub.readyCompleter!.isCompleted) { - // EOSE is the boundary between replay and live delivery. Flush any - // replay events before resolving subscribe(), so callers that begin a - // one-shot query immediately afterwards cannot classify a delayed batch - // callback as having arrived during that query. - _flushBufferedEventsNow(); liveSub.readyCompleter!.complete(); liveSub.readyCompleter = null; } @@ -717,6 +725,7 @@ class RelaySessionNotifier extends Notifier { readyCompleter.complete(); liveSub.readyCompleter = null; } + liveSub.onStatusChanged?.call(RelaySubscriptionStatus.retrying); if (liveSub.closedRetryTimer != null) return; final attempt = liveSub.closedRetryAttempt; @@ -965,35 +974,3 @@ final relaySessionProvider = NotifierProvider( RelaySessionNotifier.new, ); - -String buildNip98AuthHeader({ - required String method, - required String url, - required List bodyBytes, - required String? nsec, -}) { - if (nsec == null || nsec.isEmpty) { - throw Exception('Cannot query relay: no signing key available'); - } - final privkeyHex = nostr.Nip19.decode(payload: nsec).data; - if (privkeyHex.isEmpty) { - throw Exception('Invalid nsec'); - } - final payloadHash = SHA256Digest() - .process(Uint8List.fromList(bodyBytes)) - .map((byte) => byte.toRadixString(16).padLeft(2, '0')) - .join(); - final event = nostr.Event.from( - kind: 27235, - content: '', - tags: [ - ['u', url], - ['method', method.toUpperCase()], - ['payload', payloadHash], - ['nonce', const Uuid().v4()], - ], - secretKey: privkeyHex, - verify: false, - ); - return 'Nostr ${base64.encode(utf8.encode(event.toJson()))}'; -} diff --git a/mobile/lib/shared/relay/relay_session_auth.dart b/mobile/lib/shared/relay/relay_session_auth.dart new file mode 100644 index 00000000000..9c6c79ae02d --- /dev/null +++ b/mobile/lib/shared/relay/relay_session_auth.dart @@ -0,0 +1,33 @@ +part of 'relay_session.dart'; + +String buildNip98AuthHeader({ + required String method, + required String url, + required List bodyBytes, + required String? nsec, +}) { + if (nsec == null || nsec.isEmpty) { + throw Exception('Cannot query relay: no signing key available'); + } + final privkeyHex = nostr.Nip19.decode(payload: nsec).data; + if (privkeyHex.isEmpty) { + throw Exception('Invalid nsec'); + } + final payloadHash = SHA256Digest() + .process(Uint8List.fromList(bodyBytes)) + .map((byte) => byte.toRadixString(16).padLeft(2, '0')) + .join(); + final event = nostr.Event.from( + kind: 27235, + content: '', + tags: [ + ['u', url], + ['method', method.toUpperCase()], + ['payload', payloadHash], + ['nonce', const Uuid().v4()], + ], + secretKey: privkeyHex, + verify: false, + ); + return 'Nostr ${base64.encode(utf8.encode(event.toJson()))}'; +} diff --git a/mobile/lib/shared/relay/relay_session_types.dart b/mobile/lib/shared/relay/relay_session_types.dart new file mode 100644 index 00000000000..fe3ac8c5821 --- /dev/null +++ b/mobile/lib/shared/relay/relay_session_types.dart @@ -0,0 +1,25 @@ +import 'package:flutter/foundation.dart'; + +import 'relay_socket.dart'; + +enum SessionStatus { disconnected, connecting, connected, reconnecting } + +typedef RelaySocketFactory = + RelaySocket Function({ + required String wsUrl, + required String? nsec, + required void Function(List message) onMessage, + required void Function() onConnected, + required void Function(Object? error) onDisconnected, + }); + +@immutable +class SessionState { + final SessionStatus status; + final int reconnectAttempt; + + const SessionState({required this.status, this.reconnectAttempt = 0}); +} + +/// Recovery lifecycle for a live relay subscription. +enum RelaySubscriptionStatus { ready, retrying } diff --git a/mobile/test/features/channels/channel_detail_page_test.dart b/mobile/test/features/channels/channel_detail_page_test.dart index 195dcc25ad8..3aabad062ec 100644 --- a/mobile/test/features/channels/channel_detail_page_test.dart +++ b/mobile/test/features/channels/channel_detail_page_test.dart @@ -13,6 +13,7 @@ import 'package:http/http.dart' as http; import 'package:http/testing.dart' as http_testing; import 'package:lucide_icons_flutter/lucide_icons.dart'; import 'package:nostr/nostr.dart' as nostr; +import 'package:pointycastle/digests/sha256.dart'; import 'package:scrollable_positioned_list/scrollable_positioned_list.dart'; import 'package:buzz/features/channels/channel.dart'; import 'package:buzz/features/channels/channel_detail_page.dart'; @@ -200,7 +201,12 @@ Widget _buildTestable({ required List messages, List typing = const [], Map users = const {}, - _FakeUserCacheNotifier? userCacheNotifier, + Set? knownAgentPubkeys, + Future> Function()? loadChannelBotPubkeys, + bool watchChannelMembershipUpdates = false, + Future> Function()? loadAgentDirectory, + Future> Function()? loadAgentOwners, + UserCacheNotifier? userCacheNotifier, List members = const [], List huddleMembers = const [], _MutableHuddleMembersNotifier? huddleMembersNotifier, @@ -280,16 +286,24 @@ Widget _buildTestable({ ), if (huddleMembersNotifier != null) _mutableHuddleMembersProvider.overrideWith(() => huddleMembersNotifier), - channelBotPubkeysProvider( - _channelId, - ).overrideWith((ref) async => const {}), + if (!watchChannelMembershipUpdates) + channelBotPubkeysProvider(_channelId).overrideWith( + (ref) async => loadChannelBotPubkeys?.call() ?? const {}, + ), channelBotPubkeysProvider(_huddleChannelId).overrideWith( (ref) async => { for (final member in huddleMembers) if (member.isBot) member.pubkey.toLowerCase(), }, ), - agentOwnersProvider.overrideWith((ref) async => const {}), + agentOwnersProvider.overrideWith( + (ref) async => loadAgentOwners?.call() ?? const {}, + ), + agentDirectoryProvider.overrideWith( + (ref) async => loadAgentDirectory?.call() ?? const [], + ), + if (knownAgentPubkeys != null) + knownAgentPubkeysProvider.overrideWithValue(knownAgentPubkeys), if (directoryUsers != null) relayDirectoryUsersProvider.overrideWith((ref) async => directoryUsers), if (createChannelActions != null) @@ -327,8 +341,12 @@ Widget _buildTestable({ ), mediaHttpClientProvider.overrideWithValue(mediaClient), ], - if (relaySessionNotifier != null) - relaySessionProvider.overrideWith(() => relaySessionNotifier), + if (relaySessionNotifier != null || + (resolvedChannel.isDm && + resolvedChannel.participantPubkeys.toSet().length == 2)) + relaySessionProvider.overrideWith( + () => relaySessionNotifier ?? _IdentityUpdateRelaySession(), + ), if (relayConfigNotifier != null) relayConfigProvider.overrideWith(() => relayConfigNotifier), if (huddleMediaFactory != null) @@ -420,45 +438,847 @@ Widget _buildNavigationTestable({ ); } -/// Finder that searches for text within RichText spans. [find.text] only -/// matches the top-level text property; this also searches nested TextSpans. -Finder findRichText(String text) { - return find.byWidgetPredicate((widget) { - if (widget is RichText) { - return widget.text.toPlainText().contains(text); - } - return false; - }, description: 'RichText containing "$text"'); -} +/// Finder that searches for text within RichText spans. [find.text] only +/// matches the top-level text property; this also searches nested TextSpans. +Finder findRichText(String text) { + return find.byWidgetPredicate((widget) { + if (widget is RichText) { + return widget.text.toPlainText().contains(text); + } + return false; + }, description: 'RichText containing "$text"'); +} + +double? effectiveFontSizeForText( + InlineSpan span, + String text, [ + TextStyle? inheritedStyle, +]) { + if (span is! TextSpan) return null; + final effectiveStyle = inheritedStyle?.merge(span.style) ?? span.style; + if ((span.text ?? '').contains(text)) return effectiveStyle?.fontSize; + for (final child in span.children ?? const []) { + final size = effectiveFontSizeForText(child, text, effectiveStyle); + if (size != null) return size; + } + return null; +} + +void main() { + setUp(() async { + SharedPreferences.setMockInitialValues({}); + _testPrefs = await SharedPreferences.getInstance(); + }); + + group('ChannelDetailPage', () { + testWidgets('uses the shared 32px masked presence avatar in DM headers', ( + tester, + ) async { + 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: const [], + channel: dmChannel, + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + }, + ), + ); + await tester.pumpAndSettle(); + + final avatarFinder = find.byKey(const ValueKey('dm-header-avatar')); + final avatar = tester.widget(avatarFinder); + expect(tester.getSize(avatarFinder), const Size.square(32)); + expect(avatar.geometry, AvatarBadgeMaskGeometry.presenceDot); + expect(avatar.badge, isNotNull); + expect( + find.descendant(of: avatarFinder, matching: find.byType(ClipPath)), + findsOneWidget, + ); + final name = tester.widget( + find.byKey(const ValueKey('dm-header-name')), + ); + final presence = tester.widget( + find.byKey(const ValueKey('dm-header-presence')), + ); + expect(name.style?.fontSize, 16); + expect(name.style?.fontWeight, FontWeight.w500); + expect(presence.style?.fontSize, 14); + expect(presence.style?.fontWeight, FontWeight.w400); + expect(find.byTooltip('View members'), findsNothing); + expect(find.byTooltip('Start Huddle'), findsOneWidget); + }); + + testWidgets('hides the Huddle action in a one-to-one agent DM', ( + tester, + ) async { + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message with an agent', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + users: const { + 'agent': UserProfile( + pubkey: 'agent', + displayName: 'Agent', + ownerPubkey: 'owner', + ), + }, + ), + ); + await tester.pumpAndSettle(); + + expect(find.byKey(const ValueKey('channel-huddle-button')), findsNothing); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('hides the Huddle action for a channel bot DM', (tester) async { + final dmChannel = Channel( + id: _channelId, + name: 'Bot DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message with a channel bot', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Bot'], + participantPubkeys: const ['self', 'bot'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + loadChannelBotPubkeys: () async => const {'bot'}, + users: const {'bot': UserProfile(pubkey: 'bot', displayName: 'Bot')}, + ), + ); + await tester.pumpAndSettle(); + + expect(find.byKey(const ValueKey('channel-huddle-button')), findsNothing); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps the Huddle action hidden while agent identity loads', ( + tester, + ) async { + final directoryCompleter = Completer>(); + 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: const [], + channel: dmChannel, + loadAgentDirectory: () => directoryCompleter.future, + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + }, + ), + ); + await tester.pump(); + + expect(find.byKey(const ValueKey('channel-huddle-button')), findsNothing); + expect(find.byTooltip('Start Huddle'), findsNothing); + + directoryCompleter.complete(const []); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsOneWidget); + }); + + testWidgets('preloads DM participant profiles without a member snapshot', ( + tester, + ) async { + final preloadedPubkeys = []; + final userCache = _FakeUserCacheNotifier( + const {}, + preload: (pubkeys) async { + preloadedPubkeys.addAll(pubkeys); + return true; + }, + ); + final dmChannel = Channel( + id: _channelId, + name: 'Human 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: const [], + channel: dmChannel, + userCacheNotifier: userCache, + ), + ); + await tester.pumpAndSettle(); + + expect(preloadedPubkeys, containsAll(const ['self', 'alice'])); + expect(find.byTooltip('Start Huddle'), findsOneWidget); + }); + + testWidgets('force-refreshes cached profiles before enabling Huddle', ( + tester, + ) async { + late final _FakeUserCacheNotifier userCache; + userCache = _FakeUserCacheNotifier( + const { + 'agent': UserProfile(pubkey: 'agent', displayName: 'Cached Human'), + }, + preload: (_) async { + await Future.delayed(Duration.zero); + userCache.replace( + const UserProfile( + pubkey: 'agent', + displayName: 'Agent', + ownerPubkey: 'owner', + ), + ); + return true; + }, + ); + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + userCacheNotifier: userCache, + ), + ); + await tester.pumpAndSettle(); + + expect(userCache.state['agent']?.ownerPubkey, 'owner'); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps Huddle hidden when live owner profile beats refresh', ( + tester, + ) async { + final owner = nostr.Keys.generate(); + final agent = nostr.Keys.generate(); + final profileRefresh = Completer>(); + final relaySession = _IdentityUpdateRelaySession( + profileRefresh: profileRefresh.future, + ); + final userCache = UserCacheNotifier(); + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: ['self', agent.public], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + userCacheNotifier: userCache, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: agent.public, + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pump(); + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.emitProfile( + _profileEvent( + id: 'newer-agent', + pubkey: agent.public, + createdAt: 2, + name: 'Agent', + tags: [_authTag(owner, agent.public)], + ), + ); + profileRefresh.complete([ + _profileEvent( + id: 'older-human', + pubkey: agent.public, + createdAt: 1, + name: 'Human', + ), + ]); + await tester.pumpAndSettle(); + + expect(userCache.state[agent.public]?.ownerPubkey, owner.public); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps Huddle hidden while a verified owner profile loads', ( + tester, + ) async { + final profilePreloadCompleter = Completer(); + final userCache = _FakeUserCacheNotifier( + const {}, + preload: (_) => profilePreloadCompleter.future, + ); + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + userCacheNotifier: userCache, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'agent', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pump(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + + userCache.replace( + const UserProfile( + pubkey: 'agent', + displayName: 'Agent', + ownerPubkey: 'owner', + ), + ); + profilePreloadCompleter.complete(true); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('rechecks verified owner profiles after reconnect', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final reconnectPreloadCompleter = Completer(); + var memberPreloadCount = 0; + var blockMemberPreload = false; + final userCache = _FakeUserCacheNotifier( + const {}, + preload: (pubkeys) { + if (pubkeys.length == 1) return Future.value(true); + memberPreloadCount++; + return blockMemberPreload + ? reconnectPreloadCompleter.future + : Future.value(true); + }, + ); + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + userCacheNotifier: userCache, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'agent', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsOneWidget); + final memberPreloadsBeforeReconnect = memberPreloadCount; + blockMemberPreload = true; + + relaySession.disconnect(); + await tester.pump(); + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.connect(); + await tester.pump(); + expect(memberPreloadCount, greaterThan(memberPreloadsBeforeReconnect)); + expect(find.byTooltip('Start Huddle'), findsNothing); + + userCache.replace( + const UserProfile( + pubkey: 'agent', + displayName: 'Agent', + ownerPubkey: 'owner', + ), + ); + reconnectPreloadCompleter.complete(true); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + await tester.pump(const Duration(milliseconds: 500)); + }); + + testWidgets('keeps directory-only agent Huddle hidden after disconnect', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + loadAgentDirectory: () async => const [ + AgentDirectoryEntry(pubkey: 'agent'), + ], + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'agent', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.disconnect(); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps bot-role-only Huddle hidden after disconnect', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Bot DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Bot'], + participantPubkeys: const ['self', 'bot'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + loadChannelBotPubkeys: () async => const {'bot'}, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'bot', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.disconnect(); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps Huddle hidden until bot-role replay reaches EOSE', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Human 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: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + watchChannelMembershipUpdates: true, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsOneWidget); + + relaySession.beginMembershipReplay(); + await tester.pump(); + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.emitReplayedMembership( + NostrEvent( + id: 'membership-self', + pubkey: 'relay', + createdAt: 1, + kind: 39002, + tags: const [ + ['d', _channelId], + ['p', 'self'], + ], + content: '', + sig: 'sig', + ), + ); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.emitReplayedMembership( + NostrEvent( + id: 'membership-bot', + pubkey: 'relay', + createdAt: 2, + kind: 39002, + tags: const [ + ['d', _channelId], + ['p', 'self'], + ['p', 'alice', '', 'bot'], + ], + content: '', + sig: 'sig', + ), + ); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsNothing); + + relaySession.finishMembershipReplay(); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps the Huddle action hidden when identity loading fails', ( + tester, + ) async { + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + loadAgentOwners: () => Future.error('identity unavailable'), + disableRetries: true, + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps the Huddle action hidden when member preload fails', ( + tester, + ) async { + final dmChannel = Channel( + id: _channelId, + name: 'Agent DM', + channelType: 'dm', + visibility: 'private', + description: 'Direct message', + createdBy: 'self', + createdAt: DateTime(2025), + memberCount: 2, + participants: const ['Self', 'Agent'], + participantPubkeys: const ['self', 'agent'], + isMember: true, + ); + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + loadMembers: () => Future.error('members unavailable'), + disableRetries: true, + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('hides Huddle when a participant becomes an agent live', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Human 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, + ); + + var directoryLoadCount = 0; + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + loadAgentDirectory: () async { + directoryLoadCount++; + return directoryLoadCount == 1 + ? const [] + : const [AgentDirectoryEntry(pubkey: 'alice')]; + }, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + + expect(relaySession.identityFilter?.kinds, const [0, 10100]); + expect(relaySession.identityFilter?.authors, contains('alice')); + expect(relaySession.identityFilter?.limit, 100); + expect(find.byTooltip('Start Huddle'), findsOneWidget); + + relaySession.emitAgentProfile(pubkey: 'alice'); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('keeps Huddle hidden while identity replay retries', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Human 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: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsOneWidget); -double? effectiveFontSizeForText( - InlineSpan span, - String text, [ - TextStyle? inheritedStyle, -]) { - if (span is! TextSpan) return null; - final effectiveStyle = inheritedStyle?.merge(span.style) ?? span.style; - if ((span.text ?? '').contains(text)) return effectiveStyle?.fontSize; - for (final child in span.children ?? const []) { - final size = effectiveFontSizeForText(child, text, effectiveStyle); - if (size != null) return size; - } - return null; -} + relaySession.retryIdentitySubscription(); + await tester.pump(); + expect(find.byTooltip('Start Huddle'), findsNothing); -void main() { - setUp(() async { - SharedPreferences.setMockInitialValues({}); - _testPrefs = await SharedPreferences.getInstance(); - }); + relaySession.emitAgentProfile(pubkey: 'alice'); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsNothing); - group('ChannelDetailPage', () { - testWidgets('uses the shared 32px masked presence avatar in DM headers', ( + relaySession.readyIdentitySubscription(); + await tester.pumpAndSettle(); + expect(find.byTooltip('Start Huddle'), findsNothing); + }); + + testWidgets('queries DM participants directly for agent identity', ( tester, ) async { + final relaySession = _IdentityUpdateRelaySession(); final dmChannel = Channel( id: _channelId, - name: 'DM', + name: 'Human DM', channelType: 'dm', visibility: 'private', description: 'Direct message', @@ -474,35 +1294,132 @@ void main() { _buildTestable( messages: const [], channel: dmChannel, - users: const { - 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), - }, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], ), ); await tester.pumpAndSettle(); - final avatarFinder = find.byKey(const ValueKey('dm-header-avatar')); - final avatar = tester.widget(avatarFinder); - expect(tester.getSize(avatarFinder), const Size.square(32)); - expect(avatar.geometry, AvatarBadgeMaskGeometry.presenceDot); - expect(avatar.badge, isNotNull); + expect(relaySession.directIdentityFilter?.kinds, const [10100]); expect( - find.descendant(of: avatarFinder, matching: find.byType(ClipPath)), - findsOneWidget, + relaySession.directIdentityFilter?.authors, + containsAll(const ['self', 'alice']), ); - final name = tester.widget( - find.byKey(const ValueKey('dm-header-name')), + expect(relaySession.directIdentityFilter?.limit, 2); + }); + + testWidgets('hides Huddle for an agent found by direct DM lookup', ( + tester, + ) async { + final relaySession = _IdentityUpdateRelaySession() + ..directIdentityProfiles = const [ + NostrEvent( + id: 'old-agent-profile', + pubkey: 'alice', + createdAt: 1, + kind: 10100, + tags: [], + content: '{"name":"Agent"}', + sig: 'sig', + ), + ]; + final dmChannel = Channel( + id: _channelId, + name: 'Agent 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, ); - final presence = tester.widget( - find.byKey(const ValueKey('dm-header-presence')), + + await tester.pumpWidget( + _buildTestable( + messages: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), ); - expect(name.style?.fontSize, 16); - expect(name.style?.fontWeight, FontWeight.w500); - expect(presence.style?.fontSize, 14); - expect(presence.style?.fontWeight, FontWeight.w400); - expect(find.byTooltip('View members'), findsNothing); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsNothing); }); + testWidgets( + 'keeps Huddle hidden if the live identity subscription closes', + (tester) async { + final relaySession = _IdentityUpdateRelaySession(); + final dmChannel = Channel( + id: _channelId, + name: 'Human 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: const [], + channel: dmChannel, + relaySessionNotifier: relaySession, + members: [ + ChannelMember( + pubkey: 'self', + role: 'member', + joinedAt: DateTime(2025), + ), + ChannelMember( + pubkey: 'alice', + role: 'member', + joinedAt: DateTime(2025), + ), + ], + ), + ); + await tester.pumpAndSettle(); + + expect(find.byTooltip('Start Huddle'), findsOneWidget); + + relaySession.closeIdentitySubscription(); + await tester.pump(); + + expect(find.byTooltip('Start Huddle'), findsNothing); + }, + ); + testWidgets('keeps the Members action for group DMs', (tester) async { final dmChannel = Channel( id: _channelId, @@ -522,6 +1439,7 @@ void main() { _buildTestable( messages: const [], channel: dmChannel, + knownAgentPubkeys: const {'alice'}, users: const { 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), 'bob': UserProfile(pubkey: 'bob', displayName: 'Bob'), @@ -531,6 +1449,7 @@ void main() { await tester.pumpAndSettle(); expect(find.byTooltip('View members'), findsOneWidget); + expect(find.byTooltip('Start Huddle'), findsOneWidget); }); testWidgets( @@ -12295,6 +13214,142 @@ class _ReconnectingRelaySession extends RelaySessionNotifier { } } +class _IdentityUpdateRelaySession extends RelaySessionNotifier { + _IdentityUpdateRelaySession({this.profileRefresh}); + + final Future>? profileRefresh; + NostrFilter? identityFilter; + NostrFilter? directIdentityFilter; + List directIdentityProfiles = const []; + void Function(NostrEvent)? _identityListener; + void Function(String message)? _identityClosedListener; + void Function(RelaySubscriptionStatus status)? _identityStatusListener; + void Function(NostrEvent)? _membershipListener; + void Function(RelaySubscriptionStatus status)? _membershipStatusListener; + NostrEvent? membershipSnapshot; + + @override + SessionState build() => const SessionState(status: SessionStatus.connected); + + @override + Future> fetchHistory( + NostrFilter filter, { + Duration timeout = const Duration(seconds: 8), + }) async { + if (filter.kinds.length == 1 && filter.kinds.single == 0) { + return profileRefresh ?? const []; + } + if (filter.kinds.contains(10100) && filter.kinds.length == 1) { + directIdentityFilter = filter; + return directIdentityProfiles; + } + return membershipSnapshot == null ? const [] : [membershipSnapshot!]; + } + + @override + Future subscribeWithStatus( + NostrFilter filter, + void Function(NostrEvent) onEvent, { + void Function(String message)? onClosed, + required void Function(RelaySubscriptionStatus status) onStatusChanged, + }) async { + if (filter.kinds.contains(10100)) { + identityFilter = filter; + _identityListener = onEvent; + _identityClosedListener = onClosed; + _identityStatusListener = onStatusChanged; + onStatusChanged(RelaySubscriptionStatus.ready); + return () { + if (identical(_identityListener, onEvent)) { + _identityListener = null; + _identityClosedListener = null; + _identityStatusListener = null; + } + }; + } + _membershipListener = onEvent; + _membershipStatusListener = onStatusChanged; + onStatusChanged(RelaySubscriptionStatus.ready); + return () { + if (identical(_membershipListener, onEvent)) { + _membershipListener = null; + _membershipStatusListener = null; + } + }; + } + + @override + Future subscribe( + NostrFilter filter, + void Function(NostrEvent) onEvent, { + void Function(String message)? onClosed, + }) async { + if (filter.kinds.contains(10100)) { + identityFilter = filter; + _identityListener = onEvent; + _identityClosedListener = onClosed; + } + return () { + if (identical(_identityListener, onEvent)) { + _identityListener = null; + _identityClosedListener = null; + } + }; + } + + void emitProfile(NostrEvent event) { + _identityListener?.call(event); + } + + void emitAgentProfile({required String pubkey}) { + _identityListener?.call( + NostrEvent( + id: 'agent-profile-$pubkey', + pubkey: pubkey, + createdAt: DateTime.now().millisecondsSinceEpoch ~/ 1000, + kind: 10100, + tags: const [], + content: '{"name":"Agent"}', + sig: 'sig', + ), + ); + } + + void closeIdentitySubscription() { + _identityClosedListener?.call('unsupported filter'); + } + + void retryIdentitySubscription() { + _identityStatusListener?.call(RelaySubscriptionStatus.retrying); + } + + void readyIdentitySubscription() { + _identityStatusListener?.call(RelaySubscriptionStatus.ready); + } + + void beginMembershipReplay() { + _membershipStatusListener?.call(RelaySubscriptionStatus.retrying); + } + + void emitReplayedMembership(NostrEvent event) { + membershipSnapshot = event; + _membershipListener?.call(event); + _membershipStatusListener?.call(RelaySubscriptionStatus.retrying); + } + + void finishMembershipReplay() { + _membershipStatusListener?.call(RelaySubscriptionStatus.ready); + } + + void disconnect() { + state = const SessionState(status: SessionStatus.disconnected); + } + + void connect() { + state = const SessionState(status: SessionStatus.connected); + } +} + class _HuddleReactionRelaySession extends RelaySessionNotifier { NostrFilter? reactionFilter; void Function(NostrEvent)? _reactionListener; @@ -12437,9 +13492,46 @@ class _FakeChannelMutesNotifier extends ChannelMutesNotifier { } } +NostrEvent _profileEvent({ + required String id, + required String pubkey, + required int createdAt, + required String name, + List> tags = const [], +}) => NostrEvent( + id: id, + pubkey: pubkey, + createdAt: createdAt, + kind: 0, + tags: tags, + content: jsonEncode({'name': name}), + sig: 'sig', +); + +List _authTag(nostr.Keys owner, String agentPubkey) { + final digest = SHA256Digest().process( + Uint8List.fromList( + utf8.encode('nostr:agent-auth:${agentPubkey.toLowerCase()}:'), + ), + ); + final message = digest + .map((byte) => byte.toRadixString(16).padLeft(2, '0')) + .join(); + return [ + 'auth', + owner.public, + '', + nostr.Schnorr.sign(secretKey: owner.secret, message: message), + ]; +} + class _FakeUserCacheNotifier extends UserCacheNotifier { final Map _users; - _FakeUserCacheNotifier(this._users); + final Future Function(List)? _preload; + _FakeUserCacheNotifier( + this._users, { + Future Function(List)? preload, + }) : _preload = preload; @override Map build() => _users; @@ -12447,6 +13539,13 @@ class _FakeUserCacheNotifier extends UserCacheNotifier { @override UserProfile? get(String pubkey) => _users[pubkey.toLowerCase()]; + @override + Future preload(List pubkeys) => + _preload?.call(pubkeys) ?? Future.value(true); + + @override + Future refresh(List pubkeys) => preload(pubkeys); + void replace(UserProfile profile) { state = {...state, profile.pubkey.toLowerCase(): profile}; } diff --git a/mobile/test/features/channels/reaction_row_test.dart b/mobile/test/features/channels/reaction_row_test.dart index 44bb450670b..1f3ef745ef1 100644 --- a/mobile/test/features/channels/reaction_row_test.dart +++ b/mobile/test/features/channels/reaction_row_test.dart @@ -90,7 +90,7 @@ class _FakeUserCacheNotifier extends UserCacheNotifier { Map build() => _profiles; @override - Future preload(Iterable pubkeys) async {} + Future preload(List pubkeys) async => true; } void main() { diff --git a/mobile/test/shared/mentions/agent_identity_provider_test.dart b/mobile/test/shared/mentions/agent_identity_provider_test.dart index f584aff0419..c5355f26a43 100644 --- a/mobile/test/shared/mentions/agent_identity_provider_test.dart +++ b/mobile/test/shared/mentions/agent_identity_provider_test.dart @@ -29,8 +29,8 @@ void main() { }); await relaySession.subscribed; expect(relaySession.liveFilters.single.kinds, const [39002]); - expect(relaySession.liveFilters.single.tags['#h'], [_channelId]); - expect(relaySession.liveFilters.single.tags['#d'], isNull); + expect(relaySession.liveFilters.single.tags['#d'], [_channelId]); + expect(relaySession.liveFilters.single.tags['#h'], isNull); relaySession.emit(_membershipEvent(role: 'member')); await _pumpEventQueue(); @@ -76,6 +76,96 @@ void main() { ); }); + test('surfaces bot-role subscription setup failure', () async { + final relaySession = _MembershipRelaySessionNotifier([ + _membershipEvent(role: 'member'), + ], subscribeError: StateError('subscription unavailable')); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => relaySession)], + ); + addTearDown(container.dispose); + final keepAlive = container.listen( + channelMembershipUpdateProvider(_channelId), + (_, _) {}, + fireImmediately: true, + ); + addTearDown(keepAlive.close); + + await _pumpEventQueue(); + + final state = container.read(channelMembershipUpdateProvider(_channelId)); + expect(state.isReady, isFalse); + expect(state.error, isA()); + }); + + test('surfaces terminal bot-role subscription closure', () async { + final relaySession = _MembershipRelaySessionNotifier([ + _membershipEvent(role: 'member'), + ]); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => relaySession)], + ); + addTearDown(container.dispose); + final keepAlive = container.listen( + channelMembershipUpdateProvider(_channelId), + (_, _) {}, + fireImmediately: true, + ); + addTearDown(keepAlive.close); + + await relaySession.subscribed; + await _pumpEventQueue(); + expect( + container.read(channelMembershipUpdateProvider(_channelId)).isReady, + isTrue, + ); + + relaySession.closeSubscription('unsupported filter'); + await _pumpEventQueue(); + + final state = container.read(channelMembershipUpdateProvider(_channelId)); + expect(state.isReady, isFalse); + expect(state.error, isA()); + }); + + test('fails closed while bot-role subscription retries', () async { + final relaySession = _MembershipRelaySessionNotifier([ + _membershipEvent(role: 'member'), + ]); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => relaySession)], + ); + addTearDown(container.dispose); + final keepAlive = container.listen( + channelMembershipUpdateProvider(_channelId), + (_, _) {}, + fireImmediately: true, + ); + addTearDown(keepAlive.close); + + await relaySession.subscribed; + await _pumpEventQueue(); + expect( + container.read(channelMembershipUpdateProvider(_channelId)).isReady, + isTrue, + ); + + relaySession.setSubscriptionStatus(RelaySubscriptionStatus.retrying); + await _pumpEventQueue(); + expect( + container.read(channelMembershipUpdateProvider(_channelId)).isReady, + isFalse, + ); + + relaySession.setSubscriptionStatus(RelaySubscriptionStatus.ready); + await _pumpEventQueue(); + final recovered = container.read( + channelMembershipUpdateProvider(_channelId), + ); + expect(recovered.isReady, isTrue); + expect(recovered.error, isNull); + }); + test('disposes the live role subscription without consumers', () async { final relaySession = _MembershipRelaySessionNotifier([ _membershipEvent(role: 'bot'), @@ -162,13 +252,14 @@ Future _pumpEventQueue() async { class _MembershipRelaySessionNotifier extends RelaySessionNotifier { final List _memberships; + final Object? subscribeError; final List liveFilters = []; final List<_LiveSubscription> _subscriptions = []; final Completer _subscribed = Completer(); var unsubscribeCount = 0; var _membershipIndex = 0; - _MembershipRelaySessionNotifier(this._memberships); + _MembershipRelaySessionNotifier(this._memberships, {this.subscribeError}); Future get subscribed => _subscribed.future; @@ -184,14 +275,22 @@ class _MembershipRelaySessionNotifier extends RelaySessionNotifier { } @override - Future subscribe( + Future subscribeWithStatus( NostrFilter filter, void Function(NostrEvent) onEvent, { void Function(String message)? onClosed, + required void Function(RelaySubscriptionStatus status) onStatusChanged, }) async { + if (subscribeError case final error?) throw error; liveFilters.add(filter); - final subscription = _LiveSubscription(filter, onEvent); + final subscription = _LiveSubscription( + filter, + onEvent, + onClosed, + onStatusChanged, + ); _subscriptions.add(subscription); + onStatusChanged(RelaySubscriptionStatus.ready); if (!_subscribed.isCompleted) _subscribed.complete(); return () { unsubscribeCount++; @@ -206,13 +305,32 @@ class _MembershipRelaySessionNotifier extends RelaySessionNotifier { } } } + + void closeSubscription(String message) { + for (final subscription in List.of(_subscriptions)) { + subscription.onClosed?.call(message); + } + } + + void setSubscriptionStatus(RelaySubscriptionStatus status) { + for (final subscription in List.of(_subscriptions)) { + subscription.onStatusChanged?.call(status); + } + } } class _LiveSubscription { final NostrFilter filter; final void Function(NostrEvent) onEvent; + final void Function(String message)? onClosed; + final void Function(RelaySubscriptionStatus status)? onStatusChanged; - const _LiveSubscription(this.filter, this.onEvent); + const _LiveSubscription( + this.filter, + this.onEvent, + this.onClosed, + this.onStatusChanged, + ); } bool _matches(NostrFilter filter, NostrEvent event) { diff --git a/mobile/test/shared/profile/user_cache_provider_test.dart b/mobile/test/shared/profile/user_cache_provider_test.dart new file mode 100644 index 00000000000..9a270861076 --- /dev/null +++ b/mobile/test/shared/profile/user_cache_provider_test.dart @@ -0,0 +1,226 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:typed_data'; + +import 'package:buzz/shared/profile/user_cache_provider.dart'; +import 'package:buzz/shared/relay/relay.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:hooks_riverpod/hooks_riverpod.dart'; +import 'package:nostr/nostr.dart' as nostr; +import 'package:pointycastle/digests/sha256.dart'; + +void main() { + test('preload reports a profile batch failure', () async { + final container = ProviderContainer( + overrides: [ + relaySessionProvider.overrideWith(_FailingProfileSession.new), + ], + ); + addTearDown(container.dispose); + + final succeeded = await container.read(userCacheProvider.notifier).preload( + const ['agent'], + ); + + expect(succeeded, isFalse); + }); + + test('refresh queries profiles that are already cached', () async { + final session = _RecordingProfileSession(); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => session)], + ); + addTearDown(container.dispose); + final cache = container.read(userCacheProvider.notifier); + cache.cacheProfileEvent( + _profileEvent(id: 'cached-profile', createdAt: 1, name: 'Cached Human'), + ); + + final succeeded = await cache.refresh(const ['AGENT']); + + expect(succeeded, isTrue); + expect(session.requestedFilter?.kinds, const [0]); + expect(session.requestedFilter?.authors, const ['agent']); + expect(session.requestedFilter?.limit, 1); + }); + + test('older refresh cannot overwrite a newer live profile', () async { + final refreshCompleter = Completer>(); + final session = _RecordingProfileSession(result: refreshCompleter.future); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => session)], + ); + addTearDown(container.dispose); + final cache = container.read(userCacheProvider.notifier); + final owner = nostr.Keys.generate(); + final agent = nostr.Keys.generate(); + final refresh = cache.refresh([agent.public]); + + cache.cacheProfileEvent( + _profileEvent( + id: 'newer-agent', + pubkey: agent.public, + createdAt: 2, + name: 'Agent', + tags: [_authTag(owner, agent.public)], + ), + ); + refreshCompleter.complete([ + _profileEvent( + id: 'older-human', + pubkey: agent.public, + createdAt: 1, + name: 'Human', + ), + ]); + + expect(await refresh, isTrue); + expect(cache.state[agent.public]?.displayName, 'Agent'); + expect(cache.state[agent.public]?.ownerPubkey, owner.public); + }); + + test('newer refresh can remove obsolete owner attribution', () async { + final owner = nostr.Keys.generate(); + final agent = nostr.Keys.generate(); + final session = _RecordingProfileSession( + result: Future.value([ + _profileEvent( + id: 'newer-human', + pubkey: agent.public, + createdAt: 2, + name: 'Human', + ), + ]), + ); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => session)], + ); + addTearDown(container.dispose); + final cache = container.read(userCacheProvider.notifier); + cache.cacheProfileEvent( + _profileEvent( + id: 'older-agent', + pubkey: agent.public, + createdAt: 1, + name: 'Agent', + tags: [_authTag(owner, agent.public)], + ), + ); + + expect(await cache.refresh([agent.public]), isTrue); + expect(cache.state[agent.public]?.displayName, 'Human'); + expect(cache.state[agent.public]?.ownerPubkey, isNull); + }); + + test('non-profile history cannot poison profile order', () async { + final session = _RecordingProfileSession( + results: [ + Future.value([ + _profileEvent( + id: 'non-profile-newer', + createdAt: 3, + name: 'Ignored', + kind: 1, + ), + ]), + Future.value([ + _profileEvent(id: 'valid-older', createdAt: 2, name: 'Valid'), + ]), + ], + ); + final container = ProviderContainer( + overrides: [relaySessionProvider.overrideWith(() => session)], + ); + addTearDown(container.dispose); + final cache = container.read(userCacheProvider.notifier); + + expect(await cache.refresh(const ['agent']), isTrue); + expect(cache.state['agent'], isNull); + expect(await cache.refresh(const ['agent']), isTrue); + expect(cache.state['agent']?.displayName, 'Valid'); + }); + + test('same-second profile tie keeps the lowest event id', () { + final container = ProviderContainer(); + addTearDown(container.dispose); + final cache = container.read(userCacheProvider.notifier); + + cache.cacheProfileEvent( + _profileEvent(id: 'b', createdAt: 1, name: 'Larger ID'), + ); + cache.cacheProfileEvent( + _profileEvent(id: 'a', createdAt: 1, name: 'Lower ID'), + ); + cache.cacheProfileEvent( + _profileEvent(id: 'c', createdAt: 1, name: 'Later Larger ID'), + ); + + expect(cache.state['agent']?.displayName, 'Lower ID'); + }); +} + +NostrEvent _profileEvent({ + required String id, + required int createdAt, + required String name, + String pubkey = 'agent', + List> tags = const [], + int kind = 0, +}) => NostrEvent( + id: id, + pubkey: pubkey, + createdAt: createdAt, + kind: kind, + tags: tags, + content: jsonEncode({'name': name}), + sig: 'sig', +); + +List _authTag(nostr.Keys owner, String agentPubkey) { + final digest = SHA256Digest().process( + Uint8List.fromList( + utf8.encode('nostr:agent-auth:${agentPubkey.toLowerCase()}:'), + ), + ); + final message = digest + .map((byte) => byte.toRadixString(16).padLeft(2, '0')) + .join(); + final signature = nostr.Schnorr.sign( + secretKey: owner.secret, + message: message, + ); + return ['auth', owner.public, '', signature]; +} + +class _RecordingProfileSession extends RelaySessionNotifier { + _RecordingProfileSession({ + Future>? result, + List>>? results, + }) : _results = [...?results, ?result]; + + final List>> _results; + NostrFilter? requestedFilter; + + @override + SessionState build() => const SessionState(status: SessionStatus.connected); + + @override + Future> fetchHistory( + NostrFilter filter, { + Duration timeout = const Duration(seconds: 8), + }) async { + requestedFilter = filter; + return _results.isEmpty ? const [] : _results.removeAt(0); + } +} + +class _FailingProfileSession extends RelaySessionNotifier { + @override + SessionState build() => const SessionState(status: SessionStatus.connected); + + @override + Future> fetchHistory( + NostrFilter filter, { + Duration timeout = const Duration(seconds: 8), + }) => Future.error('profile unavailable'); +} diff --git a/mobile/test/shared/relay/relay_session_test.dart b/mobile/test/shared/relay/relay_session_test.dart index 7332075c3f6..c6ae6e9bd8a 100644 --- a/mobile/test/shared/relay/relay_session_test.dart +++ b/mobile/test/shared/relay/relay_session_test.dart @@ -749,6 +749,61 @@ void main() { }, ); + test('retryable CLOSED reports retrying until replay is ready', () async { + final timers = <_ManualTimer>[]; + final socket = _RecordingRelaySocket(); + final deliveredEvents = []; + final statuses = []; + final session = RelaySessionNotifier( + retryTimerFactory: (duration, callback) { + final timer = _ManualTimer(duration, callback); + timers.add(timer); + return timer; + }, + ); + session.debugAttachSocketForTest(socket); + + final subscribe = session.subscribeWithStatus( + _channelFilter, + deliveredEvents.add, + onStatusChanged: (status) { + if (status == RelaySubscriptionStatus.ready) { + expect(deliveredEvents, hasLength(statuses.isEmpty ? 0 : 1)); + } + statuses.add(status); + }, + ); + session.debugHandleMessage(['EOSE', 'l-1']); + final unsubscribe = await subscribe; + expect(statuses, [RelaySubscriptionStatus.ready]); + + session.debugHandleMessage(['CLOSED', 'l-1', 'error: relay overloaded']); + expect(statuses, [ + RelaySubscriptionStatus.ready, + RelaySubscriptionStatus.retrying, + ]); + + timers.single.fire(); + await Future.delayed(Duration.zero); + session.debugHandleMessage([ + 'EVENT', + 'l-1', + _event(createdAt: 30).toJson(), + ]); + expect(statuses, [ + RelaySubscriptionStatus.ready, + RelaySubscriptionStatus.retrying, + ]); + + session.debugHandleMessage(['EOSE', 'l-1']); + expect(statuses, [ + RelaySubscriptionStatus.ready, + RelaySubscriptionStatus.retrying, + RelaySubscriptionStatus.ready, + ]); + unsubscribe(); + }); + test('CLOSED retries back off and reset after EOSE', () async { final timers = <_ManualTimer>[]; final socket = _RecordingRelaySocket();