From b1c04d11226397823db4833d9cedbfc94135c8fe Mon Sep 17 00:00:00 2001 From: Aiden Almazan Date: Wed, 5 Aug 2026 18:47:26 -0700 Subject: [PATCH] fix(mobile): probe socket liveness on resume Detect half-open WebSocket sessions after short app background periods with a bounded REQ/EOSE probe. Reconnect and replay live subscriptions when the probe times out or errors. Co-authored-by: Aiden Almazan Signed-off-by: Aiden Almazan --- mobile/lib/shared/relay/relay_session.dart | 32 ++- .../test/shared/relay/relay_session_test.dart | 208 ++++++++++++++++++ 2 files changed, 239 insertions(+), 1 deletion(-) diff --git a/mobile/lib/shared/relay/relay_session.dart b/mobile/lib/shared/relay/relay_session.dart index 1c5c305b4c..49bb43121c 100644 --- a/mobile/lib/shared/relay/relay_session.dart +++ b/mobile/lib/shared/relay/relay_session.dart @@ -93,12 +93,14 @@ class RelaySessionNotifier extends Notifier { RelayRateLimitGate? rateLimitGate, RelayTimerFactory retryTimerFactory = Timer.new, Future Function(Duration) replayDelay = Future.delayed, + Duration resumeProbeTimeout = const Duration(seconds: 3), }) : _httpClient = httpClient, _socketFactory = socketFactory, _now = now ?? DateTime.now, _rateLimitGate = rateLimitGate ?? RelayRateLimitGate(), _retryTimerFactory = retryTimerFactory, - _replayDelay = replayDelay; + _replayDelay = replayDelay, + _resumeProbeTimeout = resumeProbeTimeout; final http.Client? _httpClient; final RelaySocketFactory _socketFactory; @@ -106,6 +108,7 @@ class RelaySessionNotifier extends Notifier { final RelayRateLimitGate _rateLimitGate; final RelayTimerFactory _retryTimerFactory; final Future Function(Duration) _replayDelay; + final Duration _resumeProbeTimeout; static const _baseReconnectDelayMs = 1000; static const _maxReconnectDelayMs = 30000; @@ -397,6 +400,25 @@ class RelaySessionNotifier extends Notifier { await _connect(config); } + /// Liveness probe after resume: a minimal REQ that resolves on EOSE. On + /// timeout or error the socket is presumed dead and we reconnect (live + /// subscriptions replay with the usual since-skew). + Future _verifyConnectionAfterResume(int connectionGeneration) async { + try { + await fetchHistory( + const NostrFilter(kinds: [39000], limit: 1), + timeout: _resumeProbeTimeout, + ); + } catch (_) { + if (_disposed || _paused) return; + if (connectionGeneration != _connectionGeneration || !_socketConnected) { + return; + } + if (state.status != SessionStatus.connected) return; + await reconnect(); + } + } + /// Called by the app lifecycle provider when the app goes to background. void onAppPaused() { _backgroundedAt = _now(); @@ -427,6 +449,14 @@ class RelaySessionNotifier extends Notifier { _now().difference(backgroundedAt) >= _backgroundGraceDuration; if (!backgroundedLongEnoughToRequireReconnect && state.status == SessionStatus.connected) { + // Even after a short background stint the transport can be half-open: + // iOS may drop the network on screen lock or rebind NAT on a + // Wi-Fi/cellular switch without a close frame ever reaching us. If we + // trust the flag here, the session is a zombie — live subscriptions + // stay silent until the user force-kills the app. Verify with a cheap + // REQ/EOSE round-trip instead; any failure forces a reconnect, which + // replays live subscriptions. + unawaited(_verifyConnectionAfterResume(_connectionGeneration)); return; } diff --git a/mobile/test/shared/relay/relay_session_test.dart b/mobile/test/shared/relay/relay_session_test.dart index 826e234400..277542f54f 100644 --- a/mobile/test/shared/relay/relay_session_test.dart +++ b/mobile/test/shared/relay/relay_session_test.dart @@ -546,6 +546,197 @@ void main() { }, ); + test( + 'resume within the grace period probes the socket and keeps it when the ' + 'relay answers', + () async { + final sockets = <_ProbeRecordingRelaySocket>[]; + final keychain = nostr.Keys.generate(); + var now = DateTime(2026, 8, 2, 12); + final session = RelaySessionNotifier( + now: () => now, + socketFactory: + ({ + required wsUrl, + required nsec, + required onMessage, + required onConnected, + required onDisconnected, + }) { + final socket = _ProbeRecordingRelaySocket( + wsUrl: wsUrl, + nsec: nsec, + onMessage: onMessage, + onConnected: onConnected, + onDisconnected: onDisconnected, + ); + sockets.add(socket); + return socket; + }, + ); + final container = ProviderContainer( + overrides: [ + relaySessionProvider.overrideWith(() => session), + relayConfigProvider.overrideWith( + () => _FakeRelayConfigNotifier( + baseUrl: 'https://relay.example', + nsec: keychain.nsec, + ), + ), + authProvider.overrideWith(() => _AuthenticatedAuthNotifier()), + ], + ); + addTearDown(container.dispose); + await container.read(authProvider.future); + final subscription = container.listen(relaySessionProvider, (_, _) {}); + addTearDown(subscription.close); + await Future.delayed(Duration.zero); + sockets.single.connectSuccessfully(); + + session.onAppPaused(); + now = now.add(const Duration(seconds: 4)); + session.onAppResumed(); + await Future.delayed(Duration.zero); + + // The resume path issued a liveness REQ; answer it with EOSE. + final probeReq = sockets.single.sentFrames.lastWhere( + (frame) => frame.isNotEmpty && frame.first == 'REQ', + ); + session.debugHandleMessage(['EOSE', probeReq[1]]); + await Future.delayed(Duration.zero); + + expect(sockets, hasLength(1)); + expect(sockets.single.disposeCalls, 0); + expect(session.state.status, SessionStatus.connected); + }, + ); + + test( + 'resume within the grace period reconnects when the socket is half-open ' + '(probe never answered)', + () async { + final sockets = <_ProbeRecordingRelaySocket>[]; + final keychain = nostr.Keys.generate(); + var now = DateTime(2026, 8, 2, 12); + final session = RelaySessionNotifier( + now: () => now, + resumeProbeTimeout: const Duration(milliseconds: 20), + socketFactory: + ({ + required wsUrl, + required nsec, + required onMessage, + required onConnected, + required onDisconnected, + }) { + final socket = _ProbeRecordingRelaySocket( + wsUrl: wsUrl, + nsec: nsec, + onMessage: onMessage, + onConnected: onConnected, + onDisconnected: onDisconnected, + ); + sockets.add(socket); + return socket; + }, + ); + final container = ProviderContainer( + overrides: [ + relaySessionProvider.overrideWith(() => session), + relayConfigProvider.overrideWith( + () => _FakeRelayConfigNotifier( + baseUrl: 'https://relay.example', + nsec: keychain.nsec, + ), + ), + authProvider.overrideWith(() => _AuthenticatedAuthNotifier()), + ], + ); + addTearDown(container.dispose); + await container.read(authProvider.future); + final subscription = container.listen(relaySessionProvider, (_, _) {}); + addTearDown(subscription.close); + await Future.delayed(Duration.zero); + sockets.single.connectSuccessfully(); + + session.onAppPaused(); + now = now.add(const Duration(seconds: 4)); + session.onAppResumed(); + + // Probe times out against the dead transport, forcing a reconnect. + await Future.delayed(const Duration(milliseconds: 60)); + + expect(sockets, hasLength(2)); + expect(sockets.first.disposeCalls, 1); + expect(session.state.status, SessionStatus.reconnecting); + }, + ); + + test( + 'stale resume probe cannot reconnect a replacement connection', + () async { + final sockets = <_ProbeRecordingRelaySocket>[]; + final keychain = nostr.Keys.generate(); + var now = DateTime(2026, 8, 2, 12); + final session = RelaySessionNotifier( + now: () => now, + resumeProbeTimeout: const Duration(milliseconds: 20), + socketFactory: + ({ + required wsUrl, + required nsec, + required onMessage, + required onConnected, + required onDisconnected, + }) { + final socket = _ProbeRecordingRelaySocket( + wsUrl: wsUrl, + nsec: nsec, + onMessage: onMessage, + onConnected: onConnected, + onDisconnected: onDisconnected, + ); + sockets.add(socket); + return socket; + }, + ); + final container = ProviderContainer( + overrides: [ + relaySessionProvider.overrideWith(() => session), + relayConfigProvider.overrideWith( + () => _FakeRelayConfigNotifier( + baseUrl: 'https://relay.example', + nsec: keychain.nsec, + ), + ), + authProvider.overrideWith(() => _AuthenticatedAuthNotifier()), + ], + ); + addTearDown(container.dispose); + await container.read(authProvider.future); + final subscription = container.listen(relaySessionProvider, (_, _) {}); + addTearDown(subscription.close); + await Future.delayed(Duration.zero); + sockets.single.connectSuccessfully(); + + session.onAppPaused(); + now = now.add(const Duration(seconds: 4)); + session.onAppResumed(); + await Future.delayed(Duration.zero); + + // Replace the connection before the old probe times out. + await session.reconnect(); + expect(sockets, hasLength(2)); + sockets.last.connectSuccessfully(); + + await Future.delayed(const Duration(milliseconds: 60)); + + expect(sockets, hasLength(2)); + expect(sockets.last.disposeCalls, 0); + expect(session.state.status, SessionStatus.connected); + }, + ); + test('delivers the same live event to each matching subscription', () async { final session = RelaySessionNotifier(); final firstEvents = []; @@ -1281,6 +1472,23 @@ class _ControlledRelaySocket extends RelaySocket { void disconnectWith(Object? error) => _disconnected(error); } +class _ProbeRecordingRelaySocket extends _ControlledRelaySocket { + final List> sentFrames = []; + + _ProbeRecordingRelaySocket({ + required super.wsUrl, + required super.nsec, + required super.onMessage, + required super.onConnected, + required super.onDisconnected, + }); + + @override + void send(List payload) { + sentFrames.add(payload); + } +} + const _channelId = '11111111-1111-4111-8111-111111111111'; class _FakeRelayConfigNotifier extends RelayConfigNotifier {