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