diff options
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/common.h | 6 | ||||
| -rw-r--r-- | src/libsink/meson.build | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 193 | ||||
| -rw-r--r-- | src/libsink/sink.h | 16 | ||||
| -rw-r--r-- | src/libsink/sink2.c | 688 | ||||
| -rw-r--r-- | src/libsink/sink2.h | 106 |
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); |