From 45770e1ad695f653281bde01d3059c1f2161d2d5 Mon Sep 17 00:00:00 2001 From: edde746 <86283021+edde746@users.noreply.github.com> Date: Sat, 2 May 2026 22:31:57 +0200 Subject: [PATCH] fix: unblock watch together resume close #961 --- .../providers/watch_together_provider.dart | 1 + .../services/watch_together_sync_manager.dart | 35 ++- .../watch_together_sync_manager_test.dart | 229 ++++++++++++++++++ 3 files changed, 256 insertions(+), 9 deletions(-) create mode 100644 test/watch_together/watch_together_sync_manager_test.dart diff --git a/lib/watch_together/providers/watch_together_provider.dart b/lib/watch_together/providers/watch_together_provider.dart index c11cbc6d..4d0a9702 100644 --- a/lib/watch_together/providers/watch_together_provider.dart +++ b/lib/watch_together/providers/watch_together_provider.dart @@ -432,6 +432,7 @@ class WatchTogetherProvider with ChangeNotifier { final disconnectedName = _participants.where((p) => p.peerId == peerId).map((p) => p.displayName).firstOrNull; _participants.removeWhere((p) => p.peerId == peerId); + unawaited(_syncManager?.handlePeerDisconnected(peerId)); // If host disconnected unexpectedly, start grace period for reconnection. // Skip if the host already sent a deliberate leave message. diff --git a/lib/watch_together/services/watch_together_sync_manager.dart b/lib/watch_together/services/watch_together_sync_manager.dart index fe7c2872..10e8f0d6 100644 --- a/lib/watch_together/services/watch_together_sync_manager.dart +++ b/lib/watch_together/services/watch_together_sync_manager.dart @@ -121,6 +121,9 @@ class WatchTogetherSyncManager { _player = player; _lastKnownPlaying = player.state.playing; _lastKnownRate = player.state.rate; + if (player.state.playing) { + _firstPlayCompleted = true; + } _setupPlayerSubscriptions(); _setupMessageSubscription(); @@ -174,6 +177,14 @@ class WatchTogetherSyncManager { appLogger.d('WatchTogether: Initialized $otherCount existing participants (host=${_session.isHost})'); } + /// Remove readiness tracking for a peer that dropped at the relay level. + Future handlePeerDisconnected(String peerId) async { + if (_peerReady.remove(peerId) != null) { + appLogger.d('WatchTogether: Removed disconnected peer readiness: $peerId'); + await _resumeDeferredPlayIfReady(_playerAttachmentGeneration); + } + } + /// Detach the player and stop sync void detachPlayer() { _playerAttachmentGeneration++; @@ -631,6 +642,7 @@ class WatchTogetherSyncManager { case SyncMessageType.leave: if (message.peerId != null) { _peerReady.remove(message.peerId); + await _resumeDeferredPlayIfReady(queuedAttachmentGeneration); } break; @@ -668,15 +680,7 @@ class WatchTogetherSyncManager { _peerReady[message.peerId!] = message.bufferingState ?? false; appLogger.d('WatchTogether: Peer ${message.peerId} player ready: ${message.bufferingState}'); - if (_deferredPlay && isAllReady) { - _setDeferredPlay(false); - _firstPlayCompleted = true; - final pos = _deferredPlayPosition; - _deferredPlayPosition = null; - await _applyRemotePlay(position: pos, expectedAttachmentGeneration: queuedAttachmentGeneration); - // Broadcast play to all peers now that everyone is ready - _broadcastPlayPause(true); - } + await _resumeDeferredPlayIfReady(queuedAttachmentGeneration); } break; @@ -714,12 +718,25 @@ class WatchTogetherSyncManager { ); if (!didPlay) return false; + _firstPlayCompleted = true; _lastKnownPlaying = true; return true; }, ); } + Future _resumeDeferredPlayIfReady(int expectedAttachmentGeneration) async { + if (!_deferredPlay || !isAllReady) return; + + _setDeferredPlay(false); + _firstPlayCompleted = true; + final pos = _deferredPlayPosition; + _deferredPlayPosition = null; + await _applyRemotePlay(position: pos, expectedAttachmentGeneration: expectedAttachmentGeneration); + // Broadcast play to all peers now that everyone is ready. + _broadcastPlayPause(true); + } + /// Apply remote pause command Future _applyRemotePause({int? expectedAttachmentGeneration}) async { return _runGuardedRemoteAction( diff --git a/test/watch_together/watch_together_sync_manager_test.dart b/test/watch_together/watch_together_sync_manager_test.dart new file mode 100644 index 00000000..e596272d --- /dev/null +++ b/test/watch_together/watch_together_sync_manager_test.dart @@ -0,0 +1,229 @@ +import 'dart:async'; + +import 'package:flutter_test/flutter_test.dart'; +import 'package:plezy/mpv/mpv.dart'; +import 'package:plezy/watch_together/models/sync_message.dart'; +import 'package:plezy/watch_together/models/watch_session.dart'; +import 'package:plezy/watch_together/services/watch_together_peer_service.dart'; +import 'package:plezy/watch_together/services/watch_together_sync_manager.dart'; + +void main() { + group('WatchTogetherSyncManager deferred play', () { + test('does not re-enter initial load gate after attaching an already-playing player', () async { + final peerService = _FakeWatchTogetherPeerService(peerId: 'host'); + final player = _FakePlayer(playing: true, position: const Duration(minutes: 3)); + final manager = _hostManager(peerService); + final deferredStates = []; + manager.onDeferredPlayChanged = deferredStates.add; + + manager.initializeParticipants(['host', 'guest']); + manager.attachPlayer(player); + + await player.emitPlaying(false); + await player.emitPlaying(true); + + expect(deferredStates, isNot(contains(true))); + expect(player.state.playing, isTrue); + + manager.dispose(); + await player.dispose(); + await peerService.close(); + }); + + test('remote play completion prevents a later local resume from using the initial load gate', () async { + final peerService = _FakeWatchTogetherPeerService(peerId: 'guest'); + final player = _FakePlayer(playing: false, position: const Duration(seconds: 10)); + final manager = _guestManager(peerService, controlMode: ControlMode.anyone); + final deferredStates = []; + manager.onDeferredPlayChanged = deferredStates.add; + + manager.initializeParticipants(['guest', 'host', 'other']); + manager.attachPlayer(player); + + peerService.emit(SyncMessage.playerReady(peerId: 'other', ready: false)); + await _settle(); + + peerService.emit(SyncMessage.play(peerId: 'host', position: const Duration(seconds: 20))); + await _settle(); + expect(player.state.playing, isTrue); + + await player.emitPlaying(false); + await player.emitPlaying(true); + + expect(deferredStates, isNot(contains(true))); + expect(player.state.playing, isTrue); + + manager.dispose(); + await player.dispose(); + await peerService.close(); + }); + + test('removing a disconnected not-ready peer resumes deferred play', () async { + final peerService = _FakeWatchTogetherPeerService(peerId: 'host'); + final player = _FakePlayer(playing: false, position: const Duration(minutes: 5)); + final manager = _hostManager(peerService); + final deferredStates = []; + manager.onDeferredPlayChanged = deferredStates.add; + + manager.initializeParticipants(['host', 'guest']); + manager.attachPlayer(player); + + await player.emitPlaying(true); + + expect(deferredStates, contains(true)); + expect(player.state.playing, isFalse); + + await manager.handlePeerDisconnected('guest'); + + expect(deferredStates, containsAllInOrder([true, false])); + expect(player.state.playing, isTrue); + expect(peerService.broadcasts.where((m) => m.type == SyncMessageType.play), isNotEmpty); + + manager.dispose(); + await player.dispose(); + await peerService.close(); + }); + }); +} + +WatchTogetherSyncManager _hostManager(_FakeWatchTogetherPeerService peerService) { + return WatchTogetherSyncManager( + peerService: peerService, + session: const WatchSession( + sessionId: 'ROOM1', + role: SessionRole.host, + controlMode: ControlMode.hostOnly, + state: SessionState.connected, + hostPeerId: 'host', + ), + displayName: 'Host', + ); +} + +WatchTogetherSyncManager _guestManager(_FakeWatchTogetherPeerService peerService, {required ControlMode controlMode}) { + return WatchTogetherSyncManager( + peerService: peerService, + session: WatchSession( + sessionId: 'ROOM1', + role: SessionRole.guest, + controlMode: controlMode, + state: SessionState.connected, + hostPeerId: 'host', + ), + displayName: 'Guest', + ); +} + +Future _settle() async { + await Future.delayed(Duration.zero); + await Future.delayed(Duration.zero); +} + +class _FakeWatchTogetherPeerService extends WatchTogetherPeerService { + _FakeWatchTogetherPeerService({required this.peerId}) : super(customBaseUrl: 'http://localhost'); + + final String peerId; + final StreamController _messages = StreamController.broadcast(); + final List broadcasts = []; + final Map> sentMessages = {}; + + @override + String? get myPeerId => peerId; + + @override + Stream get onMessageReceived => _messages.stream; + + @override + void broadcast(SyncMessage message) { + broadcasts.add(message); + } + + @override + void sendTo(String peerId, SyncMessage message) { + sentMessages.putIfAbsent(peerId, () => []).add(message); + } + + void emit(SyncMessage message) { + _messages.add(message); + } + + Future close() => _messages.close(); +} + +class _FakePlayer implements Player { + _FakePlayer({bool playing = false, Duration position = Duration.zero}) + : _state = PlayerState(playing: playing, buffering: false, position: position); + + PlayerState _state; + bool _disposed = false; + + final StreamController _playingController = StreamController.broadcast(); + final StreamController _bufferingController = StreamController.broadcast(); + final StreamController _rateController = StreamController.broadcast(); + + @override + PlayerState get state => _state; + + @override + PlayerStreams get streams => PlayerStreams( + playing: _playingController.stream, + completed: const Stream.empty(), + buffering: _bufferingController.stream, + position: const Stream.empty(), + duration: const Stream.empty(), + seekable: const Stream.empty(), + buffer: const Stream.empty(), + volume: const Stream.empty(), + rate: _rateController.stream, + tracks: const Stream.empty(), + track: const Stream.empty(), + log: const Stream.empty(), + error: const Stream.empty(), + audioDevice: const Stream.empty(), + audioDevices: const Stream>.empty(), + bufferRanges: const Stream>.empty(), + playbackRestart: const Stream.empty(), + backendSwitched: const Stream.empty(), + ); + + Future emitPlaying(bool value) async { + _state = _state.copyWith(playing: value); + _playingController.add(value); + await _settle(); + } + + @override + Future play() async { + _state = _state.copyWith(playing: true); + } + + @override + Future pause() async { + _state = _state.copyWith(playing: false); + } + + @override + Future seek(Duration position) async { + _state = _state.copyWith(position: position); + } + + @override + Future setRate(double rate) async { + _state = _state.copyWith(rate: rate); + _rateController.add(rate); + } + + @override + bool get disposed => _disposed; + + @override + Future dispose() async { + _disposed = true; + await _playingController.close(); + await _bufferingController.close(); + await _rateController.close(); + } + + @override + dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); +}