summaryrefslogtreecommitdiff
path: root/src/libsink
diff options
context:
space:
mode:
Diffstat (limited to 'src/libsink')
-rw-r--r--src/libsink/meson.build5
-rw-r--r--src/libsink/sink.c530
-rw-r--r--src/libsink/sink.h30
-rw-r--r--src/libsink/sink2.c688
-rw-r--r--src/libsink/sink2.h106
5 files changed, 241 insertions, 1118 deletions
diff --git a/src/libsink/meson.build b/src/libsink/meson.build
index f212b0c..0e6dfb3 100644
--- a/src/libsink/meson.build
+++ b/src/libsink/meson.build
@@ -1,2 +1,3 @@
-libsink_src = ['sink2.c']
-libsink = declare_dependency(sources: libsink_src)
+libsink_src = ['sink.c']
+libsink_deps = [bimu_client]
+libsink = declare_dependency(sources: libsink_src, dependencies: libsink_deps)
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 3692066..5addc0f 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -1,5 +1,7 @@
#include <al/log.h>
+#include "../tree/common.h"
+
#include "sink.h"
enum {
@@ -15,20 +17,22 @@ enum {
enum {
BUFFER_INIT = 0,
+ BUFFER_QUEUED,
BUFFER_CONFIGURED,
BUFFER_SET_OR_BUFFERED,
BUFFER_ADDED,
};
enum {
- ADD_BUFFER = 0,
- REMOVE_BUFFER,
- SET_BUFFERED,
START,
STOP,
TOGGLE_PAUSE,
SEEK,
- CLOSE
+ ENTRY_BUFFERED, // Currently set entry is buffered.
+ CLOSE,
+ // Internal.
+ CORK,
+ UNCORK
};
#define ENTRY_AUDIO_BUFFER_HELD(entry) \
@@ -41,60 +45,15 @@ enum {
#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: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
- }
-#endif
- }
- break;
- }
- case REMOVE_BUFFER: {
- switch (cmd->value.i) {
- case CAMU_SINK_AUDIO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
- break;
- }
-#ifndef CAMU_SINK_NO_VIDEO
- case CAMU_SINK_VIDEO: {
- struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
- sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- break;
- }
-#endif
- }
- break;
- }
- case SET_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 START: {
- struct camu_sink *sink = (struct camu_sink *)cmd->opaque;
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PAUSED) {
@@ -114,7 +73,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
}
case STOP: {
- struct camu_sink *sink = (struct camu_sink *)cmd->opaque;
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PLAYING) {
@@ -160,12 +118,19 @@ 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;
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_seek(&((struct camu_sink_local *)entry)->runner, cmd->value.f);
+ (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;
- default:
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_SET_BUFFERED, CAMU_SINK_VIDEO, NULL);
break;
+#endif
}
break;
}
@@ -173,6 +138,16 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
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;
+ }
}
}
@@ -200,21 +175,16 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
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 = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_AUDIO
});
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_AUDIO
});
-#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state == BUFFER_INIT || entry->video.state == BUFFER_ADDED) {
-#endif
+ if (ENTRY_VIDEO_EMPTY(entry) || entry->video.state == BUFFER_ADDED) {
camu_clock_resume(&entry->clock);
-#ifndef CAMU_SINK_NO_VIDEO
}
-#endif
state = BUFFER_ADDED;
} else if (state == BUFFER_CONFIGURED) {
state = BUFFER_SET_OR_BUFFERED;
@@ -231,29 +201,22 @@ static void audio_buffer_callback(void *userdata, u8 op)
add_audio_if_set_and_buffered(entry);
aki_mutex_unlock(&entry->sink->mutex);
break;
- case CAMU_BUFFER_STOP:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->audio.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
break;
- case CAMU_BUFFER_CONTINUE:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->audio.stream);
- break;
- default:
- 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,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_AUDIO
});
break;
case CAMU_BUFFER_EOF: {
@@ -267,12 +230,11 @@ static void audio_buffer_callback(void *userdata, u8 op)
// 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
+ .value.i = CAMU_SINK_AUDIO
});
#ifndef CAMU_SINK_NO_VIDEO
queue_cmd(entry->sink, (struct camu_sink_cmd){
- .op = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_VIDEO
});
#endif
@@ -289,15 +251,14 @@ static void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
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 = SET_BUFFERED,
+ .op = ENTRY_BUFFERED,
.value.i = CAMU_SINK_VIDEO
});
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = START,
- .value.i = CAMU_SINK_VIDEO,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_VIDEO
});
- if (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_ADDED) {
+ if (ENTRY_AUDIO_EMPTY(entry) || entry->audio.state == BUFFER_ADDED) {
camu_clock_resume(&entry->clock);
}
state = BUFFER_ADDED;
@@ -316,29 +277,22 @@ static void video_buffer_callback(void *userdata, u8 op)
add_video_if_set_and_buffered(entry);
aki_mutex_unlock(&entry->sink->mutex);
break;
- case CAMU_BUFFER_STOP:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_stop((struct bmu_local_stream *)entry->video.stream);
- break;
- default:
- break;
- }
+ case CAMU_BUFFER_CORK:
+ queue_cmd(entry->sink, (struct camu_sink_cmd){
+ .op = CORK,
+ .opaque = entry
+ });
break;
- case CAMU_BUFFER_CONTINUE:
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stream_continue((struct bmu_local_stream *)entry->video.stream);
- break;
- default:
- 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,
- .opaque = entry->sink
+ .value.i = CAMU_SINK_VIDEO
});
break;
}
@@ -349,18 +303,39 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
switch (op) {
+ case BIMU_CLIENT_SET: {
+ struct bmu_seek_req *req = (struct bmu_seek_req *)opaque;
+ camu_clock_set(&entry->clock, req->base, req->start);
+ camu_audio_buffer_reset(&entry->audio.buf);
+ break;
+ }
+ // 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);
- entry->audio.state = BUFFER_CONFIGURED;
+ 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);
- entry->video.state = BUFFER_CONFIGURED;
+ 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
}
@@ -384,11 +359,14 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
}
#endif
al_free(frame);
+ break;
}
break;
}
case BIMU_CLIENT_SEEK: {
- f64 seek_req = *(f64 *)opaque;
+ // Must ensure no more data from the previous
+ // stream position is sent during or after this call to seek.
+ struct bmu_seek_req *req = (struct bmu_seek_req *)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);
@@ -415,7 +393,7 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
#ifndef CAMU_SINK_NO_VIDEO
}
#endif
- camu_clock_seek(&entry->clock, seek_req);
+ camu_clock_seek(&entry->clock, req->base, req->start);
camu_audio_buffer_reset(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
if (!single_frame) {
@@ -443,32 +421,12 @@ static void client_callback(void *userdata, u8 op, struct bmu_client_stream *str
break;
}
case BIMU_CLIENT_CLOSED: {
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_close(&((struct camu_sink_local *)entry)->runner);
- break;
- default:
- break;
- }
+ bmu_client_free(&entry->client);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
#endif
- 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);
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- al_free((struct camu_sink_local *)entry);
- break;
- default:
- break;
- }
+ al_free(entry);
al_log_debug("sink", "Entry closed.");
break;
}
@@ -498,37 +456,160 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop,
return true;
}
-bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port)
+static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id)
{
- (void)sink;
- (void)addr;
- (void)port;
+ 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;
+
+ 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) {
+ struct camu_sink_entry *current = sink->current;
+ aki_mutex_unlock(&sink->mutex);
+ if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = TOGGLE_PAUSE,
- .opaque = sink->current
+ .opaque = current
});
}
- aki_mutex_unlock(&sink->mutex);
}
void camu_sink_seek(struct camu_sink *sink, f64 pos)
{
aki_mutex_lock(&sink->mutex);
- if (sink->current) {
+ struct camu_sink_entry *current = sink->current;
+ aki_mutex_unlock(&sink->mutex);
+ if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = SEEK,
.value.f = pos,
- .opaque = sink->current
+ .opaque = current
});
}
- aki_mutex_unlock(&sink->mutex);
}
void camu_sink_skip(struct camu_sink *sink, s32 n)
@@ -548,20 +629,6 @@ void camu_sink_return_current(struct camu_sink *sink)
aki_mutex_unlock(&sink->mutex);
}
-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;
- }
-}
-
void camu_sink_stop(struct camu_sink *sink)
{
aki_mutex_lock(&sink->mutex);
@@ -572,8 +639,7 @@ void camu_sink_stop(struct camu_sink *sink)
aki_mutex_unlock(&sink->mutex);
queue_cmd(sink, (struct camu_sink_cmd){
.op = STOP,
- .value.i = CAMU_SINK_AUDIO,
- .opaque = sink
+ .value.i = CAMU_SINK_AUDIO
});
queue_cmd(sink, (struct camu_sink_cmd){
.op = CLOSE
@@ -583,150 +649,14 @@ void camu_sink_stop(struct camu_sink *sink)
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) {
- switch (entry->type) {
- case CAMU_SINK_LOCAL:
- bmu_local_stop(&((struct camu_sink_local *)entry)->runner);
- break;
- default:
- break;
- }
}
+ */
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);
}
-
-// 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)) {
- 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 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;
- }
- local->e.state = ENTRY_LOADED;
- if (repeat) {
- queue_cmd(sink, (struct camu_sink_cmd){
- .op = SEEK,
- .value.f = 0,
- .opaque = &local->e
- });
- aki_mutex_unlock(&sink->mutex);
- return true;
- }
-
- local->entry = centry;
- if (!bmu_local_init(&local->runner, centry)) {
- al_free(local);
- aki_mutex_unlock(&sink->mutex);
- return false;
- }
-
- camu_sink_buffer_internal(sink, &local->e, unique_id);
-
- struct camu_renderer *renderer = NULL;
-#ifndef CAMU_SINK_NO_VIDEO
- renderer = sink->video.renderer;
-#endif
- 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_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(&((struct camu_sink_local *)entry)->runner);
- al_array_remove_at_iter(sink->entries, i);
- }
- }
-}
-
-void camu_sink_local_set(struct camu_sink *sink, str *unique_id)
-{
- aki_mutex_lock(&sink->mutex);
- struct camu_sink_local *local = entry_from_unique_id(sink, unique_id);
- struct camu_sink_entry *previous = sink->current;
- sink->current = &local->e;
- add_audio_if_set_and_buffered(&local->e);
-#ifndef CAMU_SINK_NO_VIDEO
- add_video_if_set_and_buffered(&local->e);
-#endif
- if (previous) {
- if (!camu_clock_is_paused(&previous->clock)) {
- camu_clock_pause(&previous->clock);
- }
- remove_entry_buffers(sink, previous);
- }
- 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_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(&local->e);
-#endif
-}
-
-void camu_sink_local_unload(struct camu_sink *sink, str *unique_id)
-{
- aki_mutex_lock(&sink->mutex);
- 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 efef97b..95c280f 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -4,20 +4,21 @@
#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 "../cache/entry.h"
-
#include "../buffer/clock.h"
#include "../buffer/audio.h"
#ifndef CAMU_SINK_NO_VIDEO
#include "../buffer/video.h"
#endif
-#include "../bimu/local.h"
+#include "../bimu/client.h"
+
+#include "common.h"
enum {
CAMU_SINK_AUDIO = 0,
@@ -27,24 +28,18 @@ 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 bmu_client client;
struct camu_clock clock;
struct {
u8 state;
@@ -61,12 +56,6 @@ 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;
@@ -80,6 +69,8 @@ enum {
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;
@@ -113,8 +104,3 @@ 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);
-
-bool camu_sink_local_buffer(struct camu_sink *sink, str *unique_id, struct cch_entry *entry);
-void camu_sink_local_set(struct camu_sink *sink, str *unique_id);
-void camu_sink_local_swap(struct camu_sink *sink, str *unique_id);
-void camu_sink_local_unload(struct camu_sink *sink, str *unique_id);
diff --git a/src/libsink/sink2.c b/src/libsink/sink2.c
deleted file mode 100644
index e9b51f2..0000000
--- a/src/libsink/sink2.c
+++ /dev/null
@@ -1,688 +0,0 @@
-#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
deleted file mode 100644
index 95c280f..0000000
--- a/src/libsink/sink2.h
+++ /dev/null
@@ -1,106 +0,0 @@
-#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);