fix: harden watch together teardown sync

This commit is contained in:
edde746
2026-03-06 12:52:28 +01:00
parent c9a1bc97c4
commit cf8f1849aa
@@ -1,5 +1,7 @@
import 'dart:async';
import 'package:flutter/services.dart';
import '../../mpv/mpv.dart';
import '../../utils/app_logger.dart';
import '../models/sync_message.dart';
@@ -30,6 +32,8 @@ class WatchTogetherSyncManager {
// Stream subscriptions (cancelled together in detachPlayer)
final List<StreamSubscription<dynamic>> _subscriptions = [];
Future<void> _messageQueue = Future.value();
int _playerAttachmentGeneration = 0;
// Position sync timer (host broadcasts position periodically)
Timer? _positionSyncTimer;
@@ -104,6 +108,7 @@ class WatchTogetherSyncManager {
detachPlayer();
}
_playerAttachmentGeneration++;
_player = player;
_lastKnownPlaying = player.state.playing;
_lastKnownRate = player.state.rate;
@@ -151,6 +156,11 @@ class WatchTogetherSyncManager {
/// Detach the player and stop sync
void detachPlayer() {
_playerAttachmentGeneration++;
_player = null;
_isRemoteAction = false;
_setSyncing(false);
// Announce that our player is no longer ready
if (_peerService.myPeerId != null) {
_peerService.broadcast(SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: false));
@@ -168,80 +178,111 @@ class WatchTogetherSyncManager {
_hasClockOffset = false;
_pendingPingTimestamp = null;
for (final subscription in _subscriptions) {
subscription.cancel();
}
final subscriptions = List<StreamSubscription<dynamic>>.of(_subscriptions);
_subscriptions.clear();
for (final subscription in subscriptions) {
unawaited(subscription.cancel());
}
_positionSyncTimer?.cancel();
_positionSyncTimer = null;
_player = null;
appLogger.d('WatchTogether: Player detached');
}
/// Set up subscriptions to player streams
void _setupPlayerSubscriptions() {
// Listen to playing state changes
_subscriptions.add(_player!.streams.playing.listen((isPlaying) async {
if (_isRemoteAction) return;
if (isPlaying == _lastKnownPlaying) return;
_lastKnownPlaying = isPlaying;
_subscriptions.add(
_player!.streams.playing.listen((isPlaying) async {
if (_isRemoteAction) return;
if (isPlaying == _lastKnownPlaying) return;
final player = _player;
final attachmentGeneration = _playerAttachmentGeneration;
if (player == null || !_isPlayerAttachmentCurrent(player, attachmentGeneration)) 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;
_lastKnownPlaying = isPlaying;
if (isPlaying && !isAllReady && !_firstPlayCompleted) {
// Defer until all peers have loaded video (initial sync only)
_deferredPlay = true;
_deferredPlayPosition = player.state.position;
_isRemoteAction = true;
try {
final didPause = await _runGuardedPlayerCommand(
actionName: 'deferred play pause',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.pause(),
);
if (!didPause) return;
_lastKnownPlaying = false;
} finally {
_isRemoteAction = false;
}
_broadcastPlayPause(true);
return;
}
_broadcastPlayPause(true);
return;
}
if (!isPlaying) _deferredPlay = false;
_broadcastPlayPause(isPlaying);
}));
if (!isPlaying) _deferredPlay = false;
_broadcastPlayPause(isPlaying);
}),
);
// Listen to buffering state changes
_subscriptions.add(_player!.streams.buffering.listen((isBuffering) async {
if (_isRemoteAction) return;
_subscriptions.add(
_player!.streams.buffering.listen((isBuffering) async {
if (_isRemoteAction) return;
// Announce ready when we stop buffering for the first time (video loaded)
if (!isBuffering && !_hasAnnouncedReady) {
_hasAnnouncedReady = true;
_peerReady[_peerService.myPeerId!] = true;
_peerService.broadcast(SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: true));
appLogger.d('WatchTogether: Video loaded, announcing player ready');
// Announce ready when we stop buffering for the first time (video loaded)
if (!isBuffering && !_hasAnnouncedReady) {
_hasAnnouncedReady = true;
_peerReady[_peerService.myPeerId!] = true;
_peerService.broadcast(SyncMessage.playerReady(peerId: _peerService.myPeerId!, ready: true));
appLogger.d('WatchTogether: Video loaded, announcing player ready');
if (_session.isHost) {
_sendSessionConfig();
if (_session.isHost) {
_sendSessionConfig();
}
}
}
// Broadcast for UI (peer buffering indicators) — no playback control
_peerService.broadcast(SyncMessage.buffering(isBuffering, peerId: _peerService.myPeerId));
}));
// Broadcast for UI (peer buffering indicators) — no playback control
_peerService.broadcast(SyncMessage.buffering(isBuffering, peerId: _peerService.myPeerId));
}),
);
// Listen to rate changes
_subscriptions.add(_player!.streams.rate.listen((rate) {
if (_isRemoteAction) return;
_subscriptions.add(
_player!.streams.rate.listen((rate) {
if (_isRemoteAction) return;
if (rate != _lastKnownRate) {
_lastKnownRate = rate;
if (_canControl()) {
_peerService.broadcast(SyncMessage.rate(rate, peerId: _peerService.myPeerId));
if (rate != _lastKnownRate) {
_lastKnownRate = rate;
if (_canControl()) {
_peerService.broadcast(SyncMessage.rate(rate, peerId: _peerService.myPeerId));
}
}
}
}));
}),
);
}
/// Set up subscription to incoming sync messages
void _setupMessageSubscription() {
_subscriptions.add(_peerService.onMessageReceived.listen(_handleMessage));
_subscriptions.add(
_peerService.onMessageReceived.listen((message) {
final queuedAttachmentGeneration = _playerAttachmentGeneration;
_messageQueue = _messageQueue.then((_) => _handleMessage(message, queuedAttachmentGeneration)).catchError((
Object error,
StackTrace stackTrace,
) {
appLogger.e(
'WatchTogether: Failed to handle ${message.type.name} message',
error: error,
stackTrace: stackTrace,
);
});
}),
);
}
/// Start periodic position sync (host only)
@@ -341,6 +382,126 @@ class WatchTogetherSyncManager {
return message.peerId == _session.hostPeerId;
}
bool _isRecoverablePlayerException(PlatformException error) {
return error.code == 'COMMAND_FAILED' || error.code == 'NOT_INITIALIZED';
}
bool _isPlayerAttachmentCurrent(Player player, int attachmentGeneration) {
return identical(_player, player) && _playerAttachmentGeneration == attachmentGeneration;
}
bool _isPlayerAttachmentUsable(Player player, int attachmentGeneration) {
return _isPlayerAttachmentCurrent(player, attachmentGeneration) && !player.disposed;
}
void _handleRecoverableRemoteActionFailure(
String actionName,
Object error, {
required Player player,
required int attachmentGeneration,
}) {
appLogger.w('WatchTogether: Remote $actionName skipped because player became unavailable', error: error);
if (_isPlayerAttachmentCurrent(player, attachmentGeneration)) {
detachPlayer();
}
}
Future<bool> _runGuardedPlayerCommand({
required String actionName,
required Player player,
required int attachmentGeneration,
required Future<void> Function(Player player) command,
}) async {
if (!_isPlayerAttachmentUsable(player, attachmentGeneration)) {
if (_isPlayerAttachmentCurrent(player, attachmentGeneration)) {
_handleRecoverableRemoteActionFailure(
actionName,
StateError('Player became unavailable'),
player: player,
attachmentGeneration: attachmentGeneration,
);
}
return false;
}
try {
await command(player);
} on StateError catch (e) {
_handleRecoverableRemoteActionFailure(actionName, e, player: player, attachmentGeneration: attachmentGeneration);
return false;
} on PlatformException catch (e) {
if (_isRecoverablePlayerException(e)) {
_handleRecoverableRemoteActionFailure(
actionName,
e,
player: player,
attachmentGeneration: attachmentGeneration,
);
return false;
}
rethrow;
}
if (!_isPlayerAttachmentUsable(player, attachmentGeneration)) {
if (_isPlayerAttachmentCurrent(player, attachmentGeneration)) {
_handleRecoverableRemoteActionFailure(
actionName,
StateError('Player became unavailable'),
player: player,
attachmentGeneration: attachmentGeneration,
);
}
return false;
}
return true;
}
Future<bool> _runGuardedRemoteAction({
required String actionName,
required Future<bool> Function(Player player, int attachmentGeneration) action,
int? expectedAttachmentGeneration,
}) async {
final player = _player;
if (player == null) return false;
final attachmentGeneration = _playerAttachmentGeneration;
if (expectedAttachmentGeneration != null && expectedAttachmentGeneration != attachmentGeneration) {
return false;
}
if (!_isPlayerAttachmentUsable(player, attachmentGeneration)) {
_handleRecoverableRemoteActionFailure(
actionName,
StateError('Player became unavailable'),
player: player,
attachmentGeneration: attachmentGeneration,
);
return false;
}
_isRemoteAction = true;
try {
return await action(player, attachmentGeneration);
} on StateError catch (e) {
_handleRecoverableRemoteActionFailure(actionName, e, player: player, attachmentGeneration: attachmentGeneration);
return false;
} on PlatformException catch (e) {
if (_isRecoverablePlayerException(e)) {
_handleRecoverableRemoteActionFailure(
actionName,
e,
player: player,
attachmentGeneration: attachmentGeneration,
);
return false;
}
rethrow;
} finally {
_isRemoteAction = false;
}
}
/// Broadcast play/pause state
void _broadcastPlayPause(bool isPlaying) {
if (!_canControl()) {
@@ -367,7 +528,7 @@ class WatchTogetherSyncManager {
}
/// Handle incoming sync messages
void _handleMessage(SyncMessage message) async {
Future<void> _handleMessage(SyncMessage message, int queuedAttachmentGeneration) async {
// Ignore our own messages
if (message.peerId == _peerService.myPeerId) {
return;
@@ -412,7 +573,7 @@ class WatchTogetherSyncManager {
appLogger.d('WatchTogether: Ignoring play from non-host in hostOnly mode');
break;
}
await _applyRemotePlay(position: message.position);
await _applyRemotePlay(position: message.position, expectedAttachmentGeneration: queuedAttachmentGeneration);
break;
case SyncMessageType.pause:
@@ -421,7 +582,7 @@ class WatchTogetherSyncManager {
break;
}
_deferredPlay = false;
await _applyRemotePause();
await _applyRemotePause(expectedAttachmentGeneration: queuedAttachmentGeneration);
break;
case SyncMessageType.seek:
@@ -430,7 +591,7 @@ class WatchTogetherSyncManager {
break;
}
if (message.position != null) {
await _applyRemoteSeek(message.position!);
await _applyRemoteSeek(message.position!, expectedAttachmentGeneration: queuedAttachmentGeneration);
}
break;
@@ -440,18 +601,26 @@ class WatchTogetherSyncManager {
case SyncMessageType.positionSync:
if (message.position != null) {
_checkAndCorrectDrift(message.position!, message.timestamp);
await _checkAndCorrectDrift(message.position!, message.timestamp, queuedAttachmentGeneration);
}
// Reconcile play/pause state if host sent it and we diverged
// This provides eventual consistency for play/pause state
if (message.isPlaying != null && _player != null && !_session.isHost) {
final localPlaying = _player!.state.playing;
final player = _player;
if (message.isPlaying != null && player != null && !_session.isHost) {
if (!_isPlayerAttachmentCurrent(player, queuedAttachmentGeneration)) {
break;
}
final localPlaying = player.state.playing;
if (message.isPlaying! && !localPlaying) {
appLogger.d('WatchTogether: Play/pause state diverged, syncing to host (playing)');
await _applyRemotePlay(position: message.position);
await _applyRemotePlay(
position: message.position,
expectedAttachmentGeneration: queuedAttachmentGeneration,
);
} else if (!message.isPlaying! && localPlaying) {
appLogger.d('WatchTogether: Play/pause state diverged, syncing to host (paused)');
await _applyRemotePause();
await _applyRemotePause(expectedAttachmentGeneration: queuedAttachmentGeneration);
}
}
break;
@@ -462,7 +631,7 @@ class WatchTogetherSyncManager {
break;
}
if (message.rate != null) {
await _applyRemoteRate(message.rate!);
await _applyRemoteRate(message.rate!, expectedAttachmentGeneration: queuedAttachmentGeneration);
}
break;
@@ -477,7 +646,7 @@ class WatchTogetherSyncManager {
break;
case SyncMessageType.sessionConfig:
await _handleSessionConfig(message);
await _handleSessionConfig(message, queuedAttachmentGeneration);
break;
case SyncMessageType.ping:
@@ -513,7 +682,10 @@ class WatchTogetherSyncManager {
if (_deferredPlay && isAllReady) {
_deferredPlay = false;
_firstPlayCompleted = true;
await _applyRemotePlay(position: _deferredPlayPosition);
await _applyRemotePlay(
position: _deferredPlayPosition,
expectedAttachmentGeneration: queuedAttachmentGeneration,
);
_deferredPlayPosition = null;
}
}
@@ -529,61 +701,108 @@ class WatchTogetherSyncManager {
}
}
/// Shared helper for applying remote actions with proper guarding
Future<void> _applyRemoteAction(Future<void> Function() action) async {
if (_player == null) return;
_isRemoteAction = true;
try {
await action();
} on StateError catch (e) {
appLogger.w('WatchTogether: Player disposed during remote action', error: e);
detachPlayer();
} finally {
_isRemoteAction = false;
}
}
/// Apply remote play command
Future<void> _applyRemotePlay({Duration? position}) async {
Future<bool> _applyRemotePlay({Duration? position, int? expectedAttachmentGeneration}) async {
appLogger.d('WatchTogether: Applying remote PLAY${position != null ? ' at ${position.inSeconds}s' : ''}');
await _applyRemoteAction(() async {
if (position != null) {
await _player!.seek(position);
}
await _player!.play();
_lastKnownPlaying = true;
});
return _runGuardedRemoteAction(
actionName: 'play',
expectedAttachmentGeneration: expectedAttachmentGeneration,
action: (player, attachmentGeneration) async {
if (position != null) {
final didSeek = await _runGuardedPlayerCommand(
actionName: 'play seek',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.seek(position),
);
if (!didSeek) return false;
}
final didPlay = await _runGuardedPlayerCommand(
actionName: 'play',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.play(),
);
if (!didPlay) return false;
_lastKnownPlaying = true;
return true;
},
);
}
/// Apply remote pause command
Future<void> _applyRemotePause() async {
Future<bool> _applyRemotePause({int? expectedAttachmentGeneration}) async {
appLogger.d('WatchTogether: Applying remote PAUSE');
await _applyRemoteAction(() async {
await _player!.pause();
_lastKnownPlaying = false;
});
return _runGuardedRemoteAction(
actionName: 'pause',
expectedAttachmentGeneration: expectedAttachmentGeneration,
action: (player, attachmentGeneration) async {
final didPause = await _runGuardedPlayerCommand(
actionName: 'pause',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.pause(),
);
if (!didPause) return false;
_lastKnownPlaying = false;
return true;
},
);
}
/// Apply remote seek command
Future<void> _applyRemoteSeek(Duration position) async {
Future<bool> _applyRemoteSeek(Duration position, {int? expectedAttachmentGeneration}) async {
appLogger.d('WatchTogether: Applying remote SEEK to ${position.inSeconds}s');
await _applyRemoteAction(() => _player!.seek(position));
return _runGuardedRemoteAction(
actionName: 'seek',
expectedAttachmentGeneration: expectedAttachmentGeneration,
action: (player, attachmentGeneration) {
return _runGuardedPlayerCommand(
actionName: 'seek',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.seek(position),
);
},
);
}
/// Apply remote rate change
Future<void> _applyRemoteRate(double rate) async {
Future<bool> _applyRemoteRate(double rate, {int? expectedAttachmentGeneration}) async {
appLogger.d('WatchTogether: Applying remote RATE: $rate');
await _applyRemoteAction(() async {
await _player!.setRate(rate);
_lastKnownRate = rate;
});
return _runGuardedRemoteAction(
actionName: 'rate',
expectedAttachmentGeneration: expectedAttachmentGeneration,
action: (player, attachmentGeneration) async {
final didSetRate = await _runGuardedPlayerCommand(
actionName: 'rate',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.setRate(rate),
);
if (!didSetRate) return false;
_lastKnownRate = rate;
return true;
},
);
}
/// Check and correct position drift
void _checkAndCorrectDrift(Duration remotePosition, int remoteTimestamp) {
if (_player == null || _session.isHost) return;
Future<void> _checkAndCorrectDrift(
Duration remotePosition,
int remoteTimestamp,
int expectedAttachmentGeneration,
) async {
if (_session.isHost) return;
final localPosition = _player!.state.position;
final player = _player;
if (player == null || !_isPlayerAttachmentCurrent(player, expectedAttachmentGeneration)) return;
final localPosition = player.state.position;
final now = DateTime.now().millisecondsSinceEpoch;
// Translate host's timestamp to our local time frame using clock offset
@@ -595,11 +814,11 @@ class WatchTogetherSyncManager {
final networkDelay = _hasClockOffset ? rawDelay.clamp(0, 5000) : 0;
// Estimate where remote should be now, accounting for playback time elapsed
Duration estimatedRemoteNow = remotePosition;
if (_player!.state.playing && networkDelay > 0) {
var estimatedRemoteNow = remotePosition;
if (player.state.playing && networkDelay > 0) {
// If playing, account for time elapsed during network transit
// Multiply by rate in case playback speed is different
estimatedRemoteNow = remotePosition + Duration(milliseconds: (networkDelay * _player!.state.rate).round());
estimatedRemoteNow = remotePosition + Duration(milliseconds: (networkDelay * player.state.rate).round());
}
final drift = (localPosition - estimatedRemoteNow).abs();
@@ -608,14 +827,28 @@ class WatchTogetherSyncManager {
// Excessive drift - force sync with indicator
appLogger.w('WatchTogether: Excessive drift (${drift.inSeconds}s), force syncing');
_setSyncing(true);
_applyRemoteSeek(estimatedRemoteNow);
final didSeek = await _applyRemoteSeek(
estimatedRemoteNow,
expectedAttachmentGeneration: expectedAttachmentGeneration,
);
if (!didSeek) {
_setSyncing(false);
return;
}
_syncingTimer?.cancel();
_syncingTimer = Timer(const Duration(milliseconds: 500), () => _setSyncing(false));
} else if (drift > maxAllowedDrift) {
// Normal drift correction
appLogger.d('WatchTogether: Drift correction (${drift.inMilliseconds}ms)');
_setSyncing(true);
_applyRemoteSeek(estimatedRemoteNow);
final didSeek = await _applyRemoteSeek(
estimatedRemoteNow,
expectedAttachmentGeneration: expectedAttachmentGeneration,
);
if (!didSeek) {
_setSyncing(false);
return;
}
_syncingTimer?.cancel();
_syncingTimer = Timer(const Duration(milliseconds: 300), () => _setSyncing(false));
}
@@ -645,7 +878,7 @@ class WatchTogetherSyncManager {
}
/// Handle session config from host
Future<void> _handleSessionConfig(SyncMessage message) async {
Future<void> _handleSessionConfig(SyncMessage message, int expectedAttachmentGeneration) async {
if (_session.isHost) return; // Host doesn't need to process config
appLogger.d('WatchTogether: Received session config');
@@ -661,43 +894,68 @@ class WatchTogetherSyncManager {
onSessionConfigReceived?.call(message.controlMode!);
}
// Sync to host's current state
if (_player == null) return;
_isRemoteAction = true;
try {
// Always seek to host's position first
if (message.position != null) {
await _player!.seek(message.position!);
}
// Match playback rate
if (message.rate != null) {
await _player!.setRate(message.rate!);
_lastKnownRate = message.rate!;
}
// Match play/pause state (prefer isPlaying, fall back to legacy bufferingState encoding)
final hostIsPlaying = message.isPlaying ?? (message.bufferingState == false);
if (hostIsPlaying) {
// 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);
await _runGuardedRemoteAction(
actionName: 'session config',
expectedAttachmentGeneration: expectedAttachmentGeneration,
action: (player, attachmentGeneration) async {
// Always seek to host's position first
if (message.position != null) {
final didSeek = await _runGuardedPlayerCommand(
actionName: 'session config seek',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.seek(message.position!),
);
if (!didSeek) return false;
}
} else {
await _player!.pause();
_lastKnownPlaying = false;
}
} on StateError catch (e) {
appLogger.w('WatchTogether: Player disposed during session config apply', error: e);
detachPlayer();
} finally {
_isRemoteAction = false;
}
// Match playback rate
if (message.rate != null) {
final didSetRate = await _runGuardedPlayerCommand(
actionName: 'session config rate',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.setRate(message.rate!),
);
if (!didSetRate) return false;
_lastKnownRate = message.rate!;
}
// Match play/pause state (prefer isPlaying, fall back to legacy bufferingState encoding)
final hostIsPlaying = message.isPlaying ?? (message.bufferingState == false);
if (hostIsPlaying) {
// Host was playing — defer until our video is loaded
_deferredPlay = true;
_deferredPlayPosition = message.position;
if (_hasAnnouncedReady) {
_deferredPlay = false;
_firstPlayCompleted = true;
final didPlay = await _runGuardedPlayerCommand(
actionName: 'session config play',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.play(),
);
if (!didPlay) return false;
_lastKnownPlaying = true;
}
} else {
final didPause = await _runGuardedPlayerCommand(
actionName: 'session config pause',
player: player,
attachmentGeneration: attachmentGeneration,
command: (player) => player.pause(),
);
if (!didPause) return false;
_lastKnownPlaying = false;
}
return true;
},
);
}
/// Set syncing state and notify listeners