summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorWilliam Manley <will@williammanley.net>2015-03-13 13:56:13 (GMT)
committerWim Taymans <wtaymans@redhat.com>2015-03-13 19:05:00 (GMT)
commita297b0545f664c82351bbae1fd3f943d068b5031 (patch)
tree502423dc31b58b30ae8b393ac0d3de8b786b7911
parenta19ac4b85c9dbd930bf68058e8581b7e09a640ea (diff)
downloadgst-plugins-base-a297b0545f664c82351bbae1fd3f943d068b5031.tar.gz
gst-plugins-base-a297b0545f664c82351bbae1fd3f943d068b5031.tar.xz
socketsrc: Add `connection-closed-by-peer` signal
This provides notification that the socket in use was closed by the peer and gives an opportunity to replace it with a new one which is not closed, allowing reading from many sockets in order. I use this in pulsevideo to implement reconnection logic to handle the pulsevideo service dieing, such that is can be restarted without disrupting downstream. Fixes https://bugzilla.gnome.org/show_bug.cgi?id=739546
-rw-r--r--gst/tcp/gstsocketsrc.c50
-rw-r--r--gst/tcp/gstsocketsrc.h3
-rw-r--r--tests/check/pipelines/tcp.c67
3 files changed, 115 insertions, 5 deletions
diff --git a/gst/tcp/gstsocketsrc.c b/gst/tcp/gstsocketsrc.c
index 578a86c..219f653 100644
--- a/gst/tcp/gstsocketsrc.c
+++ b/gst/tcp/gstsocketsrc.c
@@ -70,6 +70,14 @@ enum
PROP_SOCKET,
};
+enum
+{
+ CONNECTION_CLOSED_BY_PEER,
+ LAST_SIGNAL
+};
+
+static guint gst_socket_src_signals[LAST_SIGNAL] = { 0 };
+
#define gst_socket_src_parent_class parent_class
G_DEFINE_TYPE (GstSocketSrc, gst_socket_src, GST_TYPE_PUSH_SRC);
@@ -86,6 +94,8 @@ static void gst_socket_src_set_property (GObject * object, guint prop_id,
static void gst_socket_src_get_property (GObject * object, guint prop_id,
GValue * value, GParamSpec * pspec);
+#define SWAP(a, b) do { GSocket* _swap_tmp = a; a = b; b = _swap_tmp; } while (0);
+
static void
gst_socket_src_class_init (GstSocketSrcClass * klass)
{
@@ -108,6 +118,11 @@ gst_socket_src_class_init (GstSocketSrcClass * klass)
"The socket to receive packets from", G_TYPE_SOCKET,
G_PARAM_READWRITE | G_PARAM_STATIC_STRINGS));
+ gst_socket_src_signals[CONNECTION_CLOSED_BY_PEER] =
+ g_signal_new ("connection-closed-by-peer", G_TYPE_FROM_CLASS (klass),
+ G_SIGNAL_RUN_FIRST, G_STRUCT_OFFSET (GstSocketSrcClass,
+ connection_closed_by_peer), NULL, NULL, NULL, G_TYPE_NONE, 0);
+
gst_element_class_add_pad_template (gstelement_class,
gst_static_pad_template_get (&srctemplate));
@@ -167,14 +182,43 @@ gst_socket_src_fill (GstPushSrc * psrc, GstBuffer * outbuf)
GST_LOG_OBJECT (src, "asked for a buffer");
+retry:
gst_buffer_map (outbuf, &map, GST_MAP_READWRITE);
rret = g_socket_receive_with_blocking (socket, (gchar *) map.data,
map.size, TRUE, src->cancellable, &err);
gst_buffer_unmap (outbuf, &map);
if (rret == 0) {
- GST_DEBUG_OBJECT (src, "Connection closed");
- ret = GST_FLOW_EOS;
+ GSocket *tmp = NULL;
+ GST_DEBUG_OBJECT (src, "Received EOS on socket %p fd %i", socket,
+ g_socket_get_fd (socket));
+
+ /* We've hit EOS but we'll send this signal to allow someone to change
+ * our socket before we send EOS downstream. */
+ g_signal_emit (src, gst_socket_src_signals[CONNECTION_CLOSED_BY_PEER], 0);
+
+ GST_OBJECT_LOCK (src);
+
+ if (src->socket)
+ tmp = g_object_ref (src->socket);
+
+ GST_OBJECT_UNLOCK (src);
+
+ /* Do this dance with tmp to avoid unreffing with the lock held */
+ if (tmp != NULL && tmp != socket) {
+ SWAP (socket, tmp);
+ g_clear_object (&tmp);
+
+ GST_INFO_OBJECT (src, "New socket available after EOS %p fd %i: Retrying",
+ socket, g_socket_get_fd (socket));
+
+ /* retry with our new socket: */
+ goto retry;
+ } else {
+ g_clear_object (&tmp);
+ GST_INFO_OBJECT (src, "Forwarding EOS downstream");
+ ret = GST_FLOW_EOS;
+ }
} else if (rret < 0) {
if (g_error_matches (err, G_IO_ERROR, G_IO_ERROR_CANCELLED)) {
ret = GST_FLOW_FLUSHING;
@@ -210,8 +254,6 @@ no_socket:
}
}
-#define SWAP(a, b) do { GSocket* tmp = a; a = b; b = tmp; } while (0);
-
static void
gst_socket_src_set_property (GObject * object, guint prop_id,
const GValue * value, GParamSpec * pspec)
diff --git a/gst/tcp/gstsocketsrc.h b/gst/tcp/gstsocketsrc.h
index ba08a74..e628b3a 100644
--- a/gst/tcp/gstsocketsrc.h
+++ b/gst/tcp/gstsocketsrc.h
@@ -54,6 +54,9 @@ struct _GstSocketSrc {
struct _GstSocketSrcClass {
GstPushSrcClass parent_class;
+
+ /* signals */
+ void (*connection_closed_by_peer) (GstElement*);
};
GType gst_socket_src_get_type (void);
diff --git a/tests/check/pipelines/tcp.c b/tests/check/pipelines/tcp.c
index 71cdba1..997e4bd 100644
--- a/tests/check/pipelines/tcp.c
+++ b/tests/check/pipelines/tcp.c
@@ -202,7 +202,70 @@ GST_START_TEST (test_that_tcpserversink_and_tcpclientsrc_are_symmetrical)
GST_END_TEST;
-static Suite *
+static void
+on_connection_closed (GstElement * socketsrc, gpointer user_data)
+{
+ GSocket *socket = (GSocket *) user_data;
+
+ g_object_set (socketsrc, "socket", socket, NULL);
+}
+
+GST_START_TEST (test_that_we_can_provide_new_socketsrc_sockets_during_signal)
+{
+ GSocket *sockets[4] = { NULL, NULL };
+
+ GstPipeline *pipeline = NULL;
+ GstAppSink *appsink = NULL;
+ GstElement *socketsrc = NULL;
+ GstSample *sample = NULL;
+
+ socketsrc = gst_check_setup_element ("socketsrc");
+
+ fail_unless (g_socketpair (G_SOCKET_FAMILY_UNIX,
+ G_SOCKET_TYPE_STREAM | SOCK_CLOEXEC, G_SOCKET_PROTOCOL_DEFAULT,
+ &sockets[0], NULL));
+
+ fail_unless (g_socket_send (sockets[0], "hello", 5, NULL, NULL) == 5);
+ fail_unless (g_socket_shutdown (sockets[0], FALSE, TRUE, NULL));
+
+ fail_unless (g_socketpair (G_SOCKET_FAMILY_UNIX,
+ G_SOCKET_TYPE_STREAM | SOCK_CLOEXEC, G_SOCKET_PROTOCOL_DEFAULT,
+ &sockets[2], NULL));
+ fail_unless (g_socket_send (sockets[2], "goodbye", 7, NULL, NULL) == 7);
+ fail_unless (g_socket_shutdown (sockets[2], FALSE, TRUE, NULL));
+
+ g_object_set (socketsrc, "socket", sockets[1], NULL);
+
+ g_signal_connect (socketsrc, "connection-closed-by-peer",
+ G_CALLBACK (on_connection_closed), sockets[3]);
+
+ pipeline = (GstPipeline *) gst_pipeline_new (NULL);
+ appsink = GST_APP_SINK (gst_check_setup_element ("appsink"));
+ gst_bin_add_many (GST_BIN (pipeline), socketsrc, GST_ELEMENT (appsink), NULL);
+ fail_unless (gst_element_link_many (socketsrc, GST_ELEMENT (appsink), NULL));
+
+ gst_element_set_state (GST_ELEMENT (pipeline), GST_STATE_PLAYING);
+
+ fail_unless ((sample = gst_app_sink_pull_sample (appsink)) != NULL);
+ gst_buffer_memcmp (gst_sample_get_buffer (sample), 0, "hello", 5);
+ gst_sample_unref (sample);
+
+ fail_unless ((sample = gst_app_sink_pull_sample (appsink)) != NULL);
+ gst_buffer_memcmp (gst_sample_get_buffer (sample), 0, "goodbye", 7);
+ gst_sample_unref (sample);
+
+ fail_unless (NULL == gst_app_sink_pull_sample (appsink));
+ fail_unless (gst_app_sink_is_eos (appsink));
+
+ gst_element_set_state (GST_ELEMENT (pipeline), GST_STATE_NULL);
+ g_clear_object (&sockets[0]);
+ g_clear_object (&sockets[1]);
+ g_clear_object (&sockets[2]);
+ g_clear_object (&sockets[3]);
+ gst_object_unref (pipeline);
+}
+
+GST_END_TEST static Suite *
socketintegrationtest_suite (void)
{
Suite *s = suite_create ("socketintegrationtest");
@@ -215,6 +278,8 @@ socketintegrationtest_suite (void)
test_that_tcpclientsink_and_tcpserversrc_are_symmetrical);
tcase_add_test (tc_chain,
test_that_tcpserversink_and_tcpclientsrc_are_symmetrical);
+ tcase_add_test (tc_chain,
+ test_that_we_can_provide_new_socketsrc_sockets_during_signal);
return s;
}