summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/common.h6
-rw-r--r--src/libsink/meson.build2
-rw-r--r--src/libsink/sink.c193
-rw-r--r--src/libsink/sink.h16
-rw-r--r--src/libsink/sink2.c688
-rw-r--r--src/libsink/sink2.h106
6 files changed, 945 insertions, 66 deletions
diff --git a/src/libsink/common.h b/src/libsink/common.h
new file mode 100644
index 0000000..1c701c8
--- /dev/null
+++ b/src/libsink/common.h
@@ -0,0 +1,6 @@
+#pragma once
+
+enum {
+ CAMU_SINK_CMD_BUFFER = 0,
+ CAMU_SINK_CMD_SET
+};
diff --git a/src/libsink/meson.build b/src/libsink/meson.build
index 5cbe36d..f212b0c 100644
--- a/src/libsink/meson.build
+++ b/src/libsink/meson.build
@@ -1,2 +1,2 @@
-libsink_src = ['sink.c']
+libsink_src = ['sink2.c']
libsink = declare_dependency(sources: libsink_src)
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index bae6c77..3692066 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -32,12 +32,12 @@ enum {
};
#define ENTRY_AUDIO_BUFFER_HELD(entry) \
- (al_atomic_bool_load(&entry->audio.buf.ref, AL_ATOMIC_RELAXED))
+ (al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED))
#ifdef CAMU_SINK_NO_VIDEO
#define ENTRY_BUFFERS_HELD(entry) ENTRY_AUDIO_BUFFER_HELD(entry)
#else
#define ENTRY_VIDEO_BUFFER_HELD(entry) \
- (al_atomic_bool_load(&entry->video.buf.ref, AL_ATOMIC_RELAXED))
+ (al_atomic_bool_load(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED))
#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry))
#endif
@@ -160,7 +160,13 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
}
case SEEK: {
struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- bmu_local_seek(&entry->runner, cmd->value.f);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_seek(&((struct camu_sink_local *)entry)->runner, cmd->value.f);
+ break;
+ default:
+ break;
+ }
break;
}
case CLOSE: {
@@ -176,9 +182,8 @@ static void queue_signal_callback(void *userdata)
u32 size;
struct camu_sink_cmd cmd;
do {
- camu_queue_size(sink->queue, size);
- if (!size) break;
- camu_queue_pop(sink->queue, cmd);
+ camu_queue_try_pop(sink->queue, size, cmd);
+ if (size == 0) break;
handle_sink_cmd(sink, &cmd);
} while (1);
}
@@ -227,10 +232,22 @@ static void audio_buffer_callback(void *userdata, u8 op)
aki_mutex_unlock(&entry->sink->mutex);
break;
case CAMU_BUFFER_STOP:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->audio.stream);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_stream_stop((struct bmu_local_stream *)entry->audio.stream);
+ break;
+ default:
+ break;
+ }
break;
case CAMU_BUFFER_CONTINUE:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->audio.stream);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_stream_continue((struct bmu_local_stream *)entry->audio.stream);
+ break;
+ default:
+ break;
+ }
break;
case CAMU_BUFFER_PAUSED:
queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -245,15 +262,20 @@ static void audio_buffer_callback(void *userdata, u8 op)
entry->audio.state = BUFFER_SET_OR_BUFFERED;
aki_mutex_unlock(&entry->sink->mutex);
if (ret != CAMU_SINK_BUFFERS_SWAPPED) {
+ // 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.
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_AUDIO,
.opaque = entry->sink
});
+#ifndef CAMU_SINK_NO_VIDEO
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = SET_BUFFERED,
.value.i = CAMU_SINK_VIDEO
});
+#endif
}
break;
}
@@ -295,10 +317,22 @@ static void video_buffer_callback(void *userdata, u8 op)
aki_mutex_unlock(&entry->sink->mutex);
break;
case CAMU_BUFFER_STOP:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->video.stream);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_stream_stop((struct bmu_local_stream *)entry->video.stream);
+ break;
+ default:
+ break;
+ }
break;
case CAMU_BUFFER_CONTINUE:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->video.stream);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_stream_continue((struct bmu_local_stream *)entry->video.stream);
+ break;
+ default:
+ break;
+ }
break;
case CAMU_BUFFER_EOF:
queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -409,14 +443,32 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
break;
}
case BIMU_CLIENT_CLOSED: {
- bmu_local_close(&entry->runner);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_close(&((struct camu_sink_local *)entry)->runner);
+ break;
+ default:
+ break;
+ }
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
#endif
- cch_entry_free(&entry->entry);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ cch_entry_free(&((struct camu_sink_local *)entry)->entry);
+ break;
+ default:
+ break;
+ }
al_str_free(&entry->unique_id);
- al_free(entry);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ al_free((struct camu_sink_local *)entry);
+ break;
+ default:
+ break;
+ }
al_log_debug("sink", "Entry closed.");
break;
}
@@ -533,92 +585,109 @@ void camu_sink_close(struct camu_sink *sink)
aki_mutex_lock(&sink->mutex);
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
- bmu_local_stop(&entry->runner);
+ switch (entry->type) {
+ case CAMU_SINK_LOCAL:
+ bmu_local_stop(&((struct camu_sink_local *)entry)->runner);
+ break;
+ default:
+ break;
+ }
}
- aki_mutex_unlock(&sink->mutex);
al_array_free(sink->entries);
+ aki_mutex_unlock(&sink->mutex);
aki_signal_stop(&sink->signal);
camu_queue_free(sink->queue);
aki_mutex_destroy(&sink->mutex);
}
-static struct camu_sink_entry *entry_from_unique_id(struct camu_sink *sink, str *unique_id)
+// Local bimu band-aid compat.
+
+static struct camu_sink_local *entry_from_unique_id(struct camu_sink *sink, str *unique_id)
{
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
if (al_str_eq(&entry->unique_id, unique_id)) {
- aki_mutex_unlock(&sink->mutex);
- return entry;
+ return (struct camu_sink_local *)entry;
}
}
return NULL;
}
+static void camu_sink_buffer_internal(struct camu_sink *sink, struct camu_sink_entry *entry, str *unique_id)
+{
+ entry->sink = sink;
+ al_str_clone(&entry->unique_id, unique_id);
+
+ camu_clock_init(&entry->clock);
+
+ entry->audio.state = BUFFER_INIT;
+ camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
+ entry->audio.buf.callback = audio_buffer_callback;
+ entry->audio.buf.userdata = entry;
+
+ struct camu_renderer *renderer = NULL;
+#ifndef CAMU_SINK_NO_VIDEO
+ renderer = sink->video.renderer;
+ entry->video.state = BUFFER_INIT;
+ camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
+ entry->video.buf.callback = video_buffer_callback;
+ entry->video.buf.userdata = entry;
+#endif
+
+ al_array_push(sink->entries, entry);
+}
+
bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *centry)
{
aki_mutex_lock(&sink->mutex);
- bool add_entry = false;
- struct camu_sink_entry *entry = entry_from_unique_id(sink, unique_id);
- if (!entry) {
- entry = al_alloc_object(struct camu_sink_entry);
- add_entry = true;
+ bool repeat = true;
+ struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
+ if (!local) {
+ local = al_alloc_object(struct camu_sink_local);
+ local->e.type = CAMU_SINK_LOCAL;
+ repeat = false;
}
- entry->state = ENTRY_LOADED;
- if (!add_entry) {
+ local->e.state = ENTRY_LOADED;
+ if (repeat) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = SEEK,
.value.f = 0,
- .opaque = entry
+ .opaque = &local->e
});
aki_mutex_unlock(&sink->mutex);
return true;
}
- entry->entry = centry;
- if (!bmu_local_init(&entry->runner, centry)) {
- al_free(entry);
+ local->entry = centry;
+ if (!bmu_local_init(&local->runner, centry)) {
+ al_free(local);
aki_mutex_unlock(&sink->mutex);
return false;
}
- entry->sink = sink;
- al_str_clone(&entry->unique_id, unique_id);
-
- camu_clock_init(&entry->clock);
-
- entry->audio.state = BUFFER_INIT;
- camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
- entry->audio.buf.callback = audio_buffer_callback;
- entry->audio.buf.userdata = entry;
+ camu_sink_buffer_internal(sink, &local->e, unique_id);
struct camu_renderer *renderer = NULL;
#ifndef CAMU_SINK_NO_VIDEO
renderer = sink->video.renderer;
- entry->video.state = BUFFER_INIT;
- camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
- entry->video.buf.callback = video_buffer_callback;
- entry->video.buf.userdata = entry;
#endif
-
- entry->runner.callback = client_callback;
- entry->runner.userdata = entry;
- bmu_local_prepare_clients(&entry->runner, renderer);
- bmu_local_run(&entry->runner);
-
- if (add_entry) al_array_push(sink->entries, entry);
+ local->runner.callback = client_callback;
+ local->runner.userdata = &local->e;
+ bmu_local_prepare_clients(&local->runner, renderer);
+ bmu_local_run(&local->runner);
aki_mutex_unlock(&sink->mutex);
return true;
}
-static void maybe_cleanup_entries(struct camu_sink *sink)
+static void maybe_cleanup_local_entries(struct camu_sink *sink)
{
struct camu_sink_entry *entry;
al_array_foreach(sink->entries, i, entry) {
if (entry->state == ENTRY_DISREGUARDED && !ENTRY_BUFFERS_HELD(entry)) {
- bmu_local_stop(&entry->runner);
+ bmu_local_stop(&((struct camu_sink_local *)entry)->runner);
al_array_remove_at_iter(sink->entries, i);
}
}
@@ -627,12 +696,12 @@ static void maybe_cleanup_entries(struct camu_sink *sink)
void camu_sink_local_set(struct camu_sink *sink, str *unique_id)
{
aki_mutex_lock(&sink->mutex);
- struct camu_sink_entry *entry = entry_from_unique_id(sink, unique_id);
+ struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
struct camu_sink_entry *previous = sink->current;
- sink->current = entry;
- add_audio_if_set_and_buffered(entry);
+ sink->current = &local->e;
+ add_audio_if_set_and_buffered(&local->e);
#ifndef CAMU_SINK_NO_VIDEO
- add_video_if_set_and_buffered(entry);
+ add_video_if_set_and_buffered(&local->e);
#endif
if (previous) {
if (!camu_clock_is_paused(&previous->clock)) {
@@ -640,24 +709,24 @@ void camu_sink_local_set(struct camu_sink *sink, str *unique_id)
}
remove_entry_buffers(sink, previous);
}
- maybe_cleanup_entries(sink);
+ maybe_cleanup_local_entries(sink);
aki_mutex_unlock(&sink->mutex);
}
void camu_sink_local_swap(struct camu_sink *sink, str *unique_id)
{
- struct camu_sink_entry *entry = entry_from_unique_id(sink, unique_id);
- sink->current = entry;
- add_audio_if_set_and_buffered(entry);
+ struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
+ sink->current = &local->e;
+ add_audio_if_set_and_buffered(&local->e);
#ifndef CAMU_SINK_NO_VIDEO
- add_video_if_set_and_buffered(entry);
+ add_video_if_set_and_buffered(&local->e);
#endif
}
void camu_sink_local_unload(struct camu_sink *sink, str *unique_id)
{
aki_mutex_lock(&sink->mutex);
- struct camu_sink_entry *entry = entry_from_unique_id(sink, unique_id);
- if (entry && entry->state == ENTRY_LOADED) entry->state = ENTRY_DISREGUARDED;
+ struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
+ if (local && local->e.state == ENTRY_LOADED) local->e.state = ENTRY_DISREGUARDED;
aki_mutex_unlock(&sink->mutex);
}
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index 92c18d9..efef97b 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -27,20 +27,24 @@ enum {
};
enum {
+ CAMU_SINK_REMOTE = 0,
+ CAMU_SINK_LOCAL
+};
+
+enum {
CAMU_SINK_ADD_BUFFER = 0,
CAMU_SINK_REMOVE_BUFFER,
CAMU_SINK_SWAP_BUFFER,
- CAMU_SINK_SET_BUFFERED,
CAMU_SINK_START,
CAMU_SINK_STOP,
+ CAMU_SINK_ENTRY_BUFFERED,
CAMU_SINK_EXIT
};
struct camu_sink_entry {
+ u8 type;
u8 state;
str unique_id;
- struct cch_entry *entry;
- struct bmu_local runner;
struct camu_clock clock;
struct {
u8 state;
@@ -57,6 +61,12 @@ struct camu_sink_entry {
struct camu_sink *sink;
};
+struct camu_sink_local {
+ struct camu_sink_entry e;
+ struct cch_entry *entry;
+ struct bmu_local runner;
+};
+
struct camu_sink_cmd {
u8 op;
union { s64 i; f64 f; } value;
diff --git a/src/libsink/sink2.c b/src/libsink/sink2.c
new file mode 100644
index 0000000..e9b51f2
--- /dev/null
+++ b/src/libsink/sink2.c
@@ -0,0 +1,688 @@
+#include <al/log.h>
+
+#include "../tree/common.h"
+
+#include "sink2.h"
+
+enum {
+ SINK_EMPTY = 0,
+ SINK_PAUSED,
+ SINK_PLAYING
+};
+
+enum {
+ ENTRY_LOADED = 0,
+ ENTRY_DISREGUARDED
+};
+
+enum {
+ BUFFER_INIT = 0,
+ BUFFER_QUEUED,
+ BUFFER_CONFIGURED,
+ BUFFER_SET_OR_BUFFERED,
+ BUFFER_ADDED,
+};
+
+enum {
+ // TODO: Add and remove buffer can be done from
+ // any thread. Does this always need to be the case?
+ // Is it beneficial for this to be the case?
+ ADD_BUFFER = 0,
+ REMOVE_BUFFER,
+ START,
+ STOP,
+ TOGGLE_PAUSE,
+ SEEK,
+ ENTRY_BUFFERED, // Currently set entry is buffered.
+ CLOSE,
+ // Internal.
+ CORK,
+ UNCORK
+};
+
+#define ENTRY_AUDIO_BUFFER_HELD(entry) \
+ (al_atomic_bool_load(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED))
+#ifdef CAMU_SINK_NO_VIDEO
+#define ENTRY_BUFFERS_HELD(entry) ENTRY_AUDIO_BUFFER_HELD(entry)
+#else
+#define ENTRY_VIDEO_BUFFER_HELD(entry) \
+ (al_atomic_bool_load(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED))
+#define ENTRY_BUFFERS_HELD(entry) (ENTRY_AUDIO_BUFFER_HELD(entry) || ENTRY_VIDEO_BUFFER_HELD(entry))
+#endif
+
+#define ENTRY_VIDEO_EMPTY(entry) \
+ (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_QUEUED)
+#define ENTRY_AUDIO_EMPTY(entry) \
+ (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED)
+
+static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
+{
+ switch (cmd->op) {
+ case ADD_BUFFER: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ break;
+#endif
+ }
+ break;
+ }
+ case REMOVE_BUFFER: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ break;
+#endif
+ }
+ break;
+ }
+ case START: {
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ if (sink->audio.state == SINK_PAUSED) {
+ sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
+ sink->audio.state = SINK_PLAYING;
+ }
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ if (sink->video.state == SINK_PAUSED) {
+ sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
+ sink->video.state = SINK_PLAYING;
+ }
+ break;
+#endif
+ }
+ break;
+ }
+ case STOP: {
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ if (sink->audio.state == SINK_PLAYING) {
+ sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_AUDIO, NULL);
+ sink->audio.state = SINK_PAUSED;
+ }
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ if (sink->video.state == SINK_PLAYING) {
+ sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
+ sink->video.state = SINK_PAUSED;
+ }
+ break;
+#endif
+ }
+ break;
+ }
+ case TOGGLE_PAUSE: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ if (camu_clock_is_paused(&entry->clock)) {
+ camu_clock_resume(&entry->clock);
+ if (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
+ if (sink->video.state == SINK_PAUSED) {
+ sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
+ sink->video.state = SINK_PLAYING;
+ }
+#endif
+ } else {
+ camu_clock_pause(&entry->clock);
+#ifndef CAMU_SINK_NO_VIDEO
+ if (sink->video.state == SINK_PLAYING) {
+ sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
+ sink->video.state = SINK_PAUSED;
+ }
+#endif
+ }
+ break;
+ }
+ case SEEK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ (void)entry;
+ break;
+ }
+ case ENTRY_BUFFERED: {
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_AUDIO, NULL);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
+ break;
+#endif
+ }
+ break;
+ }
+ case CLOSE: {
+ sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
+ return;
+ }
+ case CORK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ bmu_client_cork(&entry->client, true);
+ break;
+ }
+ case UNCORK: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ bmu_client_cork(&entry->client, false);
+ break;
+ }
+ }
+}
+
+static void queue_signal_callback(void *userdata)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ u32 size;
+ struct camu_sink_cmd cmd;
+ do {
+ camu_queue_try_pop(sink->queue, size, cmd);
+ if (size == 0) break;
+ handle_sink_cmd(sink, &cmd);
+ } while (1);
+}
+
+static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
+{
+ camu_queue_push(sink->queue, cmd);
+ aki_signal_send(&sink->signal);
+}
+
+static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
+{
+ u8 state = entry->audio.state;
+ if (state == BUFFER_SET_OR_BUFFERED) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = ENTRY_BUFFERED,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = START,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) {
+ camu_clock_resume(&entry->clock);
+ }
+ state = BUFFER_ADDED;
+ } else if (state == BUFFER_CONFIGURED) {
+ state = BUFFER_SET_OR_BUFFERED;
+ }
+ entry->audio.state = state;
+}
+
+static void audio_buffer_callback(void *userdata, u8 op)
+{
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ switch (op) {
+ case CAMU_BUFFER_BUFFERED:
+ aki_mutex_lock(&entry->sink->mutex);
+ add_audio_if_set_and_buffered(entry);
+ aki_mutex_unlock(&entry->sink->mutex);
+ break;
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
+ break;
+ case CAMU_BUFFER_UNCORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = UNCORK,
+ .opaque = entry
+ });
+ break;
+ case CAMU_BUFFER_PAUSED:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ 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) {
+ // 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.
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_AUDIO
+ });
+#ifndef CAMU_SINK_NO_VIDEO
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = ENTRY_BUFFERED,
+ .value.i = CAMU_SINK_VIDEO
+ });
+#endif
+ }
+ break;
+ }
+ }
+}
+
+#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) {
+ 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 = ENTRY_BUFFERED,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = START,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) {
+ camu_clock_resume(&entry->clock);
+ }
+ 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;
+ switch (op) {
+ case CAMU_BUFFER_BUFFERED:
+ aki_mutex_lock(&entry->sink->mutex);
+ add_video_if_set_and_buffered(entry);
+ aki_mutex_unlock(&entry->sink->mutex);
+ break;
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
+ break;
+ case CAMU_BUFFER_UNCORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = UNCORK,
+ .opaque = entry
+ });
+ break;
+ case CAMU_BUFFER_EOF:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_VIDEO
+ });
+ break;
+ }
+}
+#endif
+
+static void client_callback(void *userdata, u8 op, struct bmu_client_stream *stream, void *opaque)
+{
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
+ switch (op) {
+ // A lot of logic here assumes that no data will be sent
+ // until all active streams are configured. This is extremely important,
+ // if it does not hold true many confusing errors will arise.
+ case BIMU_CLIENT_CONFIGURE: {
+ switch (stream->type) {
+ case BIMU_STREAM_AUDIO:
+ entry->audio.stream = stream;
+ camu_audio_buffer_configure(&entry->audio.buf, &stream->stream);
+ aki_mutex_lock(&entry->sink->mutex);
+ if (entry->audio.state == BUFFER_QUEUED) {
+ entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ } else {
+ entry->audio.state = BUFFER_CONFIGURED;
+ }
+ aki_mutex_unlock(&entry->sink->mutex);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case BIMU_STREAM_VIDEO:
+ entry->video.stream = stream;
+ camu_video_buffer_configure(&entry->video.buf, &stream->stream);
+ aki_mutex_lock(&entry->sink->mutex);
+ if (entry->video.state == BUFFER_QUEUED) {
+ entry->video.state = BUFFER_SET_OR_BUFFERED;
+ } else {
+ entry->video.state = BUFFER_CONFIGURED;
+ }
+ aki_mutex_unlock(&entry->sink->mutex);
+ break;
+#endif
+ }
+ break;
+ }
+ case BIMU_CLIENT_DATA: {
+ struct camu_frame *frame = (struct camu_frame *)opaque;
+ switch (stream->type) {
+ case BIMU_STREAM_AUDIO:
+ camu_audio_buffer_push(&entry->audio.buf, frame);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case BIMU_STREAM_VIDEO:
+ camu_video_buffer_push(&entry->video.buf, frame);
+ break;
+#endif
+ default:
+#ifdef HAVE_FFMPEG
+ if (frame->type == CAMU_FFMPEG_COMPAT) {
+ av_frame_free(&frame->av.frame);
+ }
+#endif
+ al_free(frame);
+ break;
+ }
+ break;
+ }
+ case BIMU_CLIENT_SEEK: {
+ // Must ensure no more data from the previous
+ // stream position is sent during or after this call to seek.
+ f64 seek_pos = *(f64 *)opaque;
+ aki_mutex_lock(&entry->sink->mutex);
+ if (entry->audio.state == BUFFER_ADDED) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ entry->audio.state = BUFFER_SET_OR_BUFFERED;
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
+ if (entry->video.state == BUFFER_ADDED && !single_frame) {
+ entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ entry->video.state = BUFFER_SET_OR_BUFFERED;
+ }
+#endif
+ aki_mutex_unlock(&entry->sink->mutex);
+#ifndef CAMU_SINK_NO_VIDEO
+ if (single_frame) {
+ while (ENTRY_AUDIO_BUFFER_HELD(entry)) {
+ aki_thread_sleep(AKI_TS_FROM_USEC(200));
+ }
+ } else {
+#endif
+ while (ENTRY_BUFFERS_HELD(entry)) {
+ aki_thread_sleep(AKI_TS_FROM_USEC(200));
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ }
+#endif
+ camu_clock_seek(&entry->clock, seek_pos);
+ camu_audio_buffer_reset(&entry->audio.buf);
+#ifndef CAMU_SINK_NO_VIDEO
+ if (!single_frame) {
+ camu_video_buffer_reset(&entry->video.buf);
+ }
+#endif
+ break;
+ }
+ case BIMU_CLIENT_EOF: {
+ switch (stream->type) {
+ case BIMU_STREAM_AUDIO: {
+ camu_audio_buffer_flush(&entry->audio.buf);
+ break;
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ case BIMU_STREAM_VIDEO: {
+ bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf);
+ if (!single_frame) {
+ camu_video_buffer_flush(&entry->video.buf);
+ }
+ break;
+ }
+#endif
+ }
+ break;
+ }
+ case BIMU_CLIENT_CLOSED: {
+ // close/free client.
+ camu_audio_buffer_free(&entry->audio.buf);
+#ifndef CAMU_SINK_NO_VIDEO
+ camu_video_buffer_free(&entry->video.buf);
+#endif
+ al_free(entry);
+ al_log_debug("sink", "Entry closed.");
+ break;
+ }
+ }
+}
+
+bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
+ struct camu_mixer *mixer
+#ifndef CAMU_SINK_NO_VIDEO
+ , struct camu_renderer *renderer
+#endif
+ )
+{
+ 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);
+ sink->current = NULL;
+ al_array_init(sink->entries);
+ sink->audio.mixer = mixer;
+ sink->audio.state = SINK_PAUSED;
+#ifndef CAMU_SINK_NO_VIDEO
+ sink->video.renderer = renderer;
+ sink->video.state = SINK_PAUSED;
+#endif
+ return true;
+}
+
+static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id)
+{
+ struct camu_sink_entry *entry;
+ al_array_foreach(sink->entries, i, entry) {
+ if (entry->client.node_id == node_id) return entry;
+ }
+ return NULL;
+}
+
+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);
+
+ struct camu_sink_entry *entry = al_alloc_object(struct camu_sink_entry);
+ entry->sink = sink;
+ al_array_push(sink->entries, entry);
+
+ entry->state = ENTRY_LOADED;
+ camu_clock_init(&entry->clock);
+
+ entry->audio.state = BUFFER_INIT;
+ camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer);
+ entry->audio.buf.callback = audio_buffer_callback;
+ entry->audio.buf.userdata = entry;
+
+ struct camu_renderer *renderer = NULL;
+#ifndef CAMU_SINK_NO_VIDEO
+ renderer = sink->video.renderer;
+ entry->video.state = BUFFER_INIT;
+ camu_video_buffer_init(&entry->video.buf, &entry->clock, renderer);
+ entry->video.buf.callback = video_buffer_callback;
+ entry->video.buf.userdata = entry;
+#endif
+
+ entry->client.callback = client_callback;
+ entry->client.userdata = entry;
+ bmu_client_connect(&entry->client, sink->loop, &addr, port, node_id);
+
+ 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)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ (void)conn;
+ (void)rpacket;
+
+ 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);
+ struct camu_sink_entry *previous = sink->current;
+ sink->current = entry;
+ if (entry->audio.state == BUFFER_INIT) {
+ entry->audio.state = BUFFER_QUEUED;
+ } else {
+ add_audio_if_set_and_buffered(entry);
+ }
+#ifndef CAMU_SINK_NO_VIDEO
+ if (entry->video.state == BUFFER_INIT) {
+ entry->video.state = BUFFER_QUEUED;
+ } else {
+ add_video_if_set_and_buffered(entry);
+ }
+#endif
+ if (previous) {
+ camu_clock_pause(&previous->clock);
+ remove_entry_buffers(sink, previous);
+ }
+ aki_mutex_unlock(&sink->mutex);
+
+ aki_packet_free(packet);
+ 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 }
+};
+
+static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
+{
+ struct camu_sink *sink = (struct camu_sink *)userdata;
+ sink->conn = conn;
+ struct aki_packet *packet = aki_rpc_get_packet(&sink->client, TREE_CMD_IDENTIFY);
+ aki_packet_write_u8(packet, TREE_SINK);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+}
+
+static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn)
+{
+ (void)userdata;
+ (void)conn;
+}
+
+bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
+{
+ aki_rpc_init(&sink->client, AKI_SOCKET_TCP, connection_callback,
+ connection_closed_callback, sink);
+ for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) {
+ commands[i].userdata = sink;
+ al_assert(sink->callback);
+ aki_rpc_add_command(&sink->client, &commands[i]);
+ }
+ return aki_rpc_connect(&sink->client, sink->loop, addr, port);
+}
+
+void camu_sink_toggle_pause(struct camu_sink *sink)
+{
+ aki_mutex_lock(&sink->mutex);
+ if (sink->current) {
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = TOGGLE_PAUSE,
+ .opaque = sink->current
+ });
+ }
+ aki_mutex_unlock(&sink->mutex);
+}
+
+void camu_sink_seek(struct camu_sink *sink, f64 pos)
+{
+ aki_mutex_lock(&sink->mutex);
+ if (sink->current) {
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = SEEK,
+ .value.f = pos,
+ .opaque = sink->current
+ });
+ }
+ aki_mutex_unlock(&sink->mutex);
+}
+
+void camu_sink_skip(struct camu_sink *sink, s32 n)
+{
+ (void)sink;
+ (void)n;
+}
+
+struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink)
+{
+ aki_mutex_lock(&sink->mutex);
+ return sink->current;
+}
+
+void camu_sink_return_current(struct camu_sink *sink)
+{
+ aki_mutex_unlock(&sink->mutex);
+}
+
+void camu_sink_stop(struct camu_sink *sink)
+{
+ aki_mutex_lock(&sink->mutex);
+ if (sink->current) {
+ remove_entry_buffers(sink, sink->current);
+ sink->current = NULL;
+ }
+ aki_mutex_unlock(&sink->mutex);
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = STOP,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = CLOSE
+ });
+}
+
+void camu_sink_close(struct camu_sink *sink)
+{
+ aki_mutex_lock(&sink->mutex);
+ /*
+ struct camu_sink_entry *entry;
+ al_array_foreach(sink->entries, i, entry) {
+ }
+ */
+ al_array_free(sink->entries);
+ aki_mutex_unlock(&sink->mutex);
+ aki_signal_stop(&sink->signal);
+ camu_queue_free(sink->queue);
+ aki_mutex_destroy(&sink->mutex);
+}
diff --git a/src/libsink/sink2.h b/src/libsink/sink2.h
new file mode 100644
index 0000000..95c280f
--- /dev/null
+++ b/src/libsink/sink2.h
@@ -0,0 +1,106 @@
+#pragma once
+
+//#define CAMU_SINK_NO_VIDEO
+
+#include <al/str.h>
+#include <al/array.h>
+#include <aki/rpc2.h>
+#include <aki/signal.h>
+#include <aki/timer.h>
+
+#include "../util/queue.h"
+
+#include "../buffer/clock.h"
+#include "../buffer/audio.h"
+#ifndef CAMU_SINK_NO_VIDEO
+#include "../buffer/video.h"
+#endif
+
+#include "../bimu/client.h"
+
+#include "common.h"
+
+enum {
+ CAMU_SINK_AUDIO = 0,
+#ifndef CAMU_SINK_NO_VIDEO
+ CAMU_SINK_VIDEO
+#endif
+};
+
+enum {
+ CAMU_SINK_ADD_BUFFER = 0,
+ CAMU_SINK_REMOVE_BUFFER,
+ CAMU_SINK_SWAP_BUFFER,
+ CAMU_SINK_SET_BUFFERED,
+ CAMU_SINK_START,
+ CAMU_SINK_STOP,
+ CAMU_SINK_EXIT
+};
+
+struct camu_sink_entry {
+ u8 state;
+ struct bmu_client client;
+ struct camu_clock clock;
+ struct {
+ u8 state;
+ struct camu_audio_buffer buf;
+ struct bmu_client_stream *stream;
+ } audio;
+#ifndef CAMU_SINK_NO_VIDEO
+ struct {
+ u8 state;
+ struct camu_video_buffer buf;
+ struct bmu_client_stream *stream;
+ } video;
+#endif
+ struct camu_sink *sink;
+};
+
+struct camu_sink_cmd {
+ u8 op;
+ union { s64 i; f64 f; } value;
+ void *opaque;
+};
+
+enum {
+ CAMU_SINK_OK = 0,
+ CAMU_SINK_BUFFERS_SWAPPED = 1
+};
+
+struct camu_sink {
+ struct aki_event_loop *loop;
+ struct aki_rpc client;
+ struct aki_rpc_connection *conn;
+ struct aki_signal signal;
+ queue(struct camu_sink_cmd) queue;
+ struct aki_mutex mutex;
+ struct camu_sink_entry *current;
+ array(struct camu_sink_entry *) entries;
+ struct {
+ u8 state;
+ struct camu_mixer *mixer;
+ } audio;
+#ifndef CAMU_SINK_NO_VIDEO
+ struct {
+ u8 state;
+ struct camu_renderer *renderer;
+ } video;
+#endif
+ u8 (*callback)(void *, u8, u8, void *);
+ void *userdata;
+};
+
+bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
+ struct camu_mixer *mixer
+#ifndef CAMU_SINK_NO_VIDEO
+ , struct camu_renderer *renderer
+#endif
+ );
+bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port);
+void camu_sink_toggle_pause(struct camu_sink *sink);
+void camu_sink_seek(struct camu_sink *sink, f64 pos);
+void camu_sink_skip(struct camu_sink *sink, s32 n);
+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);