diff --git a/mobile/lib/features/channels/channel_management_actions.dart b/mobile/lib/features/channels/channel_management_actions.dart index e0879eb0222..ad59f71be93 100644 --- a/mobile/lib/features/channels/channel_management_actions.dart +++ b/mobile/lib/features/channels/channel_management_actions.dart @@ -77,10 +77,14 @@ class ChannelActions { return _refreshChannelsAndRead(channelId); } + /// Reports each acknowledged write before checking continuation scope. + /// [onAccepted] records irreversible outcomes; it must not mutate scope caches. + /// The community and enclosing operation scope also fence queued writes. Future addMembers({ required String channelId, required List pubkeys, String role = 'member', + ValueChanged? onAccepted, }) async { final normalizedRole = role.trim(); if (normalizedRole.isEmpty) { @@ -100,18 +104,23 @@ class ChannelActions { // add, not be recorded as this pubkey's rejection. _ensureCommunityValid(); try { - await _signedEventRelay.submit( - kind: 9000, - content: '', - tags: [ - ['h', channelId], - ['p', pubkey], - ['role', normalizedRole], - ], + await withRelayPublicationGuard( + _ensureCommunityValid, + () => _signedEventRelay.submit( + kind: 9000, + content: '', + tags: [ + ['h', channelId], + ['p', pubkey], + ['role', normalizedRole], + ], + ), ); } catch (error) { failures[pubkey] = _relayErrorMessage(error); + continue; } + onAccepted?.call(pubkey); } _ensureCommunityValid(); _ref.invalidate(channelMembersProvider(channelId)); diff --git a/mobile/lib/features/channels/compose_bar/compose_bar_widget.dart b/mobile/lib/features/channels/compose_bar/compose_bar_widget.dart index b335aa6ecd2..e23356e2f01 100644 --- a/mobile/lib/features/channels/compose_bar/compose_bar_widget.dart +++ b/mobile/lib/features/channels/compose_bar/compose_bar_widget.dart @@ -82,6 +82,7 @@ class ComposeBar extends HookConsumerWidget { final uploadingCount = useState(0); final uploadProgress = useState(0.0); final uploadGeneration = useRef(0); + final invitationStarted = useState(false); final activeUploadCancellation = useRef(null); final voiceNote = _useComposerVoiceNote( context: context, @@ -473,6 +474,7 @@ class ComposeBar extends HookConsumerWidget { uploadingCount.value > 0) { return; } + invitationStarted.value = false; final attempt = Object(); authorizationAttempt.value = attempt; isSending.value = true; @@ -536,7 +538,10 @@ class ComposeBar extends HookConsumerWidget { for (final entry in mentionMap.value.entries) if (hasMention(text, entry.key)) entry.value, ]; - final outgoing = _OutgoingMentions(selectedMentions); + final outgoing = _OutgoingMentions( + selectedMentions, + '${ref.read(relayConfigProvider).baseUrl} / $channelId', + ); final intendedAgentKeys = { for (final mention in selectedMentions) if (mention.isAgent) mention.pubkey.toLowerCase(), @@ -580,6 +585,8 @@ class ComposeBar extends HookConsumerWidget { Future addMentionedNonMembers() async { final keys = intendedAgentKeys.intersection(outgoing.pubkeys.toSet()); await authorize(keys, prepare: true); + ensureAuthorizationCurrent(); + invitationStarted.value = true; await outgoing.addNonMembers( channelActions, scan: scan, @@ -646,6 +653,7 @@ class ComposeBar extends HookConsumerWidget { final delivery = onSend; unawaited(() async { var retainedForRetry = false; + var delivered = false; try { final uploaded = []; for (var index = 0; index < queuedAttachments.length; index++) { @@ -685,6 +693,7 @@ class ComposeBar extends HookConsumerWidget { outgoing.pubkeys, mediaTags: [...payload.mediaTags, ...outgoing.referenceTags], ); + delivered = true; } on _ComposeAuthorizationCancelled { // Keep the newer draft without displaying a false access error. } catch (error) { @@ -706,6 +715,8 @@ class ComposeBar extends HookConsumerWidget { focusNode.requestFocus(); } } finally { + if (!delivered) outgoing.reportIncomplete(messenger); + if (ownsSource()) invitationStarted.value = false; final sourceRetainsFiles = preparingAgents && context.mounted && @@ -747,6 +758,7 @@ class ComposeBar extends HookConsumerWidget { } finally { if (context.mounted && authorizationAttempt.value == attempt) { authorizationAttempt.value = null; + if (uploadingCount.value == 0) invitationStarted.value = false; isSending.value = false; } } @@ -1065,6 +1077,7 @@ class ComposeBar extends HookConsumerWidget { visible: hasPendingUploads, progress: uploadProgress.value, reducedMotion: reducedMotion, + cancelLabel: invitationStarted.value ? 'Stop remaining' : 'Cancel', onCancel: () { activeUploadCancellation.value?.cancel(); uploadGeneration.value += 1; diff --git a/mobile/lib/features/channels/compose_bar/draft_lifecycle.dart b/mobile/lib/features/channels/compose_bar/draft_lifecycle.dart index 3e085381b80..d412288d01a 100644 --- a/mobile/lib/features/channels/compose_bar/draft_lifecycle.dart +++ b/mobile/lib/features/channels/compose_bar/draft_lifecycle.dart @@ -19,6 +19,7 @@ Future _sendTextOnlyDraft({ required ComposeBarOnSend onSend, required ScaffoldMessengerState? messenger, }) async { + var delivered = false; TextEditingValue? clearedDraftText; Map? clearedDraftMentions; int? clearedDraftRevision; @@ -55,6 +56,7 @@ Future _sendTextOnlyDraft({ outgoing.pubkeys, mediaTags: [...payload.mediaTags, ...outgoing.referenceTags], ); + delivered = true; } on _ComposeAuthorizationCancelled { restoreClearedDraft(); } on StateError { @@ -64,9 +66,13 @@ Future _sendTextOnlyDraft({ // The caller runs unawaited, so surface publish failures and restore the // sent draft unless the user has already started a new one. restoreClearedDraft(); - messenger?.showSnackBar( - SnackBar(content: Text(_composeSendErrorMessage(error))), - ); + if (outgoing.acceptedInvitations == 0) { + messenger?.showSnackBar( + SnackBar(content: Text(_composeSendErrorMessage(error))), + ); + } + } finally { + if (!delivered) outgoing.reportIncomplete(messenger); } } diff --git a/mobile/lib/features/channels/compose_bar/helpers.dart b/mobile/lib/features/channels/compose_bar/helpers.dart index fa4347a9f29..42d143c8ed0 100644 --- a/mobile/lib/features/channels/compose_bar/helpers.dart +++ b/mobile/lib/features/channels/compose_bar/helpers.dart @@ -342,7 +342,8 @@ Future<_NonMemberMentionChoice?> _promptNonMemberMention( content: Text( canInvite ? '${names.join(', ')} $verb not in this channel. Invite them to ' - 'the channel, or send without inviting them.' + 'the channel, or send without inviting them. Invitations take effect ' + 'immediately and remain if the message is stopped or fails.' : '${names.join(', ')} $verb not in this channel. ' '$privateChannelAddDeniedMessage You can still send without ' 'inviting them.', @@ -448,6 +449,7 @@ Future<_NonMemberAddOutcome> _addMentionedNonMembers( required List humanPubkeys, required bool canAddMembers, required VoidCallback ensureCurrent, + required VoidCallback onAccepted, }) async { final pending = [ for (final pubkey in agentPubkeys) ([pubkey], 'bot'), @@ -473,6 +475,7 @@ Future<_NonMemberAddOutcome> _addMentionedNonMembers( channelId: channelId, pubkeys: pubkeys, role: role, + onAccepted: (_) => onAccepted(), ); ensureCurrent(); } on _ComposeAuthorizationCancelled { @@ -574,9 +577,26 @@ class _OutgoingMentions { final List> referenceTags = []; List _invitedHumanPubkeys = const []; bool _inviteAgents = false; + int acceptedInvitations = 0; + final String sourceDestination; + + void reportIncomplete(ScaffoldMessengerState? messenger) { + if (acceptedInvitations == 0) return; + messenger?.showSnackBar( + SnackBar( + content: Text( + 'Message not sent. $acceptedInvitations invitation(s) completed and remain ' + 'in effect in $sourceDestination. Review channel members before retrying. ' + 'Check your draft; attachments may need reattaching after leaving.', + ), + ), + ); + } - _OutgoingMentions(List selectedMentions) - : pubkeys = LinkedHashSet.from( + _OutgoingMentions( + List selectedMentions, + this.sourceDestination, + ) : pubkeys = LinkedHashSet.from( selectedMentions.map((candidate) => candidate.pubkey.toLowerCase()), ).toList(); @@ -623,6 +643,7 @@ class _OutgoingMentions { humanPubkeys: _invitedHumanPubkeys, canAddMembers: scan.canAddMembers, ensureCurrent: ensureCurrent, + onAccepted: () => acceptedInvitations++, ); if (outcome.notAdded.isNotEmpty) { throw Exception('Message not sent. ${outcome.errors.join(' ')}'); diff --git a/mobile/lib/features/channels/compose_bar/upload_progress_pill.dart b/mobile/lib/features/channels/compose_bar/upload_progress_pill.dart index 23f381b82be..0ded8431558 100644 --- a/mobile/lib/features/channels/compose_bar/upload_progress_pill.dart +++ b/mobile/lib/features/channels/compose_bar/upload_progress_pill.dart @@ -5,12 +5,14 @@ class _UploadProgressMotion extends StatelessWidget { final double progress; final bool reducedMotion; final VoidCallback onCancel; + final String cancelLabel; const _UploadProgressMotion({ required this.visible, required this.progress, required this.reducedMotion, required this.onCancel, + this.cancelLabel = 'Cancel', }); @override @@ -47,6 +49,7 @@ class _UploadProgressMotion extends StatelessWidget { progress: progress, reducedMotion: reducedMotion, onCancel: onCancel, + cancelLabel: cancelLabel, ), ) : const SizedBox.shrink( @@ -62,11 +65,13 @@ class _UploadProgressPill extends HookConsumerWidget { final double progress; final bool reducedMotion; final VoidCallback onCancel; + final String cancelLabel; const _UploadProgressPill({ required this.progress, required this.reducedMotion, required this.onCancel, + this.cancelLabel = 'Cancel', }); @override @@ -172,7 +177,7 @@ class _UploadProgressPill extends HookConsumerWidget { vertical: Grid.quarter, ), child: Text( - 'Cancel', + cancelLabel, style: context.textTheme.labelMedium ?.copyWith( color: context.colors.onSurface, diff --git a/mobile/lib/features/channels/send_message_provider.dart b/mobile/lib/features/channels/send_message_provider.dart index 757275c60ac..f8966991056 100644 --- a/mobile/lib/features/channels/send_message_provider.dart +++ b/mobile/lib/features/channels/send_message_provider.dart @@ -97,15 +97,18 @@ class SendMessage { _ensureDeliveryValid(); NostrEvent? localMessage; try { - await _signedEventRelay.submit( - kind: EventKind.streamMessage, - content: content, - tags: tags, - onSigned: (event) { - localMessage = event; - _markLocalMessageForAnimation(channelId, event.id); - _addLocalMessage(channelId, event); - }, + await withRelayPublicationGuard( + _ensureDeliveryValid, + () => _signedEventRelay.submit( + kind: EventKind.streamMessage, + content: content, + tags: tags, + onSigned: (event) { + localMessage = event; + _markLocalMessageForAnimation(channelId, event.id); + _addLocalMessage(channelId, event); + }, + ), ); final event = localMessage; if (event != null) _completeLocalMessage(channelId, event.id); diff --git a/mobile/lib/shared/mentions/agent_publication.dart b/mobile/lib/shared/mentions/agent_publication.dart index a7fb1e24732..2f46f0bc687 100644 --- a/mobile/lib/shared/mentions/agent_publication.dart +++ b/mobile/lib/shared/mentions/agent_publication.dart @@ -67,3 +67,30 @@ Future authorizeAgentMentions( throw Exception(message); } } + +/// Session-bound all-key evidence reader. Saved agent keys are denial-only taint. +typedef SelectedMentionAuthorizationReader = + Future> Function( + Set keys, + Set priorAgentKeys, + String viewer, + String channelId, + bool Function() isCurrent, + void Function(Map) onProfileEvidence, + ); + +/// Read-only publication seam; does not populate suggestion/profile caches. +final selectedMentionAuthorizationReaderProvider = + Provider((ref) { + final session = ref.watch(relaySessionProvider.notifier); + return (keys, prior, viewer, channel, current, observed) => + readSelectedMentionAuthorization( + session, + keys, + viewer: viewer, + channelId: channel, + priorAgentKeys: prior, + isCurrent: current, + onProfileEvidence: observed, + ); + }); diff --git a/mobile/lib/shared/mentions/selected_mention_authorization.dart b/mobile/lib/shared/mentions/selected_mention_authorization.dart index c7ef3710887..152426234d6 100644 --- a/mobile/lib/shared/mentions/selected_mention_authorization.dart +++ b/mobile/lib/shared/mentions/selected_mention_authorization.dart @@ -36,7 +36,8 @@ class SelectedMentionAuthorization { /// owner policy produces a deny-all entry, never runtime fallback. final AgentDirectoryEntry? agent; - const SelectedMentionAuthorization._(this.kind, this.isMember, this.agent); + /// Construct evidence, not a publication or role-write capability. + const SelectedMentionAuthorization(this.kind, this.isMember, this.agent); /// True when this key's evidence must be routed through the agent /// authorization evaluator rather than the ordinary flow: agent or @@ -70,6 +71,7 @@ readSelectedMentionAuthorization( required String channelId, Set priorAgentKeys = const {}, required bool Function() isCurrent, + void Function(Map)? onProfileEvidence, }) async { void check() { if (!isCurrent()) throw StateError('Mention authorization scope changed'); @@ -131,6 +133,7 @@ readSelectedMentionAuthorization( checkCurrent: check, onProfileEvidence: (value) => profiles = value, ); + onProfileEvidence?.call(profiles); check(); final latestRuntime = {}; for (final event in runtime) { @@ -186,7 +189,7 @@ readSelectedMentionAuthorization( : const [], ); } - result[key] = SelectedMentionAuthorization._( + result[key] = SelectedMentionAuthorization( kind, members.containsKey(key), agent, diff --git a/mobile/lib/shared/profile/user_cache_provider.dart b/mobile/lib/shared/profile/user_cache_provider.dart index 91131c888d0..e54083aa269 100644 --- a/mobile/lib/shared/profile/user_cache_provider.dart +++ b/mobile/lib/shared/profile/user_cache_provider.dart @@ -31,6 +31,10 @@ class UserCacheNotifier extends Notifier> { return {}; } + /// Latest observed profile revision; read-only fencing for publication. + ({int createdAt, String eventId})? profileEventOrder(String key) => + _profileEventOrders[key]; + /// Request a profile for [pubkey]. Returns immediately from cache if /// available, otherwise schedules a batch fetch. UserProfile? get(String pubkey) { diff --git a/mobile/lib/shared/relay/relay_rate_limit_gate.dart b/mobile/lib/shared/relay/relay_rate_limit_gate.dart index 92ea7d85f01..350cc132025 100644 --- a/mobile/lib/shared/relay/relay_rate_limit_gate.dart +++ b/mobile/lib/shared/relay/relay_rate_limit_gate.dart @@ -21,6 +21,9 @@ class RelayRateLimitGate { final DateTime Function() _now; final RelayTimerFactory _timerFactory; + + /// Monotonic revision of admitted capacity pauses, including extensions. + int epoch = 0; DateTime? _expiresAt; Timer? _timer; Completer? _completer; @@ -53,6 +56,7 @@ class RelayRateLimitGate { final currentExpiry = _expiresAt; if (currentExpiry != null && !newExpiry.isAfter(currentExpiry)) return; + epoch++; _expiresAt = newExpiry; _timer?.cancel(); _completer ??= Completer(); diff --git a/mobile/lib/shared/relay/relay_session.dart b/mobile/lib/shared/relay/relay_session.dart index b761f23c662..b0698416932 100644 --- a/mobile/lib/shared/relay/relay_session.dart +++ b/mobile/lib/shared/relay/relay_session.dart @@ -24,6 +24,26 @@ export 'relay_session_types.dart'; part 'relay_session_auth.dart'; +final _publicationGuardKey = Object(); + +/// Carries a synchronous local operation fence through async delivery adapters. +/// The fence runs after transport backpressure and before socket enqueue only; +/// it never changes how an already transmitted event's ACK is accounted for. +/// Nested scopes retain their enclosing operation fence. Other callers need not +/// opt in. This is not atomic authorization against unobserved remote changes. +T withRelayPublicationGuard(void Function() check, T Function() deliver) { + final enclosing = Zone.current[_publicationGuardKey] as void Function()?; + return runZoned( + deliver, + zoneValues: { + _publicationGuardKey: () { + enclosing?.call(); + check(); + }, + }, + ); +} + class _HistorySubscription { final List events = []; final Completer> completer; @@ -327,6 +347,8 @@ class RelaySessionNotifier extends Notifier { return () => _unsubscribe(subId); } + /// Publishes after backpressure; the operation scope may synchronously abort. + /// Once enqueued, ACKs settle normally even if the operation later expires. Future publish( NostrEvent event, { Duration timeout = const Duration(seconds: 8), @@ -337,6 +359,9 @@ class RelaySessionNotifier extends Notifier { throw StateError('Relay session is not connected'); } + // No await between the operation fence and the actual socket enqueue. + // This is local observed currentness, not network-atomic authorization. + (Zone.current[_publicationGuardKey] as void Function()?)?.call(); final completer = Completer(); final timer = Timer(timeout, () { @@ -381,7 +406,22 @@ class RelaySessionNotifier extends Notifier { void debugDispose() => _dispose(); @visibleForTesting - void debugSupersedeConnection() => _connectionGeneration++; + void debugSupersedeConnection() => _supersedeConnection(); + + int _supersedeConnection() { + return ++_connectionGeneration; + } + + /// Selected evidence cannot cross socket replacement or a capacity pause. + bool Function() retainPublicationEvidence() { + final generation = _connectionGeneration; + final epoch = _rateLimitGate.epoch; + final admitted = !_rateLimitGate.isActive; + return () => + admitted && + _isActiveConnection(generation) && + epoch == _rateLimitGate.epoch; + } @visibleForTesting void debugHandleDisconnected([Object? error]) { @@ -494,7 +534,7 @@ class RelaySessionNotifier extends Notifier { Future _connect(RelayConfig config) async { if (_disposed) return; - final generation = ++_connectionGeneration; + final generation = _supersedeConnection(); state = SessionState( status: _hasConnectedOnce ? SessionStatus.reconnecting diff --git a/mobile/lib/shared/relay/signed_event_relay.dart b/mobile/lib/shared/relay/signed_event_relay.dart index b0106639eb7..9e75d12790f 100644 --- a/mobile/lib/shared/relay/signed_event_relay.dart +++ b/mobile/lib/shared/relay/signed_event_relay.dart @@ -28,7 +28,8 @@ class SignedEventRelay { /// Sign and submit an event. Returns the relay's OK response as a [NostrEvent] /// whose `content` field contains the OK message (e.g. `"response:{...}"` - /// for command kinds). + /// for command kinds). Publication retains any [withRelayPublicationGuard] + /// scope through transport waits, without discarding accepted ACKs. Future submit({ required int kind, required String content, diff --git a/mobile/test/features/channels/channel_detail_page_test.dart b/mobile/test/features/channels/channel_detail_page_test.dart index 17d83c804d8..1aac770eaa0 100644 --- a/mobile/test/features/channels/channel_detail_page_test.dart +++ b/mobile/test/features/channels/channel_detail_page_test.dart @@ -14561,8 +14561,12 @@ class _FakeChannelActions extends ChannelActions { required String channelId, required List pubkeys, String role = 'member', + ValueChanged? onAccepted, }) async { await onAddMembers?.call(channelId, pubkeys); + for (final pubkey in pubkeys) { + onAccepted?.call(pubkey); + } } @override diff --git a/mobile/test/features/channels/compose_bar_test.dart b/mobile/test/features/channels/compose_bar_test.dart index 7b9b6d29ff8..39078101d00 100644 --- a/mobile/test/features/channels/compose_bar_test.dart +++ b/mobile/test/features/channels/compose_bar_test.dart @@ -29,6 +29,9 @@ import 'package:buzz/shared/theme/theme.dart'; import 'package:buzz/shared/widgets/anchored_popover_menu.dart'; import 'package:buzz/shared/widgets/mobile_tab_footer_backdrop.dart'; import 'package:shared_preferences/shared_preferences.dart'; +// Shared fixture prerequisite; production-observation rows land in the child PR. +// ignore: unused_import +import '../../shared/mentions/agent_policy_test.dart' show signed; part 'compose_bar_test/publication_tests.dart'; part 'compose_bar_test/send_lifecycle_tests.dart'; @@ -180,6 +183,10 @@ Widget _buildComposeBar({ Future>? membersFuture, Future> Function()? membersLoader, AgentAuthorizationReader? authorizationReader, + SelectedMentionAuthorizationReader? selectedReader, + RelayRateLimitGate? rateLimitGate, + http.Client? relayHttpClient, + void Function(NostrEvent)? beforePublish, List relayAgents = const [], List channels = const [], List cachedMembers = const [], @@ -202,6 +209,14 @@ Widget _buildComposeBar({ }) { return ProviderScope( overrides: [ + if (rateLimitGate != null || relayHttpClient != null) + relaySessionProvider.overrideWith( + () => _BeforePublishSession( + rateLimitGate: rateLimitGate, + httpClient: relayHttpClient, + beforePublish: beforePublish, + ), + ), customEmojiListProvider.overrideWithValue(customEmoji), mediaUploadServiceProvider.overrideWithValue(uploadService), if (voiceNoteRecorderFactory != null) @@ -230,6 +245,49 @@ Widget _buildComposeBar({ ), ], ), + if (relayHttpClient == null || selectedReader != null) + selectedMentionAuthorizationReaderProvider.overrideWithValue( + selectedReader ?? + (keys, prior, viewer, channel, current, observed) async { + final roster = + await (membersLoader?.call() ?? + membersFuture ?? + Future.value(members)); + final agentKeys = { + ...prior, + for (final agent in relayAgents) + if (keys.contains(agent.pubkey)) agent.pubkey, + }; + final agents = agentKeys.isEmpty + ? [] + : await (authorizationReader?.call( + agentKeys, + viewer, + channel, + current, + ) ?? + Future.value([ + for (final key in agentKeys) + AgentDirectoryEntry( + pubkey: key, + respondTo: 'anyone', + ownerPubkey: viewer, + channelIds: [channel], + ), + ])); + return { + for (final key in keys) + key: SelectedMentionAuthorization( + agentKeys.contains(key) + ? SelectedMentionKind.agent + : SelectedMentionKind.ordinary, + roster.any((member) => member.pubkey == key) || + channels.any((c) => c.id == channel && c.isDm), + agents.where((agent) => agent.pubkey == key).firstOrNull, + ), + }; + }, + ), agentDirectoryProvider.overrideWith((ref) async => relayAgents), agentOwnersProvider.overrideWith((ref) async => const {}), relayClientProvider.overrideWithValue( @@ -579,6 +637,33 @@ class _SwitchableRelayConfigNotifier extends RelayConfigNotifier { RelayConfig build() => initial; } +// Shared fixture prerequisite; production-observation rows land in the child PR. +// ignore: unused_element +http.Client _selectedRosterClient( + String authority, + NostrEvent Function() rosterEvent, { + List Function()? extraEvents, +}) => http_testing.MockClient((request) async { + if (request.method == 'GET') { + return http.Response(jsonEncode({'self': authority}), 200); + } + final filters = jsonDecode(request.body) as List; + return http.Response( + jsonEncode([ + if (filters.any((f) => (f['kinds'] as List).contains(39002))) + rosterEvent().toJson(), + for (final event in extraEvents?.call() ?? []) + if (filters.any( + (f) => + (f['kinds'] as List).contains(event.kind) && + (f['authors'] as List).contains(event.pubkey), + )) + event.toJson(), + ]), + 200, + ); +}); + class _RecordingRelaySocket extends RelaySocket { final List> events; final void Function(List message) handleMessage; @@ -5674,3 +5759,22 @@ Channel _makeChannel({required String name, required String channelType}) { memberCount: 5, ); } + +// Observe the real signing→session boundary without supplying a validity guard. +class _BeforePublishSession extends RelaySessionNotifier { + _BeforePublishSession({ + super.rateLimitGate, + super.httpClient, + this.beforePublish, + }); + final void Function(NostrEvent)? beforePublish; + + @override + Future publish( + NostrEvent event, { + Duration timeout = const Duration(seconds: 8), + }) { + beforePublish?.call(event); + return super.publish(event, timeout: timeout); + } +} diff --git a/mobile/test/features/channels/compose_bar_test/send_lifecycle_tests.dart b/mobile/test/features/channels/compose_bar_test/send_lifecycle_tests.dart index dc427104ac8..ed067fc61f9 100644 --- a/mobile/test/features/channels/compose_bar_test/send_lifecycle_tests.dart +++ b/mobile/test/features/channels/compose_bar_test/send_lifecycle_tests.dart @@ -84,7 +84,12 @@ void sendLifecycleTests() { } }); } - for (final action in ['revisit', 'cancel', 'partial refusal']) { + for (final action in [ + 'revisit', + 'cancel', + 'partial refusal', + 'accepted scope switch', + ]) { testWidgets('$action cannot finish an old membership batch', ( tester, ) async { @@ -98,6 +103,9 @@ void sendLifecycleTests() { Widget build({String? thread}) => _buildComposeBar( uploadService: service, currentPubkey: signer.public, + relayConfig: () => _SwitchableRelayConfigNotifier( + RelayConfig(baseUrl: 'https://relay.example', nsec: signer.nsec), + ), relayAgents: [ _testAgent('b' * 64), AgentDirectoryEntry( @@ -126,6 +134,11 @@ void sendLifecycleTests() { if (event['kind'] == 9000) await gate.future; }, onEventAcknowledged: (event) { + if (action == 'accepted scope switch' && event['kind'] == 9000) { + container + .read(relayConfigProvider.notifier) + .update(baseUrl: 'https://other.example', nsec: signer.nsec); + } if (action == 'partial refusal' && event['kind'] == 9000) { session.debugAttachSocketForTest( _RecordingRelaySocket( @@ -138,10 +151,11 @@ void sendLifecycleTests() { }, ), ); - if (action == 'cancel') { + if (action == 'cancel' || action == 'revisit') { await _openAttachmentMenu(tester); await tester.tap(find.text('Video')); await tester.pumpAndSettle(); + expect(find.byTooltip('Remove attachment'), findsOneWidget); } await _expandComposer(tester); await tester.enterText(find.byType(TextField), '@hel'); @@ -155,11 +169,14 @@ void sendLifecycleTests() { await tester.tap(find.byIcon(LucideIcons.arrowUp)); await tester.pump(); await tester.pump(const Duration(milliseconds: 300)); + expect(find.textContaining('Invitations take effect'), findsOneWidget); await tester.tap(find.text('Invite')); await tester.pump(); expect(events.where((e) => e['kind'] == 9000), hasLength(1)); if (action == 'cancel') { await tester.pump(const Duration(milliseconds: 300)); + expect(find.text('Cancel'), findsNothing); + expect(find.text('Stop remaining'), findsOneWidget); await tester.tap(find.byKey(const ValueKey('compose-upload-cancel'))); } else if (action == 'revisit') { await tester.pumpWidget(build(thread: 'other')); @@ -167,9 +184,33 @@ void sendLifecycleTests() { } gate.complete(); await tester.pumpAndSettle(); + if (action == 'accepted scope switch') { + await tester.pump(const Duration(seconds: 5)); + await tester.pumpAndSettle(); + } + expect( + find.textContaining('1 invitation(s) completed and remain in effect'), + findsOneWidget, + ); + expect(find.textContaining('Your draft is kept'), findsNothing); + expect( + find.textContaining('https://relay.example / channel-1'), + findsOneWidget, + ); + expect( + find.textContaining('attachments may need reattaching'), + findsOneWidget, + ); + if (action == 'revisit') { + expect(find.byTooltip('Remove attachment'), findsNothing); + } await tester.pump(const Duration(seconds: 5)); await tester.pumpAndSettle(); expect(sends, 0); + if (action == 'accepted scope switch') { + expect(events.where((e) => e['kind'] == 9000), hasLength(1)); + return; + } await tester.tap(find.text('@Helper Bot @Other Bot')); await tester.pumpAndSettle(); expect( diff --git a/mobile/test/shared/relay/relay_enqueue_fence_test.dart b/mobile/test/shared/relay/relay_enqueue_fence_test.dart new file mode 100644 index 00000000000..25e15202956 --- /dev/null +++ b/mobile/test/shared/relay/relay_enqueue_fence_test.dart @@ -0,0 +1,191 @@ +import 'package:buzz/features/channels/channel_management_provider.dart'; +import 'package:buzz/shared/relay/relay.dart'; +import 'package:buzz/features/channels/send_message_provider.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:hooks_riverpod/hooks_riverpod.dart'; +import 'package:nostr/nostr.dart' as nostr; + +void main() { + // Each row binds to the production object's own validity callback — + // ChannelActions.isCommunityValid for kind:9000 and + // SendMessage.isDeliveryValid for messages — exactly as + // channelActionsProvider/sendMessageProvider wire them. There is no + // test-authored outer guard here, so removing either production + // `withRelayPublicationGuard` wrapper must fail its own rows. + for (final kind in [9000, EventKind.streamMessage]) { + for (final cancel in [false, true]) { + test( + 'queued kind:$kind cancellation=$cancel via the owning guard', + () async { + final gate = RelayRateLimitGate(); + final session = RelaySessionNotifier(rateLimitGate: gate); + final socket = _Socket(); + session.debugAttachSocketForTest(socket); + final relay = SignedEventRelay( + session: session, + nsec: nostr.Keys.generate().nsec, + ); + var valid = true; + var accepted = 0; + final container = ProviderContainer(); + addTearDown(container.dispose); + final actions = container.read( + Provider( + (ref) => ChannelActions( + ref: ref, + session: session, + signedEventRelay: relay, + currentPubkey: relay.pubkey, + isCommunityValid: () => valid, + ), + ), + ); + gate.activate(10); + // Real production signing/session/backpressure and ChannelActions for + // kind9000. No query/selected-reader or publication recorder substitutes. + final pending = kind == 9000 + ? actions.addMembers( + channelId: 'channel', + pubkeys: ['target'], + onAccepted: (_) => accepted++, + ) + : SendMessage( + signedEventRelay: relay, + fetchMembers: (_) async => const [], + readUserCache: () => const {}, + addLocalMessage: (_, _) {}, + completeLocalMessage: (_, _) => accepted++, + removeLocalMessage: (_, _) {}, + isDeliveryValid: () => valid, + )(channelId: 'channel', content: '', mentionPubkeys: const []); + Object? error; + final settled = pending.then( + (_) {}, + onError: (Object e) { + error = e; + }, + ); + expect(socket.messages, isEmpty); + // Invalidate the owning production object's own callback while the + // real rate-limit gate still holds the publish. + valid = !cancel; + gate.reset(); + await Future.delayed(Duration.zero); + final sent = socket.messages + .where((p) => p.first == 'EVENT') + .toList(); + if (sent.isNotEmpty) { + // Invalidation after the socket write must not erase accepted ACKs. + valid = false; + session.debugHandleMessage([ + 'OK', + (sent.single[1] as Map)['id'], + true, + '', + ]); + } + await settled; + expect(sent, hasLength(cancel ? 0 : 1)); + expect(accepted, cancel ? 0 : 1); + if (kind == 9000) { + // The membership object reports its own fence after the loop; the + // irreversible write above still counts as accepted. + expect(error, isA()); + } else { + // A message send that reached the socket is completed, not + // cancelled, even though the operation expired after the enqueue. + expect(error, cancel ? isA() : isNull); + } + }, + ); + } + } + + // The composer wraps ChannelActions in its own operation fence, so the + // enqueue fence must consult the enclosing scope and the production + // object's own callback. This is the only test-authored outer guard, and + // it exists explicitly for that nesting contract. + test('nested operation fences compose at actual enqueue', () async { + for (final outerInvalid in [true, false]) { + final gate = RelayRateLimitGate(); + final session = RelaySessionNotifier(rateLimitGate: gate); + final socket = _Socket(); + session.debugAttachSocketForTest(socket); + final relay = SignedEventRelay( + session: session, + nsec: nostr.Keys.generate().nsec, + ); + var outerValid = true; + var innerValid = true; + var accepted = 0; + final container = ProviderContainer(); + addTearDown(container.dispose); + final actions = container.read( + Provider( + (ref) => ChannelActions( + ref: ref, + session: session, + signedEventRelay: relay, + currentPubkey: relay.pubkey, + isCommunityValid: () => innerValid, + ), + ), + ); + void ensureOuter() { + if (!outerValid) throw StateError('outer operation expired'); + } + + gate.activate(10); + final pending = withRelayPublicationGuard( + ensureOuter, + () => actions.addMembers( + channelId: 'channel', + pubkeys: ['target'], + onAccepted: (_) => accepted++, + ), + ); + Object? error; + final settled = pending.then( + (_) {}, + onError: (Object e) { + error = e; + }, + ); + expect(socket.messages, isEmpty); + if (outerInvalid) { + outerValid = false; // Enclosing operation expires during the wait. + } else { + innerValid = false; // Production object's own fence expires. + } + gate.reset(); + await Future.delayed(Duration.zero); + expect(socket.messages.where((p) => p.first == 'EVENT'), isEmpty); + await settled; + expect(accepted, 0); + if (outerInvalid) { + // The enclosing fence's failure is recorded per-pubkey and the + // membership object surfaces it as an AddMembersException. + final failure = error as AddMembersException; + expect(failure.failures['target'], contains('outer operation expired')); + } else { + // The inner fence's failure aborts the whole add before the + // per-pubkey failure report. + expect(error, isA()); + } + } + }); +} + +class _Socket extends RelaySocket { + _Socket() + : super( + wsUrl: 'wss://relay.example', + nsec: null, + onMessage: (_) {}, + onConnected: () {}, + onDisconnected: (_) {}, + ); + final messages = >[]; + @override + void send(List payload) => messages.add(payload); +}