summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-12-07 23:47:36 -0500
committerAndrew Opalach <andrew@akon.city> 2024-12-07 23:47:36 -0500
commitd81c406f62e996840d86333b81c41d0ae4aff347 (patch)
tree9e6baf260ec56d14e8931b61936b0d043cb49847 /src/libsink
parentd4da3d8a644156b9f24b560eba58a87b912c1fef (diff)
downloadcamu-d81c406f62e996840d86333b81c41d0ae4aff347.tar.gz
camu-d81c406f62e996840d86333b81c41d0ae4aff347.tar.bz2
camu-d81c406f62e996840d86333b81c41d0ae4aff347.zip
Synced list reference 1
This marks a point where the list behavior is at least moderately robust for the skip operation. It includes fixes for multiple deep-rooted issues found by testing with input simulation. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/common.h2
-rw-r--r--src/libsink/sink.c529
-rw-r--r--src/libsink/sink.h1
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(&current->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(&current->client, pos, 0);
-#else
- lia_client_seek(&current->client, pos, at);
+ at = 0;
#endif
+ lia_client_seek(&current->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);