diff --git a/lib/watch_together/services/watch_together_sync_manager.dart b/lib/watch_together/services/watch_together_sync_manager.dart index de820719..f279a37b 100644 --- a/lib/watch_together/services/watch_together_sync_manager.dart +++ b/lib/watch_together/services/watch_together_sync_manager.dart @@ -42,25 +42,22 @@ class WatchTogetherSyncManager { static const Duration positionSyncInterval = Duration(seconds: 3); static const Duration excessiveDrift = Duration(seconds: 10); - // Buffering debounce constants - static const Duration _bufferingDebounceDelay = Duration(milliseconds: 500); + // Peer readiness state (peer ID -> hasPlayerReady) + final Map _peerReady = {}; - // Buffering debounce timer - prevents false pauses from brief buffering events - Timer? _bufferingDebounceTimer; + // Whether play is deferred until all peers are ready (initial load gate) + bool _deferredPlay = false; - // Peers with pending (debounced) buffering state - final Map _pendingBufferingState = {}; + // Position to seek to when deferred play triggers + Duration? _deferredPlayPosition; + + // Whether the first coordinated play has completed (after this, late joiners catch up via positionSync) + bool _firstPlayCompleted = false; // Track last known state to avoid duplicate broadcasts bool _lastKnownPlaying = false; double _lastKnownRate = 1.0; - // Track if we were playing before a peer started buffering (for auto-resume) - bool _wasPlayingBeforeBuffering = false; - - // Position to seek to when auto-resuming deferred playback - Duration? _pendingPlayPosition; - // Whether we've announced our player as ready (first buffering: false) bool _hasAnnouncedReady = false; @@ -68,12 +65,6 @@ class WatchTogetherSyncManager { SessionConfigCallback? onSessionConfigReceived; SyncStateCallback? onSyncStateChanged; - /// Participants' buffering states (peer ID -> isBuffering) - final Map _participantBuffering = {}; - - /// Participants' ready states (peer ID -> hasPlayerReady) - final Map _participantReady = {}; - WatchTogetherSyncManager({ required WatchTogetherPeerService peerService, required WatchSession session, @@ -90,26 +81,10 @@ class WatchTogetherSyncManager { /// Whether this manager has a player attached bool get hasPlayer => _player != null; - /// Whether any participant (including local player) is currently buffering - bool get isAnyBuffering => _participantBuffering.values.any((b) => b) || (_player?.state.buffering ?? false); - - /// Whether all participants have their player attached and ready - /// Returns true if: - /// - We're alone (no other peers tracked) - /// - All tracked participants have sent playerReady(true) + /// Whether all tracked peers have their player ready bool get isAllReady { - // If no other peers are tracked, we're ready (solo viewing) - if (_participantBuffering.isEmpty) { - return true; - } - // All peers in _participantBuffering must also be in _participantReady with value true - for (final peerId in _participantBuffering.keys) { - final ready = _participantReady[peerId]; - if (ready != true) { - return false; // Peer hasn't sent ready yet or sent ready(false) - } - } - return true; + if (_peerReady.isEmpty) return true; + return _peerReady.values.every((ready) => ready); } /// Whether sync is in progress (for UI indicator) @@ -155,15 +130,11 @@ class WatchTogetherSyncManager { if (peerId != _peerService.myPeerId) { if (_session.isHost) { // Host waits for each peer to load their video before allowing play. - _participantBuffering[peerId] = true; - _participantReady[peerId] = false; + _peerReady[peerId] = false; } else { // Guests use optimistic defaults — the host coordinates readiness - // and will broadcast pause/play as needed. Pessimistic defaults - // cause a deadlock because the host's playerReady/buffering - // messages arrive before the sync manager subscribes. - _participantBuffering[peerId] = false; - _participantReady[peerId] = true; + // and will broadcast pause/play as needed. + _peerReady[peerId] = true; } } } @@ -176,25 +147,24 @@ class WatchTogetherSyncManager { // Announce that our player is no longer ready if (_peerService.myPeerId != null) { _peerService.broadcast(SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: false)); - _participantReady[_peerService.myPeerId!] = false; + _peerReady[_peerService.myPeerId!] = false; } _hasAnnouncedReady = false; + _deferredPlay = false; + _deferredPlayPosition = null; _playingSubscription?.cancel(); _bufferingSubscription?.cancel(); _rateSubscription?.cancel(); _messageSubscription?.cancel(); _positionSyncTimer?.cancel(); - _bufferingDebounceTimer?.cancel(); _playingSubscription = null; _bufferingSubscription = null; _rateSubscription = null; _messageSubscription = null; _positionSyncTimer = null; - _bufferingDebounceTimer = null; - _pendingBufferingState.clear(); _player = null; appLogger.d('WatchTogether: Player detached'); } @@ -203,30 +173,27 @@ class WatchTogetherSyncManager { void _setupPlayerSubscriptions() { // Listen to playing state changes _playingSubscription = _player!.streams.playing.listen((isPlaying) async { - if (_isRemoteAction) return; // Skip if this change was caused by a remote action + if (_isRemoteAction) return; + if (isPlaying == _lastKnownPlaying) return; + _lastKnownPlaying = isPlaying; - if (isPlaying != _lastKnownPlaying) { - _lastKnownPlaying = isPlaying; - - // If trying to play, check if all peers are ready first - if (isPlaying && (!isAllReady || isAnyBuffering)) { - appLogger.d('WatchTogether: Deferring local play - waiting for all peers to be ready'); - _wasPlayingBeforeBuffering = true; - _pendingPlayPosition = _player?.state.position; - _isRemoteAction = true; - try { - await _player!.pause(); - _lastKnownPlaying = false; - } finally { - _isRemoteAction = false; - } - // Still broadcast so peers know we want to play - _broadcastPlayPause(true); - return; + if (isPlaying && !isAllReady && !_firstPlayCompleted) { + // Defer until all peers have loaded video (initial sync only) + _deferredPlay = true; + _deferredPlayPosition = _player?.state.position; + _isRemoteAction = true; + try { + await _player!.pause(); + _lastKnownPlaying = false; + } finally { + _isRemoteAction = false; } - - _broadcastPlayPause(isPlaying); + _broadcastPlayPause(true); + return; } + + if (!isPlaying) _deferredPlay = false; + _broadcastPlayPause(isPlaying); }); // Listen to buffering state changes @@ -236,24 +203,17 @@ class WatchTogetherSyncManager { // Announce ready when we stop buffering for the first time (video loaded) if (!isBuffering && !_hasAnnouncedReady) { _hasAnnouncedReady = true; - _participantReady[_peerService.myPeerId!] = true; + _peerReady[_peerService.myPeerId!] = true; _peerService.broadcast(SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: true)); appLogger.d('WatchTogether: Video loaded, announcing player ready'); - // If host, send session config now that video is loaded with correct position if (_session.isHost) { _sendSessionConfig(); } } + // Broadcast for UI (peer buffering indicators) — no playback control _peerService.broadcast(SyncMessage.buffering(isBuffering, peerId: _peerService.myPeerId)); - - // Check for auto-resume when LOCAL buffering stops - // This fixes the bug where subtitle loading would pause both users but - // only resume would trigger from remote buffering messages, not local - if (!isBuffering) { - await _checkAutoResume(); - } }); // Listen to rate changes @@ -387,7 +347,7 @@ class WatchTogetherSyncManager { appLogger.d('WatchTogether: Ignoring pause from non-host in hostOnly mode'); break; } - _wasPlayingBeforeBuffering = false; // User intentionally paused, don't auto-resume + _deferredPlay = false; await _applyRemotePause(); break; @@ -402,43 +362,7 @@ class WatchTogetherSyncManager { break; case SyncMessageType.buffering: - if (_player == null) { - appLogger.d('WatchTogether: Ignoring buffering message, player not attached'); - break; - } - if (message.peerId != null && message.bufferingState != null) { - final peerId = message.peerId!; - final isBuffering = message.bufferingState!; - - if (isBuffering) { - // Peer started buffering - use debounce to avoid false pauses - _pendingBufferingState[peerId] = true; - - // Cancel existing debounce timer if any - _bufferingDebounceTimer?.cancel(); - _bufferingDebounceTimer = Timer(_bufferingDebounceDelay, () async { - // Check if still pending after debounce delay - if (_pendingBufferingState[peerId] == true) { - _participantBuffering[peerId] = true; - _pendingBufferingState.remove(peerId); - - // Auto-pause when any peer starts buffering (sustained) - if (isAnyBuffering && _player != null && _player!.state.playing) { - _wasPlayingBeforeBuffering = true; - appLogger.d('WatchTogether: Peer buffering (sustained), pausing playback'); - await _applyRemotePause(); - } - } - }); - } else { - // Peer stopped buffering - cancel pending debounce and update immediately - _pendingBufferingState.remove(peerId); - _participantBuffering[peerId] = false; - - // Auto-resume when all peers stop buffering AND all ready - await _checkAutoResume(); - } - } + // Buffering state used for UI only, not playback control break; case SyncMessageType.positionSync: @@ -449,12 +373,10 @@ class WatchTogetherSyncManager { // This provides eventual consistency for play/pause state if (message.isPlaying != null && _player != null && !_session.isHost) { final localPlaying = _player!.state.playing; - if (message.isPlaying! && !localPlaying && !isAnyBuffering && isAllReady) { - // Host is playing but we're paused - sync up + if (message.isPlaying! && !localPlaying) { appLogger.d('WatchTogether: Play/pause state diverged, syncing to host (playing)'); await _applyRemotePlay(position: message.position); } else if (!message.isPlaying! && localPlaying) { - // Host is paused but we're playing - sync up appLogger.d('WatchTogether: Play/pause state diverged, syncing to host (paused)'); await _applyRemotePause(); } @@ -477,8 +399,7 @@ class WatchTogetherSyncManager { case SyncMessageType.leave: if (message.peerId != null) { - _participantBuffering.remove(message.peerId); - _participantReady.remove(message.peerId); + _peerReady.remove(message.peerId); } break; @@ -506,11 +427,15 @@ class WatchTogetherSyncManager { case SyncMessageType.playerReady: if (message.peerId != null) { - _participantReady[message.peerId!] = message.bufferingState ?? false; + _peerReady[message.peerId!] = message.bufferingState ?? false; appLogger.d('WatchTogether: Peer ${message.peerId} player ready: ${message.bufferingState}'); - // If we were waiting to play and all are now ready, start playback - await _checkAutoResume(); + if (_deferredPlay && isAllReady) { + _deferredPlay = false; + _firstPlayCompleted = true; + await _applyRemotePlay(position: _deferredPlayPosition); + _deferredPlayPosition = null; + } } break; @@ -528,26 +453,9 @@ class WatchTogetherSyncManager { Future _applyRemotePlay({Duration? position}) async { if (_player == null) return; - // If not all participants have their player ready, defer play - if (!isAllReady) { - appLogger.d('WatchTogether: Deferring play - waiting for all players to be ready'); - _wasPlayingBeforeBuffering = true; - if (position != null) _pendingPlayPosition = position; - return; // Will trigger when all players send playerReady - } - - // If anyone is buffering, defer play until all ready - if (isAnyBuffering) { - appLogger.d('WatchTogether: Deferring play - waiting for all peers to stop buffering'); - _wasPlayingBeforeBuffering = true; - if (position != null) _pendingPlayPosition = position; - return; // Auto-resume will trigger when buffering clears - } - appLogger.d('WatchTogether: Applying remote PLAY${position != null ? ' at ${position.inSeconds}s' : ''}'); _isRemoteAction = true; try { - // Seek to position first if provided if (position != null) { await _player!.seek(position); } @@ -594,21 +502,6 @@ class WatchTogetherSyncManager { } } - /// Check if conditions are met to auto-resume playback - /// Called when local or remote buffering stops, or when a peer becomes ready - Future _checkAutoResume() async { - if (_player == null) return; - if (!isAllReady) return; - if (isAnyBuffering) return; - if (_player!.state.playing) return; - if (!_wasPlayingBeforeBuffering) return; - - _wasPlayingBeforeBuffering = false; - appLogger.d('WatchTogether: All conditions met, auto-resuming playback'); - await _applyRemotePlay(position: _pendingPlayPosition); - _pendingPlayPosition = null; - } - /// Apply remote rate change Future _applyRemoteRate(double rate) async { if (_player == null) return; @@ -664,13 +557,9 @@ class WatchTogetherSyncManager { if (message.peerId != null) { if (_session.isHost) { - // Host waits for new peer to load their video before allowing play. - _participantBuffering[message.peerId!] = true; - _participantReady[message.peerId!] = false; - } else if (!_participantBuffering.containsKey(message.peerId!)) { - // Guests use optimistic defaults for peers they haven't seen yet. - _participantBuffering[message.peerId!] = false; - _participantReady[message.peerId!] = true; + _peerReady[message.peerId!] = false; + } else if (!_peerReady.containsKey(message.peerId!)) { + _peerReady[message.peerId!] = true; } } @@ -680,16 +569,8 @@ class WatchTogetherSyncManager { if (_hasAnnouncedReady) { _sendSessionConfig(toPeerId: message.peerId); - // Re-send our ready and buffering state so the new peer doesn't - // get stuck waiting for updates that were broadcast before it joined. _peerService.sendTo(message.peerId!, SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: true)); } - if (_player != null) { - _peerService.sendTo( - message.peerId!, - SyncMessage.buffering(_player!.state.buffering, peerId: _peerService.myPeerId), - ); - } // Send host's join info so guest adds host to their participants list _peerService.sendTo( message.peerId!, @@ -705,12 +586,9 @@ class WatchTogetherSyncManager { appLogger.d('WatchTogether: Received session config'); // The host only sends sessionConfig after its player is ready, so we - // can safely mark it as ready and not buffering. This prevents the - // guest from being permanently stuck waiting for ready/buffering - // messages that were broadcast before it joined. + // can safely mark it as ready. if (message.peerId != null) { - _participantReady[message.peerId!] = true; - _participantBuffering[message.peerId!] = false; + _peerReady[message.peerId!] = true; } // Update control mode @@ -736,14 +614,13 @@ class WatchTogetherSyncManager { // Match play/pause state (bufferingState is reused: false = playing) if (message.bufferingState == false) { - // Host was playing - defer play until all ready - _wasPlayingBeforeBuffering = true; - _pendingPlayPosition = message.position; - // Check if we can play now - if (isAllReady && !isAnyBuffering) { + // Host was playing — defer until our video is loaded + _deferredPlay = true; + _deferredPlayPosition = message.position; + if (_hasAnnouncedReady) { + _deferredPlay = false; + _firstPlayCompleted = true; await _applyRemotePlay(position: message.position); - } else { - appLogger.d('WatchTogether: Host was playing but deferring until all ready'); } } else { await _player!.pause(); @@ -808,9 +685,7 @@ class WatchTogetherSyncManager { /// Dispose resources void dispose() { detachPlayer(); - _participantBuffering.clear(); - _participantReady.clear(); - _pendingBufferingState.clear(); + _peerReady.clear(); _hasAnnouncedReady = false; } }