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(); }