fix: watch together stream use after close

This commit is contained in:
edde746
2026-03-04 00:20:48 +01:00
parent a9e4091500
commit efac1a9f28
@@ -48,6 +48,10 @@ class WatchTogetherPeerService with KeepaliveMixin {
@override @override
Duration get pongTimeout => const Duration(seconds: 30); 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 of peer IDs when a new peer connects
Stream<String> get onPeerConnected => _peerConnectedController.stream; Stream<String> get onPeerConnected => _peerConnectedController.stream;
@@ -105,7 +109,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
}, },
onError: (error) { onError: (error) {
appLogger.e('WatchTogether: WebSocket error', error: error); appLogger.e('WatchTogether: WebSocket error', error: error);
_errorController.add( _safeAdd(_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 != null && !setupCompleter.isCompleted) {
@@ -134,7 +138,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
switch (type) { switch (type) {
case 'created': case 'created':
appLogger.d('WatchTogether: Room created: ${msg['sessionId']}'); appLogger.d('WatchTogether: Room created: ${msg['sessionId']}');
_connectionStateController.add(true); _safeAdd(_connectionStateController, true);
if (setupCompleter != null && !setupCompleter.isCompleted) { if (setupCompleter != null && !setupCompleter.isCompleted) {
setupCompleter.complete(); setupCompleter.complete();
} }
@@ -144,9 +148,9 @@ class WatchTogetherPeerService with KeepaliveMixin {
appLogger.d('WatchTogether: Joined room ${msg['sessionId']} with peers: $peers'); appLogger.d('WatchTogether: Joined room ${msg['sessionId']} with peers: $peers');
for (final peerId in peers) { for (final peerId in peers) {
_connectedPeers.add(peerId); _connectedPeers.add(peerId);
_peerConnectedController.add(peerId); _safeAdd(_peerConnectedController, peerId);
} }
_connectionStateController.add(true); _safeAdd(_connectionStateController, true);
if (setupCompleter != null && !setupCompleter.isCompleted) { if (setupCompleter != null && !setupCompleter.isCompleted) {
setupCompleter.complete(); setupCompleter.complete();
} }
@@ -155,16 +159,16 @@ class WatchTogetherPeerService with KeepaliveMixin {
final peerId = msg['peerId'] as String; final peerId = msg['peerId'] as String;
appLogger.d('WatchTogether: Peer joined: $peerId'); appLogger.d('WatchTogether: Peer joined: $peerId');
_connectedPeers.add(peerId); _connectedPeers.add(peerId);
_peerConnectedController.add(peerId); _safeAdd(_peerConnectedController, peerId);
_connectionStateController.add(true); _safeAdd(_connectionStateController, true);
case 'peerLeft': case 'peerLeft':
final peerId = msg['peerId'] as String; final peerId = msg['peerId'] as String;
appLogger.d('WatchTogether: Peer left: $peerId'); appLogger.d('WatchTogether: Peer left: $peerId');
_connectedPeers.remove(peerId); _connectedPeers.remove(peerId);
_peerDisconnectedController.add(peerId); _safeAdd(_peerDisconnectedController, peerId);
if (_connectedPeers.isEmpty) { if (_connectedPeers.isEmpty) {
_connectionStateController.add(false); _safeAdd(_connectionStateController, false);
} }
case 'message': case 'message':
@@ -175,7 +179,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
final payloadStr = payload is String ? payload : jsonEncode(payload); final payloadStr = payload is String ? payload : jsonEncode(payload);
final syncMsg = SyncMessage.fromJson(payloadStr); final syncMsg = SyncMessage.fromJson(payloadStr);
appLogger.d('WatchTogether: Received ${syncMsg.type} from $from'); appLogger.d('WatchTogether: Received ${syncMsg.type} from $from');
_messageReceivedController.add(syncMsg); _safeAdd(_messageReceivedController, syncMsg);
} catch (e) { } catch (e) {
appLogger.e('WatchTogether: Failed to parse sync message payload', error: e); appLogger.e('WatchTogether: Failed to parse sync message payload', error: e);
} }
@@ -186,7 +190,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
final message = msg['message'] as String? ?? 'Unknown error'; final message = msg['message'] as String? ?? 'Unknown error';
appLogger.e('WatchTogether: Server error: $code - $message'); appLogger.e('WatchTogether: Server error: $code - $message');
final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message'); final error = PeerError(type: PeerErrorType.serverError, message: '$code: $message');
_errorController.add(error); _safeAdd(_errorController, error);
if (setupCompleter != null && !setupCompleter.isCompleted) { if (setupCompleter != null && !setupCompleter.isCompleted) {
setupCompleter.completeError(error); setupCompleter.completeError(error);
} }
@@ -232,10 +236,10 @@ class WatchTogetherPeerService with KeepaliveMixin {
// Notify peers lost // Notify peers lost
for (final peerId in _connectedPeers.toList()) { for (final peerId in _connectedPeers.toList()) {
_peerDisconnectedController.add(peerId); _safeAdd(_peerDisconnectedController, peerId);
} }
_connectedPeers.clear(); _connectedPeers.clear();
_connectionStateController.add(false); _safeAdd(_connectionStateController, false);
// Attempt to reconnect if we had a session // Attempt to reconnect if we had a session
if (_sessionId != null) { if (_sessionId != null) {
@@ -247,7 +251,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
void _attemptReconnect() { void _attemptReconnect() {
if (_reconnectAttempts >= _maxReconnectAttempts) { if (_reconnectAttempts >= _maxReconnectAttempts) {
appLogger.e('WatchTogether: Max reconnect attempts reached'); appLogger.e('WatchTogether: Max reconnect attempts reached');
_errorController.add( _safeAdd(_errorController,
const PeerError( const PeerError(
type: PeerErrorType.connectionFailed, type: PeerErrorType.connectionFailed,
message: 'Lost connection to relay after multiple reconnect attempts', message: 'Lost connection to relay after multiple reconnect attempts',
@@ -398,7 +402,7 @@ class WatchTogetherPeerService with KeepaliveMixin {
_isHost = false; _isHost = false;
_reconnectAttempts = 0; _reconnectAttempts = 0;
_connectionStateController.add(false); _safeAdd(_connectionStateController, false);
} }
/// Dispose all resources /// Dispose all resources