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);
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)
{
"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));
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;
}
}
-#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)
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");
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;
}