From c193cd4dcbfed65fd6d63b808625f37642abb40d Mon Sep 17 00:00:00 2001 From: npub1tquskdu6yc4h8l7xxtceculxw600grekeq0xg2ukqfrwl7vrzg3quz3gmp <58390b379a262b73ffc632f19c73e6769ef40f36c81e642b960246eff9831222@buzz.block.builderlab.xyz> Date: Sun, 2 Aug 2026 10:53:12 -0700 Subject: [PATCH] fix(mobile): recover stale relay sessions Co-authored-by: Tom Brow Signed-off-by: Tom Brow --- mobile/lib/shared/relay/relay_session.dart | 23 +++- mobile/lib/shared/relay/relay_socket.dart | 12 +- mobile/pubspec.lock | 2 +- mobile/pubspec.yaml | 1 + .../test/shared/relay/relay_session_test.dart | 119 +++++++++++++++++- .../relay/relay_socket_liveness_test.dart | 111 ++++++++++++++++ 6 files changed, 260 insertions(+), 8 deletions(-) create mode 100644 mobile/test/shared/relay/relay_socket_liveness_test.dart diff --git a/mobile/lib/shared/relay/relay_session.dart b/mobile/lib/shared/relay/relay_session.dart index d34c5405e8..0bb211a2b0 100644 --- a/mobile/lib/shared/relay/relay_session.dart +++ b/mobile/lib/shared/relay/relay_session.dart @@ -78,17 +78,21 @@ class RelaySessionNotifier extends Notifier { RelaySessionNotifier({ http.Client? httpClient, RelaySocketFactory socketFactory = RelaySocket.new, + DateTime Function()? now, }) : _httpClient = httpClient, - _socketFactory = socketFactory; + _socketFactory = socketFactory, + _now = now ?? DateTime.now; final http.Client? _httpClient; final RelaySocketFactory _socketFactory; + final DateTime Function() _now; static const _baseReconnectDelayMs = 1000; static const _maxReconnectDelayMs = 30000; static const _eventBatchMs = 16; static const _reconnectReplaySkewSeconds = 5; static const _maxRecentDeliveryKeys = 5000; + static const _backgroundGraceDuration = Duration(seconds: 5); RelaySocket? _socket; final Map _historySubscriptions = {}; @@ -99,6 +103,7 @@ class RelaySessionNotifier extends Notifier { Timer? _reconnectTimer; Timer? _flushTimer; Timer? _backgroundGraceTimer; + DateTime? _backgroundedAt; int _reconnectDelayMs = _baseReconnectDelayMs; int _subIdCounter = 0; bool _disposed = false; @@ -315,8 +320,9 @@ class RelaySessionNotifier extends Notifier { /// Called by the app lifecycle provider when the app goes to background. void onAppPaused() { + _backgroundedAt = _now(); _backgroundGraceTimer?.cancel(); - _backgroundGraceTimer = Timer(const Duration(seconds: 5), _pauseNow); + _backgroundGraceTimer = Timer(_backgroundGraceDuration, _pauseNow); } void _pauseNow() { @@ -331,12 +337,18 @@ class RelaySessionNotifier extends Notifier { /// Called by the app lifecycle provider when the app returns to foreground. void onAppResumed() { _paused = false; + final backgroundedAt = _backgroundedAt; + _backgroundedAt = null; _backgroundGraceTimer?.cancel(); _backgroundGraceTimer = null; - // If still connected, nothing to do — the socket survived the background - // grace window. - if (state.status == SessionStatus.connected) return; + final backgroundedLongEnoughToRequireReconnect = + backgroundedAt != null && + _now().difference(backgroundedAt) >= _backgroundGraceDuration; + if (!backgroundedLongEnoughToRequireReconnect && + state.status == SessionStatus.connected) { + return; + } // Cancel any in-flight reconnect backoff timer so we reconnect immediately // instead of waiting for the (possibly large) exponential delay. @@ -650,6 +662,7 @@ class RelaySessionNotifier extends Notifier { _reconnectTimer?.cancel(); _flushTimer?.cancel(); _backgroundGraceTimer?.cancel(); + _backgroundedAt = null; _cancelAllHistory(null); _rejectAllPending(null); _recentDeliveryKeys.clear(); diff --git a/mobile/lib/shared/relay/relay_socket.dart b/mobile/lib/shared/relay/relay_socket.dart index 267b030391..5b23279e81 100644 --- a/mobile/lib/shared/relay/relay_socket.dart +++ b/mobile/lib/shared/relay/relay_socket.dart @@ -3,6 +3,7 @@ import 'dart:convert'; import 'package:flutter/foundation.dart'; import 'package:nostr/nostr.dart' as nostr; +import 'package:web_socket_channel/io.dart'; import 'package:web_socket_channel/web_socket_channel.dart'; import 'nostr_models.dart'; @@ -30,6 +31,12 @@ Exception classifyRelayAuthFailure(String message) { } class RelaySocket { + /// Interval for sending a ping and awaiting its pong before disconnecting. + static const pingInterval = Duration(seconds: 30); + + @visibleForTesting + static Duration debugPingInterval = pingInterval; + final String _wsUrl; final String? _nsec; final void Function(List message) _onMessage; @@ -63,7 +70,10 @@ class RelaySocket { _state = SocketState.connecting; try { - _channel = WebSocketChannel.connect(Uri.parse(_wsUrl)); + _channel = IOWebSocketChannel.connect( + Uri.parse(_wsUrl), + pingInterval: debugPingInterval, + ); await _channel!.ready; } catch (e) { _state = SocketState.disconnected; diff --git a/mobile/pubspec.lock b/mobile/pubspec.lock index 05ccc2ea0f..6287e4c86c 100644 --- a/mobile/pubspec.lock +++ b/mobile/pubspec.lock @@ -274,7 +274,7 @@ packages: source: hosted version: "0.3.5+2" crypto: - dependency: transitive + dependency: "direct dev" description: name: crypto sha256: c8ea0233063ba03258fbcf2ca4d6dadfefe14f02fab57702265467a19f27fadf diff --git a/mobile/pubspec.yaml b/mobile/pubspec.yaml index 42b7935582..41d2a0aeb8 100644 --- a/mobile/pubspec.yaml +++ b/mobile/pubspec.yaml @@ -47,6 +47,7 @@ dev_dependencies: flutter_test: sdk: flutter flutter_lints: ^6.0.0 + crypto: ^3.0.7 custom_lint: ^0.8.0 riverpod_lint: ^3.1.0 mocktail: ^1.0.4 diff --git a/mobile/test/shared/relay/relay_session_test.dart b/mobile/test/shared/relay/relay_session_test.dart index 4647c05e07..643fbe37e9 100644 --- a/mobile/test/shared/relay/relay_session_test.dart +++ b/mobile/test/shared/relay/relay_session_test.dart @@ -260,6 +260,120 @@ void main() { expect(session.state.status, SessionStatus.disconnected); }); + test( + 'resume reconnects a stale connected session after a long pause', + () async { + final sockets = <_ControlledRelaySocket>[]; + 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 = _ControlledRelaySocket( + 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(minutes: 5)); + session.onAppResumed(); + await Future.delayed(Duration.zero); + + expect(sockets, hasLength(2)); + expect(sockets.first.disposeCalls, 1); + expect(session.state.status, SessionStatus.reconnecting); + }, + ); + + test( + 'resume keeps a connected session within the background grace period', + () async { + final sockets = <_ControlledRelaySocket>[]; + 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 = _ControlledRelaySocket( + 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); + + expect(sockets, hasLength(1)); + expect(sockets.single.disposeCalls, 0); + expect(session.state.status, SessionStatus.connected); + }, + ); + test('delivers the same live event to each matching subscription', () async { final session = RelaySessionNotifier(); final firstEvents = []; @@ -402,6 +516,7 @@ class _AuthenticatedAuthNotifier extends AuthNotifier { class _ControlledRelaySocket extends RelaySocket { final void Function() _connected; final void Function(Object? error) _disconnected; + int disposeCalls = 0; _ControlledRelaySocket({ required super.wsUrl, @@ -416,7 +531,9 @@ class _ControlledRelaySocket extends RelaySocket { Future connect() async {} @override - void dispose() {} + void dispose() { + disposeCalls++; + } void connectSuccessfully() => _connected(); diff --git a/mobile/test/shared/relay/relay_socket_liveness_test.dart b/mobile/test/shared/relay/relay_socket_liveness_test.dart new file mode 100644 index 0000000000..835ebc2b61 --- /dev/null +++ b/mobile/test/shared/relay/relay_socket_liveness_test.dart @@ -0,0 +1,111 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:buzz/shared/relay/relay_socket.dart'; +import 'package:crypto/crypto.dart'; +import 'package:flutter_test/flutter_test.dart'; + +/// A server that completes the WS handshake then never speaks again: no pongs, +/// no close frame. Only a client-side ping timeout can notice. +Future _silentAfterHandshakeServer() async { + final server = await ServerSocket.bind(InternetAddress.loopbackIPv4, 0); + server.listen((client) { + client.listen( + (data) { + final match = RegExp( + r'Sec-WebSocket-Key: (.*)\r\n', + caseSensitive: false, + ).firstMatch(String.fromCharCodes(data)); + if (match == null) return; + final accept = base64.encode( + sha1 + .convert( + utf8.encode( + '${match.group(1)!.trim()}258EAFA5-E914-47DA-95CA-C5AB0DC85B11', + ), + ) + .bytes, + ); + client.write( + 'HTTP/1.1 101 Switching Protocols\r\n' + 'Upgrade: websocket\r\nConnection: Upgrade\r\n' + 'Sec-WebSocket-Accept: $accept\r\n\r\n', + ); + }, + onError: (_) {}, + onDone: () {}, + ); + }); + return server; +} + +void main() { + const testPingInterval = Duration(milliseconds: 150); + + setUp(() { + RelaySocket.debugPingInterval = testPingInterval; + }); + + tearDown(() { + RelaySocket.debugPingInterval = RelaySocket.pingInterval; + }); + + test('detects a peer that stops answering pings', () async { + final server = await _silentAfterHandshakeServer(); + + final disconnected = Completer(); + final socket = RelaySocket( + wsUrl: 'ws://127.0.0.1:${server.port}', + nsec: null, + onMessage: (_) {}, + onConnected: () {}, + onDisconnected: (error) { + if (!disconnected.isCompleted) disconnected.complete(error); + }, + ); + unawaited(socket.connect()); + + var detected = true; + try { + await disconnected.future.timeout(testPingInterval * 4); + } on TimeoutException { + detected = false; + } + + expect( + detected, + isTrue, + reason: + 'RelaySocket must surface an unanswered ping through onDisconnected', + ); + + socket.dispose(); + await server.close(); + }); + + test('keeps an idle but healthy peer connected', () async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + server.transform(WebSocketTransformer()).listen((ws) { + // A healthy relay answers pings without sending application data. + ws.listen((_) {}, onError: (_) {}, onDone: () {}); + }); + + var tornDown = false; + final socket = RelaySocket( + wsUrl: 'ws://127.0.0.1:${server.port}', + nsec: null, + onMessage: (_) {}, + onConnected: () {}, + onDisconnected: (_) => tornDown = true, + ); + unawaited(socket.connect()); + + await Future.delayed(testPingInterval * 4); + + expect(tornDown, isFalse); + + await socket.disconnect(); + await server.close(force: true); + }); +}