refactor(watch-together): share relay setup
This commit is contained in:
@@ -44,6 +44,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
|
|
||||||
WebSocketChannel? _channel;
|
WebSocketChannel? _channel;
|
||||||
StreamSubscription? _channelSubscription;
|
StreamSubscription? _channelSubscription;
|
||||||
|
Completer<void>? _setupCompleter;
|
||||||
final Set<String> _connectedPeers = {};
|
final Set<String> _connectedPeers = {};
|
||||||
String? _sessionId;
|
String? _sessionId;
|
||||||
String? _myPeerId;
|
String? _myPeerId;
|
||||||
@@ -123,31 +124,54 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
return channel;
|
return channel;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Connect, listen, and send a room setup announcement.
|
||||||
|
Future<Completer<void>> _connectAndAnnounce(String type) async {
|
||||||
|
final channel = await _connectToRelay();
|
||||||
|
_channel = channel;
|
||||||
|
|
||||||
|
_listenToChannel(channel);
|
||||||
|
startKeepalive();
|
||||||
|
|
||||||
|
return _announce(type);
|
||||||
|
}
|
||||||
|
|
||||||
|
Completer<void> _announce(String type) {
|
||||||
|
final completer = Completer<void>();
|
||||||
|
_setupCompleter = completer;
|
||||||
|
_sendRaw({'type': type, 'sessionId': _sessionId, 'peerId': _myPeerId});
|
||||||
|
return completer;
|
||||||
|
}
|
||||||
|
|
||||||
/// Listen on the channel stream and route incoming server messages.
|
/// Listen on the channel stream and route incoming server messages.
|
||||||
void _listenToChannel(WebSocketChannel channel, {Completer<void>? setupCompleter}) {
|
void _listenToChannel(WebSocketChannel channel) {
|
||||||
_channelSubscription?.cancel();
|
_channelSubscription?.cancel();
|
||||||
_channelSubscription = channel.stream.listen(
|
_channelSubscription = channel.stream.listen(
|
||||||
(data) {
|
(data) {
|
||||||
|
if (!identical(_channel, channel)) return;
|
||||||
resetPongTimer();
|
resetPongTimer();
|
||||||
_handleServerMessage(data as String, setupCompleter: setupCompleter);
|
_handleServerMessage(data as String);
|
||||||
},
|
},
|
||||||
onError: (error) {
|
onError: (error) {
|
||||||
|
if (!identical(_channel, channel)) return;
|
||||||
appLogger.e('WatchTogether: WebSocket error', error: error);
|
appLogger.e('WatchTogether: WebSocket error', error: error);
|
||||||
_safeAdd(
|
_safeAdd(
|
||||||
_errorController,
|
_errorController,
|
||||||
PeerError(type: PeerErrorType.serverError, message: 'WebSocket error: $error', originalError: error),
|
PeerError(type: PeerErrorType.serverError, message: 'WebSocket error: $error', originalError: error),
|
||||||
);
|
);
|
||||||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||||||
setupCompleter.completeError(error);
|
completer.completeError(error);
|
||||||
|
_setupCompleter = null;
|
||||||
}
|
}
|
||||||
_handleWebSocketClosed();
|
_handleWebSocketClosed();
|
||||||
},
|
},
|
||||||
onDone: () {
|
onDone: () {
|
||||||
|
if (!identical(_channel, channel)) return;
|
||||||
appLogger.w('WatchTogether: WebSocket closed');
|
appLogger.w('WatchTogether: WebSocket closed');
|
||||||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||||||
setupCompleter.completeError(
|
completer.completeError(
|
||||||
const PeerError(type: PeerErrorType.connectionFailed, message: 'WebSocket closed before setup completed'),
|
const PeerError(type: PeerErrorType.connectionFailed, message: 'WebSocket closed before setup completed'),
|
||||||
);
|
);
|
||||||
|
_setupCompleter = null;
|
||||||
}
|
}
|
||||||
_handleWebSocketClosed();
|
_handleWebSocketClosed();
|
||||||
},
|
},
|
||||||
@@ -155,7 +179,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Handle an incoming server message (JSON string).
|
/// Handle an incoming server message (JSON string).
|
||||||
void _handleServerMessage(String raw, {Completer<void>? setupCompleter}) {
|
void _handleServerMessage(String raw) {
|
||||||
try {
|
try {
|
||||||
final msg = jsonDecode(raw) as Map<String, dynamic>;
|
final msg = jsonDecode(raw) as Map<String, dynamic>;
|
||||||
final type = msg['type'] as String?;
|
final type = msg['type'] as String?;
|
||||||
@@ -164,8 +188,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
case 'created':
|
case 'created':
|
||||||
appLogger.d('WatchTogether: Room created: ${msg['sessionId']}');
|
appLogger.d('WatchTogether: Room created: ${msg['sessionId']}');
|
||||||
_safeAdd(_connectionStateController, true);
|
_safeAdd(_connectionStateController, true);
|
||||||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||||||
setupCompleter.complete();
|
completer.complete();
|
||||||
|
_setupCompleter = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
case 'joined':
|
case 'joined':
|
||||||
@@ -176,8 +201,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
_safeAdd(_peerConnectedController, peerId);
|
_safeAdd(_peerConnectedController, peerId);
|
||||||
}
|
}
|
||||||
_safeAdd(_connectionStateController, true);
|
_safeAdd(_connectionStateController, true);
|
||||||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||||||
setupCompleter.complete();
|
completer.complete();
|
||||||
|
_setupCompleter = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
case 'peerJoined':
|
case 'peerJoined':
|
||||||
@@ -220,8 +246,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
appLogger.e('WatchTogether: Server error: $code - $message');
|
appLogger.e('WatchTogether: Server error: $code - $message');
|
||||||
final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message', serverCode: code);
|
final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message', serverCode: code);
|
||||||
_safeAdd(_errorController, error);
|
_safeAdd(_errorController, error);
|
||||||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||||||
setupCompleter.completeError(error);
|
completer.completeError(error);
|
||||||
|
_setupCompleter = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
case 'pong':
|
case 'pong':
|
||||||
@@ -299,26 +326,17 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
_reconnectTimer?.cancel();
|
_reconnectTimer?.cancel();
|
||||||
_reconnectTimer = Timer(delay, () async {
|
_reconnectTimer = Timer(delay, () async {
|
||||||
try {
|
try {
|
||||||
final channel = await _connectToRelay();
|
|
||||||
_channel = channel;
|
|
||||||
|
|
||||||
final completer = Completer<void>();
|
|
||||||
_listenToChannel(channel, setupCompleter: completer);
|
|
||||||
startKeepalive();
|
|
||||||
|
|
||||||
// Always try join first — the room may still have peers (e.g. host
|
// Always try join first — the room may still have peers (e.g. host
|
||||||
// reconnecting while guests remain). Fall back to create only if
|
// reconnecting while guests remain). Fall back to create only if
|
||||||
// the room no longer exists and we were the host.
|
// the room no longer exists and we were the host.
|
||||||
_sendRaw({'type': 'join', 'sessionId': _sessionId, 'peerId': _myPeerId});
|
final completer = await _connectAndAnnounce('join');
|
||||||
|
|
||||||
try {
|
try {
|
||||||
await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect');
|
await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect');
|
||||||
} on PeerError catch (e) {
|
} on PeerError catch (e) {
|
||||||
if (_isHost && e.serverCode == 'room_not_found') {
|
if (_isHost && e.serverCode == 'room_not_found') {
|
||||||
appLogger.d('WatchTogether: Room gone, re-creating as host');
|
appLogger.d('WatchTogether: Room gone, re-creating as host');
|
||||||
final createCompleter = Completer<void>();
|
final createCompleter = _announce('create');
|
||||||
_listenToChannel(channel, setupCompleter: createCompleter);
|
|
||||||
_sendRaw({'type': 'create', 'sessionId': _sessionId, 'peerId': _myPeerId});
|
|
||||||
await createCompleter.future.namedTimeout(
|
await createCompleter.future.namedTimeout(
|
||||||
const Duration(seconds: 10),
|
const Duration(seconds: 10),
|
||||||
operation: 'WatchTogether reconnect create',
|
operation: 'WatchTogether reconnect create',
|
||||||
@@ -357,14 +375,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
_reconnectAttempts = 0;
|
_reconnectAttempts = 0;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
final channel = await _connectToRelay();
|
final completer = await _connectAndAnnounce('create');
|
||||||
_channel = channel;
|
|
||||||
|
|
||||||
final completer = Completer<void>();
|
|
||||||
_listenToChannel(channel, setupCompleter: completer);
|
|
||||||
startKeepalive();
|
|
||||||
|
|
||||||
_sendRaw({'type': 'create', 'sessionId': _sessionId, 'peerId': _myPeerId});
|
|
||||||
|
|
||||||
await completer.future.timeout(
|
await completer.future.timeout(
|
||||||
const Duration(seconds: 10),
|
const Duration(seconds: 10),
|
||||||
@@ -394,14 +405,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
_reconnectAttempts = 0;
|
_reconnectAttempts = 0;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
final channel = await _connectToRelay();
|
final completer = await _connectAndAnnounce('join');
|
||||||
_channel = channel;
|
|
||||||
|
|
||||||
final completer = Completer<void>();
|
|
||||||
_listenToChannel(channel, setupCompleter: completer);
|
|
||||||
startKeepalive();
|
|
||||||
|
|
||||||
_sendRaw({'type': 'join', 'sessionId': _sessionId, 'peerId': _myPeerId});
|
|
||||||
|
|
||||||
await completer.future.timeout(
|
await completer.future.timeout(
|
||||||
const Duration(seconds: 10),
|
const Duration(seconds: 10),
|
||||||
@@ -447,6 +451,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
|
|||||||
appLogger.d('WatchTogether: channel close ignored', error: e);
|
appLogger.d('WatchTogether: channel close ignored', error: e);
|
||||||
}
|
}
|
||||||
_channel = null;
|
_channel = null;
|
||||||
|
_setupCompleter = null;
|
||||||
|
|
||||||
_connectedPeers.clear();
|
_connectedPeers.clear();
|
||||||
_sessionId = null;
|
_sessionId = null;
|
||||||
|
|||||||
@@ -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<void> Function(int connection, WebSocket socket, Map<String, dynamic> message);
|
||||||
|
|
||||||
|
class _RelayServer {
|
||||||
|
_RelayServer._(this._server, this._handler);
|
||||||
|
|
||||||
|
final HttpServer _server;
|
||||||
|
final _MessageHandler _handler;
|
||||||
|
final List<WebSocket> sockets = [];
|
||||||
|
final List<List<Map<String, dynamic>>> 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<String, dynamic>();
|
||||||
|
relay.messages[connection].add(message);
|
||||||
|
relay._handler(connection, socket, message);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
return relay;
|
||||||
|
}
|
||||||
|
|
||||||
|
void send(WebSocket socket, Map<String, dynamic> message) => socket.add(jsonEncode(message));
|
||||||
|
|
||||||
|
Future<void> close() async {
|
||||||
|
for (final socket in sockets) {
|
||||||
|
await socket.close();
|
||||||
|
}
|
||||||
|
await _server.close(force: true);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
void main() {
|
||||||
|
final services = <WatchTogetherPeerService>[];
|
||||||
|
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 = <String>[];
|
||||||
|
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<void>();
|
||||||
|
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<PeerError>()
|
||||||
|
.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<PeerError>()
|
||||||
|
.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 = <String>[];
|
||||||
|
final peerSeen = Completer<void>();
|
||||||
|
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<void>.delayed(const Duration(milliseconds: 50));
|
||||||
|
|
||||||
|
expect(peerEvents, ['guest-1']);
|
||||||
|
expect(relay.messages.single.where((message) => message['type'] == 'create'), hasLength(1));
|
||||||
|
});
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user