-
Notifications
You must be signed in to change notification settings - Fork 8
feat(push): PR-0b — app lifecycle service and resume hydration #464
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
Merged
Changes from all commits
Commits
Show all changes
3 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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<void> 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'); | ||
| }), | ||
| ); | ||
| } | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,143 @@ | ||
| import 'dart:async'; | ||
|
|
||
| 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/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. | ||
| /// | ||
| /// **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<void> Function(ProviderContainer container); | ||
|
|
||
| /// 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, 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<ResyncOutcome> Function()? resync, | ||
| List<Hydrator>? hydrators, | ||
| Stream<Object?> Function()? updates, | ||
| this.settleQuiet = const Duration(milliseconds: 1500), | ||
| this.settleMax = const Duration(seconds: 10), | ||
| }) : _resync = resync ?? nostr_api.resync, | ||
| _hydrators = hydrators ?? defaultHydrators, | ||
| _updates = updates ?? _bridgeTradeUpdates; | ||
|
|
||
| final ProviderContainer container; | ||
| final Future<ResyncOutcome> Function() _resync; | ||
| final List<Hydrator> _hydrators; | ||
| final Stream<Object?> 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<void> 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'); | ||
| } | ||
| await _hydrateAll(); | ||
| await _waitForReplayToSettle(); | ||
| await _hydrateAll(); | ||
| } | ||
|
|
||
| Future<void> _hydrateAll() async { | ||
| for (final hydrate in _hydrators) { | ||
| try { | ||
| await hydrate(container); | ||
| } catch (e, st) { | ||
| debugPrint('[lifecycle] hydration failed: $e\n$st'); | ||
| } | ||
| } | ||
| } | ||
|
|
||
| /// Completes [settleQuiet] after the last update, or at [settleMax]. | ||
| Future<void> _waitForReplayToSettle() async { | ||
| final settled = Completer<void>(); | ||
| void done() { | ||
| if (!settled.isCompleted) settled.complete(); | ||
| } | ||
|
|
||
| Timer quiet = Timer(settleQuiet, done); | ||
| final cap = Timer(settleMax, done); | ||
| StreamSubscription<Object?>? 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(); | ||
|
grunch marked this conversation as resolved.
|
||
| } | ||
| } | ||
|
|
||
| /// The bridge's trade updates, as a stream, for the settle window. | ||
| Stream<TradeUpdate> _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: | ||
| /// trades first, since chat rooms and disputes are read off the trade list. | ||
| final List<Hydrator> defaultHydrators = [ | ||
| hydrateTrades, | ||
| hydrateChatRooms, | ||
| hydrateDisputes, | ||
| hydrateNotifications, | ||
| ]; | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.