From b100eb175e0cfa59a15c4125f26ab5474c2347ec Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Sat, 27 Jan 2024 20:39:25 -0500 Subject: wip - Mostly work in the client Signed-off-by: Andrew Opalach --- src/libsink/common.h | 3 +- src/libsink/sink.c | 247 ++++++++++++++++++++++++++++++++++----------------- src/libsink/sink.h | 2 + 3 files changed, 169 insertions(+), 83 deletions(-) (limited to 'src/libsink') diff --git a/src/libsink/common.h b/src/libsink/common.h index 1c701c8..48e43e8 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -2,5 +2,6 @@ enum { CAMU_SINK_CMD_BUFFER = 0, - CAMU_SINK_CMD_SET + CAMU_SINK_CMD_SET, + CAMU_SINK_CMD_QUEUE, }; diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 3eecf5f..134be31 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -1,9 +1,12 @@ #include #include "../tree/common.h" + #include "../bimu/common.h" #include "../bimu/handler.h" +#include "../buffer/common.h" + #include "sink.h" #include "common.h" @@ -31,15 +34,18 @@ enum { START, STOP, TOGGLE_PAUSE, + SKIP, RESEEK, SEEK, SET_BUFFERED, // Currently set entry is buffered. CLOSE, // Internal. CORK, - UNCORK + UNCORK, }; +#define BIMU_DELAY_IGNORE 0 + #define ENTRY_AUDIO_BUFFER_HELD(entry) \ (al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED)) #ifdef CAMU_SINK_NO_VIDEO @@ -55,6 +61,36 @@ enum { #define ENTRY_AUDIO_EMPTY(entry) \ (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED) +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 void add_audio_if_set_and_buffered(struct camu_sink_entry *entry); +#ifndef CAMU_SINK_NO_VIDEO +static void add_video_if_set_and_buffered(struct camu_sink_entry *entry); +#endif + +static void swap_buffers_internal(struct camu_sink *sink) +{ + remove_entry_buffers(sink, sink->current); + sink->current = sink->queued; + sink->queued = NULL; + add_audio_if_set_and_buffered(sink->current); +#ifndef CAMU_SINK_NO_VIDEO + add_video_if_set_and_buffered(sink->current); +#endif +} + static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) { switch (cmd->op) { @@ -122,6 +158,12 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) } break; } + case SKIP: { + struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_SKIP); + aki_packet_write_s32(packet, cmd->value.i); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } case RESEEK: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; bmu_client_reseek(&entry->client); @@ -146,6 +188,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) // break; // } case CLOSE: { + aki_signal_stop(&sink->signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); return; } @@ -198,11 +241,16 @@ static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd) aki_signal_send(&sink->signal); } -static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) +void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->audio.state; if (state == BUFFER_SET_OR_BUFFERED) { + // TODO: should_resume is not robust. +#ifndef CAMU_SINK_NO_VIDEO bool should_resume = ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED; +#else + bool should_resume = true; +#endif if (should_resume && !camu_clock_calc_tick(&entry->clock)) { queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = RESEEK, @@ -230,6 +278,40 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) entry->audio.state = state; } +#ifndef CAMU_SINK_NO_VIDEO +void add_video_if_set_and_buffered(struct camu_sink_entry *entry) +{ + u8 state = entry->video.state; + if (state == BUFFER_SET_OR_BUFFERED) { + 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 (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; + } + entry->video.state = state; +} +#endif + static void audio_buffer_callback(void *userdata, u8 op) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; @@ -260,15 +342,16 @@ static void audio_buffer_callback(void *userdata, u8 op) }); 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) { - if (1) { - // 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. + struct camu_sink *sink = entry->sink; + bool swapped = false; + aki_mutex_lock(&sink->mutex); + if (sink->queued) { + swap_buffers_internal(sink); + swapped = true; + al_log_info("sink", "Audio buffers swapped (gapless)."); + } + aki_mutex_unlock(&sink->mutex); + if (!swapped) { queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = STOP, .value.i = CAMU_SINK_AUDIO @@ -286,38 +369,6 @@ static void audio_buffer_callback(void *userdata, u8 op) } #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) { - 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 (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; - } - entry->video.state = state; -} - static void video_buffer_callback(void *userdata, u8 op) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; @@ -388,7 +439,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream 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->delay = BIMU_DELAY_IGNORE; req->ts -= req->delay; bool late = ts <= req->ts; if (late) {} @@ -502,6 +553,7 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, aki_signal_start(&sink->signal, sink->loop); camu_queue_init(sink->queue); aki_mutex_init(&sink->mutex); + sink->queued = NULL; sink->current = NULL; al_array_init(sink->entries); sink->audio.mixer = mixer; @@ -526,27 +578,19 @@ 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_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) +static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *sink, str *addr, s32 port, u32 node_id) { - 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 = entry_from_node_id(sink, node_id); + if (entry) return entry; - struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry); - entry->sink = sink; + entry = al_alloc_object(struct camu_sink_entry); al_array_push(sink->entries, entry); - + entry->sink = sink; entry->state = ENTRY_LOADED; entry->audio.state = BUFFER_INIT; @@ -567,26 +611,32 @@ static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *c entry->client.callback = client_callback; entry->client.userdata = entry; - bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id); + bmu_client_connect(&entry->client, sink->loop, addr, port, node_id); + + return entry; +} + +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); + + aki_mutex_lock(&sink->mutex); + ensure_entry_buffered_internal(sink, &addr, port, node_id); + aki_mutex_unlock(&sink->mutex); 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) @@ -595,12 +645,17 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn (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); + 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 *entry = ensure_entry_buffered_internal(sink, &addr, port, node_id); + if (entry == sink->current) goto out; struct camu_sink_entry *previous = sink->current; sink->current = entry; + if (sink->queued) sink->queued = NULL; if (entry->audio.state == BUFFER_INIT) { entry->audio.state = BUFFER_QUEUED; } else { @@ -617,16 +672,35 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn camu_clock_pause(&previous->clock); remove_entry_buffers(sink, previous); } - out: aki_mutex_unlock(&sink->mutex); aki_packet_free(packet); return false; } +static bool queue_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); + + aki_mutex_lock(&sink->mutex); + sink->queued = ensure_entry_buffered_internal(sink, &addr, port, node_id); + aki_mutex_unlock(&sink->mutex); + + 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 } + { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL }, + { .op = CAMU_SINK_CMD_QUEUE, .callback = queue_command_callback, .userdata = NULL } }; static void connection_callback(void *userdata, struct aki_rpc_connection *conn) @@ -640,7 +714,8 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn) static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { - (void)userdata; + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)sink; (void)conn; } @@ -685,8 +760,10 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent) void camu_sink_skip(struct camu_sink *sink, s32 n) { - (void)sink; - (void)n; + queue_cmd(sink, (struct camu_sink_cmd){ + .op = SKIP, + .value.i = n + }); } struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink) @@ -719,15 +796,21 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { + // TODO: Make sure no commands can come in after this. + aki_rpc_disconnect(&sink->client); aki_mutex_lock(&sink->mutex); struct camu_sink_entry *entry; al_array_foreach(sink->entries, i, entry) { bmu_client_close(&entry->client); + al_array_remove_at_iter(sink->entries, i); } aki_mutex_unlock(&sink->mutex); +} + +void camu_sink_free(struct camu_sink *sink) +{ al_array_free(sink->entries); - aki_rpc_disconnect(&sink->client); - aki_signal_stop(&sink->signal); + aki_rpc_free(&sink->client); camu_queue_free(sink->queue); aki_mutex_destroy(&sink->mutex); } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index e3550bb..40d20d4 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -73,6 +73,7 @@ struct camu_sink { struct aki_signal signal; queue(struct camu_sink_cmd) queue; struct aki_mutex mutex; + struct camu_sink_entry *queued; struct camu_sink_entry *current; array(struct camu_sink_entry *) entries; struct { @@ -103,3 +104,4 @@ 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); +void camu_sink_free(struct camu_sink *sink); -- cgit v1.2.3-101-g0448