refactor: simplify watch together buffering/ready system
This commit is contained in:
@@ -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<String, bool> _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<String, bool> _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<String, bool> _participantBuffering = {};
|
||||
|
||||
/// Participants' ready states (peer ID -> hasPlayerReady)
|
||||
final Map<String, bool> _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<void> _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<void> _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<void> _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;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user