From 7877a16cdcc81eb10982c6b159f216516266a7d8 Mon Sep 17 00:00:00 2001 From: grunch Date: Sun, 13 Sep 2026 14:40:10 -0300 Subject: [PATCH 1/2] =?UTF-8?q?feat(push):=20PR-0b=20=E2=80=94=20app=20lif?= =?UTF-8?q?ecycle=20service=20and=20resume=20hydration?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit AppLifecycleService turns paused/hidden → resumed into one call, with a latch (an inactive flap never suspended anything) and a debounce. ResumeResync runs resync() in Rust, then every hydrator: trades, chat rooms, disputes and notifications re-read their state from the bridge through the same path cold start uses. Disputes gain that path here: a dispute the peer opened while the process was suspended is listed on resume without a restart (the v1 bug of mobile#675). docs/PUSH_NOTIFICATIONS.md Phase 0, T0.2 + T0.3. Closes #308. --- lib/core/app_bootstrap.dart | 9 + lib/core/lifecycle/app_lifecycle_service.dart | 91 +++++++++ lib/core/lifecycle/resume_resync.dart | 66 +++++++ .../chat/providers/chat_providers.dart | 100 +++++----- .../providers/disputes_providers.dart | 94 +++++++++- .../providers/notifications_provider.dart | 49 +++-- .../trades/providers/trades_providers.dart | 10 + .../lifecycle/app_lifecycle_service_test.dart | 155 ++++++++++++++++ .../lifecycle/hydration_regression_test.dart | 172 ++++++++++++++++++ test/core/lifecycle/resume_resync_test.dart | 70 +++++++ 10 files changed, 747 insertions(+), 69 deletions(-) create mode 100644 lib/core/lifecycle/app_lifecycle_service.dart create mode 100644 lib/core/lifecycle/resume_resync.dart create mode 100644 test/core/lifecycle/app_lifecycle_service_test.dart create mode 100644 test/core/lifecycle/hydration_regression_test.dart create mode 100644 test/core/lifecycle/resume_resync_test.dart diff --git a/lib/core/app_bootstrap.dart b/lib/core/app_bootstrap.dart index dbceec15..cba94416 100644 --- a/lib/core/app_bootstrap.dart +++ b/lib/core/app_bootstrap.dart @@ -13,6 +13,8 @@ import 'package:mostro/core/font_licenses.dart'; import 'package:mostro/core/mostro_defaults.dart'; import 'package:mostro/core/services/identity_service.dart'; import 'package:mostro/core/test_environment.dart'; +import 'package:mostro/core/lifecycle/app_lifecycle_service.dart'; +import 'package:mostro/core/lifecycle/resume_resync.dart'; import 'package:mostro/core/web/bridge_probe.dart'; import 'package:mostro/features/settings/providers/settings_provider.dart'; import 'package:mostro/features/settings/widgets/mostro_node_selector.dart'; @@ -209,6 +211,13 @@ Future bootstrapAndRun({List seedRelays = const []}) async { _consumeBondSlashed(bondSlashedStream, container); _consumeBondClaims(bondClaimStream, container); + // Resume = resync in Rust, then re-hydrate every notifier from the bridge + // (issue #308, docs/PUSH_NOTIFICATIONS.md §10). Attached before runApp so + // the first suspension is observed too. + AppLifecycleService( + onResume: ResumeResync(container: container).run, + ).attach(); + runApp( UncontrolledProviderScope(container: container, child: const MostroApp()), ); diff --git a/lib/core/lifecycle/app_lifecycle_service.dart b/lib/core/lifecycle/app_lifecycle_service.dart new file mode 100644 index 00000000..a8777281 --- /dev/null +++ b/lib/core/lifecycle/app_lifecycle_service.dart @@ -0,0 +1,91 @@ +import 'dart:async'; + +import 'package:flutter/widgets.dart'; + +/// Turns the OS lifecycle into two calls, and nothing else. +/// +/// The process is frozen wholesale while the app is in the background: the +/// relay sockets die, and whatever the daemon or a peer sent meanwhile is only +/// on the relays. Coming back is therefore one operation — `resync()` in +/// Rust, then every notifier re-reads its truth from the bridge — and this +/// observer only decides *when* that operation runs (issue #308, +/// docs/PUSH_NOTIFICATIONS.md §10). The work itself lives in [onResume], a +/// plain async function tests call directly. +/// +/// Two rules keep it from firing for nothing: +/// +/// - **A latch.** Only a real `paused` (or `hidden`, which is what web and +/// desktop deliver) arms the next `resumed`. The `inactive → resumed` flap +/// of a permission dialog, a share sheet or an app-switcher peek never +/// suspended anything and is ignored. +/// - **A debounce.** A resume that flaps within [debounce] runs once, at the +/// end; the Rust side coalesces concurrent passes too, but there is no +/// reason to ask twice. +/// +/// Every platform gets the same path: there is no `dart:io` platform gate to +/// fake in a host test, and a desktop window that is never hidden simply +/// never arms the latch. +class AppLifecycleService with WidgetsBindingObserver { + AppLifecycleService({ + required this.onResume, + this.onPause, + this.debounce = const Duration(milliseconds: 300), + }); + + /// Runs after a suspension ends. Exceptions are caught and logged: a + /// failed resync must never take the UI down. + final Future Function() onResume; + + /// Runs when a suspension starts. Optional: on this architecture the work + /// belongs on the resume side. + final void Function()? onPause; + + final Duration debounce; + + bool _suspended = false; + Timer? _pending; + + /// Whether a suspension is in progress, i.e. the next `resumed` will fire + /// [onResume]. Exposed for tests. + @visibleForTesting + bool get suspended => _suspended; + + void attach() => WidgetsBinding.instance.addObserver(this); + + void detach() { + WidgetsBinding.instance.removeObserver(this); + _pending?.cancel(); + _pending = null; + } + + @override + void didChangeAppLifecycleState(AppLifecycleState state) { + switch (state) { + case AppLifecycleState.paused: + case AppLifecycleState.hidden: + if (!_suspended) { + _suspended = true; + _pending?.cancel(); + _pending = null; + onPause?.call(); + } + case AppLifecycleState.resumed: + if (!_suspended) return; + _suspended = false; + _pending?.cancel(); + _pending = Timer(debounce, _fireResume); + case AppLifecycleState.inactive: + case AppLifecycleState.detached: + break; + } + } + + void _fireResume() { + _pending = null; + unawaited( + onResume().catchError((Object e, StackTrace st) { + debugPrint('[lifecycle] resume handler failed: $e\n$st'); + }), + ); + } +} diff --git a/lib/core/lifecycle/resume_resync.dart b/lib/core/lifecycle/resume_resync.dart new file mode 100644 index 00000000..02ae7ca0 --- /dev/null +++ b/lib/core/lifecycle/resume_resync.dart @@ -0,0 +1,66 @@ +import 'package:flutter/foundation.dart'; +import 'package:flutter_riverpod/flutter_riverpod.dart'; + +import 'package:mostro/features/chat/providers/chat_providers.dart'; +import 'package:mostro/features/disputes/providers/disputes_providers.dart'; +import 'package:mostro/features/notifications/providers/notifications_provider.dart'; +import 'package:mostro/features/trades/providers/trades_providers.dart'; +import 'package:mostro/src/rust/api/nostr.dart' as nostr_api; +import 'package:mostro/src/rust/api/types.dart'; + +/// One hydration hook: re-read a feature's protocol state from the bridge. +/// +/// **Streams are for live updates; queries are for hydration. Resume always +/// re-hydrates.** A notifier fed only by incremental events silently loses +/// everything that happened while the process was suspended, and a +/// hand-maintained list of per-feature refresh calls is exactly what let v1 +/// ship the dispute-chat bug (MostroP2P/mobile#675). So a feature exposes the +/// same code path it uses at cold start, and resume runs all of them. +typedef Hydrator = Future Function(ProviderContainer container); + +/// The resume routine: `resync()` in Rust, then every hydrator, in order. +/// +/// Pure enough to test with fakes: the bridge call and the hook list are +/// injected. Each step is isolated — a hydrator that throws is logged and +/// the next one still runs, and a failed resync still hydrates (whatever is +/// already on disk is newer than what the notifiers hold). +class ResumeResync { + ResumeResync({ + required this.container, + Future Function()? resync, + List? hydrators, + }) : _resync = resync ?? nostr_api.resync, + _hydrators = hydrators ?? defaultHydrators; + + final ProviderContainer container; + final Future Function() _resync; + final List _hydrators; + + Future run() async { + try { + final outcome = await _resync(); + debugPrint( + '[lifecycle] resync: online=${outcome.online} ' + 'flushed=${outcome.flushed} coalesced=${outcome.coalesced}', + ); + } catch (e) { + debugPrint('[lifecycle] resync failed: $e'); + } + for (final hydrate in _hydrators) { + try { + await hydrate(container); + } catch (e, st) { + debugPrint('[lifecycle] hydration failed: $e\n$st'); + } + } + } +} + +/// Every notifier that holds protocol-derived state, in dependency order: +/// trades first, since chat rooms and disputes are read off the trade list. +final List defaultHydrators = [ + hydrateTrades, + hydrateChatRooms, + hydrateDisputes, + hydrateNotifications, +]; diff --git a/lib/features/chat/providers/chat_providers.dart b/lib/features/chat/providers/chat_providers.dart index d678ef27..af7e9930 100644 --- a/lib/features/chat/providers/chat_providers.dart +++ b/lib/features/chat/providers/chat_providers.dart @@ -23,14 +23,14 @@ class ChatRoomState { this.lastMessageIsOwn = false, this.lastMessageAt = 0, this.unreadCount = 0, - }) : assert( - peerIconIndex >= 0 && peerIconIndex <= 36, - 'peerIconIndex must be 0–36, got $peerIconIndex', - ), - assert( - peerColorHue >= 0 && peerColorHue <= 359, - 'peerColorHue must be 0–359, got $peerColorHue', - ); + }) : assert( + peerIconIndex >= 0 && peerIconIndex <= 36, + 'peerIconIndex must be 0–36, got $peerIconIndex', + ), + assert( + peerColorHue >= 0 && peerColorHue <= 359, + 'peerColorHue must be 0–359, got $peerColorHue', + ); /// The trade / order ID that identifies this chat room. final String orderId; @@ -85,9 +85,10 @@ class ChatRoomState { peerIconIndex: peerIconIndex ?? this.peerIconIndex, peerColorHue: peerColorHue ?? this.peerColorHue, isSelling: isSelling ?? this.isSelling, - lastMessage: identical(lastMessage, _noChange) - ? this.lastMessage - : lastMessage as String?, + lastMessage: + identical(lastMessage, _noChange) + ? this.lastMessage + : lastMessage as String?, lastMessageIsOwn: lastMessageIsOwn ?? this.lastMessageIsOwn, lastMessageAt: lastMessageAt ?? this.lastMessageAt, unreadCount: unreadCount ?? this.unreadCount, @@ -180,8 +181,8 @@ class ChatRoomsNotifier extends StateNotifier> { /// Sorted order is provided by [sortedChatRoomsProvider]. final chatRoomsNotifierProvider = StateNotifierProvider>( - (_) => ChatRoomsNotifier(), -); + (_) => ChatRoomsNotifier(), + ); /// Chat rooms sorted by [ChatRoomState.lastMessageAt] descending (newest first). final sortedChatRoomsProvider = Provider>((ref) { @@ -202,8 +203,7 @@ final chatCountProvider = Provider((ref) { /// Maps orderId → last-read unix timestamp (seconds). /// /// In-memory only. Sembast persistence deferred to a future phase. -final chatReadStatusProvider = - StateProvider>((_) => const {}); +final chatReadStatusProvider = StateProvider>((_) => const {}); // ── Trade → ChatRoom bridge ─────────────────────────────────────────────────── @@ -211,9 +211,7 @@ final chatReadStatusProvider = /// /// Returns `null` when [TradeInfo.counterpartyPubkey] is empty — meaning the /// peer identity has not been exchanged yet and there is no chat room to show. -Future tradeInfoToChatRoom( - rust_types.TradeInfo trade, -) async { +Future tradeInfoToChatRoom(rust_types.TradeInfo trade) async { final peerPubkey = trade.counterpartyPubkey; if (peerPubkey.isEmpty) return null; @@ -238,9 +236,10 @@ Future tradeInfoToChatRoom( // Peer messages only: the room preview and unread badge describe the // buyer<->seller conversation, not the dispute channel that shares the // order key (PR #254 review). - msgs = (await messages_api.getMessages(tradeId: trade.order.id)) - .where((m) => m.messageType == rust_types.MessageType.peer) - .toList(); + msgs = + (await messages_api.getMessages( + tradeId: trade.order.id, + )).where((m) => m.messageType == rust_types.MessageType.peer).toList(); } catch (_) { msgs = const []; } @@ -268,16 +267,15 @@ Future tradeInfoToChatRoom( /// FutureProvider that converts the full trade list into [ChatRoomState]s. /// /// Only trades with a known peer pubkey are included. Sorted newest-message first. -final chatRoomsFromTradesProvider = - FutureProvider>((ref) async { +final chatRoomsFromTradesProvider = FutureProvider>(( + ref, +) async { final trades = await ref.watch(rawTradesProvider.future); // Resolve all rooms concurrently — each call does async I/O (NymIdentity // lookup + message history fetch) so parallel execution is significantly // faster than the previous sequential await loop. - final results = await Future.wait( - trades.map((t) => tradeInfoToChatRoom(t)), - ); + final results = await Future.wait(trades.map((t) => tradeInfoToChatRoom(t))); final rooms = results.whereType().toList(); rooms.sort((a, b) => b.lastMessageAt.compareTo(a.lastMessageAt)); @@ -288,26 +286,40 @@ final chatRoomsFromTradesProvider = /// /// Used by [ChatRoomScreen] to update its local list in real-time without /// polling. -final incomingMessageProvider = - StreamProvider.autoDispose.family( - (ref, tradeId) async* { - final stream = await messages_api.onNewMessage(tradeId: tradeId); - while (true) { - final msg = await stream.next(); - if (msg == null) break; - yield msg; - } - }, -); +final incomingMessageProvider = StreamProvider.autoDispose + .family((ref, tradeId) async* { + final stream = await messages_api.onNewMessage(tradeId: tradeId); + while (true) { + final msg = await stream.next(); + if (msg == null) break; + yield msg; + } + }); /// FutureProvider that loads the full message history for a trade once. /// /// [ChatRoomScreen] seeds its local state from this, then appends live /// updates via [incomingMessageProvider]. -final messageHistoryProvider = - FutureProvider.autoDispose.family, String>( - // Peer messages only — this feeds the peer chat room (PR #254 review). - (ref, tradeId) async => (await messages_api.getMessages(tradeId: tradeId)) - .where((m) => m.messageType == rust_types.MessageType.peer) - .toList(), -); +final messageHistoryProvider = FutureProvider.autoDispose + .family, String>( + // Peer messages only — this feeds the peer chat room (PR #254 review). + (ref, tradeId) async => + (await messages_api.getMessages(tradeId: tradeId)) + .where((m) => m.messageType == rust_types.MessageType.peer) + .toList(), + ); + +// ── Hydration (resume) ──────────────────────────────────────────────────────── + +/// Rebuild the chat rooms from the trade list and the persisted messages — +/// the same path `ChatRoomsScreen` runs on init — and upsert them, so a room +/// added live while the fetch ran survives. Assumes the trade list was +/// hydrated first (`defaultHydrators` orders it so). +Future hydrateChatRooms(ProviderContainer container) async { + container.invalidate(chatRoomsFromTradesProvider); + final rooms = await container.read(chatRoomsFromTradesProvider.future); + final notifier = container.read(chatRoomsNotifierProvider.notifier); + for (final room in rooms) { + notifier.upsertRoom(room); + } +} diff --git a/lib/features/disputes/providers/disputes_providers.dart b/lib/features/disputes/providers/disputes_providers.dart index 9f982c5a..0eaaa96f 100644 --- a/lib/features/disputes/providers/disputes_providers.dart +++ b/lib/features/disputes/providers/disputes_providers.dart @@ -1,6 +1,11 @@ import 'package:flutter/foundation.dart'; import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:mostro/features/trades/providers/trades_providers.dart'; +import 'package:mostro/shared/utils/platform_int64.dart'; +import 'package:mostro/src/rust/api/disputes.dart' as disputes_api; +import 'package:mostro/src/rust/api/types.dart' as rust_types; + // ── Dispute models ──────────────────────────────────────────────────────────── /// Dispute lifecycle status, matching the Rust `DisputeStatus` enum. @@ -61,8 +66,8 @@ class DisputeItem { this.peerIconIndex = 0, this.peerColorHue = 180, this.isSelling = false, - }) : assert(peerIconIndex >= 0 && peerIconIndex <= 36), - assert(peerColorHue >= 0 && peerColorHue <= 359); + }) : assert(peerIconIndex >= 0 && peerIconIndex <= 36), + assert(peerColorHue >= 0 && peerColorHue <= 359); final String id; final String tradeId; @@ -154,8 +159,8 @@ class DisputeNotifier extends StateNotifier> { /// Empty until bridge events are integrated (Phase 12+). final disputeNotifierProvider = StateNotifierProvider>( - (_) => DisputeNotifier(), -); + (_) => DisputeNotifier(), + ); /// All disputes sorted newest-first. /// @@ -179,19 +184,90 @@ final disputeUnreadCountProvider = Provider((ref) { }); /// Look up a single dispute by its ID. -final disputeByIdProvider = - Provider.family((ref, id) { - return ref.watch(disputeNotifierProvider).where((d) => d.id == id).firstOrNull; +final disputeByIdProvider = Provider.family((ref, id) { + return ref + .watch(disputeNotifierProvider) + .where((d) => d.id == id) + .firstOrNull; }); /// Look up a dispute by its associated trade ID. /// /// Used by [TradeDetailScreen] to resolve the correct `disputeId` before /// navigating to [DisputeChatScreen]. -final disputeByTradeIdProvider = - Provider.family((ref, tradeId) { +final disputeByTradeIdProvider = Provider.family(( + ref, + tradeId, +) { return ref .watch(disputeNotifierProvider) .where((d) => d.tradeId == tradeId) .firstOrNull; }); + +// ── Hydration (resume) ──────────────────────────────────────────────────────── + +/// The bridge's dispute record, as the list shows it. The peer's handle and +/// side are the row's own concern (they come from the trade, not the +/// dispute) and stay whatever the UI already set. +DisputeItem disputeItemFromRust(rust_types.Dispute dispute) => DisputeItem( + id: dispute.id, + tradeId: dispute.tradeId, + status: switch (dispute.status) { + rust_types.DisputeStatus.open => DisputeStatus.open, + rust_types.DisputeStatus.inReview => DisputeStatus.inReview, + rust_types.DisputeStatus.resolved => DisputeStatus.resolved, + }, + initiatedByMe: dispute.initiatedByMe, + openedAt: platformInt64ToInt(dispute.openedAt), + reason: dispute.reason, + adminPubkey: dispute.adminPubkey, + resolution: switch (dispute.resolution) { + null => null, + rust_types.DisputeResolution.fundsToBuyer => DisputeResolution.fundsToBuyer, + rust_types.DisputeResolution.fundsToSeller => + DisputeResolution.fundsToSeller, + rust_types.DisputeResolution.cooperativeCancel => + DisputeResolution.cooperativeCancel, + }, + resolvedAt: + dispute.resolvedAt == null + ? null + : platformInt64ToInt(dispute.resolvedAt), + isRead: dispute.isRead, +); + +/// Trade statuses that can carry a dispute record on the bridge. +const _disputedStatuses = { + rust_types.OrderStatus.dispute, + rust_types.OrderStatus.settledByAdmin, + rust_types.OrderStatus.canceledByAdmin, + rust_types.OrderStatus.completedByAdmin, +}; + +/// Re-read every dispute the bridge knows for the trades that can have one +/// and upsert it: the read flag the UI manages survives (`upsert` keeps it), +/// and a dispute opened by the peer while the process was suspended appears +/// without a restart — the shape of the v1 bug this exists to prevent +/// (MostroP2P/mobile#675). Assumes the trade list was hydrated first. +Future hydrateDisputes( + ProviderContainer container, { + Future Function({required String tradeId})? getDispute, +}) async { + final lookup = getDispute ?? disputes_api.getDispute; + final trades = await container.read(rawTradesProvider.future); + final notifier = container.read(disputeNotifierProvider.notifier); + for (final trade in trades) { + if (!_disputedStatuses.contains(trade.order.status)) continue; + final rust_types.Dispute? dispute; + try { + dispute = await lookup(tradeId: trade.order.id); + } catch (e) { + debugPrint( + '[disputes] hydrate: getDispute(${trade.order.id}) failed: $e', + ); + continue; + } + if (dispute != null) notifier.upsert(disputeItemFromRust(dispute)); + } +} diff --git a/lib/features/notifications/providers/notifications_provider.dart b/lib/features/notifications/providers/notifications_provider.dart index 01a75d0b..b19b2a49 100644 --- a/lib/features/notifications/providers/notifications_provider.dart +++ b/lib/features/notifications/providers/notifications_provider.dart @@ -18,10 +18,11 @@ import 'package:mostro/features/notifications/providers/sembast_factory_io.dart' class SembastNotificationsStore { SembastNotificationsStore({DatabaseFactory? factory, String? path}) - : _factoryOverride = factory, - _pathOverride = path; + : _factoryOverride = factory, + _pathOverride = path; static const _dbName = 'notifications.db'; + /// The int-keyed store this feature shipped with. static const _legacyStoreName = 'notifications'; @@ -37,6 +38,7 @@ class SembastNotificationsStore { Database? _db; Completer? _opening; + /// Keyed by notification id. Before this, the store used auto-incrementing /// integer keys and every write looked its record up with a `Finder` — a /// full-store scan per save, and O(n) scans for an O(n) bulk update. @@ -103,7 +105,9 @@ class SembastNotificationsStore { final db = await _open(); final records = await _store.find(db); return records - .map((r) => NotificationModel.fromJson(Map.from(r.value))) + .map( + (r) => NotificationModel.fromJson(Map.from(r.value)), + ) .toList(); } @@ -126,7 +130,10 @@ class SembastNotificationsStore { }); } - Future _upsert(DatabaseClient client, NotificationModel notification) async { + Future _upsert( + DatabaseClient client, + NotificationModel notification, + ) async { final json = Map.from(notification.toJson()) ..removeWhere((_, v) => v == null); await _store.record(notification.id).put(client, json); @@ -142,7 +149,8 @@ class SembastNotificationsStore { Future saveIfUnprocessed(NotificationModel notification) async { final db = await _open(); return db.transaction((txn) async { - final already = await _processed.record(notification.id).get(txn) ?? false; + final already = + await _processed.record(notification.id).get(txn) ?? false; if (already) return false; await _processed.record(notification.id).put(txn, true); await _upsert(txn, notification); @@ -185,14 +193,14 @@ final sembastNotificationsStoreProvider = Provider( /// one record and preserves the user's read/delete state. Single source of /// truth for the list, the bell, and every producer (listeners and push path). final notificationsProvider = - StateNotifierProvider>( - (ref) { - final store = ref.watch(sembastNotificationsStoreProvider); - final notifier = NotificationsNotifier(store: store); - notifier.loadInitialData(); - return notifier; - }, -); + StateNotifierProvider>(( + ref, + ) { + final store = ref.watch(sembastNotificationsStoreProvider); + final notifier = NotificationsNotifier(store: store); + notifier.loadInitialData(); + return notifier; + }); /// Count of unread notifications. final unreadNotificationCountProvider = Provider( @@ -221,8 +229,9 @@ class NotificationsNotifier extends StateNotifier> { for (final n in state) { byId[n.id] = n; } - state = byId.values.toList() - ..sort((a, b) => b.timestamp.compareTo(a.timestamp)); + state = + byId.values.toList() + ..sort((a, b) => b.timestamp.compareTo(a.timestamp)); } catch (e) { debugPrint('NotificationsNotifier: failed to load from Sembast: $e'); } @@ -339,4 +348,12 @@ class NotificationsNotifier extends StateNotifier> { ); add(notification); } -} \ No newline at end of file +} + +// ── Hydration (resume) ──────────────────────────────────────────────────────── + +/// Merge the persisted notifications back into state — the cold-start load, +/// which keeps whatever was added live. The cards themselves come from the +/// Rust streams the resync replays, deduplicated by `addIfNew`. +Future hydrateNotifications(ProviderContainer container) => + container.read(notificationsProvider.notifier).loadInitialData(); diff --git a/lib/features/trades/providers/trades_providers.dart b/lib/features/trades/providers/trades_providers.dart index 1844412d..37633731 100644 --- a/lib/features/trades/providers/trades_providers.dart +++ b/lib/features/trades/providers/trades_providers.dart @@ -108,3 +108,13 @@ void refreshTrades(WidgetRef ref) => ref.invalidate(rawTradesProvider); final orderBookNotificationCountProvider = Provider( (ref) => ref.watch(needsActionCountProvider), ); + +// ── Hydration (resume) ──────────────────────────────────────────────────────── + +/// Re-read the trade list from the bridge. The cold-start path is the same +/// query, so this only drops the cache; every watcher refetches. Called by +/// the resume routine after `resync()` (lib/core/lifecycle/resume_resync.dart). +Future hydrateTrades(ProviderContainer container) async { + container.invalidate(rawTradesProvider); + await container.read(rawTradesProvider.future); +} diff --git a/test/core/lifecycle/app_lifecycle_service_test.dart b/test/core/lifecycle/app_lifecycle_service_test.dart new file mode 100644 index 00000000..b27d61df --- /dev/null +++ b/test/core/lifecycle/app_lifecycle_service_test.dart @@ -0,0 +1,155 @@ +import 'package:flutter/widgets.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:mostro/core/lifecycle/app_lifecycle_service.dart'; + +void main() { + const debounce = Duration(milliseconds: 50); + + late int resumes; + late int pauses; + late AppLifecycleService service; + + setUp(() { + resumes = 0; + pauses = 0; + }); + + AppLifecycleService attach( + WidgetTester tester, { + Future Function()? onResume, + }) { + service = AppLifecycleService( + onResume: onResume ?? () async => resumes++, + onPause: () => pauses++, + debounce: debounce, + )..attach(); + addTearDown(service.detach); + return service; + } + + Future deliver( + WidgetTester tester, + List states, + ) async { + for (final state in states) { + tester.binding.handleAppLifecycleStateChanged(state); + } + await tester.pump(debounce * 2); + } + + testWidgets('a paused → resumed suspension fires the resume handler once', ( + tester, + ) async { + // Arrange + attach(tester); + + // Act + await deliver(tester, [ + AppLifecycleState.inactive, + AppLifecycleState.paused, + AppLifecycleState.resumed, + ]); + + // Assert + expect(pauses, 1); + expect(resumes, 1); + }); + + testWidgets('hidden arms the latch too — what web and desktop deliver', ( + tester, + ) async { + attach(tester); + + await deliver(tester, [ + AppLifecycleState.hidden, + AppLifecycleState.resumed, + ]); + + expect(resumes, 1); + }); + + testWidgets( + 'an inactive → resumed flap never suspended anything and is ignored', + (tester) async { + // A permission dialog, a share sheet, an app-switcher peek. + attach(tester); + + await deliver(tester, [ + AppLifecycleState.inactive, + AppLifecycleState.resumed, + AppLifecycleState.inactive, + AppLifecycleState.resumed, + ]); + + expect(pauses, 0); + expect(resumes, 0); + }, + ); + + testWidgets('a resume that flaps within the debounce runs once, at the end', ( + tester, + ) async { + attach(tester); + + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.paused); + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.resumed); + await tester.pump(const Duration(milliseconds: 10)); + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.paused); + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.resumed); + await tester.pump(const Duration(milliseconds: 10)); + expect(resumes, 0, reason: 'still inside the debounce'); + await tester.pump(debounce * 2); + + expect(resumes, 1); + expect(pauses, 2); + }); + + testWidgets('a second paused while suspended does not re-arm or re-pause', ( + tester, + ) async { + attach(tester); + + await deliver(tester, [ + AppLifecycleState.paused, + AppLifecycleState.paused, + AppLifecycleState.resumed, + ]); + + expect(pauses, 1); + expect(resumes, 1); + }); + + testWidgets('a throwing resume handler is logged, not rethrown', ( + tester, + ) async { + attach(tester, onResume: () async => throw StateError('boom')); + + await deliver(tester, [ + AppLifecycleState.paused, + AppLifecycleState.resumed, + ]); + + expect(tester.takeException(), isNull); + }); + + testWidgets('detach stops a pending resume', (tester) async { + final s = attach(tester); + + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.paused); + tester.binding.handleAppLifecycleStateChanged(AppLifecycleState.resumed); + s.detach(); + await tester.pump(debounce * 2); + + expect(resumes, 0); + }); + + test('the service exposes its latch for tests', () { + final s = AppLifecycleService(onResume: () async {}); + expect(s.suspended, isFalse); + s.didChangeAppLifecycleState(AppLifecycleState.paused); + expect(s.suspended, isTrue); + s.didChangeAppLifecycleState(AppLifecycleState.resumed); + expect(s.suspended, isFalse); + s.detach(); + }); +} diff --git a/test/core/lifecycle/hydration_regression_test.dart b/test/core/lifecycle/hydration_regression_test.dart new file mode 100644 index 00000000..33b6871c --- /dev/null +++ b/test/core/lifecycle/hydration_regression_test.dart @@ -0,0 +1,172 @@ +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:mostro/core/lifecycle/resume_resync.dart'; +import 'package:mostro/features/chat/providers/chat_providers.dart'; +import 'package:mostro/features/disputes/providers/disputes_providers.dart'; +import 'package:mostro/features/trades/providers/trades_providers.dart'; +import 'package:mostro/src/rust/api/types.dart' as rust_types; + +import '../../support/fake_trades.dart'; + +/// The v1 dispute-chat bug, encoded (MostroP2P/mobile#675, issue #308): a +/// dispute the peer opened and a trade that moved while the process was +/// suspended must be visible after resume, without a restart. The bridge is +/// a mutable fake: what it holds after "suspension" is what the notifiers +/// must show after the resume routine ran. +void main() { + late List bridgeTrades; + late Map bridgeDisputes; + late ProviderContainer container; + + ChatRoomState room(String orderId) => ChatRoomState( + orderId: orderId, + peerPubkey: 'peer-$orderId', + peerHandle: 'Peer', + peerIconIndex: 0, + peerColorHue: 0, + isSelling: false, + ); + + setUp(() { + bridgeTrades = [fakeTrade(id: 't1', status: rust_types.OrderStatus.active)]; + bridgeDisputes = {}; + container = ProviderContainer( + overrides: [ + rawTradesProvider.overrideWith((ref) async => bridgeTrades.toList()), + chatRoomsFromTradesProvider.overrideWith((ref) async { + final trades = await ref.watch(rawTradesProvider.future); + return [for (final t in trades) room(t.order.id)]; + }), + ], + ); + addTearDown(container.dispose); + }); + + Future lookup({required String tradeId}) async => + bridgeDisputes[tradeId]; + + ResumeResync routine() => ResumeResync( + container: container, + resync: + () async => const rust_types.ResyncOutcome( + online: true, + flushed: 0, + coalesced: false, + ), + hydrators: [ + hydrateTrades, + hydrateChatRooms, + (c) => hydrateDisputes(c, getDispute: lookup), + ], + ); + + test( + 'a trade that moved while suspended shows its new status after resume', + () async { + // Arrange — cold start read the trade as active. + expect( + (await container.read(rawTradesProvider.future)).single.order.status, + rust_types.OrderStatus.active, + ); + + // "Suspended": the daemon released the escrow. + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.success), + ]; + + // Act + await routine().run(); + + // Assert + expect( + (await container.read(rawTradesProvider.future)).single.order.status, + rust_types.OrderStatus.success, + ); + }, + ); + + test( + 'a dispute the peer opened while suspended is listed after resume', + () async { + // Arrange — nothing disputed at cold start. + container.read(disputeNotifierProvider); + expect(container.read(disputeNotifierProvider), isEmpty); + + // "Suspended": the peer opened a dispute; the trade row moved with it. + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.dispute), + ]; + bridgeDisputes['order-t1'] = const rust_types.Dispute( + id: 'd1', + tradeId: 'order-t1', + status: rust_types.DisputeStatus.inReview, + initiatedByMe: false, + adminPubkey: 'solver', + openedAt: 1234, + isRead: false, + ); + + // Act + await routine().run(); + + // Assert + final listed = container.read(disputeNotifierProvider).single; + expect(listed.id, 'd1'); + expect(listed.status, DisputeStatus.inReview); + expect(listed.adminPubkey, 'solver'); + expect(listed.initiatedByMe, isFalse); + }, + ); + + test('re-hydrating a dispute keeps the read flag the user set', () async { + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.dispute), + ]; + bridgeDisputes['order-t1'] = const rust_types.Dispute( + id: 'd1', + tradeId: 'order-t1', + status: rust_types.DisputeStatus.open, + initiatedByMe: true, + openedAt: 1234, + isRead: false, + ); + await routine().run(); + container.read(disputeNotifierProvider.notifier).markRead('d1'); + + await routine().run(); + + expect(container.read(disputeNotifierProvider).single.isRead, isTrue); + }); + + test( + 'a chat room for a trade that revealed its peer while suspended appears', + () async { + expect(container.read(chatRoomsNotifierProvider), isEmpty); + + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.active), + fakeTrade(id: 't2', status: rust_types.OrderStatus.active), + ]; + + await routine().run(); + + expect( + container.read(chatRoomsNotifierProvider).map((r) => r.orderId), + unorderedEquals(['order-t1', 'order-t2']), + ); + }, + ); + + test( + 'running the routine twice over an unchanged bridge changes nothing', + () async { + await routine().run(); + final before = container.read(chatRoomsNotifierProvider); + + await routine().run(); + + expect(container.read(chatRoomsNotifierProvider).length, before.length); + expect((await container.read(rawTradesProvider.future)).length, 1); + }, + ); +} diff --git a/test/core/lifecycle/resume_resync_test.dart b/test/core/lifecycle/resume_resync_test.dart new file mode 100644 index 00000000..ee12cbcd --- /dev/null +++ b/test/core/lifecycle/resume_resync_test.dart @@ -0,0 +1,70 @@ +import 'package:flutter_riverpod/flutter_riverpod.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:mostro/core/lifecycle/resume_resync.dart'; +import 'package:mostro/src/rust/api/types.dart'; + +void main() { + const ok = ResyncOutcome(online: true, flushed: 0, coalesced: false); + + late ProviderContainer container; + late List calls; + + setUp(() { + container = ProviderContainer(); + addTearDown(container.dispose); + calls = []; + }); + + Hydrator hook(String name, {bool throws = false}) => (c) async { + expect(identical(c, container), isTrue); + calls.add(name); + if (throws) throw StateError('$name failed'); + }; + + test('resync runs first, then every hydrator in order', () async { + // Arrange + final routine = ResumeResync( + container: container, + resync: () async { + calls.add('resync'); + return ok; + }, + hydrators: [hook('trades'), hook('chat'), hook('disputes')], + ); + + // Act + await routine.run(); + + // Assert + expect(calls, ['resync', 'trades', 'chat', 'disputes']); + }); + + test('a hydrator that throws does not stop the ones after it', () async { + final routine = ResumeResync( + container: container, + resync: () async => ok, + hydrators: [hook('trades'), hook('chat', throws: true), hook('disputes')], + ); + + await routine.run(); + + expect(calls, ['trades', 'chat', 'disputes']); + }); + + test('a failed resync still hydrates — disk is newer than memory', () async { + final routine = ResumeResync( + container: container, + resync: () async => throw StateError('no bridge'), + hydrators: [hook('trades')], + ); + + await expectLater(routine.run(), completes); + + expect(calls, ['trades']); + }); + + test('the default hook list covers every protocol-state notifier', () { + // Trades first: chat rooms and disputes are read off the trade list. + expect(defaultHydrators, hasLength(4)); + }); +} From c3d0b38c68b111658a84a3aea862f7912cb7cd3d Mon Sep 17 00:00:00 2001 From: grunch Date: Sun, 13 Sep 2026 15:54:43 -0300 Subject: [PATCH 2/2] fix(push): review round 1 on PR-0b (Codex) Hydrate twice: once for what is on disk, once more when the relay replay settles (quiet after the last trade update, capped). Query every trade that can own a dispute, not only rows whose status says so. Chat rooms merge by timestamp so a snapshot never rolls back a live room. A notice deleted while a resume load is reading stays deleted. --- lib/core/lifecycle/resume_resync.dart | 89 +++++++++++++- .../chat/providers/chat_providers.dart | 21 +++- .../chat/screens/chat_rooms_screen.dart | 6 +- .../providers/disputes_providers.dart | 20 ++-- .../providers/notifications_provider.dart | 33 +++++- .../lifecycle/hydration_regression_test.dart | 99 +++++++++++++++- test/core/lifecycle/resume_resync_test.dart | 106 ++++++++++++++++- .../chat_rooms_upsert_if_newer_test.dart | 55 +++++++++ .../notifications_resume_load_test.dart | 110 ++++++++++++++++++ 9 files changed, 510 insertions(+), 29 deletions(-) create mode 100644 test/features/chat/providers/chat_rooms_upsert_if_newer_test.dart create mode 100644 test/features/notifications/notifications_resume_load_test.dart diff --git a/lib/core/lifecycle/resume_resync.dart b/lib/core/lifecycle/resume_resync.dart index 02ae7ca0..0feae82a 100644 --- a/lib/core/lifecycle/resume_resync.dart +++ b/lib/core/lifecycle/resume_resync.dart @@ -1,3 +1,5 @@ +import 'dart:async'; + import 'package:flutter/foundation.dart'; import 'package:flutter_riverpod/flutter_riverpod.dart'; @@ -6,6 +8,7 @@ import 'package:mostro/features/disputes/providers/disputes_providers.dart'; import 'package:mostro/features/notifications/providers/notifications_provider.dart'; import 'package:mostro/features/trades/providers/trades_providers.dart'; import 'package:mostro/src/rust/api/nostr.dart' as nostr_api; +import 'package:mostro/src/rust/api/orders.dart' as orders_api; import 'package:mostro/src/rust/api/types.dart'; /// One hydration hook: re-read a feature's protocol state from the bridge. @@ -18,24 +21,50 @@ import 'package:mostro/src/rust/api/types.dart'; /// same code path it uses at cold start, and resume runs all of them. typedef Hydrator = Future Function(ProviderContainer container); -/// The resume routine: `resync()` in Rust, then every hydrator, in order. +/// The resume routine: `resync()` in Rust, then every hydrator, in order — +/// twice. +/// +/// `resync()` returns once the subscriptions are re-issued, not once the +/// relays have replayed what was missed: the events the app was suspended +/// for arrive over the next seconds and are written to the database as they +/// land. So the hydrators run **immediately** (what is on disk is already +/// newer than what the notifiers hold) and **once more when the replay +/// settles**: after [settleQuiet] without a trade update following the last +/// one, or after [settleMax] whatever is still arriving. A trade update +/// already refreshes the trade list on its own; the second pass is for the +/// notifiers derived from it — chat rooms, disputes — which nothing else +/// re-reads. /// -/// Pure enough to test with fakes: the bridge call and the hook list are -/// injected. Each step is isolated — a hydrator that throws is logged and -/// the next one still runs, and a failed resync still hydrates (whatever is -/// already on disk is newer than what the notifiers hold). +/// Pure enough to test with fakes: the bridge call, the hook list and the +/// update stream are injected. Each step is isolated — a hydrator that throws +/// is logged and the next one still runs, and a failed resync still hydrates. class ResumeResync { ResumeResync({ required this.container, Future Function()? resync, List? hydrators, + Stream Function()? updates, + this.settleQuiet = const Duration(milliseconds: 1500), + this.settleMax = const Duration(seconds: 10), }) : _resync = resync ?? nostr_api.resync, - _hydrators = hydrators ?? defaultHydrators; + _hydrators = hydrators ?? defaultHydrators, + _updates = updates ?? _bridgeTradeUpdates; final ProviderContainer container; final Future Function() _resync; final List _hydrators; + final Stream Function() _updates; + + /// Replay is considered settled this long after its last trade update. + final Duration settleQuiet; + /// The second pass runs at the latest this long after the first, so a + /// stream that never goes quiet (or never emits) still gets one. + final Duration settleMax; + + /// Runs the routine to completion: resync, hydrate, wait for the replay to + /// settle, hydrate again. Awaiting it is optional; the lifecycle service + /// does not. Future run() async { try { final outcome = await _resync(); @@ -46,6 +75,12 @@ class ResumeResync { } catch (e) { debugPrint('[lifecycle] resync failed: $e'); } + await _hydrateAll(); + await _waitForReplayToSettle(); + await _hydrateAll(); + } + + Future _hydrateAll() async { for (final hydrate in _hydrators) { try { await hydrate(container); @@ -54,6 +89,48 @@ class ResumeResync { } } } + + /// Completes [settleQuiet] after the last update, or at [settleMax]. + Future _waitForReplayToSettle() async { + final settled = Completer(); + void done() { + if (!settled.isCompleted) settled.complete(); + } + + Timer quiet = Timer(settleQuiet, done); + final cap = Timer(settleMax, done); + StreamSubscription? sub; + try { + sub = _updates().listen( + (_) { + quiet.cancel(); + quiet = Timer(settleQuiet, done); + }, + onError: (Object e) { + debugPrint('[lifecycle] trade update stream failed: $e'); + done(); + }, + onDone: done, + ); + } catch (e) { + debugPrint('[lifecycle] trade update stream unavailable: $e'); + done(); + } + await settled.future; + quiet.cancel(); + cap.cancel(); + await sub?.cancel(); + } +} + +/// The bridge's trade updates, as a stream, for the settle window. +Stream _bridgeTradeUpdates() async* { + final stream = await orders_api.onTradeUpdated(); + while (true) { + final update = await stream.next(); + if (update == null) break; + yield update; + } } /// Every notifier that holds protocol-derived state, in dependency order: diff --git a/lib/features/chat/providers/chat_providers.dart b/lib/features/chat/providers/chat_providers.dart index af7e9930..ccdb28f3 100644 --- a/lib/features/chat/providers/chat_providers.dart +++ b/lib/features/chat/providers/chat_providers.dart @@ -132,6 +132,21 @@ class ChatRoomsNotifier extends StateNotifier> { } } + /// Upsert a room read from storage without rolling back a live one. + /// + /// A snapshot is built room by room, and a message folded into a room after + /// its own history was read but before the whole list is ready is newer + /// than what the snapshot holds for it. Such a room keeps its live preview, + /// time and unread count; a snapshot at or past the live room's time, and a + /// room the list does not have yet, are taken as read. + void upsertIfNewer(ChatRoomState room) { + final existing = state.indexWhere((r) => r.orderId == room.orderId); + if (existing >= 0 && state[existing].lastMessageAt > room.lastMessageAt) { + return; + } + upsertRoom(room); + } + /// Folds a message that arrived for [orderId] into its room's preview, /// time and unread count, so the list stays live without a reload. /// @@ -312,14 +327,14 @@ final messageHistoryProvider = FutureProvider.autoDispose // ── Hydration (resume) ──────────────────────────────────────────────────────── /// Rebuild the chat rooms from the trade list and the persisted messages — -/// the same path `ChatRoomsScreen` runs on init — and upsert them, so a room -/// added live while the fetch ran survives. Assumes the trade list was +/// the same path `ChatRoomsScreen` runs on init — and merge them in, so a room +/// added or updated live while the fetch ran survives. Assumes the trade list was /// hydrated first (`defaultHydrators` orders it so). Future hydrateChatRooms(ProviderContainer container) async { container.invalidate(chatRoomsFromTradesProvider); final rooms = await container.read(chatRoomsFromTradesProvider.future); final notifier = container.read(chatRoomsNotifierProvider.notifier); for (final room in rooms) { - notifier.upsertRoom(room); + notifier.upsertIfNewer(room); } } diff --git a/lib/features/chat/screens/chat_rooms_screen.dart b/lib/features/chat/screens/chat_rooms_screen.dart index 7f13bea5..17ae5d28 100644 --- a/lib/features/chat/screens/chat_rooms_screen.dart +++ b/lib/features/chat/screens/chat_rooms_screen.dart @@ -49,11 +49,11 @@ class _ChatRoomsScreenState extends ConsumerState { try { final rooms = await ref.read(chatRoomsFromTradesProvider.future); if (!mounted) return; - // Upsert rather than replace: a room added concurrently by - // ChatRoomScreen.upsertRoom (a message landing mid-fetch) survives. + // Merge rather than replace: a room added or updated concurrently by + // ChatRoomScreen (a message landing mid-fetch) survives. final notifier = ref.read(chatRoomsNotifierProvider.notifier); for (final room in rooms) { - notifier.upsertRoom(room); + notifier.upsertIfNewer(room); } } catch (e) { debugPrint('[chat] syncRoomsFromTrades failed: $e'); diff --git a/lib/features/disputes/providers/disputes_providers.dart b/lib/features/disputes/providers/disputes_providers.dart index 0eaaa96f..fdf44b67 100644 --- a/lib/features/disputes/providers/disputes_providers.dart +++ b/lib/features/disputes/providers/disputes_providers.dart @@ -237,15 +237,19 @@ DisputeItem disputeItemFromRust(rust_types.Dispute dispute) => DisputeItem( isRead: dispute.isRead, ); -/// Trade statuses that can carry a dispute record on the bridge. -const _disputedStatuses = { - rust_types.OrderStatus.dispute, - rust_types.OrderStatus.settledByAdmin, - rust_types.OrderStatus.canceledByAdmin, - rust_types.OrderStatus.completedByAdmin, +/// Trade statuses under which the bridge cannot hold a dispute: the trade +/// ended without one. Everything else is queried, because a dispute's record +/// and the row's status are written by different arms — `admin-took-dispute` +/// creates an `InReview` dispute without touching the status, so the row can +/// still read `active`, `fiatSent` or `inProgress` while a dispute exists. +const _undisputableStatuses = { + rust_types.OrderStatus.success, + rust_types.OrderStatus.canceled, + rust_types.OrderStatus.expired, + rust_types.OrderStatus.cooperativelyCanceled, }; -/// Re-read every dispute the bridge knows for the trades that can have one +/// Re-read every dispute the bridge knows for the trades that can own one /// and upsert it: the read flag the UI manages survives (`upsert` keeps it), /// and a dispute opened by the peer while the process was suspended appears /// without a restart — the shape of the v1 bug this exists to prevent @@ -258,7 +262,7 @@ Future hydrateDisputes( final trades = await container.read(rawTradesProvider.future); final notifier = container.read(disputeNotifierProvider.notifier); for (final trade in trades) { - if (!_disputedStatuses.contains(trade.order.status)) continue; + if (_undisputableStatuses.contains(trade.order.status)) continue; final rust_types.Dispute? dispute; try { dispute = await lookup(tradeId: trade.order.id); diff --git a/lib/features/notifications/providers/notifications_provider.dart b/lib/features/notifications/providers/notifications_provider.dart index b19b2a49..edde145a 100644 --- a/lib/features/notifications/providers/notifications_provider.dart +++ b/lib/features/notifications/providers/notifications_provider.dart @@ -214,18 +214,37 @@ class NotificationsNotifier extends StateNotifier> { final SembastNotificationsStore? store; + /// Ids deleted while a load was reading the store. The snapshot the load + /// returns still holds them, so the merge must not bring them back. + final Set _deletedDuringLoad = {}; + + /// Loads in flight; deletions are only tracked while this is non-zero. + int _loadsInFlight = 0; + + /// Set by [deleteAll] while a load is in flight: the whole snapshot that + /// load returns predates the wipe and is discarded. + bool _wipedDuringLoad = false; + /// Load persisted notifications into state. Called once on construction - /// when a [store] is provided. + /// when a [store] is provided, and again on every resume (the resync + /// hydration, lib/core/lifecycle/resume_resync.dart). /// /// Merges the persisted snapshot with whatever is already in state, keyed by /// id, so a delayed load never drops (or overwrites with a stale copy) a /// notification added live while the load was in flight. Records added this - /// session win on conflict. + /// session win on conflict. A record the user deleted while the load was + /// reading is not resurrected: [delete] and [deleteAll] note the removal, + /// and the merge skips it. Future loadInitialData() async { if (store == null) return; + _loadsInFlight++; try { final loaded = await store!.loadAll(); - final byId = {for (final n in loaded) n.id: n}; + if (_wipedDuringLoad) return; + final byId = { + for (final n in loaded) + if (!_deletedDuringLoad.contains(n.id)) n.id: n, + }; for (final n in state) { byId[n.id] = n; } @@ -234,6 +253,12 @@ class NotificationsNotifier extends StateNotifier> { ..sort((a, b) => b.timestamp.compareTo(a.timestamp)); } catch (e) { debugPrint('NotificationsNotifier: failed to load from Sembast: $e'); + } finally { + _loadsInFlight--; + if (_loadsInFlight == 0) { + _deletedDuringLoad.clear(); + _wipedDuringLoad = false; + } } } @@ -302,6 +327,7 @@ class NotificationsNotifier extends StateNotifier> { } Future delete(String id) async { + if (_loadsInFlight > 0) _deletedDuringLoad.add(id); state = state.where((n) => n.id != id).toList(); try { await store?.deleteRecord(id); @@ -311,6 +337,7 @@ class NotificationsNotifier extends StateNotifier> { } Future deleteAll() async { + if (_loadsInFlight > 0) _wipedDuringLoad = true; state = []; try { await store?.deleteAll(); diff --git a/test/core/lifecycle/hydration_regression_test.dart b/test/core/lifecycle/hydration_regression_test.dart index 33b6871c..6e77aff1 100644 --- a/test/core/lifecycle/hydration_regression_test.dart +++ b/test/core/lifecycle/hydration_regression_test.dart @@ -1,3 +1,5 @@ +import 'dart:async'; + import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:flutter_test/flutter_test.dart'; import 'package:mostro/core/lifecycle/resume_resync.dart'; @@ -45,8 +47,11 @@ void main() { Future lookup({required String tradeId}) async => bridgeDisputes[tradeId]; - ResumeResync routine() => ResumeResync( + ResumeResync routine({Stream Function()? updates}) => ResumeResync( container: container, + updates: updates ?? () => const Stream.empty(), + settleQuiet: const Duration(milliseconds: 30), + settleMax: const Duration(milliseconds: 200), resync: () async => const rust_types.ResyncOutcome( online: true, @@ -118,6 +123,98 @@ void main() { }, ); + test( + 'a dispute the solver took while the row still reads fiat-sent is listed', + () async { + // `admin-took-dispute` creates the InReview record without touching the + // trade status: the row may still read active, fiat-sent or in-progress. + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.fiatSent), + ]; + bridgeDisputes['order-t1'] = const rust_types.Dispute( + id: 'd1', + tradeId: 'order-t1', + status: rust_types.DisputeStatus.inReview, + initiatedByMe: false, + adminPubkey: 'solver', + openedAt: 1234, + isRead: false, + ); + + await routine().run(); + + expect(container.read(disputeNotifierProvider).single.id, 'd1'); + }, + ); + + test('a trade that finished without a dispute is not queried', () async { + var lookups = 0; + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.success), + fakeTrade(id: 't2', status: rust_types.OrderStatus.canceled), + ]; + final counting = ResumeResync( + container: container, + resync: + () async => const rust_types.ResyncOutcome( + online: true, + flushed: 0, + coalesced: false, + ), + hydrators: [ + hydrateTrades, + (c) => hydrateDisputes( + c, + getDispute: ({required tradeId}) async { + lookups++; + return null; + }, + ), + ], + updates: () => const Stream.empty(), + settleMax: const Duration(milliseconds: 50), + ); + + await counting.run(); + + expect(lookups, 0); + }); + + test( + 'events that land during the replay, after resync returned, are picked up', + () async { + // resync() returns once the subscriptions are re-issued; the events + // arrive over the next seconds. The first pass sees the old row; the + // replay writes the new one and emits a trade update; the second pass + // sees it. + final updates = StreamController(); + addTearDown(updates.close); + final run = routine(updates: () => updates.stream).run(); + await Future.delayed(const Duration(milliseconds: 10)); + expect( + (await container.read(rawTradesProvider.future)).single.order.status, + rust_types.OrderStatus.active, + reason: 'the first pass ran against the pre-replay row', + ); + + bridgeTrades = [ + fakeTrade(id: 't1', status: rust_types.OrderStatus.dispute), + ]; + bridgeDisputes['order-t1'] = const rust_types.Dispute( + id: 'd1', + tradeId: 'order-t1', + status: rust_types.DisputeStatus.open, + initiatedByMe: false, + openedAt: 1234, + isRead: false, + ); + updates.add(null); + await run; + + expect(container.read(disputeNotifierProvider).single.id, 'd1'); + }, + ); + test('re-hydrating a dispute keeps the read flag the user set', () async { bridgeTrades = [ fakeTrade(id: 't1', status: rust_types.OrderStatus.dispute), diff --git a/test/core/lifecycle/resume_resync_test.dart b/test/core/lifecycle/resume_resync_test.dart index ee12cbcd..c598ae03 100644 --- a/test/core/lifecycle/resume_resync_test.dart +++ b/test/core/lifecycle/resume_resync_test.dart @@ -1,3 +1,5 @@ +import 'dart:async'; + import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:flutter_test/flutter_test.dart'; import 'package:mostro/core/lifecycle/resume_resync.dart'; @@ -21,7 +23,10 @@ void main() { if (throws) throw StateError('$name failed'); }; - test('resync runs first, then every hydrator in order', () async { + /// No replay to wait for: the settle window closes at once. + Stream noUpdates() => const Stream.empty(); + + test('resync runs first, then every hydrator in order, twice', () async { // Arrange final routine = ResumeResync( container: container, @@ -30,13 +35,22 @@ void main() { return ok; }, hydrators: [hook('trades'), hook('chat'), hook('disputes')], + updates: noUpdates, ); // Act await routine.run(); - // Assert - expect(calls, ['resync', 'trades', 'chat', 'disputes']); + // Assert — once for what is on disk, once more when the replay settled. + expect(calls, [ + 'resync', + 'trades', + 'chat', + 'disputes', + 'trades', + 'chat', + 'disputes', + ]); }); test('a hydrator that throws does not stop the ones after it', () async { @@ -44,11 +58,12 @@ void main() { container: container, resync: () async => ok, hydrators: [hook('trades'), hook('chat', throws: true), hook('disputes')], + updates: noUpdates, ); await routine.run(); - expect(calls, ['trades', 'chat', 'disputes']); + expect(calls.take(3), ['trades', 'chat', 'disputes']); }); test('a failed resync still hydrates — disk is newer than memory', () async { @@ -56,11 +71,92 @@ void main() { container: container, resync: () async => throw StateError('no bridge'), hydrators: [hook('trades')], + updates: noUpdates, ); await expectLater(routine.run(), completes); - expect(calls, ['trades']); + expect(calls, ['trades', 'trades']); + }); + + group('the second pass waits for the replay to settle', () { + const quiet = Duration(milliseconds: 40); + const max = Duration(milliseconds: 300); + + test('it runs after a quiet window following the last update', () async { + final updates = StreamController(); + addTearDown(updates.close); + final passes = []; + final routine = ResumeResync( + container: container, + resync: () async => ok, + hydrators: [(c) async => passes.add(DateTime.now())], + updates: () => updates.stream, + settleQuiet: quiet, + settleMax: max, + ); + + final run = routine.run(); + // Replay: three updates 15 ms apart keep the window open. + for (var i = 0; i < 3; i++) { + await Future.delayed(const Duration(milliseconds: 15)); + updates.add(i); + } + final lastUpdate = DateTime.now(); + await run; + + expect(passes, hasLength(2)); + expect( + passes[1].difference(lastUpdate), + greaterThanOrEqualTo(quiet - const Duration(milliseconds: 5)), + reason: 'the second pass came only once the replay went quiet', + ); + }); + + test( + 'a replay that never goes quiet still gets its pass at the cap', + () async { + final updates = StreamController(); + addTearDown(updates.close); + final routine = ResumeResync( + container: container, + resync: () async => ok, + hydrators: [hook('h')], + updates: () => updates.stream, + settleQuiet: quiet, + settleMax: max, + ); + final chatter = Timer.periodic( + const Duration(milliseconds: 10), + (_) => updates.add(null), + ); + addTearDown(chatter.cancel); + + final started = DateTime.now(); + await routine.run(); + + expect(calls, ['h', 'h']); + expect(DateTime.now().difference(started), lessThan(max * 3)); + }, + ); + + test( + 'an update stream that fails does not block the second pass', + () async { + final routine = ResumeResync( + container: container, + resync: () async => ok, + hydrators: [hook('h')], + updates: () => Stream.error(StateError('no bridge')), + settleQuiet: quiet, + settleMax: max, + ); + + await routine.run(); + + expect(calls, ['h', 'h']); + }, + ); }); test('the default hook list covers every protocol-state notifier', () { diff --git a/test/features/chat/providers/chat_rooms_upsert_if_newer_test.dart b/test/features/chat/providers/chat_rooms_upsert_if_newer_test.dart new file mode 100644 index 00000000..81179060 --- /dev/null +++ b/test/features/chat/providers/chat_rooms_upsert_if_newer_test.dart @@ -0,0 +1,55 @@ +import 'package:flutter_test/flutter_test.dart'; +import 'package:mostro/features/chat/providers/chat_providers.dart'; + +ChatRoomState _room({ + String id = 'o1', + int at = 100, + String? last = 'hola', + int unread = 0, +}) => ChatRoomState( + orderId: id, + peerPubkey: 'peer', + peerHandle: 'used-jaguar', + peerIconIndex: 0, + peerColorHue: 0, + isSelling: true, + lastMessage: last, + lastMessageAt: at, + unreadCount: unread, +); + +void main() { + group('ChatRoomsNotifier.upsertIfNewer', () { + test('a room the list does not have is added', () { + final notifier = ChatRoomsNotifier(); + + notifier.upsertIfNewer(_room()); + + expect(notifier.state.single.orderId, 'o1'); + }); + + test('a snapshot older than the live room leaves it alone', () { + // A message folded in while the snapshot was being built. + final notifier = + ChatRoomsNotifier() + ..setRooms([_room(at: 200, last: 'nuevo', unread: 1)]); + + notifier.upsertIfNewer(_room(at: 100, last: 'hola')); + + final room = notifier.state.single; + expect(room.lastMessage, 'nuevo'); + expect(room.lastMessageAt, 200); + expect(room.unreadCount, 1); + }); + + test('a snapshot at or past the live room replaces it', () { + final notifier = ChatRoomsNotifier()..setRooms([_room(at: 100)]); + + notifier.upsertIfNewer(_room(at: 100, last: 'same time, re-read')); + expect(notifier.state.single.lastMessage, 'same time, re-read'); + + notifier.upsertIfNewer(_room(at: 300, last: 'later')); + expect(notifier.state.single.lastMessage, 'later'); + }); + }); +} diff --git a/test/features/notifications/notifications_resume_load_test.dart b/test/features/notifications/notifications_resume_load_test.dart new file mode 100644 index 00000000..dcf46f09 --- /dev/null +++ b/test/features/notifications/notifications_resume_load_test.dart @@ -0,0 +1,110 @@ +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:mostro/features/notifications/models/notification_model.dart'; +import 'package:mostro/features/notifications/providers/notifications_provider.dart'; +import 'package:sembast/sembast_memory.dart'; + +NotificationModel _note(String id) => NotificationModel( + id: id, + type: NotificationType.system, + title: 'title-$id', + message: 'message-$id', + timestamp: DateTime.utc(2026, 1, 1), +); + +/// A store whose read can be held open, so a deletion can land between the +/// snapshot being taken and the merge publishing it — the resume race. +class _HeldStore extends SembastNotificationsStore { + _HeldStore({required super.factory, required super.path}); + + Completer? hold; + + @override + Future> loadAll() async { + final gate = hold; + if (gate != null) await gate.future; + return super.loadAll(); + } +} + +void main() { + late _HeldStore store; + late NotificationsNotifier notifier; + + setUp(() async { + final factory = newDatabaseFactoryMemory(); + store = _HeldStore( + factory: factory, + path: 'resume-${DateTime.now().microsecondsSinceEpoch}.db', + ); + notifier = NotificationsNotifier(store: store); + await notifier.loadInitialData(); + await notifier.add(_note('a')); + await notifier.add(_note('b')); + expect(notifier.state.map((n) => n.id), unorderedEquals(['a', 'b'])); + }); + + test( + 'a notice deleted while a resume load is reading stays deleted', + () async { + // Arrange — the load's snapshot is taken before the delete... + store.hold = Completer(); + final load = notifier.loadInitialData(); + await Future.delayed(Duration.zero); + + // Act — ...the user deletes, then the snapshot returns. + await notifier.delete('a'); + store.hold!.complete(); + await load; + + // Assert + expect(notifier.state.map((n) => n.id), ['b']); + }, + ); + + test( + 'a clear-all while a resume load is reading wins over the snapshot', + () async { + store.hold = Completer(); + final load = notifier.loadInitialData(); + await Future.delayed(Duration.zero); + + await notifier.deleteAll(); + store.hold!.complete(); + await load; + + expect(notifier.state, isEmpty); + }, + ); + + test( + 'a notice added while a resume load is reading survives the merge', + () async { + store.hold = Completer(); + final load = notifier.loadInitialData(); + await Future.delayed(Duration.zero); + + await notifier.add(_note('c')); + store.hold!.complete(); + await load; + + expect(notifier.state.map((n) => n.id), unorderedEquals(['a', 'b', 'c'])); + }, + ); + + test( + 'a delete after the load finished is not remembered by the next load', + () async { + await notifier.delete('a'); + await notifier.loadInitialData(); + expect(notifier.state.map((n) => n.id), ['b']); + + // The next load starts clean: a fresh record with the deleted id is a + // different notice and comes back. + await notifier.add(_note('a')); + await notifier.loadInitialData(); + expect(notifier.state.map((n) => n.id), unorderedEquals(['a', 'b'])); + }, + ); +}