summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-11-08 14:53:40 -0500
committerAndrew Opalach <andrew@akon.city> 2024-11-08 14:53:40 -0500
commit5e3641e5e692c3f2f644a4bb809c88727cb8bee9 (patch)
treefbac2007fdcebda13406350fdcef977910427078 /src/libsink
parentf56abfafcd4fa722b807278b138f805112cd953e (diff)
downloadcamu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.tar.gz
camu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.tar.bz2
camu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.zip
Command queue for portal and list, work on server
Most of the server stuff can undoubtedly be simplified. I'm still working that out. Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/sink.c104
-rw-r--r--src/libsink/sink.h8
2 files changed, 81 insertions, 31 deletions
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 2c8790b..899023d 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -91,6 +91,8 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e
entry->audio.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) {
entry->audio.state = BUFFER_CONFIGURED;
+ } else if (entry->audio.state == BUFFER_QUEUED) {
+ entry->audio.state = BUFFER_INIT;
}
}
@@ -102,6 +104,9 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e
entry->video.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->video.state == BUFFER_SET_OR_BUFFERED) {
entry->video.state = BUFFER_CONFIGURED;
+ } else if (entry->video.state == BUFFER_QUEUED) {
+ // This can be hit when skipping through entries very fast.
+ entry->video.state = BUFFER_INIT;
}
}
#endif
@@ -276,7 +281,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case CLOSE: {
- aki_signal_stop(&sink->signal);
+ aki_signal_stop(&sink->queue_signal);
sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
return;
}
@@ -298,7 +303,7 @@ static void queue_signal_callback(void *userdata)
static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
{
camu_queue_push(sink->queue, cmd);
- aki_signal_send(&sink->signal);
+ aki_signal_send(&sink->queue_signal);
}
static void maybe_remove_previous(struct camu_sink *sink)
@@ -310,13 +315,32 @@ 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 *entry, struct camu_sink_entry *current)
+{
+ al_assert(entry != current);
+ struct camu_sink_entry *rentry;
+ al_array_foreach_rev(sink->previous, i, rentry) {
+ if (rentry == current) {
+ // If the entry we are about to add is in previous,
+ // remove it immediately.
+ remove_entry_buffers(sink, current);
+ al_array_remove_at(sink->previous, i);
+ }
+ }
+ // Don't accept duplicates.
+ al_array_foreach_rev(sink->previous, i, rentry) {
+ if (rentry == entry) return;
+ }
+ al_array_push(sink->previous, entry);
+}
+
void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
- u8 state = entry->audio.state;
- al_assert(state != BUFFER_ADDED);
- if (state == BUFFER_CONFIGURED) {
- state = BUFFER_SET_OR_BUFFERED;
- } else if (state == BUFFER_SET_OR_BUFFERED) {
+ al_assert(entry->audio.state != BUFFER_ADDED);
+ 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.
@@ -324,24 +348,24 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
if (VIDEO_READY_OR_EMPTY(entry)) {
maybe_remove_previous(entry->sink);
}
+#ifndef CAMU_SINK_LOCAL
camu_audio_buffer_unpause(&entry->audio.buf);
+#endif
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
.value.i = CAMU_SINK_AUDIO
});
- state = BUFFER_ADDED;
}
- entry->audio.state = state;
}
#ifndef CAMU_SINK_NO_VIDEO
void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
{
- u8 state = entry->video.state;
- al_assert(state != BUFFER_ADDED);
- if (state == BUFFER_CONFIGURED) {
- state = BUFFER_SET_OR_BUFFERED;
- } else if (state == BUFFER_SET_OR_BUFFERED) {
+ al_assert(entry->video.state != BUFFER_ADDED);
+ 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);
if (AUDIO_READY_OR_EMPTY(entry)) {
maybe_remove_previous(entry->sink);
@@ -351,9 +375,7 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
.op = single_frame ? STOP : START,
.value.i = CAMU_SINK_VIDEO
});
- state = BUFFER_ADDED;
}
- entry->video.state = state;
}
#endif
@@ -416,7 +438,7 @@ static void video_buffer_callback(void *userdata, u8 op)
case CAMU_BUFFER_BUFFERED: {
bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
aki_mutex_lock(&sink->mutex);
- if (single_frame || !entry->ended) {
+ if (!entry->ended || single_frame) {
add_video_if_set_and_buffered(entry);
}
aki_mutex_unlock(&sink->mutex);
@@ -477,16 +499,21 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
frames -= sink->video.renderer->get_latency(sink->video.renderer);
camu_video_buffer_set_latency(&entry->video.buf, -frames);
}
+#else
+ (void)sink;
+ (void)entry;
#endif
#else
// To sync clients with differing audio latencies our only option is to factor the mixer
// 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)) {
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);
}
+#endif
camu_audio_buffer_set_latency(&entry->audio.buf, audio);
#endif
}
@@ -650,10 +677,10 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
)
{
sink->loop = loop;
- aki_signal_init(&sink->signal, queue_signal_callback, sink);
- aki_signal_start(&sink->signal, sink->loop);
- camu_queue_init(sink->queue);
aki_mutex_init(&sink->mutex);
+ aki_signal_init(&sink->queue_signal, queue_signal_callback, sink);
+ aki_signal_start(&sink->queue_signal, sink->loop);
+ camu_queue_init(sink->queue);
sink->queued = NULL;
sink->current = NULL;
al_array_init(sink->previous);
@@ -720,7 +747,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *entry)
}
} else {
if (sink->current) {
- al_array_push(sink->previous, sink->current);
+ maybe_add_to_previous(sink, sink->current, entry);
}
}
set_or_queue_entry(entry);
@@ -832,8 +859,6 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
entry->ended = ended;
- if (entry == sink->current) goto out;
-
if (op == LIANA_SINK_BUFFER) {
goto out;
} else if (op == LIANA_SINK_BUFFER_AND_QUEUE) {
@@ -844,10 +869,14 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
#ifdef CAMU_SINK_LOCAL
(void)at;
(void)pause;
- if (sink->current && !camu_clock_is_paused(&sink->current->clock)) {
- camu_clock_pause(&sink->current->clock, 0);
+ if (sink->current) {
+ if (!camu_clock_is_paused(&sink->current->clock)) {
+ camu_clock_pause(&sink->current->clock, 0);
+ }
+ maybe_add_to_previous(sink, sink->current, entry);
}
- switch_to(sink, entry);
+ set_or_queue_entry(entry);
+ sink->current = entry;
// This will resume a user paused stream.
camu_clock_resume(&entry->clock, 0);
#else
@@ -858,7 +887,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
sink->target = NULL;
}
if (sink->current) {
- al_array_push(sink->previous, sink->current);
+ maybe_add_to_previous(sink, sink->current, entry);
}
set_or_queue_entry(entry);
sink->current = entry;
@@ -871,7 +900,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
if (entry == sink->current) {
camu_audio_buffer_unpause(&entry->audio.buf);
} else if (sink->current) {
- al_array_push(sink->previous, sink->current);
+ maybe_add_to_previous(sink, sink->current, entry);
}
camu_clock_resume(&entry->clock, at);
set_or_queue_entry(entry);
@@ -1028,6 +1057,14 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
aki_rpc_connection_command(sink->conn, packet, idd_callback, sink);
}
+static void reconnect_timer_callback(void *userdata, struct aki_timer *timer)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)timer;
+ aki_rpc_reconnect(&sink->client, &sink->addr, sink->port);
+ aki_timer_stop(&sink->reconnect_timer);
+}
+
static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
@@ -1035,12 +1072,15 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
al_assert(sink->conn == conn);
sink->conn = NULL;
}
+ aki_timer_again(&sink->reconnect_timer);
}
bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name)
{
al_str_clone(&sink->name, name);
sink->type = type;
+ aki_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink);
+ aki_timer_set_repeat(&sink->reconnect_timer, AKI_TS_FROM_USEC(1000000));
aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink);
for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) {
commands[i].userdata = sink;
@@ -1050,7 +1090,9 @@ bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str
if (!aki_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) {
return false;
}
- aki_rpc_connect(&sink->client, addr, port);
+ al_str_clone(&sink->addr, addr);
+ sink->port = port;
+ aki_rpc_connect(&sink->client, &sink->addr, sink->port);
return true;
}
@@ -1142,6 +1184,8 @@ void camu_sink_stop(struct camu_sink *sink)
void camu_sink_close(struct camu_sink *sink)
{
+ aki_timer_stop(&sink->reconnect_timer);
+ aki_timer_disable(&sink->reconnect_timer);
if (sink->conn) aki_rpc_conn_disconnect(sink->conn);
struct camu_sink_entry *entry;
al_array_foreach_rev(sink->entries, i, entry) {
@@ -1160,4 +1204,6 @@ void camu_sink_free(struct camu_sink *sink)
aki_rpc_free(&sink->client);
camu_queue_free(sink->queue);
aki_mutex_destroy(&sink->mutex);
+ al_str_free(&sink->addr);
+ al_str_free(&sink->name);
}
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index b830230..0b55ad6 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -4,6 +4,7 @@
#include <al/array.h>
#include <aki/rpc2.h>
#include <aki/signal.h>
+#include <aki/timer.h>
#include "../util/queue.h"
@@ -68,11 +69,14 @@ struct camu_sink {
struct aki_event_loop *loop;
str name;
u8 type;
+ str addr;
+ u16 port;
struct aki_rpc client;
struct aki_rpc_connection *conn;
- struct aki_signal signal;
- queue(struct camu_sink_cmd) queue;
struct aki_mutex mutex;
+ struct aki_timer reconnect_timer;
+ struct aki_signal queue_signal;
+ queue(struct camu_sink_cmd) queue;
str default_list;
struct camu_sink_entry *current;
struct camu_sink_entry *queued;