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