From baa31742cb80aeef0814d4baf72b6c0561f48cfd Mon Sep 17 00:00:00 2001 From: edde746 <86283021+edde746@users.noreply.github.com> Date: Sun, 6 Sep 2026 13:29:37 +0200 Subject: [PATCH] fix(player): bind live TV clock generations to the source the load reports Live-TV clock generations were matched to mpv sources first-in-first-out: every start-file popped the oldest registered open, on the assumption of exactly one start-file per open in dispatch order. Opens without a generation on the same player, an Android loadfile rejected silently by nativeCommand, and the independent delivery of the command ack and the start-file event all broke that, so a seek that reopened the stream could calibrate against the wrong source. The loadfile reply now carries mpv's playlist_entry_id on Android, Apple, Linux and Windows, PlayerNative.open resolves with it, and the live session binds each generation to that id explicitly. Source events that land before the reply are buffered per id and replayed on binding; a rejected or unreachable load fails its generation instead of leaving a phantom; opens with no generation are invisible to clock binding. Android now reports a rejected mpv command as COMMAND_FAILED like the other cores. --- .../com/edde746/plezy/mpv/MpvPlayerCore.kt | 26 ++-- .../com/edde746/plezy/mpv/MpvPlayerPlugin.kt | 18 ++- .../edde746/plezy/mpv/MpvPlayerPluginTest.kt | 16 +++ android/libmpv/src/main/cpp/main.cpp | 33 ++++- .../com/edde746/plezy/libmpv/MpvPlayer.kt | 15 ++- lib/mpv/player/player.dart | 4 + lib/mpv/player/player_native.dart | 21 ++- .../video_player/live_tv_session_state.dart | 102 ++++++++++----- lib/screens/video_player/parts/live_tv.dart | 41 ++++-- .../video_player/parts/playback_services.dart | 6 - linux/runner/mpv/mpv_player.cc | 7 +- linux/runner/mpv/mpv_player.h | 2 +- linux/runner/mpv/mpv_player_lifecycle_test.cc | 76 ++++++++++- linux/runner/mpv/mpv_plugin.cc | 13 +- .../apple/MpvPlayer/MpvPlayerCoreBase.swift | 56 +++++++- .../MpvPlayer/MpvPlayerPluginShared.swift | 7 +- shared/mpv/mpv_player_common.h | 65 +++++++++- shared/mpv/mpv_player_common_test.cpp | 90 +++++++++++++ test/mpv/player_open_test.dart | 41 ++++++ .../live_tv_session_state_test.dart | 121 +++++++++++++++++- windows/runner/mpv/mpv_player.cpp | 7 +- windows/runner/mpv/mpv_player.h | 2 +- .../mpv/mpv_player_property_contract_test.cpp | 81 ++++++++++++ windows/runner/mpv/mpv_plugin.cpp | 15 ++- 24 files changed, 760 insertions(+), 105 deletions(-) diff --git a/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerCore.kt b/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerCore.kt index 6ce02b496..4fdbd8d88 100644 --- a/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerCore.kt +++ b/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerCore.kt @@ -1882,21 +1882,29 @@ class MpvPlayerCore private constructor( } fun command(args: Array, onComplete: ((Boolean) -> Unit)? = null) { + commandForSource(args) { onComplete?.invoke(it.isSuccess) } + } + + /** + * Runs an mpv command on the ordered writer. Completes on the main thread with the playlist + * entry id a `loadfile` created (null for every other command), or with the failure mpv + * reported — a rejected load never starts a source, so it must not be reported as one. + */ + fun commandForSource(args: Array, onComplete: (Result) -> Unit) { if (!isInitialized || disposing || args.isEmpty() || !scope.isActive) { - onComplete?.invoke(false) + onComplete(Result.failure(IllegalStateException("MPV player unavailable"))) return } scope.launch(mpvWriteDispatcher) { - var success = false - try { - player?.command(*args) - success = true + val outcome: Result = try { + val p = player ?: throw IllegalStateException("MPV player unavailable") + Result.success(p.command(*args)) } catch (e: Exception) { Log.w(TAG, "command failed", e) - } finally { - withContext(NonCancellable + Dispatchers.Main) { - onComplete?.invoke(success) - } + Result.failure(e) + } + withContext(NonCancellable + Dispatchers.Main) { + onComplete(outcome) } } } diff --git a/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerPlugin.kt b/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerPlugin.kt index 686195dd5..88925f57c 100644 --- a/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerPlugin.kt +++ b/android/app/src/main/kotlin/com/edde746/plezy/mpv/MpvPlayerPlugin.kt @@ -474,12 +474,18 @@ open class MpvPlayerPlugin( result.error("NOT_INITIALIZED", "Player not initialized", null) return } - core.command(args.toTypedArray()) { success -> - if (success) { - result.success(null) - } else { - result.error("COMMAND_FAILED", "mpv command failed", args) - } + // `loadfile` answers with the playlist entry it created so Dart can tie + // the load to that source's start-file/playback-restart/end-file events; + // every other command answers null. + core.commandForSource(args.toTypedArray()) { outcome -> + outcome.fold( + onSuccess = { playlistEntryId -> + result.success(playlistEntryId?.let { mapOf("playlistEntryId" to it) }) + }, + onFailure = { error -> + result.error("COMMAND_FAILED", error.message ?: "mpv command failed", args) + } + ) } } diff --git a/android/app/src/test/kotlin/com/edde746/plezy/mpv/MpvPlayerPluginTest.kt b/android/app/src/test/kotlin/com/edde746/plezy/mpv/MpvPlayerPluginTest.kt index caa2f79a5..193791739 100644 --- a/android/app/src/test/kotlin/com/edde746/plezy/mpv/MpvPlayerPluginTest.kt +++ b/android/app/src/test/kotlin/com/edde746/plezy/mpv/MpvPlayerPluginTest.kt @@ -49,6 +49,22 @@ class MpvPlayerPluginTest { assertNull(result.successValue) } + @Test + fun commandWithoutNativePlayerReportsFailureInsteadOfSilentSuccess() { + // A load that never reached mpv produces no source; answering success + // would leave Dart waiting on a start-file that never comes. + val plugin = MpvPlayerPlugin() + installCore(plugin, testCore(null)) + val result = RecordingResult() + + plugin.onMethodCall(MethodCall("command", mapOf("args" to listOf("loadfile", "x", "replace"))), result) + awaitCompletion(result) + + assertEquals("COMMAND_FAILED", result.errorCode) + assertEquals(1, result.completionCount) + assertNull(result.successValue) + } + @Test fun audioSpdifCodecsWithoutContextAnswersEmptySoMpvDecodes() { // mpv force-passthroughs every codec named in audio-spdif with no decode fallback, so diff --git a/android/libmpv/src/main/cpp/main.cpp b/android/libmpv/src/main/cpp/main.cpp index 209a5e7f9..cccf731f8 100644 --- a/android/libmpv/src/main/cpp/main.cpp +++ b/android/libmpv/src/main/cpp/main.cpp @@ -6,6 +6,7 @@ #include #include #include +#include #include #include #include @@ -30,7 +31,7 @@ jni_func(void, nativeDestroy); jni_func(jint, nativeSetLogLevel, jstring level); -jni_func(void, nativeCommand, jobjectArray jarray); +jni_func(jlong, nativeCommand, jobjectArray jarray); jni_func(void, nativeHookContinue, jlong id); }; @@ -143,14 +144,18 @@ jni_func(jint, nativeSetLogLevel, jstring jlevel) { return result; } -jni_func(void, nativeCommand, jobjectArray jarray) { - CHECK_MPV_INIT(); +// Runs a command synchronously. Returns the negative mpv error on failure, +// the `playlist_entry_id` a `loadfile` created (always > 0), or 0 for a +// command that succeeded without one. `die` (a Java RuntimeException) is +// reserved for a missing core, matching the other entry points. +jni_func(jlong, nativeCommand, jobjectArray jarray) { + CHECK_MPV_INIT_RET(MPV_ERROR_UNINITIALIZED); const char* arguments[128] = {0}; int len = env->GetArrayLength(jarray); if (len >= (int)ARRAYLEN(arguments)) { die("too many command arguments"); - return; + return MPV_ERROR_INVALID_PARAMETER; } std::vector storage; @@ -162,7 +167,25 @@ jni_func(void, nativeCommand, jobjectArray jarray) { env->DeleteLocalRef(jarg); } - mpv_command(g_mpv, arguments); + mpv_node result{}; + const int status = mpv_command_ret(g_mpv, arguments, &result); + if (status < 0) { + ALOGE("mpv_command(%s) returned error %s", len > 0 ? arguments[0] : "", mpv_error_string(status)); + return status; + } + + jlong playlist_entry_id = 0; + const mpv_node_list* map = result.format == MPV_FORMAT_NODE_MAP ? result.u.list : nullptr; + if (map && map->keys && map->values) { + for (int i = 0; i < map->num; ++i) { + if (map->keys[i] && strcmp(map->keys[i], "playlist_entry_id") == 0 && map->values[i].format == MPV_FORMAT_INT64) { + playlist_entry_id = (jlong)map->values[i].u.int64; + break; + } + } + } + mpv_free_node_contents(&result); + return playlist_entry_id; } jni_func(void, nativeHookContinue, jlong id) { diff --git a/android/libmpv/src/main/java/com/edde746/plezy/libmpv/MpvPlayer.kt b/android/libmpv/src/main/java/com/edde746/plezy/libmpv/MpvPlayer.kt index 960dace54..d89e9f715 100644 --- a/android/libmpv/src/main/java/com/edde746/plezy/libmpv/MpvPlayer.kt +++ b/android/libmpv/src/main/java/com/edde746/plezy/libmpv/MpvPlayer.kt @@ -156,7 +156,8 @@ class MpvPlayer private constructor() : AutoCloseable { @JvmStatic private external fun nativeDestroy() - @JvmStatic private external fun nativeCommand(cmd: Array) + /** Negative mpv error, the playlist entry id a `loadfile` created, or 0 when the command returned none. */ + @JvmStatic private external fun nativeCommand(cmd: Array): Long @JvmStatic private external fun nativeSetLogLevel(level: String): Int @@ -256,9 +257,17 @@ class MpvPlayer private constructor() : AutoCloseable { // Commands - suspend fun command(vararg args: String) { + /** + * Runs an mpv command. `loadfile` returns the id of the playlist entry it created — the + * `sourceId` carried by that source's start-file / playback-restart / end-file events; every + * other command returns null. A command mpv rejects throws [MpvException]: a rejected load + * never produces a source, so the caller must not wait for one. + */ + suspend fun command(vararg args: String): Long? { checkNotClosed() - withContext(Dispatchers.IO) { nativeCommand(args) } + val status = withContext(Dispatchers.IO) { nativeCommand(args) } + if (status < 0) throw MpvException("Command '${args.firstOrNull() ?: ""}' failed: error $status") + return if (status > 0) status else null } /** Called on the core's ordered IO writer, without suspending between writes. */ diff --git a/lib/mpv/player/player.dart b/lib/mpv/player/player.dart index 38cf45436..b6ecc762f 100644 --- a/lib/mpv/player/player.dart +++ b/lib/mpv/player/player.dart @@ -74,6 +74,10 @@ abstract class Player { /// /// [media] - The media source to open. /// [play] - Whether to start playback immediately (default: true). + /// + /// Backends that can identify the source they started resolve with its id + /// (mpv: the playlist entry id carried by that source's stream events, see + /// `PlayerNative.open`); the base contract promises nothing. Future open( Media media, { bool play = true, diff --git a/lib/mpv/player/player_native.dart b/lib/mpv/player/player_native.dart index a01bf21ea..b29f4abb6 100644 --- a/lib/mpv/player/player_native.dart +++ b/lib/mpv/player/player_native.dart @@ -370,17 +370,25 @@ class PlayerNative extends PlayerBase { return ('fdclose://$fd', fd); } + /// Opens [media] and resolves with the id of the mpv playlist entry the + /// `loadfile` created — the `sourceId` that entry's `start-file`, + /// `playback-restart` and `end-file` events carry — or null when the load + /// never reached mpv (core unavailable) or the core did not report one. + /// + /// Method-channel contract: the `command` reply for `loadfile` is + /// `{'playlistEntryId': int}` on every native core; every other command + /// replies null. @override - Future open( + Future open( Media media, { bool play = true, bool isLive = false, List? externalSubtitles, Duration? timelineDuration, }) async { - if (_nativeCoreUnavailable) return; + if (_nativeCoreUnavailable) return null; await _ensureInitialized(); - if (_nativeCoreUnavailable) return; + if (_nativeCoreUnavailable) return null; // `loadfile replace` (below) clears the native playlist, dropping any // gapless entry armed via setNext — settle its content-fd claim first. // No transition is surfaced: the caller is replacing playback anyway. @@ -451,7 +459,11 @@ class PlayerNative extends PlayerBase { loadfileArgs.addAll(['-1', loadfileOption]); } if (audioOnly) _expectOpenFileLoad = true; - await command(loadfileArgs); + // The core can be torn down while the awaits above were suspended; the + // `command` path makes the same re-check before dispatching. + if (_nativeCoreUnavailable) return null; + final loadfileReply = await invoke('command', {'args': loadfileArgs}); + final playlistEntryId = loadfileReply?['playlistEntryId']; // mpv's pause property survives loadfile; in-place reloads pause the old // file before resolving, so explicitly unpause for the replacement. Set @@ -460,6 +472,7 @@ class PlayerNative extends PlayerBase { if (play) { await setProperty('pause', 'no'); } + return playlistEntryId is int ? playlistEntryId : null; } @override diff --git a/lib/screens/video_player/live_tv_session_state.dart b/lib/screens/video_player/live_tv_session_state.dart index 1f71a72f1..7f29bd800 100644 --- a/lib/screens/video_player/live_tv_session_state.dart +++ b/lib/screens/video_player/live_tv_session_state.dart @@ -16,6 +16,15 @@ class _LiveClockOpen { bool canceled = false; } +/// What the player has reported so far about one MPV source that no clock +/// open has claimed yet. The `loadfile` reply that names the source and the +/// source's own events travel on independent channels, so readiness or +/// failure can land before the open learns its id. +class _UnclaimedSourceEvents { + Duration? readyPosition; + bool failed = false; +} + /// Mutable runtime state for one live TV playback: the current /// [LiveTvPlaybackSession] protocol handle, the timeline heartbeat /// machinery, the capture buffer used for time-shifting, and the @@ -72,10 +81,15 @@ class LiveTvSessionState { int? _latestClockGeneration; int? activeClockSourceId; double? pendingStreamEpoch; - final List<_LiveClockOpen> _unboundClockOpens = []; final Map _clockOpensBySource = {}; final Map _clockOpensByGeneration = {}; + /// Recent source events no open has claimed, keyed by source id in arrival + /// order and bounded so opens that never register (live-edge re-opens, + /// VOD on the same player) cannot grow it. + final Map _unclaimedSources = {}; + static const int _maxUnclaimedSources = 8; + /// Fallback level for live TV stream errors (mirrors Plex web client /// behavior). 0 = directStream+directStreamAudio, 1 = no directStream, /// 2 = no DS + no DS audio. @@ -91,45 +105,62 @@ class LiveTvSessionState { /// The player route is closed instead of attempting to reuse that session. bool exitOnResume = false; - /// Register an offset-based MPV open before dispatching `loadfile`. - /// - /// Opens and MPV START_FILE events are ordered. Canceled entries stay in the - /// unbound queue until their START_FILE arrives so a late predecessor cannot - /// steal the next open's source identity. + /// Register an offset-based MPV open before dispatching `loadfile`. The + /// returned generation is the handle the caller binds to the source id the + /// load reports ([bindClockOpen]); until then the open is unbound and no + /// source event can reach it. Every earlier open is superseded. int beginClockOpen(int targetEpoch) { - final previousOpens = <_LiveClockOpen>{ - ..._clockOpensByGeneration.values, - ..._unboundClockOpens, - ..._clockOpensBySource.values, - }; + final previousOpens = <_LiveClockOpen>{..._clockOpensByGeneration.values, ..._clockOpensBySource.values}; for (final open in previousOpens) { open.canceled = true; if (!open.result.isCompleted) open.result.complete(false); } - _clockOpensBySource.removeWhere((_, open) => open.canceled); + _clockOpensByGeneration.clear(); + _clockOpensBySource.clear(); final open = _LiveClockOpen(generation: ++_nextClockGeneration, targetEpoch: targetEpoch); _clockOpensByGeneration[open.generation] = open; - _unboundClockOpens.add(open); _latestClockGeneration = open.generation; pendingStreamEpoch = targetEpoch.toDouble(); return open.generation; } - /// Bind the oldest dispatched live open to the source MPV started next. - bool bindClockSource(PlayerSourceStarted source) { - if (_unboundClockOpens.isEmpty) return false; - final open = _unboundClockOpens.removeAt(0); - open.sourceId = source.sourceId; - if (open.canceled) return false; - _clockOpensBySource[source.sourceId] = open; + /// Bind [generation] to the MPV source its `loadfile` reply named. Source + /// events that already arrived for [sourceId] apply immediately, so a + /// readiness or failure report that beat the reply is not lost. Returns + /// whether the open is still live and now keyed by the source. + bool bindClockOpen(int generation, int sourceId) { + final open = _clockOpensByGeneration[generation]; + if (open == null) return false; + open.sourceId = sourceId; + if (open.canceled) { + _unclaimedSources.remove(sourceId); + return false; + } + _clockOpensBySource[sourceId] = open; + final observed = _unclaimedSources.remove(sourceId); + if (observed == null) return true; + if (observed.failed) { + _failClockOpen(open); + return false; + } + final readyPosition = observed.readyPosition; + if (readyPosition != null) { + calibrateClockSource(PlayerSourceReady(sourceId: sourceId, position: readyPosition)); + } return true; } /// Calibrate epoch time against the first decoded position of [source]. + /// Readiness of a source no open has claimed yet is kept for a later + /// [bindClockOpen]. bool calibrateClockSource(PlayerSourceReady source) { final open = _clockOpensBySource.remove(source.sourceId); - if (open == null || open.canceled || open.generation != _latestClockGeneration) return false; + if (open == null) { + _unclaimedSource(source.sourceId).readyPosition = source.position; + return false; + } + if (open.canceled || open.generation != _latestClockGeneration) return false; streamStartEpoch = open.targetEpoch - source.position.inMilliseconds / 1000.0; activeClockSourceId = source.sourceId; @@ -141,7 +172,22 @@ class LiveTvSessionState { void failClockSource(PlayerSourceFailed source) { final open = _clockOpensBySource.remove(source.sourceId); - if (open != null) _failClockOpen(open); + if (open == null) { + _unclaimedSource(source.sourceId).failed = true; + return; + } + _failClockOpen(open); + } + + _UnclaimedSourceEvents _unclaimedSource(int sourceId) { + final existing = _unclaimedSources.remove(sourceId); + if (existing != null) { + return _unclaimedSources[sourceId] = existing; + } + if (_unclaimedSources.length >= _maxUnclaimedSources) { + _unclaimedSources.remove(_unclaimedSources.keys.first); + } + return _unclaimedSources[sourceId] = _UnclaimedSourceEvents(); } void failClockOpen(int generation) { @@ -160,7 +206,6 @@ class LiveTvSessionState { void _failClockOpen(_LiveClockOpen open) { open.canceled = true; - _unboundClockOpens.remove(open); final sourceId = open.sourceId; if (sourceId != null && identical(_clockOpensBySource[sourceId], open)) { _clockOpensBySource.remove(sourceId); @@ -175,14 +220,13 @@ class LiveTvSessionState { if (!open.result.isCompleted) open.result.complete(false); } + /// Resolves true once the open's source is calibrated, false when it is + /// superseded, failed or timed out. A timed-out open stays registered so a + /// late `loadfile` reply or readiness event can still bind and calibrate it. Future clockOpenResult(int generation) { final open = _clockOpensByGeneration[generation]; if (open == null) return Future.value(false); - return open.result.future.whenComplete(() { - if (identical(_clockOpensByGeneration[generation], open)) { - _clockOpensByGeneration.remove(generation); - } - }); + return open.result.future; } void cancelClockOpens() { @@ -190,8 +234,8 @@ class LiveTvSessionState { for (final generation in generations) { failClockOpen(generation); } - _unboundClockOpens.clear(); _clockOpensBySource.clear(); + _unclaimedSources.clear(); _latestClockGeneration = null; pendingStreamEpoch = null; } diff --git a/lib/screens/video_player/parts/live_tv.dart b/lib/screens/video_player/parts/live_tv.dart index a291d7b84..31d557891 100644 --- a/lib/screens/video_player/parts/live_tv.dart +++ b/lib/screens/video_player/parts/live_tv.dart @@ -235,8 +235,10 @@ extension _VideoPlayerLiveTvMethods on VideoPlayerScreenState { /// Re-opens the current session's live stream at [streamUrl]. /// /// Offset-based MPV opens register their requested absolute [targetEpoch] - /// before `loadfile`. When [awaitClock] is true, success means the new - /// source's first rendered player position has been mapped to that epoch. + /// before `loadfile` and bind that registration to the source id the load + /// reports, so only that source's events can calibrate it. When [awaitClock] + /// is true, success means the new source's first rendered player position + /// has been mapped to that epoch. Future _openLiveStream( Player player, String streamUrl, { @@ -246,21 +248,32 @@ extension _VideoPlayerLiveTvMethods on VideoPlayerScreenState { bool applyOptions = true, }) async { _live.streamGeneration++; - final clockGeneration = targetEpoch != null && player is PlayerNative ? _live.beginClockOpen(targetEpoch) : null; - final clockResult = clockGeneration == null ? null : _live.clockOpenResult(clockGeneration); - try { + final media = Media(streamUrl, headers: const {'Accept-Language': 'en'}); + final playNow = play ?? automotivePlaybackAllowedNow(); + if (targetEpoch == null || player is! PlayerNative) { if (applyOptions) await _setLiveStreamOptions(player); - await player.open( - Media(streamUrl, headers: const {'Accept-Language': 'en'}), - play: play ?? automotivePlaybackAllowedNow(), - isLive: true, - ); - } catch (_) { - if (clockGeneration != null) _live.failClockOpen(clockGeneration); - rethrow; + await player.open(media, play: playNow, isLive: true); + return true; } - if (clockResult == null || clockGeneration == null) return true; + final clockGeneration = _live.beginClockOpen(targetEpoch); + final clockResult = _live.clockOpenResult(clockGeneration); + final int? sourceId; + try { + if (applyOptions) await _setLiveStreamOptions(player); + sourceId = await player.open(media, play: playNow, isLive: true); + } catch (_) { + _live.failClockOpen(clockGeneration); + rethrow; + } + if (sourceId == null) { + // The load never reached mpv (core unavailable), so no source will ever + // report for this open. + _live.failClockOpen(clockGeneration); + return false; + } + _live.bindClockOpen(clockGeneration, sourceId); + if (!awaitClock) { unawaited(clockResult); return true; diff --git a/lib/screens/video_player/parts/playback_services.dart b/lib/screens/video_player/parts/playback_services.dart index 029c3a1d7..1eed14903 100644 --- a/lib/screens/video_player/parts/playback_services.dart +++ b/lib/screens/video_player/parts/playback_services.dart @@ -146,12 +146,6 @@ extension _VideoPlayerPlaybackServiceMethods on VideoPlayerScreenState { } if (widget.isLive) { - _playerStreamSubscriptions.add( - currentPlayer.streams.sourceStarted.listen((source) { - if (!mounted || player != currentPlayer) return; - _live.bindClockSource(source); - }), - ); _playerStreamSubscriptions.add( currentPlayer.streams.sourceReady.listen((source) { if (!mounted || player != currentPlayer) return; diff --git a/linux/runner/mpv/mpv_player.cc b/linux/runner/mpv/mpv_player.cc index 245e0db43..00c419e16 100644 --- a/linux/runner/mpv/mpv_player.cc +++ b/linux/runner/mpv/mpv_player.cc @@ -702,6 +702,9 @@ void MpvPlayer::Dispose() { for (auto& callback : cancelled.status) { callback(MPV_ERROR_UNINITIALIZED); } + for (auto& callback : cancelled.commands) { + callback(MPV_ERROR_UNINITIALIZED, nullptr); + } for (auto& callback : cancelled.properties) { callback(-1, ""); } @@ -754,7 +757,7 @@ void MpvPlayer::Command(const std::vector& args) { CommandAsync(arg void MpvPlayer::CommandAsync(const std::vector& args, CommandCallback callback) { if (disposed_ || !mpv_) { - if (callback) callback(MPV_ERROR_UNINITIALIZED); + if (callback) callback(MPV_ERROR_UNINITIALIZED, nullptr); return; } @@ -1040,7 +1043,7 @@ void MpvPlayer::TryAudioReload(const char* reason, int attempt, uint64_t request LogRecovery("issuing ao-reload (reason=" + std::string(reason) + ", attempt " + std::to_string(attempt) + ")"); const std::string reason_copy = reason; auto callback_context = callback_context_; - CommandAsync({"ao-reload"}, [callback_context, reason_copy, attempt, request_generation](int error) { + CommandAsync({"ao-reload"}, [callback_context, reason_copy, attempt, request_generation](int error, const mpv_node*) { auto lease = callback_context->Acquire(); if (!lease) return; MpvPlayer* player = lease.player(); diff --git a/linux/runner/mpv/mpv_player.h b/linux/runner/mpv/mpv_player.h index 04ac5c2d5..4d3a52d67 100644 --- a/linux/runner/mpv/mpv_player.h +++ b/linux/runner/mpv/mpv_player.h @@ -127,7 +127,7 @@ class MpvPlayer { /// Callback types for async mpv requests. using StatusCallback = plezy::mpv_common::StatusCallback; - using CommandCallback = StatusCallback; + using CommandCallback = plezy::mpv_common::CommandCallback; using GetPropertyCallback = plezy::mpv_common::GetPropertyCallback; /// Executes an mpv command asynchronously to prevent UI blocking. diff --git a/linux/runner/mpv/mpv_player_lifecycle_test.cc b/linux/runner/mpv/mpv_player_lifecycle_test.cc index 75a8c71ea..92c63bab8 100644 --- a/linux/runner/mpv/mpv_player_lifecycle_test.cc +++ b/linux/runner/mpv/mpv_player_lifecycle_test.cc @@ -46,6 +46,10 @@ class MpvPlayerLifecycleTestPeer { player.pending_requests_.RegisterStatus(std::move(callback)); } + static uint64_t RegisterPendingCommand(MpvPlayer& player, MpvPlayer::CommandCallback callback) { + return player.pending_requests_.RegisterCommand(std::move(callback)); + } + static int PendingSourceCount(MpvPlayer& player) { std::lock_guard lock(player.source_mutex_); return (player.wakeup_source_id_ != 0 ? 1 : 0) + (player.redraw_source_id_ != 0 ? 1 : 0) + @@ -432,14 +436,83 @@ void TestUnavailableCommandFails() { MpvPlayer player; int callback_count = 0; int status = MPV_ERROR_SUCCESS; + mpv_node sentinel{}; + const mpv_node* reported_result = &sentinel; - player.CommandAsync({"stop"}, [&](int error) { + player.CommandAsync({"stop"}, [&](int error, const mpv_node* result) { ++callback_count; status = error; + reported_result = result; }); Check(callback_count == 1, "a command without an mpv handle must complete exactly once"); Check(status == MPV_ERROR_UNINITIALIZED, "a command without an mpv handle must fail as uninitialized"); + Check(reported_result == nullptr, "a failed command must not carry a result node"); +} + +// The `loadfile` reply names the playlist entry mpv created for the load — +// the source id its start-file/playback-restart/end-file events carry — so the +// Dart side can bind a load to its source instead of guessing by arrival order. +void TestCommandReplyCarriesPlaylistEntryId() { + MpvPlayer player; + constexpr int64_t kEntryId = 8000000004LL; + + int callback_count = 0; + int status = MPV_ERROR_SUCCESS; + int64_t reported_entry_id = 0; + bool reported_entry = false; + const uint64_t request_id = + MpvPlayerLifecycleTestPeer::RegisterPendingCommand(player, [&](int error, const mpv_node* result) { + ++callback_count; + status = error; + reported_entry = plezy::mpv_common::PlaylistEntryIdFromCommandResult(result, &reported_entry_id); + }); + + const char* keys[] = {"playlist_entry_id"}; + mpv_node values[1]{}; + values[0].format = MPV_FORMAT_INT64; + values[0].u.int64 = kEntryId; + mpv_node_list map{}; + map.num = 1; + map.keys = const_cast(keys); + map.values = values; + mpv_event_command command{}; + command.result.format = MPV_FORMAT_NODE_MAP; + command.result.u.list = ↦ + mpv_event reply{}; + reply.event_id = MPV_EVENT_COMMAND_REPLY; + reply.reply_userdata = request_id; + reply.data = &command; + MpvPlayerLifecycleTestPeer::HandleEvent(player, &reply); + + Check(callback_count == 1, "a command reply must complete its request exactly once"); + Check(status == MPV_ERROR_SUCCESS, "a successful command reply changed its status"); + Check(reported_entry, "the loadfile reply must expose the playlist entry id"); + Check(reported_entry_id == kEntryId, "the playlist entry id lost signed 64-bit precision"); + + // A reply without a result map (every non-loadfile command) answers no id. + int64_t unexpected_entry_id = 0; + bool unexpected_entry = true; + const uint64_t plain_request_id = + MpvPlayerLifecycleTestPeer::RegisterPendingCommand(player, [&](int, const mpv_node* result) { + unexpected_entry = plezy::mpv_common::PlaylistEntryIdFromCommandResult(result, &unexpected_entry_id); + }); + mpv_event_command plain_command{}; + plain_command.result.format = MPV_FORMAT_NONE; + reply.reply_userdata = plain_request_id; + reply.data = &plain_command; + MpvPlayerLifecycleTestPeer::HandleEvent(player, &reply); + Check(!unexpected_entry, "a command without a result map must not report a playlist entry id"); + + // A failed reply must not expose whatever the event's result slot holds. + const mpv_node* failed_result = &command.result; + const uint64_t failed_request_id = MpvPlayerLifecycleTestPeer::RegisterPendingCommand( + player, [&](int, const mpv_node* result) { failed_result = result; }); + reply.reply_userdata = failed_request_id; + reply.error = MPV_ERROR_COMMAND; + reply.data = &command; + MpvPlayerLifecycleTestPeer::HandleEvent(player, &reply); + Check(failed_result == nullptr, "a failed command reply must not carry a result node"); } void TestUnavailablePropertyWriteFails() { @@ -721,6 +794,7 @@ int main() { mpv::TestUnavailablePropertyWriteFails(); mpv::TestNodeConversionRejectsMalformedPayloads(); mpv::TestUnavailableCommandFails(); + mpv::TestCommandReplyCarriesPlaylistEntryId(); mpv::TestPendingPropertyWriteFailsOnDispose(); mpv::TestQueuedSourcesAreRetired(context); mpv::TestNativeLeaseBlocksDispose(); diff --git a/linux/runner/mpv/mpv_plugin.cc b/linux/runner/mpv/mpv_plugin.cc index d252c7b12..cbb1b906a 100644 --- a/linux/runner/mpv/mpv_plugin.cc +++ b/linux/runner/mpv/mpv_plugin.cc @@ -1138,13 +1138,22 @@ static void mpv_plugin_handle_method_call(FlMethodChannel* channel, FlMethodCall } } g_object_ref(method_call); - self->player->CommandAsync(command_args, [method_call](int error) { + self->player->CommandAsync(command_args, [method_call](int error, const mpv_node* command_result) { g_autoptr(FlMethodResponse) async_response = nullptr; if (error < 0) { async_response = FL_METHOD_RESPONSE(fl_method_error_response_new("COMMAND_FAILED", "MPV command failed", nullptr)); } else { - async_response = FL_METHOD_RESPONSE(fl_method_success_response_new(nullptr)); + // `loadfile` answers with the playlist entry it created so Dart can + // tie the load to that source's start-file/playback-restart/end-file + // events; every other command answers null. + int64_t playlist_entry_id = 0; + g_autoptr(FlValue) reply = nullptr; + if (plezy::mpv_common::PlaylistEntryIdFromCommandResult(command_result, &playlist_entry_id)) { + reply = fl_value_new_map(); + fl_value_set_string_take(reply, "playlistEntryId", fl_value_new_int(playlist_entry_id)); + } + async_response = FL_METHOD_RESPONSE(fl_method_success_response_new(reply)); } fl_method_call_respond(method_call, async_response, nullptr); g_object_unref(method_call); diff --git a/shared/apple/MpvPlayer/MpvPlayerCoreBase.swift b/shared/apple/MpvPlayer/MpvPlayerCoreBase.swift index 65138e937..5da247fe0 100644 --- a/shared/apple/MpvPlayer/MpvPlayerCoreBase.swift +++ b/shared/apple/MpvPlayer/MpvPlayerCoreBase.swift @@ -203,6 +203,9 @@ class MpvPlayerCoreBase: NSObject { private enum PendingRequest { case void((Result) -> Void) + /// A command reply. The payload is the playlist entry `loadfile` created + /// (the `sourceId` that entry's events carry); nil for every other command. + case command((Result) -> Void) case getProperty((Result) -> Void) } @@ -757,16 +760,20 @@ class MpvPlayerCoreBase: NSObject { commandAsync(args) { _ in } } - func commandAsync(_ args: [String], completion: @escaping (Result) -> Void) { + /// Runs an mpv command. `loadfile` completes with the id of the playlist + /// entry it created — the `sourceId` of that source's start-file, + /// playback-restart and end-file events — so a caller can bind the load to + /// its source; every other command completes with nil. + func commandAsync(_ args: [String], completion: @escaping (Result) -> Void) { guard !args.isEmpty else { - completeOnMain { completion(.success(())) } + completeOnMain { completion(.success(nil)) } return } var cargs: [UnsafeMutablePointer?] = args.map { strdup($0) } cargs.append(nil) - submitAsyncRequest(.void(completion)) { mpv, requestId in + submitAsyncRequest(.command(completion)) { mpv, requestId in cargs.withUnsafeBufferPointer { buffer in var constPointers = buffer.map { UnsafePointer($0) } return mpv_command_async(mpv, requestId, &constPointers) @@ -967,6 +974,8 @@ class MpvPlayerCoreBase: NSObject { switch request { case .void(let completion): completion(.failure(error)) + case .command(let completion): + completion(.failure(error)) case .getProperty(let completion): completion(.failure(error)) } @@ -1020,6 +1029,8 @@ class MpvPlayerCoreBase: NSObject { switch request { case .void(let completion): completion(.failure(error)) + case .command(let completion): + completion(.failure(error)) case .getProperty(let completion): completion(.failure(error)) } @@ -1036,6 +1047,8 @@ class MpvPlayerCoreBase: NSObject { switch request { case .void(let completion): completion(.failure(error)) + case .command(let completion): + completion(.failure(error)) case .getProperty(let completion): completion(.failure(error)) } @@ -1051,6 +1064,41 @@ class MpvPlayerCoreBase: NSObject { } } + private func completeCommandRequest(_ event: mpv_event) { + guard case .command(let completion) = takeRequest(event.reply_userdata) else { return } + + if event.error < 0 { + let error = mpvError(event.error) + DispatchQueue.main.async { + completion(.failure(error)) + } + return + } + + // The result node belongs to the event: read it here, on the event queue. + var playlistEntryId: Int64? + if let commandPointer = event.data?.assumingMemoryBound(to: mpv_event_command.self) { + playlistEntryId = Self.playlistEntryId(in: commandPointer.pointee.result) + } + DispatchQueue.main.async { + completion(.success(playlistEntryId)) + } + } + + /// The `playlist_entry_id` of a `loadfile` result map; nil for any other + /// command result. + private static func playlistEntryId(in result: mpv_node) -> Int64? { + guard result.format == MPV_FORMAT_NODE_MAP, let list = result.u.list else { return nil } + let map = list.pointee + guard map.num > 0, let keys = map.keys, let values = map.values else { return nil } + for index in 0..; +// `result` is the command's return node (owned by the reply event, valid only +// for the duration of the callback) or nullptr when the command failed or +// returned nothing. +using CommandCallback = std::function; using GetPropertyCallback = std::function; static constexpr char kSetPropertyFailedCode[] = "SET_PROPERTY_FAILED"; @@ -46,6 +50,7 @@ inline std::string SetPropertyErrorDescription(int status) { struct CancelledRequests { std::vector status; + std::vector commands; std::vector properties; }; @@ -67,6 +72,22 @@ class AsyncRequestRegistry { return callback; } + uint64_t RegisterCommand(CommandCallback callback) { + std::lock_guard lock(mutex_); + const uint64_t request_id = next_id_++; + commands_[request_id] = std::move(callback); + return request_id; + } + + CommandCallback TakeCommand(uint64_t request_id) { + std::lock_guard lock(mutex_); + auto it = commands_.find(request_id); + if (it == commands_.end()) return nullptr; + auto callback = std::move(it->second); + commands_.erase(it); + return callback; + } + uint64_t RegisterProperty(GetPropertyCallback callback) { std::lock_guard lock(mutex_); const uint64_t request_id = next_id_++; @@ -92,6 +113,12 @@ class AsyncRequestRegistry { cancelled.status.push_back(std::move(request.second)); } } + cancelled.commands.reserve(commands_.size()); + for (auto& request : commands_) { + if (request.second) { + cancelled.commands.push_back(std::move(request.second)); + } + } cancelled.properties.reserve(properties_.size()); for (auto& request : properties_) { if (request.second) { @@ -99,6 +126,7 @@ class AsyncRequestRegistry { } } status_.clear(); + commands_.clear(); properties_.clear(); return cancelled; } @@ -106,6 +134,7 @@ class AsyncRequestRegistry { private: uint64_t next_id_ = 1; std::map status_; + std::map commands_; std::map properties_; std::mutex mutex_; }; @@ -117,7 +146,7 @@ class AsyncRequestRegistry { // handle guard stays with the callers: they own the lifecycle flags and the // error code they report when the handle is gone. inline void SubmitCommandAsync( - mpv_handle* mpv, AsyncRequestRegistry& requests, const std::vector& args, StatusCallback callback) { + mpv_handle* mpv, AsyncRequestRegistry& requests, const std::vector& args, CommandCallback callback) { std::vector c_args; c_args.reserve(args.size() + 1); for (const auto& arg : args) { @@ -125,15 +154,31 @@ inline void SubmitCommandAsync( } c_args.push_back(nullptr); - const uint64_t request_id = callback ? requests.RegisterStatus(std::move(callback)) : 0; + const uint64_t request_id = callback ? requests.RegisterCommand(std::move(callback)) : 0; // mpv_command_async returns immediately. const int result = mpv_command_async(mpv, request_id, c_args.data()); if (result < 0) { - auto pending = requests.TakeStatus(request_id); - if (pending) pending(result); + auto pending = requests.TakeCommand(request_id); + if (pending) pending(result, nullptr); } } +// The playlist entry a `loadfile` created, read from the command's result map. +// mpv reports it as `playlist_entry_id`; other commands return no map, so a +// caller gets false for them and answers with no id. +inline bool PlaylistEntryIdFromCommandResult(const mpv_node* result, int64_t* playlist_entry_id) { + if (!result || !playlist_entry_id || result->format != MPV_FORMAT_NODE_MAP) return false; + const mpv_node_list* map = result->u.list; + if (!map || map->num <= 0 || !map->keys || !map->values) return false; + for (int i = 0; i < map->num; i++) { + if (!map->keys[i] || strcmp(map->keys[i], "playlist_entry_id") != 0) continue; + if (map->values[i].format != MPV_FORMAT_INT64) return false; + *playlist_entry_id = map->values[i].u.int64; + return true; + } + return false; +} + inline void SubmitSetPropertyAsync( mpv_handle* mpv, AsyncRequestRegistry& requests, const std::string& name, const std::string& value, StatusCallback callback) { @@ -163,7 +208,17 @@ inline void SubmitGetPropertyAsync( template inline bool DispatchReplyEvent(AsyncRequestRegistry& requests, const mpv_event* event, const Sanitizer& sanitize) { switch (event->event_id) { - case MPV_EVENT_COMMAND_REPLY: + case MPV_EVENT_COMMAND_REPLY: { + CommandCallback callback = requests.TakeCommand(event->reply_userdata); + if (callback) { + const mpv_node* result = nullptr; + if (event->error >= 0 && event->data) { + result = &static_cast(event->data)->result; + } + callback(event->error, result); + } + return true; + } case MPV_EVENT_SET_PROPERTY_REPLY: { StatusCallback callback = requests.TakeStatus(event->reply_userdata); if (callback) { diff --git a/shared/mpv/mpv_player_common_test.cpp b/shared/mpv/mpv_player_common_test.cpp index 189b7cba2..2073c02a0 100644 --- a/shared/mpv/mpv_player_common_test.cpp +++ b/shared/mpv/mpv_player_common_test.cpp @@ -20,30 +20,119 @@ using plezy::mpv_common::AudioReloadReason; void TestRequestRegistry() { plezy::mpv_common::AsyncRequestRegistry registry; bool status_called = false; + bool command_called = false; bool property_called = false; const auto status_id = registry.RegisterStatus([&](int error) { status_called = error == -7; }); + mpv_node command_node{}; + const auto command_id = registry.RegisterCommand( + [&](int error, const mpv_node* result) { command_called = error == -9 && result == &command_node; }); const auto property_id = registry.RegisterProperty( [&](int error, const std::string& value) { property_called = error == -8 && value == "value"; }); auto status = registry.TakeStatus(status_id); + auto command = registry.TakeCommand(command_id); auto property = registry.TakeProperty(property_id); assert(status); + assert(command); assert(property); status(-7); + command(-9, &command_node); property(-8, "value"); assert(status_called); + assert(command_called); assert(property_called); assert(!registry.TakeStatus(status_id)); + assert(!registry.TakeCommand(command_id)); assert(!registry.TakeProperty(property_id)); + // Ids are one namespace: a command id must not surface as another type. + assert(!registry.TakeStatus(command_id)); registry.RegisterStatus([](int) {}); + registry.RegisterCommand([](int, const mpv_node*) {}); registry.RegisterProperty([](int, const std::string&) {}); auto cancelled = registry.CancelAll(); assert(cancelled.status.size() == 1); + assert(cancelled.commands.size() == 1); assert(cancelled.properties.size() == 1); } +// The command reply hands the registered callback mpv's result node only for +// a successful reply; `PlaylistEntryIdFromCommandResult` then reads the entry +// `loadfile` created and refuses anything that is not that map. +void TestCommandReplyResultContract() { + using namespace plezy::mpv_common; + const auto sanitize = [](const char* value) { return std::string(value); }; + + AsyncRequestRegistry registry; + const mpv_node* seen_result = nullptr; + int seen_error = 0; + const auto id = registry.RegisterCommand([&](int error, const mpv_node* result) { + seen_error = error; + seen_result = result; + }); + + const char* keys[] = {"playlist_entry_id"}; + mpv_node values[1]{}; + values[0].format = MPV_FORMAT_INT64; + values[0].u.int64 = 9000000005LL; + mpv_node_list map{}; + map.num = 1; + map.keys = const_cast(keys); + map.values = values; + mpv_event_command command{}; + command.result.format = MPV_FORMAT_NODE_MAP; + command.result.u.list = ↦ + mpv_event reply{}; + reply.event_id = MPV_EVENT_COMMAND_REPLY; + reply.reply_userdata = id; + reply.data = &command; + assert(DispatchReplyEvent(registry, &reply, sanitize)); + assert(seen_error == 0); + assert(seen_result == &command.result); + int64_t entry_id = 0; + assert(PlaylistEntryIdFromCommandResult(seen_result, &entry_id)); + assert(entry_id == 9000000005LL); + + // A failed reply never exposes the result slot. + const auto failed_id = registry.RegisterCommand([&](int error, const mpv_node* result) { + seen_error = error; + seen_result = result; + }); + reply.reply_userdata = failed_id; + reply.error = MPV_ERROR_COMMAND; + assert(DispatchReplyEvent(registry, &reply, sanitize)); + assert(seen_error == MPV_ERROR_COMMAND); + assert(seen_result == nullptr); + + // A SET_PROPERTY_REPLY with a command's id completes nothing: the types are + // distinct registries sharing one id space. + const auto stranded_id = registry.RegisterCommand([&](int, const mpv_node*) { assert(false); }); + mpv_event property_reply{}; + property_reply.event_id = MPV_EVENT_SET_PROPERTY_REPLY; + property_reply.reply_userdata = stranded_id; + assert(DispatchReplyEvent(registry, &property_reply, sanitize)); + assert(registry.TakeCommand(stranded_id)); + + // Only an INT64 `playlist_entry_id` inside a map counts. + assert(!PlaylistEntryIdFromCommandResult(nullptr, &entry_id)); + mpv_node none{}; + none.format = MPV_FORMAT_NONE; + assert(!PlaylistEntryIdFromCommandResult(&none, &entry_id)); + const char* other_keys[] = {"other"}; + mpv_node_list other_map{}; + other_map.num = 1; + other_map.keys = const_cast(other_keys); + other_map.values = values; + mpv_node other{}; + other.format = MPV_FORMAT_NODE_MAP; + other.u.list = &other_map; + assert(!PlaylistEntryIdFromCommandResult(&other, &entry_id)); + values[0].format = MPV_FORMAT_DOUBLE; + values[0].u.double_ = 1.0; + assert(!PlaylistEntryIdFromCommandResult(&command.result, &entry_id)); +} + void TestConcurrentRequestCompletion() { for (int iteration = 0; iteration < 200; ++iteration) { plezy::mpv_common::AsyncRequestRegistry registry; @@ -448,6 +537,7 @@ void TestHdrHelpers() { int main() { TestRequestRegistry(); + TestCommandReplyResultContract(); TestConcurrentRequestCompletion(); TestSetPropertyResultContract(); TestPropertyObservationRegistry(); diff --git a/test/mpv/player_open_test.dart b/test/mpv/player_open_test.dart index 8cc937bbe..492f55fcc 100644 --- a/test/mpv/player_open_test.dart +++ b/test/mpv/player_open_test.dart @@ -955,6 +955,47 @@ void main() { ); }); + test('MPV open resolves with the playlist entry id the loadfile reply names', () async { + await withMockPlayerChannels( + methodChannelName: 'com.plezy/mpv_player', + eventChannelName: 'com.plezy/mpv_player/events', + methodHandler: (call) { + switch (call.method) { + case 'initialize': + return Future.value(true); + case 'command': + final args = Map.from(call.arguments as Map)['args'] as List; + return Future.value(args.first == 'loadfile' ? {'playlistEntryId': 17} : null); + default: + return Future.value(null); + } + }, + testBody: () async { + final player = PlayerNative(); + try { + expect(await player.open(Media('https://example.test/live.m3u8'), isLive: true), 17); + } finally { + await player.dispose(); + } + }, + ); + }); + + test('MPV open resolves null when the core does not name the source', () async { + await withMockPlayerChannels( + methodChannelName: 'com.plezy/mpv_player', + eventChannelName: 'com.plezy/mpv_player/events', + testBody: () async { + final player = PlayerNative(); + try { + expect(await player.open(Media('https://example.test/live.m3u8'), isLive: true), isNull); + } finally { + await player.dispose(); + } + }, + ); + }); + test('MPV source readiness carries the first rendered non-zero clock position once', () async { await withMockPlayerChannels( methodChannelName: 'com.plezy/mpv_player', diff --git a/test/screens/video_player/live_tv_session_state_test.dart b/test/screens/video_player/live_tv_session_state_test.dart index aebc19a77..e0b2f7cb2 100644 --- a/test/screens/video_player/live_tv_session_state_test.dart +++ b/test/screens/video_player/live_tv_session_state_test.dart @@ -64,7 +64,7 @@ void main() { final result = state.clockOpenResult(generation); expect(state.epochForPosition(const Duration(seconds: 52)), 1093); - expect(state.bindClockSource(const PlayerSourceStarted(7)), isTrue); + expect(state.bindClockOpen(generation, 7), isTrue); expect(state.calibrateClockSource(const PlayerSourceReady(sourceId: 7, position: Duration(seconds: 47))), isTrue); expect(await result, isTrue); @@ -79,8 +79,8 @@ void main() { final secondGeneration = state.beginClockOpen(1070); final secondResult = state.clockOpenResult(secondGeneration); - expect(state.bindClockSource(const PlayerSourceStarted(11)), isFalse); - expect(state.bindClockSource(const PlayerSourceStarted(12)), isTrue); + expect(state.bindClockOpen(firstGeneration, 11), isFalse); + expect(state.bindClockOpen(secondGeneration, 12), isTrue); expect( state.calibrateClockSource(const PlayerSourceReady(sourceId: 11, position: Duration(seconds: 52))), isFalse, @@ -95,11 +95,104 @@ void main() { expect(state.epochForPosition(const Duration(seconds: 45)), 1075); }); + test('readiness that beats the loadfile reply calibrates on bind', () async { + final state = LiveTvSessionState(null); + final generation = state.beginClockOpen(1093); + final result = state.clockOpenResult(generation); + + expect( + state.calibrateClockSource(const PlayerSourceReady(sourceId: 7, position: Duration(seconds: 47))), + isFalse, + ); + expect(state.epochForPosition(const Duration(seconds: 52)), 1093); + + expect(state.bindClockOpen(generation, 7), isTrue); + expect(await result, isTrue); + expect(state.streamStartEpoch, 1046); + expect(state.epochForPosition(const Duration(seconds: 52)), 1098); + }); + + test('a failure that beats the loadfile reply fails the open on bind', () async { + final state = LiveTvSessionState(null); + final generation = state.beginClockOpen(1093); + final result = state.clockOpenResult(generation); + + state.failClockSource(const PlayerSourceFailed(7)); + expect(state.bindClockOpen(generation, 7), isFalse); + + expect(await result, isFalse); + expect(state.pendingStreamEpoch, isNull); + expect( + state.calibrateClockSource(const PlayerSourceReady(sourceId: 7, position: Duration(seconds: 47))), + isFalse, + ); + }); + + test('a rejected first load cannot claim the second seek\'s source', () async { + final state = LiveTvSessionState(null)..streamStartEpoch = 1000; + final firstGeneration = state.beginClockOpen(1085); + final firstResult = state.clockOpenResult(firstGeneration); + final secondGeneration = state.beginClockOpen(1075); + final secondResult = state.clockOpenResult(secondGeneration); + + // The second load's source reports before either reply lands. + state.calibrateClockSource(const PlayerSourceReady(sourceId: 12, position: Duration(seconds: 40))); + // mpv rejected the first loadfile: no source ever existed for it. + state.failClockOpen(firstGeneration); + expect(await firstResult, isFalse); + expect(state.epochForPosition(const Duration(seconds: 40)), 1075); + + expect(state.bindClockOpen(secondGeneration, 12), isTrue); + expect(await secondResult, isTrue); + expect(state.streamStartEpoch, 1035); + expect(state.epochForPosition(const Duration(seconds: 45)), 1080); + }); + + test('an unregistered open on the same player is invisible to clock binding', () async { + final state = LiveTvSessionState(null); + final firstGeneration = state.beginClockOpen(1085); + expect(state.bindClockOpen(firstGeneration, 11), isTrue); + state.calibrateClockSource(const PlayerSourceReady(sourceId: 11, position: Duration(seconds: 10))); + expect(state.streamStartEpoch, 1075); + + // A live-edge re-open registers no clock generation; its source reports + // and is never claimed. + state.calibrateClockSource(const PlayerSourceReady(sourceId: 12, position: Duration(seconds: 3))); + expect(state.streamStartEpoch, 1075); + + final secondGeneration = state.beginClockOpen(1070); + final secondResult = state.clockOpenResult(secondGeneration); + expect(state.bindClockOpen(secondGeneration, 13), isTrue); + expect(state.epochForPosition(const Duration(seconds: 5)), 1070); + expect( + state.calibrateClockSource(const PlayerSourceReady(sourceId: 13, position: Duration(seconds: 40))), + isTrue, + ); + expect(await secondResult, isTrue); + expect(state.streamStartEpoch, 1030); + }); + + test('unclaimed source reports are bounded to the most recent sources', () async { + final state = LiveTvSessionState(null); + final generation = state.beginClockOpen(1093); + final result = state.clockOpenResult(generation); + + state.calibrateClockSource(const PlayerSourceReady(sourceId: 1, position: Duration(seconds: 47))); + for (var sourceId = 2; sourceId <= 9; sourceId++) { + state.calibrateClockSource(PlayerSourceReady(sourceId: sourceId, position: Duration.zero)); + } + + expect(state.bindClockOpen(generation, 1), isTrue); + expect(state.epochForPosition(const Duration(seconds: 52)), 1093); + expect(state.calibrateClockSource(const PlayerSourceReady(sourceId: 1, position: Duration(seconds: 47))), isTrue); + expect(await result, isTrue); + }); + test('a calibration timeout keeps the target until late readiness arrives', () async { final state = LiveTvSessionState(null); final generation = state.beginClockOpen(1093); final result = state.clockOpenResult(generation); - state.bindClockSource(const PlayerSourceStarted(7)); + state.bindClockOpen(generation, 7); state.timeoutClockOpen(generation); @@ -109,11 +202,25 @@ void main() { expect(state.epochForPosition(const Duration(seconds: 52)), 1098); }); + test('a timed-out open still binds when its loadfile reply arrives late', () async { + final state = LiveTvSessionState(null); + final generation = state.beginClockOpen(1093); + final result = state.clockOpenResult(generation); + + state.timeoutClockOpen(generation); + expect(await result, isFalse); + state.calibrateClockSource(const PlayerSourceReady(sourceId: 7, position: Duration(seconds: 47))); + expect(state.epochForPosition(const Duration(seconds: 52)), 1093); + + expect(state.bindClockOpen(generation, 7), isTrue); + expect(state.epochForPosition(const Duration(seconds: 52)), 1098); + }); + test('a zero-based source preserves the existing epoch mapping', () async { final state = LiveTvSessionState(null); final generation = state.beginClockOpen(1093); final result = state.clockOpenResult(generation); - state.bindClockSource(const PlayerSourceStarted(7)); + state.bindClockOpen(generation, 7); state.calibrateClockSource(const PlayerSourceReady(sourceId: 7, position: Duration.zero)); expect(await result, isTrue); @@ -132,7 +239,7 @@ void main() { final generation = state.beginClockOpen(targetEpoch); final result = state.clockOpenResult(generation); final currentSourceId = ++sourceId; - state.bindClockSource(PlayerSourceStarted(currentSourceId)); + state.bindClockOpen(generation, currentSourceId); if (requestedEpochs.length == 1) { state.calibrateClockSource( PlayerSourceReady(sourceId: currentSourceId, position: const Duration(seconds: 52)), @@ -198,7 +305,7 @@ void main() { state.streamGeneration++; final generation = state.beginClockOpen(990); final result = state.clockOpenResult(generation); - state.bindClockSource(const PlayerSourceStarted(3)); + state.bindClockOpen(generation, 3); state.calibrateClockSource(const PlayerSourceReady(sourceId: 3, position: Duration.zero)); expect(await result, isTrue); diff --git a/windows/runner/mpv/mpv_player.cpp b/windows/runner/mpv/mpv_player.cpp index eda14439f..b6de0bb9a 100644 --- a/windows/runner/mpv/mpv_player.cpp +++ b/windows/runner/mpv/mpv_player.cpp @@ -658,6 +658,9 @@ void MpvPlayer::Dispose() { for (auto& callback : cancelled.status) { callback(MPV_ERROR_UNINITIALIZED); } + for (auto& callback : cancelled.commands) { + callback(MPV_ERROR_UNINITIALIZED, nullptr); + } for (auto& callback : cancelled.properties) { callback(-1, ""); } @@ -690,7 +693,7 @@ void MpvPlayer::Command(const std::vector& args) { CommandAsync(arg void MpvPlayer::CommandAsync(const std::vector& args, CommandCallback callback) { if (!mpv_) { - if (callback) callback(0); + if (callback) callback(0, nullptr); return; } @@ -798,7 +801,7 @@ void MpvPlayer::LogHdrPipelineOnce() { void MpvPlayer::TryAudioReload(const char* reason, int attempt, uint64_t request_generation) { LogRecovery("issuing ao-reload (reason=" + std::string(reason) + ", attempt " + std::to_string(attempt) + ")"); const std::string reason_copy = reason; - CommandAsync({"ao-reload"}, [this, reason_copy, attempt, request_generation](int error) { + CommandAsync({"ao-reload"}, [this, reason_copy, attempt, request_generation](int error, const mpv_node*) { audio_recovery_.CompleteReload(request_generation); LogRecovery( "ao-reload completed (reason=" + reason_copy + ", attempt " + std::to_string(attempt) + diff --git a/windows/runner/mpv/mpv_player.h b/windows/runner/mpv/mpv_player.h index 474bbe91b..3021be722 100644 --- a/windows/runner/mpv/mpv_player.h +++ b/windows/runner/mpv/mpv_player.h @@ -54,7 +54,7 @@ class MpvPlayer { // Callback types for async mpv requests. using StatusCallback = plezy::mpv_common::StatusCallback; - using CommandCallback = StatusCallback; + using CommandCallback = plezy::mpv_common::CommandCallback; using GetPropertyCallback = plezy::mpv_common::GetPropertyCallback; // Executes an mpv command asynchronously to prevent UI blocking. diff --git a/windows/runner/mpv/mpv_player_property_contract_test.cpp b/windows/runner/mpv/mpv_player_property_contract_test.cpp index 36df26a78..55eef6375 100644 --- a/windows/runner/mpv/mpv_player_property_contract_test.cpp +++ b/windows/runner/mpv/mpv_player_property_contract_test.cpp @@ -23,6 +23,10 @@ class MpvPlayerPropertyContractTestPeer { static void RegisterPendingPropertyRead(MpvPlayer& player, MpvPlayer::GetPropertyCallback callback) { player.pending_requests_.RegisterProperty(std::move(callback)); } + + static uint64_t RegisterPendingCommand(MpvPlayer& player, MpvPlayer::CommandCallback callback) { + return player.pending_requests_.RegisterCommand(std::move(callback)); + } static void RegisterObservedNode(MpvPlayer& player, const std::string& name, int id) { player.observed_properties_.Register(name, "node", id); } @@ -412,6 +416,82 @@ void TestPendingPropertyWriteFailsOnDispose() { Check(callback_count == 1, "repeated dispose must not complete a property write twice"); } +// The `loadfile` reply names the playlist entry mpv created for the load — +// the source id its start-file/playback-restart/end-file events carry — so the +// Dart side can bind a load to its source instead of guessing by arrival order. +void TestCommandReplyCarriesPlaylistEntryId() { + MpvPlayer player; + constexpr int64_t kEntryId = 8000000004LL; + + int callback_count = 0; + int status = MPV_ERROR_SUCCESS; + int64_t reported_entry_id = 0; + bool reported_entry = false; + const uint64_t request_id = + MpvPlayerPropertyContractTestPeer::RegisterPendingCommand(player, [&](int error, const mpv_node* result) { + ++callback_count; + status = error; + reported_entry = plezy::mpv_common::PlaylistEntryIdFromCommandResult(result, &reported_entry_id); + }); + + const char* keys[] = {"playlist_entry_id"}; + mpv_node values[1]{}; + values[0].format = MPV_FORMAT_INT64; + values[0].u.int64 = kEntryId; + mpv_node_list map{}; + map.num = 1; + map.keys = const_cast(keys); + map.values = values; + mpv_event_command command{}; + command.result.format = MPV_FORMAT_NODE_MAP; + command.result.u.list = ↦ + mpv_event reply{}; + reply.event_id = MPV_EVENT_COMMAND_REPLY; + reply.reply_userdata = request_id; + reply.data = &command; + MpvPlayerPropertyContractTestPeer::HandleEvent(player, &reply); + + Check(callback_count == 1, "a command reply must complete its request exactly once"); + Check(status == MPV_ERROR_SUCCESS, "a successful command reply changed its status"); + Check(reported_entry, "the loadfile reply must expose the playlist entry id"); + Check(reported_entry_id == kEntryId, "the playlist entry id lost signed 64-bit precision"); + + // A reply without a result map (every non-loadfile command) answers no id. + int64_t unexpected_entry_id = 0; + bool unexpected_entry = true; + const uint64_t plain_request_id = + MpvPlayerPropertyContractTestPeer::RegisterPendingCommand(player, [&](int, const mpv_node* result) { + unexpected_entry = plezy::mpv_common::PlaylistEntryIdFromCommandResult(result, &unexpected_entry_id); + }); + mpv_event_command plain_command{}; + plain_command.result.format = MPV_FORMAT_NONE; + reply.reply_userdata = plain_request_id; + reply.data = &plain_command; + MpvPlayerPropertyContractTestPeer::HandleEvent(player, &reply); + Check(!unexpected_entry, "a command without a result map must not report a playlist entry id"); + + // A failed reply must not expose whatever the event's result slot holds. + const mpv_node* failed_result = &command.result; + const uint64_t failed_request_id = MpvPlayerPropertyContractTestPeer::RegisterPendingCommand( + player, [&](int, const mpv_node* result) { failed_result = result; }); + reply.reply_userdata = failed_request_id; + reply.error = MPV_ERROR_COMMAND; + reply.data = &command; + MpvPlayerPropertyContractTestPeer::HandleEvent(player, &reply); + Check(failed_result == nullptr, "a failed command reply must not carry a result node"); + + // Dispose completes a pending command as uninitialized, with no result. + int cancelled_status = MPV_ERROR_SUCCESS; + const mpv_node* cancelled_result = &command.result; + MpvPlayerPropertyContractTestPeer::RegisterPendingCommand(player, [&](int error, const mpv_node* result) { + cancelled_status = error; + cancelled_result = result; + }); + player.Dispose(); + Check(cancelled_status == MPV_ERROR_UNINITIALIZED, "dispose must cancel a pending command as uninitialized"); + Check(cancelled_result == nullptr, "a cancelled command must not carry a result node"); +} + void TestPendingRequestTypesRemainDistinctOnDispose() { MpvPlayer player; int write_count = 0; @@ -976,6 +1056,7 @@ int main() { mpv::TestSourceQualifiedEventPayloads(); mpv::TestUnavailablePropertyWriteFails(); mpv::TestPendingPropertyWriteFailsOnDispose(); + mpv::TestCommandReplyCarriesPlaylistEntryId(); mpv::TestPendingRequestTypesRemainDistinctOnDispose(); mpv::TestInnerSubclassOwnershipIsSerializedAndDetached(); mpv::TestTimedOutSubclassDetachCanBeAdopted(); diff --git a/windows/runner/mpv/mpv_plugin.cpp b/windows/runner/mpv/mpv_plugin.cpp index 7950c5249..1d6ceffaf 100644 --- a/windows/runner/mpv/mpv_plugin.cpp +++ b/windows/runner/mpv/mpv_plugin.cpp @@ -235,11 +235,22 @@ void MpvPlayerPlugin::HandleMethodCall( auto result_ptr = std::make_shared>>(std::move(result)); std::string cmd_name = command_args.empty() ? "unknown" : command_args[0]; - player_->CommandAsync(command_args, [this, result_ptr, cmd_name](int error) { - PostToPlatformThread([result_ptr, cmd_name, error]() { + player_->CommandAsync(command_args, [this, result_ptr, cmd_name](int error, const mpv_node* command_result) { + // `loadfile` answers with the playlist entry it created so Dart can tie + // the load to that source's start-file/playback-restart/end-file events; + // every other command answers null. The node belongs to the reply event, + // so the id is read here, before the hop to the platform thread. + int64_t playlist_entry_id = 0; + const bool has_playlist_entry = + error >= 0 && plezy::mpv_common::PlaylistEntryIdFromCommandResult(command_result, &playlist_entry_id); + PostToPlatformThread([result_ptr, cmd_name, error, has_playlist_entry, playlist_entry_id]() { if (error < 0) { (*result_ptr) ->Error("COMMAND_FAILED", "MPV command failed: " + cmd_name + " (error " + std::to_string(error) + ")"); + } else if (has_playlist_entry) { + flutter::EncodableMap reply; + reply[flutter::EncodableValue("playlistEntryId")] = flutter::EncodableValue(playlist_entry_id); + (*result_ptr)->Success(flutter::EncodableValue(reply)); } else { (*result_ptr)->Success(); }