summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/common.h3
-rw-r--r--src/libsink/sink.c247
-rw-r--r--src/libsink/sink.h2
3 files changed, 169 insertions, 83 deletions
diff --git a/src/libsink/common.h b/src/libsink/common.h
index 1c701c8..48e43e8 100644
--- a/src/libsink/common.h
+++ b/src/libsink/common.h
@@ -2,5 +2,6 @@
enum {
CAMU_SINK_CMD_BUFFER = 0,
- CAMU_SINK_CMD_SET
+ CAMU_SINK_CMD_SET,
+ CAMU_SINK_CMD_QUEUE,
};
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 3eecf5f..134be31 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,9 +1,12 @@
#include <al/log.h>
#include "../tree/common.h"
+
#include "../bimu/common.h"
#include "../bimu/handler.h"
+#include "../buffer/common.h"
+
#include "sink.h"
#include "common.h"
@@ -31,15 +34,18 @@ enum {
START,
STOP,
TOGGLE_PAUSE,
+ SKIP,
RESEEK,
SEEK,
SET_BUFFERED, // Currently set entry is buffered.
CLOSE,
// Internal.
CORK,
- UNCORK
+ UNCORK,
};
+#define BIMU_DELAY_IGNORE 0
+
#define ENTRY_AUDIO_BUFFER_HELD(entry) \
(al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED))
#ifdef CAMU_SINK_NO_VIDEO
@@ -55,6 +61,36 @@ enum {
#define ENTRY_AUDIO_EMPTY(entry) \
(entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED)
+static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
+{
+#ifndef CAMU_SINK_NO_VIDEO
+ 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;
+ }
+#endif
+ 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;
+ }
+}
+
+static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry);
+#ifndef CAMU_SINK_NO_VIDEO
+static void add_video_if_set_and_buffered(struct camu_sink_entry *entry);
+#endif
+
+static void swap_buffers_internal(struct camu_sink *sink)
+{
+ remove_entry_buffers(sink, sink->current);
+ sink->current = sink->queued;
+ sink->queued = NULL;
+ add_audio_if_set_and_buffered(sink->current);
+#ifndef CAMU_SINK_NO_VIDEO
+ add_video_if_set_and_buffered(sink->current);
+#endif
+}
+
static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
{
switch (cmd->op) {
@@ -122,6 +158,12 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
break;
}
+ case SKIP: {
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_SKIP);
+ aki_packet_write_s32(packet, cmd->value.i);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
case RESEEK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
bmu_client_reseek(&entry->client);
@@ -146,6 +188,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
// break;
// }
case CLOSE: {
+ aki_signal_stop(&sink->signal);
sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
return;
}
@@ -198,11 +241,16 @@ static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
aki_signal_send(&sink->signal);
}
-static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
+void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
{
u8 state = entry->audio.state;
if (state == BUFFER_SET_OR_BUFFERED) {
+ // TODO: should_resume is not robust.
+#ifndef CAMU_SINK_NO_VIDEO
bool should_resume = ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED;
+#else
+ bool should_resume = true;
+#endif
if (should_resume && !camu_clock_calc_tick(&entry->clock)) {
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = RESEEK,
@@ -230,6 +278,40 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
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;
+ if (state == BUFFER_SET_OR_BUFFERED) {
+ bool should_resume = ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED;
+ if (should_resume && !camu_clock_calc_tick(&entry->clock)) {
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = RESEEK,
+ .opaque = entry
+ });
+ return;
+ }
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = START,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ if (should_resume) {
+ camu_clock_resume(&entry->clock);
+ entry->state = ENTRY_BUFFERED;
+ }
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ //queue_cmd(entry->sink, (struct camu_sink_cmd){
+ // .op = SET_BUFFERED,
+ // .value.i = CAMU_SINK_VIDEO
+ //});
+ state = BUFFER_ADDED;
+ } else if (state == BUFFER_CONFIGURED) {
+ state = BUFFER_SET_OR_BUFFERED;
+ }
+ entry->video.state = state;
+}
+#endif
+
static void audio_buffer_callback(void *userdata, u8 op)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
@@ -260,15 +342,16 @@ static void audio_buffer_callback(void *userdata, u8 op)
});
break;
case CAMU_BUFFER_EOF: {
- //aki_mutex_lock(&entry->sink->mutex);
- //u8 ret = entry->sink->callback(entry->sink->userdata, CAMU_SINK_SWAP_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- //entry->audio.state = BUFFER_SET_OR_BUFFERED;
- //aki_mutex_unlock(&entry->sink->mutex);
- //if (ret != CAMU_SINK_BUFFERS_SWAPPED) {
- if (1) {
- // TODO: Can cause popping. Should run a timer
- // and if any action would queue_cmd(AUDIO_START) before the
- // the timer, don't queue the stop.
+ struct camu_sink *sink = entry->sink;
+ bool swapped = false;
+ aki_mutex_lock(&sink->mutex);
+ if (sink->queued) {
+ swap_buffers_internal(sink);
+ swapped = true;
+ al_log_info("sink", "Audio buffers swapped (gapless).");
+ }
+ aki_mutex_unlock(&sink->mutex);
+ if (!swapped) {
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_AUDIO
@@ -286,38 +369,6 @@ static void audio_buffer_callback(void *userdata, u8 op)
}
#ifndef CAMU_SINK_NO_VIDEO
-static void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
-{
- u8 state = entry->video.state;
- if (state == BUFFER_SET_OR_BUFFERED) {
- bool should_resume = ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED;
- if (should_resume && !camu_clock_calc_tick(&entry->clock)) {
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = RESEEK,
- .opaque = entry
- });
- return;
- }
- queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = START,
- .value.i = CAMU_SINK_VIDEO
- });
- if (should_resume) {
- camu_clock_resume(&entry->clock);
- entry->state = ENTRY_BUFFERED;
- }
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- //queue_cmd(entry->sink, (struct camu_sink_cmd){
- // .op = SET_BUFFERED,
- // .value.i = CAMU_SINK_VIDEO
- //});
- state = BUFFER_ADDED;
- } else if (state == BUFFER_CONFIGURED) {
- state = BUFFER_SET_OR_BUFFERED;
- }
- entry->video.state = state;
-}
-
static void video_buffer_callback(void *userdata, u8 op)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
@@ -388,7 +439,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_vcr_stream *stream
case BIMU_CLIENT_SET: {
u64 ts = aki_get_timestamp();
struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
- req->delay = 0; // Ignore delay.
+ req->delay = BIMU_DELAY_IGNORE;
req->ts -= req->delay;
bool late = ts <= req->ts;
if (late) {}
@@ -502,6 +553,7 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
aki_signal_start(&sink->signal, sink->loop);
camu_queue_init(sink->queue);
aki_mutex_init(&sink->mutex);
+ sink->queued = NULL;
sink->current = NULL;
al_array_init(sink->entries);
sink->audio.mixer = mixer;
@@ -526,27 +578,19 @@ static void buffer_timer_callback(void *userdata, struct aki_timer *timer)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
(void)timer;
- aki_mutex_lock(&entry->sink->mutex);
- aki_mutex_unlock(&entry->sink->mutex);
+ //aki_mutex_lock(&entry->sink->mutex);
+ //aki_mutex_unlock(&entry->sink->mutex);
aki_timer_stop(&entry->timer);
}
-static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
- struct aki_packet *packet, struct aki_packet *rpacket)
+static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *sink, str *addr, s32 port, u32 node_id)
{
- struct camu_sink *sink = (struct camu_sink *)userdata;
- (void)conn;
- (void)rpacket;
-
- str addr;
- aki_packet_read_str(packet, &addr);
- s32 port = aki_packet_read_s32(packet);
- u16 node_id = aki_packet_read_u16(packet);
+ struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
+ if (entry) return entry;
- struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry);
- entry->sink = sink;
+ entry = al_alloc_object(struct camu_sink_entry);
al_array_push(sink->entries, entry);
-
+ entry->sink = sink;
entry->state = ENTRY_LOADED;
entry->audio.state = BUFFER_INIT;
@@ -567,26 +611,32 @@ static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *c
entry->client.callback = client_callback;
entry->client.userdata = entry;
- bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id);
+ bmu_client_connect(&entry->client, sink->loop, addr, port, node_id);
+
+ return entry;
+}
+
+static bool buffer_command_callback(void *userdata, struct aki_rpc_connection *conn,
+ struct aki_packet *packet, struct aki_packet *rpacket)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)conn;
+ (void)rpacket;
+
+ str addr;
+ aki_packet_read_str(packet, &addr);
+ s32 port = aki_packet_read_s32(packet);
+ u16 node_id = aki_packet_read_u16(packet);
+
+ aki_mutex_lock(&sink->mutex);
+ ensure_entry_buffered_internal(sink, &addr, port, node_id);
+ aki_mutex_unlock(&sink->mutex);
aki_packet_free(packet);
return false;
}
-static void remove_entry_buffers(struct camu_sink *sink, struct camu_sink_entry *entry)
-{
-#ifndef CAMU_SINK_NO_VIDEO
- 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;
- }
-#endif
- 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;
- }
-}
static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
@@ -595,12 +645,17 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
(void)conn;
(void)rpacket;
+ str addr;
+ aki_packet_read_str(packet, &addr);
+ s32 port = aki_packet_read_s32(packet);
u16 node_id = aki_packet_read_u16(packet);
+
aki_mutex_lock(&sink->mutex);
- struct camu_sink_entry *entry = entry_from_node_id(sink, node_id);
- if (!entry) goto out;
+ struct camu_sink_entry *entry = ensure_entry_buffered_internal(sink, &addr, port, node_id);
+ if (entry == sink->current) goto out;
struct camu_sink_entry *previous = sink->current;
sink->current = entry;
+ if (sink->queued) sink->queued = NULL;
if (entry->audio.state == BUFFER_INIT) {
entry->audio.state = BUFFER_QUEUED;
} else {
@@ -617,16 +672,35 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn
camu_clock_pause(&previous->clock);
remove_entry_buffers(sink, previous);
}
-
out:
aki_mutex_unlock(&sink->mutex);
aki_packet_free(packet);
return false;
}
+static bool queue_command_callback(void *userdata, struct aki_rpc_connection *conn,
+ struct aki_packet *packet, struct aki_packet *rpacket)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)conn;
+ (void)rpacket;
+
+ str addr;
+ aki_packet_read_str(packet, &addr);
+ s32 port = aki_packet_read_s32(packet);
+ u16 node_id = aki_packet_read_u16(packet);
+
+ aki_mutex_lock(&sink->mutex);
+ sink->queued = ensure_entry_buffered_internal(sink, &addr, port, node_id);
+ aki_mutex_unlock(&sink->mutex);
+
+ return false;
+}
+
static struct aki_rpc_command commands[] = {
{ .op = CAMU_SINK_CMD_BUFFER, .callback = buffer_command_callback, .userdata = NULL },
- { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL }
+ { .op = CAMU_SINK_CMD_SET, .callback = set_command_callback, .userdata = NULL },
+ { .op = CAMU_SINK_CMD_QUEUE, .callback = queue_command_callback, .userdata = NULL }
};
static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
@@ -640,7 +714,8 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn)
{
- (void)userdata;
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)sink;
(void)conn;
}
@@ -685,8 +760,10 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent)
void camu_sink_skip(struct camu_sink *sink, s32 n)
{
- (void)sink;
- (void)n;
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = SKIP,
+ .value.i = n
+ });
}
struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
@@ -719,15 +796,21 @@ void camu_sink_stop(struct camu_sink *sink)
void camu_sink_close(struct camu_sink *sink)
{
+ // TODO: Make sure no commands can come in after this.
+ aki_rpc_disconnect(&sink->client);
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
bmu_client_close(&entry->client);
+ al_array_remove_at_iter(sink->entries, i);
}
aki_mutex_unlock(&sink->mutex);
+}
+
+void camu_sink_free(struct camu_sink *sink)
+{
al_array_free(sink->entries);
- aki_rpc_disconnect(&sink->client);
- aki_signal_stop(&sink->signal);
+ aki_rpc_free(&sink->client);
camu_queue_free(sink->queue);
aki_mutex_destroy(&sink->mutex);
}
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index e3550bb..40d20d4 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -73,6 +73,7 @@ struct camu_sink {
struct aki_signal signal;
queue(struct camu_sink_cmd) queue;
struct aki_mutex mutex;
+ struct camu_sink_entry *queued;
struct camu_sink_entry *current;
array(struct camu_sink_entry *) entries;
struct {
@@ -103,3 +104,4 @@ struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink);
void camu_sink_return_current(struct camu_sink *sink);
void camu_sink_stop(struct camu_sink *sink);
void camu_sink_close(struct camu_sink *sink);
+void camu_sink_free(struct camu_sink *sink);