diff options
Diffstat (limited to 'src/libsink')
| -rw-r--r-- | src/libsink/common.h | 8 | ||||
| -rw-r--r-- | src/libsink/meson.build | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 791 | ||||
| -rw-r--r-- | src/libsink/sink.h | 32 |
4 files changed, 559 insertions, 274 deletions
diff --git a/src/libsink/common.h b/src/libsink/common.h index a5ad25c..deae228 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -1,6 +1,12 @@ #pragma once +#define CAMU_SINK_LOCAL 0 + enum { CAMU_SINK_SET = 0, - CAMU_SINK_QUEUE + CAMU_SINK_BUFFER, + CAMU_SINK_BUFFER_AND_QUEUE, + CAMU_SINK_CLEAR, + CAMU_SINK_PAUSE, + CAMU_SINK_SEEK }; diff --git a/src/libsink/meson.build b/src/libsink/meson.build index 83cfdd7..adc06d2 100644 --- a/src/libsink/meson.build +++ b/src/libsink/meson.build @@ -1,3 +1,3 @@ libsink_src = ['sink.c'] -libsink_deps = [shrub_client] +libsink_deps = [liana_client] libsink = declare_dependency(sources: libsink_src, dependencies: libsink_deps) diff --git a/src/libsink/sink.c b/src/libsink/sink.c index c6727ab..859b717 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -2,13 +2,11 @@ #include "../server/common.h" -#include "../shrub/common.h" -#include "../shrub/handler.h" +#include "../liana/list.h" +#include "../liana/handler.h" #include "../buffer/common.h" -#include "../list/list.h" - #include "sink.h" #include "common.h" @@ -23,7 +21,8 @@ enum { BUFFER_QUEUED, BUFFER_CONFIGURED, BUFFER_SET_OR_BUFFERED, - BUFFER_ADDED + BUFFER_ADDED, + BUFFER_REMOVED }; enum { @@ -32,21 +31,13 @@ enum { TOGGLE_PAUSE, SEEK, SKIP, - FINISHED, + SHUFFLE, + END, RESEEK, - SET_BUFFERED, // Currently set entry is buffered. CLOSE }; -#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_MAX_AGE 7 #ifdef CAMU_SINK_NO_VIDEO #define ENTRY_VIDEO_READY_OR_EMPTY(entry) true @@ -57,18 +48,67 @@ enum { #define ENTRY_AUDIO_READY_OR_EMPTY(entry) \ (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED || entry->audio.state == BUFFER_ADDED) +#if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED +#define BLOCKING_SLEEP(delay) aki_thread_sleep(delay) +#else +#define BLOCKING_SLEEP(delay) aki_event_loop_sleep(sink->loop, delay) +#endif + +static bool entry_audio_buffer_held(struct camu_sink_entry *entry) +{ +#ifdef CAMU_MIXER_THREADED + return al_atomic_load(u8)(&(entry)->audio.buf.ref, AL_ATOMIC_RELAXED) == 1; +#else + (void)entry; + return false; +#endif +} + +#ifndef CAMU_SINK_NO_VIDEO +static bool entry_video_buffer_held(struct camu_sink_entry *entry) +{ +#ifdef CAMU_SCREEN_THREADED + return al_atomic_load(u8)(&(entry)->video.buf.ref, AL_ATOMIC_RELAXED) == 1; +#else + (void)entry; + return false; +#endif +} +#endif + +static bool entry_buffers_held(struct camu_sink_entry *entry) +{ +#ifndef CAMU_SINK_NO_VIDEO + return entry_audio_buffer_held(entry) || entry_video_buffer_held(entry); +#else + return entry_audio_buffer_held(entry); +#endif +} + +static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + 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; + } +} + +#ifndef CAMU_SINK_NO_VIDEO +static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + 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 + 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; - } + remove_entry_video_buffer(sink, entry); #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; - } + remove_entry_audio_buffer(sink, entry); } static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry); @@ -76,21 +116,57 @@ static void add_audio_if_set_and_buffered(struct camu_sink_entry *entry); static void add_video_if_set_and_buffered(struct camu_sink_entry *entry); #endif -static void swap_buffers_internal(struct camu_sink *sink) +static void set_or_queue_entry(struct camu_sink_entry *entry) { - remove_entry_buffers(sink, sink->current); - sink->current = sink->queued; - sink->queued = NULL; - add_audio_if_set_and_buffered(sink->current); + if (entry->audio.state == BUFFER_INIT) { + entry->audio.state = BUFFER_QUEUED; + } else { + add_audio_if_set_and_buffered(entry); + } #ifndef CAMU_SINK_NO_VIDEO - add_video_if_set_and_buffered(sink->current); + if (entry->video.state == BUFFER_INIT) { + entry->video.state = BUFFER_QUEUED; + } else { + add_video_if_set_and_buffered(entry); + } #endif } +#if CAMU_SINK_LOCAL +static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + if (camu_clock_is_paused(&entry->clock)) { + camu_clock_resume(&entry->clock, 0); + bool no_audio = entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED; + if (!no_audio && 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 + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); + if (!single_frame && 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, 0); +#ifndef CAMU_SINK_NO_VIDEO + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); + if (!single_frame && sink->video.state == SINK_PLAYING) { + sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL); + sink->video.state = SINK_PAUSED; + } +#endif + } +} +#endif + static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) { switch (cmd->op) { case START: { + aki_mutex_lock(&sink->mutex); switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PAUSED) { @@ -107,9 +183,11 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; #endif } + aki_mutex_unlock(&sink->mutex); break; } case STOP: { + aki_mutex_lock(&sink->mutex); switch (cmd->value.i) { case CAMU_SINK_AUDIO: if (sink->audio.state == SINK_PLAYING) { @@ -126,90 +204,69 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; #endif } + aki_mutex_unlock(&sink->mutex); break; } case SKIP: { + if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_u8(packet, CAMU_SKIP); - s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID; + aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_u8(packet, CAMU_LIST_SKIP); + //s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; + s32 sequence = LIANA_SEQUENCE_ANY; aki_packet_write_s32(packet, sequence); aki_packet_write_s32(packet, cmd->value.i); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } + case SHUFFLE: { + if (!sink->conn) return; + struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); + aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_u8(packet, CAMU_LIST_SHUFFLE); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } case TOGGLE_PAUSE: { - struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; -#ifdef CAMU_SINK_LOCAL_PAUSE - 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, -1.0); -#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 - } +#if CAMU_SINK_LOCAL + sink_local_pause(sink, (struct camu_sink_entry *)cmd->opaque); #else - f64 base = camu_clock_get_base_pts(&entry->clock); - if (!camu_clock_is_paused(&entry->clock)) { - f64 pts = camu_clock_get_pts(&entry->clock, 0.0); - if (pts > base) base = pts; - } + if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_u8(packet, CAMU_TOGGLE_PAUSE); - s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID; + aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE); + s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, base * 1000000.0); + aki_packet_write_f64(packet, cmd->value.f); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); #endif break; } case SEEK: { + if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_u8(packet, CAMU_SEEK); - s32 sequence = sink->current ? sink->current->sequence : CAMU_SEQUENCE_INVALID; + aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_u8(packet, CAMU_LIST_SEEK); + s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; aki_packet_write_s32(packet, sequence); aki_packet_write_f64(packet, cmd->value.f); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } - case FINISHED: { + case END: { + if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_u8(packet, CAMU_FINISHED); + aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_u8(packet, CAMU_LIST_END); aki_packet_write_s32(packet, cmd->value.i); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } case RESEEK: { struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; - shrb_client_reseek(&entry->client); + lia_client_reseek(&entry->client); 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 CLOSE: { aki_signal_stop(&sink->signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); @@ -250,20 +307,15 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) u8 state = entry->audio.state; if (state == BUFFER_SET_OR_BUFFERED) { bool can_resume = ENTRY_VIDEO_READY_OR_EMPTY(entry); - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = START, - .value.i = CAMU_SINK_AUDIO - }); camu_audio_buffer_unpause(&entry->audio.buf); entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); if (can_resume) { - camu_clock_ready(&entry->clock); maybe_remove_previous(entry->sink); } - //queue_cmd(entry->sink, (struct camu_sink_cmd){ - // .op = SET_BUFFERED, - // .value.i = CAMU_SINK_AUDIO - //}); + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = START, + .value.i = CAMU_SINK_AUDIO + }); state = BUFFER_ADDED; } else if (state == BUFFER_CONFIGURED) { state = BUFFER_SET_OR_BUFFERED; @@ -277,19 +329,15 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry) u8 state = entry->video.state; if (state == BUFFER_SET_OR_BUFFERED) { bool can_resume = ENTRY_AUDIO_READY_OR_EMPTY(entry); - queue_cmd(entry->sink, (struct camu_sink_cmd){ - .op = START, - .value.i = CAMU_SINK_VIDEO - }); entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); if (can_resume) { - camu_clock_ready(&entry->clock); maybe_remove_previous(entry->sink); } - //queue_cmd(entry->sink, (struct camu_sink_cmd){ - // .op = SET_BUFFERED, - // .value.i = CAMU_SINK_VIDEO - //}); + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); + queue_cmd(entry->sink, (struct camu_sink_cmd){ + .op = single_frame ? STOP : START, + .value.i = CAMU_SINK_VIDEO + }); state = BUFFER_ADDED; } else if (state == BUFFER_CONFIGURED) { state = BUFFER_SET_OR_BUFFERED; @@ -309,10 +357,10 @@ static void audio_buffer_callback(void *userdata, u8 op) aki_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_CORK: - shrb_vcr_stream_cork(entry->audio.stream); + lia_vcr_cork(entry->audio.track); break; case CAMU_BUFFER_UNCORK: - shrb_vcr_stream_uncork(entry->audio.stream); + lia_vcr_uncork(entry->audio.track); break; case CAMU_BUFFER_PAUSED: aki_mutex_lock(&sink->mutex); @@ -325,28 +373,18 @@ static void audio_buffer_callback(void *userdata, u8 op) aki_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_EOF: { - bool swapped = false; + lia_vcr_cork(entry->audio.track); aki_mutex_lock(&sink->mutex); + camu_clock_end(&entry->clock); + remove_entry_audio_buffer(sink, entry); if (sink->queued) { - al_log_info("sink", "Swapping audio buffers (gapless)."); - swap_buffers_internal(sink); - swapped = true; + set_or_queue_entry(sink->queued); + sink->current = sink->queued; + sink->queued = NULL; } aki_mutex_unlock(&sink->mutex); - if (!swapped) { - 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 = SET_BUFFERED, - // .value.i = CAMU_SINK_VIDEO - //}); -#endif - } queue_cmd(sink, (struct camu_sink_cmd){ - .op = FINISHED, + .op = END, .value.i = entry->sequence }); break; @@ -366,22 +404,39 @@ static void video_buffer_callback(void *userdata, u8 op) aki_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_CORK: - shrb_vcr_stream_cork(entry->video.stream); + lia_vcr_cork(entry->video.track); break; case CAMU_BUFFER_UNCORK: - shrb_vcr_stream_uncork(entry->video.stream); + lia_vcr_uncork(entry->video.track); break; case CAMU_BUFFER_EOF: { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = STOP, - .value.i = CAMU_SINK_VIDEO - }); + lia_vcr_cork(entry->video.track); + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); + bool swapped = false; aki_mutex_lock(&sink->mutex); - bool no_audio = entry->audio.state == BUFFER_INIT; + if (!single_frame) { + remove_entry_video_buffer(sink, entry); + } + bool run_queue = !single_frame && (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED); + if (run_queue) { + camu_clock_end(&entry->clock); + if (sink->queued) { + set_or_queue_entry(sink->queued); + sink->current = sink->queued; + sink->queued = NULL; + swapped = true; + } + } aki_mutex_unlock(&sink->mutex); - if (no_audio) { + if (!swapped) { queue_cmd(sink, (struct camu_sink_cmd){ - .op = FINISHED, + .op = STOP, + .value.i = CAMU_SINK_VIDEO + }); + } + if (run_queue) { + queue_cmd(sink, (struct camu_sink_cmd){ + .op = END, .value.i = entry->sequence }); } @@ -391,124 +446,114 @@ static void video_buffer_callback(void *userdata, u8 op) } #endif -static void client_callback(void *userdata, u8 op, struct shrb_vcr_stream *stream, void *opaque) +static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + // Entry has both audio and video configured. +#ifndef CAMU_SINK_NO_VIDEO + if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED && + entry->video.state != BUFFER_INIT && entry->video.state != BUFFER_QUEUED) { + camu_video_buffer_set_latency(&entry->video.buf, -camu_mixer_get_latency(sink->audio.mixer)); + } +#else + (void)sink; + (void)entry; +#endif +} + +static void client_callback(void *userdata, u8 op, struct camu_codec_stream *stream, void *opaque) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; struct camu_sink *sink = entry->sink; switch (op) { - case SHRUB_CLIENT_CONFIGURE: { + case LIANA_CLIENT_CONFIGURE: { // No data can be sent until all selected streams are configured. switch (stream->type) { - case SHRUB_STREAM_AUDIO: - entry->audio.stream = stream; - camu_audio_buffer_configure(&entry->audio.buf, &stream->stream); - aki_mutex_lock(&entry->sink->mutex); + case CAMU_STREAM_AUDIO: + entry->audio.track = (struct lia_vcr_track *)opaque; + if (!camu_audio_buffer_configure(&entry->audio.buf, stream)) { + lia_client_disconnect(&entry->client); + } + aki_mutex_lock(&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); + evaluate_latency(sink, entry); + aki_mutex_unlock(&sink->mutex); break; #ifndef CAMU_SINK_NO_VIDEO - case SHRUB_STREAM_VIDEO: - entry->video.stream = stream; - camu_video_buffer_configure(&entry->video.buf, &stream->stream); - aki_mutex_lock(&entry->sink->mutex); + case CAMU_STREAM_VIDEO: + entry->video.track = (struct lia_vcr_track *)opaque; + if (!camu_video_buffer_configure(&entry->video.buf, stream)) { + lia_client_disconnect(&entry->client); + } + aki_mutex_lock(&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); + evaluate_latency(sink, entry); + aki_mutex_unlock(&sink->mutex); break; #endif } break; } - case SHRUB_CLIENT_SET: { - struct shrb_seek_req *req = (struct shrb_seek_req *)opaque; - req->delay = SHRUB_DELAY_IGNORE; - camu_clock_set(&entry->clock, req); - if (req->paused) req->base = req->paused_at; - break; - } - case SHRUB_CLIENT_PAUSE: { - f64 pts = *(f64 *)opaque; - aki_mutex_lock(&entry->sink->mutex); - camu_clock_arm_pause(&entry->clock, pts); - entry->audio.armed = false; - aki_mutex_unlock(&entry->sink->mutex); - break; - } - case SHRUB_CLIENT_RESUME: { - aki_mutex_lock(&entry->sink->mutex); - struct shrb_resume_req *req = (struct shrb_resume_req *)opaque; - if (sink->queued) { - camu_clock_offset(&sink->queued->clock, req->offset); - } - camu_clock_arm_resume(&entry->clock, req->start); - if (sink->audio.state == SINK_PAUSED) { - sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); - sink->audio.state = SINK_PLAYING; - } else { - entry->audio.armed = true; - } - aki_mutex_unlock(&entry->sink->mutex); - break; - } - case SHRUB_CLIENT_DATA: { - struct camu_frame *frame = (struct camu_frame *)opaque; + case LIANA_CLIENT_DATA: { + struct camu_codec_frame *frame = (struct camu_codec_frame *)opaque; switch (stream->type) { - case SHRUB_STREAM_AUDIO: + case CAMU_STREAM_AUDIO: camu_audio_buffer_push(&entry->audio.buf, frame); break; #ifndef CAMU_SINK_NO_VIDEO - case SHRUB_STREAM_VIDEO: + case CAMU_STREAM_VIDEO: camu_video_buffer_push(&entry->video.buf, frame); break; #endif default: - camu_frame_discard(frame); + camu_codec_frame_discard(frame); break; } break; } - case SHRUB_CLIENT_REMOVE_BUFFERS: { - aki_mutex_lock(&entry->sink->mutex); + case LIANA_CLIENT_REMOVE_BUFFERS: { + aki_mutex_lock(&sink->mutex); if (entry->audio.state == BUFFER_ADDED) { - entry->sink->callback(entry->sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + sink->callback(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); + bool reconnect = *(bool *)opaque; + bool keep_video = reconnect && camu_video_buffer_is_single_frame(&entry->video.buf); + if (!keep_video && 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 - 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)); } + aki_mutex_unlock(&sink->mutex); #ifndef CAMU_SINK_NO_VIDEO - camu_video_buffer_reset(&entry->video.buf); + while (keep_video ? entry_audio_buffer_held(entry) : entry_buffers_held(entry)) { + BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000)); + } + if (!keep_video) camu_video_buffer_reset(&entry->video.buf); +#else + while (entry_buffers_held(entry)) { + BLOCKING_SLEEP(AKI_TS_FROM_USEC(2000)); } #endif camu_audio_buffer_reset(&entry->audio.buf); break; } - case SHRUB_CLIENT_EOF: { + case LIANA_CLIENT_EOF: { switch (stream->type) { - case SHRUB_STREAM_AUDIO: { + case CAMU_STREAM_AUDIO: { camu_audio_buffer_flush(&entry->audio.buf); break; } #ifndef CAMU_SINK_NO_VIDEO - case SHRUB_STREAM_VIDEO: { + case CAMU_STREAM_VIDEO: { bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); if (!single_frame) { camu_video_buffer_flush(&entry->video.buf); @@ -519,19 +564,47 @@ static void client_callback(void *userdata, u8 op, struct shrb_vcr_stream *strea } break; } - case SHRUB_CLIENT_CLOSED: { - shrb_client_free(&entry->client); + case LIANA_CLIENT_CLOSED: { + aki_mutex_lock(&sink->mutex); + if (entry == sink->current) sink->current = NULL; + if (entry == sink->queued) sink->queued = NULL; + lia_client_free(&entry->client); camu_audio_buffer_free(&entry->audio.buf); #ifndef CAMU_SINK_NO_VIDEO camu_video_buffer_free(&entry->video.buf); #endif + bool removed = false; + struct camu_sink_entry *rentry; + al_array_foreach(sink->entries, i, rentry) { + if (rentry == entry) { + al_array_remove_at(sink->entries, i); + removed = true; + al_log_debug("sink", "Entry closed by disconnect."); + break; + } + } al_free(entry); - al_log_debug("sink", "Entry closed."); + if (!removed) { + al_log_debug("sink", "Entry closed by cleanup."); + } + aki_mutex_unlock(&sink->mutex); break; } } } +static void mixer_callback(void *userdata, u8 op) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + if (op == CAMU_MIXER_EMPTY) { + al_log_debug("sink", "Mixer empty."); + queue_cmd(sink, (struct camu_sink_cmd){ + .op = STOP, + .value.i = CAMU_SINK_AUDIO + }); + } +} + bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, struct camu_mixer *mixer #ifndef CAMU_SINK_NO_VIDEO @@ -548,6 +621,9 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, sink->current = NULL; al_array_init(sink->previous); al_array_init(sink->entries); + sink->lru = 0; + mixer->callback = mixer_callback; + mixer->userdata = sink; sink->audio.mixer = mixer; sink->audio.state = SINK_PAUSED; #ifndef CAMU_SINK_NO_VIDEO @@ -557,25 +633,82 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, return true; } -static struct camu_sink_entry *entry_from_node_id(struct camu_sink *sink, u16 node_id) +static s32 lru_compare(const void *a, const void *b) { + struct camu_sink_entry *aa = *((struct camu_sink_entry **)a); + struct camu_sink_entry *bb = *((struct camu_sink_entry **)b); + if (aa->lru > bb->lru) return -1; + else if (aa->lru < bb->lru) return 1; + return 0; +} + +static void maybe_cleanup_old_entries(struct camu_sink *sink) +{ + al_array_sort(sink->entries, struct camu_sink_entry *, lru_compare); + // We check size <= MAX_AGE in the loop because sink->lru + // is not indicative of the amount of entries we have loaded. + // There are various reasons for this but the most obvious is + // that it's incremented for buffer and queue operations. + // + // Handle sink->lru wrapping. + // 0 65532 65533 65534 65535 + // 0 1 65533 65534 65535 + // 0 1 2 65534 65535 + // 0 1 2 3 65535 + // 0 1 2 3 4 struct camu_sink_entry *entry; - al_array_foreach(sink->entries, i, entry) { - if (entry->client.node_id == node_id) return entry; + al_array_foreach_rev(sink->entries, i, entry) { + if (sink->entries.size <= ENTRY_MAX_AGE) return; + if (entry->lru > sink->lru && (UINT16_MAX - (entry->lru - 1)) + sink->lru >= ENTRY_MAX_AGE) { + al_array_remove_at(sink->entries, i); + lia_client_disconnect(&entry->client); + } + } + if (sink->lru >= ENTRY_MAX_AGE) { + al_array_foreach_rev(sink->entries, i, entry) { + if (sink->entries.size <= ENTRY_MAX_AGE) return; + if (sink->lru - entry->lru >= ENTRY_MAX_AGE) { + al_array_remove_at(sink->entries, i); + lia_client_disconnect(&entry->client); + } + } + } +} + +static void clock_callback(void *userdata, u8 op) +{ + struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; + struct camu_sink *sink = entry->sink; + if (op == CAMU_CLOCK_PAUSED) { + aki_mutex_lock(&sink->mutex); + if (sink->target) { + if (sink->current) { + al_array_push(sink->previous, sink->current); + } + set_or_queue_entry(sink->target); + sink->current = sink->target; + sink->target = NULL; + } + aki_mutex_unlock(&sink->mutex); } - return NULL; } -static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink *sink, - str *addr, s32 port, u32 node_id) +static struct camu_sink_entry *get_entry_from_id(struct camu_sink *sink, u16 id, bool *created) { - struct camu_sink_entry *entry = entry_from_node_id(sink, node_id); - if (entry) return entry; + struct camu_sink_entry *entry; + al_array_foreach(sink->entries, i, entry) { + if (entry->id == id) { + *created = false; + return entry; + } + } entry = al_alloc_object(struct camu_sink_entry); - al_array_push(sink->entries, entry); + entry->id = id; entry->sink = sink; + camu_clock_init(&entry->clock, clock_callback, entry); + entry->audio.state = BUFFER_INIT; entry->audio.armed = false; camu_audio_buffer_init(&entry->audio.buf, &entry->clock, sink->audio.mixer); @@ -584,15 +717,17 @@ static struct camu_sink_entry *ensure_entry_buffered_internal(struct camu_sink * #ifndef CAMU_SINK_NO_VIDEO entry->video.state = BUFFER_INIT; - camu_video_buffer_init(&entry->video.buf, &entry->clock, 0.0, sink->video.renderer); + camu_video_buffer_init(&entry->video.buf, &entry->clock, sink->video.renderer); entry->video.buf.callback = video_buffer_callback; entry->video.buf.userdata = entry; #endif entry->client.callback = client_callback; entry->client.userdata = entry; - // Assume we can error out at any point after connect is called (even from within connect() itself). - shrb_client_connect(&entry->client, sink->loop, addr, port, node_id); + + al_array_push(sink->entries, entry); + + *created = true; return entry; } @@ -604,100 +739,237 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn (void)conn; (void)rpacket; + aki_mutex_lock(&sink->mutex); + + u8 op = aki_packet_read_u8(packet); +#if CAMU_SINK_LOCAL + if (op == CAMU_SINK_CLEAR) { + if (sink->current) { + remove_entry_buffers(sink, sink->current); + sink->current = NULL; + } + goto out; + } +#else + if (op == CAMU_SINK_CLEAR) { al_assert(false); } // Unimplemented. +#endif + str addr; aki_packet_read_str(packet, &addr); - s32 port = aki_packet_read_s32(packet); + u16 port = aki_packet_read_u16(packet); u16 node_id = aki_packet_read_u16(packet); s32 sequence = aki_packet_read_s32(packet); - u64 start = aki_packet_read_u64(packet); - u8 paused = aki_packet_read_u8(packet); - u64 paused_at = aki_packet_read_u64(packet); + u64 at = aki_packet_read_u64(packet); + u64 seek_pos = aki_packet_read_u64(packet); + u8 pause = aki_packet_read_u8(packet); - aki_mutex_lock(&sink->mutex); - struct camu_sink_entry *entry = ensure_entry_buffered_internal(sink, &addr, port, node_id); + bool created; + struct camu_sink_entry *entry = get_entry_from_id(sink, node_id, &created); entry->sequence = sequence; - if (entry == sink->current) goto out; - if (!paused) { - // This will only have an effect if the clock is already set. - camu_clock_arm_resume(&entry->clock, start + SHRUB_BASE_DELAY); + sink->lru = al_u16_inc_wrap(sink->lru); + entry->lru = sink->lru; + + if (created) { + camu_clock_set(&entry->clock, seek_pos / 1000000.0); + struct camu_renderer *renderer = NULL; +#ifndef CAMU_SINK_NO_VIDEO + renderer = sink->video.renderer; +#endif + lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, seek_pos, renderer); + } + + if (op == CAMU_SINK_BUFFER) { + goto out; + } else if (op == CAMU_SINK_BUFFER_AND_QUEUE) { + sink->queued = entry; + goto out; } + +#if CAMU_SINK_LOCAL + (void)at; + (void)pause; if (sink->current) { - camu_clock_pause(&sink->current->clock, paused_at / 1000000.0); + if (!camu_clock_is_paused(&sink->current->clock)) { + camu_clock_pause(&sink->current->clock, 0); + } al_array_push(sink->previous, sink->current); } + set_or_queue_entry(entry); sink->current = entry; - if (sink->queued) sink->queued = NULL; - if (entry->audio.state == BUFFER_INIT) { - entry->audio.state = BUFFER_QUEUED; - } else { - add_audio_if_set_and_buffered(entry); + // This will resume a user paused stream. + camu_clock_resume(&entry->clock, 0); +#else + switch (pause) { + case LIANA_PAUSE_NONE: + if (sink->current) { + al_array_push(sink->previous, sink->current); + } + set_or_queue_entry(entry); + sink->current = entry; + break; + case LIANA_PAUSE_RESUME: + camu_clock_resume(&entry->clock, at); + if (sink->current) { + al_array_push(sink->previous, sink->current); + } + set_or_queue_entry(entry); + sink->current = entry; + break; + case LIANA_PAUSE_PAUSE: + if (sink->current) { + if (camu_clock_is_ended(&sink->current->clock)) { + // Server thought we weren't done, be we are. + al_array_push(sink->previous, sink->current); + set_or_queue_entry(entry); + sink->current = entry; + } else { + al_printf("ay\n"); + sink->target = entry; + camu_clock_pause(&sink->current->clock, at); + } + } + break; + case LIANA_PAUSE_BOTH: { + camu_clock_resume(&entry->clock, at); + struct camu_sink_entry *prev_target = sink->target; + if (prev_target) { + if (prev_target == entry) { + sink->target = NULL; + } else { + sink->target = entry; + } + camu_clock_pause(&prev_target->clock, at); + } else if (sink->current) { + if (camu_clock_is_ended(&sink->current->clock)) { + al_array_push(sink->previous, sink->current); + set_or_queue_entry(entry); + sink->current = entry; + } else { + sink->target = entry; + camu_clock_pause(&sink->current->clock, at); + } + } + break; } -#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 + out: aki_mutex_unlock(&sink->mutex); + if (op != CAMU_SINK_BUFFER) { + maybe_cleanup_old_entries(sink); + } + aki_packet_free(packet); return false; } -static bool queue_command_callback(void *userdata, struct aki_rpc_connection *conn, +static bool pause_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { struct camu_sink *sink = (struct camu_sink *)userdata; (void)conn; (void)rpacket; - str addr; - aki_packet_read_str(packet, &addr); - s32 port = aki_packet_read_s32(packet); - u16 node_id = aki_packet_read_u16(packet); + aki_mutex_lock(&sink->mutex); + struct camu_sink_entry *current = sink->current; + if (!current) goto out; +#if CAMU_SINK_LOCAL + sink_local_pause(sink, current); +#else s32 sequence = aki_packet_read_s32(packet); + u64 at = aki_packet_read_u64(packet); + u8 pause = aki_packet_read_u8(packet); + if (current->sequence == sequence) { + if (pause == LIANA_PAUSE_PAUSE) { + camu_clock_pause(&sink->current->clock, at); + current->audio.armed = false; + } else { + camu_clock_resume(&sink->current->clock, at); + if (sink->audio.state == SINK_PAUSED) { + sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); + sink->audio.state = SINK_PLAYING; + } else { + current->audio.armed = true; + } + } + } +#endif +out: + aki_mutex_unlock(&sink->mutex); + + aki_packet_free(packet); + return false; +} + +static bool seek_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; + + s32 sequence = aki_packet_read_s32(packet); + u64 at = aki_packet_read_u64(packet); + u64 pos = aki_packet_read_u64(packet); aki_mutex_lock(&sink->mutex); - sink->queued = ensure_entry_buffered_internal(sink, &addr, port, node_id); - sink->queued->sequence = sequence; + struct camu_sink_entry *current = sink->current; aki_mutex_unlock(&sink->mutex); + if (current->sequence == sequence) { + // TODO: Thread-safety here. + lia_client_seek(¤t->client, pos); + camu_clock_set(¤t->clock, pos / 1000000.0); + camu_clock_resume(¤t->clock, at); + } aki_packet_free(packet); - return false; } static struct aki_rpc_command commands[] = { { .op = CAMU_SINK_SET, .callback = set_command_callback, .userdata = NULL }, - { .op = CAMU_SINK_QUEUE, .callback = queue_command_callback, .userdata = NULL } + { .op = CAMU_SINK_PAUSE, .callback = pause_command_callback, .userdata = NULL }, + { .op = CAMU_SINK_SEEK, .callback = seek_command_callback, .userdata = NULL } }; +static void idd_callback(void *userdata, struct aki_packet *packet) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)sink; + aki_packet_free(packet); +} + 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, CAMU_SERVER_IDENTIFY); aki_packet_write_u8(packet, CAMU_SINK); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + aki_packet_write_str(packet, &sink->name); + aki_rpc_connection_command(sink->conn, packet, idd_callback, sink); } static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; - (void)sink; - (void)conn; + if (sink->conn) { + al_assert(sink->conn == conn); + sink->conn = NULL; + } } -bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port) +bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name) { + al_str_clone(&sink->name, name); + sink->type = type; aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; al_assert(sink->callback); aki_rpc_add_command(&sink->client, &commands[i]); } - if (!aki_rpc_prepare_client(&sink->client, AKI_SOCKET_TCP, CAMU_MULTIPLEX_RPC)) { + if (!aki_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) { return false; } aki_rpc_connect(&sink->client, addr, port); @@ -715,14 +987,6 @@ void camu_sink_return_current(struct camu_sink *sink) aki_mutex_unlock(&sink->mutex); } -void camu_sink_add(struct camu_sink *sink, str *line) -{ - struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_u8(packet, CAMU_ADD); - aki_packet_write_str(packet, line); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); -} - void camu_sink_skip(struct camu_sink *sink, s32 n) { queue_cmd(sink, (struct camu_sink_cmd){ @@ -731,14 +995,26 @@ void camu_sink_skip(struct camu_sink *sink, s32 n) }); } +void camu_sink_shuffle(struct camu_sink *sink) +{ + queue_cmd(sink, (struct camu_sink_cmd){ + .op = SHUFFLE + }); +} + void camu_sink_toggle_pause(struct camu_sink *sink) { aki_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; aki_mutex_unlock(&sink->mutex); + f64 pts = -1.0; + if (!camu_clock_is_paused(¤t->clock)) { + pts = camu_clock_get_pts(¤t->clock, 0.0); + } if (current) { queue_cmd(sink, (struct camu_sink_cmd){ .op = TOGGLE_PAUSE, + .value.f = pts, .opaque = current }); } @@ -788,19 +1064,20 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { - // TODO: Make sure no commands can come in after this. - aki_rpc_conn_disconnect(sink->client.conn); - aki_mutex_lock(&sink->mutex); + if (sink->conn) aki_rpc_conn_disconnect(sink->conn); struct camu_sink_entry *entry; - al_array_foreach(sink->entries, i, entry) { - shrb_client_close(&entry->client); + al_array_foreach_rev(sink->entries, i, entry) { + al_array_remove_at(sink->entries, i); + lia_client_disconnect(&entry->client); } - sink->entries.size = 0; - aki_mutex_unlock(&sink->mutex); } void camu_sink_free(struct camu_sink *sink) { + struct camu_sink_entry *entry; + al_array_foreach(sink->entries, i, entry) { + al_free(entry); + } al_array_free(sink->entries); aki_rpc_free(&sink->client); camu_queue_free(sink->queue); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index b789c9a..b9249cd 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -1,7 +1,5 @@ #pragma once -//#define CAMU_SINK_LOCAL_PAUSE - #include <al/str.h> #include <al/array.h> #include <aki/rpc2.h> @@ -15,7 +13,7 @@ #include "../buffer/video.h" #endif -#include "../shrub/client.h" +#include "../liana/client.h" enum { CAMU_SINK_AUDIO = 0, @@ -28,27 +26,32 @@ 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 }; +enum { + CAMU_SINK_OK = 0 +}; + struct camu_sink_entry { + u16 id; s32 sequence; + u16 lru; struct camu_clock clock; - struct shrb_client client; + struct lia_client client; struct { u8 state; bool armed; struct camu_audio_buffer buf; - struct shrb_vcr_stream *stream; + struct lia_vcr_track *track; } audio; #ifndef CAMU_SINK_NO_VIDEO struct { u8 state; struct camu_video_buffer buf; - struct shrb_vcr_stream *stream; + struct lia_vcr_track *track; } video; #endif struct camu_sink *sink; @@ -60,22 +63,21 @@ struct camu_sink_cmd { void *opaque; }; -enum { - CAMU_SINK_OK = 0, - CAMU_SINK_BUFFERS_SWAPPED = 1 -}; - struct camu_sink { struct aki_event_loop *loop; + str name; + u8 type; 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 *queued; struct camu_sink_entry *current; + struct camu_sink_entry *queued; + struct camu_sink_entry *target; array(struct camu_sink_entry *) previous; array(struct camu_sink_entry *) entries; + u16 lru; struct { u8 state; struct camu_mixer *mixer; @@ -96,11 +98,11 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, , struct camu_renderer *renderer #endif ); -bool camu_sink_connect(struct camu_sink *sink, str *addr, s32 port); +bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name); struct camu_sink_entry *camu_sink_get_current(struct camu_sink *sink); void camu_sink_return_current(struct camu_sink *sink); -void camu_sink_add(struct camu_sink *sink, str *line); void camu_sink_skip(struct camu_sink *sink, s32 n); +void camu_sink_shuffle(struct camu_sink *sink); void camu_sink_toggle_pause(struct camu_sink *sink); void camu_sink_seek(struct camu_sink *sink, f64 pos); void camu_sink_reseek(struct camu_sink *sink); |