diff options
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/common.h | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 529 | ||||
| -rw-r--r-- | src/libsink/sink.h | 1 |
3 files changed, 335 insertions, 197 deletions
diff --git a/src/libsink/common.h b/src/libsink/common.h index 708b1bb..0bb3ded 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -1,6 +1,6 @@ #pragma once -//#define CAMU_SINK_LOCAL +#define CAMU_SINK_LOCAL enum { CAMU_SINK_SET = 0, diff --git a/src/libsink/sink.c b/src/libsink/sink.c index ede3889..fe56221 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -14,6 +14,7 @@ #include "../render/renderer_libplacebo.h" #endif +// Requested state of the sinks outputs. enum { SINK_EMPTY = 0, SINK_PAUSED, @@ -22,12 +23,19 @@ enum { enum { BUFFER_INIT = 0, + // Set but not configured. BUFFER_QUEUED, + // Ready to receive data. BUFFER_CONFIGURED, + // The next call to add can add the buffer. BUFFER_SET_OR_BUFFERED, + // Treat the buffer like it's added, even though it might not be. BUFFER_ADDED, + // Effectively SET_OR_BUFFERED but not addable until after a reset. + BUFFER_ENDED }; +// Command queue commands. enum { START, STOP, @@ -35,28 +43,43 @@ enum { SEEK, SKIP, SHUFFLE, - END, RESEEK, + UNSET, + END, CLOSE }; +// Number of entries to keep buffered at one time. #define ENTRY_MAX_AGE 5 +// If a buffer is still INIT or QUEUED after an entry is configured, it's "empty". #define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED) +#define BUFFER_NOT_EMPTY(buf) (!BUFFER_EMPTY(buf)) #ifdef CAMU_SINK_NO_VIDEO #define VIDEO_ADDED_OR_EMPTY(entry) true #else -#define VIDEO_ADDED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state == BUFFER_ADDED) +#define VIDEO_ADDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->video)) +#endif +#define AUDIO_ADDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->audio)) + +#ifdef CAMU_SINK_NO_VIDEO +#define VIDEO_ENDED_OR_EMPTY(entry) true +#else +#define VIDEO_ENDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->video)) #endif -#define AUDIO_ADDED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->audio) || (entry)->audio.state == BUFFER_ADDED) +#define AUDIO_ENDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->audio)) #ifdef CAMU_SINK_NO_VIDEO #define VIDEO_REMOVED_OR_EMPTY(entry) true #else -#define VIDEO_REMOVED_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state != BUFFER_ADDED) +#define VIDEO_REMOVED_OR_EMPTY(entry) ((entry)->video.state != BUFFER_ADDED) +#endif +#define AUDIO_REMOVED_OR_EMPTY(entry) ((entry)->audio.state != BUFFER_ADDED) + +#ifndef CAMU_SINK_NO_VIDEO +#define ENTRY_IS_SINGLE_FRAME(entry) camu_video_buffer_is_single_frame(&(entry)->video.buf) #endif -#define AUDIO_REMOVED_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) @@ -64,7 +87,7 @@ enum { #define BLOCKING_SLEEP(delay) aki_event_loop_sleep(sink->loop, delay) #endif -static bool entry_audio_buffer_held(struct camu_sink_entry *entry) +static inline bool entry_audio_buffer_held(struct camu_sink_entry *entry) { #ifdef CAMU_MIXER_THREADED return al_atomic_load(u8)(&entry->audio.buf.ref, AL_ATOMIC_RELAXED) == 1; @@ -75,7 +98,7 @@ static bool entry_audio_buffer_held(struct camu_sink_entry *entry) } #ifndef CAMU_SINK_NO_VIDEO -static bool entry_video_buffer_held(struct camu_sink_entry *entry) +static inline bool entry_video_buffer_held(struct camu_sink_entry *entry) { #ifdef CAMU_SCREEN_THREADED return al_atomic_load(u8)(&entry->video.buf.ref, AL_ATOMIC_RELAXED) == 1; @@ -86,7 +109,7 @@ static bool entry_video_buffer_held(struct camu_sink_entry *entry) } #endif -static bool entry_buffers_held(struct camu_sink_entry *entry) +static inline bool entry_buffers_held(struct camu_sink_entry *entry) { #ifndef CAMU_SINK_NO_VIDEO return entry_audio_buffer_held(entry) || entry_video_buffer_held(entry); @@ -97,7 +120,8 @@ static bool entry_buffers_held(struct camu_sink_entry *entry) static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_entry *entry) { - al_assert(!entry->ended); + al_assert(!entry->ended && entry->audio.state != BUFFER_ENDED); + al_assert(entry->audio.state != BUFFER_INIT); 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; @@ -111,6 +135,9 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e #ifndef CAMU_SINK_NO_VIDEO static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_entry *entry) { + // Don't assert !entry->ended here because of single frame handling. + al_assert(entry->video.state != BUFFER_ENDED); + al_assert(entry->video.state != BUFFER_INIT); 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; @@ -123,13 +150,20 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e } #endif +// It's possible for some of an entries buffers to be ENDED while others are still +// ADDED and playing. We handle that by making remove_entry_buffers() and +// add_audio/video_if_set_and_buffered() no-ops for ENDED buffers. static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry) { al_assert(!entry->ended); + if (entry->audio.state != BUFFER_ENDED) { + remove_entry_audio_buffer(sink, entry); + } #ifndef CAMU_SINK_NO_VIDEO - remove_entry_video_buffer(sink, entry); + if (entry->video.state != BUFFER_ENDED) { + remove_entry_video_buffer(sink, entry); + } #endif - remove_entry_audio_buffer(sink, entry); } static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry); @@ -160,14 +194,12 @@ 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 = BUFFER_EMPTY(&entry->audio); - if (!no_audio && sink->audio.state == SINK_PAUSED) { + if (!BUFFER_EMPTY(&entry->audio) && 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 - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (!single_frame && sink->video.state == SINK_PAUSED) { + if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) { sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL); sink->video.state = SINK_PLAYING; } @@ -175,8 +207,7 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent } else { camu_clock_pause(&entry->clock, 0); #ifndef CAMU_SINK_NO_VIDEO - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (!single_frame && sink->video.state == SINK_PLAYING) { + if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) { sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL); sink->video.state = SINK_PAUSED; } @@ -280,6 +311,19 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } + case RESEEK: { + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; + lia_client_reseek(&entry->client); + break; + } + case UNSET: { + if (!sink->conn) return; + struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + aki_packet_write_str(packet, &sink->default_list); + aki_packet_write_u8(packet, CAMU_LIST_UNSET); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } case END: { if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); @@ -289,11 +333,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } - case RESEEK: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - lia_client_reseek(&entry->client); - break; - } case CLOSE: { aki_signal_stop(&sink->queue_signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); @@ -341,8 +380,8 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, { sink->loop = loop; aki_mutex_init(&sink->mutex); - aki_signal_init(&sink->queue_signal, queue_signal_callback, sink); - aki_signal_start(&sink->queue_signal, sink->loop); + aki_signal_init(&sink->queue_signal, sink->loop, queue_signal_callback, sink); + aki_signal_start(&sink->queue_signal); camu_queue_init(sink->queue); sink->queued = NULL; sink->current = NULL; @@ -369,30 +408,44 @@ 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 *previous, struct camu_sink_entry *target) +// Due to the looseness of the previous queue, we may have to explicitly remove an entry +// at a point if it becomes incorrect to attempt removing it's buffers. +// An obvious example of this is at the point an entry gets freed. See LIANA_CLIENT_REMOVE_BUFFERS. +static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_entry *entry) { - al_assert(previous != target); - // If our target is ended, remove previous immediately. - if (target->ended) { - remove_entry_buffers(sink, previous); - return; + struct camu_sink_entry *rentry; + al_array_foreach_rev(sink->previous, i, rentry) { + if (rentry == entry) { + al_array_remove_at(sink->previous, i); + } } +} + +static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target) +{ + al_assert(previous != target && !previous->ended); + struct camu_sink_entry *rentry; al_array_foreach_rev(sink->previous, i, rentry) { - // If the entry we are about to add is in previous, run the queue now. + // If the entry we are about to set is in previous, run the queue. if (rentry == target) { maybe_remove_previous(sink); - return; + break; } } - // Don't accept duplicates. - al_array_foreach_rev(sink->previous, i, rentry) { - if (rentry == previous) return; + + // Every call to maybe_add_to_previous() should map to a remove_entry_buffers(). + // We can take a shortcut here because pushing an entry to previous is + // pointless if it's not currently added. + if (!AUDIO_ADDED_OR_EMPTY(previous) && !VIDEO_ADDED_OR_EMPTY(previous)) { + remove_entry_buffers(sink, previous); + return; } + al_array_push(sink->previous, previous); } -static s32 lru_compare(const void *a, const void *b) +static s32 entry_lru_compare(const void *a, const void *b) { struct camu_sink_entry *aa = *((struct camu_sink_entry **)a); struct camu_sink_entry *bb = *((struct camu_sink_entry **)b); @@ -403,28 +456,34 @@ static s32 lru_compare(const void *a, const void *b) static void maybe_cleanup_old_entries(struct camu_sink *sink) { - // We check size <= MAX_AGE in the loop because sink->lru - // is not indicative of the amount of entries we have loaded. - // There are various reasons for this but the most obvious is - // that it's incremented for buffer and queue operations. - // - // Handle sink->lru wrapping. - // 0 65532 65533 65534 65535 - // 0 1 65533 65534 65535 - // 0 1 2 65534 65535 - // 0 1 2 3 65535 - // 0 1 2 3 4 - al_array_sort(sink->entries, struct camu_sink_entry *, lru_compare); + al_array_sort(sink->entries, struct camu_sink_entry *, entry_lru_compare); + + // We check size <= MAX_AGE in the loops because sink->lru is + // not indicative of the amount of entries we have loaded. + // The most obvious reason for that is that it's incremented for + // buffer and queue operations. + struct camu_sink_entry *entry; + // Handle sink->lru wrapping. This has to happen in a step + // before the no wrapping case. al_array_foreach_rev(sink->entries, i, entry) { - if (entry->lru > sink->lru && (UINT16_MAX - (entry->lru - 1)) + sink->lru >= ENTRY_MAX_AGE) { + // 0 65532 65533 65534 65535 + // 0 1 65533 65534 65535 + // 0 1 2 65534 65535 + // 0 1 2 3 65535 + // 0 1 2 3 4 + if (entry->lru > sink->lru && ((UINT16_MAX - entry->lru) + 1) + sink->lru >= ENTRY_MAX_AGE) { al_array_remove_at(sink->entries, i); lia_client_disconnect(&entry->client); } if (sink->entries.size <= ENTRY_MAX_AGE) return; } + + // Checking sink->lru >= ENTRY_MAX_AGE should guarantee + // we don't have to consider wrapping here. if (sink->lru >= ENTRY_MAX_AGE) { al_array_foreach_rev(sink->entries, i, entry) { + al_assert(sink->lru >= entry->lru); if (sink->lru - entry->lru >= ENTRY_MAX_AGE) { al_array_remove_at(sink->entries, i); lia_client_disconnect(&entry->client); @@ -437,20 +496,42 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { al_assert(!entry->ended && entry->audio.state != BUFFER_ADDED); + if (entry->audio.state == BUFFER_ENDED) { + al_log_warn("sink", "Tried to add an ended audio buffer."); + return; + } 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. - // Same but reversed in add_video_if_set_and_buffered(). - if (VIDEO_ADDED_OR_EMPTY(entry)) { - maybe_remove_previous(entry->sink); - } #ifndef CAMU_SINK_LOCAL camu_audio_buffer_unpause(&entry->audio.buf); #endif + if (VIDEO_ADDED_OR_EMPTY(entry)) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); +#ifndef CAMU_SINK_NO_VIDEO + bool stop_video = false; + // Single frames are unconditionally added in add_video_if_set_and_buffered(). + if (!ENTRY_IS_SINGLE_FRAME(entry)) { + if (entry->video.state == BUFFER_ADDED) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); + } else { + // Entry has no video. + stop_video = true; + } + } +#endif + maybe_remove_previous(entry->sink); +#ifndef CAMU_SINK_NO_VIDEO + // This must come after maybe_remove_previous(). + if (stop_video) { + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = STOP, + .value.i = CAMU_SINK_VIDEO + }); + } +#endif + } queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO @@ -461,16 +542,26 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) #ifndef CAMU_SINK_NO_VIDEO void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { - al_assert(!entry->ended && entry->video.state != BUFFER_ADDED); + al_assert(entry->video.state != BUFFER_ADDED); + if (entry->video.state == BUFFER_ENDED) { + al_log_warn("sink", "Tried to add an ended video buffer."); + return; + } 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); + bool single_frame = ENTRY_IS_SINGLE_FRAME(entry); if (AUDIO_ADDED_OR_EMPTY(entry)) { + if (entry->audio.state == BUFFER_ADDED) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + } + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); maybe_remove_previous(entry->sink); + } else if (single_frame) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); } - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); + // Stopping video here is needed if skipping from a video to an image. queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = single_frame ? STOP : START, .value.i = CAMU_SINK_VIDEO @@ -479,29 +570,62 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry) } #endif +static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) +{ + if (sink->current) { + struct camu_sink_entry *current = sink->current; + al_assert(current != target); + if (!current->ended) { + if (target->ended) { + remove_entry_buffers(sink, current); + } else { + maybe_add_to_previous(sink, current, target); + } + } else { +#ifndef CAMU_SINK_NO_VIDEO + if (ENTRY_IS_SINGLE_FRAME(current)) { + remove_entry_video_buffer(sink, current); + } +#endif + al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current)); + } + } + + if (!target->ended) { + add_or_queue_entry(target); +#ifndef CAMU_SINK_NO_VIDEO + } else if (ENTRY_IS_SINGLE_FRAME(target)) { + add_video_if_set_and_buffered(target); +#endif + } + + sink->current = target; +} + +#ifndef CAMU_SINK_LOCAL +static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *target, u64 at) +{ + if (!sink->current || sink->current->ended) { + switch_to(sink, target); + } else { + sink->target = target; + camu_clock_pause(&sink->current->clock, at); + } +} +#endif + static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink_entry *entry) { al_log_info("sink", "Entry ended."); entry->ended = true; - // TODO: Can it make sense for this entry to be in previous? - maybe_remove_previous(sink); + maybe_remove_from_previous(sink, entry); if (sink->target) { - if (!sink->target->ended) { - add_or_queue_entry(sink->target); - } - sink->current = sink->target; + switch_to(sink, sink->target); sink->target = NULL; return true; - } else if (sink->queued) { - // TODO: What is the right behavior if sink->queued - // and sink->target are both set. - if (!sink->queued->ended) { - add_or_queue_entry(sink->queued); - } - sink->current = sink->queued; - sink->queued = NULL; + }/* else if (sink->queued) { return true; - } + }*/ return false; } @@ -535,11 +659,16 @@ static void audio_buffer_callback(void *userdata, u8 op) lia_vcr_cork(entry->audio.track); aki_mutex_lock(&sink->mutex); al_log_info("sink", "Audio EOF."); - remove_entry_audio_buffer(sink, entry); - bool run_queue = VIDEO_REMOVED_OR_EMPTY(entry); + if (entry->audio.state == BUFFER_ADDED) { + remove_entry_audio_buffer(sink, entry); + } + // This assert likely doesn't matter, but should be kept + // if it doesn't unnecessarialy trip. + al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED); + entry->audio.state = BUFFER_ENDED; + bool run_queue = VIDEO_ENDED_OR_EMPTY(entry); #ifndef CAMU_SINK_NO_VIDEO - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - run_queue = run_queue || single_frame; + run_queue = run_queue || ENTRY_IS_SINGLE_FRAME(entry); #endif if (run_queue) { end_entry_and_advance_queue(sink, entry); @@ -576,16 +705,20 @@ static void video_buffer_callback(void *userdata, u8 op) break; case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->video.track); - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); bool swapped = false; + bool run_queue = false; aki_mutex_lock(&sink->mutex); al_log_info("sink", "Video EOF."); - if (!single_frame) { - remove_entry_video_buffer(sink, entry); - } - bool run_queue = !single_frame && AUDIO_REMOVED_OR_EMPTY(entry); - if (run_queue) { - swapped = end_entry_and_advance_queue(sink, entry); + if (!ENTRY_IS_SINGLE_FRAME(entry)) { + if (entry->video.state == BUFFER_ADDED) { + remove_entry_video_buffer(sink, entry); + } + al_assert(entry->video.state == BUFFER_SET_OR_BUFFERED); + entry->video.state = BUFFER_ENDED; + if (AUDIO_ENDED_OR_EMPTY(entry)) { + swapped = end_entry_and_advance_queue(sink, entry); + run_queue = true; + } } aki_mutex_unlock(&sink->mutex); if (!swapped) { @@ -606,36 +739,14 @@ static void video_buffer_callback(void *userdata, u8 op) } #endif -static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) +static void maybe_unset_current(struct camu_sink *sink) { if (sink->current) { - struct camu_sink_entry *current = sink->current; - if (current->ended) { - bool single_frame = camu_video_buffer_is_single_frame(¤t->video.buf); - if (single_frame) { - remove_entry_video_buffer(sink, current); - } - al_assert(current->audio.state != BUFFER_ADDED && current->video.state != BUFFER_ADDED); - } else { - maybe_add_to_previous(sink, current, target); + if (!sink->current->ended) { + remove_entry_buffers(sink, sink->current); } - } - if (!target->ended) { - add_or_queue_entry(target); - } - sink->current = target; -} - -static void set_target_and_pause(struct camu_sink *sink, struct camu_sink_entry *target, u64 at) -{ - if (!sink->current || sink->current->ended) { - if (!target->ended) { - add_or_queue_entry(target); - } - sink->current = target; - } else { - sink->target = target; - camu_clock_pause(&sink->current->clock, at); + // TODO: A current-less state is not properly handled. + sink->current = NULL; } } @@ -645,9 +756,11 @@ static void clock_callback(void *userdata, u8 op) struct camu_sink *sink = entry->sink; if (op == CAMU_CLOCK_PAUSED) { aki_mutex_lock(&sink->mutex); - if (sink->target) { - switch_to(sink, sink->target); - sink->target = NULL; + if (entry == sink->current) { + if (sink->target) { + switch_to(sink, sink->target); + sink->target = NULL; + } } aki_mutex_unlock(&sink->mutex); } @@ -657,8 +770,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent { #ifdef CAMU_SINK_LOCAL #ifndef CAMU_SINK_NO_VIDEO - // Entry has both audio and video configured. - if (!BUFFER_EMPTY(&entry->audio) && !BUFFER_EMPTY(&entry->video)) { + if (BUFFER_NOT_EMPTY(&entry->audio) && BUFFER_NOT_EMPTY(&entry->video)) { f64 audio = camu_mixer_get_latency(sink->audio.mixer); s32 frames = audio / entry->video.buf.avg_frame_duration; frames -= sink->video.renderer->get_latency(sink->video.renderer); @@ -673,7 +785,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent // 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)) { + if (BUFFER_NOT_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); @@ -694,40 +806,32 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str case CAMU_STREAM_AUDIO: entry->audio.track = (struct lia_vcr_track *)opaque; if (!camu_audio_buffer_configure(&entry->audio.buf, stream, sink->audio.mixer)) { + al_log_warn("sink", "Audio buffer failed to configure."); lia_client_disconnect(&entry->client); } aki_mutex_lock(&sink->mutex); - if (entry->audio.state == BUFFER_QUEUED) { - entry->audio.state = BUFFER_SET_OR_BUFFERED; - } else { - entry->audio.state = BUFFER_CONFIGURED; - } - if (VIDEO_ADDED_OR_EMPTY(entry)) { - evaluate_latency(sink, entry); - } + entry->audio.state = entry->audio.state == BUFFER_QUEUED ? + BUFFER_SET_OR_BUFFERED : BUFFER_CONFIGURED; + if (VIDEO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry); aki_mutex_unlock(&sink->mutex); break; #ifndef CAMU_SINK_NO_VIDEO case CAMU_STREAM_VIDEO: entry->video.track = (struct lia_vcr_track *)opaque; if (!camu_video_buffer_configure(&entry->video.buf, stream, sink->video.renderer)) { + al_log_warn("sink", "Video buffer failed to configure."); lia_client_disconnect(&entry->client); } aki_mutex_lock(&sink->mutex); - if (entry->video.state == BUFFER_QUEUED) { - entry->video.state = BUFFER_SET_OR_BUFFERED; - } else { - entry->video.state = BUFFER_CONFIGURED; - } - if (AUDIO_ADDED_OR_EMPTY(entry)) { - evaluate_latency(sink, entry); - } + entry->video.state = entry->video.state == BUFFER_QUEUED ? + BUFFER_SET_OR_BUFFERED : BUFFER_CONFIGURED; + if (AUDIO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry); aki_mutex_unlock(&sink->mutex); break; case CAMU_STREAM_SUBTITLE: #ifdef CAMU_HAVE_FFMPEG if (!camu_video_buffer_configure_subtitles(&entry->video.buf, stream->av.stream->codecpar)) { - // TODO + // TODO: } #endif break; @@ -748,17 +852,23 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str struct camu_codec_frame *frame = (struct camu_codec_frame *)opaque; switch (stream->type) { case CAMU_STREAM_AUDIO: - camu_audio_buffer_push(&entry->audio.buf, frame); + if (BUFFER_NOT_EMPTY(&entry->audio)) { + camu_audio_buffer_push(&entry->audio.buf, frame); + return; + } break; #ifndef CAMU_SINK_NO_VIDEO case CAMU_STREAM_VIDEO: - camu_video_buffer_push(&entry->video.buf, frame); + if (BUFFER_NOT_EMPTY(&entry->video)) { + camu_video_buffer_push(&entry->video.buf, frame); + return; + } break; #endif default: - camu_codec_frame_discard(frame); break; } + camu_codec_frame_discard(frame); break; } case LIANA_CLIENT_SUBTITLE: { @@ -776,32 +886,54 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str break; } case LIANA_CLIENT_REMOVE_BUFFERS: { + bool reconnect = *(bool *)opaque; + aki_mutex_lock(&sink->mutex); - if (entry->audio.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + al_assert(!(reconnect && entry->ended)); + + // There are 2 reasons this entry might be in previous. + // 1. It's getting cleaned up before any entries set after it are buffered. + // 2. It got added to previous then seeked. + maybe_remove_from_previous(sink, entry); + + if (entry->audio.state == BUFFER_ENDED) { entry->audio.state = BUFFER_SET_OR_BUFFERED; + } else if (BUFFER_NOT_EMPTY(&entry->audio)) { + remove_entry_audio_buffer(sink, entry); } + #ifndef CAMU_SINK_NO_VIDEO - bool reconnect = *(bool *)opaque; - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - bool keep_video = reconnect && single_frame; - if (!keep_video && entry->video.state == BUFFER_ADDED) { - sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); + bool ignore_video = BUFFER_EMPTY(&entry->video); + if (entry->video.state == BUFFER_ENDED) { entry->video.state = BUFFER_SET_OR_BUFFERED; + } else if (BUFFER_NOT_EMPTY(&entry->video)) { + ignore_video = reconnect && ENTRY_IS_SINGLE_FRAME(entry); + if (!ignore_video) { + remove_entry_video_buffer(sink, entry); + } } #endif aki_mutex_unlock(&sink->mutex); + + while ( // Block until buffers are removed. #ifndef CAMU_SINK_NO_VIDEO - while (keep_video ? entry_audio_buffer_held(entry) : entry_buffers_held(entry)) { - BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000)); - } - if (reconnect && !single_frame) camu_video_buffer_reset(&entry->video.buf); + ignore_video ? entry_audio_buffer_held(entry) : entry_buffers_held(entry) #else - while (entry_buffers_held(entry)) { - BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000)); - } + entry_buffers_held(entry) +#endif + ) { BLOCKING_SLEEP(AKI_TS_FROM_USEC(2500)); } + + if (reconnect) { + if (BUFFER_NOT_EMPTY(&entry->audio)) { + camu_audio_buffer_reset(&entry->audio.buf); + } +#ifndef CAMU_SINK_NO_VIDEO + if (BUFFER_NOT_EMPTY(&entry->video) && !ignore_video) { + camu_video_buffer_reset(&entry->video.buf); + } #endif - camu_audio_buffer_reset(&entry->audio.buf); + } + break; } case LIANA_CLIENT_RESUME_AT: { @@ -813,13 +945,15 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str case LIANA_CLIENT_EOF: { switch (stream->type) { case CAMU_STREAM_AUDIO: { - camu_audio_buffer_flush(&entry->audio.buf); + if (BUFFER_NOT_EMPTY(&entry->audio)) { + camu_audio_buffer_flush(&entry->audio.buf); + } break; } #ifndef CAMU_SINK_NO_VIDEO case CAMU_STREAM_VIDEO: { - bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); - if (!single_frame) { + // Single frames are immediately flushed inside the buffer. + if (BUFFER_NOT_EMPTY(&entry->video) && !ENTRY_IS_SINGLE_FRAME(entry)) { camu_video_buffer_flush(&entry->video.buf); } break; @@ -831,12 +965,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str case LIANA_CLIENT_CLOSED: { aki_mutex_lock(&sink->mutex); if (entry == sink->target) { - if (sink->current) { - if (!sink->current->ended) { - remove_entry_buffers(sink, sink->current); - } - sink->current = NULL; - } + maybe_unset_current(sink); sink->target = NULL; } else if (entry == sink->current) { // current's buffers will already be removed. @@ -845,9 +974,11 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str switch_to(sink, sink->target); sink->target = NULL; } - } else if (entry == sink->queued) { + // If current was never fully added we need to do this here. + maybe_remove_previous(sink); + }/* else if (entry == sink->queued) { sink->queued = NULL; - } + }*/ lia_client_free(&entry->client); camu_audio_buffer_free(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO @@ -923,20 +1054,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn aki_mutex_lock(&sink->mutex); u8 op = aki_packet_read_u8(packet); -#ifdef CAMU_SINK_LOCAL - // TODO: look over this. if (op == LIANA_SINK_UNSET) { - if (sink->current) { - if (!sink->current->ended) { - remove_entry_buffers(sink, sink->current); - } - sink->current = NULL; - } + maybe_unset_current(sink); goto out; } -#else - if (op == LIANA_SINK_UNSET) { al_assert(false); } // Unimplemented. -#endif str addr; aki_packet_read_str(packet, &addr); @@ -946,7 +1067,9 @@ 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 previous_ended = aki_packet_read_bool(packet); bool ended = aki_packet_read_bool(packet); + (void)previous_ended; (void)ended; bool created; @@ -966,41 +1089,48 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (op == LIANA_SINK_BUFFER) { goto out; - } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { + }/* else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { sink->queued = entry; goto out; - } + }*/ #ifdef CAMU_SINK_LOCAL (void)at; (void)pause; - if (sink->current && !sink->current->ended && !camu_clock_is_paused(&sink->current->clock)) { + if (sink->current && !sink->current->ended && !camu_clock_is_armed(&sink->current->clock)) { camu_clock_pause(&sink->current->clock, 0); } switch_to(sink, entry); - // This will resume a user paused stream. - camu_clock_resume(&entry->clock, 0); + if (!entry->ended && camu_clock_is_armed(&entry->clock)) { + // This will resume user-paused entries, but whatever. + camu_clock_resume(&entry->clock, 0); + } #else - // PAUSE_NONE and PAUSE_PAUSE mean the server expects the entry - // being set to be ended. Not acting accordingly here is the - // only place where local entry->ended and the servers expectation - // being mismatched can cause issues. switch (pause) { - case LIANA_PAUSE_NONE: + case LIANA_PAUSE_NONE: { + struct camu_sink_entry *target = sink->target; if (sink->target) { al_log_warn("sink", "Ignoring target on NONE."); + if (!camu_clock_is_armed(&target->clock)) { // TMP + camu_clock_pause(&target->clock, at); + } sink->target = NULL; } if (entry != sink->current) { - if (camu_clock_is_paused(&entry->clock)) { // TMP + if (camu_clock_is_armed(&entry->clock)) { // TMP camu_clock_resume(&entry->clock, at); } switch_to(sink, entry); } break; - case LIANA_PAUSE_RESUME: + } + case LIANA_PAUSE_RESUME: { + struct camu_sink_entry *target = sink->target; if (sink->target) { al_log_warn("sink", "Ignoring target on RESUME."); + if (!camu_clock_is_armed(&target->clock)) { // TMP + camu_clock_pause(&target->clock, at); + } sink->target = NULL; } camu_clock_resume(&entry->clock, at); @@ -1010,6 +1140,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn camu_audio_buffer_unpause(&entry->audio.buf); } break; + } case LIANA_PAUSE_PAUSE: { struct camu_sink_entry *prev_target = sink->target; if (prev_target) { @@ -1024,14 +1155,14 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn // The targets clock is only guaranteed to be resumed // if it was set in the PAUSE_BOTH case. I feel like // this needs to be simpler. - if (!camu_clock_is_paused(&prev_target->clock)) { // TMP + if (!camu_clock_is_armed(&prev_target->clock)) { // TMP camu_clock_pause(&prev_target->clock, at); } } else { - if (camu_clock_is_paused(&entry->clock)) { // TMP + if (camu_clock_is_armed(&entry->clock)) { // TMP camu_clock_resume(&entry->clock, at); } - set_target_and_pause(sink, entry, at); + pause_and_swap_to(sink, entry, at); } break; } @@ -1046,11 +1177,11 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn sink->target = entry; camu_audio_buffer_unpause(&entry->audio.buf); } - if (!camu_clock_is_paused(&prev_target->clock)) { // TMP + if (!camu_clock_is_armed(&prev_target->clock)) { // TMP camu_clock_pause(&prev_target->clock, at); } } else { - set_target_and_pause(sink, entry, at); + pause_and_swap_to(sink, entry, at); } break; } @@ -1067,6 +1198,7 @@ out: return false; } +// TODO: pause and seek shouldn't depend on the entry being current. static bool pause_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { @@ -1129,11 +1261,9 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con if (current->sequence == sequence) { current->ended = false; #ifdef CAMU_SINK_LOCAL - (void)at; - lia_client_seek(¤t->client, pos, 0); -#else - lia_client_seek(¤t->client, pos, at); + at = 0; #endif + lia_client_seek(¤t->client, pos, at); } out: @@ -1272,6 +1402,13 @@ void camu_sink_reseek(struct camu_sink *sink) }); } +void camu_sink_unset(struct camu_sink *sink) +{ + queue_cmd(sink, (struct camu_sink_cmd){ + .op = UNSET + }); +} + void camu_sink_offset_volume(struct camu_sink *sink, f64 amount) { camu_mixer_offset_volume(sink->audio.mixer, amount); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index a9a6da0..f544a36 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -112,6 +112,7 @@ void camu_sink_shuffle(struct camu_sink *sink); void camu_sink_toggle_pause(struct camu_sink *sink); void camu_sink_seek(struct camu_sink *sink, f64 pos); void camu_sink_reseek(struct camu_sink *sink); +void camu_sink_unset(struct camu_sink *sink); void camu_sink_offset_volume(struct camu_sink *sink, f64 amount); void camu_sink_stop(struct camu_sink *sink); void camu_sink_close(struct camu_sink *sink); |