diff options
| -rw-r--r-- | .gitignore | 2 | ||||
| -rw-r--r-- | meson.build | 3 | ||||
| -rwxr-xr-x | scripts/run.sh | 4 | ||||
| -rwxr-xr-x | scripts/run_debug.sh | 4 | ||||
| -rwxr-xr-x | scripts/run_tree.sh | 4 | ||||
| -rw-r--r-- | src/bimu/client.c | 24 | ||||
| -rw-r--r-- | src/bimu/client.h | 1 | ||||
| -rw-r--r-- | src/bimu/handler.h | 7 | ||||
| -rw-r--r-- | src/bimu/meson.build | 12 | ||||
| -rw-r--r-- | src/bimu/server.c | 35 | ||||
| -rw-r--r-- | src/buffer/audio.c | 10 | ||||
| -rw-r--r-- | src/buffer/clock.c | 73 | ||||
| -rw-r--r-- | src/buffer/clock.h | 9 | ||||
| -rw-r--r-- | src/buffer/video.c | 1 | ||||
| -rw-r--r-- | src/fruits/cmv/meson.build | 10 | ||||
| -rw-r--r-- | src/fruits/cmv/tui.c | 165 | ||||
| -rw-r--r-- | src/fruits/cmv/tui.h | 28 | ||||
| -rw-r--r-- | src/fruits/sink/meson.build | 10 | ||||
| -rw-r--r-- | src/fruits/sink/sink.c (renamed from src/fruits/cmv/cmv2.c) | 164 | ||||
| -rw-r--r-- | src/libsink/meson.build | 5 | ||||
| -rw-r--r-- | src/libsink/sink.c | 530 | ||||
| -rw-r--r-- | src/libsink/sink.h | 30 | ||||
| -rw-r--r-- | src/libsink/sink2.c | 688 | ||||
| -rw-r--r-- | src/libsink/sink2.h | 106 | ||||
| -rw-r--r-- | src/tree/list.c | 78 | ||||
| -rw-r--r-- | src/tree/list.h | 1 | ||||
| -rw-r--r-- | src/tree/meson.build | 2 | ||||
| -rw-r--r-- | src/tree/tree.c | 31 | ||||
| -rw-r--r-- | src/tree/tree.h | 9 |
29 files changed, 467 insertions, 1579 deletions
@@ -12,6 +12,8 @@ subprojects/c89atomic/ subprojects/c89atomic.wrap subprojects/curl.wrap subprojects/ffmpeg-6.1/ +subprojects/glfw3/ +subprojects/glfw3.wrap subprojects/jansson.wrap subprojects/libakiyo subprojects/libalabaster diff --git a/meson.build b/meson.build index ef9da30..3bc60fb 100644 --- a/meson.build +++ b/meson.build @@ -26,9 +26,8 @@ subdir('src/cache') subdir('src/bimu') subdir('src/mixer') subdir('src/render') -subdir('src/fruits/cap') subdir('src/libclient') subdir('src/libsink') subdir('src/tree') -subdir('src/fruits/cmv') +subdir('src/fruits/sink') subdir('src/fruits/cmc') diff --git a/scripts/run.sh b/scripts/run.sh index abd024a..f2bbaa0 100755 --- a/scripts/run.sh +++ b/scripts/run.sh @@ -1,5 +1,3 @@ #! /usr/bin/env sh -export PYTHONDONTWRITEBYTECODE=1 -export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages export mesa_glthread=true -rm -f /tmp/cmv_sock && ./src/fruits/cmv/cmv "$@" +./src/fruits/sink/sink "$@" diff --git a/scripts/run_debug.sh b/scripts/run_debug.sh index b05a533..bd9a667 100755 --- a/scripts/run_debug.sh +++ b/scripts/run_debug.sh @@ -1,4 +1,2 @@ #! /usr/bin/env sh -export PYTHONDONTWRITEBYTECODE=1 -export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages -rm -f /tmp/cmv_sock && gdb -ex run --args ./src/fruits/cmv/cmv $@ +gdb -ex run --args ./src/fruits/cmv/cmv "$@" diff --git a/scripts/run_tree.sh b/scripts/run_tree.sh index 4359c29..274ec5c 100755 --- a/scripts/run_tree.sh +++ b/scripts/run_tree.sh @@ -1,4 +1,4 @@ #! /usr/bin/env sh export PYTHONDONTWRITEBYTECODE=1 -export PYTHONPATH=$HOME/c/shoki/vendor/vendor/lib/python3.11/site-packages -gdb -ex run --args ./src/tree/tree $@ +export PYTHONPATH=$HOME/c/shoki/vendor/vendor +rm -f /tmp/tree_sock && gdb -ex run --args ./src/tree/tree $@ diff --git a/src/bimu/client.c b/src/bimu/client.c index d35e3a1..b05d950 100644 --- a/src/bimu/client.c +++ b/src/bimu/client.c @@ -144,6 +144,19 @@ static void send_subscribe_packet(struct bmu_client *client) aki_packet_stream_send_packet(&client->control, packet); } +static void parse_start_packet(struct bmu_client *client, struct aki_packet *packet) +{ + u64 seek_pos = aki_packet_read_u64(packet); + u64 ats = aki_packet_read_u64(packet); + u8 paused = aki_packet_read_u8(packet); + (void)paused; + struct bmu_seek_req req = { + .base = seek_pos, + .start = ats + }; + client->callback(client->userdata, BIMU_CLIENT_SET, NULL, &req); +} + static void control_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) { @@ -157,8 +170,13 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st send_subscribe_packet(client); break; case BIMU_CONTROL_START: + parse_start_packet(client, packet); connect_data(client); break; + case BIMU_CONTROL_PAUSE: + break; + case BIMU_CONTROL_RESUME: + break; default: break; } @@ -186,3 +204,9 @@ void bmu_client_cork(struct bmu_client *client, bool cork) aki_packet_stream_cork(&client->data, cork); client->corked = cork; } + +void bmu_client_free(struct bmu_client *client) +{ + aki_packet_stream_free(&client->control); + aki_packet_stream_free(&client->data); +} diff --git a/src/bimu/client.h b/src/bimu/client.h index 54c903a..464ee03 100644 --- a/src/bimu/client.h +++ b/src/bimu/client.h @@ -31,3 +31,4 @@ struct bmu_client { void bmu_client_connect(struct bmu_client *client, struct aki_event_loop *loop, str *addr, s32 port, u16 node_id); void bmu_client_cork(struct bmu_client *client, bool cork); +void bmu_client_free(struct bmu_client *client); diff --git a/src/bimu/handler.h b/src/bimu/handler.h index 3881288..f9e8440 100644 --- a/src/bimu/handler.h +++ b/src/bimu/handler.h @@ -17,13 +17,18 @@ struct bmu_server_handler { }; enum { - BIMU_CLIENT_CONFIGURE = 0, + BIMU_CLIENT_SET = 0, + BIMU_CLIENT_CONFIGURE, BIMU_CLIENT_DATA, BIMU_CLIENT_SEEK, BIMU_CLIENT_EOF, BIMU_CLIENT_CLOSED }; +struct bmu_seek_req { + u64 base, start; +}; + struct bmu_client_stream; struct bmu_client_handler { bool (*init)(struct bmu_client_handler *, struct camu_renderer *, struct bmu_client_stream *); diff --git a/src/bimu/meson.build b/src/bimu/meson.build index 086114d..dce7f7b 100644 --- a/src/bimu/meson.build +++ b/src/bimu/meson.build @@ -1,8 +1,12 @@ -bimu_src = [ - 'local.c', +bimu_server_src = [ 'server.c', + 'handlers/codec_server.c' +] + +bimu_client_src = [ 'client.c', - 'handlers/codec_server.c', 'handlers/codec_client.c' ] -bimu = declare_dependency(sources: bimu_src) + +bimu_server = declare_dependency(sources: bimu_server_src) +bimu_client = declare_dependency(sources: bimu_client_src) diff --git a/src/bimu/server.c b/src/bimu/server.c index f572290..565f0b1 100644 --- a/src/bimu/server.c +++ b/src/bimu/server.c @@ -10,14 +10,16 @@ static void control_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) { - (void)userdata; + struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata; (void)stream; + aki_packet_stream_free(connection->control); } static void data_connection_closed_callback(void *userdata, struct aki_packet_stream *stream) { - (void)userdata; + struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata; (void)stream; + aki_packet_stream_free(connection->data); } static void default_packet_sent_callback(void *userdata, struct aki_packet *packet) @@ -43,8 +45,11 @@ static u8 packet_pool_callback(void *userdata, struct aki_packet *packet) { (void)userdata; struct bmu_node_connection *connection = (struct bmu_node_connection *)packet->userdata; - aki_packet_stream_send_packet(connection->data, packet); - return AKI_PACKET_POOL_KEEP; + if (connection->data->connected) { + aki_packet_stream_send_packet(connection->data, packet); + return AKI_PACKET_POOL_KEEP; + } + return AKI_PACKET_POOL_RETURN; } static void send_resume_packet(struct bmu_node_connection *connection, u64 seek_pos, u64 ts) @@ -69,11 +74,11 @@ static void send_start_packet(struct bmu_node_connection *connection) struct aki_packet *packet = aki_packet_create(); aki_packet_write_u8(packet, BIMU_CONTROL_START); u64 ts = aki_get_timestamp(); - u64 ats = ts + START_DELAY; + u64 ats = ts - START_DELAY; s64 seek_pos = connection->node->seek_pos; if (!connection->node->paused) { if (connection->node->start >= 0) { - if ((s64)ts < (connection->node->start - DELAY_GIVE)) { + if ((ats - DELAY_GIVE) < (u64)connection->node->start) { ats = connection->node->start; } else { seek_pos += ats - connection->node->start; @@ -93,6 +98,13 @@ static void send_start_packet(struct bmu_node_connection *connection) aki_packet_stream_send_packet(connection->control, packet); } +static void send_seek_packet(struct bmu_node_connection *connection) +{ + struct aki_packet *packet = aki_packet_create(); + aki_packet_write_u8(packet, BIMU_CONTROL_SEEK); + aki_packet_stream_send_packet(connection->control, packet); +} + static void control_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) { struct bmu_node_connection *connection = (struct bmu_node_connection *)userdata; @@ -119,7 +131,7 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st node->start = -1; } else if (paused == BIMU_PAUSED_AND_RUNNING) { node->paused = false; - node->start = aki_get_timestamp() + START_DELAY; + node->start = aki_get_timestamp() - START_DELAY; } struct bmu_node_connection *rconnection; al_array_foreach(node->connections, i, rconnection) { @@ -149,9 +161,7 @@ static void control_packet_callback(void *userdata, struct aki_packet_stream *st case BIMU_CONTROL_SEEK: { u64 pos = aki_packet_read_u64(packet); connection->seek_pos = pos; - struct aki_packet *npacket = aki_packet_create(); - aki_packet_write_u8(npacket, BIMU_CONTROL_SEEK); - aki_packet_stream_send_packet(connection->control, npacket); + send_seek_packet(connection); break; } } @@ -206,9 +216,9 @@ static void maybe_start_handler(struct bmu_node_connection *connection) if (connection->paused == BIMU_NOT_PAUSED) { if (node->start == -1) { node->paused = false; - node->start = aki_get_timestamp() + START_DELAY; + node->start = aki_get_timestamp() - START_DELAY; } - //send_resume_packet(connection, connection->seek_pos, node->start); + send_resume_packet(connection, connection->seek_pos, node->start); connection->paused = BIMU_PLAYING; } else if (connection->paused == BIMU_PAUSED) { connection->paused = BIMU_PAUSED_AND_RUNNING; @@ -314,7 +324,6 @@ u16 bmu_server_create_node(struct bmu_server *server, struct cch_entry *entry) node->start = -1; node->seek_pos = 0; aki_packet_pool_init(&node->pool, 16, server->loop, packet_pool_callback, node); - //aki_packet_pool_init(&node->pool, 1400, server->loop, packet_pool_callback, node); node->connection_id = 1; al_array_init(node->connections); node->server = server; diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 9799a95..ffb10dc 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -70,10 +70,9 @@ bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *cl { buf->mixer = mixer; buf->clock = clock; - buf->pts = camu_mixer_get_latency(buf->mixer); buf->buffered = false; buf->pause = PAUSE_PRE; - buf->ignore_desync = true; + buf->ignore_desync = false; #ifdef CAMU_AUDIO_BUFFER_FADE buf->fade_period = 0; buf->fade_offset = 0; @@ -198,7 +197,7 @@ void camu_audio_buffer_reset(struct camu_audio_buffer *buf) { al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED); al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED); - buf->pts = camu_clock_get_base_pts(buf->clock); + buf->pts = camu_clock_get_base_pts(buf->clock) + camu_mixer_get_latency(buf->mixer); buf->buffered = false; buf->pause = PAUSE_PRE; al_ring_buffer_reset(&buf->rb); @@ -259,7 +258,8 @@ static void handle_fade(struct camu_audio_buffer *buf, u8 *data, size_t size, bo size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req) { - if (camu_clock_is_paused(buf->clock)) { + f64 pts = camu_clock_get_pts(buf->clock); + if (pts < 0.0) { #ifdef CAMU_AUDIO_BUFFER_FADE if (buf->pause == PAUSE_PLAYING) { // Fade out. Signified by >0 fade_period and PAUSE_FADING. @@ -299,7 +299,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re #ifdef CAMU_AUDIO_BUFFER_FADE if (buf->pause != PAUSE_FADING) { #endif - f64 pts = camu_clock_get_pts(buf->clock) - buf->pts; + pts -= buf->pts; // We can assume a call to audio_buffer_read will happen before // the user NEEDS the data, so, we can't accurately assess the sync here. if (UNLIKELY(!buf->ignore_desync && (buf->pause == PAUSE_PRE || buf->pause == PAUSE_UNPAUSED))) { diff --git a/src/buffer/clock.c b/src/buffer/clock.c index c9b18ce..a1c3e83 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -1,55 +1,68 @@ +#include <al/lib.h> + #include "clock.h" -void camu_clock_init(struct camu_clock *clock) +void camu_clock_set(struct camu_clock *clock, u64 base, u64 start) { - f64 tick = aki_get_tick(); - al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED); - al_atomic_f64_store(&clock->start, tick, AL_ATOMIC_RELAXED); - al_atomic_f64_store(&clock->base, 0.0, AL_ATOMIC_RELAXED); + al_atomic_f64_store(&clock->base, base / 1000000.f, AL_ATOMIC_RELAXED); + al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELAXED); + al_atomic_u64_store(&clock->start, start, AL_ATOMIC_RELAXED); + camu_clock_resume(clock); } -void camu_clock_pause(struct camu_clock *clock) +void camu_clock_resume(struct camu_clock *clock) { - f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED); - if (pause != -1.0) return; - al_atomic_f64_store(&clock->pause, aki_get_tick(), AL_ATOMIC_RELAXED); + u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE); + if (start == 0L) return; + u64 ts = aki_get_timestamp(); + al_assert(start < ts); + f64 diff = (ts - start) / 1000000.f; + f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED); + if (tick == 0.0) { + tick = aki_get_tick() - al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED); + } + al_atomic_f64_store(&clock->tick, tick + diff, AL_ATOMIC_RELAXED); + al_atomic_u64_store(&clock->start, 0L, AL_ATOMIC_RELEASE); } -void camu_clock_resume(struct camu_clock *clock) +static f64 get_pts_internal(struct camu_clock *clock) { - f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED); - if (pause == -1.0) return; - f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED); - f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED); - if (start == 0.0) start = (pause - base); - al_atomic_f64_store(&clock->pause, -1.0, AL_ATOMIC_RELAXED); - al_atomic_f64_store(&clock->start, aki_get_tick() - (pause - start), AL_ATOMIC_RELAXED); + f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED); + return aki_get_tick() - tick; } -void camu_clock_seek(struct camu_clock *clock, f64 pos) +void camu_clock_pause(struct camu_clock *clock) { - al_atomic_f64_store(&clock->start, 0.0, AL_ATOMIC_RELAXED); - al_atomic_f64_store(&clock->base, pos, AL_ATOMIC_RELAXED); - f64 tick = aki_get_tick(); - al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED); + u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE); + if (start != 0L) return; + al_atomic_f64_store(&clock->base, get_pts_internal(clock), AL_ATOMIC_RELAXED); + al_atomic_u64_store(&clock->start, aki_get_timestamp(), AL_ATOMIC_RELEASE); +} + +void camu_clock_seek(struct camu_clock *clock, u64 base, u64 start) +{ + (void)clock; + (void)base; + (void)start; } bool camu_clock_is_paused(struct camu_clock *clock) { - f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED); - return pause != -1.0; + u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED); + return start != 0L; } f64 camu_clock_get_base_pts(struct camu_clock *clock) { - f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED); - return base; + return al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED); } f64 camu_clock_get_pts(struct camu_clock *clock) { - f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED); - f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED); - if (pause != -1.f) return pause - start; - return aki_get_tick() - start; + u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED); + if (start != 0L) return -1.0; + f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED); + f64 pts = get_pts_internal(clock); + if (pts < base) return -1.0; + return pts; } diff --git a/src/buffer/clock.h b/src/buffer/clock.h index 9f34807..e39a963 100644 --- a/src/buffer/clock.h +++ b/src/buffer/clock.h @@ -4,13 +4,14 @@ #include <al/atomic.h> struct camu_clock { - atomic_f64 base, start, pause; + atomic_f64 base, tick; + atomic_u64 start; }; -void camu_clock_init(struct camu_clock *clock); -void camu_clock_pause(struct camu_clock *clock); +void camu_clock_set(struct camu_clock *clock, u64 base, u64 start); void camu_clock_resume(struct camu_clock *clock); -void camu_clock_seek(struct camu_clock *clock, f64 pos); +void camu_clock_pause(struct camu_clock *clock); +void camu_clock_seek(struct camu_clock *clock, u64 base, u64 start); bool camu_clock_is_paused(struct camu_clock *clock); f64 camu_clock_get_base_pts(struct camu_clock *clock); f64 camu_clock_get_pts(struct camu_clock *clock); diff --git a/src/buffer/video.c b/src/buffer/video.c index 6c8d6e3..3d3961b 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -137,6 +137,7 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf) bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out) { f64 pts = buf->single_frame ? 0.0 : camu_clock_get_pts(buf->clock); + if (pts < 0.0) pts = camu_clock_get_base_pts(buf->clock); u8 ret = buf->queue->read(buf->queue, pts, out); if (ret == CAMU_QUEUE_EOF || (ret == CAMU_QUEUE_OK && buf->single_frame)) { buf->callback(buf->userdata, CAMU_BUFFER_EOF); diff --git a/src/fruits/cmv/meson.build b/src/fruits/cmv/meson.build deleted file mode 100644 index 1c9ed2d..0000000 --- a/src/fruits/cmv/meson.build +++ /dev/null @@ -1,10 +0,0 @@ -cmv_src = ['cmv2.c'] -cmv_deps = [common_deps, buffer, cache, av, screen, render, mixer, bimu, libsink, cap] - -use_tui = true -if use_tui - cmv_src += ['tui.c'] - cmv_deps += [dependency('notcurses')] -endif - -executable('cmv', cmv_src, dependencies: cmv_deps) diff --git a/src/fruits/cmv/tui.c b/src/fruits/cmv/tui.c deleted file mode 100644 index 8315a13..0000000 --- a/src/fruits/cmv/tui.c +++ /dev/null @@ -1,165 +0,0 @@ -#include <aki/file.h> -#include <al/log.h> -#include <sys/ioctl.h> -#include <notcurses/direct.h> - -#include "tui.h" - -static struct cmv_tui *tui_global; - -static void update_term_size(struct cmv_tui *tui) -{ - struct winsize w; - ioctl(tui->fd, TIOCGWINSZ, &w); - tui->width = w.ws_col; - tui->height = w.ws_row; -} - -static void handle_winch(s32 sig) -{ - (void)sig; - struct cmv_tui *tui = tui_global; - aki_mutex_lock(&tui->mutex); - tui->update_size = true; - aki_mutex_unlock(&tui->mutex); -} - -static void input_poll_callback(void *userdata, s32 revents) -{ - struct cmv_tui *tui = (struct cmv_tui *)userdata; - (void)revents; - struct ncinput input; - u32 ret; - while (1) { - ret = ncdirect_get_nblock(tui->dir, &input); - if (ret == (u32)-1 || ret == 0) break; - if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) { - } - } -} - -bool tui_init(struct cmv_tui *tui, struct aki_event_loop *loop) -{ - tui->update_size = true; - tui->width = 0; - tui->height = 0; - tui->prev_lines = 0; - - tui->fd = fileno(stdout); - tui->out = fdopen(tui->fd, "w"); - if (!tui->out) return false; - tui->buffer = al_malloc(256); - - //if (!(tui->dir = ncdirect_init(NULL, tui->out, 0))) { - // fclose(tui->out); - // return false; - //} - - //aki_poll_init(&tui->poll, ncdirect_inputready_fd(tui->dir), AKI_POLL_READ, input_poll_callback, tui); - //aki_poll_start(&tui->poll, loop); - - tui_global = tui; - aki_mutex_init(&tui->mutex); - signal(SIGWINCH, handle_winch); - - // Disable buffer. - setbuf(tui->out, NULL); - - // Disable echo. - tcgetattr(tui->fd, &tui->term); - tui->term.c_lflag &= ~ECHO; - tcsetattr(tui->fd, 0, &tui->term); - - al_queue_init(&tui->log_buffer, 256); - - return true; -} - -void tui_push_log_msg(struct cmv_tui *tui, char *msg) -{ - al_queue_push(&tui->log_buffer, msg); -} - -static void draw_now_playing(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock) -{ - char *buffer = tui->buffer; - f64 pts = camu_clock_get_pts(clock); - f64 duration = bmu_local_get_duration(runner); - s32 text = 0; - s64 minute = (s64)(pts / 60L); - s64 hour = minute / 60L; - minute -= hour * 60L; - text += al_sprintf(buffer + text, "["); - if (hour > 0) text += al_sprintf(buffer + text, "%.2ld:", hour); - text += al_sprintf(buffer + text, "%.2ld:%.2ld/", minute, (s64)pts % 60L); - minute = (s64)(duration / 60L); - hour = minute / 60L; - minute -= hour * 60L; - if (hour > 0) text += al_sprintf(buffer + text, "%.2ld:", hour); - text += al_sprintf(buffer + text, "%.2ld:%.2ld", minute, (s64)duration % 60L); - text += al_sprintf(buffer + text, "]"); - s32 parts = 0; - if (pts > 0.0 && duration != 0.0) { - if (pts > duration) pts = duration; - f64 percent = pts / duration; - parts = ((tui->width - text) * percent); - } else if (duration == 0.0) { - parts = (tui->width - text); - } - al_memset(buffer + text, '-', parts); - buffer[text + parts] = '\0'; - fprintf(tui->out, "%s", buffer); -} - -void tui_draw(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock) -{ - aki_mutex_lock(&tui->mutex); - if (tui->update_size) { - update_term_size(tui); - tui->update_size = false; - } - aki_mutex_unlock(&tui->mutex); - - // Erase line. - fprintf(tui->out, "\r\033[K"); - for (u32 i = 0; i < tui->prev_lines; i++) { - // Up one, erase line. - fprintf(tui->out, "\033[A\r\033[K"); - } - - char *msg; - while (al_queue_pop(&tui->log_buffer, (void **)&msg)) { - fprintf(tui->out, "%s\n", msg); - al_free(msg); - } - - // Simply printing each log message with a newline will leave - // a one character high space in the bottom of the terminal window. - // We use that to display a very basic status line (like mpv). - // Trying to use more than that one line space without a more complete - // terminal interface gets scuffed very fast. - - tui->prev_lines = 0; - - if (runner) { - draw_now_playing(tui, runner, clock); - } -} - -void tui_close(struct cmv_tui *tui) -{ - if (!tui->out) return; - al_free(tui->buffer); - char *msg; - while (al_queue_pop(&tui->log_buffer, (void **)&msg)) { - al_free(msg); - } - al_queue_free(&tui->log_buffer); - aki_mutex_destroy(&tui->mutex); - fprintf(tui->out, "\n"); - fflush(tui->out); - tui->term.c_lflag |= ECHO; - tcsetattr(tui->fd, 0, &tui->term); - //ncdirect_stop(tui->dir); - fclose(tui->out); -} diff --git a/src/fruits/cmv/tui.h b/src/fruits/cmv/tui.h deleted file mode 100644 index 6898539..0000000 --- a/src/fruits/cmv/tui.h +++ /dev/null @@ -1,28 +0,0 @@ -#pragma once - -#include <al/queue.h> -#include <aki/event_loop.h> -#include <termios.h> - -#include "../../bimu/local.h" -#include "../../buffer/video.h" - -struct cmv_tui { - s32 fd; - FILE *out; - char *buffer; - struct ncdirect *dir; - struct aki_poll poll; - struct aki_mutex mutex; - struct termios term; - bool update_size; - u32 width; - u32 height; - u32 prev_lines; - queue log_buffer; -}; - -bool tui_init(struct cmv_tui *tui, struct aki_event_loop *loop); -void tui_push_log_msg(struct cmv_tui *tui, char *msg); -void tui_draw(struct cmv_tui *tui, struct bmu_local *runner, struct camu_clock *clock); -void tui_close(struct cmv_tui *tui); diff --git a/src/fruits/sink/meson.build b/src/fruits/sink/meson.build new file mode 100644 index 0000000..c409857 --- /dev/null +++ b/src/fruits/sink/meson.build @@ -0,0 +1,10 @@ +sink_src = ['sink.c'] +sink_deps = [common_deps, buffer, cache, av, screen, render, mixer, libsink] + +use_tui = false +if use_tui + sink_src += [] + sink_deps += [dependency('notcurses')] +endif + +executable('sink', sink_src, dependencies: sink_deps) diff --git a/src/fruits/cmv/cmv2.c b/src/fruits/sink/sink.c index 89e3427..d0455a4 100644 --- a/src/fruits/cmv/cmv2.c +++ b/src/fruits/sink/sink.c @@ -1,11 +1,5 @@ -#define CMV_USE_TUI 0 -#define CMV_USE_SOCKET 0 - #include <al/log.h> #include <aki/event_loop.h> -#if CMV_USE_SOCKET -#include <aki/line_processor.h> -#endif #include "../../codec/libav/common.h" @@ -15,17 +9,11 @@ #include "../../render/renderer_libplacebo.h" //#include "../../render/renderer_tiger.h" -#include "../../libsink/sink2.h" +#include "../../libsink/sink.h" #include "../../tree/common.h" #include "../../shoki/src/search.h" -#include "../cap/cap.h" - -#if CMV_USE_TUI -#include "tui.h" -#endif - struct cmv { s32 quit; struct camu_screen scr; @@ -33,19 +21,8 @@ struct cmv { struct camu_mixer mixer; struct aki_event_loop loop; struct camu_sink sink; - struct cap_runner cap; -#if CMV_USE_SOCKET - struct aki_socket socket; - struct aki_line_processor pro; -#endif -#if CMV_USE_TUI - struct cmv_tui tui; - struct aki_timer timer; -#endif }; -static str *default_list = al_str_c("default"); - static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque) { struct cmv *c = (struct cmv *)userdata; @@ -85,20 +62,11 @@ static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque) case CAMU_SINK_SWAP_BUFFER: switch (type) { case CAMU_SINK_AUDIO: { - struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque; - camu_mixer_remove_buffer(&c->mixer, buf); - if (cap_list_set_completed(&c->cap, default_list)) { - al_log_debug("cmv", "Audio buffers swapped (gapless)."); - return CAMU_SINK_BUFFERS_SWAPPED; - } else { - al_log_debug("cmv", "Audio buffer removed."); - } break; } } break; case CAMU_SINK_SET_BUFFERED: { - cap_list_pump(&c->cap, default_list, true); break; } case CAMU_SINK_START: @@ -131,10 +99,6 @@ static u8 sink_callback(void *userdata, u8 op, u8 type, void *opaque) } break; case CAMU_SINK_EXIT: -#if CMV_USE_SOCKET - aki_line_processor_stop(&c->pro); - aki_socket_close(&c->socket); -#endif camu_sink_close(&c->sink); aki_event_loop_break(&c->loop); break; @@ -147,10 +111,8 @@ static void screen_callback(void *userdata, u8 op, f64 float0) struct cmv *c = (struct cmv *)userdata; switch (op) { case CAMU_SCREEN_NEXT: - cap_list_skip(&c->cap, default_list, 1); break; case CAMU_SCREEN_PREVIOUS: - cap_list_skip(&c->cap, default_list, -1); break; case CAMU_SCREEN_TOGGLE_PAUSE: camu_sink_toggle_pause(&c->sink); @@ -164,76 +126,6 @@ static void screen_callback(void *userdata, u8 op, f64 float0) } } -/* -static bool cap_callback(void *userdata, u8 op, str *name, str *unique_id, void *opaque) -{ - struct cmv *c = (struct cmv *)userdata; - (void)name; - switch (op) { - case CAP_BUFFER: - return camu_sink_local_buffer(&c->sink, unique_id, (struct cch_entry *)opaque); - case CAP_SET: - camu_sink_local_set(&c->sink, unique_id); - al_log_info("cmv", "Now playing: %.*s.", AL_STR_PRINTF(unique_id)); - break; - case CAP_SWAP: - camu_sink_local_swap(&c->sink, unique_id); - break; - case CAP_UNLOAD: - camu_sink_local_unload(&c->sink, unique_id); - break; - } - return true; -} -*/ - -#if CMV_USE_SOCKET -static u8 line_callback(void *userdata, str *line) -{ - struct cmv *c = (struct cmv *)userdata; - if (al_str_eq(line, al_str_c(";NEXT"))) { - cap_list_skip(&c->cap, default_list, 1); - } else if (al_str_eq(line, al_str_c(";PREV"))) { - cap_list_skip(&c->cap, default_list, -1); - } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) { - cap_list_shuffle(&c->cap, default_list); - } else { - cap_list_add(&c->cap, default_list, line); - } - return AKI_LINE_PROCESSOR_CONTINUE; -} -#endif - -#if CMV_USE_TUI -static void timer_callback(void *userdata, struct aki_timer *timer) -{ - struct cmv *c = (struct cmv *)userdata; - struct camu_sink_entry *entry = camu_sink_get_current(&c->sink); - if (entry) { - tui_draw(&c->tui, &((struct camu_sink_local *)entry)->runner, &entry->clock); - } else { - tui_draw(&c->tui, NULL, NULL); - } - camu_sink_return_current(&c->sink); - aki_timer_again(timer); -} - -static s32 log_callback(void *userdata, char *s) -{ - struct cmv *c = (struct cmv *)userdata; - tui_push_log_msg(&c->tui, s); - return 0; -} - -#ifdef HAVE_FFMPEG -static void lav_log_callback(void *userdata, s32 level, const char *fmt, va_list args) -{ - (void)userdata; - if (level < AV_LOG_VERBOSE) al_logv("warn", "lav_internal", (char *)fmt, args); -} -#endif -#endif - static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) { struct cmv *c = (struct cmv *)userdata; @@ -258,16 +150,6 @@ s32 main(s32 argc, char *argv[]) aki_event_loop_init(&c.loop); -#if CMV_USE_TUI - if (!tui_init(&c.tui, &c.loop)) { - goto err; - }; - al_set_print(log_callback, &c); -#ifdef HAVE_FFMPEG - camu_lav_set_log_callback(lav_log_callback); -#endif -#endif - c.scr.callback = screen_callback; c.scr.userdata = &c; if (!camu_screen_init(&c.scr) || !camu_screen_create_window(&c.scr, "cmv")) { @@ -280,10 +162,6 @@ s32 main(s32 argc, char *argv[]) } c.renderer->render(c.renderer, &c.scr); -#if CAP_USE_PYTHON - bool py_init = argc == 1 ? sho_python_init() : false; -#endif - camu_mixer_init(&c.mixer, (struct camu_audio *)&audio_plugin_miniaudio); c.mixer.audio->configure_stream(c.mixer.audio, NULL); @@ -292,33 +170,6 @@ s32 main(s32 argc, char *argv[]) c.sink.userdata = &c; camu_sink_connect(&c.sink, al_str_c("127.0.0.1"), TREE_PORT); - /* - cap_init(&c.cap, cap_callback, &c); - cap_make_list(&c.cap, default_list); - for (s32 i = 1; i < argc; i++) { - cap_list_add(&c.cap, default_list, al_str_c(argv[i])); - } - */ - -#if CMV_USE_SOCKET - c.socket.type = AKI_SOCKET_UNIX; - aki_socket_init(&c.socket); - aki_socket_set_blocking(&c.socket, false); - c.pro.callback = line_callback; - c.pro.userdata = &c; - aki_line_processor_init(&c.pro, al_str_c("\n")); - aki_line_processor_open_socket(&c.pro, &c.socket); - if (aki_socket_listen(&c.socket, al_str_c("/tmp/cmv_sock"), 0)) { - aki_line_processor_run(&c.pro, &c.loop); - } -#endif - -#if CMV_USE_TUI - aki_timer_init(&c.timer, &c.loop, timer_callback, &c); - aki_timer_set_repeat(&c.timer, 0.05); - aki_timer_again(&c.timer); -#endif - struct aki_thread thread0; aki_thread_create(&thread0, event_loop_thread, &c); @@ -340,23 +191,10 @@ s32 main(s32 argc, char *argv[]) c.renderer->free(&c.renderer); camu_screen_close(&c.scr); - //cap_close(&c.cap); - -#if CAP_USE_PYTHON - if (py_init) sho_python_close(); -#endif - -#if CMV_USE_TUI - tui_close(&c.tui); -#endif - aki_common_close(); return EXIT_SUCCESS; err: -#if CMV_USE_TUI - tui_close(&c.tui); -#endif aki_event_loop_destroy(&c.loop); aki_common_close(); return EXIT_FAILURE; 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); diff --git a/src/tree/list.c b/src/tree/list.c index 016c0de..8102e05 100644 --- a/src/tree/list.c +++ b/src/tree/list.c @@ -1,6 +1,7 @@ #ifdef AKIYO_HAS_CURL #include "../cache/handlers/http.h" #endif +#include "../cache/handlers/file.h" #include "../libsink/common.h" @@ -17,9 +18,30 @@ void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name) list->tree = tree; } +static void send_buffer_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry) +{ + struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER); + aki_packet_write_str(packet, al_str_c("127.0.0.1")); + aki_packet_write_s32(packet, TREE_STREAM_PORT); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); +} + +static void send_set_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry) +{ + struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); +} + void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink) { al_array_push(list->sinks, sink); + if (list->set == list->current) { + struct tree_list_entry *entry = al_array_at(list->entries, list->current); + send_buffer_cmd(list, sink, entry); + send_set_cmd(list, sink, entry); + } } void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink) @@ -62,51 +84,55 @@ bool tree_list_add(struct tree_list *list, str *unique_id, u32 index) return true; } -void tree_list_skip(struct tree_list *list, s32 n) +bool tree_list_add_external(struct tree_list *list, str *path) { - s32 size = (s32)list->entries.size; - if (list->current + n < 0 || list->current + n >= size) { - return; + struct tree_server *tree = list->tree; + + struct tree_list_entry *entry = al_alloc_object(struct tree_list_entry); + + entry->entry = cch_handler_file_create(path); + if (!entry->entry) { + al_free(entry); + return false; } - list->current += n; + + entry->node_id = bmu_server_create_node(&tree->streams.server, entry->entry); + + entry->buffer_requested = false; + + al_array_push(list->entries, entry); + tree_list_pump(list); -} -static void send_buffer_cmd(struct tree_list *list, struct tree_list_entry *entry) -{ - struct aki_packet *packet; - struct tree_sink *sink; - al_array_foreach(list->sinks, i, sink) { - packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER); - aki_packet_write_str(packet, al_str_c("127.0.0.1")); - aki_packet_write_s32(packet, TREE_STREAM_PORT); - aki_packet_write_u16(packet, entry->node_id); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - } + return true; } -static void send_set_cmd(struct tree_list *list, struct tree_list_entry *entry) +void tree_list_skip(struct tree_list *list, s32 n) { - struct aki_packet *packet; - struct tree_sink *sink; - al_array_foreach(list->sinks, i, sink) { - packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET); - aki_packet_write_u16(packet, entry->node_id); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + s32 size = (s32)list->entries.size; + if (list->current + n < 0 || list->current + n >= size) { + return; } + list->current += n; + tree_list_pump(list); } void tree_list_pump(struct tree_list *list) { s32 size = (s32)list->entries.size; if (size <= list->current) return; + struct tree_sink *sink; if (list->set != list->current) { struct tree_list_entry *entry = al_array_at(list->entries, list->current); if (!entry->buffer_requested) { - send_buffer_cmd(list, entry); + al_array_foreach(list->sinks, i, sink) { + send_buffer_cmd(list, sink, entry); + } entry->buffer_requested = true; } - send_set_cmd(list, entry); + al_array_foreach(list->sinks, i, sink) { + send_set_cmd(list, sink, entry); + } list->set = list->current; } } diff --git a/src/tree/list.h b/src/tree/list.h index ac7f8ee..3defe9f 100644 --- a/src/tree/list.h +++ b/src/tree/list.h @@ -27,5 +27,6 @@ void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name) void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink); void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink); bool tree_list_add(struct tree_list *list, str *unique_id, u32 index); +bool tree_list_add_external(struct tree_list *list, str *path); void tree_list_skip(struct tree_list *list, s32 n); void tree_list_pump(struct tree_list *list); diff --git a/src/tree/meson.build b/src/tree/meson.build index e9b8670..e566687 100644 --- a/src/tree/meson.build +++ b/src/tree/meson.build @@ -3,5 +3,5 @@ tree_src = [ 'resource_manager.c', 'list.c' ] -tree_deps = [common_deps, shoki, cache, cap, bimu, av] +tree_deps = [common_deps, shoki, cache, bimu_server, av] executable('tree', sources: tree_src, dependencies: tree_deps) diff --git a/src/tree/tree.c b/src/tree/tree.c index 014397f..e671d9c 100644 --- a/src/tree/tree.c +++ b/src/tree/tree.c @@ -270,11 +270,11 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection al_array_foreach(tree->sinks, i, sink) { if (sink->conn == conn) { al_log_info("tree", "Sink removed."); - cleanup_sink(sink); al_array_remove_at_iter(tree->sinks, i); struct tree_user *user = al_array_at(tree->users, 0); struct tree_list *list = al_array_at(user->lists, 0); tree_list_remove_sink(list, sink); + cleanup_sink(sink); break; } } @@ -337,6 +337,22 @@ static bool open_db(struct tree_server *tree, str *path) return true; } +#if TREE_USE_SOCKET +static u8 line_callback(void *userdata, str *line) +{ + struct tree_server *tree = (struct tree_server *)userdata; + if (al_str_eq(line, al_str_c(";NEXT"))) { + } else if (al_str_eq(line, al_str_c(";PREV"))) { + } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) { + } else { + struct tree_user *user = al_array_at(tree->users, 0); + struct tree_list *list = al_array_at(user->lists, 0); + tree_list_add_external(list, line); + } + return AKI_LINE_PROCESSOR_CONTINUE; +} +#endif + static void sigint_handler(s32 signum) { (void)signum; @@ -381,6 +397,19 @@ s32 main(void) bmu_server_init(&tree.streams.server); bmu_server_listen(&tree.streams.server, &tree.loop, al_str_c("0.0.0.0"), TREE_STREAM_PORT); +#if TREE_USE_SOCKET + tree.socket.type = AKI_SOCKET_UNIX; + aki_socket_init(&tree.socket); + aki_socket_set_blocking(&tree.socket, false); + tree.pro.callback = line_callback; + tree.pro.userdata = &tree; + aki_line_processor_init(&tree.pro, al_str_c("\n")); + aki_line_processor_open_socket(&tree.pro, &tree.socket); + if (aki_socket_listen(&tree.socket, al_str_c("/tmp/tree_sock"), 0)) { + aki_line_processor_run(&tree.pro, &tree.loop); + } +#endif + aki_event_loop_run(&tree.loop); if (py_init) sho_python_close(); diff --git a/src/tree/tree.h b/src/tree/tree.h index 478306a..2b61e3b 100644 --- a/src/tree/tree.h +++ b/src/tree/tree.h @@ -1,9 +1,14 @@ #pragma once +#define TREE_USE_SOCKET 1 + #include <aki/rpc2.h> #include <sho/post.h> #include <sho/post_cache.h> #include <sho/search.h> +#if TREE_USE_SOCKET +#include <aki/line_processor.h> +#endif #include "../bimu/server.h" @@ -47,4 +52,8 @@ struct tree_server { struct { struct bmu_server server; } streams; +#if TREE_USE_SOCKET + struct aki_socket socket; + struct aki_line_processor pro; +#endif }; |