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