From e6c1c0afc69ee13465bb68a6bdeb35b6876146d4 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Sat, 6 Jan 2024 19:54:21 -0500 Subject: Initial synced playback Signed-off-by: Andrew Opalach --- src/libsink/meson.build | 5 +- src/libsink/sink.c | 530 ++++++++++++++++--------------------- src/libsink/sink.h | 30 +-- src/libsink/sink2.c | 688 ------------------------------------------------ src/libsink/sink2.h | 106 -------- 5 files changed, 241 insertions(+), 1118 deletions(-) delete mode 100644 src/libsink/sink2.c delete mode 100644 src/libsink/sink2.h (limited to 'src/libsink') diff --git a/src/libsink/meson.build b/src/libsink/meson.build index f212b0c..0e6dfb3 100644 --- a/src/libsink/meson.build +++ b/src/libsink/meson.build @@ -1,2 +1,3 @@ -libsink_src = ['sink2.c'] -libsink = declare_dependency(sources: libsink_src) +libsink_src = ['sink.c'] +libsink_deps = [bimu_client] +libsink = declare_dependency(sources: libsink_src, dependencies: libsink_deps) diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 3692066..5addc0f 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -1,5 +1,7 @@ #include +#include "../tree/common.h" + #include "sink.h" enum { @@ -15,20 +17,22 @@ enum { enum { BUFFER_INIT = 0, + BUFFER_QUEUED, BUFFER_CONFIGURED, BUFFER_SET_OR_BUFFERED, BUFFER_ADDED, }; enum { - ADD_BUFFER = 0, - REMOVE_BUFFER, - SET_BUFFERED, START, STOP, TOGGLE_PAUSE, SEEK, - CLOSE + ENTRY_BUFFERED, // Currently set entry is buffered. + CLOSE, + // Internal. + CORK, + UNCORK }; #define ENTRY_AUDIO_BUFFER_HELD(entry) \ @@ -41,60 +45,15 @@ enum { #define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry)) #endif +#define ENTRY_VIDEO_EMPTY(entry) \ + (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED) +#define ENTRY_AUDIO_EMPTY(entry) \ + (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED) + static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) { switch (cmd->op) { - case ADD_BUFFER: { - switch (cmd->value.i) { - case CAMU_SINK_AUDIO: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - break; - } -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - break; - } -#endif - } - break; - } - case REMOVE_BUFFER: { - switch (cmd->value.i) { - case CAMU_SINK_AUDIO: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - break; - } -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - break; - } -#endif - } - break; - } - 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 START: { - struct camu_sink *sink = (struct camu_sink *)cmd->opaque; switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PAUSED) { @@ -114,7 +73,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case STOP: { - struct camu_sink *sink = (struct camu_sink *)cmd->opaque; switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PLAYING) { @@ -160,12 +118,19 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) } case SEEK: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_seek(&((struct camu_sink_local *)entry)->runner, cmd->value.f); + (void)entry; + break; + } + case ENTRY_BUFFERED: { + switch (cmd->value.i) { + case CAMU_SINK_AUDIO: + sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL); break; - default: +#ifndef CAMU_SINK_NO_VIDEO + case CAMU_SINK_VIDEO: + sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL); break; +#endif } break; } @@ -173,6 +138,16 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) 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); + break; + } } } @@ -200,21 +175,16 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) 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 = SET_BUFFERED, + .op = ENTRY_BUFFERED, .value.i = CAMU_SINK_AUDIO }); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, - .value.i = CAMU_SINK_AUDIO, - .opaque = entry->sink + .value.i = CAMU_SINK_AUDIO }); -#ifndef CAMU_SINK_NO_VIDEO - if (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_ADDED) { -#endif + if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) { camu_clock_resume(&entry->clock); -#ifndef CAMU_SINK_NO_VIDEO } -#endif state = BUFFER_ADDED; } else if (state == BUFFER_CONFIGURED) { state = BUFFER_SET_OR_BUFFERED; @@ -231,29 +201,22 @@ static void audio_buffer_callback(void *userdata, u8 op) add_audio_if_set_and_buffered(entry); aki_mutex_unlock(&entry->sink->mutex); break; - case CAMU_BUFFER_STOP: - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_stream_stop((struct bmu_local_stream *)entry->audio.stream); - break; - default: - break; - } + case CAMU_BUFFER_CORK: + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = CORK, + .opaque = entry + }); break; - case CAMU_BUFFER_CONTINUE: - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_stream_continue((struct bmu_local_stream *)entry->audio.stream); - break; - default: - break; - } + case CAMU_BUFFER_UNCORK: + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = UNCORK, + .opaque = entry + }); break; case CAMU_BUFFER_PAUSED: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = STOP, - .value.i = CAMU_SINK_AUDIO, - .opaque = entry->sink + .value.i = CAMU_SINK_AUDIO }); break; case CAMU_BUFFER_EOF: { @@ -267,12 +230,11 @@ static void audio_buffer_callback(void *userdata, u8 op) // the timer, don't queue the stop. queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = STOP, - .value.i = CAMU_SINK_AUDIO, - .opaque = entry->sink + .value.i = CAMU_SINK_AUDIO }); #ifndef CAMU_SINK_NO_VIDEO queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = SET_BUFFERED, + .op = ENTRY_BUFFERED, .value.i = CAMU_SINK_VIDEO }); #endif @@ -289,15 +251,14 @@ static void add_video_if_set_and_buffered(struct camu_sink_entry *entry) 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 = SET_BUFFERED, + .op = ENTRY_BUFFERED, .value.i = CAMU_SINK_VIDEO }); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, - .value.i = CAMU_SINK_VIDEO, - .opaque = entry->sink + .value.i = CAMU_SINK_VIDEO }); - if (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_ADDED) { + if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) { camu_clock_resume(&entry->clock); } state = BUFFER_ADDED; @@ -316,29 +277,22 @@ static void video_buffer_callback(void *userdata, u8 op) add_video_if_set_and_buffered(entry); aki_mutex_unlock(&entry->sink->mutex); break; - case CAMU_BUFFER_STOP: - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_stream_stop((struct bmu_local_stream *)entry->video.stream); - break; - default: - break; - } + case CAMU_BUFFER_CORK: + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = CORK, + .opaque = entry + }); break; - case CAMU_BUFFER_CONTINUE: - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_stream_continue((struct bmu_local_stream *)entry->video.stream); - break; - default: - break; - } + case CAMU_BUFFER_UNCORK: + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = UNCORK, + .opaque = entry + }); break; case CAMU_BUFFER_EOF: queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = STOP, - .value.i = CAMU_SINK_VIDEO, - .opaque = entry->sink + .value.i = CAMU_SINK_VIDEO }); break; } @@ -349,18 +303,39 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str { 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); + break; + } + // A lot of logic here assumes that no data will be sent + // until all active streams are configured. This is extremely important, + // if it does not hold true many confusing errors will arise. case BIMU_CLIENT_CONFIGURE: { switch (stream->type) { case BIMU_STREAM_AUDIO: entry->audio.stream = stream; camu_audio_buffer_configure(&entry->audio.buf, &stream->stream); - entry->audio.state = BUFFER_CONFIGURED; + aki_mutex_lock(&entry->sink->mutex); + if (entry->audio.state == BUFFER_QUEUED) { + entry->audio.state = BUFFER_SET_OR_BUFFERED; + } else { + entry->audio.state = BUFFER_CONFIGURED; + } + aki_mutex_unlock(&entry->sink->mutex); break; #ifndef CAMU_SINK_NO_VIDEO case BIMU_STREAM_VIDEO: entry->video.stream = stream; camu_video_buffer_configure(&entry->video.buf, &stream->stream); - entry->video.state = BUFFER_CONFIGURED; + aki_mutex_lock(&entry->sink->mutex); + if (entry->video.state == BUFFER_QUEUED) { + entry->video.state = BUFFER_SET_OR_BUFFERED; + } else { + entry->video.state = BUFFER_CONFIGURED; + } + aki_mutex_unlock(&entry->sink->mutex); break; #endif } @@ -384,11 +359,14 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str } #endif al_free(frame); + break; } break; } case BIMU_CLIENT_SEEK: { - f64 seek_req = *(f64 *)opaque; + // Must ensure no more data from the previous + // stream position is sent during or after this call to seek. + struct bmu_seek_req *req = (struct bmu_seek_req *)opaque; aki_mutex_lock(&entry->sink->mutex); if (entry->audio.state == BUFFER_ADDED) { entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); @@ -415,7 +393,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str #ifndef CAMU_SINK_NO_VIDEO } #endif - camu_clock_seek(&entry->clock, seek_req); + camu_clock_seek(&entry->clock, req->base, req->start); camu_audio_buffer_reset(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO if (!single_frame) { @@ -443,32 +421,12 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str break; } case BIMU_CLIENT_CLOSED: { - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_close(&((struct camu_sink_local *)entry)->runner); - break; - default: - break; - } + bmu_client_free(&entry->client); camu_audio_buffer_free(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO camu_video_buffer_free(&entry->video.buf); #endif - switch (entry->type) { - case CAMU_SINK_LOCAL: - cch_entry_free(&((struct camu_sink_local *)entry)->entry); - break; - default: - break; - } - al_str_free(&entry->unique_id); - switch (entry->type) { - case CAMU_SINK_LOCAL: - al_free((struct camu_sink_local *)entry); - break; - default: - break; - } + al_free(entry); al_log_debug("sink", "Entry closed."); break; } @@ -498,37 +456,160 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, return true; } -bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port) +static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id) { - (void)sink; - (void)addr; - (void)port; + struct camu_sink_entry *entry; + al_array_foreach(sink->entries, i, entry) { + if (entry->client.node_id == node_id) return entry; + } + return NULL; +} + +static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)conn; + (void)rpacket; + + str addr; + aki_packet_read_str(packet, &addr); + s32 port = aki_packet_read_s32(packet); + u16 node_id = aki_packet_read_u16(packet); + + struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry); + entry->sink = sink; + al_array_push(sink->entries, entry); + + entry->state = ENTRY_LOADED; + + entry->audio.state = BUFFER_INIT; + camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer); + entry->audio.buf.callback = audio_buffer_callback; + entry->audio.buf.userdata = entry; + + struct camu_renderer *renderer = NULL; +#ifndef CAMU_SINK_NO_VIDEO + renderer = sink->video.renderer; + entry->video.state = BUFFER_INIT; + camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer); + entry->video.buf.callback = video_buffer_callback; + entry->video.buf.userdata = entry; +#endif + + 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; } +static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry) +{ +#ifndef CAMU_SINK_NO_VIDEO + if (entry->video.state == BUFFER_ADDED) { + sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); + entry->video.state = BUFFER_SET_OR_BUFFERED; + } +#endif + if (entry->audio.state == BUFFER_ADDED) { + sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + entry->audio.state = BUFFER_SET_OR_BUFFERED; + } +} + +static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)conn; + (void)rpacket; + + 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); + struct camu_sink_entry *previous = sink->current; + sink->current = entry; + if (entry->audio.state == BUFFER_INIT) { + entry->audio.state = BUFFER_QUEUED; + } else { + add_audio_if_set_and_buffered(entry); + } +#ifndef CAMU_SINK_NO_VIDEO + if (entry->video.state == BUFFER_INIT) { + entry->video.state = BUFFER_QUEUED; + } else { + add_video_if_set_and_buffered(entry); + } +#endif + if (previous) { + camu_clock_pause(&previous->clock); + remove_entry_buffers(sink, previous); + } + aki_mutex_unlock(&sink->mutex); + + aki_packet_free(packet); + return false; +} + +static struct aki_rpc_command commands[] = { + { .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL }, + { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL } +}; + +static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + sink->conn = conn; + struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY); + aki_packet_write_u8(packet, TREE_SINK); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); +} + +static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) +{ + (void)userdata; + (void)conn; +} + +bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port) +{ + aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback, + connection_closed_callback, sink); + for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { + commands[i].userdata = sink; + al_assert(sink->callback); + aki_rpc_add_command(&sink->client, &commands[i]); + } + return aki_rpc_connect(&sink->client, sink->loop, addr, port); +} + void camu_sink_toggle_pause(struct camu_sink *sink) { aki_mutex_lock(&sink->mutex); - if (sink->current) { + struct camu_sink_entry *current = sink->current; + aki_mutex_unlock(&sink->mutex); + if (current) { queue_cmd(sink, (struct camu_sink_cmd){ .op = TOGGLE_PAUSE, - .opaque = sink->current + .opaque = current }); } - aki_mutex_unlock(&sink->mutex); } void camu_sink_seek(struct camu_sink *sink, f64 pos) { aki_mutex_lock(&sink->mutex); - if (sink->current) { + struct camu_sink_entry *current = sink->current; + aki_mutex_unlock(&sink->mutex); + if (current) { queue_cmd(sink, (struct camu_sink_cmd){ .op = SEEK, .value.f = pos, - .opaque = sink->current + .opaque = current }); } - aki_mutex_unlock(&sink->mutex); } void camu_sink_skip(struct camu_sink *sink, s32 n) @@ -548,20 +629,6 @@ void camu_sink_return_current(struct camu_sink *sink) aki_mutex_unlock(&sink->mutex); } -static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry) -{ -#ifndef CAMU_SINK_NO_VIDEO - if (entry->video.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - entry->video.state = BUFFER_SET_OR_BUFFERED; - } -#endif - if (entry->audio.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - entry->audio.state = BUFFER_SET_OR_BUFFERED; - } -} - void camu_sink_stop(struct camu_sink *sink) { aki_mutex_lock(&sink->mutex); @@ -572,8 +639,7 @@ void camu_sink_stop(struct camu_sink *sink) aki_mutex_unlock(&sink->mutex); queue_cmd(sink, (struct camu_sink_cmd){ .op = STOP, - .value.i = CAMU_SINK_AUDIO, - .opaque = sink + .value.i = CAMU_SINK_AUDIO }); queue_cmd(sink, (struct camu_sink_cmd){ .op = CLOSE @@ -583,150 +649,14 @@ 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) { - switch (entry->type) { - case CAMU_SINK_LOCAL: - bmu_local_stop(&((struct camu_sink_local *)entry)->runner); - break; - default: - break; - } } + */ al_array_free(sink->entries); aki_mutex_unlock(&sink->mutex); aki_signal_stop(&sink->signal); camu_queue_free(sink->queue); aki_mutex_destroy(&sink->mutex); } - -// Local bimu band-aid compat. - -static struct camu_sink_local *entry_from_unique_id(struct camu_sink *sink, str *unique_id) -{ - struct camu_sink_entry *entry; - al_array_foreach(sink->entries, i, entry) { - if (al_str_eq(&entry->unique_id, unique_id)) { - return (struct camu_sink_local *)entry; - } - } - return NULL; -} - -static void camu_sink_buffer_internal(struct camu_sink *sink, struct camu_sink_entry *entry, str *unique_id) -{ - entry->sink = sink; - al_str_clone(&entry->unique_id, unique_id); - - camu_clock_init(&entry->clock); - - entry->audio.state = BUFFER_INIT; - camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer); - entry->audio.buf.callback = audio_buffer_callback; - entry->audio.buf.userdata = entry; - - struct camu_renderer *renderer = NULL; -#ifndef CAMU_SINK_NO_VIDEO - renderer = sink->video.renderer; - entry->video.state = BUFFER_INIT; - camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer); - entry->video.buf.callback = video_buffer_callback; - entry->video.buf.userdata = entry; -#endif - - al_array_push(sink->entries, entry); -} - -bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *centry) -{ - aki_mutex_lock(&sink->mutex); - - bool repeat = true; - struct camu_sink_local *local = entry_from_unique_id(sink, unique_id); - if (!local) { - local = al_alloc_object(struct camu_sink_local); - local->e.type = CAMU_SINK_LOCAL; - repeat = false; - } - local->e.state = ENTRY_LOADED; - if (repeat) { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = SEEK, - .value.f = 0, - .opaque = &local->e - }); - aki_mutex_unlock(&sink->mutex); - return true; - } - - local->entry = centry; - if (!bmu_local_init(&local->runner, centry)) { - al_free(local); - aki_mutex_unlock(&sink->mutex); - return false; - } - - camu_sink_buffer_internal(sink, &local->e, unique_id); - - struct camu_renderer *renderer = NULL; -#ifndef CAMU_SINK_NO_VIDEO - renderer = sink->video.renderer; -#endif - local->runner.callback = client_callback; - local->runner.userdata = &local->e; - bmu_local_prepare_clients(&local->runner, renderer); - bmu_local_run(&local->runner); - - aki_mutex_unlock(&sink->mutex); - - return true; -} - -static void maybe_cleanup_local_entries(struct camu_sink *sink) -{ - struct camu_sink_entry *entry; - al_array_foreach(sink->entries, i, entry) { - if (entry->state == ENTRY_DISREGUARDED && !ENTRY_BUFFERS_HELD(entry)) { - bmu_local_stop(&((struct camu_sink_local *)entry)->runner); - al_array_remove_at_iter(sink->entries, i); - } - } -} - -void camu_sink_local_set(struct camu_sink *sink, str *unique_id) -{ - aki_mutex_lock(&sink->mutex); - struct camu_sink_local *local = entry_from_unique_id(sink, unique_id); - struct camu_sink_entry *previous = sink->current; - sink->current = &local->e; - add_audio_if_set_and_buffered(&local->e); -#ifndef CAMU_SINK_NO_VIDEO - add_video_if_set_and_buffered(&local->e); -#endif - if (previous) { - if (!camu_clock_is_paused(&previous->clock)) { - camu_clock_pause(&previous->clock); - } - remove_entry_buffers(sink, previous); - } - maybe_cleanup_local_entries(sink); - aki_mutex_unlock(&sink->mutex); -} - -void camu_sink_local_swap(struct camu_sink *sink, str *unique_id) -{ - struct camu_sink_local *local = entry_from_unique_id(sink, unique_id); - sink->current = &local->e; - add_audio_if_set_and_buffered(&local->e); -#ifndef CAMU_SINK_NO_VIDEO - add_video_if_set_and_buffered(&local->e); -#endif -} - -void camu_sink_local_unload(struct camu_sink *sink, str *unique_id) -{ - aki_mutex_lock(&sink->mutex); - struct camu_sink_local *local = entry_from_unique_id(sink, unique_id); - if (local && local->e.state == ENTRY_LOADED) local->e.state = ENTRY_DISREGUARDED; - aki_mutex_unlock(&sink->mutex); -} diff --git a/src/libsink/sink.h b/src/libsink/sink.h index efef97b..95c280f 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -4,20 +4,21 @@ #include #include +#include #include #include #include "../util/queue.h" -#include "../cache/entry.h" - #include "../buffer/clock.h" #include "../buffer/audio.h" #ifndef CAMU_SINK_NO_VIDEO #include "../buffer/video.h" #endif -#include "../bimu/local.h" +#include "../bimu/client.h" + +#include "common.h" enum { CAMU_SINK_AUDIO = 0, @@ -26,25 +27,19 @@ enum { #endif }; -enum { - CAMU_SINK_REMOTE = 0, - CAMU_SINK_LOCAL -}; - enum { CAMU_SINK_ADD_BUFFER = 0, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_SWAP_BUFFER, + CAMU_SINK_SET_BUFFERED, CAMU_SINK_START, CAMU_SINK_STOP, - CAMU_SINK_ENTRY_BUFFERED, CAMU_SINK_EXIT }; struct camu_sink_entry { - u8 type; u8 state; - str unique_id; + struct bmu_client client; struct camu_clock clock; struct { u8 state; @@ -61,12 +56,6 @@ struct camu_sink_entry { struct camu_sink *sink; }; -struct camu_sink_local { - struct camu_sink_entry e; - struct cch_entry *entry; - struct bmu_local runner; -}; - struct camu_sink_cmd { u8 op; union { s64 i; f64 f; } value; @@ -80,6 +69,8 @@ enum { struct camu_sink { struct aki_event_loop *loop; + struct aki_rpc client; + struct aki_rpc_connection *conn; struct aki_signal signal; queue(struct camu_sink_cmd) queue; struct aki_mutex mutex; @@ -113,8 +104,3 @@ struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink); void camu_sink_return_current(struct camu_sink *sink); void camu_sink_stop(struct camu_sink *sink); void camu_sink_close(struct camu_sink *sink); - -bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *entry); -void camu_sink_local_set(struct camu_sink *sink, str *unique_id); -void camu_sink_local_swap(struct camu_sink *sink, str *unique_id); -void camu_sink_local_unload(struct camu_sink *sink, str *unique_id); diff --git a/src/libsink/sink2.c b/src/libsink/sink2.c deleted file mode 100644 index e9b51f2..0000000 --- a/src/libsink/sink2.c +++ /dev/null @@ -1,688 +0,0 @@ -#include - -#include "../tree/common.h" - -#include "sink2.h" - -enum { - SINK_EMPTY = 0, - SINK_PAUSED, - SINK_PLAYING -}; - -enum { - ENTRY_LOADED = 0, - ENTRY_DISREGUARDED -}; - -enum { - BUFFER_INIT = 0, - BUFFER_QUEUED, - BUFFER_CONFIGURED, - BUFFER_SET_OR_BUFFERED, - BUFFER_ADDED, -}; - -enum { - // TODO: Add and remove buffer can be done from - // any thread. Does this always need to be the case? - // Is it beneficial for this to be the case? - ADD_BUFFER = 0, - REMOVE_BUFFER, - START, - STOP, - TOGGLE_PAUSE, - SEEK, - ENTRY_BUFFERED, // Currently set entry is buffered. - CLOSE, - // Internal. - CORK, - UNCORK -}; - -#define ENTRY_AUDIO_BUFFER_HELD(entry) \ - (al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED)) -#ifdef CAMU_SINK_NO_VIDEO -#define ENTRY_BUFFERS_HELD(entry) ENTRY_AUDIO_BUFFER_HELD(entry) -#else -#define ENTRY_VIDEO_BUFFER_HELD(entry) \ - (al_atomic_bool_load(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED)) -#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry)) -#endif - -#define ENTRY_VIDEO_EMPTY(entry) \ - (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED) -#define ENTRY_AUDIO_EMPTY(entry) \ - (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED) - -static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) -{ - switch (cmd->op) { - case ADD_BUFFER: { - 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_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - break; -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: - sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - break; -#endif - } - break; - } - case REMOVE_BUFFER: { - 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_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - break; -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - break; -#endif - } - break; - } - case START: { - switch (cmd->value.i) { - case CAMU_SINK_AUDIO: - if (sink->audio.state == SINK_PAUSED) { - sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); - sink->audio.state = SINK_PLAYING; - } - break; -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: - if (sink->video.state == SINK_PAUSED) { - sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL); - sink->video.state = SINK_PLAYING; - } - break; -#endif - } - break; - } - case STOP: { - switch (cmd->value.i) { - case CAMU_SINK_AUDIO: - if (sink->audio.state == SINK_PLAYING) { - sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_AUDIO, NULL); - sink->audio.state = SINK_PAUSED; - } - break; -#ifndef CAMU_SINK_NO_VIDEO - case CAMU_SINK_VIDEO: - if (sink->video.state == SINK_PLAYING) { - sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL); - sink->video.state = SINK_PAUSED; - } - break; -#endif - } - break; - } - case TOGGLE_PAUSE: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - if (camu_clock_is_paused(&entry->clock)) { - camu_clock_resume(&entry->clock); - if (sink->audio.state == SINK_PAUSED) { - sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); - sink->audio.state = SINK_PLAYING; - } -#ifndef CAMU_SINK_NO_VIDEO - if (sink->video.state == SINK_PAUSED) { - sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL); - sink->video.state = SINK_PLAYING; - } -#endif - } else { - camu_clock_pause(&entry->clock); -#ifndef CAMU_SINK_NO_VIDEO - if (sink->video.state == SINK_PLAYING) { - sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL); - sink->video.state = SINK_PAUSED; - } -#endif - } - break; - } - case SEEK: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - (void)entry; - break; - } - case ENTRY_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; - 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); - break; - } - } -} - -static void queue_signal_callback(void *userdata) -{ - struct camu_sink *sink = (struct camu_sink *)userdata; - u32 size; - struct camu_sink_cmd cmd; - do { - camu_queue_try_pop(sink->queue, size, cmd); - if (size == 0) break; - handle_sink_cmd(sink, &cmd); - } while (1); -} - -static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd) -{ - camu_queue_push(sink->queue, cmd); - aki_signal_send(&sink->signal); -} - -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 - }); - 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) { - camu_clock_resume(&entry->clock); - } - state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } - entry->audio.state = state; -} - -static void audio_buffer_callback(void *userdata, u8 op) -{ - struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; - switch (op) { - case CAMU_BUFFER_BUFFERED: - aki_mutex_lock(&entry->sink->mutex); - add_audio_if_set_and_buffered(entry); - aki_mutex_unlock(&entry->sink->mutex); - break; - case CAMU_BUFFER_CORK: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = CORK, - .opaque = entry - }); - break; - case CAMU_BUFFER_UNCORK: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = UNCORK, - .opaque = entry - }); - break; - case CAMU_BUFFER_PAUSED: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = STOP, - .value.i = CAMU_SINK_AUDIO - }); - break; - case CAMU_BUFFER_EOF: { - aki_mutex_lock(&entry->sink->mutex); - u8 ret = entry->sink->callback(entry->sink->userdata, CAMU_SINK_SWAP_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - entry->audio.state = BUFFER_SET_OR_BUFFERED; - aki_mutex_unlock(&entry->sink->mutex); - if (ret != CAMU_SINK_BUFFERS_SWAPPED) { - // TODO: Can cause popping. Should run a timer - // and if any action would queue_cmd(AUDIO_START) before the - // the timer, don't queue the stop. - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = STOP, - .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 - }); -#endif - } - break; - } - } -} - -#ifndef CAMU_SINK_NO_VIDEO -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 - }); - 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) { - camu_clock_resume(&entry->clock); - } - state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } - entry->video.state = state; -} - -static void video_buffer_callback(void *userdata, u8 op) -{ - struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; - switch (op) { - case CAMU_BUFFER_BUFFERED: - aki_mutex_lock(&entry->sink->mutex); - add_video_if_set_and_buffered(entry); - aki_mutex_unlock(&entry->sink->mutex); - break; - case CAMU_BUFFER_CORK: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = CORK, - .opaque = entry - }); - break; - case CAMU_BUFFER_UNCORK: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = UNCORK, - .opaque = entry - }); - break; - case CAMU_BUFFER_EOF: - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = STOP, - .value.i = CAMU_SINK_VIDEO - }); - break; - } -} -#endif - -static void client_callback(void *userdata, u8 op, struct bmu_client_stream *stream, void *opaque) -{ - struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; - switch (op) { - // A lot of logic here assumes that no data will be sent - // until all active streams are configured. This is extremely important, - // if it does not hold true many confusing errors will arise. - case BIMU_CLIENT_CONFIGURE: { - switch (stream->type) { - case BIMU_STREAM_AUDIO: - entry->audio.stream = stream; - camu_audio_buffer_configure(&entry->audio.buf, &stream->stream); - aki_mutex_lock(&entry->sink->mutex); - if (entry->audio.state == BUFFER_QUEUED) { - entry->audio.state = BUFFER_SET_OR_BUFFERED; - } else { - entry->audio.state = BUFFER_CONFIGURED; - } - aki_mutex_unlock(&entry->sink->mutex); - break; -#ifndef CAMU_SINK_NO_VIDEO - case BIMU_STREAM_VIDEO: - entry->video.stream = stream; - camu_video_buffer_configure(&entry->video.buf, &stream->stream); - aki_mutex_lock(&entry->sink->mutex); - if (entry->video.state == BUFFER_QUEUED) { - entry->video.state = BUFFER_SET_OR_BUFFERED; - } else { - entry->video.state = BUFFER_CONFIGURED; - } - aki_mutex_unlock(&entry->sink->mutex); - break; -#endif - } - break; - } - case BIMU_CLIENT_DATA: { - struct camu_frame *frame = (struct camu_frame *)opaque; - switch (stream->type) { - case BIMU_STREAM_AUDIO: - camu_audio_buffer_push(&entry->audio.buf, frame); - break; -#ifndef CAMU_SINK_NO_VIDEO - case BIMU_STREAM_VIDEO: - camu_video_buffer_push(&entry->video.buf, frame); - break; -#endif - default: -#ifdef HAVE_FFMPEG - if (frame->type == CAMU_FFMPEG_COMPAT) { - av_frame_free(&frame->av.frame); - } -#endif - al_free(frame); - break; - } - break; - } - case BIMU_CLIENT_SEEK: { - // Must ensure no more data from the previous - // stream position is sent during or after this call to seek. - f64 seek_pos = *(f64 *)opaque; - aki_mutex_lock(&entry->sink->mutex); - if (entry->audio.state == BUFFER_ADDED) { - entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - entry->audio.state = BUFFER_SET_OR_BUFFERED; - } -#ifndef CAMU_SINK_NO_VIDEO - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (entry->video.state == BUFFER_ADDED && !single_frame) { - entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - entry->video.state = BUFFER_SET_OR_BUFFERED; - } -#endif - aki_mutex_unlock(&entry->sink->mutex); -#ifndef CAMU_SINK_NO_VIDEO - if (single_frame) { - while (ENTRY_AUDIO_BUFFER_HELD(entry)) { - aki_thread_sleep(AKI_TS_FROM_USEC(200)); - } - } else { -#endif - while (ENTRY_BUFFERS_HELD(entry)) { - aki_thread_sleep(AKI_TS_FROM_USEC(200)); - } -#ifndef CAMU_SINK_NO_VIDEO - } -#endif - camu_clock_seek(&entry->clock, seek_pos); - camu_audio_buffer_reset(&entry->audio.buf); -#ifndef CAMU_SINK_NO_VIDEO - if (!single_frame) { - camu_video_buffer_reset(&entry->video.buf); - } -#endif - break; - } - case BIMU_CLIENT_EOF: { - switch (stream->type) { - case BIMU_STREAM_AUDIO: { - camu_audio_buffer_flush(&entry->audio.buf); - break; - } -#ifndef CAMU_SINK_NO_VIDEO - case BIMU_STREAM_VIDEO: { - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (!single_frame) { - camu_video_buffer_flush(&entry->video.buf); - } - break; - } -#endif - } - break; - } - case BIMU_CLIENT_CLOSED: { - // close/free client. - camu_audio_buffer_free(&entry->audio.buf); -#ifndef CAMU_SINK_NO_VIDEO - camu_video_buffer_free(&entry->video.buf); -#endif - al_free(entry); - al_log_debug("sink", "Entry closed."); - break; - } - } -} - -bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, - struct camu_mixer *mixer -#ifndef CAMU_SINK_NO_VIDEO - , struct camu_renderer *renderer -#endif - ) -{ - sink->loop = loop; - aki_signal_init(&sink->signal, queue_signal_callback, sink); - aki_signal_start(&sink->signal, sink->loop); - camu_queue_init(sink->queue); - aki_mutex_init(&sink->mutex); - sink->current = NULL; - al_array_init(sink->entries); - sink->audio.mixer = mixer; - sink->audio.state = SINK_PAUSED; -#ifndef CAMU_SINK_NO_VIDEO - sink->video.renderer = renderer; - sink->video.state = SINK_PAUSED; -#endif - return true; -} - -static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id) -{ - struct camu_sink_entry *entry; - al_array_foreach(sink->entries, i, entry) { - if (entry->client.node_id == node_id) return entry; - } - return NULL; -} - -static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) -{ - struct camu_sink *sink = (struct camu_sink *)userdata; - (void)conn; - (void)rpacket; - - str addr; - aki_packet_read_str(packet, &addr); - s32 port = aki_packet_read_s32(packet); - u16 node_id = aki_packet_read_u16(packet); - - struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry); - entry->sink = sink; - al_array_push(sink->entries, entry); - - entry->state = ENTRY_LOADED; - camu_clock_init(&entry->clock); - - entry->audio.state = BUFFER_INIT; - camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer); - entry->audio.buf.callback = audio_buffer_callback; - entry->audio.buf.userdata = entry; - - struct camu_renderer *renderer = NULL; -#ifndef CAMU_SINK_NO_VIDEO - renderer = sink->video.renderer; - entry->video.state = BUFFER_INIT; - camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer); - entry->video.buf.callback = video_buffer_callback; - entry->video.buf.userdata = entry; -#endif - - 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; -} - -static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry) -{ -#ifndef CAMU_SINK_NO_VIDEO - if (entry->video.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - entry->video.state = BUFFER_SET_OR_BUFFERED; - } -#endif - if (entry->audio.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - entry->audio.state = BUFFER_SET_OR_BUFFERED; - } -} - -static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn, - struct aki_packet *packet, struct aki_packet *rpacket) -{ - struct camu_sink *sink = (struct camu_sink *)userdata; - (void)conn; - (void)rpacket; - - 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); - struct camu_sink_entry *previous = sink->current; - sink->current = entry; - if (entry->audio.state == BUFFER_INIT) { - entry->audio.state = BUFFER_QUEUED; - } else { - add_audio_if_set_and_buffered(entry); - } -#ifndef CAMU_SINK_NO_VIDEO - if (entry->video.state == BUFFER_INIT) { - entry->video.state = BUFFER_QUEUED; - } else { - add_video_if_set_and_buffered(entry); - } -#endif - if (previous) { - camu_clock_pause(&previous->clock); - remove_entry_buffers(sink, previous); - } - aki_mutex_unlock(&sink->mutex); - - aki_packet_free(packet); - return false; -} - -static struct aki_rpc_command commands[] = { - { .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL }, - { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL } -}; - -static void connection_callback(void *userdata, struct aki_rpc_connection *conn) -{ - struct camu_sink *sink = (struct camu_sink *)userdata; - sink->conn = conn; - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY); - aki_packet_write_u8(packet, TREE_SINK); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); -} - -static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) -{ - (void)userdata; - (void)conn; -} - -bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port) -{ - aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback, - connection_closed_callback, sink); - for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { - commands[i].userdata = sink; - al_assert(sink->callback); - aki_rpc_add_command(&sink->client, &commands[i]); - } - return aki_rpc_connect(&sink->client, sink->loop, addr, port); -} - -void camu_sink_toggle_pause(struct camu_sink *sink) -{ - aki_mutex_lock(&sink->mutex); - if (sink->current) { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = TOGGLE_PAUSE, - .opaque = sink->current - }); - } - aki_mutex_unlock(&sink->mutex); -} - -void camu_sink_seek(struct camu_sink *sink, f64 pos) -{ - aki_mutex_lock(&sink->mutex); - if (sink->current) { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = SEEK, - .value.f = pos, - .opaque = sink->current - }); - } - aki_mutex_unlock(&sink->mutex); -} - -void camu_sink_skip(struct camu_sink *sink, s32 n) -{ - (void)sink; - (void)n; -} - -struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink) -{ - aki_mutex_lock(&sink->mutex); - return sink->current; -} - -void camu_sink_return_current(struct camu_sink *sink) -{ - aki_mutex_unlock(&sink->mutex); -} - -void camu_sink_stop(struct camu_sink *sink) -{ - aki_mutex_lock(&sink->mutex); - if (sink->current) { - remove_entry_buffers(sink, sink->current); - sink->current = NULL; - } - aki_mutex_unlock(&sink->mutex); - queue_cmd(sink, (struct camu_sink_cmd){ - .op = STOP, - .value.i = CAMU_SINK_AUDIO - }); - queue_cmd(sink, (struct camu_sink_cmd){ - .op = CLOSE - }); -} - -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) { - } - */ - al_array_free(sink->entries); - aki_mutex_unlock(&sink->mutex); - aki_signal_stop(&sink->signal); - camu_queue_free(sink->queue); - aki_mutex_destroy(&sink->mutex); -} diff --git a/src/libsink/sink2.h b/src/libsink/sink2.h deleted file mode 100644 index 95c280f..0000000 --- a/src/libsink/sink2.h +++ /dev/null @@ -1,106 +0,0 @@ -#pragma once - -//#define CAMU_SINK_NO_VIDEO - -#include -#include -#include -#include -#include - -#include "../util/queue.h" - -#include "../buffer/clock.h" -#include "../buffer/audio.h" -#ifndef CAMU_SINK_NO_VIDEO -#include "../buffer/video.h" -#endif - -#include "../bimu/client.h" - -#include "common.h" - -enum { - CAMU_SINK_AUDIO = 0, -#ifndef CAMU_SINK_NO_VIDEO - CAMU_SINK_VIDEO -#endif -}; - -enum { - CAMU_SINK_ADD_BUFFER = 0, - CAMU_SINK_REMOVE_BUFFER, - CAMU_SINK_SWAP_BUFFER, - CAMU_SINK_SET_BUFFERED, - CAMU_SINK_START, - CAMU_SINK_STOP, - CAMU_SINK_EXIT -}; - -struct camu_sink_entry { - u8 state; - struct bmu_client client; - struct camu_clock clock; - struct { - u8 state; - struct camu_audio_buffer buf; - struct bmu_client_stream *stream; - } audio; -#ifndef CAMU_SINK_NO_VIDEO - struct { - u8 state; - struct camu_video_buffer buf; - struct bmu_client_stream *stream; - } video; -#endif - struct camu_sink *sink; -}; - -struct camu_sink_cmd { - u8 op; - union { s64 i; f64 f; } value; - void *opaque; -}; - -enum { - CAMU_SINK_OK = 0, - CAMU_SINK_BUFFERS_SWAPPED = 1 -}; - -struct camu_sink { - struct aki_event_loop *loop; - struct aki_rpc client; - struct aki_rpc_connection *conn; - struct aki_signal signal; - queue(struct camu_sink_cmd) queue; - struct aki_mutex mutex; - struct camu_sink_entry *current; - array(struct camu_sink_entry *) entries; - struct { - u8 state; - struct camu_mixer *mixer; - } audio; -#ifndef CAMU_SINK_NO_VIDEO - struct { - u8 state; - struct camu_renderer *renderer; - } video; -#endif - u8 (*callback)(void *, u8, u8, void *); - void *userdata; -}; - -bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, - struct camu_mixer *mixer -#ifndef CAMU_SINK_NO_VIDEO - , struct camu_renderer *renderer -#endif - ); -bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port); -void camu_sink_toggle_pause(struct camu_sink *sink); -void camu_sink_seek(struct camu_sink *sink, f64 pos); -void camu_sink_skip(struct camu_sink *sink, s32 n); -struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink); -void camu_sink_return_current(struct camu_sink *sink); -void camu_sink_stop(struct camu_sink *sink); -void camu_sink_close(struct camu_sink *sink); -- cgit v1.2.3-101-g0448