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