From 72eb0c6381f9406a4e45d23f65373d4936770433 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Sun, 27 Oct 2024 15:50:21 -0400 Subject: More synced list, another vcr fix Signed-off-by: Andrew Opalach --- src/libsink/common.h | 10 +--- src/libsink/sink.c | 155 ++++++++++++++++++++++++++++----------------------- src/libsink/sink.h | 1 + 3 files changed, 88 insertions(+), 78 deletions(-) (limited to 'src/libsink') diff --git a/src/libsink/common.h b/src/libsink/common.h index 0b1a6c9..c29e494 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -1,13 +1,9 @@ #pragma once -#define CAMU_SINK_LOCAL 1 +#define CAMU_SINK_LOCAL 0 enum { - CAMU_SINK_CLEAR = 0, - CAMU_SINK_SET, - CAMU_SINK_BUFFER, - CAMU_SINK_BUFFER_AND_QUEUE, + CAMU_SINK_SET = 0, CAMU_SINK_PAUSE, - CAMU_SINK_SEEK, - CAMU_SINK_DURATION + CAMU_SINK_SEEK }; diff --git a/src/libsink/sink.c b/src/libsink/sink.c index b9554c4..a2262bc 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -22,7 +22,6 @@ enum { BUFFER_CONFIGURED, BUFFER_SET_OR_BUFFERED, BUFFER_ADDED, - BUFFER_REMOVED }; enum { @@ -39,14 +38,14 @@ enum { #define ENTRY_MAX_AGE 7 +#define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED) + #ifdef CAMU_SINK_NO_VIDEO -#define ENTRY_VIDEO_READY_OR_EMPTY(entry) true +#define VIDEO_READY_OR_EMPTY(entry) true #else -#define ENTRY_VIDEO_READY_OR_EMPTY(entry) \ - ((entry)->video.state == BUFFER_INIT || (entry)->video.state == BUFFER_QUEUED || (entry)->video.state == BUFFER_ADDED) +#define VIDEO_READY_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state == BUFFER_ADDED) #endif -#define ENTRY_AUDIO_READY_OR_EMPTY(entry) \ - ((entry)->audio.state == BUFFER_INIT || (entry)->audio.state == BUFFER_QUEUED || (entry)->audio.state == BUFFER_ADDED) +#define AUDIO_READY_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->audio) || (entry)->audio.state == BUFFER_ADDED) #if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED #define BLOCKING_SLEEP(delay) aki_thread_sleep(delay) @@ -127,6 +126,7 @@ static void set_or_queue_entry(struct camu_sink_entry *entry) } else { add_audio_if_set_and_buffered(entry); } + #ifndef CAMU_SINK_NO_VIDEO if (entry->video.state == BUFFER_INIT) { entry->video.state = BUFFER_QUEUED; @@ -141,7 +141,7 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent { if (camu_clock_is_paused(&entry->clock)) { camu_clock_resume(&entry->clock, 0); - bool no_audio = entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED; + bool no_audio = BUFFER_EMPTY(&entry->audio); if (!no_audio && sink->audio.state == SINK_PAUSED) { sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); sink->audio.state = SINK_PLAYING; @@ -313,22 +313,23 @@ static void maybe_remove_previous(struct camu_sink *sink) void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->audio.state; - if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) { + al_assert(state != BUFFER_ADDED); + if (state == BUFFER_CONFIGURED) { + state = BUFFER_SET_OR_BUFFERED; + } else if (state == BUFFER_SET_OR_BUFFERED) { + 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. + // Same but reversed in add_video_if_set_and_buffered(). + if (VIDEO_READY_OR_EMPTY(entry)) { + maybe_remove_previous(entry->sink); + } camu_audio_buffer_unpause(&entry->audio.buf); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO }); - } - if (state == BUFFER_SET_OR_BUFFERED) { - bool can_resume = ENTRY_VIDEO_READY_OR_EMPTY(entry); - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - if (can_resume) { - maybe_remove_previous(entry->sink); - } state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; } entry->audio.state = state; } @@ -337,23 +338,20 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->video.state; - if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) { + al_assert(state != BUFFER_ADDED); + if (state == BUFFER_CONFIGURED) { + state = BUFFER_SET_OR_BUFFERED; + } else if (state == BUFFER_SET_OR_BUFFERED) { + 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); + } bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = single_frame ? STOP : START, .value.i = CAMU_SINK_VIDEO }); - } - if (state == BUFFER_SET_OR_BUFFERED) { - bool can_resume = ENTRY_AUDIO_READY_OR_EMPTY(entry); - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - if (can_resume) { - maybe_remove_previous(entry->sink); - } - state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; } entry->video.state = state; } @@ -366,7 +364,9 @@ static void audio_buffer_callback(void *userdata, u8 op) switch (op) { case CAMU_BUFFER_BUFFERED: aki_mutex_lock(&sink->mutex); - add_audio_if_set_and_buffered(entry); + if (!entry->ended) { + add_audio_if_set_and_buffered(entry); + } aki_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_CORK: @@ -388,8 +388,10 @@ static void audio_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->audio.track); aki_mutex_lock(&sink->mutex); - camu_clock_end(&entry->clock); remove_entry_audio_buffer(sink, entry); + camu_clock_end(&entry->clock); + // EOF on the audio buffer unconditionally + // advances the queue. if (sink->queued) { set_or_queue_entry(sink->queued); sink->current = sink->queued; @@ -411,11 +413,15 @@ static void video_buffer_callback(void *userdata, u8 op) struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; struct camu_sink *sink = entry->sink; switch (op) { - case CAMU_BUFFER_BUFFERED: + case CAMU_BUFFER_BUFFERED: { + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); aki_mutex_lock(&sink->mutex); - add_video_if_set_and_buffered(entry); + if (single_frame || !entry->ended) { + add_video_if_set_and_buffered(entry); + } aki_mutex_unlock(&sink->mutex); break; + } case CAMU_BUFFER_CORK: lia_vcr_cork(entry->video.track); break; @@ -430,7 +436,8 @@ static void video_buffer_callback(void *userdata, u8 op) if (!single_frame) { remove_entry_video_buffer(sink, entry); } - bool run_queue = !single_frame && (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED); + // If this entry has no audio, advance the queue. + bool run_queue = !single_frame && BUFFER_EMPTY(&entry->audio); if (run_queue) { camu_clock_end(&entry->clock); if (sink->queued) { @@ -463,10 +470,9 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent { // Entry has both audio and video configured. #ifndef CAMU_SINK_NO_VIDEO - if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED && - entry->video.state != BUFFER_INIT && entry->video.state != BUFFER_QUEUED) { + if (!BUFFER_EMPTY(&entry->audio) && !BUFFER_EMPTY(&entry->video)) { camu_video_buffer_set_latency(&entry->video.buf, -camu_mixer_get_latency(sink->audio.mixer)); - } else if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED) { + } else if (!BUFFER_EMPTY(&entry->audio)) { #if !CAMU_SINK_LOCAL camu_audio_buffer_set_latency(&entry->audio.buf, -camu_mixer_get_latency(sink->audio.mixer)); #endif @@ -698,6 +704,21 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) } } +static void switch_to(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + if (entry->ended || camu_clock_is_ended(&entry->clock)) { + if (sink->current) { + remove_entry_buffers(sink, sink->current); + } + } else { + if (sink->current) { + al_array_push(sink->previous, sink->current); + } + } + set_or_queue_entry(entry); + sink->current = entry; +} + static void clock_callback(void *userdata, u8 op) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; @@ -705,11 +726,7 @@ static void clock_callback(void *userdata, u8 op) if (op == CAMU_CLOCK_PAUSED) { aki_mutex_lock(&sink->mutex); if (sink->target) { - if (sink->current) { - al_array_push(sink->previous, sink->current); - } - set_or_queue_entry(sink->target); - sink->current = sink->target; + switch_to(sink, sink->target); sink->target = NULL; } aki_mutex_unlock(&sink->mutex); @@ -766,7 +783,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn u8 op = aki_packet_read_u8(packet); #if CAMU_SINK_LOCAL - if (op == CAMU_SINK_CLEAR) { + if (op == LIANA_SINK_UNSET) { if (sink->current) { remove_entry_buffers(sink, sink->current); sink->current = NULL; @@ -774,7 +791,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn goto out; } #else - if (op == CAMU_SINK_CLEAR) { al_assert(false); } // Unimplemented. + if (op == LIANA_SINK_UNSET) { al_assert(false); } // Unimplemented. #endif str addr; @@ -785,6 +802,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn u64 at = aki_packet_read_u64(packet); u64 seek_pos = aki_packet_read_u64(packet); u8 pause = aki_packet_read_u8(packet); + bool ended = aki_packet_read_bool(packet); bool created; struct camu_sink_entry *entry = get_entry_from_id(sink, node_id, &created); @@ -794,6 +812,9 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (created) { camu_clock_set(&entry->clock, seek_pos / 1000000.0); + if (ended) { + camu_clock_end(&entry->clock); + } struct camu_renderer *renderer = NULL; #ifndef CAMU_SINK_NO_VIDEO renderer = sink->video.renderer; @@ -801,9 +822,11 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, seek_pos, renderer); } - if (op == CAMU_SINK_BUFFER) { + entry->ended = ended; + + if (op == LIANA_SINK_BUFFER) { goto out; - } else if (op == CAMU_SINK_BUFFER_AND_QUEUE) { + } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { sink->queued = entry; goto out; } @@ -811,14 +834,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn #if CAMU_SINK_LOCAL (void)at; (void)pause; - if (sink->current) { - if (!camu_clock_is_paused(&sink->current->clock)) { - camu_clock_pause(&sink->current->clock, 0); - } - al_array_push(sink->previous, sink->current); + if (sink->current && !camu_clock_is_paused(&sink->current->clock)) { + camu_clock_pause(&sink->current->clock, 0); } - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); // This will resume a user paused stream. camu_clock_resume(&entry->clock, 0); #else @@ -849,20 +868,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn sink->current = entry; break; case LIANA_PAUSE_PAUSE: - if (camu_clock_is_ended(&entry->clock)) { - // Server didn't know this entry was ended, but it is. - queue_cmd(sink, (struct camu_sink_cmd){ - .op = END, - .value.i = entry->sequence - }); - camu_clock_pause(&sink->current->clock, at); - sink->current = entry; - } else if (sink->current) { + if (sink->current) { if (camu_clock_is_ended(&sink->current->clock)) { // Server thought we weren't done, be we are. - al_array_push(sink->previous, sink->current); - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); } else { sink->target = entry; camu_clock_pause(&sink->current->clock, at); @@ -882,9 +891,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn camu_clock_pause(&prev_target->clock, at); } else if (sink->current) { if (camu_clock_is_ended(&sink->current->clock)) { - al_array_push(sink->previous, sink->current); - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); } else { sink->target = entry; camu_clock_pause(&sink->current->clock, at); @@ -897,7 +904,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn out: aki_mutex_unlock(&sink->mutex); - if (op != CAMU_SINK_BUFFER) { + if (op != LIANA_SINK_BUFFER) { maybe_cleanup_old_entries(sink); } @@ -912,15 +919,17 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co (void)conn; (void)rpacket; + s32 sequence = aki_packet_read_s32(packet); + u64 at = aki_packet_read_u64(packet); + u8 pause = aki_packet_read_u8(packet); + aki_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; if (!current) goto out; + #if CAMU_SINK_LOCAL sink_local_pause(sink, current); #else - s32 sequence = aki_packet_read_s32(packet); - u64 at = aki_packet_read_u64(packet); - u8 pause = aki_packet_read_u8(packet); if (current->sequence == sequence) { if (pause == LIANA_PAUSE_PAUSE) { camu_clock_pause(&sink->current->clock, at); @@ -957,7 +966,10 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con aki_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; aki_mutex_unlock(&sink->mutex); + if (!current) goto out; + if (current->sequence == sequence) { + current->ended = false; #if CAMU_SINK_LOCAL (void)at; lia_client_seek(¤t->client, pos, 0); @@ -966,6 +978,7 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con #endif } +out: aki_packet_free(packet); return false; } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 6d560f6..b830230 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -41,6 +41,7 @@ struct camu_sink_entry { u16 lru; struct camu_clock clock; struct lia_client client; + bool ended; struct { u8 state; bool armed; -- cgit v1.2.3-101-g0448