summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-11-30 14:32:51 -0500
committerAndrew Opalach <andrew@akon.city> 2025-11-30 14:32:51 -0500
commitc8412bbedae0fce38db96833732e8ce904721e4c (patch)
tree4611046186c714513d042966ce97825df0fecf9e /src/libsink
parent0d6d13425015d78606232874498327cabcb0e4e2 (diff)
downloadcamu-c8412bbedae0fce38db96833732e8ce904721e4c.tar.gz
camu-c8412bbedae0fce38db96833732e8ce904721e4c.tar.bz2
camu-c8412bbedae0fce38db96833732e8ce904721e4c.zip
Build cleanup and fixes from sink testing
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/desktop.c9
-rw-r--r--src/libsink/input_simulator.c10
-rw-r--r--src/libsink/sink.c342
-rw-r--r--src/libsink/sink.h6
4 files changed, 155 insertions, 212 deletions
diff --git a/src/libsink/desktop.c b/src/libsink/desktop.c
index 8582cd4..af6fc50 100644
--- a/src/libsink/desktop.c
+++ b/src/libsink/desktop.c
@@ -208,14 +208,16 @@ bool camu_desktop_open(struct camu_desktop *c, const char *window_name)
c->scr.callback = screen_callback;
c->scr.userdata = c;
if (!camu_screen_init(&c->scr) || !camu_screen_create_window(&c->scr, window_name)) {
+ log_error("Failed to create window.");
return false;
}
#if defined CAMU_RENDERER_MOMO
c->renderer = camu_renderer_momo_create();
#elif defined CAMU_RENDERER_LIBPLACEBO
- c->renderer = camu_renderer_lp_create();
+ c->renderer = camu_renderer_libplacebo_create();
#endif
if (!camu_screen_create_renderer(&c->scr, c->renderer)) {
+ log_error("Failed to create renderer.");
return false;
}
// Need to render twice for the window to show early.
@@ -252,7 +254,10 @@ bool camu_desktop_tick(struct camu_desktop *c)
{
bool force;
if (camu_screen_tick(&c->scr, &force) && !c->should_quit) {
- c->renderer->render(c->renderer, &c->scr, force);
+ if (!c->renderer->render(c->renderer, &c->scr, force)) {
+ log_error("render() failed, forcing exit.");
+ c->should_quit = 1;
+ }
}
return !c->should_quit;
}
diff --git a/src/libsink/input_simulator.c b/src/libsink/input_simulator.c
index 19786af..5b9d6e6 100644
--- a/src/libsink/input_simulator.c
+++ b/src/libsink/input_simulator.c
@@ -31,23 +31,23 @@ static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata)
switch (al_random_int(0, MARK - 1)) {
case SKIP: {
s32 n = al_random_int(1, 5);
- log_debug("SKIP (n: %d).", n);
+ log_info("SKIP (n: %d).", n);
camu_sink_skip(sink, n);
break;
}
case BACKSKIP: {
s32 n = al_random_int(-5, -1);
- log_debug("BACKSKIP (n: %d).", n);
+ log_info("BACKSKIP (n: %d).", n);
camu_sink_skip(sink, n);
break;
}
case SHUFFLE: {
- log_debug("SHUFFLE.");
+ log_info("SHUFFLE.");
camu_sink_shuffle(sink);
break;
}
case TOGGLE_PAUSE: {
- log_debug("TOGGLE_PAUSE.");
+ log_info("TOGGLE_PAUSE.");
camu_sink_toggle_pause(sink);
break;
}
@@ -56,7 +56,7 @@ static nn_thread_result NNWT_THREADCALL input_simulation_thread(void *userdata)
if (pos < 0.005) pos = 0.0;
if (pos > 0.995) pos = 1.0;
else if (pos > 0.99) pos = 0.9999;
- log_debug("SEEK (pos: %f).", pos);
+ log_info("SEEK (pos: %f).", pos);
camu_sink_seek(sink, &pos, CAMU_SEEK_PERCENT);
break;
}
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index e9b4df4..fccbc3e 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -10,7 +10,7 @@
#include "sink.h"
#include "common.h"
-#define CAMU_SINK_LOCAL
+//#define CAMU_SINK_LOCAL
//#define CAMU_SINK_ONESHOT
// Requested state of the sinks outputs.
@@ -71,8 +71,7 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX);
#define CONNECTION_NUMBER(id) ((id >> 48) & 0xffff)
// Only sink->target can ever be 0xb00b.
-#define OxbOOb ((struct camu_sink_entry *)0xb00b)
-#define ENTRY_IS_VALID(entry) ((entry) && (entry) != OxbOOb)
+#define ENTRY_IS_VALID(entry) ((entry) && (entry) != (struct camu_sink_entry *)0xb00b)
// printf format for entries.
#ifdef AL_DEBUG
@@ -87,14 +86,19 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX);
#define VIDEO_STATE(entry) ((entry)->video.state)
// If a buffer is still INIT or QUEUED after the entry is configured, it's "empty".
-// Also, if a buffer errors, it will be detached in CLIENT_REMOVE_BUFFERS.
-// Empty is a state that cannot change while a buffer is being used (push()/read()).
+// An ERRORED buffer will be DETACHED after CLIENT_REMOVE_BUFFERS and is then considered empty.
+// Empty is a state that cannot change while a buffer could be in use (push()/read()).
#define AUDIO_EMPTY(entry) (AUDIO_STATE(entry) <= BUFFER_DETACHED)
#define VIDEO_EMPTY(entry) (VIDEO_STATE(entry) <= BUFFER_DETACHED)
#define AUDIO_ENDED(entry) (AUDIO_STATE(entry) >= BUFFER_ENDED)
#define VIDEO_ENDED(entry) (VIDEO_STATE(entry) >= BUFFER_ENDED)
-#define ENTRY_ENDED(entry) ((AUDIO_EMPTY(entry) || AUDIO_ENDED(entry)) && (VIDEO_EMPTY(entry) || VIDEO_ENDED(entry)))
+/*
+#define ENTRY_ENDED(entry) \
+ ((AUDIO_ENDED(entry) && (VIDEO_ENDED(entry) || VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry))) || \
+ (AUDIO_EMPTY(entry) && VIDEO_ENDED(entry)))
+*/
+#define ENTRY_ENDED(entry) entry->ended
// IGNORED = ENDED or EMPTY.
#define AUDIO_ADDED_OR_IGNORED(entry) (AUDIO_STATE(entry) >= BUFFER_ADDED || AUDIO_EMPTY(entry))
@@ -113,6 +117,8 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX);
#define BLOCKING_SLEEP(delay) nn_event_loop_sleep(sink->loop, delay)
#endif
+#define CMD(cmd, ...) ((struct camu_sink_cmd){ .op = cmd, __VA_ARGS__ })
+
// Functions that might happen on separate threads.
// CAMU_MIXER_THREADED:
// audio_buffer_callback()
@@ -127,7 +133,7 @@ AL_STATIC_ASSERT(max_age_lt_lru, ENTRY_MAX_AGE, <, SINK_LRU_MAX);
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;
+ return atomic_load(u8)(&entry->audio.buf.ref, AL_ATOMIC_RELAXED) == 1;
#else
(void)entry;
return false;
@@ -137,19 +143,30 @@ static inline bool entry_audio_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;
+ return atomic_load(u8)(&entry->video.buf.ref, AL_ATOMIC_RELAXED) == 1;
#else
(void)entry;
return false;
#endif
}
-static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
+static void queue_cmds(struct camu_sink *sink, u32 count, ...)
{
- camu_queue_push(sink->queue, cmd);
+ camu_queue_lock(sink->queue);
+ va_list cmds;
+ va_start(cmds, count);
+ for (u32 i = 0; i < count; i++) {
+ camu_queue_push(sink->queue, va_arg(cmds, struct camu_sink_cmd));
+ }
+ camu_queue_unlock(sink->queue);
nn_signal_send(&sink->queue_signal);
}
+static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
+{
+ queue_cmds(sink, 1, cmd);
+}
+
static void refresh_video_output(struct camu_sink *sink)
{
#ifndef CAMU_SINK_NO_VIDEO
@@ -165,11 +182,7 @@ static inline void add_entry_audio_buffer(struct camu_sink_entry *entry)
#ifdef CAMU_MIXER_THREADED_START_STOP
sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
#else
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = ADD_BUFFER,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry
- });
+ queue_cmd(sink, CMD(ADD_BUFFER, .v.u = CAMU_SINK_AUDIO, .opaque = entry));
#endif
}
@@ -194,11 +207,7 @@ static void remove_entry_audio_buffer(struct camu_sink_entry *entry)
#ifdef CAMU_MIXER_THREADED_START_STOP
sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
#else
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = REMOVE_BUFFER,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry
- });
+ queue_cmd(sink, CMD(REMOVE_BUFFER, .v.u = CAMU_SINK_AUDIO, .opaque = entry));
#endif
break;
case BUFFER_SET_OR_BUFFERED:
@@ -243,7 +252,6 @@ static void remove_entry_video_buffer(struct camu_sink_entry *entry)
static void remove_entry_buffers(struct camu_sink_entry *entry)
{
log_trace("remove_entry_buffers("ENTRY_FMT"), audio_state: %hhu, video_state: %hhu.", ENTRY_ARG(entry), AUDIO_STATE(entry), VIDEO_STATE(entry));
- al_assert(!entry->ended);
if (!AUDIO_ENDED(entry)) remove_entry_audio_buffer(entry);
if (!VIDEO_ENDED(entry)) remove_entry_video_buffer(entry);
}
@@ -256,14 +264,15 @@ static void add_or_queue_entry(struct camu_sink_entry *entry)
log_trace("add_or_queue_entry("ENTRY_FMT"), audio_state: %hhu, video_state: %hhu.", ENTRY_ARG(entry), AUDIO_STATE(entry), VIDEO_STATE(entry));
al_assert(AUDIO_STATE(entry) != BUFFER_QUEUED);
al_assert(VIDEO_STATE(entry) != BUFFER_QUEUED);
+ // Either buffer could be DETACHED.
if (AUDIO_STATE(entry) == BUFFER_INIT) {
AUDIO_STATE(entry) = BUFFER_QUEUED;
- } else if (!AUDIO_ENDED(entry)) {
+ } else if (!AUDIO_ENDED_OR_EMPTY(entry)) {
add_audio_if_set_and_buffered(entry);
}
if (VIDEO_STATE(entry) == BUFFER_INIT) {
VIDEO_STATE(entry) = BUFFER_QUEUED;
- } else if (!VIDEO_ENDED(entry)) {
+ } else if (!VIDEO_ENDED_OR_EMPTY(entry)) {
add_video_if_set_and_buffered(entry);
}
}
@@ -279,14 +288,14 @@ static void maybe_disconnect_entry(struct camu_sink_entry *entry)
}
#ifdef CAMU_SINK_LOCAL
-static void local_pause(struct camu_sink *sink, struct camu_sink_entry *entry)
+static void local_entry_pause(struct camu_sink *sink, struct camu_sink_entry *entry)
{
if (!camu_clock_is_paused(&entry->clock)) {
entry->paused = true;
camu_clock_pause(&entry->clock, 0);
log_info("Clock paused.");
// Audio stop will be handled by a BUFFER_PAUSED callback.
- if (!VIDEO_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) {
+ if (!VIDEO_ENDED_OR_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) {
#ifndef CAMU_SINK_NO_VIDEO
sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
#endif
@@ -296,13 +305,13 @@ static void local_pause(struct camu_sink *sink, struct camu_sink_entry *entry)
entry->paused = false;
camu_clock_resume(&entry->clock, 0);
log_info("Clock resumed.");
- if (!VIDEO_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) {
+ if (!VIDEO_ENDED_OR_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) {
#ifndef CAMU_SINK_NO_VIDEO
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
#endif
sink->video.state = SINK_PLAYING;
}
- if (!AUDIO_EMPTY(entry) && sink->audio.state == SINK_PAUSED) {
+ if (!AUDIO_ENDED_OR_EMPTY(entry) && sink->audio.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
sink->audio.state = SINK_PLAYING;
}
@@ -310,6 +319,12 @@ static void local_pause(struct camu_sink *sink, struct camu_sink_entry *entry)
}
#endif
+// sink->current could be NULL.
+static inline struct camu_sink_entry *get_entry_for_command(struct camu_sink *sink)
+{
+ return ENTRY_IS_VALID(sink->target) ? sink->target : sink->current;
+}
+
static inline s32 get_sequence_for_command(struct camu_sink_entry *entry)
{
// SEQUENCE_ANY resolves order on the server.
@@ -329,7 +344,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
{
switch (cmd->op) {
case START: {
- switch (cmd->value.i) {
+ switch (cmd->v.u) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
@@ -348,7 +363,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case STOP: {
- switch (cmd->value.i) {
+ switch (cmd->v.u) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PLAYING) {
sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_AUDIO, NULL);
@@ -368,7 +383,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
case ADD_BUFFER: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
+ switch (cmd->v.u) {
case CAMU_SINK_AUDIO:
sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
break;
@@ -382,7 +397,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
case REMOVE_BUFFER: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- switch (cmd->value.i) {
+ switch (cmd->v.u) {
case CAMU_SINK_AUDIO:
sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
break;
@@ -395,7 +410,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case CLEAR_BUFFERS: {
- switch (cmd->value.i) {
+ switch (cmd->v.u) {
case CAMU_SINK_AUDIO:
sink->callback(sink->userdata, CAMU_SINK_CLEAR, CAMU_SINK_AUDIO, NULL);
break;
@@ -424,21 +439,21 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
nn_packet_write_str(packet, &sink->default_list);
nn_packet_write_u8(packet, CAMU_LIST_SKIP);
nn_packet_write_s32(packet, get_sequence_for_command(entry));
- nn_packet_write_s32(packet, (s32)cmd->value.i);
+ nn_packet_write_s32(packet, (s32)cmd->v.i);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
case TOGGLE_PAUSE: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
#ifdef CAMU_SINK_LOCAL
- if (!entry->held) local_pause(sink, entry);
+ if (!entry->held) local_entry_pause(sink, entry);
#else
if (!sink->conn) return;
struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
nn_packet_write_str(packet, &sink->default_list);
nn_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE);
nn_packet_write_s32(packet, get_sequence_for_command(entry));
- nn_packet_write_f64(packet, cmd->value.f);
+ nn_packet_write_f64(packet, cmd->v.f);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
#endif
break;
@@ -451,7 +466,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
nn_packet_write_u8(packet, CAMU_LIST_SEEK);
nn_packet_write_s32(packet, entry->sequence);
nn_packet_write_u32(packet, REMOTE_ENTRY_ID(entry->id));
- nn_packet_write_u64(packet, cmd->value.u);
+ nn_packet_write_u64(packet, cmd->v.u);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -476,7 +491,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
nn_packet_write_str(packet, &sink->default_list);
nn_packet_write_u8(packet, CAMU_LIST_END);
nn_packet_write_u32(packet, REMOTE_ENTRY_ID(entry->id));
- nn_packet_write_u32(packet, (u32)cmd->value.u);
+ nn_packet_write_u32(packet, (u32)cmd->v.u);
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
@@ -486,11 +501,11 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
static void queue_signal_callback(void *userdata)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
- u32 size;
+ u32 count;
struct camu_sink_cmd cmd;
for (;;) {
- camu_queue_try_pop(sink->queue, size, cmd);
- if (size == 0) break;
+ camu_queue_try_pop(sink->queue, count, cmd);
+ if (count == 0) break;
handle_sink_cmd(sink, &cmd);
}
}
@@ -498,7 +513,8 @@ static void queue_signal_callback(void *userdata)
static void mixer_callback(void *userdata, u8 op)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
- // @TODO: Isn't the idea of MIXER_EMPTY to not STOP audio when we know the mixer will be empty?
+ // @TODO: Isn't the idea of MIXER_EMPTY to not explicitly STOP the audio in places we
+ // expect the mixer to be empty (rely on MIXER_EMPTY).
if (op == CAMU_MIXER_EMPTY) {
log_info("Mixer empty.");
// Regardless of if we are checking an entry's state here, we have to sync with
@@ -506,10 +522,7 @@ static void mixer_callback(void *userdata, u8 op)
nn_mutex_lock(&sink->lock);
// This feels a bit too loose.
if (!(sink->current && AUDIO_STATE(sink->current) == BUFFER_ADDED)) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_AUDIO));
}
nn_mutex_unlock(&sink->lock);
}
@@ -568,10 +581,7 @@ static void maybe_remove_previous(struct camu_sink *sink)
// Cleanup entries from old connections. This is especially important to
// keep reseek() from being overly wasteful.
if (CONNECTION_NUMBER(previous->id) != sink->connection_number) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = EJECT_ENTRY,
- .opaque = previous
- });
+ queue_cmd(sink, CMD(EJECT_ENTRY, .opaque = previous));
}
}
sink->previous.count = 0;
@@ -600,10 +610,9 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry
{
log_trace("maybe_add_to_previous("ENTRY_FMT", "ENTRY_FMT").", ENTRY_ARG(previous), ENTRY_ARG(target));
al_assert(previous != target);
- al_assert(!previous->ended);
- // If none of the entry's buffers are added, we don't care about adding it to previous.
- bool dangling_target = target == OxbOOb;
- if (dangling_target || target->ended || (AUDIO_STATE(previous) != BUFFER_ADDED && VIDEO_STATE(previous) != BUFFER_ADDED)) {
+ // If none of an entry's buffers are added, we don't care about adding it to previous.
+ bool dangling_target = target == (struct camu_sink_entry *)0xb00b;
+ if (dangling_target || ENTRY_ENDED(target) || (AUDIO_STATE(previous) != BUFFER_ADDED && VIDEO_STATE(previous) != BUFFER_ADDED)) {
remove_entry_buffers(previous);
return;
}
@@ -613,18 +622,14 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry
static void after_add_entry(struct camu_sink_entry *entry, bool skip_audio, bool skip_video)
{
maybe_remove_previous(entry->sink);
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = (skip_video || entry->paused) ? STOP : START,
- .value.i = CAMU_SINK_VIDEO
- });
+ queue_cmds(entry->sink, 2,
+ CMD((skip_video || entry->paused) ? STOP : START, .v.u = CAMU_SINK_VIDEO),
+ CMD((skip_audio || entry->paused) ? STOP : START, .v.u = CAMU_SINK_AUDIO)
+ );
if (VIDEO_ENDED_OR_EMPTY(entry)) {
// Clear the screen if skipping from a video to an audio-only entry.
refresh_video_output(entry->sink);
}
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = (skip_audio || entry->paused) ? STOP : START,
- .value.i = CAMU_SINK_AUDIO
- });
}
// Call this after setting state to ADDED because this entry might be in previous.
@@ -641,7 +646,7 @@ static void do_add_entry(struct camu_sink_entry *entry)
void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
- al_assert(!entry->ended);
+ al_assert(!ENTRY_ENDED(entry));
al_assert(AUDIO_STATE(entry) != BUFFER_INIT);
al_assert(AUDIO_STATE(entry) != BUFFER_QUEUED);
al_assert(AUDIO_STATE(entry) != BUFFER_ADDED);
@@ -662,7 +667,7 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
// Single frame entries will be added/removed with ended set.
- if (entry->ended) al_assert(VIDEO_IS_SINGLE_FRAME(entry));
+ if (ENTRY_ENDED(entry)) al_assert(VIDEO_IS_SINGLE_FRAME(entry));
al_assert(VIDEO_STATE(entry) != BUFFER_INIT);
al_assert(VIDEO_STATE(entry) != BUFFER_QUEUED);
al_assert(VIDEO_STATE(entry) != BUFFER_ADDED);
@@ -690,56 +695,34 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
{
struct camu_sink_entry *current = sink->current;
- log_trace("switch_to("ENTRY_FMT"(ended: %s)), current: "ENTRY_FMT"(ended: %s).",
- ENTRY_ARG(target), BOOLSTR(ENTRY_IS_VALID(target) ? target->ended : false),
- ENTRY_ARG(current), BOOLSTR(current ? current->ended : false));
+ log_trace("switch_to("ENTRY_FMT"), current: "ENTRY_FMT".", ENTRY_ARG(target), ENTRY_ARG(current));
- bool ensure_removed = false;
if (current) {
struct camu_sink_entry *suspended = sink->suspended;
al_assert(current != target);
- ensure_removed = suspended || current->ended;
if (suspended) {
al_assert(suspended == current);
- // A buffer's state being QUEUED should be impossible while it's
- // entry is suspended.
+ // It should be impossible for a buffer to be QUEUED while it's entry is suspended.
al_assert(AUDIO_STATE(suspended) != BUFFER_QUEUED);
al_assert(VIDEO_STATE(suspended) != BUFFER_QUEUED);
sink->suspended = NULL;
log_warn("Unset suspended entry as a substitute for remove.");
- } else if (!current->ended) {
+ } else {
maybe_add_to_previous(sink, current, target);
}
}
- // @TODO: Cleanup stop_video? and log_trace lengths.
- bool stop_video = false;
- bool dangling_target = target == OxbOOb;
+ bool dangling_target = target == (struct camu_sink_entry *)0xb00b;
+ bool stop_video = dangling_target;
if (dangling_target) {
log_trace("Ignored dangling target.");
} else {
- if (!target->ended) {
- remove_previous_if_contains(sink, target);
- add_or_queue_entry(target);
- } else {
- log_trace("Target ended in switch_to().");
- if (!VIDEO_ENDED_OR_EMPTY(target) && VIDEO_IS_SINGLE_FRAME(target)) {
- add_video_if_set_and_buffered(target);
- }
- stop_video = true;
- }
- }
-
- if (ensure_removed) {
- if (!VIDEO_ENDED_OR_EMPTY(current) && VIDEO_IS_SINGLE_FRAME(current)) {
- remove_entry_video_buffer(current);
- }
- al_assert(AUDIO_STATE(current) != BUFFER_ADDED);
- al_assert(VIDEO_STATE(current) != BUFFER_ADDED);
+ remove_previous_if_contains(sink, target);
+ add_or_queue_entry(target);
+ stop_video = VIDEO_ENDED_OR_EMPTY(target) || VIDEO_IS_SINGLE_FRAME(target);
}
if (dangling_target) {
- stop_video = true;
sink->current = NULL;
} else {
target->audio.ignore_paused = false;
@@ -750,10 +733,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
}
if (stop_video) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_VIDEO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO));
refresh_video_output(sink);
}
}
@@ -781,17 +761,19 @@ static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *ta
al_assert(target != current);
al_assert(!sink->target);
// This is extra verbose because the order is important.
- // 1. sink->target has to be set before calling clock_pause().
- // 2. current must still be paused even if it's ended.
- // 3. In the immediate case, switch_to() has to come last.
- if (current && !current->ended) {
+ // 1. sink->target has to be set before calling clock_pause().
+ if (current && !ENTRY_ENDED(current)) {
sink->target = target;
}
+ // 2. current must still be paused, even if it's ended.
if (current) {
current->audio.ignore_paused = true;
camu_clock_pause(&current->clock, at);
}
- if (!current || current->ended) switch_to(sink, target);
+ // 3. In the immediate swap case, switch_to() has to come last.
+ if (!current || ENTRY_ENDED(current)) {
+ switch_to(sink, target);
+ }
}
#endif
@@ -805,16 +787,14 @@ static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink
#endif
// This entry's buffers cannot be added again until after a reset.
entry->ended = true;
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = END,
- .value.u = entry->reset_token,
- .opaque = entry
- });
+ al_assert(ENTRY_ENDED(entry));
+ queue_cmd(sink, CMD(END, .v.u = entry->reset_token, .opaque = entry));
#ifdef LIANA_LIST_SCUFFED_LOOP
log_info("Looping.");
return true;
#endif
if (sink->target) {
+ log_info("Buffers swapped on end() (Gapless if queued).");
switch_to(sink, sink->target);
sink->target = NULL;
return true;
@@ -843,17 +823,18 @@ static void audio_buffer_callback(void *userdata, u8 op)
nn_mutex_lock(&sink->lock);
if (!entry->audio.ignore_paused && entry->paused) {
log_info("Audio buffer paused.");
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_AUDIO));
}
nn_mutex_unlock(&sink->lock);
break;
case CAMU_BUFFER_EOF:
case CAMU_BUFFER_ERRORED: {
bool error = op == CAMU_BUFFER_ERRORED;
- log_debug(error ? "Audio buffer errored." : "Audio EOF.");
+ if (error) {
+ log_error("Audio buffer errored.");
+ } else {
+ log_debug("Audio EOF.");
+ }
nn_mutex_lock(&sink->lock);
// EOF and ERRORED come from the outputs read() thread. So, having threaded
// outputs means anything could have happened while waiting on the lock above.
@@ -894,7 +875,11 @@ static void video_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_EOF:
case CAMU_BUFFER_ERRORED: {
bool error = op == CAMU_BUFFER_ERRORED;
- log_debug(error ? "Video buffer errored." : "Video EOF.");
+ if (error) {
+ log_error("Video buffer errored.");
+ } else {
+ log_debug("Video EOF.");
+ }
nn_mutex_lock(&sink->lock);
if (!AUDIO_EMPTY(entry)) {
camu_audio_buffer_set_no_video(&entry->audio.buf, true);
@@ -915,10 +900,7 @@ static void video_buffer_callback(void *userdata, u8 op)
}
nn_mutex_unlock(&sink->lock);
if (!swapped) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_VIDEO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO));
}
break;
}
@@ -937,10 +919,7 @@ static void clock_callback(void *userdata, u8 op)
switch_to(sink, sink->target);
sink->target = NULL;
} else if (entry->paused) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_VIDEO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO));
}
}
nn_mutex_unlock(&sink->lock);
@@ -953,20 +932,11 @@ static void evaluate_and_set_buffer_params(struct camu_sink *sink, struct camu_s
f64 avg_frame_duration = entry->video.buf.avg_frame_duration;
#ifdef CAMU_SINK_LOCAL
if (!AUDIO_EMPTY(entry) && !ignore_video) {
- // Delay either the audio or video so we can start the clock as
- // soon as possible while keeping A/V sync.
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
struct camu_renderer *renderer = sink->video.renderer;
f64 video = renderer->get_latency(renderer) * avg_frame_duration;
- // If the audio buffer is delayed by less than the audio latency, data will be skipped.
- f64 base = -audio;
- if (audio > video) {
- camu_audio_buffer_set_latency(&entry->audio.buf, base);
- camu_video_buffer_set_latency(&entry->video.buf, base - (audio - video));
- } else if (video > audio) {
- camu_audio_buffer_set_latency(&entry->audio.buf, base);
- camu_video_buffer_set_latency(&entry->video.buf, base + (video - audio));
- }
+ camu_audio_buffer_set_latency(&entry->audio.buf, -audio);
+ camu_video_buffer_set_latency(&entry->video.buf, -video);
}
// If we're local we don't have to worry about syncing audio-only entries.
camu_audio_buffer_set_ignore_desync(&entry->audio.buf, ignore_video);
@@ -1013,6 +983,7 @@ static void remove_from_queue_by_opaque(struct camu_sink *sink, void *opaque)
static void client_callback(void *userdata, u8 op, struct camu_codec_stream *stream, void *opaque)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ al_assert(ENTRY_IS_VALID(entry));
struct camu_sink *sink = entry->sink;
switch (op) {
case LIANA_CLIENT_CONFIGURE: {
@@ -1106,13 +1077,8 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
log_trace("remove_buffers("ENTRY_FMT", %s, %s), entry == current: %s.", ENTRY_ARG(entry), BOOLSTR(rec->reconnect), BOOLSTR(rec->unconfigured), BOOLSTR(entry == sink->current));
if (entry == sink->current) {
- if (!entry->ended) {
- // It should only be possible for a buffer to be INIT if entry is not current or ended.
- // Ended case is: (audio or video buffer empty) -> entry switched off of -> entry ended -> switch_to()'d.
- // - The empty buffer doesn't get re-QUEUED because the entry is ended.
- al_assert(AUDIO_STATE(entry) != BUFFER_INIT);
- al_assert(VIDEO_STATE(entry) != BUFFER_INIT);
- }
+ al_assert(AUDIO_STATE(entry) != BUFFER_INIT);
+ al_assert(VIDEO_STATE(entry) != BUFFER_INIT);
if (sink->target) {
// This is necessary to avoid re-adding an entry with an in-between clock state.
// See note in clock.c::camu_clock_seek().
@@ -1140,7 +1106,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
// Ignore unconfigured entries.
if (rec->unconfigured) {
- al_assert(!entry->ended);
al_assert(AUDIO_EMPTY(entry) && VIDEO_EMPTY(entry));
nn_mutex_unlock(&sink->lock);
return;
@@ -1309,7 +1274,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
remove_from_queue_by_opaque(sink, entry);
if (entry == sink->target) {
- sink->target = OxbOOb;
+ sink->target = (struct camu_sink_entry *)0xb00b;
log_warn("Attempting to handle a disconnected target.");
} else if (entry == sink->current) {
if (sink->suspended) {
@@ -1322,10 +1287,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
} else {
sink->current = NULL;
if (removed) { // Don't stop video on exit.
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_VIDEO
- });
+ queue_cmd(sink, CMD(STOP, .v.u = CAMU_SINK_VIDEO));
}
}
// If current was never fully added we need to call this here.
@@ -1401,9 +1363,12 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
u8 op = nn_packet_read_u8(packet);
if (op == LIANA_SINK_UNSET) {
nn_mutex_lock(&sink->lock);
- if (sink->current) {
- maybe_disconnect_entry(sink->current);
+ struct camu_sink_entry *current = sink->current;
+ nn_mutex_unlock(&sink->lock);
+ if (current) {
+ maybe_disconnect_entry(current);
}
+ nn_mutex_lock(&sink->lock);
goto out;
}
@@ -1499,7 +1464,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
case LIANA_PAUSE_PAUSE:
if (prev_target) {
al_assert(current);
- al_assert(!current->ended);
+ al_assert(!ENTRY_ENDED(current));
al_assert(prev_target != entry);
if (entry == current) {
// pause_and_swap_to() negated.
@@ -1507,7 +1472,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
} else {
sink->target = entry;
}
- if (prev_target != OxbOOb) {
+ if (prev_target != (struct camu_sink_entry *)0xb00b) {
camu_clock_pause(&prev_target->clock, at);
}
} else {
@@ -1518,7 +1483,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
case LIANA_PAUSE_BOTH:
if (prev_target) {
al_assert(current);
- al_assert(!current->ended);
+ al_assert(!ENTRY_ENDED(current));
al_assert(prev_target != entry);
if (entry == current) {
sink->target = NULL;
@@ -1526,7 +1491,7 @@ static bool set_command_callback(void *userdata, struct nn_rpc_connection *conn,
} else {
sink->target = entry;
}
- if (prev_target != OxbOOb) {
+ if (prev_target != (struct camu_sink_entry *)0xb00b) {
camu_clock_pause(&prev_target->clock, at);
}
} else {
@@ -1562,12 +1527,14 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
if (!entry) goto out;
- al_assert(entry->sequence == sequence);
nn_mutex_lock(&sink->lock);
- log_trace("pause("ENTRY_FMT"), %s, audio_state: %hhu, video_state: %hhu.", ENTRY_ARG(entry), lia_pause_op_name(pause), AUDIO_STATE(entry), VIDEO_STATE(entry));
+ // As long as the list discards skips with a non-current sequence, this should hold true.
+ al_assert(entry->sequence == sequence);
+ log_trace("pause("ENTRY_FMT"), %s, audio_state: %hhu, video_state: %hhu.",
+ ENTRY_ARG(entry), lia_pause_op_name(pause), AUDIO_STATE(entry), VIDEO_STATE(entry));
#ifdef CAMU_SINK_LOCAL
(void)at;
- if (!entry->held) local_pause(sink, entry);
+ if (!entry->held) local_entry_pause(sink, entry);
#else
switch (pause) {
case LIANA_PAUSE_PAUSE: {
@@ -1582,18 +1549,12 @@ static bool pause_command_callback(void *userdata, struct nn_rpc_connection *con
entry->paused = false;
camu_clock_resume(&entry->clock, at);
log_info("Clock resumed.");
- if (!VIDEO_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry)) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_VIDEO
- });
+ if (!VIDEO_ENDED_OR_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry)) {
+ queue_cmd(sink, CMD(START, .v.u = CAMU_SINK_VIDEO));
}
- if (!AUDIO_EMPTY(entry)) {
+ if (!AUDIO_ENDED_OR_EMPTY(entry)) {
camu_audio_buffer_resync(&entry->audio.buf);
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_AUDIO
- });
+ queue_cmd(sink, CMD(START, .v.u = CAMU_SINK_AUDIO));
}
break;
}
@@ -1615,7 +1576,6 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn
u64 id = LOCAL_ENTRY_ID(sink, nn_packet_read_u32(packet));
s32 sequence = nn_packet_read_s32(packet);
- (void)sequence;
u64 at = nn_packet_read_u64(packet);
u64 pos = nn_packet_read_u64(packet);
u32 reset_token = nn_packet_read_u32(packet);
@@ -1623,7 +1583,8 @@ static bool seek_command_callback(void *userdata, struct nn_rpc_connection *conn
struct camu_sink_entry *entry = get_entry_from_id(sink, id);
if (!entry) goto out;
nn_mutex_lock(&sink->lock);
- log_trace("seek("ENTRY_FMT"), reset_token: %u.", ENTRY_ARG(entry), reset_token);
+ entry->sequence = sequence;
+ log_trace("seek("ENTRY_FMT", %.2f), reset_token: %u.", ENTRY_ARG(entry), pos / 1000000.0, reset_token);
entry->reset_token = reset_token;
nn_mutex_unlock(&sink->lock);
// The rest of the seek is handled in CLIENT_REMOVE_BUFFERS/RESUME_AT/RECONNECTED.
@@ -1761,22 +1722,12 @@ void camu_sink_return_current(struct camu_sink *sink)
nn_mutex_unlock(&sink->lock);
}
-// sink->current could be NULL.
-static inline struct camu_sink_entry *get_entry_for_command(struct camu_sink *sink)
-{
- return ENTRY_IS_VALID(sink->target) ? sink->target : sink->current;
-}
-
void camu_sink_skip(struct camu_sink *sink, s32 n)
{
nn_mutex_lock(&sink->lock);
struct camu_sink_entry *current = get_entry_for_command(sink);
nn_mutex_unlock(&sink->lock);
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SKIP,
- .value.i = n,
- .opaque = current
- });
+ queue_cmd(sink, CMD(SKIP, .v.i = n, .opaque = current));
}
void camu_sink_toggle_pause(struct camu_sink *sink)
@@ -1785,11 +1736,8 @@ void camu_sink_toggle_pause(struct camu_sink *sink)
struct camu_sink_entry *current = get_entry_for_command(sink);
nn_mutex_unlock(&sink->lock);
if (!current) return;
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = TOGGLE_PAUSE,
- .value.f = camu_clock_get_pts(&current->clock, 0.0, false),
- .opaque = current
- });
+ f64 pts = camu_clock_get_pts(&current->clock, 0.0, false);
+ queue_cmd(sink, CMD(TOGGLE_PAUSE, .v.f = pts, .opaque = current));
}
void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode)
@@ -1811,19 +1759,19 @@ void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode)
switch (mode) {
case CAMU_SEEK_POS: {
u64 pos = *(u64 *)value;
- cmd.value.u = pos;
+ cmd.v.u = pos;
break;
}
case CAMU_SEEK_RELATIVE: {
f64 offset = *(f64 *)value;
pts = MAX(pts + offset, 0.0);
- cmd.value.u = (u64)(pts * 1000000);
+ cmd.v.u = (u64)(pts * 1000000);
break;
}
case CAMU_SEEK_PERCENT: {
// There's probably a way to lose less precision here.
f64 percent = *(f64 *)value;
- cmd.value.u = (u64)(duration * percent);
+ cmd.v.u = (u64)(duration * percent);
break;
}
}
@@ -1832,31 +1780,21 @@ void camu_sink_seek(struct camu_sink *sink, void *value, u8 mode)
void camu_sink_reseek(struct camu_sink *sink)
{
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = RESEEK
- });
+ queue_cmd(sink, CMD(RESEEK));
}
void camu_sink_shuffle(struct camu_sink *sink)
{
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SHUFFLE
- });
+ queue_cmd(sink, CMD(SHUFFLE));
}
void camu_sink_stop(struct camu_sink *sink)
{
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = STOP,
- .value.i = CAMU_SINK_AUDIO
- });
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = CLEAR_BUFFERS,
- .value.i = CAMU_SINK_AUDIO
- });
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = CLOSE
- });
+ queue_cmds(sink, 3,
+ CMD(STOP, .v.u = CAMU_SINK_AUDIO),
+ CMD(CLEAR_BUFFERS, .v.u = CAMU_SINK_AUDIO),
+ CMD(CLOSE)
+ );
}
void camu_sink_close(struct camu_sink *sink)
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 5958829..1e8bbf1 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -56,7 +56,7 @@ struct camu_sink_entry {
bool paused;
struct {
u8 state;
- // Don't stop the audio output on BUFFER_PAUSED from this entry.
+ // Don't stop audio on a BUFFER_PAUSED from this entry.
bool ignore_paused;
struct camu_audio_buffer buf;
struct lia_vcr_track *track;
@@ -71,7 +71,7 @@ struct camu_sink_entry {
struct camu_sink_cmd {
u8 op;
- union { s64 i; u64 u; f64 f; } value;
+ union { s64 i; u64 u; f64 f; } v;
void *opaque;
};
@@ -90,7 +90,7 @@ struct camu_sink {
queue(struct camu_sink_cmd) queue;
struct camu_sink_entry *current;
struct camu_sink_entry *target;
- // If an entry that was current was removed for a reconnect, it's "suspended".
+ // If current was removed for a reconnect, it's "suspended".
struct camu_sink_entry *suspended;
array(struct camu_sink_entry *) previous;
array(struct camu_sink_entry *) entries;