From 5e3641e5e692c3f2f644a4bb809c88727cb8bee9 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Fri, 8 Nov 2024 14:53:40 -0500 Subject: Command queue for portal and list, work on server Most of the server stuff can undoubtedly be simplified. I'm still working that out. Signed-off-by: Andrew Opalach --- src/libsink/sink.c | 104 ++++++++++++++++++++++++++++++++++++++--------------- src/libsink/sink.h | 8 +++-- 2 files changed, 81 insertions(+), 31 deletions(-) (limited to 'src/libsink') diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 2c8790b..899023d 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -91,6 +91,8 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e entry->audio.state = BUFFER_SET_OR_BUFFERED; } else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) { entry->audio.state = BUFFER_CONFIGURED; + } else if (entry->audio.state == BUFFER_QUEUED) { + entry->audio.state = BUFFER_INIT; } } @@ -102,6 +104,9 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e entry->video.state = BUFFER_SET_OR_BUFFERED; } else if (entry->video.state == BUFFER_SET_OR_BUFFERED) { entry->video.state = BUFFER_CONFIGURED; + } else if (entry->video.state == BUFFER_QUEUED) { + // This can be hit when skipping through entries very fast. + entry->video.state = BUFFER_INIT; } } #endif @@ -276,7 +281,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case CLOSE: { - aki_signal_stop(&sink->signal); + aki_signal_stop(&sink->queue_signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); return; } @@ -298,7 +303,7 @@ static void queue_signal_callback(void *userdata) static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd) { camu_queue_push(sink->queue, cmd); - aki_signal_send(&sink->signal); + aki_signal_send(&sink->queue_signal); } static void maybe_remove_previous(struct camu_sink *sink) @@ -310,13 +315,32 @@ static void maybe_remove_previous(struct camu_sink *sink) sink->previous.size = 0; } +static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *entry, struct camu_sink_entry *current) +{ + al_assert(entry != current); + struct camu_sink_entry *rentry; + al_array_foreach_rev(sink->previous, i, rentry) { + if (rentry == current) { + // If the entry we are about to add is in previous, + // remove it immediately. + remove_entry_buffers(sink, current); + al_array_remove_at(sink->previous, i); + } + } + // Don't accept duplicates. + al_array_foreach_rev(sink->previous, i, rentry) { + if (rentry == entry) return; + } + al_array_push(sink->previous, entry); +} + void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { - u8 state = entry->audio.state; - al_assert(state != BUFFER_ADDED); - if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } else if (state == BUFFER_SET_OR_BUFFERED) { + al_assert(entry->audio.state != BUFFER_ADDED); + if (entry->audio.state == BUFFER_CONFIGURED) { + entry->audio.state = BUFFER_SET_OR_BUFFERED; + } else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) { + entry->audio.state = BUFFER_ADDED; entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); // It's possible for this entry's video buffer to have been added and removed by EOF // before this point. This needs to be a consideration for keeping sync. @@ -324,24 +348,24 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) if (VIDEO_READY_OR_EMPTY(entry)) { maybe_remove_previous(entry->sink); } +#ifndef CAMU_SINK_LOCAL camu_audio_buffer_unpause(&entry->audio.buf); +#endif queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO }); - state = BUFFER_ADDED; } - 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; - al_assert(state != BUFFER_ADDED); - if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } else if (state == BUFFER_SET_OR_BUFFERED) { + al_assert(entry->video.state != BUFFER_ADDED); + if (entry->video.state == BUFFER_CONFIGURED) { + entry->video.state = BUFFER_SET_OR_BUFFERED; + } else if (entry->video.state == BUFFER_SET_OR_BUFFERED) { + entry->video.state = BUFFER_ADDED; entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); if (AUDIO_READY_OR_EMPTY(entry)) { maybe_remove_previous(entry->sink); @@ -351,9 +375,7 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry) .op = single_frame ? STOP : START, .value.i = CAMU_SINK_VIDEO }); - state = BUFFER_ADDED; } - entry->video.state = state; } #endif @@ -416,7 +438,7 @@ static void video_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_BUFFERED: { bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); aki_mutex_lock(&sink->mutex); - if (single_frame || !entry->ended) { + if (!entry->ended || single_frame) { add_video_if_set_and_buffered(entry); } aki_mutex_unlock(&sink->mutex); @@ -477,16 +499,21 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent frames -= sink->video.renderer->get_latency(sink->video.renderer); camu_video_buffer_set_latency(&entry->video.buf, -frames); } +#else + (void)sink; + (void)entry; #endif #else // To sync clients with differing audio latencies our only option is to factor the mixer // latency directly into the audio buffer. f64 audio = camu_mixer_get_latency(sink->audio.mixer); +#ifndef CAMU_SINK_NO_VIDEO if (!BUFFER_EMPTY(&entry->video)) { s32 frames = audio / entry->video.buf.avg_frame_duration; frames += sink->video.renderer->get_latency(sink->video.renderer); camu_video_buffer_set_latency(&entry->video.buf, frames); } +#endif camu_audio_buffer_set_latency(&entry->audio.buf, audio); #endif } @@ -650,10 +677,10 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, ) { 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); + aki_signal_init(&sink->queue_signal, queue_signal_callback, sink); + aki_signal_start(&sink->queue_signal, sink->loop); + camu_queue_init(sink->queue); sink->queued = NULL; sink->current = NULL; al_array_init(sink->previous); @@ -720,7 +747,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *entry) } } else { if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } } set_or_queue_entry(entry); @@ -832,8 +859,6 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn entry->ended = ended; - if (entry == sink->current) goto out; - if (op == LIANA_SINK_BUFFER) { goto out; } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { @@ -844,10 +869,14 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn #ifdef CAMU_SINK_LOCAL (void)at; (void)pause; - if (sink->current && !camu_clock_is_paused(&sink->current->clock)) { - camu_clock_pause(&sink->current->clock, 0); + if (sink->current) { + if (!camu_clock_is_paused(&sink->current->clock)) { + camu_clock_pause(&sink->current->clock, 0); + } + maybe_add_to_previous(sink, sink->current, entry); } - switch_to(sink, entry); + set_or_queue_entry(entry); + sink->current = entry; // This will resume a user paused stream. camu_clock_resume(&entry->clock, 0); #else @@ -858,7 +887,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn sink->target = NULL; } if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } set_or_queue_entry(entry); sink->current = entry; @@ -871,7 +900,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (entry == sink->current) { camu_audio_buffer_unpause(&entry->audio.buf); } else if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } camu_clock_resume(&entry->clock, at); set_or_queue_entry(entry); @@ -1028,6 +1057,14 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn) aki_rpc_connection_command(sink->conn, packet, idd_callback, sink); } +static void reconnect_timer_callback(void *userdata, struct aki_timer *timer) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)timer; + aki_rpc_reconnect(&sink->client, &sink->addr, sink->port); + aki_timer_stop(&sink->reconnect_timer); +} + static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; @@ -1035,12 +1072,15 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection al_assert(sink->conn == conn); sink->conn = NULL; } + aki_timer_again(&sink->reconnect_timer); } bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name) { al_str_clone(&sink->name, name); sink->type = type; + aki_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); + aki_timer_set_repeat(&sink->reconnect_timer, AKI_TS_FROM_USEC(1000000)); aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; @@ -1050,7 +1090,9 @@ bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str if (!aki_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) { return false; } - aki_rpc_connect(&sink->client, addr, port); + al_str_clone(&sink->addr, addr); + sink->port = port; + aki_rpc_connect(&sink->client, &sink->addr, sink->port); return true; } @@ -1142,6 +1184,8 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { + aki_timer_stop(&sink->reconnect_timer); + aki_timer_disable(&sink->reconnect_timer); if (sink->conn) aki_rpc_conn_disconnect(sink->conn); struct camu_sink_entry *entry; al_array_foreach_rev(sink->entries, i, entry) { @@ -1160,4 +1204,6 @@ void camu_sink_free(struct camu_sink *sink) aki_rpc_free(&sink->client); camu_queue_free(sink->queue); aki_mutex_destroy(&sink->mutex); + al_str_free(&sink->addr); + al_str_free(&sink->name); } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index b830230..0b55ad6 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -4,6 +4,7 @@ #include #include #include +#include #include "../util/queue.h" @@ -68,11 +69,14 @@ struct camu_sink { struct aki_event_loop *loop; str name; u8 type; + str addr; + u16 port; struct aki_rpc client; struct aki_rpc_connection *conn; - struct aki_signal signal; - queue(struct camu_sink_cmd) queue; struct aki_mutex mutex; + struct aki_timer reconnect_timer; + struct aki_signal queue_signal; + queue(struct camu_sink_cmd) queue; str default_list; struct camu_sink_entry *current; struct camu_sink_entry *queued; -- cgit v1.2.3-101-g0448