Introduces shared seams for paginated views, D-pad reorder, media control routing, async singletons and the device method channel, then points the open-coded copies at them. Also removes unused models and duplicated provider/server plumbing, folds the twice-implemented artifact store in the server, and factors the repeated Flutter toolchain prologue in CI into a composite action.
873 lines
30 KiB
Dart
873 lines
30 KiB
Dart
import 'dart:async';
|
||
import 'dart:convert';
|
||
import 'dart:math';
|
||
|
||
import 'package:uuid/uuid.dart';
|
||
import '../../utils/future_extensions.dart';
|
||
import 'package:web_socket_channel/web_socket_channel.dart';
|
||
|
||
import '../../i18n/strings.g.dart';
|
||
import '../../services/base_peer_service.dart';
|
||
import '../../utils/app_logger.dart';
|
||
import '../models/sync_message.dart';
|
||
import 'relay_protocol.g.dart';
|
||
import 'watch_together_relay_endpoint.dart';
|
||
|
||
// Re-export so existing callers that import from here keep working.
|
||
export '../../services/base_peer_service.dart' show PeerError, PeerErrorType;
|
||
|
||
class _PinnedHostChangedError extends PeerError {
|
||
const _PinnedHostChangedError()
|
||
: super(type: PeerErrorType.serverError, message: 'Relay returned an invalid joined response');
|
||
}
|
||
|
||
/// Service for managing Watch Together connections via a WebSocket relay
|
||
///
|
||
/// This service handles:
|
||
/// - Creating sessions (as host)
|
||
/// - Joining sessions (as guest)
|
||
/// - Sending/receiving sync messages through the relay server
|
||
/// - Reconnection on WebSocket drops
|
||
class WatchTogetherPeerService with KeepaliveMixin {
|
||
final WatchTogetherRelayEndpoint endpoint;
|
||
final WebSocketChannel Function(Uri uri) _channelFactory;
|
||
|
||
/// Test synchronization point after relay setup succeeds but before the
|
||
/// reconnect is published to consumers.
|
||
final Future<void> Function()? debugReconnectSetupSucceededBarrier;
|
||
|
||
WatchTogetherPeerService({
|
||
WatchTogetherRelayEndpoint? endpoint,
|
||
this.debugReconnectSetupSucceededBarrier,
|
||
WebSocketChannel Function(Uri uri)? debugChannelFactory,
|
||
}) : endpoint = endpoint ?? WatchTogetherRelayEndpoint.defaultEndpoint,
|
||
_channelFactory = debugChannelFactory ?? ((uri) => WebSocketChannel.connect(uri));
|
||
static const int _relayProtocolVersion = RelayProtocol.protocolVersion;
|
||
|
||
WebSocketChannel? _channel;
|
||
StreamSubscription? _channelSubscription;
|
||
Completer<void>? _setupCompleter;
|
||
String? _setupRequestType;
|
||
final Set<String> _connectedPeers = {};
|
||
String? _sessionId;
|
||
String? _myPeerId;
|
||
bool _isHost = false;
|
||
String? _reconnectToken;
|
||
String? _hostPeerId;
|
||
|
||
// Stream controllers for events
|
||
final _peerConnectedController = StreamController<String>.broadcast();
|
||
final _peerDisconnectedController = StreamController<String>.broadcast();
|
||
final _messageReceivedController = StreamController<SyncMessage>.broadcast();
|
||
final _errorController = StreamController<PeerError>.broadcast();
|
||
final _connectionStateController = StreamController<bool>.broadcast();
|
||
final _sessionEndedController = StreamController<void>.broadcast();
|
||
|
||
// Reconnection state
|
||
int _reconnectAttempts = 0;
|
||
static const int _maxReconnectAttempts = 3;
|
||
Timer? _reconnectTimer;
|
||
int _connectionEpoch = 0;
|
||
bool _disposed = false;
|
||
bool _initialSetupInProgress = false;
|
||
bool _teardownInProgress = false;
|
||
Future<void>? _releaseFuture;
|
||
|
||
/// Called after a successful reconnection so the provider can re-announce join.
|
||
void Function()? onReconnected;
|
||
|
||
// Keepalive (via KeepaliveMixin)
|
||
@override
|
||
Duration get pingInterval => const Duration(seconds: 15);
|
||
@override
|
||
Duration get pongTimeout => const Duration(seconds: 30);
|
||
|
||
void _safeAdd<T>(StreamController<T> controller, T event) {
|
||
if (!controller.isClosed) controller.add(event);
|
||
}
|
||
|
||
/// Stream of peer IDs when a new peer connects
|
||
Stream<String> get onPeerConnected => _peerConnectedController.stream;
|
||
|
||
/// Stream of peer IDs when a peer disconnects
|
||
Stream<String> get onPeerDisconnected => _peerDisconnectedController.stream;
|
||
|
||
/// Stream of sync messages received from peers
|
||
Stream<SyncMessage> get onMessageReceived => _messageReceivedController.stream;
|
||
|
||
/// Stream of errors
|
||
Stream<PeerError> get onError => _errorController.stream;
|
||
|
||
/// Stream of connection state changes (true = connected, false = disconnected)
|
||
Stream<bool> get onConnectionStateChanged => _connectionStateController.stream;
|
||
|
||
/// Emitted when the host has durably ended the relay room.
|
||
Stream<void> get onSessionEnded => _sessionEndedController.stream;
|
||
|
||
/// Current session ID (null if not in a session)
|
||
String? get sessionId => _sessionId;
|
||
|
||
/// This peer's ID
|
||
String? get myPeerId => _myPeerId;
|
||
|
||
/// Relay-declared peer ID whose messages carry host authority.
|
||
String? get hostPeerId => _hostPeerId;
|
||
|
||
/// Whether this peer is the host
|
||
bool get isHost => _isHost;
|
||
|
||
/// Whether currently connected to a session
|
||
bool get isConnected => _channel != null && _connectedPeers.isNotEmpty;
|
||
|
||
/// List of connected peer IDs
|
||
List<String> get connectedPeers => _connectedPeers.toList();
|
||
|
||
/// Generate a short, readable session ID (5 alphanumeric chars)
|
||
static String _generateSessionId() {
|
||
const chars = 'ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789';
|
||
final random = Random.secure();
|
||
return String.fromCharCodes(List.generate(5, (_) => chars.codeUnitAt(random.nextInt(chars.length))));
|
||
}
|
||
|
||
/// Mint the reconnect capability before setup so a lost setup ACK can be
|
||
/// retried without relying on server-returned state.
|
||
static String _mintReconnectToken() {
|
||
final random = Random.secure();
|
||
final bytes = List<int>.generate(32, (_) => random.nextInt(256), growable: false);
|
||
return base64Url.encode(bytes).replaceAll('=', '');
|
||
}
|
||
|
||
/// Connect to the relay WebSocket and set up the message listener.
|
||
/// Returns a completer that completes when the expected response arrives.
|
||
Future<WebSocketChannel> _connectToRelay({Duration? timeout, String operation = 'WatchTogether connect'}) async {
|
||
final channel = _channelFactory(endpoint.webSocketUri);
|
||
|
||
try {
|
||
final ready = channel.ready;
|
||
if (timeout == null) {
|
||
await ready;
|
||
} else {
|
||
await ready.namedTimeout(timeout, operation: operation);
|
||
}
|
||
return channel;
|
||
} catch (_) {
|
||
try {
|
||
unawaited(channel.sink.close());
|
||
} catch (error) {
|
||
appLogger.d('WatchTogether: pending channel close ignored', error: error);
|
||
}
|
||
rethrow;
|
||
}
|
||
}
|
||
|
||
/// Connect, listen, and send a room setup announcement.
|
||
Future<Completer<void>> _connectAndAnnounce(
|
||
String type,
|
||
int epoch, {
|
||
Duration? connectTimeout,
|
||
String connectOperation = 'WatchTogether connect',
|
||
}) async {
|
||
final channel = await _connectToRelay(timeout: connectTimeout, operation: connectOperation);
|
||
if (_disposed || epoch != _connectionEpoch || _sessionId == null) {
|
||
unawaited(channel.sink.close());
|
||
throw StateError('Watch Together connection attempt became stale');
|
||
}
|
||
_channel = channel;
|
||
|
||
_listenToChannel(channel);
|
||
startKeepalive();
|
||
|
||
return _announce(type);
|
||
}
|
||
|
||
Completer<void> _announce(String type) {
|
||
final completer = Completer<void>();
|
||
_setupCompleter = completer;
|
||
_setupRequestType = type;
|
||
final reconnectToken = _reconnectToken;
|
||
_sendRaw({
|
||
'type': type,
|
||
'sessionId': _sessionId,
|
||
'peerId': _myPeerId,
|
||
'reconnectToken': ?reconnectToken,
|
||
'protocolVersion': _relayProtocolVersion,
|
||
});
|
||
return completer;
|
||
}
|
||
|
||
/// Listen on the channel stream and route incoming server messages.
|
||
void _listenToChannel(WebSocketChannel channel) {
|
||
_channelSubscription?.cancel();
|
||
_channelSubscription = channel.stream.listen(
|
||
(data) {
|
||
if (!identical(_channel, channel)) return;
|
||
resetPongTimer();
|
||
_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 case final completer? when !completer.isCompleted) {
|
||
completer.completeError(error);
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
}
|
||
_handleWebSocketClosed();
|
||
},
|
||
onDone: () {
|
||
if (!identical(_channel, channel)) return;
|
||
appLogger.w('WatchTogether: WebSocket closed');
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
completer.completeError(
|
||
const PeerError(type: PeerErrorType.connectionFailed, message: 'WebSocket closed before setup completed'),
|
||
);
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
}
|
||
_handleWebSocketClosed();
|
||
},
|
||
);
|
||
}
|
||
|
||
static final RegExp _reconnectTokenPattern = RegExp(r'^[A-Za-z0-9_-]{43}$');
|
||
|
||
PeerError _invalidSetupResponse(String type) {
|
||
return PeerError(type: PeerErrorType.serverError, message: 'Relay returned an invalid $type response');
|
||
}
|
||
|
||
List<String> _acceptSetupResponse(Map<String, dynamic> msg, String type) {
|
||
final responseSessionId = msg['sessionId'];
|
||
final hostPeerId = msg['hostPeerId'];
|
||
final reconnectToken = msg['reconnectToken'];
|
||
final protocolVersion = msg['protocolVersion'];
|
||
if (responseSessionId is! String ||
|
||
responseSessionId != _sessionId ||
|
||
hostPeerId is! String ||
|
||
!RelayProtocol.isValidPeerId(hostPeerId) ||
|
||
(_isHost && hostPeerId != _myPeerId) ||
|
||
reconnectToken is! String ||
|
||
reconnectToken != _reconnectToken ||
|
||
!_reconnectTokenPattern.hasMatch(reconnectToken) ||
|
||
protocolVersion != _relayProtocolVersion) {
|
||
throw _invalidSetupResponse(type);
|
||
}
|
||
|
||
final establishedHostPeerId = _hostPeerId;
|
||
if (!_isHost && _reconnectToken != null && establishedHostPeerId != null && hostPeerId != establishedHostPeerId) {
|
||
throw const _PinnedHostChangedError();
|
||
}
|
||
|
||
final rawPeers = msg['peers'];
|
||
final peers = <String>[];
|
||
if (rawPeers != null) {
|
||
if (rawPeers is! List) throw _invalidSetupResponse(type);
|
||
for (final peerId in rawPeers) {
|
||
if (peerId is! String || !RelayProtocol.isValidPeerId(peerId)) {
|
||
throw _invalidSetupResponse(type);
|
||
}
|
||
peers.add(peerId);
|
||
}
|
||
}
|
||
|
||
_hostPeerId = hostPeerId;
|
||
_reconnectToken = reconnectToken;
|
||
return peers;
|
||
}
|
||
|
||
void _acceptTeardownResponse(Map<String, dynamic> msg, String type) {
|
||
if (msg['sessionId'] != _sessionId || msg['protocolVersion'] != _relayProtocolVersion) {
|
||
throw _invalidSetupResponse(type);
|
||
}
|
||
if (type == RelayProtocol.left && msg['peerId'] != _myPeerId) {
|
||
throw _invalidSetupResponse(type);
|
||
}
|
||
}
|
||
|
||
void _failSetup(PeerError error) {
|
||
_safeAdd(_errorController, error);
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.completeError(error);
|
||
}
|
||
}
|
||
|
||
void _rejectAdmittedGuestSetup(_PinnedHostChangedError error) {
|
||
_safeAdd(_errorController, error);
|
||
final rejectedSetup = _setupCompleter;
|
||
if (rejectedSetup == null || rejectedSetup.isCompleted) return;
|
||
|
||
final leaveCompleter = _announce(RelayProtocol.leave);
|
||
unawaited(() async {
|
||
try {
|
||
await leaveCompleter.future.namedTimeout(
|
||
const Duration(seconds: 10),
|
||
operation: 'WatchTogether rejected reconnect leave',
|
||
);
|
||
} catch (releaseError) {
|
||
appLogger.d('WatchTogether: rejected reconnect leave ignored', error: releaseError);
|
||
} finally {
|
||
if (identical(_setupCompleter, leaveCompleter)) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
}
|
||
if (!rejectedSetup.isCompleted) rejectedSetup.completeError(error);
|
||
}
|
||
}());
|
||
}
|
||
|
||
bool _isExhaustedGuestReconnectRoomNotFound(String code) =>
|
||
code == RelayProtocol.roomNotFoundCode &&
|
||
!_isHost &&
|
||
!_initialSetupInProgress &&
|
||
!_teardownInProgress &&
|
||
_hostPeerId != null &&
|
||
_setupRequestType == RelayProtocol.join &&
|
||
_reconnectAttempts >= _maxReconnectAttempts;
|
||
|
||
void _handleGuestSessionEnded() {
|
||
_teardownInProgress = true;
|
||
++_connectionEpoch;
|
||
_reconnectTimer?.cancel();
|
||
_reconnectTimer = null;
|
||
stopKeepalive();
|
||
final error = const PeerError(type: PeerErrorType.invalidSession, message: 'Watch Together session ended');
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
_safeAdd(_errorController, error);
|
||
completer.completeError(error);
|
||
}
|
||
_safeAdd(_sessionEndedController, null);
|
||
_handleWebSocketClosed();
|
||
}
|
||
|
||
/// Handle an incoming server message (JSON string).
|
||
void _handleServerMessage(String raw) {
|
||
try {
|
||
final msg = jsonDecode(raw) as Map<String, dynamic>;
|
||
final type = msg['type'] as String?;
|
||
|
||
switch (type) {
|
||
case RelayProtocol.created:
|
||
late final List<String> peers;
|
||
try {
|
||
peers = _acceptSetupResponse(msg, RelayProtocol.created);
|
||
} on _PinnedHostChangedError catch (error) {
|
||
_rejectAdmittedGuestSetup(error);
|
||
break;
|
||
} on PeerError catch (error) {
|
||
_failSetup(error);
|
||
break;
|
||
}
|
||
appLogger.d('WatchTogether: Room created: ${msg['sessionId']} with peers: $peers');
|
||
for (final peerId in peers) {
|
||
if (_connectedPeers.add(peerId)) {
|
||
_safeAdd(_peerConnectedController, peerId);
|
||
}
|
||
}
|
||
_safeAdd(_connectionStateController, true);
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.complete();
|
||
}
|
||
|
||
case RelayProtocol.joined:
|
||
late final List<String> peers;
|
||
try {
|
||
peers = _acceptSetupResponse(msg, RelayProtocol.joined);
|
||
} on _PinnedHostChangedError catch (error) {
|
||
_rejectAdmittedGuestSetup(error);
|
||
break;
|
||
} on PeerError catch (error) {
|
||
_failSetup(error);
|
||
break;
|
||
}
|
||
appLogger.d('WatchTogether: Joined room ${msg['sessionId']} with peers: $peers');
|
||
for (final peerId in peers) {
|
||
_connectedPeers.add(peerId);
|
||
_safeAdd(_peerConnectedController, peerId);
|
||
}
|
||
_safeAdd(_connectionStateController, true);
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.complete();
|
||
}
|
||
|
||
case RelayProtocol.peerJoined:
|
||
final peerId = msg['peerId'] as String;
|
||
appLogger.d('WatchTogether: Peer joined: $peerId');
|
||
_connectedPeers.add(peerId);
|
||
_safeAdd(_peerConnectedController, peerId);
|
||
_safeAdd(_connectionStateController, true);
|
||
|
||
case RelayProtocol.peerLeft:
|
||
final peerId = msg['peerId'] as String;
|
||
appLogger.d('WatchTogether: Peer left: $peerId');
|
||
_connectedPeers.remove(peerId);
|
||
_safeAdd(_peerDisconnectedController, peerId);
|
||
if (_connectedPeers.isEmpty) {
|
||
_safeAdd(_connectionStateController, false);
|
||
}
|
||
|
||
case RelayProtocol.message:
|
||
final payload = msg['payload'];
|
||
final serverFrom = msg['from'] as String?;
|
||
if (payload != null) {
|
||
try {
|
||
final payloadStr = payload is String ? payload : jsonEncode(payload);
|
||
var syncMsg = SyncMessage.fromJson(payloadStr);
|
||
// Use the server-authenticated sender ID instead of the
|
||
// self-reported peerId in the payload to prevent spoofing.
|
||
if (serverFrom != null && syncMsg.peerId != serverFrom) {
|
||
syncMsg = syncMsg.copyWith(peerId: serverFrom);
|
||
}
|
||
_safeAdd(_messageReceivedController, syncMsg);
|
||
} catch (e) {
|
||
appLogger.e('WatchTogether: Failed to parse sync message payload', error: e);
|
||
}
|
||
}
|
||
|
||
case RelayProtocol.left:
|
||
try {
|
||
_acceptTeardownResponse(msg, RelayProtocol.left);
|
||
} on PeerError catch (error) {
|
||
_failSetup(error);
|
||
break;
|
||
}
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.complete();
|
||
}
|
||
|
||
case RelayProtocol.ended:
|
||
try {
|
||
_acceptTeardownResponse(msg, RelayProtocol.ended);
|
||
} on PeerError catch (error) {
|
||
_failSetup(error);
|
||
break;
|
||
}
|
||
final expectedTeardown =
|
||
_setupRequestType == RelayProtocol.endSession || _setupRequestType == RelayProtocol.leave;
|
||
if (expectedTeardown) {
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.complete();
|
||
}
|
||
} else if (!_isHost) {
|
||
_handleGuestSessionEnded();
|
||
} else {
|
||
_failSetup(_invalidSetupResponse(RelayProtocol.ended));
|
||
}
|
||
|
||
case RelayProtocol.error:
|
||
final code = msg['code'] as String? ?? 'unknown';
|
||
final message = msg['message'] as String? ?? t.common.unknown;
|
||
if (_isExhaustedGuestReconnectRoomNotFound(code)) {
|
||
appLogger.d('WatchTogether: Room gone after reconnect retries; guest session ended');
|
||
_handleGuestSessionEnded();
|
||
break;
|
||
}
|
||
appLogger.e('WatchTogether: Server error: $code - $message');
|
||
final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message', serverCode: code);
|
||
_safeAdd(_errorController, error);
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
completer.completeError(error);
|
||
}
|
||
|
||
case RelayProtocol.pong:
|
||
// Handled by resetPongTimer() already
|
||
break;
|
||
|
||
default:
|
||
appLogger.w('WatchTogether: Unknown server message type: $type');
|
||
}
|
||
} catch (_) {
|
||
appLogger.e('WatchTogether: Failed to parse server message');
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
_failSetup(
|
||
const PeerError(type: PeerErrorType.serverError, message: 'Relay returned an invalid setup response'),
|
||
);
|
||
}
|
||
}
|
||
}
|
||
|
||
@override
|
||
void sendPing() => _sendRaw({'type': RelayProtocol.ping});
|
||
|
||
@override
|
||
void onPongTimeout() {
|
||
appLogger.w('WatchTogether: Pong timeout — closing WebSocket');
|
||
try {
|
||
_channel?.sink.close();
|
||
} catch (e) {
|
||
appLogger.d('WatchTogether: pong-timeout close ignored', error: e);
|
||
}
|
||
}
|
||
|
||
/// Send a raw JSON map to the relay.
|
||
void _sendRaw(Map<String, dynamic> msg) {
|
||
try {
|
||
_channel?.sink.add(jsonEncode(msg));
|
||
} catch (e) {
|
||
appLogger.e('WatchTogether: Failed to send message', error: e);
|
||
}
|
||
}
|
||
|
||
/// Handle the WebSocket being closed unexpectedly — attempt reconnection.
|
||
void _handleWebSocketClosed() {
|
||
final channel = _channel;
|
||
final shouldReconnect = !_initialSetupInProgress && !_teardownInProgress;
|
||
if (shouldReconnect) ++_connectionEpoch;
|
||
stopKeepalive();
|
||
unawaited(_channelSubscription?.cancel());
|
||
_channelSubscription = null;
|
||
_channel = null;
|
||
if (channel != null) unawaited(channel.sink.close());
|
||
|
||
for (final peerId in _connectedPeers.toList()) {
|
||
_safeAdd(_peerDisconnectedController, peerId);
|
||
}
|
||
_connectedPeers.clear();
|
||
_safeAdd(_connectionStateController, false);
|
||
|
||
if (shouldReconnect && !_disposed && _sessionId != null) {
|
||
_attemptReconnect(_connectionEpoch);
|
||
}
|
||
}
|
||
|
||
/// Attempt to reconnect to the relay and re-join/re-create the room.
|
||
void _attemptReconnect(int epoch) {
|
||
if (_disposed || epoch != _connectionEpoch || _sessionId == null) return;
|
||
if (_reconnectAttempts >= _maxReconnectAttempts) {
|
||
appLogger.e('WatchTogether: Max reconnect attempts reached');
|
||
_safeAdd(
|
||
_errorController,
|
||
const PeerError(
|
||
type: PeerErrorType.connectionFailed,
|
||
message: 'Lost connection to relay after multiple reconnect attempts',
|
||
),
|
||
);
|
||
return;
|
||
}
|
||
|
||
_reconnectAttempts++;
|
||
final delay = Duration(seconds: _reconnectAttempts * 2);
|
||
appLogger.d('WatchTogether: Reconnect attempt $_reconnectAttempts/$_maxReconnectAttempts in ${delay.inSeconds}s');
|
||
|
||
_reconnectTimer?.cancel();
|
||
_reconnectTimer = Timer(delay, () async {
|
||
if (_disposed || epoch != _connectionEpoch || _sessionId == null) return;
|
||
try {
|
||
final completer = await _connectAndAnnounce(RelayProtocol.join, epoch);
|
||
if (_disposed || epoch != _connectionEpoch) return;
|
||
|
||
try {
|
||
await completer.future.namedTimeout(const Duration(seconds: 10), operation: 'WatchTogether reconnect');
|
||
} on PeerError catch (e) {
|
||
if (_disposed || epoch != _connectionEpoch) return;
|
||
if (_isHost && e.serverCode == RelayProtocol.roomNotFoundCode) {
|
||
appLogger.d('WatchTogether: Room gone, re-creating as host');
|
||
final createCompleter = _announce(RelayProtocol.create);
|
||
await createCompleter.future.namedTimeout(
|
||
const Duration(seconds: 10),
|
||
operation: 'WatchTogether reconnect create',
|
||
);
|
||
} else {
|
||
rethrow;
|
||
}
|
||
}
|
||
|
||
await debugReconnectSetupSucceededBarrier?.call();
|
||
|
||
if (_disposed || epoch != _connectionEpoch) return;
|
||
_reconnectAttempts = 0;
|
||
appLogger.d('WatchTogether: Reconnected successfully');
|
||
try {
|
||
onReconnected?.call();
|
||
} catch (e) {
|
||
appLogger.e('WatchTogether: Reconnect callback failed', error: e);
|
||
}
|
||
} catch (e) {
|
||
if (_disposed || epoch != _connectionEpoch) return;
|
||
appLogger.e('WatchTogether: Reconnect failed', error: e);
|
||
_handleWebSocketClosed();
|
||
}
|
||
});
|
||
}
|
||
|
||
bool _isRetryableInitialSetupError(Object error) =>
|
||
error is TimeoutException ||
|
||
(error is PeerError &&
|
||
(error.type == PeerErrorType.connectionFailed ||
|
||
error.type == PeerErrorType.networkError ||
|
||
error.type == PeerErrorType.timeout)) ||
|
||
error is! PeerError;
|
||
|
||
Future<void> _resetTransportForInitialRetry() async {
|
||
stopKeepalive();
|
||
final subscription = _channelSubscription;
|
||
final channel = _channel;
|
||
_channelSubscription = null;
|
||
_channel = null;
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
await subscription?.cancel();
|
||
try {
|
||
await channel?.sink.close();
|
||
} catch (error) {
|
||
appLogger.d('WatchTogether: setup retry close ignored', error: error);
|
||
}
|
||
}
|
||
|
||
Future<void> _performInitialSetup(String type, int epoch, PeerError timeoutError) async {
|
||
_initialSetupInProgress = true;
|
||
try {
|
||
for (var attempt = 0; attempt < _maxReconnectAttempts; attempt++) {
|
||
try {
|
||
final completer = await _connectAndAnnounce(type, epoch);
|
||
await completer.future.timeout(const Duration(seconds: 10), onTimeout: () => throw timeoutError);
|
||
return;
|
||
} catch (error) {
|
||
if (_disposed || epoch != _connectionEpoch) rethrow;
|
||
await _resetTransportForInitialRetry();
|
||
if (!_isRetryableInitialSetupError(error) || attempt + 1 >= _maxReconnectAttempts) {
|
||
rethrow;
|
||
}
|
||
await Future<void>.delayed(Duration(milliseconds: 250 * (attempt + 1)));
|
||
}
|
||
}
|
||
} finally {
|
||
_initialSetupInProgress = false;
|
||
}
|
||
}
|
||
|
||
Future<void> _bestEffortReleaseFailedSetup(Object setupError) async {
|
||
if (!_isRetryableInitialSetupError(setupError)) return;
|
||
try {
|
||
await releaseSession();
|
||
} catch (releaseError) {
|
||
appLogger.d('WatchTogether: failed setup reservation release ignored', error: releaseError);
|
||
}
|
||
}
|
||
|
||
/// Create a new session as host
|
||
///
|
||
/// Returns the session ID that others can use to join.
|
||
/// If [sessionId] is provided, uses that instead of generating a new one.
|
||
Future<String> createSession({String? sessionId}) async {
|
||
if (_sessionId != null || _channel != null) {
|
||
await disconnect();
|
||
}
|
||
|
||
final resolvedSessionId = sessionId?.toUpperCase() ?? _generateSessionId();
|
||
if (!RelayProtocol.isValidSessionId(resolvedSessionId)) {
|
||
throw ArgumentError.value(
|
||
sessionId,
|
||
'sessionId',
|
||
'Must be 1–${RelayProtocol.maxSessionIdLength} letters, digits, _ or -',
|
||
);
|
||
}
|
||
_isHost = true;
|
||
_sessionId = resolvedSessionId;
|
||
_myPeerId = const Uuid().v4();
|
||
_reconnectToken = _mintReconnectToken();
|
||
_reconnectAttempts = 0;
|
||
final epoch = ++_connectionEpoch;
|
||
|
||
try {
|
||
await _performInitialSetup(
|
||
RelayProtocol.create,
|
||
epoch,
|
||
const PeerError(type: PeerErrorType.timeout, message: 'Timed out creating session'),
|
||
);
|
||
|
||
appLogger.d('WatchTogether: Session created: $_sessionId');
|
||
return _sessionId!;
|
||
} catch (e) {
|
||
appLogger.e('WatchTogether: Failed to create session', error: e);
|
||
await _bestEffortReleaseFailedSetup(e);
|
||
await disconnect();
|
||
rethrow;
|
||
}
|
||
}
|
||
|
||
/// Join an existing session as guest.
|
||
Future<void> joinSession(String sessionId) async {
|
||
if (_sessionId != null || _channel != null) {
|
||
await disconnect();
|
||
}
|
||
|
||
final resolvedSessionId = sessionId.toUpperCase();
|
||
if (!RelayProtocol.isValidSessionId(resolvedSessionId)) {
|
||
throw ArgumentError.value(
|
||
sessionId,
|
||
'sessionId',
|
||
'Must be 1–${RelayProtocol.maxSessionIdLength} letters, digits, _ or -',
|
||
);
|
||
}
|
||
_isHost = false;
|
||
_sessionId = resolvedSessionId;
|
||
_myPeerId = const Uuid().v4();
|
||
_reconnectToken = _mintReconnectToken();
|
||
_reconnectAttempts = 0;
|
||
final epoch = ++_connectionEpoch;
|
||
|
||
try {
|
||
await _performInitialSetup(
|
||
RelayProtocol.join,
|
||
epoch,
|
||
PeerError(type: PeerErrorType.timeout, message: t.watchTogether.failedToJoin),
|
||
);
|
||
|
||
appLogger.d('WatchTogether: Joined session: $_sessionId');
|
||
} catch (e) {
|
||
appLogger.e('WatchTogether: Failed to join session', error: e);
|
||
await _bestEffortReleaseFailedSetup(e);
|
||
await disconnect();
|
||
rethrow;
|
||
}
|
||
}
|
||
|
||
/// Broadcast a message to all connected peers
|
||
void broadcast(SyncMessage message) {
|
||
final payload = message.toJson();
|
||
_sendRaw({'type': RelayProtocol.broadcast, 'payload': payload});
|
||
}
|
||
|
||
/// Send a message to a specific peer
|
||
void sendTo(String peerId, SyncMessage message) {
|
||
if (!RelayProtocol.isValidPeerId(peerId)) {
|
||
throw ArgumentError.value(peerId, 'peerId', 'Must be 1–${RelayProtocol.maxPeerIdLength} letters, digits, _ or -');
|
||
}
|
||
final payload = message.toJson();
|
||
_sendRaw({'type': RelayProtocol.sendTo, 'to': peerId, 'payload': payload});
|
||
}
|
||
|
||
/// Explicitly release this peer's relay ownership. Hosts destroy the room;
|
||
/// guests release their reserved reconnect identity. If transport was lost,
|
||
/// authenticate a fresh connection first so an intentional exit is not
|
||
/// mistaken for a transient disconnect.
|
||
Future<void> releaseSession() {
|
||
final active = _releaseFuture;
|
||
if (active != null) return active;
|
||
final operation = _releaseSession();
|
||
_releaseFuture = operation;
|
||
return operation.whenComplete(() {
|
||
if (identical(_releaseFuture, operation)) _releaseFuture = null;
|
||
});
|
||
}
|
||
|
||
Future<void> _releaseSession() async {
|
||
if (_sessionId == null || _myPeerId == null || _reconnectToken == null) return;
|
||
|
||
_teardownInProgress = true;
|
||
_reconnectTimer?.cancel();
|
||
_reconnectTimer = null;
|
||
final epoch = ++_connectionEpoch;
|
||
try {
|
||
if (_setupCompleter case final completer? when !completer.isCompleted) {
|
||
await _resetTransportForInitialRetry();
|
||
}
|
||
for (var attempt = 0; attempt < _maxReconnectAttempts; attempt++) {
|
||
try {
|
||
if (_channel == null) {
|
||
final reconnectCompleter = await _connectAndAnnounce(
|
||
RelayProtocol.join,
|
||
epoch,
|
||
connectTimeout: const Duration(seconds: 10),
|
||
connectOperation: 'WatchTogether release reconnect',
|
||
);
|
||
await reconnectCompleter.future.namedTimeout(
|
||
const Duration(seconds: 10),
|
||
operation: 'WatchTogether release reconnect',
|
||
);
|
||
}
|
||
|
||
final releaseCompleter = _announce(_isHost ? RelayProtocol.endSession : RelayProtocol.leave);
|
||
await releaseCompleter.future.namedTimeout(
|
||
const Duration(seconds: 10),
|
||
operation: _isHost ? 'WatchTogether end session' : 'WatchTogether leave session',
|
||
);
|
||
return;
|
||
} catch (error) {
|
||
if (error is PeerError &&
|
||
(error.serverCode == RelayProtocol.roomNotFoundCode ||
|
||
error.serverCode == RelayProtocol.notInRoomCode ||
|
||
(!_isHost && error.serverCode == RelayProtocol.peerIdUnavailableCode))) {
|
||
return;
|
||
}
|
||
await _resetTransportForInitialRetry();
|
||
if (!_isRetryableInitialSetupError(error) || attempt + 1 >= _maxReconnectAttempts) {
|
||
rethrow;
|
||
}
|
||
await Future<void>.delayed(Duration(milliseconds: 250 * (attempt + 1)));
|
||
}
|
||
}
|
||
} finally {
|
||
_teardownInProgress = false;
|
||
}
|
||
}
|
||
|
||
/// Close and forget local relay state without sending a release. Established
|
||
/// intentional exits call [releaseSession] before this cleanup step.
|
||
Future<void> disconnect() async {
|
||
appLogger.d('WatchTogether: Disconnecting...');
|
||
++_connectionEpoch;
|
||
_reconnectTimer?.cancel();
|
||
_reconnectTimer = null;
|
||
stopKeepalive();
|
||
|
||
final subscription = _channelSubscription;
|
||
final channel = _channel;
|
||
_channelSubscription = null;
|
||
_channel = null;
|
||
final setupCompleter = _setupCompleter;
|
||
_setupCompleter = null;
|
||
_setupRequestType = null;
|
||
if (setupCompleter != null && !setupCompleter.isCompleted) {
|
||
setupCompleter.completeError(StateError('Watch Together connection cancelled'));
|
||
}
|
||
_connectedPeers.clear();
|
||
_sessionId = null;
|
||
_myPeerId = null;
|
||
_reconnectToken = null;
|
||
_hostPeerId = null;
|
||
_isHost = false;
|
||
_reconnectAttempts = 0;
|
||
_teardownInProgress = false;
|
||
|
||
unawaited(subscription?.cancel());
|
||
try {
|
||
await channel?.sink.close();
|
||
} catch (e) {
|
||
appLogger.d('WatchTogether: channel close ignored', error: e);
|
||
}
|
||
_safeAdd(_connectionStateController, false);
|
||
}
|
||
|
||
/// Dispose all resources.
|
||
void dispose() {
|
||
if (_disposed) return;
|
||
_disposed = true;
|
||
unawaited(disconnect());
|
||
|
||
_peerConnectedController.close();
|
||
_peerDisconnectedController.close();
|
||
_messageReceivedController.close();
|
||
_errorController.close();
|
||
_connectionStateController.close();
|
||
_sessionEndedController.close();
|
||
}
|
||
}
|