diff options
| author | 2024-01-20 18:45:29 -0500 | |
|---|---|---|
| committer | 2024-01-20 18:45:29 -0500 | |
| commit | 88ea5bdca85fc0954d1dd64309a308199c1af4a8 (patch) | |
| tree | ef9dc778b9c42ab2fecce8810f777ccc133cebef /src/libsink | |
| parent | b79f074a86525ece1394af7bc5870516a8c38ba9 (diff) | |
| download | camu-88ea5bdca85fc0954d1dd64309a308199c1af4a8.tar.gz camu-88ea5bdca85fc0954d1dd64309a308199c1af4a8.tar.bz2 camu-88ea5bdca85fc0954d1dd64309a308199c1af4a8.zip | |
wip
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/sink.c | 170 | ||||
| -rw-r--r-- | src/libsink/sink.h | 7 |
2 files changed, 126 insertions, 51 deletions
diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 58e3db3..343275b 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -1,8 +1,11 @@ #include <al/log.h> #include "../tree/common.h" +#include "../bimu/common.h" +#include "../bimu/handler.h" #include "sink.h" +#include "common.h" enum { SINK_EMPTY = 0, @@ -12,6 +15,7 @@ enum { enum { ENTRY_LOADED = 0, + ENTRY_BUFFERED, ENTRY_DISREGUARDED }; @@ -27,8 +31,9 @@ enum { START, STOP, TOGGLE_PAUSE, + RESEEK, SEEK, - ENTRY_BUFFERED, // Currently set entry is buffered. + SET_BUFFERED, // Currently set entry is buffered. CLOSE, // Internal. CORK, @@ -94,6 +99,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) case TOGGLE_PAUSE: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; if (camu_clock_is_paused(&entry->clock)) { + camu_clock_calc_tick(&entry->clock); camu_clock_resume(&entry->clock); if (sink->audio.state == SINK_PAUSED) { sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); @@ -116,36 +122,59 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) } break; } + case RESEEK: { + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; + bmu_client_reseek(&entry->client); + break; + } case SEEK: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; bmu_client_seek(&entry->client, cmd->value.f); break; } - case ENTRY_BUFFERED: { +// case SET_BUFFERED: { +// switch (cmd->value.i) { +// case CAMU_SINK_AUDIO: +// sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL); +// break; +//#ifndef CAMU_SINK_NO_VIDEO +// case CAMU_SINK_VIDEO: +// sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL); +// break; +//#endif +// } +// break; +// } + case CLOSE: { + sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); + return; + } + case CORK: { + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; switch (cmd->value.i) { case CAMU_SINK_AUDIO: - sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL); + bmu_vcr_stream_cork(entry->audio.stream); break; #ifndef CAMU_SINK_NO_VIDEO case CAMU_SINK_VIDEO: - sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL); + bmu_vcr_stream_cork(entry->video.stream); break; #endif } break; } - case CLOSE: { - sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); - return; - } - case CORK: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - bmu_client_cork(&entry->client, true); - break; - } case UNCORK: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - bmu_client_cork(&entry->client, false); + switch (cmd->value.i) { + case CAMU_SINK_AUDIO: + bmu_vcr_stream_uncork(entry->audio.stream); + break; +#ifndef CAMU_SINK_NO_VIDEO + case CAMU_SINK_VIDEO: + bmu_vcr_stream_uncork(entry->video.stream); + break; +#endif + } break; } } @@ -173,18 +202,27 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->audio.state; if (state == BUFFER_SET_OR_BUFFERED) { - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = ENTRY_BUFFERED, - .value.i = CAMU_SINK_AUDIO - }); + bool should_resume = ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED; + if (should_resume && !camu_clock_calc_tick(&entry->clock)) { + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = RESEEK, + .opaque = entry + }); + return; + } queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO }); - if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) { + if (should_resume) { camu_clock_resume(&entry->clock); + entry->state = ENTRY_BUFFERED; } + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + //queue_cmd(entry->sink, (struct camu_sink_cmd){ + // .op = SET_BUFFERED, + // .value.i = CAMU_SINK_AUDIO + //}); state = BUFFER_ADDED; } else if (state == BUFFER_CONFIGURED) { state = BUFFER_SET_OR_BUFFERED; @@ -204,12 +242,14 @@ static void audio_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_CORK: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = CORK, + .value.i = CAMU_SINK_AUDIO, .opaque = entry }); break; case CAMU_BUFFER_UNCORK: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = UNCORK, + .value.i = CAMU_SINK_AUDIO, .opaque = entry }); break; @@ -234,10 +274,10 @@ static void audio_buffer_callback(void *userdata, u8 op) .value.i = CAMU_SINK_AUDIO }); #ifndef CAMU_SINK_NO_VIDEO - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = ENTRY_BUFFERED, - .value.i = CAMU_SINK_VIDEO - }); + //queue_cmd(entry->sink, (struct camu_sink_cmd){ + // .op = SET_BUFFERED, + // .value.i = CAMU_SINK_VIDEO + //}); #endif } break; @@ -250,18 +290,27 @@ static void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->video.state; if (state == BUFFER_SET_OR_BUFFERED) { - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = ENTRY_BUFFERED, - .value.i = CAMU_SINK_VIDEO - }); + bool should_resume = ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED; + if (should_resume && !camu_clock_calc_tick(&entry->clock)) { + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = RESEEK, + .opaque = entry + }); + return; + } queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_VIDEO }); - if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) { + if (should_resume) { camu_clock_resume(&entry->clock); + entry->state = ENTRY_BUFFERED; } + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); + //queue_cmd(entry->sink, (struct camu_sink_cmd){ + // .op = SET_BUFFERED, + // .value.i = CAMU_SINK_VIDEO + //}); state = BUFFER_ADDED; } else if (state == BUFFER_CONFIGURED) { state = BUFFER_SET_OR_BUFFERED; @@ -281,12 +330,14 @@ static void video_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_CORK: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = CORK, + .value.i = CAMU_SINK_VIDEO, .opaque = entry }); break; case CAMU_BUFFER_UNCORK: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = UNCORK, + .value.i = CAMU_SINK_VIDEO, .opaque = entry }); break; @@ -300,22 +351,10 @@ static void video_buffer_callback(void *userdata, u8 op) } #endif -static void client_callback(void *userdata, u8 op, struct bmu_client_stream *stream, void *opaque) +static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream, void *opaque) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; switch (op) { - case BIMU_CLIENT_SET: { - struct bmu_seek_req *req = (struct bmu_seek_req *)opaque; - camu_clock_set(&entry->clock, req->base, req->start); - camu_audio_buffer_reset(&entry->audio.buf); -#ifndef CAMU_SINK_NO_VIDEO - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (!single_frame) { - camu_video_buffer_reset(&entry->video.buf); - } -#endif - break; - } case BIMU_CLIENT_CONFIGURE: { // No data can be sent until all active streams are configured. switch (stream->type) { @@ -346,6 +385,23 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str } break; } + case BIMU_CLIENT_SET: { + u64 ts = aki_get_timestamp(); + struct bmu_seek_req *req = (struct bmu_seek_req *)opaque; + //req->delay = 0; // Ignore delay. + req->ts -= req->delay; + bool late = ts <= req->ts; + if (late) {} + camu_clock_set(&entry->clock, req->base, req->ts, req->delay); + camu_audio_buffer_reset(&entry->audio.buf); + if (!late) { + // This is a fail-safe for the case where receiving data is extremely slow. + // It should not be depended on for checking if an entry buffers in time. + aki_timer_set_repeat(&entry->timer, AKI_TS_FROM_USEC(ts - req->ts)); + aki_timer_again(&entry->timer); + } + break; + } case BIMU_CLIENT_DATA: { struct camu_frame *frame = (struct camu_frame *)opaque; switch (stream->type) { @@ -395,6 +451,11 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str #ifndef CAMU_SINK_NO_VIDEO } #endif +#ifndef CAMU_SINK_NO_VIDEO + if (!single_frame) { + camu_video_buffer_reset(&entry->video.buf); + } +#endif break; } case BIMU_CLIENT_EOF: { @@ -417,6 +478,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str } case BIMU_CLIENT_CLOSED: { bmu_client_free(&entry->client); + aki_timer_stop(&entry->timer); camu_audio_buffer_free(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO camu_video_buffer_free(&entry->video.buf); @@ -460,6 +522,15 @@ static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 no return NULL; } +static void buffer_timer_callback(void *userdata, struct aki_timer *timer) +{ + struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; + (void)timer; + aki_mutex_lock(&entry->sink->mutex); + aki_mutex_unlock(&entry->sink->mutex); + aki_timer_stop(&entry->timer); +} + static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { @@ -492,11 +563,14 @@ static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *c entry->video.buf.userdata = entry; #endif + aki_timer_init(&entry->timer, sink->loop, buffer_timer_callback, entry); + entry->client.callback = client_callback; entry->client.userdata = entry; bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id); aki_packet_free(packet); + return false; } @@ -524,6 +598,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn u16 node_id = aki_packet_read_u16(packet); aki_mutex_lock(&sink->mutex); struct camu_sink_entry *entry = entry_from_node_id(sink, node_id); + if (!entry) goto out; struct camu_sink_entry *previous = sink->current; sink->current = entry; if (entry->audio.state == BUFFER_INIT) { @@ -542,8 +617,9 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn camu_clock_pause(&previous->clock); remove_entry_buffers(sink, previous); } - aki_mutex_unlock(&sink->mutex); +out: + aki_mutex_unlock(&sink->mutex); aki_packet_free(packet); return false; } @@ -644,13 +720,13 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { aki_mutex_lock(&sink->mutex); - /* struct camu_sink_entry *entry; al_array_foreach(sink->entries, i, entry) { + bmu_client_close(&entry->client); } - */ - al_array_free(sink->entries); aki_mutex_unlock(&sink->mutex); + al_array_free(sink->entries); + aki_rpc_disconnect(&sink->client); aki_signal_stop(&sink->signal); camu_queue_free(sink->queue); aki_mutex_destroy(&sink->mutex); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 95c280f..e3550bb 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -18,8 +18,6 @@ #include "../bimu/client.h" -#include "common.h" - enum { CAMU_SINK_AUDIO = 0, #ifndef CAMU_SINK_NO_VIDEO @@ -41,16 +39,17 @@ struct camu_sink_entry { u8 state; struct bmu_client client; struct camu_clock clock; + struct aki_timer timer; struct { u8 state; struct camu_audio_buffer buf; - struct bmu_client_stream *stream; + struct bmu_vcr_stream *stream; } audio; #ifndef CAMU_SINK_NO_VIDEO struct { u8 state; struct camu_video_buffer buf; - struct bmu_client_stream *stream; + struct bmu_vcr_stream *stream; } video; #endif struct camu_sink *sink; |