1 /* GStreamer
2  *
3  * Copyright (C) 2006 Thomas Vander Stichele <thomas at apestaart dot org>
4  *
5  * This library is free software; you can redistribute it and/or
6  * modify it under the terms of the GNU Library General Public
7  * License as published by the Free Software Foundation; either
8  * version 2 of the License, or (at your option) any later version.
9  *
10  * This library is distributed in the hope that it will be useful,
11  * but WITHOUT ANY WARRANTY; without even the implied warranty of
12  * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the GNU
13  * Library General Public License for more details.
14  *
15  * You should have received a copy of the GNU Library General Public
16  * License along with this library; if not, write to the
17  * Free Software Foundation, Inc., 51 Franklin St, Fifth Floor,
18  * Boston, MA 02110-1301, USA.
19  */
20 #ifdef HAVE_CONFIG_H
21 #include "config.h"
22 #endif
23 
24 #include <unistd.h>
25 #include <sys/ioctl.h>
26 #include <sys/socket.h>
27 #ifdef HAVE_FIONREAD_IN_SYS_FILIO
28 #include <sys/filio.h>
29 #endif
30 
31 #include <gio/gio.h>
32 #include <gst/check/gstcheck.h>
33 
34 static GstPad *mysrcpad;
35 
36 static GstStaticPadTemplate srctemplate = GST_STATIC_PAD_TEMPLATE ("src",
37     GST_PAD_SRC,
38     GST_PAD_ALWAYS,
39     GST_STATIC_CAPS ("application/x-gst-check")
40     );
41 
42 static GstElement *
setup_multisocketsink(void)43 setup_multisocketsink (void)
44 {
45   GstElement *multisocketsink;
46 
47   GST_DEBUG ("setup_multisocketsink");
48   multisocketsink = gst_check_setup_element ("multisocketsink");
49   mysrcpad = gst_check_setup_src_pad (multisocketsink, &srctemplate);
50   gst_pad_set_active (mysrcpad, TRUE);
51 
52   return multisocketsink;
53 }
54 
55 static void
cleanup_multisocketsink(GstElement * multisocketsink)56 cleanup_multisocketsink (GstElement * multisocketsink)
57 {
58   GST_DEBUG ("cleanup_multisocketsink");
59 
60   gst_check_teardown_src_pad (multisocketsink);
61   gst_check_teardown_element (multisocketsink);
62 }
63 
64 static void
wait_bytes_served(GstElement * sink,guint64 bytes)65 wait_bytes_served (GstElement * sink, guint64 bytes)
66 {
67   guint64 bytes_served = 0;
68 
69   while (bytes_served != bytes) {
70     g_object_get (sink, "bytes-served", &bytes_served, NULL);
71   }
72 }
73 
74 /* FIXME: possibly racy, since if it would write, we may not get it
75  * immediately ? */
76 #define fail_if_can_read(msg,fd) \
77 G_STMT_START { \
78   long avail; \
79 \
80   fail_if (ioctl (fd, FIONREAD, &avail) < 0, "%s: could not ioctl", msg); \
81   fail_if (avail > 0, "%s: has bytes available to read"); \
82 } G_STMT_END;
83 
84 
GST_START_TEST(test_no_clients)85 GST_START_TEST (test_no_clients)
86 {
87   GstElement *sink;
88   GstBuffer *buffer;
89   GstCaps *caps;
90 
91   sink = setup_multisocketsink ();
92 
93   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
94 
95   caps = gst_caps_from_string ("application/x-gst-check");
96   buffer = gst_buffer_new_and_alloc (4);
97   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
98   gst_caps_unref (caps);
99   fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
100 
101   GST_DEBUG ("cleaning up multisocketsink");
102   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
103   cleanup_multisocketsink (sink);
104 }
105 
106 GST_END_TEST;
107 
108 static gboolean
setup_handles(GSocket ** sinkhandle,GSocket ** srchandle)109 setup_handles (GSocket ** sinkhandle, GSocket ** srchandle)
110 {
111   GError *error = NULL;
112   gint sv[3];
113 
114 
115 //  g_assert (*sinkhandle);
116 //  g_assert (*srchandle);
117 
118   fail_if (socketpair (PF_UNIX, SOCK_STREAM, 0, sv));
119 
120   *sinkhandle = g_socket_new_from_fd (sv[1], &error);
121   fail_if (error);
122   fail_if (*sinkhandle == NULL);
123   *srchandle = g_socket_new_from_fd (sv[0], &error);
124   fail_if (error);
125   fail_if (*srchandle == NULL);
126 
127   return TRUE;
128 }
129 
130 static gboolean
read_handle_n_bytes_exactly(GSocket * srchandle,void * buf,size_t count)131 read_handle_n_bytes_exactly (GSocket * srchandle, void *buf, size_t count)
132 {
133   gssize total_read, read;
134   gchar *data = buf;
135 
136   GST_DEBUG ("reading exactly %" G_GSIZE_FORMAT " bytes", count);
137 
138   /* loop to make sure the sink has had a chance to write out all data.
139    * Depending on system load it might be written in multiple write calls,
140    * so it's possible our first read() just returns parts of the data. */
141   total_read = 0;
142   do {
143     read =
144         g_socket_receive (srchandle, data + total_read, count - total_read,
145         NULL, NULL);
146 
147     if (read == 0)              /* socket was closed */
148       return FALSE;
149 
150     if (read < 0)
151       fail ("read error");
152 
153     total_read += read;
154 
155     GST_INFO ("read %" G_GSSIZE_FORMAT " bytes, total now %" G_GSSIZE_FORMAT,
156         read, total_read);
157   }
158   while (total_read < count);
159 
160   return TRUE;
161 }
162 
163 static ssize_t
read_handle(GSocket * srchandle,void * buf,size_t count)164 read_handle (GSocket * srchandle, void *buf, size_t count)
165 {
166   gssize ret;
167 
168   ret = g_socket_receive (srchandle, buf, count, NULL, NULL);
169 
170   return ret;
171 }
172 
173 #define fail_unless_read(msg,handle,size,ref) \
174 G_STMT_START { \
175   char data[size + 1]; \
176   int nbytes; \
177 \
178   GST_DEBUG ("%s: reading %d bytes", msg, size); \
179   nbytes = read_handle (handle, data, size); \
180   data[size] = 0; \
181   GST_DEBUG ("%s: read %d bytes", msg, nbytes); \
182   fail_if (nbytes < size); \
183   fail_unless (memcmp (data, ref, size) == 0, \
184       "data read '%s' differs from '%s'", data, ref); \
185 } G_STMT_END;
186 
187 #define fail_unless_num_handles(sink,num) \
188 G_STMT_START { \
189   gint handles; \
190   g_object_get (sink, "num-handles", &handles, NULL); \
191   fail_unless (handles == num, \
192       "sink has %d handles instead of expected %d", handles, num); \
193 } G_STMT_END;
194 
GST_START_TEST(test_add_client)195 GST_START_TEST (test_add_client)
196 {
197   GstElement *sink;
198   GstBuffer *buffer;
199   GstCaps *caps;
200   gchar data[9];
201   GSocket *sinksocket, *srcsocket;
202 
203   sink = setup_multisocketsink ();
204   fail_unless (setup_handles (&sinksocket, &srcsocket));
205 
206 
207   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
208 
209   /* add the client */
210   g_signal_emit_by_name (sink, "add", sinksocket);
211 
212   caps = gst_caps_from_string ("application/x-gst-check");
213   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
214   GST_DEBUG ("Created test caps %p %" GST_PTR_FORMAT, caps, caps);
215   buffer = gst_buffer_new_and_alloc (4);
216   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
217   ASSERT_CAPS_REFCOUNT (caps, "caps", 3);
218   gst_buffer_fill (buffer, 0, "dead", 4);
219   gst_buffer_append_memory (buffer,
220       gst_memory_new_wrapped (GST_MEMORY_FLAG_READONLY, (gpointer) " good", 5,
221           0, 5, NULL, NULL));
222   fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
223 
224   GST_DEBUG ("reading");
225   fail_if (read_handle (srcsocket, data, 9) < 9);
226   fail_unless (strncmp (data, "dead good", 9) == 0);
227   wait_bytes_served (sink, 9);
228 
229   GST_DEBUG ("cleaning up multisocketsink");
230   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
231   cleanup_multisocketsink (sink);
232 
233   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
234   gst_caps_unref (caps);
235 
236   g_object_unref (srcsocket);
237   g_object_unref (sinksocket);
238 }
239 
240 GST_END_TEST;
241 
242 typedef struct
243 {
244   GSocket *sinksocket, *srcsocket;
245   GstElement *sink;
246 } TestSinkAndSocket;
247 
248 static void
setup_sink_with_socket(TestSinkAndSocket * tsas)249 setup_sink_with_socket (TestSinkAndSocket * tsas)
250 {
251   GstCaps *caps = NULL;
252 
253   tsas->sink = setup_multisocketsink ();
254   fail_unless (setup_handles (&tsas->sinksocket, &tsas->srcsocket));
255 
256   ASSERT_SET_STATE (tsas->sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
257 
258   /* add the client */
259   g_signal_emit_by_name (tsas->sink, "add", tsas->sinksocket);
260 
261   caps = gst_caps_from_string ("application/x-gst-check");
262   gst_check_setup_events (mysrcpad, tsas->sink, caps, GST_FORMAT_BYTES);
263   gst_caps_unref (caps);
264 }
265 
266 static void
teardown_sink_with_socket(TestSinkAndSocket * tsas)267 teardown_sink_with_socket (TestSinkAndSocket * tsas)
268 {
269   if (tsas->sink != NULL) {
270     ASSERT_SET_STATE (tsas->sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
271     cleanup_multisocketsink (tsas->sink);
272     tsas->sink = 0;
273   }
274   if (tsas->sinksocket != NULL) {
275     g_object_unref (tsas->sinksocket);
276     tsas->sinksocket = 0;
277   }
278   if (tsas->srcsocket != NULL) {
279     g_object_unref (tsas->srcsocket);
280     tsas->srcsocket = 0;
281   }
282 }
283 
GST_START_TEST(test_sending_buffers_with_9_gstmemories)284 GST_START_TEST (test_sending_buffers_with_9_gstmemories)
285 {
286   TestSinkAndSocket tsas = { 0 };
287   GstBuffer *buffer;
288   int i;
289   const char *numbers[9] = { "one", "two", "three", "four", "five", "six",
290     "seven", "eight", "nine"
291   };
292   const char numbers_concat[] = "onetwothreefourfivesixseveneightnine";
293   gchar data[sizeof (numbers_concat)];
294   int len = sizeof (numbers_concat) - 1;
295 
296   setup_sink_with_socket (&tsas);
297 
298   buffer = gst_buffer_new ();
299   for (i = 0; i < G_N_ELEMENTS (numbers); i++)
300     gst_buffer_append_memory (buffer,
301         gst_memory_new_wrapped (GST_MEMORY_FLAG_READONLY, (gpointer) numbers[i],
302             strlen (numbers[i]), 0, strlen (numbers[i]), NULL, NULL));
303   fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
304 
305   fail_unless (read_handle_n_bytes_exactly (tsas.srcsocket, data, len));
306   fail_unless (strncmp (data, numbers_concat, len) == 0);
307 
308   teardown_sink_with_socket (&tsas);
309 }
310 
311 GST_END_TEST;
312 
313 /* from the given two data buffers, create two streamheader buffers and
314  * some caps that match it, and store them in the given pointers
315  * returns  one ref to each of the buffers and the caps */
316 static void
gst_multisocketsink_create_streamheader(const gchar * data1,const gchar * data2,GstBuffer ** hbuf1,GstBuffer ** hbuf2,GstCaps ** caps)317 gst_multisocketsink_create_streamheader (const gchar * data1,
318     const gchar * data2, GstBuffer ** hbuf1, GstBuffer ** hbuf2,
319     GstCaps ** caps)
320 {
321   GstBuffer *buf;
322   GValue array = { 0 };
323   GValue value = { 0 };
324   GstStructure *structure;
325   guint size1 = strlen (data1);
326   guint size2 = strlen (data2);
327 
328   fail_if (hbuf1 == NULL);
329   fail_if (hbuf2 == NULL);
330   fail_if (caps == NULL);
331 
332   /* create caps with streamheader, set the caps, and push the HEADER
333    * buffers */
334   *hbuf1 = gst_buffer_new_and_alloc (size1);
335   GST_BUFFER_FLAG_SET (*hbuf1, GST_BUFFER_FLAG_HEADER);
336   gst_buffer_fill (*hbuf1, 0, data1, size1);
337   *hbuf2 = gst_buffer_new_and_alloc (size2);
338   GST_BUFFER_FLAG_SET (*hbuf2, GST_BUFFER_FLAG_HEADER);
339   gst_buffer_fill (*hbuf2, 0, data2, size2);
340 
341   g_value_init (&array, GST_TYPE_ARRAY);
342 
343   g_value_init (&value, GST_TYPE_BUFFER);
344   /* we take a copy, set it on the array (which refs it), then unref our copy */
345   buf = gst_buffer_copy (*hbuf1);
346   gst_value_set_buffer (&value, buf);
347   ASSERT_BUFFER_REFCOUNT (buf, "copied buffer", 2);
348   gst_buffer_unref (buf);
349   gst_value_array_append_value (&array, &value);
350   g_value_unset (&value);
351 
352   g_value_init (&value, GST_TYPE_BUFFER);
353   buf = gst_buffer_copy (*hbuf2);
354   gst_value_set_buffer (&value, buf);
355   ASSERT_BUFFER_REFCOUNT (buf, "copied buffer", 2);
356   gst_buffer_unref (buf);
357   gst_value_array_append_value (&array, &value);
358   g_value_unset (&value);
359 
360   *caps = gst_caps_from_string ("application/x-gst-check");
361   structure = gst_caps_get_structure (*caps, 0);
362 
363   gst_structure_set_value (structure, "streamheader", &array);
364   g_value_unset (&array);
365   ASSERT_CAPS_REFCOUNT (*caps, "streamheader caps", 1);
366 
367   /* we want to keep them around for the tests */
368   gst_buffer_ref (*hbuf1);
369   gst_buffer_ref (*hbuf2);
370 
371   GST_DEBUG ("created streamheader caps %p %" GST_PTR_FORMAT, *caps, *caps);
372 }
373 
374 
375 /* this test:
376  * - adds a first client
377  * - sets streamheader caps on the pad
378  * - pushes the HEADER buffers
379  * - pushes a buffer
380  * - verifies that the client received all the data correctly, and did not
381  *   get multiple copies of the streamheader
382  * - adds a second client
383  * - verifies that this second client receives the streamheader caps too, plus
384  * - the new buffer
385  */
GST_START_TEST(test_streamheader)386 GST_START_TEST (test_streamheader)
387 {
388   GstElement *sink;
389   GstBuffer *hbuf1, *hbuf2, *buf;
390   GstCaps *caps;
391   GSocket *socket[4];
392 
393   sink = setup_multisocketsink ();
394 
395   fail_unless (setup_handles (&socket[0], &socket[1]));
396   fail_unless (setup_handles (&socket[2], &socket[3]));
397 
398   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
399 
400   /* add the first client */
401   fail_unless_num_handles (sink, 0);
402   g_signal_emit_by_name (sink, "add", socket[0]);
403   fail_unless_num_handles (sink, 1);
404 
405   /* create caps with streamheader, set the caps, and push the HEADER
406    * buffers */
407   gst_multisocketsink_create_streamheader ("babe", "deadbeef", &hbuf1, &hbuf2,
408       &caps);
409   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 2);
410   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 2);
411   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
412   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
413   /* one is ours, two from set_caps */
414   ASSERT_CAPS_REFCOUNT (caps, "caps", 3);
415 
416   fail_unless (gst_pad_push (mysrcpad, hbuf1) == GST_FLOW_OK);
417   fail_unless (gst_pad_push (mysrcpad, hbuf2) == GST_FLOW_OK);
418   // FIXME: we can't assert on the refcount because giving away the ref
419   //        doesn't mean the refcount decreases
420   // ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 1);
421   // ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 1);
422 
423   //FIXME:
424   //fail_if_can_read ("first client", socket[1]);
425 
426   /* push a non-HEADER buffer, this should trigger the client receiving the
427    * first three buffers */
428   buf = gst_buffer_new_and_alloc (4);
429   gst_buffer_fill (buf, 0, "f00d", 4);
430   gst_pad_push (mysrcpad, buf);
431 
432   fail_unless_read ("first client", socket[1], 4, "babe");
433   fail_unless_read ("first client", socket[1], 8, "deadbeef");
434   fail_unless_read ("first client", socket[1], 4, "f00d");
435   wait_bytes_served (sink, 16);
436 
437   /* now add the second client */
438   g_signal_emit_by_name (sink, "add", socket[2]);
439   fail_unless_num_handles (sink, 2);
440   //FIXME:
441   //fail_if_can_read ("second client", socket[3]);
442 
443   /* now push another buffer, which will trigger streamheader for second
444    * client */
445   buf = gst_buffer_new_and_alloc (4);
446   gst_buffer_fill (buf, 0, "deaf", 4);
447   gst_pad_push (mysrcpad, buf);
448 
449   fail_unless_read ("first client", socket[1], 4, "deaf");
450 
451   fail_unless_read ("second client", socket[3], 4, "babe");
452   fail_unless_read ("second client", socket[3], 8, "deadbeef");
453   /* we missed the f00d buffer */
454   fail_unless_read ("second client", socket[3], 4, "deaf");
455   wait_bytes_served (sink, 36);
456 
457   GST_DEBUG ("cleaning up multisocketsink");
458 
459   fail_unless_num_handles (sink, 2);
460   g_signal_emit_by_name (sink, "remove", socket[0]);
461   fail_unless_num_handles (sink, 1);
462   g_signal_emit_by_name (sink, "remove", socket[2]);
463   fail_unless_num_handles (sink, 0);
464 
465   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
466   cleanup_multisocketsink (sink);
467 
468   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 1);
469   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 1);
470   gst_buffer_unref (hbuf1);
471   gst_buffer_unref (hbuf2);
472 
473   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
474   gst_caps_unref (caps);
475 
476   g_object_unref (socket[0]);
477   g_object_unref (socket[1]);
478   g_object_unref (socket[2]);
479   g_object_unref (socket[3]);
480 }
481 
482 GST_END_TEST;
483 
484 /* this tests changing of streamheaders
485  * - set streamheader caps on the pad
486  * - pushes the HEADER buffers
487  * - pushes a buffer
488  * - add a first client
489  * - verifies that this first client receives the first streamheader caps,
490  *   plus a new buffer
491  * - change streamheader caps
492  * - verify that the first client receives the new streamheader buffers as well
493  */
GST_START_TEST(test_change_streamheader)494 GST_START_TEST (test_change_streamheader)
495 {
496   GstElement *sink;
497   GstBuffer *hbuf1, *hbuf2, *buf;
498   GstCaps *caps;
499   GSocket *socket[4];
500 
501   sink = setup_multisocketsink ();
502 
503   fail_unless (setup_handles (&socket[0], &socket[1]));
504   fail_unless (setup_handles (&socket[2], &socket[3]));
505 
506   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
507 
508   /* create caps with streamheader, set the caps, and push the HEADER
509    * buffers */
510   gst_multisocketsink_create_streamheader ("first", "header", &hbuf1, &hbuf2,
511       &caps);
512   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
513   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
514   /* one is ours, two from set_caps */
515   ASSERT_CAPS_REFCOUNT (caps, "caps", 3);
516 
517   /* one to hold for the test and one to give away */
518   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 2);
519   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 2);
520 
521   fail_unless (gst_pad_push (mysrcpad, hbuf1) == GST_FLOW_OK);
522   fail_unless (gst_pad_push (mysrcpad, hbuf2) == GST_FLOW_OK);
523 
524   /* add the first client */
525   g_signal_emit_by_name (sink, "add", socket[0]);
526 
527   /* verify this hasn't triggered a write yet */
528   /* FIXME: possibly racy, since if it would write, we may not get it
529    * immediately ? */
530   //fail_if_can_read ("first client, no buffer", socket[1]);
531 
532   /* now push a buffer and read */
533   buf = gst_buffer_new_and_alloc (4);
534   gst_buffer_fill (buf, 0, "f00d", 4);
535   gst_pad_push (mysrcpad, buf);
536 
537   fail_unless_read ("change: first client", socket[1], 5, "first");
538   fail_unless_read ("change: first client", socket[1], 6, "header");
539   fail_unless_read ("change: first client", socket[1], 4, "f00d");
540   //wait_bytes_served (sink, 16);
541 
542   /* now add the second client */
543   g_signal_emit_by_name (sink, "add", socket[2]);
544   //fail_if_can_read ("second client, no buffer", socket[3]);
545 
546   /* change the streamheader */
547 
548   /* only we have a reference to the streamheaders now */
549   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 1);
550   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 1);
551   gst_buffer_unref (hbuf1);
552   gst_buffer_unref (hbuf2);
553 
554   /* drop our ref to the previous caps */
555   gst_caps_unref (caps);
556 
557   gst_multisocketsink_create_streamheader ("second", "header", &hbuf1, &hbuf2,
558       &caps);
559   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
560   /* one to hold for the test and one to give away */
561   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 2);
562   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 2);
563 
564   fail_unless (gst_pad_push (mysrcpad, hbuf1) == GST_FLOW_OK);
565   fail_unless (gst_pad_push (mysrcpad, hbuf2) == GST_FLOW_OK);
566 
567   /* verify neither client has new data available to read */
568   //fail_if_can_read ("first client, changed streamheader", socket[1]);
569   //fail_if_can_read ("second client, changed streamheader", socket[3]);
570 
571   /* now push another buffer, which will trigger streamheader for second
572    * client, but should also send new streamheaders to first client */
573   buf = gst_buffer_new_and_alloc (8);
574   gst_buffer_fill (buf, 0, "deadbabe", 8);
575   gst_pad_push (mysrcpad, buf);
576 
577   fail_unless_read ("first client", socket[1], 6, "second");
578   fail_unless_read ("first client", socket[1], 6, "header");
579   fail_unless_read ("first client", socket[1], 8, "deadbabe");
580 
581   /* new streamheader data */
582   fail_unless_read ("second client", socket[3], 6, "second");
583   fail_unless_read ("second client", socket[3], 6, "header");
584   /* we missed the f00d buffer */
585   fail_unless_read ("second client", socket[3], 8, "deadbabe");
586   //wait_bytes_served (sink, 36);
587 
588   GST_DEBUG ("cleaning up multisocketsink");
589   g_signal_emit_by_name (sink, "remove", socket[0]);
590   g_signal_emit_by_name (sink, "remove", socket[2]);
591   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
592 
593   /* setting to NULL should have cleared the streamheader */
594   ASSERT_BUFFER_REFCOUNT (hbuf1, "hbuf1", 1);
595   ASSERT_BUFFER_REFCOUNT (hbuf2, "hbuf2", 1);
596   gst_buffer_unref (hbuf1);
597   gst_buffer_unref (hbuf2);
598   cleanup_multisocketsink (sink);
599 
600   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
601   gst_caps_unref (caps);
602 
603   g_object_unref (socket[0]);
604   g_object_unref (socket[1]);
605   g_object_unref (socket[2]);
606   g_object_unref (socket[3]);
607 }
608 
609 GST_END_TEST;
610 
611 static GstBuffer *
gst_new_buffer(int i)612 gst_new_buffer (int i)
613 {
614   GstMapInfo info;
615   gchar *data;
616 
617   GstBuffer *buffer = gst_buffer_new_and_alloc (16);
618 
619   /* copy some id */
620   g_assert (gst_buffer_map (buffer, &info, GST_MAP_WRITE));
621   data = (gchar *) info.data;
622   g_snprintf (data, 16, "deadbee%08x", i);
623   gst_buffer_unmap (buffer, &info);
624 
625   return buffer;
626 }
627 
628 
629 /* keep 100 bytes and burst 80 bytes to clients */
GST_START_TEST(test_burst_client_bytes)630 GST_START_TEST (test_burst_client_bytes)
631 {
632   GstElement *sink;
633   GstCaps *caps;
634   GSocket *socket[6];
635   gint i;
636   guint buffers_queued;
637 
638   sink = setup_multisocketsink ();
639   /* make sure we keep at least 100 bytes at all times */
640   g_object_set (sink, "bytes-min", 100, NULL);
641   g_object_set (sink, "sync-method", 3, NULL);  /* 3 = burst */
642   g_object_set (sink, "burst-format", GST_FORMAT_BYTES, NULL);
643   g_object_set (sink, "burst-value", (guint64) 80, NULL);
644 
645   fail_unless (setup_handles (&socket[0], &socket[1]));
646   fail_unless (setup_handles (&socket[2], &socket[3]));
647   fail_unless (setup_handles (&socket[4], &socket[5]));
648 
649   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
650 
651   caps = gst_caps_from_string ("application/x-gst-check");
652   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
653   GST_DEBUG ("Created test caps %p %" GST_PTR_FORMAT, caps, caps);
654 
655   /* push buffers in, 9 * 16 bytes = 144 bytes */
656   for (i = 0; i < 9; i++) {
657     GstBuffer *buffer = gst_new_buffer (i);
658 
659     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
660   }
661 
662   /* check that at least 7 buffers (112 bytes) are in the queue */
663   g_object_get (sink, "buffers-queued", &buffers_queued, NULL);
664   fail_if (buffers_queued != 7);
665 
666   /* now add the clients */
667   fail_unless_num_handles (sink, 0);
668   g_signal_emit_by_name (sink, "add", socket[0]);
669   fail_unless_num_handles (sink, 1);
670   g_signal_emit_by_name (sink, "add_full", socket[2], 3,
671       GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 200);
672   g_signal_emit_by_name (sink, "add_full", socket[4], 3,
673       GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 50);
674   fail_unless_num_handles (sink, 3);
675 
676   /* push last buffer to make client fds ready for reading */
677   for (i = 9; i < 10; i++) {
678     GstBuffer *buffer = gst_new_buffer (i);
679 
680     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
681   }
682 
683   /* now we should only read the last 5 buffers (5 * 16 = 80 bytes) */
684   GST_DEBUG ("Reading from client 1");
685   fail_unless_read ("client 1", socket[1], 16, "deadbee00000005");
686   fail_unless_read ("client 1", socket[1], 16, "deadbee00000006");
687   fail_unless_read ("client 1", socket[1], 16, "deadbee00000007");
688   fail_unless_read ("client 1", socket[1], 16, "deadbee00000008");
689   fail_unless_read ("client 1", socket[1], 16, "deadbee00000009");
690 
691   /* second client only bursts 50 bytes = 4 buffers (we get 4 buffers since
692    * the max allows it) */
693   GST_DEBUG ("Reading from client 2");
694   fail_unless_read ("client 2", socket[3], 16, "deadbee00000006");
695   fail_unless_read ("client 2", socket[3], 16, "deadbee00000007");
696   fail_unless_read ("client 2", socket[3], 16, "deadbee00000008");
697   fail_unless_read ("client 2", socket[3], 16, "deadbee00000009");
698 
699   /* third client only bursts 50 bytes = 4 buffers, we can't send
700    * more than 50 bytes so we only get 3 buffers (48 bytes). */
701   GST_DEBUG ("Reading from client 3");
702   fail_unless_read ("client 3", socket[5], 16, "deadbee00000007");
703   fail_unless_read ("client 3", socket[5], 16, "deadbee00000008");
704   fail_unless_read ("client 3", socket[5], 16, "deadbee00000009");
705 
706   GST_DEBUG ("cleaning up multisocketsink");
707   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
708   cleanup_multisocketsink (sink);
709 
710   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
711   gst_caps_unref (caps);
712 
713   g_object_unref (socket[0]);
714   g_object_unref (socket[1]);
715   g_object_unref (socket[2]);
716   g_object_unref (socket[3]);
717   g_object_unref (socket[4]);
718   g_object_unref (socket[5]);
719 }
720 
721 GST_END_TEST;
722 
723 /* keep 100 bytes and burst 80 bytes to clients */
GST_START_TEST(test_burst_client_bytes_keyframe)724 GST_START_TEST (test_burst_client_bytes_keyframe)
725 {
726   GstElement *sink;
727   GstCaps *caps;
728   GSocket *socket[6];
729   gint i;
730   guint buffers_queued;
731 
732   sink = setup_multisocketsink ();
733   /* make sure we keep at least 100 bytes at all times */
734   g_object_set (sink, "bytes-min", 100, NULL);
735   g_object_set (sink, "sync-method", 4, NULL);  /* 4 = burst_keyframe */
736   g_object_set (sink, "burst-format", GST_FORMAT_BYTES, NULL);
737   g_object_set (sink, "burst-value", (guint64) 80, NULL);
738 
739   fail_unless (setup_handles (&socket[0], &socket[1]));
740   fail_unless (setup_handles (&socket[2], &socket[3]));
741   fail_unless (setup_handles (&socket[4], &socket[5]));
742 
743   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
744 
745   caps = gst_caps_from_string ("application/x-gst-check");
746   GST_DEBUG ("Created test caps %p %" GST_PTR_FORMAT, caps, caps);
747   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
748 
749   /* push buffers in, 9 * 16 bytes = 144 bytes */
750   for (i = 0; i < 9; i++) {
751     GstBuffer *buffer = gst_new_buffer (i);
752 
753     /* mark most buffers as delta */
754     if (i != 0 && i != 4 && i != 8)
755       GST_BUFFER_FLAG_SET (buffer, GST_BUFFER_FLAG_DELTA_UNIT);
756 
757     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
758   }
759 
760   /* check that at least 7 buffers (112 bytes) are in the queue */
761   g_object_get (sink, "buffers-queued", &buffers_queued, NULL);
762   fail_if (buffers_queued != 7);
763 
764   /* now add the clients */
765   g_signal_emit_by_name (sink, "add", socket[0]);
766   g_signal_emit_by_name (sink, "add_full", socket[2],
767       4, GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 90);
768   g_signal_emit_by_name (sink, "add_full", socket[4],
769       4, GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 50);
770 
771   /* push last buffer to make client fds ready for reading */
772   for (i = 9; i < 10; i++) {
773     GstBuffer *buffer = gst_new_buffer (i);
774 
775     GST_BUFFER_FLAG_SET (buffer, GST_BUFFER_FLAG_DELTA_UNIT);
776 
777     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
778   }
779 
780   /* now we should only read the last 6 buffers (min 5 * 16 = 80 bytes),
781    * keyframe at buffer 4 */
782   GST_DEBUG ("Reading from client 1");
783   fail_unless_read ("client 1", socket[1], 16, "deadbee00000004");
784   fail_unless_read ("client 1", socket[1], 16, "deadbee00000005");
785   fail_unless_read ("client 1", socket[1], 16, "deadbee00000006");
786   fail_unless_read ("client 1", socket[1], 16, "deadbee00000007");
787   fail_unless_read ("client 1", socket[1], 16, "deadbee00000008");
788   fail_unless_read ("client 1", socket[1], 16, "deadbee00000009");
789 
790   /* second client only bursts 50 bytes = 4 buffers, there is
791    * no keyframe above min and below max, so get one below min */
792   GST_DEBUG ("Reading from client 2");
793   fail_unless_read ("client 2", socket[3], 16, "deadbee00000008");
794   fail_unless_read ("client 2", socket[3], 16, "deadbee00000009");
795 
796   /* third client only bursts 50 bytes = 4 buffers, we can't send
797    * more than 50 bytes so we only get 2 buffers (32 bytes). */
798   GST_DEBUG ("Reading from client 3");
799   fail_unless_read ("client 3", socket[5], 16, "deadbee00000008");
800   fail_unless_read ("client 3", socket[5], 16, "deadbee00000009");
801 
802   GST_DEBUG ("cleaning up multisocketsink");
803   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
804   cleanup_multisocketsink (sink);
805 
806   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
807   gst_caps_unref (caps);
808 
809   g_object_unref (socket[0]);
810   g_object_unref (socket[1]);
811   g_object_unref (socket[2]);
812   g_object_unref (socket[3]);
813   g_object_unref (socket[4]);
814   g_object_unref (socket[5]);
815 }
816 
817 GST_END_TEST;
818 
819 
820 
821 /* keep 100 bytes and burst 80 bytes to clients */
GST_START_TEST(test_burst_client_bytes_with_keyframe)822 GST_START_TEST (test_burst_client_bytes_with_keyframe)
823 {
824   GstElement *sink;
825   GstCaps *caps;
826   GSocket *socket[6];
827   gint i;
828   guint buffers_queued;
829 
830   sink = setup_multisocketsink ();
831 
832   /* make sure we keep at least 100 bytes at all times */
833   g_object_set (sink, "bytes-min", 100, NULL);
834   g_object_set (sink, "sync-method", 5, NULL);  /* 5 = burst_with_keyframe */
835   g_object_set (sink, "burst-format", GST_FORMAT_BYTES, NULL);
836   g_object_set (sink, "burst-value", (guint64) 80, NULL);
837 
838   fail_unless (setup_handles (&socket[0], &socket[1]));
839   fail_unless (setup_handles (&socket[2], &socket[3]));
840   fail_unless (setup_handles (&socket[4], &socket[5]));
841 
842   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
843 
844   caps = gst_caps_from_string ("application/x-gst-check");
845   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
846   GST_DEBUG ("Created test caps %p %" GST_PTR_FORMAT, caps, caps);
847 
848   /* push buffers in, 9 * 16 bytes = 144 bytes */
849   for (i = 0; i < 9; i++) {
850     GstBuffer *buffer = gst_new_buffer (i);
851 
852     /* mark most buffers as delta */
853     if (i != 0 && i != 4 && i != 8)
854       GST_BUFFER_FLAG_SET (buffer, GST_BUFFER_FLAG_DELTA_UNIT);
855 
856     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
857   }
858 
859   /* check that at least 7 buffers (112 bytes) are in the queue */
860   g_object_get (sink, "buffers-queued", &buffers_queued, NULL);
861   fail_if (buffers_queued != 7);
862 
863   /* now add the clients */
864   g_signal_emit_by_name (sink, "add", socket[0]);
865   g_signal_emit_by_name (sink, "add_full", socket[2],
866       5, GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 90);
867   g_signal_emit_by_name (sink, "add_full", socket[4],
868       5, GST_FORMAT_BYTES, (guint64) 50, GST_FORMAT_BYTES, (guint64) 50);
869 
870   /* push last buffer to make client fds ready for reading */
871   for (i = 9; i < 10; i++) {
872     GstBuffer *buffer = gst_new_buffer (i);
873 
874     GST_BUFFER_FLAG_SET (buffer, GST_BUFFER_FLAG_DELTA_UNIT);
875 
876     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
877   }
878 
879   /* now we should only read the last 6 buffers (min 5 * 16 = 80 bytes),
880    * keyframe at buffer 4 */
881   GST_DEBUG ("Reading from client 1");
882   fail_unless_read ("client 1", socket[1], 16, "deadbee00000004");
883   fail_unless_read ("client 1", socket[1], 16, "deadbee00000005");
884   fail_unless_read ("client 1", socket[1], 16, "deadbee00000006");
885   fail_unless_read ("client 1", socket[1], 16, "deadbee00000007");
886   fail_unless_read ("client 1", socket[1], 16, "deadbee00000008");
887   fail_unless_read ("client 1", socket[1], 16, "deadbee00000009");
888 
889   /* second client only bursts 50 bytes = 4 buffers, there is
890    * no keyframe above min and below max, so send min */
891   GST_DEBUG ("Reading from client 2");
892   fail_unless_read ("client 2", socket[3], 16, "deadbee00000006");
893   fail_unless_read ("client 2", socket[3], 16, "deadbee00000007");
894   fail_unless_read ("client 2", socket[3], 16, "deadbee00000008");
895   fail_unless_read ("client 2", socket[3], 16, "deadbee00000009");
896 
897   /* third client only bursts 50 bytes = 4 buffers, we can't send
898    * more than 50 bytes so we only get 3 buffers (48 bytes). */
899   GST_DEBUG ("Reading from client 3");
900   fail_unless_read ("client 3", socket[5], 16, "deadbee00000007");
901   fail_unless_read ("client 3", socket[5], 16, "deadbee00000008");
902   fail_unless_read ("client 3", socket[5], 16, "deadbee00000009");
903 
904   GST_DEBUG ("cleaning up multisocketsink");
905   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
906   cleanup_multisocketsink (sink);
907 
908   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
909   gst_caps_unref (caps);
910 
911   g_object_unref (socket[0]);
912   g_object_unref (socket[1]);
913   g_object_unref (socket[2]);
914   g_object_unref (socket[3]);
915   g_object_unref (socket[4]);
916   g_object_unref (socket[5]);
917 }
918 
919 GST_END_TEST;
920 
921 /* Check that we can get data when multisocketsink is configured in next-keyframe
922  * mode */
GST_START_TEST(test_client_next_keyframe)923 GST_START_TEST (test_client_next_keyframe)
924 {
925   GstElement *sink;
926   GstCaps *caps;
927   GSocket *socket[2];
928   gint i;
929 
930   sink = setup_multisocketsink ();
931   g_object_set (sink, "sync-method", 1, NULL);  /* 1 = next-keyframe */
932 
933   fail_unless (setup_handles (&socket[0], &socket[1]));
934 
935   ASSERT_SET_STATE (sink, GST_STATE_PLAYING, GST_STATE_CHANGE_ASYNC);
936 
937   caps = gst_caps_from_string ("application/x-gst-check");
938   gst_check_setup_events (mysrcpad, sink, caps, GST_FORMAT_BYTES);
939   GST_DEBUG ("Created test caps %p %" GST_PTR_FORMAT, caps, caps);
940 
941   /* now add our client */
942   g_signal_emit_by_name (sink, "add", socket[0]);
943 
944   /* push buffers in: keyframe, then non-keyframe */
945   for (i = 0; i < 2; i++) {
946     GstBuffer *buffer = gst_new_buffer (i);
947     if (i > 0)
948       GST_BUFFER_FLAG_SET (buffer, GST_BUFFER_FLAG_DELTA_UNIT);
949 
950     fail_unless (gst_pad_push (mysrcpad, buffer) == GST_FLOW_OK);
951   }
952 
953   /* now we should be able to read some data */
954   GST_DEBUG ("Reading from client 1");
955   fail_unless_read ("client 1", socket[1], 16, "deadbee00000000");
956   fail_unless_read ("client 1", socket[1], 16, "deadbee00000001");
957 
958   GST_DEBUG ("cleaning up multisocketsink");
959   ASSERT_SET_STATE (sink, GST_STATE_NULL, GST_STATE_CHANGE_SUCCESS);
960   cleanup_multisocketsink (sink);
961 
962   ASSERT_CAPS_REFCOUNT (caps, "caps", 1);
963   gst_caps_unref (caps);
964 
965   g_object_unref (socket[0]);
966   g_object_unref (socket[1]);
967 }
968 
969 GST_END_TEST;
970 
971 /* FIXME: add test simulating chained oggs where:
972  * sync-method is burst-on-connect
973  * (when multisocketsink actually does burst-on-connect based on byte size, not
974    "last keyframe" which any frame for audio :))
975  * an old client still needs to read from before the new streamheaders
976  * a new client gets the new streamheaders
977  */
978 static Suite *
multisocketsink_suite(void)979 multisocketsink_suite (void)
980 {
981   Suite *s = suite_create ("multisocketsink");
982   TCase *tc_chain = tcase_create ("general");
983 
984   suite_add_tcase (s, tc_chain);
985   tcase_add_test (tc_chain, test_no_clients);
986   tcase_add_test (tc_chain, test_add_client);
987   tcase_add_test (tc_chain, test_sending_buffers_with_9_gstmemories);
988   tcase_add_test (tc_chain, test_streamheader);
989   tcase_add_test (tc_chain, test_change_streamheader);
990   tcase_add_test (tc_chain, test_burst_client_bytes);
991   tcase_add_test (tc_chain, test_burst_client_bytes_keyframe);
992   tcase_add_test (tc_chain, test_burst_client_bytes_with_keyframe);
993   tcase_add_test (tc_chain, test_client_next_keyframe);
994 
995   return s;
996 }
997 
998 GST_CHECK_MAIN (multisocketsink);
999