Commit 4ce68a44 authored by Jacob Walser's avatar Jacob Walser

Handle all message processing in _onBusMessage

And don't pass along message pointers
parent b67ca0e3
...@@ -43,8 +43,9 @@ VideoReceiver::VideoReceiver(QObject* parent) ...@@ -43,8 +43,9 @@ VideoReceiver::VideoReceiver(QObject* parent)
#if defined(QGC_GST_STREAMING) #if defined(QGC_GST_STREAMING)
_timer.setSingleShot(true); _timer.setSingleShot(true);
connect(&_timer, &QTimer::timeout, this, &VideoReceiver::_timeout); connect(&_timer, &QTimer::timeout, this, &VideoReceiver::_timeout);
connect(this, &VideoReceiver::recordingEOSReceived, this, &VideoReceiver::_eosCB); connect(this, &VideoReceiver::msgErrorReceived, this, &VideoReceiver::_handleError);
connect(this, &VideoReceiver::busMessage, this, &VideoReceiver::_handleBusMessage); connect(this, &VideoReceiver::msgEOSReceived, this, &VideoReceiver::_handleEOS);
connect(this, &VideoReceiver::msgStateChangedReceived, this, &VideoReceiver::_handleStateChanged);
#endif #endif
} }
...@@ -328,7 +329,12 @@ void VideoReceiver::stop() ...@@ -328,7 +329,12 @@ void VideoReceiver::stop()
GstBus* bus = gst_pipeline_get_bus(GST_PIPELINE(_pipeline)); GstBus* bus = gst_pipeline_get_bus(GST_PIPELINE(_pipeline));
GstMessage* message = gst_bus_timed_pop_filtered(bus, GST_CLOCK_TIME_NONE, (GstMessageType)(GST_MESSAGE_EOS|GST_MESSAGE_ERROR)); GstMessage* message = gst_bus_timed_pop_filtered(bus, GST_CLOCK_TIME_NONE, (GstMessageType)(GST_MESSAGE_EOS|GST_MESSAGE_ERROR));
gst_object_unref(bus); gst_object_unref(bus);
_handleBusMessage(message); if(GST_MESSAGE_TYPE(message) == GST_MESSAGE_ERROR) {
_shutdownPipeline();
qCritical() << "Error stopping pipeline!";
} else if(GST_MESSAGE_TYPE(message) == GST_MESSAGE_EOS) {
_handleEOS();
}
gst_message_unref(message); gst_message_unref(message);
} }
#endif #endif
...@@ -345,22 +351,7 @@ void VideoReceiver::setVideoSavePath(const QString & path) ...@@ -345,22 +351,7 @@ void VideoReceiver::setVideoSavePath(const QString & path)
qCDebug(VideoReceiverLog) << "New Path:" << _path; qCDebug(VideoReceiverLog) << "New Path:" << _path;
} }
#if defined(QGC_GST_STREAMING) void VideoReceiver::_shutdownPipeline() {
void VideoReceiver::_handleBusMessage(GstMessage* msg)
{
switch (GST_MESSAGE_TYPE(msg)) {
case GST_MESSAGE_ERROR:
do {
gchar* debug;
GError* error;
gst_message_parse_error(msg, &error, &debug);
g_free(debug);
qCritical() << error->message;
g_error_free(error);
} while(0);
// No break!
case GST_MESSAGE_EOS:
{
GstBus* bus = NULL; GstBus* bus = NULL;
if ((bus = gst_pipeline_get_bus(GST_PIPELINE(_pipeline))) != NULL) { if ((bus = gst_pipeline_get_bus(GST_PIPELINE(_pipeline))) != NULL) {
gst_bus_disable_sync_message_emission(bus); gst_bus_disable_sync_message_emission(bus);
...@@ -371,22 +362,41 @@ void VideoReceiver::_handleBusMessage(GstMessage* msg) ...@@ -371,22 +362,41 @@ void VideoReceiver::_handleBusMessage(GstMessage* msg)
gst_bin_remove(GST_BIN(_pipeline), _videoSink); gst_bin_remove(GST_BIN(_pipeline), _videoSink);
gst_object_unref(_pipeline); gst_object_unref(_pipeline);
_pipeline = NULL; _pipeline = NULL;
delete _sink;
_sink = NULL;
_serverPresent = false; _serverPresent = false;
_streaming = false; _streaming = false;
_recording = false; _recording = false;
_stopping = false; _stopping = false;
_running = false; _running = false;
emit recordingChanged(); emit recordingChanged();
}
#if defined(QGC_GST_STREAMING)
void VideoReceiver::_handleError() {
qCDebug(VideoReceiverLog) << "Gstreamer error!";
_shutdownPipeline();
}
#endif
#if defined(QGC_GST_STREAMING)
void VideoReceiver::_handleEOS() {
if(_stopping) {
_shutdownPipeline();
qCDebug(VideoReceiverLog) << "Stopped"; qCDebug(VideoReceiverLog) << "Stopped";
} else if(_recording && _sink->removing) {
_shutdownRecordingBranch();
} else {
qCritical() << "VideoReceiver: Unexpected EOS!";
_shutdownPipeline();
} }
break; }
case GST_MESSAGE_STATE_CHANGED: #endif
#if defined(QGC_GST_STREAMING)
void VideoReceiver::_handleStateChanged() {
_streaming = GST_STATE(_pipeline) == GST_STATE_PLAYING; _streaming = GST_STATE(_pipeline) == GST_STATE_PLAYING;
qCDebug(VideoReceiverLog) << "State changed, _streaming:" << _streaming; qCDebug(VideoReceiverLog) << "State changed, _streaming:" << _streaming;
break;
default:
break;
}
} }
#endif #endif
...@@ -396,7 +406,28 @@ gboolean VideoReceiver::_onBusMessage(GstBus* bus, GstMessage* msg, gpointer dat ...@@ -396,7 +406,28 @@ gboolean VideoReceiver::_onBusMessage(GstBus* bus, GstMessage* msg, gpointer dat
Q_UNUSED(bus) Q_UNUSED(bus)
Q_ASSERT(msg != NULL && data != NULL); Q_ASSERT(msg != NULL && data != NULL);
VideoReceiver* pThis = (VideoReceiver*)data; VideoReceiver* pThis = (VideoReceiver*)data;
pThis->busMessage(msg);
switch(GST_MESSAGE_TYPE(msg)) {
case(GST_MESSAGE_ERROR): {
gchar* debug;
GError* error;
gst_message_parse_error(msg, &error, &debug);
g_free(debug);
qCritical() << error->message;
g_error_free(error);
pThis->msgErrorReceived();
}
break;
case(GST_MESSAGE_EOS):
pThis->msgEOSReceived();
break;
case(GST_MESSAGE_STATE_CHANGED):
pThis->msgStateChangedReceived();
break;
default:
break;
}
return TRUE; return TRUE;
} }
#endif #endif
...@@ -486,10 +517,8 @@ void VideoReceiver::stopRecording(void) ...@@ -486,10 +517,8 @@ void VideoReceiver::stopRecording(void)
// -At this point all of the recoring elements have been flushed, and the video file has been finalized // -At this point all of the recoring elements have been flushed, and the video file has been finalized
// -Now we can remove the temporary pipeline and its elements // -Now we can remove the temporary pipeline and its elements
#if defined(QGC_GST_STREAMING) #if defined(QGC_GST_STREAMING)
void VideoReceiver::_eosCB(GstMessage* message) void VideoReceiver::_shutdownRecordingBranch()
{ {
Q_UNUSED(message)
gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->queue); gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->queue);
gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->mux); gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->mux);
gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->filesink); gst_bin_remove(GST_BIN(_pipelineStopRec), _sink->filesink);
...@@ -520,7 +549,7 @@ void VideoReceiver::_eosCB(GstMessage* message) ...@@ -520,7 +549,7 @@ void VideoReceiver::_eosCB(GstMessage* message)
// -Setup watch and handler for EOS event on the temporary pipeline's bus // -Setup watch and handler for EOS event on the temporary pipeline's bus
// -Send an EOS event at the beginning of that pipeline // -Send an EOS event at the beginning of that pipeline
#if defined(QGC_GST_STREAMING) #if defined(QGC_GST_STREAMING)
void VideoReceiver::_unlinkCB(GstPadProbeInfo* info) void VideoReceiver::_detachRecordingBranch(GstPadProbeInfo* info)
{ {
Q_UNUSED(info) Q_UNUSED(info)
...@@ -541,7 +570,7 @@ void VideoReceiver::_unlinkCB(GstPadProbeInfo* info) ...@@ -541,7 +570,7 @@ void VideoReceiver::_unlinkCB(GstPadProbeInfo* info)
// Add handler for EOS event // Add handler for EOS event
GstBus* bus = gst_pipeline_get_bus(GST_PIPELINE(_pipelineStopRec)); GstBus* bus = gst_pipeline_get_bus(GST_PIPELINE(_pipelineStopRec));
gst_bus_enable_sync_message_emission(bus); gst_bus_enable_sync_message_emission(bus);
g_signal_connect(bus, "sync-message", G_CALLBACK(_eosCallBack), this); g_signal_connect(bus, "sync-message", G_CALLBACK(_onBusMessage), this);
gst_object_unref(bus); gst_object_unref(bus);
if(gst_element_set_state(_pipelineStopRec, GST_STATE_PLAYING) == GST_STATE_CHANGE_FAILURE) { if(gst_element_set_state(_pipelineStopRec, GST_STATE_PLAYING) == GST_STATE_CHANGE_FAILURE) {
...@@ -556,31 +585,16 @@ void VideoReceiver::_unlinkCB(GstPadProbeInfo* info) ...@@ -556,31 +585,16 @@ void VideoReceiver::_unlinkCB(GstPadProbeInfo* info)
} }
#endif #endif
// This is only installed on the transient _pipelineStopRec in order
// to finalize a video file. It is not used for the main _pipeline.
#if defined(QGC_GST_STREAMING)
gboolean VideoReceiver::_eosCallBack(GstBus* bus, GstMessage* message, gpointer user_data)
{
Q_UNUSED(bus)
Q_ASSERT(message != NULL && user_data != NULL);
if(GST_MESSAGE_TYPE(message) == GST_MESSAGE_EOS) {
VideoReceiver* pThis = (VideoReceiver*)user_data;
pThis->recordingEOSReceived(message);
}
return FALSE;
}
#endif
#if defined(QGC_GST_STREAMING) #if defined(QGC_GST_STREAMING)
GstPadProbeReturn VideoReceiver::_unlinkCallBack(GstPad* pad, GstPadProbeInfo* info, gpointer user_data) GstPadProbeReturn VideoReceiver::_unlinkCallBack(GstPad* pad, GstPadProbeInfo* info, gpointer user_data)
{ {
Q_UNUSED(pad); Q_UNUSED(pad);
Q_ASSERT(info != NULL && user_data != NULL); Q_ASSERT(info != NULL && user_data != NULL);
VideoReceiver* pThis = (VideoReceiver*)user_data; VideoReceiver* pThis = (VideoReceiver*)user_data;
// We will only execute once // We will only act once
if(!g_atomic_int_compare_and_exchange(&pThis->_sink->removing, FALSE, TRUE)) if(g_atomic_int_compare_and_exchange(&pThis->_sink->removing, FALSE, TRUE))
return GST_PAD_PROBE_REMOVE; pThis->_detachRecordingBranch(info);
pThis->_unlinkCB(info);
return GST_PAD_PROBE_REMOVE; return GST_PAD_PROBE_REMOVE;
} }
#endif #endif
...@@ -49,8 +49,11 @@ public: ...@@ -49,8 +49,11 @@ public:
signals: signals:
void recordingChanged(); void recordingChanged();
void recordingEOSReceived(GstMessage* message); #if defined(QGC_GST_STREAMING)
void busMessage(GstMessage* message); void msgErrorReceived();
void msgEOSReceived();
void msgStateChangedReceived();
#endif
public slots: public slots:
void start (); void start ();
...@@ -63,10 +66,12 @@ public slots: ...@@ -63,10 +66,12 @@ public slots:
private slots: private slots:
#if defined(QGC_GST_STREAMING) #if defined(QGC_GST_STREAMING)
void _eosCB(GstMessage* message);
void _timeout (); void _timeout ();
void _connected (); void _connected ();
void _socketError (QAbstractSocket::SocketError socketError); void _socketError (QAbstractSocket::SocketError socketError);
void _handleError();
void _handleEOS();
void _handleStateChanged();
#endif #endif
private: private:
...@@ -88,11 +93,12 @@ private: ...@@ -88,11 +93,12 @@ private:
Sink* _sink; Sink* _sink;
GstElement* _tee; GstElement* _tee;
void _handleBusMessage(GstMessage* message);
void _unlinkCB(GstPadProbeInfo* info);
static gboolean _onBusMessage(GstBus* bus, GstMessage* message, gpointer user_data); static gboolean _onBusMessage(GstBus* bus, GstMessage* message, gpointer user_data);
static gboolean _eosCallBack(GstBus* bus, GstMessage* message, gpointer user_data);
static GstPadProbeReturn _unlinkCallBack(GstPad* pad, GstPadProbeInfo* info, gpointer user_data); static GstPadProbeReturn _unlinkCallBack(GstPad* pad, GstPadProbeInfo* info, gpointer user_data);
void _detachRecordingBranch(GstPadProbeInfo* info);
void _shutdownRecordingBranch();
void _shutdownPipeline();
#endif #endif
QString _uri; QString _uri;
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment