diff --git a/middleware/InterfacePlayerRDK.cpp b/middleware/InterfacePlayerRDK.cpp index bc36fa268b..abd318e4f2 100644 --- a/middleware/InterfacePlayerRDK.cpp +++ b/middleware/InterfacePlayerRDK.cpp @@ -1371,6 +1371,18 @@ void InterfacePlayerRDK::TearDownStream(int type) else if (mediaType == eGST_MEDIATYPE_SUBTITLE) { g_clear_object(&interfacePlayerPriv->gstPrivateContext->subtitle_sink); + pthread_mutex_lock(&stream->sourceLock); + if (stream->sinkbin) + { + MW_LOG_WARN("InterfacePlayerRDK::TearDownStream: CC sinkbin still assigned, clearing"); + g_clear_object(&stream->sinkbin); + } + if (stream->source) + { + MW_LOG_WARN("InterfacePlayerRDK::TearDownStream: CC source still assigned, clearing"); + g_clear_object(&stream->source); + } + pthread_mutex_unlock(&stream->sourceLock); } tearDownCb(false, mediaType); MW_LOG_MIL("InterfacePlayerRDK::TearDownStream: exit mediaType = %d", mediaType); @@ -2258,7 +2270,7 @@ int InterfacePlayerRDK::SetupStream(int streamId, void *playerInstance, std::st gst_element_add_pad(subtitlebin, gst_ghost_pad_new("sink", gst_element_get_static_pad(vipertransform, "sink"))); g_object_set(stream->sinkbin, "text-sink", subtitlebin, NULL); - interfacePlayerPriv->gstPrivateContext->subtitle_sink = textsink; + interfacePlayerPriv->gstPrivateContext->subtitle_sink = GST_ELEMENT(gst_object_ref(textsink)); MW_LOG_MIL("using rialtomsesubtitlesink muted=%d sink=%p", interfacePlayerPriv->gstPrivateContext->subtitleMuted, interfacePlayerPriv->gstPrivateContext->subtitle_sink); g_object_set(textsink, "mute", interfacePlayerPriv->gstPrivateContext->subtitleMuted ? TRUE : FALSE, NULL); } @@ -2327,7 +2339,7 @@ int InterfacePlayerRDK::SetupStream(int streamId, void *playerInstance, std::st MW_LOG_INFO("setting has-drm=false for clear HLS/TS playback"); g_object_set(vidsink, "has-drm", FALSE, NULL); } - interfacePlayerPriv->gstPrivateContext->video_sink = vidsink; + interfacePlayerPriv->gstPrivateContext->video_sink = GST_ELEMENT(gst_object_ref(vidsink)); // RDKEMW-18286: Set show-video-window=FALSE at sink creation time. // This is the EARLIEST possible point. The Rialto delegate will queue @@ -2357,7 +2369,7 @@ int InterfacePlayerRDK::SetupStream(int streamId, void *playerInstance, std::st { MW_LOG_INFO("Created rialtomseaudiosink : %s",GST_ELEMENT_NAME(audSink)); g_object_set(stream->sinkbin, "audio-sink", audSink, NULL); - interfacePlayerPriv->gstPrivateContext->audio_sink = audSink; + interfacePlayerPriv->gstPrivateContext->audio_sink = GST_ELEMENT(gst_object_ref(audSink)); } else { @@ -3309,16 +3321,30 @@ void InterfacePlayerPriv::SendNewSegmentEvent(int type, GstClockTime startPts ,G if (gstPrivateContext->usingRialtoSink) { GstCaps *currentCaps = gst_app_src_get_caps(GST_APP_SRC(stream->source)); - GstSample *sample = gst_sample_new (nullptr, currentCaps, &segment, nullptr); - - MW_LOG_INFO("Pushing sample with segment for mediaType[%d]. start %" G_GUINT64_FORMAT " stop %" G_GUINT64_FORMAT" rate %f applied_rate %f", mediaType, segment.start, segment.stop, segment.rate, segment.applied_rate); - if (GST_FLOW_OK != gst_app_src_push_sample(GST_APP_SRC(stream->source), sample)) + if (currentCaps != NULL) { - MW_LOG_ERR("Failed to push sample with segment for mediaType[%d]", mediaType); + GstSample *sample = gst_sample_new (nullptr, currentCaps, &segment, nullptr); + if (sample != NULL) + { + MW_LOG_INFO("Pushing sample with segment for mediaType[%d]. start %" G_GUINT64_FORMAT " stop %" G_GUINT64_FORMAT" rate %f applied_rate %f", mediaType, segment.start, segment.stop, segment.rate, segment.applied_rate); + if (GST_FLOW_OK != gst_app_src_push_sample(GST_APP_SRC(stream->source), sample)) + { + MW_LOG_ERR("Failed to push sample with segment for mediaType[%d]", mediaType); + } + gst_sample_unref(sample); + } + else + { + MW_LOG_ERR("Failed to create sample for mediaType[%d]", mediaType); + } + gst_caps_unref(currentCaps); + } + else + { + MW_LOG_WARN("Cannot push segment for mediaType[%d] - caps not yet set on appsrc", mediaType); } - gst_sample_unref(sample); - gst_caps_unref(currentCaps); } + else { MW_LOG_INFO("Sending segment event for mediaType[%d]. start %" G_GUINT64_FORMAT " stop %" G_GUINT64_FORMAT" rate %f applied_rate %f", mediaType, segment.start, segment.stop, segment.rate, segment.applied_rate); diff --git a/middleware/drm/DrmSession.cpp b/middleware/drm/DrmSession.cpp index 34087df3fb..a735cb6f59 100644 --- a/middleware/drm/DrmSession.cpp +++ b/middleware/drm/DrmSession.cpp @@ -24,12 +24,17 @@ #include "DrmSession.h" #include "PlayerLogManager.h" +#include /** * @brief Constructor for DrmSession. */ DrmSession::DrmSession(const string &keySystem) : m_keySystem(keySystem),m_OutputProtectionEnabled(false) , mContentSecurityManagerSession() + , mLifecycleMutex() + , mLifecycleCV() + , mActiveOperations(0) + , mMarkedForDestruction(false) { } @@ -40,6 +45,55 @@ DrmSession::~DrmSession() { } +/** + * @brief DELIA-70726 fix: Acquire lifecycle guard before use in decrypt(). + */ +bool DrmSession::AcquireForUse() +{ + std::lock_guard lock(mLifecycleMutex); + if (mMarkedForDestruction) + { + return false; + } + mActiveOperations++; + return true; +} + +/** + * @brief DELIA-70726 fix: Release lifecycle guard after use in decrypt(). + */ +void DrmSession::ReleaseAfterUse() +{ + std::lock_guard lock(mLifecycleMutex); + if (mActiveOperations > 0) + { + mActiveOperations--; + } + if (mActiveOperations == 0) + { + mLifecycleCV.notify_all(); + } +} + +/** + * @brief DELIA-70726 fix: Block deletion of this session until any decrypt() + * call already in progress (having acquired the guard) has finished. + */ +void DrmSession::PrepareForDestruction(uint32_t timeoutMs) +{ + std::unique_lock lock(mLifecycleMutex); + mMarkedForDestruction = true; + if (mActiveOperations > 0) + { + MW_LOG_WARN("DrmSession::PrepareForDestruction : waiting for %d in-flight decrypt operation(s) to complete before delete", mActiveOperations); + mLifecycleCV.wait_for(lock, std::chrono::milliseconds(timeoutMs), [this]() { return mActiveOperations == 0; }); + if (mActiveOperations > 0) + { + MW_LOG_ERR("DrmSession::PrepareForDestruction : timed out waiting for in-flight decrypt operation(s); proceeding with destruction"); + } + } +} + /** * @brief Get the DRM System, ie, UUID for PlayReady WideVine etc.. */ diff --git a/middleware/drm/DrmSession.h b/middleware/drm/DrmSession.h index 4ff9e78497..8c9b00075d 100755 --- a/middleware/drm/DrmSession.h +++ b/middleware/drm/DrmSession.h @@ -29,6 +29,9 @@ #include #include #include +#include +#include +#include #include "DrmUtils.h" #include "ContentSecurityManagerSession.h" @@ -66,7 +69,53 @@ class DrmSession std::string m_keySystem; bool m_OutputProtectionEnabled; ContentSecurityManagerSession mContentSecurityManagerSession; + + /* DELIA-70726 fix: + * Lifecycle guard used to prevent the DrmSession object from being + * deleted (e.g. by DrmSessionManager during DRM session slot + * reuse/eviction on back-to-back channel changes) while a GStreamer + * pipeline thread (multiqueue/decryptor) is concurrently inside + * decrypt()/verifyOutputProtection(). Without this guard, deletion of + * a session that is still referenced by an old, not-yet-fully-torn-down + * pipeline results in a use-after-free SIGSEGV inside + * OCDMSessionAdapter::verifyOutputProtection(). + */ + std::mutex mLifecycleMutex; + std::condition_variable mLifecycleCV; + int mActiveOperations; + bool mMarkedForDestruction; + public: + /** + * @fn AcquireForUse + * @brief Must be called by any external caller (e.g. the GStreamer + * decryptor element) before invoking decrypt() on a DrmSession + * obtained via a raw/cached pointer. Returns false if the + * session is already being torn down, in which case decrypt() + * MUST NOT be called on this object. + * @retval true if it is safe to call decrypt(), false otherwise. + */ + bool AcquireForUse(); + + /** + * @fn ReleaseAfterUse + * @brief Must be called exactly once for every successful AcquireForUse(), + * after the decrypt() call completes. + */ + void ReleaseAfterUse(); + + /** + * @fn PrepareForDestruction + * @brief Must be called by the owner (DrmSessionManager) before deleting + * this DrmSession. Marks the session so that any new + * AcquireForUse() calls fail fast, and blocks (bounded) until all + * in-flight decrypt() operations that already acquired the guard + * have completed, making it safe to free the object. + * @param timeoutMs maximum time to wait for in-flight operations to drain. + */ + void PrepareForDestruction(uint32_t timeoutMs = 3000); + + /** * @brief Create drm session with given init data * @param f_pbInitData : pointer to initdata diff --git a/middleware/drm/DrmSessionManager.cpp b/middleware/drm/DrmSessionManager.cpp index 61c8f404be..d0ca663676 100755 --- a/middleware/drm/DrmSessionManager.cpp +++ b/middleware/drm/DrmSessionManager.cpp @@ -106,6 +106,10 @@ void DrmSessionManager::clearSessionData() { if (drmSessionContexts != NULL && drmSessionContexts[i].drmSession != NULL) { + /* DELIA-70726 fix: block until any in-flight decrypt() on this session + * (called from a GStreamer pipeline thread via a cached raw pointer) + * has completed, before freeing the object. */ + drmSessionContexts[i].drmSession->PrepareForDestruction(); MW_SAFE_DELETE(drmSessionContexts[i].drmSession); drmSessionContexts[i] = DrmSessionContext(); } @@ -198,6 +202,8 @@ void DrmSessionManager::clearDrmSession(bool forceClearSession) if (drmSessionContexts[i].drmSession != NULL) { MW_LOG_WARN("DrmSessionManager:: Clearing failed Session Data Slot : %d", i); + /* DELIA-70726 fix: see clearSessionData() for rationale. */ + drmSessionContexts[i].drmSession->PrepareForDestruction(); MW_SAFE_DELETE(drmSessionContexts[i].drmSession); } } @@ -703,6 +709,13 @@ KeyState DrmSessionManager::getDrmSession(int &err, std::shared_ptr d MW_LOG_WARN("existing DRM session for %s has different key in slot %d", drmSessionContexts[sessionSlot].drmSession->getKeySystem().c_str(), sessionSlot); } MW_LOG_WARN("deleting existing DRM session for %s ", drmSessionContexts[sessionSlot].drmSession->getKeySystem().c_str()); + /* DELIA-70726 fix: this slot may still be referenced by a GStreamer + * decryptor element of a previous, not-yet-fully-torn-down pipeline + * (rapid/back-to-back channel change). Block here until any decrypt() + * call already in flight against this session finishes, so the delete + * below cannot race with OCDMSessionAdapter::verifyOutputProtection()/ + * decrypt() running on the old pipeline's multiqueue thread. */ + drmSessionContexts[sessionSlot].drmSession->PrepareForDestruction(); MW_SAFE_DELETE(drmSessionContexts[sessionSlot].drmSession); } this->ProfileUpdateCb(); diff --git a/middleware/drm/helper/WidevineDrmHelper.cpp b/middleware/drm/helper/WidevineDrmHelper.cpp index 0f2075915c..02f4e30ee3 100755 --- a/middleware/drm/helper/WidevineDrmHelper.cpp +++ b/middleware/drm/helper/WidevineDrmHelper.cpp @@ -24,6 +24,7 @@ #include #include +#include #include "WidevineDrmHelper.h" #include "DrmUtils.h" @@ -184,18 +185,55 @@ void WidevineDrmHelper::setDrmMetaData(const std::string& metaData) void WidevineDrmHelper::setDefaultKeyID(const std::string& cencData) { + mDefaultKeySlot = -1; std::vector defaultKeyID(cencData.begin(), cencData.end()); + // Also convert UUID string (e.g. "f3dff538-b8c9-58e4-e8cd-96cf811d32dc") to 16-byte binary + // for comparison against binary keyIDs parsed from PSSH + std::vector defaultKeyIDBinary; + std::string uuidHex; + uuidHex.reserve(cencData.size()); + for (char c : cencData) + { + if (c != '-') + { + uuidHex += c; + } + } + if (uuidHex.size() == 32) + { + defaultKeyIDBinary.reserve(16); + for (size_t i = 0; i < uuidHex.size(); i += 2) + { + char hexPair[3] = {uuidHex[i], uuidHex[i + 1], '\0'}; + char* end = nullptr; + unsigned long v = std::strtoul(hexPair, &end, 16); + if (end != hexPair + 2 || v > 0xFF) + { + MW_LOG_WARN("setDefaultKeyID: invalid hex in cencData at offset %zu", i); + defaultKeyIDBinary.clear(); + break; + } + defaultKeyIDBinary.push_back(static_cast(v)); + } + } + if(!mKeyIDs.empty()) { for(auto& it : mKeyIDs) { - if(defaultKeyID == it.second) + if(defaultKeyID == it.second || defaultKeyIDBinary == it.second) { mDefaultKeySlot = it.first; - MW_LOG_WARN("setDefaultKeyID : %s slot : %d", cencData.c_str(), mDefaultKeySlot); + MW_LOG_WARN("setDefaultKeyID : %s slot : %d", PlayerLogManager::getHexDebugStr(it.second).c_str(), mDefaultKeySlot); + break; } } } + if (mDefaultKeySlot < 0 && !mKeyIDs.empty()) + { + mDefaultKeySlot = mKeyIDs.begin()->first; + MW_LOG_WARN("setDefaultKeyID: no match found for cencData, defaulting to first slot %d", mDefaultKeySlot); + } } @@ -212,17 +250,21 @@ void WidevineDrmHelper::createInitData(std::vector& initData) const void WidevineDrmHelper::getKey(std::vector& keyID) const { MW_LOG_WARN("WidevineDrmHelper::getKey defaultkey: %d mKeyIDs.size:%zu", mDefaultKeySlot, mKeyIDs.size()); - if ((mDefaultKeySlot >= 0) && (mDefaultKeySlot < mKeyIDs.size())) + if ((mDefaultKeySlot >= 0) && (mKeyIDs.find(mDefaultKeySlot) != mKeyIDs.end())) { - keyID = this->mKeyIDs.at(mDefaultKeySlot); + keyID = mKeyIDs.at(mDefaultKeySlot); } - else if (mKeyIDs.size() > 0) + else if (!mKeyIDs.empty()) { - keyID = this->mKeyIDs.at(0); + if (mDefaultKeySlot >= 0) + { + MW_LOG_WARN("mDefaultKeySlot(%d) not found in mKeyIDs, falling back to first entry", mDefaultKeySlot); + } + keyID = mKeyIDs.begin()->second; } else { - MW_LOG_ERR("No key"); + MW_LOG_ERR("No key available - mKeyIDs is empty"); } } diff --git a/middleware/gst-plugins/drm/gst/gstcdmidecryptor.cpp b/middleware/gst-plugins/drm/gst/gstcdmidecryptor.cpp index 10add97a62..c1fa226584 100755 --- a/middleware/gst-plugins/drm/gst/gstcdmidecryptor.cpp +++ b/middleware/gst-plugins/drm/gst/gstcdmidecryptor.cpp @@ -459,6 +459,7 @@ gst_cdmidecryptor_transform_caps(GstBaseTransform * trans, GST_LOG_OBJECT(trans, "returning %" GST_PTR_FORMAT, transformedCaps); if (direction == GST_PAD_SINK && !gst_caps_is_empty(transformedCaps)) { + GstCaps* sinkCapsCopy = NULL; g_mutex_lock(&cdmidecryptor->mutex); // clean up previous caps if (cdmidecryptor->sinkCaps) @@ -467,9 +468,12 @@ gst_cdmidecryptor_transform_caps(GstBaseTransform * trans, cdmidecryptor->sinkCaps = NULL; } cdmidecryptor->sinkCaps = gst_caps_copy(transformedCaps); - g_cond_signal(&cdmidecryptor->sinkCapsCond); - g_mutex_unlock(&cdmidecryptor->mutex); - GST_DEBUG_OBJECT(trans, "Set sinkCaps to %" GST_PTR_FORMAT, cdmidecryptor->sinkCaps); + sinkCapsCopy = gst_caps_ref(cdmidecryptor->sinkCaps); // take an extra ref to keep it alive + g_cond_signal(&cdmidecryptor->sinkCapsCond); + g_mutex_unlock(&cdmidecryptor->mutex); + GST_DEBUG_OBJECT(trans, "Set sinkCaps to %" GST_PTR_FORMAT, sinkCapsCopy); + gst_caps_unref(sinkCapsCopy); // release the extra ref + } return transformedCaps; } @@ -534,12 +538,19 @@ static GstFlowReturn gst_cdmidecryptor_transform_ip( { // call decrypt even for clear samples in order to copy it to a secure buffer. If secure buffers are not supported // decrypt() call will return without doing anything - if (cdmidecryptor->drmSession != NULL && cdmidecryptor->sinkCaps != NULL) - errorCode = cdmidecryptor->drmSession->decrypt(keyIDBuffer, ivBuffer, buffer, subSampleCount, subsamplesBuffer, cdmidecryptor->sinkCaps); + /* DELIA-70726 fix: guard against the DrmSession being concurrently torn down + * (e.g. DrmSessionManager reusing/evicting the slot during a back-to-back + * channel change) while this pipeline still has buffers in flight. */ + if (cdmidecryptor->drmSession != NULL && cdmidecryptor->sinkCaps != NULL + && cdmidecryptor->drmSession->AcquireForUse()) + { + errorCode = cdmidecryptor->drmSession->decrypt(keyIDBuffer, ivBuffer, buffer, subSampleCount, subsamplesBuffer, cdmidecryptor->sinkCaps); + cdmidecryptor->drmSession->ReleaseAfterUse(); + } else - { /* If drmSession creation failed, then the call will be aborted here */ + { /* If drmSession creation failed, or is being destroyed, the call will be aborted here */ result = GST_FLOW_NOT_SUPPORTED; - GST_ERROR_OBJECT(cdmidecryptor, "drmSession or sinkCaps is **** NULL ****, returning GST_FLOW_NOT_SUPPORTED"); + GST_ERROR_OBJECT(cdmidecryptor, "drmSession or sinkCaps is **** NULL **** (or session is being destroyed), returning GST_FLOW_NOT_SUPPORTED"); } } goto free_resources; @@ -656,7 +667,20 @@ static GstFlowReturn gst_cdmidecryptor_transform_ip( result = GST_FLOW_NOT_SUPPORTED; goto free_resources; } + /* DELIA-70726 fix: guard against the DrmSession being concurrently torn down + * (e.g. DrmSessionManager reusing/evicting the slot during a back-to-back + * channel change) while this pipeline still has buffers in flight. Without + * this guard, the multiqueue/decryptor thread can call decrypt() on a + * DrmSession that is being (or has already been) freed, causing a + * use-after-free SIGSEGV inside OCDMSessionAdapter::verifyOutputProtection(). */ + if (!cdmidecryptor->drmSession->AcquireForUse()) + { + GST_ERROR_OBJECT(cdmidecryptor, "drmSession is being destroyed, aborting decrypt"); + result = GST_FLOW_NOT_SUPPORTED; + goto free_resources; + } errorCode = cdmidecryptor->drmSession->decrypt(keyIDBuffer, ivBuffer, buffer, subSampleCount, subsamplesBuffer, cdmidecryptor->sinkCaps); + cdmidecryptor->drmSession->ReleaseAfterUse(); cdmidecryptor->streamEncrypted = true; if (errorCode != 0 || cdmidecryptor->hdcpOpProtectionFailCount) diff --git a/middleware/test/utests/fakes/FakeGStreamer.cpp b/middleware/test/utests/fakes/FakeGStreamer.cpp index ee05fa5f7a..6674a442d1 100644 --- a/middleware/test/utests/fakes/FakeGStreamer.cpp +++ b/middleware/test/utests/fakes/FakeGStreamer.cpp @@ -995,18 +995,30 @@ GstPad * gst_ghost_pad_new (const gchar * name, GstPad * target) GstCaps *gst_app_src_get_caps(GstAppSrc *appsrc) { TRACE_FUNC(); + if (g_mockGStreamer != nullptr) + { + return g_mockGStreamer->gst_app_src_get_caps(appsrc); + } return NULL; } GstSample *gst_sample_new (GstBuffer * buffer, GstCaps * caps, const GstSegment * segment, GstStructure * info) { TRACE_FUNC(); + if (g_mockGStreamer != nullptr) + { + return g_mockGStreamer->gst_sample_new(buffer, caps, segment, info); + } return NULL; } GstFlowReturn gst_app_src_push_sample (GstAppSrc * appsrc, GstSample * sample) { TRACE_FUNC(); + if (g_mockGStreamer != nullptr) + { + return g_mockGStreamer->gst_app_src_push_sample(appsrc, sample); + } return GST_FLOW_OK; } diff --git a/middleware/test/utests/mocks/MockGStreamer.h b/middleware/test/utests/mocks/MockGStreamer.h index b1db26c9b5..06a4bc1a11 100644 --- a/middleware/test/utests/mocks/MockGStreamer.h +++ b/middleware/test/utests/mocks/MockGStreamer.h @@ -84,6 +84,9 @@ class MockGStreamer MOCK_METHOD(void, gst_segment_init, (GstSegment *segment, GstFormat format)); MOCK_METHOD(GstEvent *, gst_event_new_segment, (GstSegment *segment)); MOCK_METHOD(GstEvent*, gst_event_new_custom, (GstEventType type, GstStructure* structure), ()); + MOCK_METHOD(GstCaps *, gst_app_src_get_caps, (GstAppSrc *appsrc)); + MOCK_METHOD(GstSample *, gst_sample_new, (GstBuffer *buffer, GstCaps *caps, const GstSegment *segment, GstStructure *info)); + MOCK_METHOD(GstFlowReturn, gst_app_src_push_sample, (GstAppSrc *appsrc, GstSample *sample)); /* gst_app_sink_get_type diff --git a/middleware/test/utests/tests/InterfacePlayerTests/InterfacePlayerFunctionTests.cpp b/middleware/test/utests/tests/InterfacePlayerTests/InterfacePlayerFunctionTests.cpp index cd34729ad5..041f8d6c75 100644 --- a/middleware/test/utests/tests/InterfacePlayerTests/InterfacePlayerFunctionTests.cpp +++ b/middleware/test/utests/tests/InterfacePlayerTests/InterfacePlayerFunctionTests.cpp @@ -1463,6 +1463,81 @@ TEST_F(InterfacePlayerTests, SendNewSegmentEvent_VideoMediaType) mInterfacePrivatePlayer->SendNewSegmentEvent(mediaType, startPts, stopPts); //failure } +TEST_F(InterfacePlayerTests, SendNewSegmentEvent_RialtoSink_CapsNull) +{ + GstMediaType mediaType = eGST_MEDIATYPE_VIDEO; + GstClockTime startPts = 1000; + GstClockTime stopPts = 2000; + mPlayerContext->stream[mediaType].format = GST_FORMAT_ISO_BMFF; + mPlayerContext->usingRialtoSink = true; + + // gst_app_src_get_caps returns NULL - segment cannot be pushed + EXPECT_CALL(*g_mockGStreamer, gst_segment_init(_, GST_FORMAT_TIME)) + .Times(1); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_get_caps(_)) + .WillOnce(Return(nullptr)); + + // gst_sample_new and gst_app_src_push_sample should NOT be called + EXPECT_CALL(*g_mockGStreamer, gst_sample_new(_, _, _, _)) + .Times(0); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_push_sample(_, _)) + .Times(0); + + mInterfacePrivatePlayer->SendNewSegmentEvent(mediaType, startPts, stopPts); +} + +TEST_F(InterfacePlayerTests, SendNewSegmentEvent_RialtoSink_PushSampleSuccess) +{ + GstMediaType mediaType = eGST_MEDIATYPE_VIDEO; + GstClockTime startPts = 1000; + GstClockTime stopPts = 2000; + mPlayerContext->stream[mediaType].format = GST_FORMAT_ISO_BMFF; + mPlayerContext->usingRialtoSink = true; + + GstCaps fakeCaps = {}; + GstSample fakeSample = {}; + + EXPECT_CALL(*g_mockGStreamer, gst_segment_init(_, GST_FORMAT_TIME)) + .Times(1); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_get_caps(_)) + .WillOnce(Return(&fakeCaps)); + EXPECT_CALL(*g_mockGStreamer, gst_sample_new(nullptr, &fakeCaps, _, nullptr)) + .WillOnce(Return(&fakeSample)); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_push_sample(_, &fakeSample)) + .WillOnce(Return(GST_FLOW_OK)); + // gst_sample_unref and gst_caps_unref expand to gst_mini_object_unref + EXPECT_CALL(*g_mockGStreamer, gst_mini_object_unref(_)) + .Times(2); + + mInterfacePrivatePlayer->SendNewSegmentEvent(mediaType, startPts, stopPts); +} + +TEST_F(InterfacePlayerTests, SendNewSegmentEvent_RialtoSink_PushSampleFailure) +{ + GstMediaType mediaType = eGST_MEDIATYPE_VIDEO; + GstClockTime startPts = 1000; + GstClockTime stopPts = 2000; + mPlayerContext->stream[mediaType].format = GST_FORMAT_ISO_BMFF; + mPlayerContext->usingRialtoSink = true; + + GstCaps fakeCaps = {}; + GstSample fakeSample = {}; + + EXPECT_CALL(*g_mockGStreamer, gst_segment_init(_, GST_FORMAT_TIME)) + .Times(1); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_get_caps(_)) + .WillOnce(Return(&fakeCaps)); + EXPECT_CALL(*g_mockGStreamer, gst_sample_new(nullptr, &fakeCaps, _, nullptr)) + .WillOnce(Return(&fakeSample)); + EXPECT_CALL(*g_mockGStreamer, gst_app_src_push_sample(_, &fakeSample)) + .WillOnce(Return(GST_FLOW_ERROR)); + // gst_sample_unref and gst_caps_unref still called even on push failure + EXPECT_CALL(*g_mockGStreamer, gst_mini_object_unref(_)) + .Times(2); + + mInterfacePrivatePlayer->SendNewSegmentEvent(mediaType, startPts, stopPts); +} + TEST_F(InterfacePlayerTests, Queue_and_ClearProtectionEvent) { std::string formatType = "cenc";