Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions lib/core/app_bootstrap.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -209,6 +211,13 @@ Future<void> bootstrapAndRun({List<String> 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()),
);
Expand Down
91 changes: 91 additions & 0 deletions lib/core/lifecycle/app_lifecycle_service.dart
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');
}),
);
}
}
143 changes: 143 additions & 0 deletions lib/core/lifecycle/resume_resync.dart
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);
Comment thread
grunch marked this conversation as resolved.
} 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();
Comment thread
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,
];
Loading
Loading