Rework transition to playing state after seek. Trigger the transition after preroll and both audio/video has buffered the targeted seek time. Change-Id: I731da77fec101ade1ed0a5ad43c4e8cf642511bc Signed-off-by: Ievgen Mutavchi <Ievgen_Mutavchi@comcast.com>
diff --git a/src/third_party/starboard/rdk/shared/player/player_internal.cc b/src/third_party/starboard/rdk/shared/player/player_internal.cc index 3fe2993..a108110 100644 --- a/src/third_party/starboard/rdk/shared/player/player_internal.cc +++ b/src/third_party/starboard/rdk/shared/player/player_internal.cc
@@ -33,6 +33,7 @@ #include <string> #include <vector> #include <algorithm> +#include <functional> #include <cstring> #include "starboard/once.h" @@ -55,6 +56,7 @@ static constexpr int kMaxNumberOfSamplesPerWrite = 1; static const char kCustomInstantRateChangeEventName[] = "custom-instant-rate-change"; static const char kDidReceiveFirstSegmentMsgName[] = "did-receive-first-segment"; +static const char kDidReachBufferingTargetMsgName[] = "did-reach-buffering-target"; // static int Player::MaxNumberOfSamplesPerWrite() { @@ -470,7 +472,7 @@ #if GST_CHECK_VERSION(1,18,0) gst_pad_add_probe ( - pad, GST_PAD_PROBE_TYPE_EVENT_UPSTREAM, + pad, GST_PAD_PROBE_TYPE_EVENT_BOTH, [](GstPad * pad, GstPadProbeInfo * info, gpointer data) -> GstPadProbeReturn { GstEvent *event = GST_PAD_PROBE_INFO_EVENT(info); GstSegment *segment = reinterpret_cast<GstSegment*>(data); @@ -1078,6 +1080,22 @@ std::string msg_; }; +class FunctionTask: public Task { +public: + FunctionTask(std::function<void()> && fn, const char debug_msg[]) + : func_(std::move(fn)) + , msg_(debug_msg) { } + + ~FunctionTask() override {} + + void Do() override { func_(); } + + void PrintInfo() override { GST_TRACE("FunctionTask: %s", msg_); } +private: + std::function<void()> func_; + const char* msg_; +}; + class PlayerImpl : public Player { public: PlayerImpl(SbPlayer player, @@ -1223,6 +1241,8 @@ static void SetupElement(GstElement* pipeline, GstElement* element, PlayerImpl* self); + static void OnVideoBufferUnderflow(PlayerImpl* self); + bool ChangePipelineState(GstState state) const; void DispatchOnWorkerThread(Task* task) const; gint64 GetPosition() const; @@ -1253,6 +1273,8 @@ void WritePendingSamples(); void CheckBuffering(gint64 position); void ConfigureLimitedVideo(); + void SchedulePlayingStateUpdate(); + void AddBufferingProbe(GstClockTime target, int ticket); SbPlayer player_; SbWindow window_; @@ -1313,6 +1335,8 @@ SbTime buf_target_min_ts_ { kSbTimeMax }; bool need_instant_rate_change_ { false }; int need_first_segment_ack_ { static_cast<int>(MediaType::kBoth) }; + int buffering_state_ { 0 }; + mutable int playing_state_update_source_id_{0u}; }; struct PlayerRegistry @@ -1644,7 +1668,7 @@ self->SetBounds(0, bounds.x, bounds.y, bounds.w, bounds.h); } - if (is_rate_pending) { + if (is_rate_pending && GST_STATE(self->pipeline_) == GST_STATE_PLAYING) { GST_INFO("Sending pending SetRate(rate=%lf)", rate); self->SetRate(rate); } else if (is_seek_pending) { @@ -1695,6 +1719,8 @@ self->context_, kSbPlayerStatePresenting)); self->state_ = State::kPresenting; } + + self->SchedulePlayingStateUpdate(); } } break; @@ -1889,6 +1915,17 @@ } // static +void PlayerImpl::OnVideoBufferUnderflow(PlayerImpl* self) +{ + GST_WARNING("Decoder need data state = 0x%x," + " video appsrc level = %lld kb," + " audio appsrc level = %lld kb", + self->decoder_state_data_, + gst_app_src_get_current_level_bytes(GST_APP_SRC(self->video_appsrc_)) / 1024, + gst_app_src_get_current_level_bytes(GST_APP_SRC(self->audio_appsrc_)) / 1024); +} + +// static void PlayerImpl::SetupElement(GstElement* pipeline, GstElement* element, PlayerImpl* self) { @@ -1896,7 +1933,7 @@ static bool disable_wait_video = !!getenv("COBALT_AML_DISABLE_WAIT_VIDEO"); bool has_video = (self->video_codec_ != kSbMediaVideoCodecNone); if (has_video && g_str_has_prefix(GST_ELEMENT_NAME(element), "amlhalasink") && !disable_wait_video) { - g_object_set(element, "wait-video", TRUE, nullptr); + g_object_set(element, "wait-video", TRUE, "a-wait-timeout", 4000, nullptr); } else if (has_video && g_str_has_prefix(GST_ELEMENT_NAME(element), "westerossink")) { @@ -1904,6 +1941,9 @@ GST_INFO("Setting westerossink zoom-mode to 0"); g_object_set(element, "zoom-mode", 0, nullptr); } + g_signal_connect_swapped( + G_OBJECT(element), "buffer-underflow-callback", + G_CALLBACK(OnVideoBufferUnderflow), self); } else if (g_str_has_prefix(GST_ELEMENT_NAME(element), "brcmaudiosink")) { @@ -2044,26 +2084,6 @@ RecordTimestamp(sample_type, timestamp); - if (MinTimestamp(nullptr) == GST_BUFFER_TIMESTAMP(buffer) && - GST_STATE(pipeline_) <= GST_STATE_PAUSED && - (GST_STATE_PENDING(pipeline_) == GST_STATE_VOID_PENDING || - GST_STATE_PENDING(pipeline_) == GST_STATE_PAUSED) && - rate_ > .0) { - - gint64 seek_pos_ns = GST_CLOCK_TIME_NONE; - { - ::starboard::ScopedLock lock(mutex_); - if (seek_position_ != kSbTimeMax) - seek_pos_ns = seek_position_ * kSbTimeNanosecondsPerMicrosecond; - } - - if (!GST_CLOCK_TIME_IS_VALID(seek_pos_ns) || GST_BUFFER_TIMESTAMP(buffer) >= seek_pos_ns) { - GST_TRACE("Moving to playing for %" GST_TIME_FORMAT, - GST_TIME_ARGS(GST_BUFFER_TIMESTAMP(buffer))); - ChangePipelineState(GST_STATE_PLAYING); - } - } - if (sample_infos[0].drm_info) { GST_LOG("Encounterd encrypted %s sample", sample_type == kSbMediaTypeVideo ? "video" : "audio"); @@ -2294,7 +2314,6 @@ is_seek_pending_ = false; rate = rate_; - state_ = State::kPrerollAfterSeek; } GST_DEBUG("Calling seek"); @@ -2308,48 +2327,61 @@ seek_to_timestamp * kSbTimeNanosecondsPerMicrosecond, GST_SEEK_TYPE_NONE, 0)) { GST_ERROR_OBJECT(pipeline_, "Seek failed"); - ::starboard::ScopedLock lock(mutex_); DispatchOnWorkerThread(new PlayerStatusTask(player_status_func_, player_, ticket_, context_, kSbPlayerStatePresenting)); - state_ = State::kPresenting; + DispatchOnWorkerThread(new FunctionTask([this]() { + state_ = State::kPresenting; + }, "Presenting after seek failure")); } else { GST_DEBUG("Seek called with success"); + DispatchOnWorkerThread(new FunctionTask([this]() { + state_ = State::kPrerollAfterSeek; + }, "Preroll after seek")); } + + AddBufferingProbe(seek_to_timestamp * kSbTimeNanosecondsPerMicrosecond, ticket); } bool PlayerImpl::SetRate(double rate) { GST_DEBUG_OBJECT(pipeline_, "===> rate %lf (rate_ %lf), TID: %d", rate, rate_, SbThreadGetId()); + GstState state; + double old_rate; bool success = true; - { - ::starboard::ScopedLock lock(mutex_); - decoder_state_data_ = 0; - eos_data_ = 0; - } + + mutex_.Acquire(); + old_rate = rate_; + rate_ = rate; + pending_rate_ = .0; if (rate == .0) { + mutex_.Release(); ChangePipelineState(GST_STATE_PAUSED); - } else { - ChangePipelineState(GST_STATE_PLAYING); + return true; } - if (rate != .0 && (rate != 1. || need_instant_rate_change_)) { - { - ::starboard::ScopedLock lock(mutex_); - if (is_seek_pending_) { - GST_DEBUG_OBJECT(pipeline_, "Rate will be set when doing seek"); - rate_ = rate; - return true; - } - if (GST_STATE(pipeline_) < GST_STATE_PLAYING || need_first_segment_ack_) { - GST_DEBUG_OBJECT(pipeline_, "===> Set rate postponed"); - pending_rate_ = rate; - return true; - } - pending_rate_ = .0; + gst_element_get_state(pipeline_, &state, nullptr, 0); + + if (state < GST_STATE_PLAYING) + SchedulePlayingStateUpdate(); + + if (rate != 1. || need_instant_rate_change_) { + if (is_seek_pending_) { + mutex_.Release(); + GST_DEBUG_OBJECT(pipeline_, "Rate will be set when doing seek"); + return true; } + if (state < GST_STATE_PLAYING || need_first_segment_ack_) { + rate_ = old_rate; + pending_rate_ = rate; + mutex_.Release(); + GST_DEBUG_OBJECT(pipeline_, "===> Set rate postponed"); + return true; + } + need_instant_rate_change_ = ( rate != 1. ); + mutex_.Release(); #if GST_CHECK_VERSION(1,18,0) static const bool kEnableInstantRateChangeSeek = ([]()->bool { @@ -2380,15 +2412,14 @@ pipeline_, gst_event_new_custom(GST_EVENT_CUSTOM_DOWNSTREAM_OOB, s)); } - need_instant_rate_change_ = ( rate != 1. ); + mutex_.Acquire(); } - if (success) { - ::starboard::ScopedLock lock(mutex_); - rate_ = rate; - } else { + if (!success) { + rate_ = old_rate; GST_ERROR_OBJECT(pipeline_, "Set rate failed"); } + mutex_.Release(); return success; } @@ -2496,6 +2527,14 @@ } } } + else { + ::starboard::ScopedLock lock(mutex_); + guint src_id = std::exchange(playing_state_update_source_id_, 0u); + if (src_id != 0u) { + GSource* src = g_main_context_find_source_by_id(main_loop_context_, src_id); + g_source_destroy(src); + } + } GST_INFO_OBJECT(pipeline_, "Changing state to %s", gst_element_state_get_name(state)); return gst_element_set_state(pipeline_, state) != GST_STATE_CHANGE_FAILURE; @@ -2713,6 +2752,7 @@ } else if (gst_structure_has_name(structure, kDidReceiveFirstSegmentMsgName)) { if (GST_MESSAGE_SRC(message) == GST_OBJECT(audio_appsrc_) || GST_MESSAGE_SRC(message) == GST_OBJECT(video_appsrc_)) { + GST_INFO("Received '%s' message from %" GST_PTR_FORMAT, kDidReceiveFirstSegmentMsgName, GST_MESSAGE_SRC(message)); bool should_set_rate = false; double rate = 0.; auto type = GST_MESSAGE_SRC(message) == GST_OBJECT(audio_appsrc_) ? MediaType::kAudio : MediaType::kVideo; @@ -2729,6 +2769,23 @@ } } } + else if (gst_structure_has_name(structure, kDidReachBufferingTargetMsgName)) { + if (GST_MESSAGE_SRC(message) == GST_OBJECT(audio_appsrc_) || GST_MESSAGE_SRC(message) == GST_OBJECT(video_appsrc_)) { + int ticket; + if (gst_structure_get_int(structure, "ticket", &ticket)) { + GST_INFO("Received '%s' message from %" GST_PTR_FORMAT, kDidReachBufferingTargetMsgName, GST_MESSAGE_SRC(message)); + auto type = GST_MESSAGE_SRC(message) == GST_OBJECT(audio_appsrc_) ? MediaType::kAudio : MediaType::kVideo; + bool should_update_playing_state = false; + ::starboard::ScopedLock lock(mutex_); + if (ticket == ticket_ && buffering_state_ != 0) { + buffering_state_ &= ~(static_cast<int>(type)); + should_update_playing_state = (buffering_state_ == 0); + } + if (should_update_playing_state) + SchedulePlayingStateUpdate(); + } + } + } } void PlayerImpl::ConfigureLimitedVideo() { @@ -2754,6 +2811,115 @@ audio_codec_ = kSbMediaAudioCodecNone; } +void PlayerImpl::AddBufferingProbe(GstClockTime target, int ticket) { + struct BufferingProbeData { + GstClockTime target_time; + int ticket; + }; + + auto add_probe = [](GstElement* element, GstClockTime target, int ticket, GstPadProbeCallback callback) -> gulong { + BufferingProbeData* data = reinterpret_cast<BufferingProbeData*>(g_malloc0(sizeof (BufferingProbeData))); + data->target_time = target; + data->ticket = ticket; + + GstPad* pad = gst_element_get_static_pad(element, "src"); + GstPadProbeType probe_type = static_cast<GstPadProbeType>(GST_PAD_PROBE_TYPE_BUFFER | GST_PAD_PROBE_TYPE_EVENT_FLUSH); + gulong ret = gst_pad_add_probe (pad, probe_type, callback, data, g_free); + GST_DEBUG_OBJECT(element, + "Buffering probe added to %" GST_PTR_FORMAT ", target time: %" GST_TIME_FORMAT ", ticket: %d", + pad, GST_TIME_ARGS(target), ticket); + gst_object_unref(pad); + + return ret; + }; + + auto buffering_probe_callback = [](GstPad * pad, GstPadProbeInfo * info, gpointer user_data) -> GstPadProbeReturn { + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_EVENT_FLUSH) + return GST_PAD_PROBE_REMOVE; + + SB_CHECK(GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_BUFFER); + + GstBuffer* buffer = gst_pad_probe_info_get_buffer(info); + BufferingProbeData* data = reinterpret_cast<BufferingProbeData*>(user_data); + + GST_TRACE_OBJECT(pad, "Testing buffer: %" GST_PTR_FORMAT " for target time: %" GST_TIME_FORMAT " with ticket: %d", buffer, GST_TIME_ARGS(data->target_time), data->ticket); + + if ( GST_BUFFER_TIMESTAMP(buffer) > data->target_time ) { + + GstObject* parent = gst_pad_get_parent(pad); + if (!parent) { + GST_WARNING_OBJECT(pad, "Pad(%" GST_PTR_FORMAT ") has no parent", pad); + SB_DCHECK(parent); + return GST_PAD_PROBE_REMOVE; + } + + GST_DEBUG_OBJECT(parent, "Did reach target buffering time:%" GST_TIME_FORMAT ", ticket: %d, on pad: %" GST_PTR_FORMAT, GST_TIME_ARGS(data->target_time), data->ticket, pad); + GstStructure* structure = gst_structure_new(kDidReachBufferingTargetMsgName, "ticket", G_TYPE_INT, data->ticket, nullptr); + gst_element_post_message(GST_ELEMENT(parent), gst_message_new_application(parent, structure)); + gst_object_unref(parent); + + return GST_PAD_PROBE_REMOVE; + } + + return GST_PAD_PROBE_OK; + }; + + buffering_state_ = 0; + + if (audio_appsrc_) { + if (add_probe(audio_appsrc_, target, ticket, buffering_probe_callback) != 0u) + buffering_state_ |= static_cast<int>(MediaType::kAudio); + } + + if (video_appsrc_) { + if (SbDrmSystemIsValid(drm_system_)) { + target += 5 * 16 * GST_MSECOND; + } + if (add_probe(video_appsrc_, target, ticket, buffering_probe_callback) != 0u) + buffering_state_ |= static_cast<int>(MediaType::kVideo); + } +} + +void PlayerImpl::SchedulePlayingStateUpdate() { + mutex_.DCheckAcquired(); + + if (playing_state_update_source_id_ != 0u) // already scheduled + return; + + const auto update_callback = [](gpointer data) -> gboolean { + PlayerImpl* self = static_cast<PlayerImpl*>(data); + bool should_be_playing = false, can_play = false; + + { + GstState state, pending; + GstStateChangeReturn ret; + bool need_preroll; + guint src_id; + + ret = gst_element_get_state(self->pipeline_, &state, &pending, 0); + need_preroll = (state == GST_STATE_PAUSED && pending == GST_STATE_PAUSED && ret == GST_STATE_CHANGE_ASYNC); + + ::starboard::ScopedLock lock(self->mutex_); + src_id = std::exchange(self->playing_state_update_source_id_, 0u); + should_be_playing = (self->rate_ || self->pending_rate_); + can_play = (self->buffering_state_ == 0 && !need_preroll && src_id != 0u && self->buf_target_min_ts_ == kSbTimeMax); // the update was canceled when src_id == 0 + } + + GST_DEBUG_OBJECT(self->pipeline_, "Update pipeline state, should be playing: %d, can_play: %d", should_be_playing, can_play); + + if (should_be_playing && can_play) + self->ChangePipelineState(GST_STATE_PLAYING); + + return G_SOURCE_REMOVE; + }; + + GSource* src = g_idle_source_new(); + g_source_set_priority (src, G_PRIORITY_DEFAULT); + g_source_set_callback(src, update_callback, this, nullptr); + playing_state_update_source_id_ = g_source_attach(src, main_loop_context_); + g_source_unref(src); +} + } // namespace void ForceStop() {