diff --git a/Source/WebCore/platform/graphics/gstreamer/MediaPlayerPrivateGStreamerMSE.cpp b/Source/WebCore/platform/graphics/gstreamer/MediaPlayerPrivateGStreamerMSE.cpp index 51014bce808fa..2c89c8f658111 100644 --- a/Source/WebCore/platform/graphics/gstreamer/MediaPlayerPrivateGStreamerMSE.cpp +++ b/Source/WebCore/platform/graphics/gstreamer/MediaPlayerPrivateGStreamerMSE.cpp @@ -77,17 +77,25 @@ using namespace std; namespace WebCore { +struct PadProbeInformation +{ + AppendPipeline* m_appendPipeline; + const char* m_description; + gulong m_probeId; +}; + class AppendPipeline : public ThreadSafeRefCounted { public: enum AppendStage { Invalid, NotStarted, Ongoing, KeyNegotiation, DataStarve, Sampling, LastSample, Aborting }; - static const unsigned int s_dataStarvedTimeoutMsec = 1000; + static const unsigned int s_dataStarvedTimeoutMsec = 2000; static const unsigned int s_lastSampleTimeoutMsec = 250; AppendPipeline(PassRefPtr mediaSourceClient, PassRefPtr sourceBufferPrivate, MediaPlayerPrivateGStreamerMSE* playerPrivate); virtual ~AppendPipeline(); void handleElementMessage(GstMessage*); + void handleApplicationMessage(GstMessage*); gint id(); AppendStage appendStage() { return m_appendStage; } @@ -125,8 +133,14 @@ class AppendPipeline : public ThreadSafeRefCounted { void scheduleLastSampleTimer(); void cancelLastSampleTimer(); + void reportEndOfAppendDataMarkNeeded(); + void reportEndOfAppendDataMarkReceived(guint id); + private: void resetPipeline(); + void checkEndOfAppendDataMarkReceived(); + void handleEndOfAppendDataMarkNeeded(); + void handleEndOfAppendDataMarkReceived(const GstStructure*); // TODO: Hide everything and use getters/setters. private: @@ -142,8 +156,8 @@ class AppendPipeline : public ThreadSafeRefCounted { GstFlowReturn m_flowReturn; GstElement* m_pipeline; + GRefPtr m_bus; GstElement* m_appsrc; - GstElement* m_typefind; GstElement* m_qtdemux; GstElement* m_decryptor; @@ -160,6 +174,22 @@ class AppendPipeline : public ThreadSafeRefCounted { GstCaps* m_demuxerSrcPadCaps; FloatSize m_presentationSize; + // Unique id of the current append operation. Used to mark + // custom events, detect them in the sink and trigger lastSampleTimeout + // ahead of time. + + // This is the last id marked right after appending to appsrc + guint m_appendIdMarkedInSrc; + + // This is the last id received by the probe in the appsink sink pad + guint m_appendIdReceivedInSink; + + gulong m_appsinkDataEnteringProbeId; + gulong m_appsrcDataLeavingProbeId; +#ifdef DEBUG_APPEND_PIPELINE_PADS + struct PadProbeInformation m_demuxerDataEnteringPadProbeInformation; +#endif + // Some appended data are only headers and don't generate any // useful stream data for decoding. This is detected with a // timeout and reported to the upper layers, so update/updateend @@ -1190,18 +1220,28 @@ static void appendPipelineDemuxerPadRemoved(GstElement*, GstPad*, AppendPipeline static gboolean appendPipelineDemuxerConnectToAppSinkMainThread(PadInfo*); static gboolean appendPipelineDemuxerDisconnectFromAppSinkMainThread(PadInfo*); static void appendPipelineAppSinkCapsChanged(GObject*, GParamSpec*, AppendPipeline*); +static GstPadProbeReturn appendPipelineAppsinkDataEntering(GstPad*, GstPadProbeInfo*, AppendPipeline*); +static GstPadProbeReturn appendPipelineAppsrcDataLeaving(GstPad*, GstPadProbeInfo*, AppendPipeline*); +#ifdef DEBUG_APPEND_PIPELINE_PADS +static GstPadProbeReturn appendPipelinePadProbeDebugInformation(GstPad*, GstPadProbeInfo*, struct PadProbeInformation*); +#endif static GstFlowReturn appendPipelineAppSinkNewSample(GstElement*, AppendPipeline*); static gboolean appendPipelineAppSinkNewSampleMainThread(NewSampleInfo*); static void appendPipelineAppSinkEOS(GstElement*, AppendPipeline*); static gboolean appendPipelineAppSinkEOSMainThread(AppendPipeline* ap); -static gboolean appendPipelineDataStarveTimeout(AppendPipeline* ap); -static gboolean appendPipelineLastSampleTimeout(AppendPipeline* ap); +static gboolean appendPipelineDataStarveTimeout(AppendPipeline*); +static gboolean appendPipelineLastSampleTimeout(gpointer); static void appendPipelineElementMessageCallback(GstBus*, GstMessage* message, AppendPipeline* ap) { ap->handleElementMessage(message); } +static void appendPipelineApplicationMessageCallback(GstBus*, GstMessage* message, AppendPipeline* appendPipeline) +{ + appendPipeline->handleApplicationMessage(message); +} + AppendPipeline::AppendPipeline(PassRefPtr mediaSourceClient, PassRefPtr sourceBufferPrivate, MediaPlayerPrivateGStreamerMSE* playerPrivate) : m_mediaSourceClient(mediaSourceClient) , m_sourceBufferPrivate(sourceBufferPrivate) @@ -1209,6 +1249,8 @@ AppendPipeline::AppendPipeline(PassRefPtr mediaSo , m_id(0) , m_appSinkCaps(NULL) , m_demuxerSrcPadCaps(NULL) + , m_appendIdMarkedInSrc(0) + , m_appendIdReceivedInSink(0) , m_dataStarvedTimeoutTag(0) , m_lastSampleTimeoutTag(0) , m_appendStage(NotStarted) @@ -1223,11 +1265,12 @@ AppendPipeline::AppendPipeline(PassRefPtr mediaSo // The track name is still unknown at this time, though. m_pipeline = gst_pipeline_new(NULL); - GRefPtr bus = adoptGRef(gst_pipeline_get_bus(GST_PIPELINE(m_pipeline))); - gst_bus_add_signal_watch(bus.get()); - gst_bus_enable_sync_message_emission(bus.get()); + m_bus = adoptGRef(gst_pipeline_get_bus(GST_PIPELINE(m_pipeline))); + gst_bus_add_signal_watch(m_bus.get()); + gst_bus_enable_sync_message_emission(m_bus.get()); - g_signal_connect(bus.get(), "sync-message::element", G_CALLBACK(appendPipelineElementMessageCallback), this); + g_signal_connect(m_bus.get(), "sync-message::element", G_CALLBACK(appendPipelineElementMessageCallback), this); + g_signal_connect(m_bus.get(), "message::application", G_CALLBACK(appendPipelineApplicationMessageCallback), this); g_mutex_init(&m_newSampleMutex); g_cond_init(&m_newSampleCondition); @@ -1237,7 +1280,6 @@ AppendPipeline::AppendPipeline(PassRefPtr mediaSo m_decryptor = NULL; m_appsrc = gst_element_factory_make("appsrc", NULL); - m_typefind = gst_element_factory_make("typefind", NULL); m_qtdemux = gst_element_factory_make("qtdemux", NULL); { GValue val = G_VALUE_INIT; @@ -1253,6 +1295,26 @@ AppendPipeline::AppendPipeline(PassRefPtr mediaSo GRefPtr appSinkPad = adoptGRef(gst_element_get_static_pad(m_appsink, "sink")); g_signal_connect(appSinkPad.get(), "notify::caps", G_CALLBACK(appendPipelineAppSinkCapsChanged), this); +#ifdef DEBUG_APPEND_PIPELINE_PADS + m_appsinkDataEnteringProbeId = gst_pad_add_probe(appSinkPad.get(), static_cast(GST_PAD_PROBE_TYPE_BUFFER | GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM), reinterpret_cast(appendPipelineAppsinkDataEntering), this, nullptr); +#else + m_appsinkDataEnteringProbeId = gst_pad_add_probe(appSinkPad.get(), GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM, reinterpret_cast(appendPipelineAppsinkDataEntering), this, nullptr); +#endif + + GRefPtr appsrcPad = adoptGRef(gst_element_get_static_pad(m_appsrc, "src")); +#ifdef DEBUG_APPEND_PIPELINE_PADS + m_appsrcDataLeavingProbeId = gst_pad_add_probe(appsrcPad.get(), static_cast(GST_PAD_PROBE_TYPE_BUFFER | GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM), reinterpret_cast(appendPipelineAppsrcDataLeaving), this, nullptr); +#else + m_appsrcDataLeavingProbeId = gst_pad_add_probe(appsrcPad.get(), GST_PAD_PROBE_TYPE_BUFFER, reinterpret_cast(appendPipelineAppsrcDataLeaving), this, nullptr); +#endif + +#ifdef DEBUG_APPEND_PIPELINE_PADS + GRefPtr demuxerPad = adoptGRef(gst_element_get_static_pad(m_qtdemux, "sink")); + m_demuxerDataEnteringPadProbeInformation.m_appendPipeline = this; + m_demuxerDataEnteringPadProbeInformation.m_description = "demuxer data entering"; + m_demuxerDataEnteringPadProbeInformation.m_probeId = gst_pad_add_probe(demuxerPad.get(), static_cast(GST_PAD_PROBE_TYPE_BUFFER | GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM), reinterpret_cast(appendPipelinePadProbeDebugInformation), &m_demuxerDataEnteringPadProbeInformation, nullptr); +#endif + // These signals won't be connected outside of the lifetime of "this". g_signal_connect(m_qtdemux, "pad-added", G_CALLBACK(appendPipelineDemuxerPadAdded), this); g_signal_connect(m_qtdemux, "pad-removed", G_CALLBACK(appendPipelineDemuxerPadRemoved), this); @@ -1261,12 +1323,11 @@ AppendPipeline::AppendPipeline(PassRefPtr mediaSo // Add_many will take ownership of a reference. Request one ref more for ourselves. gst_object_ref(m_appsrc); - gst_object_ref(m_typefind); gst_object_ref(m_qtdemux); gst_object_ref(m_appsink); - gst_bin_add_many(GST_BIN(m_pipeline), m_appsrc, m_typefind, m_qtdemux, NULL); - gst_element_link_many(m_appsrc, m_typefind, m_qtdemux, NULL); + gst_bin_add_many(GST_BIN(m_pipeline), m_appsrc, m_qtdemux, NULL); + gst_element_link(m_appsrc, m_qtdemux); gst_element_set_state(m_pipeline, GST_STATE_READY); }; @@ -1292,10 +1353,9 @@ AppendPipeline::~AppendPipeline() cancelLastSampleTimer(); if (m_pipeline) { - GRefPtr bus = adoptGRef(gst_pipeline_get_bus(GST_PIPELINE(m_pipeline))); - ASSERT(bus); - g_signal_handlers_disconnect_by_func(bus.get(), reinterpret_cast(appendPipelineElementMessageCallback), this); - gst_bus_disable_sync_message_emission(bus.get()); + ASSERT(m_bus); + g_signal_handlers_disconnect_by_func(m_bus.get(), reinterpret_cast(appendPipelineElementMessageCallback), this); + gst_bus_disable_sync_message_emission(m_bus.get()); gst_element_set_state (m_pipeline, GST_STATE_NULL); gst_object_unref(m_pipeline); @@ -1303,16 +1363,18 @@ AppendPipeline::~AppendPipeline() } if (m_appsrc) { + GRefPtr appsrcPad = adoptGRef(gst_element_get_static_pad(m_appsrc, "src")); + gst_pad_remove_probe(appsrcPad.get(), m_appsrcDataLeavingProbeId); gst_object_unref(m_appsrc); m_appsrc = NULL; } - if (m_typefind) { - gst_object_unref(m_typefind); - m_typefind = NULL; - } - if (m_qtdemux) { +#ifdef DEBUG_APPEND_PIPELINE_PADS + GRefPtr demuxerPad = adoptGRef(gst_element_get_static_pad(m_qtdemux, "sink")); + gst_pad_remove_probe(demuxerPad.get(), m_demuxerDataEnteringPadProbeInformation.m_probeId); +#endif + g_signal_handlers_disconnect_by_func(m_qtdemux, (gpointer)appendPipelineDemuxerPadAdded, this); g_signal_handlers_disconnect_by_func(m_qtdemux, (gpointer)appendPipelineDemuxerPadRemoved, this); @@ -1330,6 +1392,8 @@ AppendPipeline::~AppendPipeline() g_signal_handlers_disconnect_by_func(m_appsink, (gpointer)appendPipelineAppSinkNewSample, this); g_signal_handlers_disconnect_by_func(m_appsink, (gpointer)appendPipelineAppSinkEOS, this); + gst_pad_remove_probe(appSinkPad.get(), m_appsinkDataEnteringProbeId); + gst_object_unref(m_appsink); m_appsink = NULL; } @@ -1389,6 +1453,35 @@ void AppendPipeline::handleElementMessage(GstMessage* message) m_playerPrivate->handleSyncMessage(message); } +void AppendPipeline::handleApplicationMessage(GstMessage* message) +{ + ASSERT(WTF::isMainThread()); + + const GstStructure* structure = gst_message_get_structure(message); + + if (gst_structure_has_name(structure, "end-of-append-data-mark-received")) { + handleEndOfAppendDataMarkReceived(structure); + return; + } + + if (gst_structure_has_name(structure, "end-of-append-data-mark-needed")) { + handleEndOfAppendDataMarkNeeded(); + return; + } + + ASSERT_NOT_REACHED(); +} + +void AppendPipeline::handleEndOfAppendDataMarkReceived(const GstStructure* structure) +{ + gst_structure_get(structure, "id", G_TYPE_UINT, &m_appendIdReceivedInSink, NULL); + ASSERT(m_appendIdReceivedInSink); + + TRACE_MEDIA_MESSAGE("received end of append id %u in the sink", m_appendIdReceivedInSink); + if (m_appendStage == Sampling || m_appendStage == Ongoing) + checkEndOfAppendDataMarkReceived(); +} + gint AppendPipeline::id() { ASSERT(WTF::isMainThread()); @@ -1446,7 +1539,7 @@ void AppendPipeline::scheduleLastSampleTimer() { if (m_lastSampleTimeoutTag) cancelLastSampleTimer(); - m_lastSampleTimeoutTag = g_timeout_add(s_lastSampleTimeoutMsec, GSourceFunc(appendPipelineLastSampleTimeout), this); + m_lastSampleTimeoutTag = g_timeout_add(s_lastSampleTimeoutMsec, appendPipelineLastSampleTimeout, this); } void AppendPipeline::cancelLastSampleTimer() @@ -1489,6 +1582,7 @@ void AppendPipeline::setAppendStage(AppendStage newAppendStage) case NotStarted: ok = true; if (m_pendingBuffer) { + TRACE_MEDIA_MESSAGE("pushing pending buffer %p", m_pendingBuffer.get()); gst_app_src_push_buffer(GST_APP_SRC(appsrc()), m_pendingBuffer.leakRef()); nextAppendStage = Ongoing; } @@ -1753,6 +1847,32 @@ void AppendPipeline::appSinkCapsChanged() gst_caps_unref(caps); } +void AppendPipeline::checkEndOfAppendDataMarkReceived() +{ + ASSERT(WTF::isMainThread()); + + if (!m_appendIdReceivedInSink || m_appendIdMarkedInSrc != m_appendIdReceivedInSink) + return; + + TRACE_MEDIA_MESSAGE("end of append data mark was received"); + + switch (m_appendStage) { + case Ongoing: + TRACE_MEDIA_MESSAGE("DataStarve"); + m_appendIdReceivedInSink = 0; + setAppendStage(DataStarve); + break; + case Sampling: + TRACE_MEDIA_MESSAGE("LastSample"); + m_appendIdReceivedInSink = 0; + setAppendStage(LastSample); + break; + default: + ERROR_MEDIA_MESSAGE("Unexpected"); + break; + } +} + void AppendPipeline::appSinkNewSample(GstSample* sample) { ASSERT(WTF::isMainThread()); @@ -1799,6 +1919,8 @@ void AppendPipeline::appSinkNewSample(GstSample* sample) m_flowReturn = GST_FLOW_OK; g_cond_signal(&m_newSampleCondition); g_mutex_unlock(&m_newSampleMutex); + + checkEndOfAppendDataMarkReceived(); } void AppendPipeline::appSinkEOS() @@ -1904,12 +2026,49 @@ void AppendPipeline::abort() GstFlowReturn AppendPipeline::pushNewBuffer(GstBuffer* buffer) { + GstFlowReturn result; + if (m_abortPending) { m_pendingBuffer = adoptGRef(buffer); - return GST_FLOW_OK; + result = GST_FLOW_OK; + } else { + setAppendStage(AppendPipeline::Ongoing); + TRACE_MEDIA_MESSAGE("pushing new buffer %p", buffer); + result = gst_app_src_push_buffer(GST_APP_SRC(appsrc()), buffer); } - setAppendStage(AppendPipeline::Ongoing); - return gst_app_src_push_buffer(GST_APP_SRC(appsrc()), buffer); + + return result; +} + +void AppendPipeline::handleEndOfAppendDataMarkNeeded() +{ + GstEvent* event = gst_event_new_custom(GST_EVENT_CUSTOM_DOWNSTREAM, gst_structure_new_empty("end-of-append-data-mark")); + m_appendIdMarkedInSrc = gst_event_get_seqnum(event); + m_appendIdReceivedInSink = 0; + + TRACE_MEDIA_MESSAGE("marking end of append with id %u", m_appendIdMarkedInSrc); + + gst_element_send_event(m_appsrc, event); + + GstBuffer* emptyBuffer = gst_buffer_new_and_alloc(0); + gst_buffer_fill(emptyBuffer, 0, nullptr, 0); + gst_app_src_push_buffer(GST_APP_SRC(appsrc()), emptyBuffer); +} + +void AppendPipeline::reportEndOfAppendDataMarkReceived(guint id) +{ + GstStructure* structure = gst_structure_new("end-of-append-data-mark-received", "id", G_TYPE_UINT, id, NULL); + GstMessage* message = gst_message_new_application(GST_OBJECT(m_appsink), structure); + gst_bus_post(m_bus.get(), message); + TRACE_MEDIA_MESSAGE("received message with id %u, re-posted to bus", id); +} + +void AppendPipeline::reportEndOfAppendDataMarkNeeded() +{ + GstStructure* structure = gst_structure_new_empty("end-of-append-data-mark-needed"); + GstMessage* message = gst_message_new_application(GST_OBJECT(m_appsrc), structure); + gst_bus_post(m_bus.get(), message); + TRACE_MEDIA_MESSAGE("received buffer going thru, re-posted to bus"); } GstFlowReturn AppendPipeline::handleNewSample(GstElement* appsink) @@ -2066,6 +2225,10 @@ void AppendPipeline::connectToAppSink(GstPad* demuxerSrcPad) break; } + // The previous mark has probably been lost because appsink was disconnected. Mark again. + TRACE_MEDIA_MESSAGE("previous append end mark lost, reinjecting"); + handleEndOfAppendDataMarkNeeded(); + g_cond_signal(&m_padAddRemoveCondition); } @@ -2129,6 +2292,103 @@ static void appendPipelineAppSinkCapsChanged(GObject*, GParamSpec*, AppendPipeli g_timeout_add(0, appSinkCapsChangedFromMainThread, ap); } +static GstPadProbeReturn appendPipelineAppsrcDataLeaving(GstPad*, GstPadProbeInfo* info, AppendPipeline* appendPipeline) +{ + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_BUFFER) { + GstBuffer* buffer = GST_PAD_PROBE_INFO_BUFFER(info); + gsize bufferSize = gst_buffer_get_size(buffer); + + TRACE_MEDIA_MESSAGE("buffer of size %" G_GSIZE_FORMAT " going thru", bufferSize); + + if (bufferSize > 0) + appendPipeline->reportEndOfAppendDataMarkNeeded(); + + return GST_PAD_PROBE_OK; + } + +#ifdef DEBUG_APPEND_PIPELINE_PADS + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM) { + GstEvent* event = GST_PAD_PROBE_INFO_EVENT(info); + if (GST_EVENT_TYPE(event) != GST_EVENT_CUSTOM_DOWNSTREAM) + return GST_PAD_PROBE_OK; + + const GstStructure* structure = gst_event_get_structure(event); + if (!gst_structure_has_name(structure, "end-of-append-data-mark")) + return GST_PAD_PROBE_OK; + + guint id = gst_event_get_seqnum(event); + TRACE_MEDIA_MESSAGE("custom downstream event id=%u", id); + + return GST_PAD_PROBE_OK; + } +#endif + + ASSERT_NOT_REACHED(); + return GST_PAD_PROBE_OK; +} + +static GstPadProbeReturn appendPipelineAppsinkDataEntering(GstPad*, GstPadProbeInfo* info, AppendPipeline* appendPipeline) +{ + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM) { + GstEvent* event = GST_PAD_PROBE_INFO_EVENT(info); + if (GST_EVENT_TYPE(event) != GST_EVENT_CUSTOM_DOWNSTREAM) + return GST_PAD_PROBE_OK; + + const GstStructure* structure = gst_event_get_structure(event); + if (!gst_structure_has_name(structure, "end-of-append-data-mark")) + return GST_PAD_PROBE_OK; + + guint id = gst_event_get_seqnum(event); + + TRACE_MEDIA_MESSAGE("id=%u", id); + + appendPipeline->reportEndOfAppendDataMarkReceived(id); + + return GST_PAD_PROBE_OK; + } + +#ifdef DEBUG_APPEND_PIPELINE_PADS + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_BUFFER) { + GstBuffer* buffer = GST_PAD_PROBE_INFO_BUFFER(info); + TRACE_MEDIA_MESSAGE("buffer of size %" G_GSIZE_FORMAT " going thru", gst_buffer_get_size(buffer)); + return GST_PAD_PROBE_OK; + } +#endif + + ASSERT_NOT_REACHED(); + return GST_PAD_PROBE_OK; +} + +#ifdef DEBUG_APPEND_PIPELINE_PADS +static GstPadProbeReturn appendPipelinePadProbeDebugInformation(GstPad*, GstPadProbeInfo* info, struct PadProbeInformation* padProbeInformation) +{ + ASSERT(GST_PAD_PROBE_INFO_TYPE(info) != static_cast(GST_PAD_PROBE_TYPE_BUFFER | GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM)); + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_BUFFER) { + GstBuffer* buffer = GST_PAD_PROBE_INFO_BUFFER(info); + TRACE_MEDIA_MESSAGE("%s: buffer of size %" G_GSIZE_FORMAT " going thru", padProbeInformation->m_description, gst_buffer_get_size(buffer)); + return GST_PAD_PROBE_OK; + } + + if (GST_PAD_PROBE_INFO_TYPE(info) & GST_PAD_PROBE_TYPE_EVENT_DOWNSTREAM) { + GstEvent* event = GST_PAD_PROBE_INFO_EVENT(info); + if (GST_EVENT_TYPE(event) != GST_EVENT_CUSTOM_DOWNSTREAM) + return GST_PAD_PROBE_OK; + + const GstStructure* structure = gst_event_get_structure(event); + if (!gst_structure_has_name(structure, "end-of-append-data-mark")) + return GST_PAD_PROBE_OK; + + guint id = gst_event_get_seqnum(event); + TRACE_MEDIA_MESSAGE("%s: custom downstream event id=%u", padProbeInformation->m_description, id); + + return GST_PAD_PROBE_OK; + } + + ASSERT_NOT_REACHED(); + return GST_PAD_PROBE_OK; +} +#endif + static void appendPipelineDemuxerPadAdded(GstElement*, GstPad* demuxerSrcPad, AppendPipeline* ap) { ap->connectToAppSinkFromAnyThread(demuxerSrcPad); @@ -2185,23 +2445,21 @@ static gboolean appendPipelineAppSinkEOSMainThread(AppendPipeline* ap) return G_SOURCE_REMOVE; } -static gboolean appendPipelineDataStarveTimeout(AppendPipeline* ap) +static gboolean appendPipelineDataStarveTimeout(AppendPipeline* appendPipeline) { - LOG_MEDIA_MESSAGE("data starve timer fired"); - if (ap->appendStage()==AppendPipeline::AppendStage::Invalid) - return G_SOURCE_REMOVE; - - ap->setAppendStage(AppendPipeline::DataStarve); + AppendPipeline::AppendStage appendStage = appendPipeline->appendStage(); + INFO_MEDIA_MESSAGE("data starve timer fired, stage %s", dumpAppendStage(appendStage)); + if (appendStage == AppendPipeline::AppendStage::Ongoing) { + WARN_MEDIA_MESSAGE("setting DataStarve because of timeout"); + appendPipeline->setAppendStage(AppendPipeline::AppendStage::DataStarve); + } return G_SOURCE_REMOVE; } -static gboolean appendPipelineLastSampleTimeout(AppendPipeline* ap) +static gboolean appendPipelineLastSampleTimeout(gpointer) { - TRACE_MEDIA_MESSAGE("last sample timer fired"); - if (ap->appendStage()==AppendPipeline::AppendStage::Invalid) - return G_SOURCE_REMOVE; - - ap->setAppendStage(AppendPipeline::LastSample); + ERROR_MEDIA_MESSAGE("last sample timer fired"); + ASSERT_NOT_REACHED(); return G_SOURCE_REMOVE; }