From 649eaa4cce8c42345b833aaa34af2ee392ab2106 Mon Sep 17 00:00:00 2001 From: edde746 <86283021+edde746@users.noreply.github.com> Date: Sun, 12 Jul 2026 05:58:05 +0200 Subject: [PATCH] refactor(watch-together): share relay setup --- .../services/watch_together_peer_service.dart | 85 ++++---- .../watch_together_peer_service_test.dart | 205 ++++++++++++++++++ 2 files changed, 250 insertions(+), 40 deletions(-) create mode 100644 test/watch_together/watch_together_peer_service_test.dart diff --git a/lib/watch_together/services/watch_together_peer_service.dart b/lib/watch_together/services/watch_together_peer_service.dart index 0b2e2ca1..a9128236 100644 --- a/lib/watch_together/services/watch_together_peer_service.dart +++ b/lib/watch_together/services/watch_together_peer_service.dart @@ -44,6 +44,7 @@ class WatchTogetherPeerService with KeepaliveMixin { WebSocketChannel? _channel; StreamSubscription? _channelSubscription; + Completer? _setupCompleter; final Set _connectedPeers = {}; String? _sessionId; String? _myPeerId; @@ -123,31 +124,54 @@ class WatchTogetherPeerService with KeepaliveMixin { return channel; } + /// Connect, listen, and send a room setup announcement. + Future> _connectAndAnnounce(String type) async { + final channel = await _connectToRelay(); + _channel = channel; + + _listenToChannel(channel); + startKeepalive(); + + return _announce(type); + } + + Completer _announce(String type) { + final completer = Completer(); + _setupCompleter = completer; + _sendRaw({'type': type, 'sessionId': _sessionId, 'peerId': _myPeerId}); + return completer; + } + /// Listen on the channel stream and route incoming server messages. - void _listenToChannel(WebSocketChannel channel, {Completer? setupCompleter}) { + void _listenToChannel(WebSocketChannel channel) { _channelSubscription?.cancel(); _channelSubscription = channel.stream.listen( (data) { + if (!identical(_channel, channel)) return; resetPongTimer(); - _handleServerMessage(data as String, setupCompleter: setupCompleter); + _handleServerMessage(data as String); }, onError: (error) { + if (!identical(_channel, channel)) return; appLogger.e('WatchTogether: WebSocket error', error: error); _safeAdd( _errorController, PeerError(type: PeerErrorType.serverError, message: 'WebSocket error: $error', originalError: error), ); - if (setupCompleter != null && !setupCompleter.isCompleted) { - setupCompleter.completeError(error); + if (_setupCompleter case final completer? when !completer.isCompleted) { + completer.completeError(error); + _setupCompleter = null; } _handleWebSocketClosed(); }, onDone: () { + if (!identical(_channel, channel)) return; appLogger.w('WatchTogether: WebSocket closed'); - if (setupCompleter != null && !setupCompleter.isCompleted) { - setupCompleter.completeError( + if (_setupCompleter case final completer? when !completer.isCompleted) { + completer.completeError( const PeerError(type: PeerErrorType.connectionFailed, message: 'WebSocket closed before setup completed'), ); + _setupCompleter = null; } _handleWebSocketClosed(); }, @@ -155,7 +179,7 @@ class WatchTogetherPeerService with KeepaliveMixin { } /// Handle an incoming server message (JSON string). - void _handleServerMessage(String raw, {Completer? setupCompleter}) { + void _handleServerMessage(String raw) { try { final msg = jsonDecode(raw) as Map; final type = msg['type'] as String?; @@ -164,8 +188,9 @@ class WatchTogetherPeerService with KeepaliveMixin { case 'created': appLogger.d('WatchTogether: Room created: ${msg['sessionId']}'); _safeAdd(_connectionStateController, true); - if (setupCompleter != null && !setupCompleter.isCompleted) { - setupCompleter.complete(); + if (_setupCompleter case final completer? when !completer.isCompleted) { + completer.complete(); + _setupCompleter = null; } case 'joined': @@ -176,8 +201,9 @@ class WatchTogetherPeerService with KeepaliveMixin { _safeAdd(_peerConnectedController, peerId); } _safeAdd(_connectionStateController, true); - if (setupCompleter != null && !setupCompleter.isCompleted) { - setupCompleter.complete(); + if (_setupCompleter case final completer? when !completer.isCompleted) { + completer.complete(); + _setupCompleter = null; } case 'peerJoined': @@ -220,8 +246,9 @@ class WatchTogetherPeerService with KeepaliveMixin { appLogger.e('WatchTogether: Server error: $code - $message'); final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message', serverCode: code); _safeAdd(_errorController, error); - if (setupCompleter != null && !setupCompleter.isCompleted) { - setupCompleter.completeError(error); + if (_setupCompleter case final completer? when !completer.isCompleted) { + completer.completeError(error); + _setupCompleter = null; } case 'pong': @@ -299,26 +326,17 @@ class WatchTogetherPeerService with KeepaliveMixin { _reconnectTimer?.cancel(); _reconnectTimer = Timer(delay, () async { try { - final channel = await _connectToRelay(); - _channel = channel; - - final completer = Completer(); - _listenToChannel(channel, setupCompleter: completer); - startKeepalive(); - // Always try join first — the room may still have peers (e.g. host // reconnecting while guests remain). Fall back to create only if // the room no longer exists and we were the host. - _sendRaw({'type': 'join', 'sessionId': _sessionId, 'peerId': _myPeerId}); + final completer = await _connectAndAnnounce('join'); try { await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect'); } on PeerError catch (e) { if (_isHost && e.serverCode == 'room_not_found') { appLogger.d('WatchTogether: Room gone, re-creating as host'); - final createCompleter = Completer(); - _listenToChannel(channel, setupCompleter: createCompleter); - _sendRaw({'type': 'create', 'sessionId': _sessionId, 'peerId': _myPeerId}); + final createCompleter = _announce('create'); await createCompleter.future.namedTimeout( const Duration(seconds: 10), operation: 'WatchTogether reconnect create', @@ -357,14 +375,7 @@ class WatchTogetherPeerService with KeepaliveMixin { _reconnectAttempts = 0; try { - final channel = await _connectToRelay(); - _channel = channel; - - final completer = Completer(); - _listenToChannel(channel, setupCompleter: completer); - startKeepalive(); - - _sendRaw({'type': 'create', 'sessionId': _sessionId, 'peerId': _myPeerId}); + final completer = await _connectAndAnnounce('create'); await completer.future.timeout( const Duration(seconds: 10), @@ -394,14 +405,7 @@ class WatchTogetherPeerService with KeepaliveMixin { _reconnectAttempts = 0; try { - final channel = await _connectToRelay(); - _channel = channel; - - final completer = Completer(); - _listenToChannel(channel, setupCompleter: completer); - startKeepalive(); - - _sendRaw({'type': 'join', 'sessionId': _sessionId, 'peerId': _myPeerId}); + final completer = await _connectAndAnnounce('join'); await completer.future.timeout( const Duration(seconds: 10), @@ -447,6 +451,7 @@ class WatchTogetherPeerService with KeepaliveMixin { appLogger.d('WatchTogether: channel close ignored', error: e); } _channel = null; + _setupCompleter = null; _connectedPeers.clear(); _sessionId = null; diff --git a/test/watch_together/watch_together_peer_service_test.dart b/test/watch_together/watch_together_peer_service_test.dart new file mode 100644 index 00000000..d53d29d9 --- /dev/null +++ b/test/watch_together/watch_together_peer_service_test.dart @@ -0,0 +1,205 @@ +import 'dart:async'; +import 'dart:convert'; +import 'dart:io'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:plezy/watch_together/services/watch_together_peer_service.dart'; + +typedef _MessageHandler = FutureOr Function(int connection, WebSocket socket, Map message); + +class _RelayServer { + _RelayServer._(this._server, this._handler); + + final HttpServer _server; + final _MessageHandler _handler; + final List sockets = []; + final List>> messages = []; + + String get baseUrl => 'http://${_server.address.address}:${_server.port}'; + + static Future<_RelayServer> start(_MessageHandler handler) async { + final server = await HttpServer.bind(InternetAddress.loopbackIPv4, 0); + final relay = _RelayServer._(server, handler); + server.listen((request) async { + final socket = await WebSocketTransformer.upgrade(request); + final connection = relay.sockets.length; + relay.sockets.add(socket); + relay.messages.add([]); + socket.listen((data) { + final message = (jsonDecode(data as String) as Map).cast(); + relay.messages[connection].add(message); + relay._handler(connection, socket, message); + }); + }); + return relay; + } + + void send(WebSocket socket, Map message) => socket.add(jsonEncode(message)); + + Future close() async { + for (final socket in sockets) { + await socket.close(); + } + await _server.close(force: true); + } +} + +void main() { + final services = []; + final relays = <_RelayServer>[]; + + WatchTogetherPeerService serviceFor(_RelayServer relay) { + final service = WatchTogetherPeerService(customBaseUrl: relay.baseUrl); + services.add(service); + return service; + } + + Future<_RelayServer> relayWith(_MessageHandler handler) async { + final relay = await _RelayServer.start(handler); + relays.add(relay); + return relay; + } + + tearDown(() async { + for (final service in services.reversed) { + await service.disconnect(); + service.dispose(); + } + services.clear(); + for (final relay in relays.reversed) { + await relay.close(); + } + relays.clear(); + }); + + test('host connects, listens, and announces create with the existing wire format', () async { + late final _RelayServer relay; + relay = await relayWith((_, socket, message) { + if (message['type'] == 'create') { + relay.send(socket, {'type': 'created', 'sessionId': message['sessionId']}); + } + }); + final service = serviceFor(relay); + + expect(await service.createSession(sessionId: 'abc12'), 'ABC12'); + expect(service.isHost, isTrue); + expect(service.myPeerId, 'wt-ABC12'); + expect(relay.messages.single, [ + {'type': 'create', 'sessionId': 'ABC12', 'peerId': 'wt-ABC12'}, + ]); + }); + + test('guest connects, listens, and announces join with the existing wire format', () async { + late final _RelayServer relay; + relay = await relayWith((_, socket, message) { + if (message['type'] == 'join') { + relay.send(socket, { + 'type': 'joined', + 'sessionId': message['sessionId'], + 'peers': ['wt-ROOM1'], + }); + } + }); + final service = serviceFor(relay); + final connectedPeers = []; + final subscription = service.onPeerConnected.listen(connectedPeers.add); + addTearDown(subscription.cancel); + + await service.joinSession('room1'); + + expect(service.isHost, isFalse); + expect(service.connectedPeers, ['wt-ROOM1']); + expect(connectedPeers, ['wt-ROOM1']); + expect(relay.messages.single, [ + {'type': 'join', 'sessionId': 'ROOM1', 'peerId': service.myPeerId}, + ]); + }); + + test('host reconnect joins first and re-creates a missing room on the same socket', () async { + late final _RelayServer relay; + relay = await relayWith((connection, socket, message) { + if (connection == 0 && message['type'] == 'create') { + relay.send(socket, {'type': 'created', 'sessionId': message['sessionId']}); + } else if (connection == 1 && message['type'] == 'join') { + relay.send(socket, {'type': 'error', 'code': 'room_not_found', 'message': 'Room not found'}); + } else if (connection == 1 && message['type'] == 'create') { + relay.send(socket, {'type': 'created', 'sessionId': message['sessionId']}); + } + }); + final service = serviceFor(relay); + final reconnected = Completer(); + var reconnectCallbacks = 0; + service.onReconnected = () { + reconnectCallbacks++; + reconnected.complete(); + }; + + await service.createSession(sessionId: 'room2'); + await relay.sockets.single.close(); + await reconnected.future.timeout(const Duration(seconds: 6)); + + expect(reconnectCallbacks, 1); + expect(relay.sockets, hasLength(2)); + expect(relay.messages[0], [ + {'type': 'create', 'sessionId': 'ROOM2', 'peerId': 'wt-ROOM2'}, + ]); + expect(relay.messages[1], [ + {'type': 'join', 'sessionId': 'ROOM2', 'peerId': 'wt-ROOM2'}, + {'type': 'create', 'sessionId': 'ROOM2', 'peerId': 'wt-ROOM2'}, + ]); + }); + + test('setup preserves typed timeout and relay errors', () async { + final timeoutRelay = await relayWith((_, _, _) {}); + final timeoutService = serviceFor(timeoutRelay); + + await expectLater( + timeoutService.createSession(sessionId: 'slow1'), + throwsA( + isA() + .having((error) => error.type, 'type', PeerErrorType.timeout) + .having((error) => error.message, 'message', 'Timed out creating session'), + ), + ); + + late final _RelayServer errorRelay; + errorRelay = await relayWith((_, socket, message) { + errorRelay.send(socket, {'type': 'error', 'code': 'room_full', 'message': 'Room is full'}); + }); + final errorService = serviceFor(errorRelay); + + await expectLater( + errorService.joinSession('full1'), + throwsA( + isA() + .having((error) => error.type, 'type', PeerErrorType.serverError) + .having((error) => error.serverCode, 'serverCode', 'room_full'), + ), + ); + }); + + test('one setup installs one listener and sends one announcement', () async { + late final _RelayServer relay; + relay = await relayWith((_, socket, message) { + if (message['type'] == 'create') { + relay.send(socket, {'type': 'created', 'sessionId': message['sessionId']}); + } + }); + final service = serviceFor(relay); + final peerEvents = []; + final peerSeen = Completer(); + final subscription = service.onPeerConnected.listen((peerId) { + peerEvents.add(peerId); + if (!peerSeen.isCompleted) peerSeen.complete(); + }); + addTearDown(subscription.cancel); + + await service.createSession(sessionId: 'once1'); + relay.send(relay.sockets.single, {'type': 'peerJoined', 'peerId': 'guest-1'}); + await peerSeen.future.timeout(const Duration(seconds: 1)); + await Future.delayed(const Duration(milliseconds: 50)); + + expect(peerEvents, ['guest-1']); + expect(relay.messages.single.where((message) => message['type'] == 'create'), hasLength(1)); + }); +}