diff options
| author | 2024-12-18 11:06:44 -0500 | |
|---|---|---|
| committer | 2024-12-18 11:06:44 -0500 | |
| commit | 3d55d2722a3129449ab1418e73abd97caa7fd2ae (patch) | |
| tree | 0cba3f876c9b339d20d8e4a38a7a5d786fa2102d | |
| parent | 692785bc9da6904cf17e986fb034730ed3d78231 (diff) | |
| download | camu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.tar.gz camu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.tar.bz2 camu-3d55d2722a3129449ab1418e73abd97caa7fd2ae.zip | |
Extend server resource loading, work on client
Also, cleanup in preparation for changing code style.
Signed-off-by: Andrew Opalach <andrew@akon.city>
35 files changed, 666 insertions, 372 deletions
@@ -41,11 +41,11 @@ ] }, "locked": { - "lastModified": 1734093295, - "narHash": "sha256-hSwgGpcZtdDsk1dnzA0xj5cNaHgN9A99hRF/mxMtwS4=", + "lastModified": 1734344598, + "narHash": "sha256-wNX3hsScqDdqKWOO87wETUEi7a/QlPVgpC/Lh5rFOuA=", "owner": "nix-community", "repo": "home-manager", - "rev": "66c5d8b62818ec4c1edb3e941f55ef78df8141a8", + "rev": "83ecd50915a09dca928971139d3a102377a8d242", "type": "github" }, "original": { @@ -57,11 +57,11 @@ }, "nixos-hardware": { "locked": { - "lastModified": 1733861262, - "narHash": "sha256-+jjPup/ByS0LEVIrBbt7FnGugJgLeG9oc+ivFASYn2U=", + "lastModified": 1734352517, + "narHash": "sha256-mfv+J/vO4nqmIOlq8Y1rRW8hVsGH3M+I2ESMjhuebDs=", "owner": "NixOS", "repo": "nixos-hardware", - "rev": "cf737e2eba82b603f54f71b10cb8fd09d22ce3f5", + "rev": "b12e314726a4226298fe82776b4baeaa7bcf3dcd", "type": "github" }, "original": { @@ -110,11 +110,11 @@ }, "nixpkgs_2": { "locked": { - "lastModified": 1733940404, - "narHash": "sha256-Pj39hSoUA86ZePPF/UXiYHHM7hMIkios8TYG29kQT4g=", + "lastModified": 1734424634, + "narHash": "sha256-cHar1vqHOOyC7f1+tVycPoWTfKIaqkoe1Q6TnKzuti4=", "owner": "NixOS", "repo": "nixpkgs", - "rev": "5d67ea6b4b63378b9c13be21e2ec9d1afc921713", + "rev": "d3c42f187194c26d9f0309a8ecc469d6c878ce33", "type": "github" }, "original": { diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 0d78382..2cedcc8 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -11,7 +11,7 @@ #define BUFFER_SIZE 9.0 #define BUFFER_MARK_MIN 4.3 // Must be a most half of the buffer size. -#define BUFFER_MARK_BUFFERED 4.0 +#define BUFFER_MARK_BUFFERED 2.5 #ifdef CAMU_AUDIO_BUFFER_FADE #define FADE_STEP(fmt) (1.75f / (fmt)->sample_rate) @@ -101,10 +101,10 @@ void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency) } #ifdef _DEBUG_ -#define OCCUPIED_SECONDS_DEBUG(buf) \ +#define BUFFERED_SECONDS_DEBUG(buf) \ camu_audio_format_bytes_to_sec(&buf->fmt.req, al_ring_buffer_occupied(&buf->rb)) #else -#define OCCUPIED_SECONDS_DEBUG(buf) 0 +#define BUFFERED_SECONDS_DEBUG(buf) 0 #endif static inline bool frame_is_late(struct camu_clock *clock, f64 base, f64 pts, f64 duration) @@ -137,7 +137,7 @@ static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32 size_t space = al_ring_buffer_space(&buf->rb); if (!buf->buffered && buf->size - space > buf->mark.buffered) { - al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf)); + al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", BUFFERED_SECONDS_DEBUG(buf)); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); buf->buffered = true; } @@ -213,7 +213,7 @@ void camu_audio_buffer_flush(struct camu_audio_buffer *buf) al_log_debug("audio_buffer", "Buffer filled by flush."); } if (!buf->buffered) { - al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf)); + al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", BUFFERED_SECONDS_DEBUG(buf)); buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED); buf->buffered = true; } diff --git a/src/cache/threaded_waits.c b/src/cache/threaded_waits.c index a93fe85..10f08f0 100644 --- a/src/cache/threaded_waits.c +++ b/src/cache/threaded_waits.c @@ -20,8 +20,7 @@ static bool wait_range_satisfied(struct cch_backing *backing, struct cch_handler return false; } -bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, - struct cch_handler_wait *wait) +bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait) { aki_mutex_lock(&handler->mutex); bool canceled = handler->disabled; diff --git a/src/cache/threaded_waits.h b/src/cache/threaded_waits.h index edbee03..ee0838d 100644 --- a/src/cache/threaded_waits.h +++ b/src/cache/threaded_waits.h @@ -6,8 +6,7 @@ #include "backing.h" void cch_threaded_waits_init(struct cch_handler *handler); -bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, - struct cch_handler_wait *wait); +bool cch_threaded_wait_for_range(struct cch_handler *handler, struct cch_backing *backing, struct cch_handler_wait *wait); void cch_threaded_waits_signal_any(struct cch_handler *handler); void cch_threaded_wait_disable(struct cch_handler_wait *wait); void cch_threaded_waits_disable_all(struct cch_handler *handler); diff --git a/src/codec/ffmpeg/common.c b/src/codec/ffmpeg/common.c index b7a232e..c147d81 100644 --- a/src/codec/ffmpeg/common.c +++ b/src/codec/ffmpeg/common.c @@ -9,30 +9,32 @@ void camu_ff_set_log_callback(void (*callback)(void *, int, const char *, va_lis av_log_set_callback(callback); } -static char *av_log_buf = NULL; -static s32 av_log_pos = 0; +static char *buf = NULL; +static s32 pos = 0; + +// If a line takes more than 2 steps to print, make sure we stay within AL_LOG_MESSAGE_SIZE. +#define CHUNK_SIZE (AL_LOG_MESSAGE_SIZE / 2) static void av_log_callback(void *userdata, int level, const char *fmt, va_list args) { (void)userdata; - al_assert(av_log_buf); + al_assert(buf); if (level < AV_LOG_DEBUG) { - av_log_pos += al_vsnprintf(&av_log_buf[av_log_pos], 512, fmt, args); - if (av_log_buf[av_log_pos - 1] == '\n' || av_log_pos >= 512) { - al_log_info("ff_log", av_log_buf); - av_log_pos = 0; + pos += al_vsnprintf(&buf[pos], CHUNK_SIZE, fmt, args); + if (buf[pos - 1] == '\n' || pos >= CHUNK_SIZE) { + al_log_info("ff_log", buf); + pos = 0; } } } void camu_ff_set_default_log_callback() { - av_log_buf = (char *)al_malloc(1024); + buf = (char *)al_malloc(AL_LOG_MESSAGE_SIZE); camu_ff_set_log_callback(av_log_callback); } - void camu_ff_free_default_log_callback() { - al_free(av_log_buf); + al_free(buf); } diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index 51eca40..47fcc0c 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -92,7 +92,7 @@ static s32 send_packet(struct camu_ff_decoder *av, AVPacket *pkt) s32 ret = avcodec_send_packet(av->codec_context, pkt); if (ret < 0 && ret != AVERROR(EAGAIN) && ret != AVERROR_EOF) { - al_log_error("ff_decoder", "Error sending packet to the decoder: (%s).", av_err2str(ret)); + al_log_error("ff_decoder", "Error sending packet to the decoder (%s).", av_err2str(ret)); } return ret; @@ -123,7 +123,7 @@ static s32 receive_frames(struct camu_ff_decoder *av) al_free(frame); // Checking for EAGAIN should prevent an infinite loop. if (ret == AVERROR(EAGAIN) || ret == AVERROR_EOF) break; - al_log_error("ff_decoder", "Error receiving packet from the decoder: (%s).", av_err2str(ret)); + al_log_error("ff_decoder", "Error receiving packet from the decoder (%s).", av_err2str(ret)); continue; } // Track pts and duration of the previous frame so we can handle multiple frames in a single packet. diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c index 24c1ba6..f756168 100644 --- a/src/codec/ffmpeg/demuxer.c +++ b/src/codec/ffmpeg/demuxer.c @@ -3,7 +3,7 @@ #include "demuxer.h" #include "avio.h" -#define DEMUX_BUF_SIZE (4096 * 4) +#define DEMUX_BUFFER_SIZE (4096 * 4) static void close_internal(struct camu_ff_demuxer *av) { @@ -28,8 +28,8 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl goto err; } - u8 *buf = (u8 *)av_malloc(DEMUX_BUF_SIZE); - av->io_context = avio_alloc_context(buf, DEMUX_BUF_SIZE, 0, handle, camu_avio_read, NULL, camu_avio_seek); + u8 *buf = (u8 *)av_malloc(DEMUX_BUFFER_SIZE); + av->io_context = avio_alloc_context(buf, DEMUX_BUFFER_SIZE, 0, handle, camu_avio_read, NULL, camu_avio_seek); if (!av->io_context) { al_log_error("ff_demuxer", "Failed to create custom io context."); goto err; @@ -70,31 +70,30 @@ static bool ff_demuxer_init(struct camu_demuxer *demux, struct cch_handle *handl } #define GUESS_STREAM_IS_IMAGE(stream) \ - (stream->duration >= 0 && stream->nb_frames <= 1 && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0)) + (stream->duration >= 0L && stream->nb_frames <= 1L && (stream->avg_frame_rate.den == 0 && stream->r_frame_rate.den > 0)) - av->duration = 0; + s64 default_duration = av->format_context->duration < 0L ? 0L : av->format_context->duration; + av->duration = 0Lu; for (u32 i = 0; i < av->format_context->nb_streams; i++) { AVStream *stream = av->format_context->streams[i]; enum AVMediaType type = stream->codecpar->codec_type; - u64 duration = 0; + u64 duration = 0Lu; if (type == AVMEDIA_TYPE_AUDIO || type == AVMEDIA_TYPE_VIDEO || type == AVMEDIA_TYPE_SUBTITLE) { // Try to detect attached images. if (type == AVMEDIA_TYPE_VIDEO && GUESS_STREAM_IS_IMAGE(stream)) { - al_log_info("ff_demuxer", "Guessing that stream #%u is an image.", i); + al_log_info("ff_demuxer", "Assuming stream #%u is an image.", i); } else { - if (stream->duration < 0) { - al_log_warn("ff_demuxer", "Stream #%u has an invalid duration (%ld).", i, stream->duration); - stream->duration = av->format_context->duration < 0 ? 0 : - av_rescale_q(av->format_context->duration, AV_TIME_BASE_Q, stream->time_base); - al_log_info("ff_demuxer", "Defaulting stream #%u to a duration of %.3fs.", + if (stream->duration < 0L) { + stream->duration = av_rescale_q(default_duration, AV_TIME_BASE_Q, stream->time_base); + al_assert(stream->duration >= 0L); + al_log_info("ff_demuxer", "Stream #%u has an invalid duration, defaulting to %.3fs.", i, stream->duration * av_q2d(stream->time_base)); - al_assert(stream->duration >= 0); } duration = (u64)av_rescale_q(stream->duration, stream->time_base, AV_TIME_BASE_Q); if (duration > av->duration) av->duration = duration; } - if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0; + if (stream->start_time == AV_NOPTS_VALUE) stream->start_time = 0L; } else if (type != AVMEDIA_TYPE_ATTACHMENT) { continue; } diff --git a/src/codec/ffmpeg/encoder.c b/src/codec/ffmpeg/encoder.c index 11aaaff..6e914c8 100644 --- a/src/codec/ffmpeg/encoder.c +++ b/src/codec/ffmpeg/encoder.c @@ -27,7 +27,7 @@ bool camu_ff_encoder_init(struct camu_ff_encoder *enc, char *encoder_name, char const AVCodec *codec = avcodec_find_encoder_by_name(encoder_name); if (!codec) { - al_log_error("ff_encoder", "Failed to find codec with name: %s.", encoder_name); + al_log_error("ff_encoder", "Failed to find codec by name %s.", encoder_name); goto err; } @@ -95,17 +95,21 @@ void camu_ff_encoder_push(struct camu_ff_encoder *enc, u8 **data, s32 sample_cou enc->frame->data[0] = data[0]; s32 ret = avcodec_send_frame(enc->codec_context, enc->frame); if (ret < 0 && ret != AVERROR_EOF) { - al_log_error("ff_encoder", "Error sending frame to the encoder: (%s).", av_err2str(ret)); + al_log_error("ff_encoder", "Error sending frame to the encoder (%s).", av_err2str(ret)); return; } AVPacket pkt = { 0 }; - ret = avcodec_receive_packet(enc->codec_context, &pkt); - if (ret < 0 && ret != AVERROR_EOF) { - al_log_error("ff_encoder", "Error receiving packet from the encoder: (%s).", av_err2str(ret)); - return; - } - enc->callback(enc->userdata, &pkt); - enc->frame->pts += SAMPLES_PER_FRAME; + do { + ret = avcodec_receive_packet(enc->codec_context, &pkt); + if (ret < 0) { + if (!(ret = AVERROR_EOF || ret == AVERROR(EAGAIN))) { + al_log_error("ff_encoder", "Error receiving packet from the encoder (%s).", av_err2str(ret)); + } + return; + } + enc->callback(enc->userdata, &pkt); + enc->frame->pts += SAMPLES_PER_FRAME; + } while (1); } void camu_ff_encoder_reset(struct camu_ff_encoder *enc) diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c index 724a489..3c62ba8 100644 --- a/src/fruits/cmc/cmc.c +++ b/src/fruits/cmc/cmc.c @@ -1,10 +1,8 @@ #include <aki/common.h> #include <aki/event_loop.h> -#include "../../libclient/client.h" -#include "../../server/common.c" #include "../../portal/src/packet_ext.h" -#include "../../portal/src/post_cache.h" +#include "../../server/common.c" #include "cmc.h" @@ -12,18 +10,8 @@ enum { CLI_OPEN_UI = 0, - CLI_ADD -}; - -struct cmc { - struct aki_event_loop loop; - u8 command; - array(str) args; - struct camu_client client; - struct camu_post_cache cache; - array(struct cmc_search *) searches; - struct cmc_ui ui; - struct aki_poll input_poll; + CLI_ADD, + CLI_SEARCH }; static struct cmc_search *get_search_by_id(struct cmc *c, s32 id) @@ -35,21 +23,20 @@ static struct cmc_search *get_search_by_id(struct cmc *c, s32 id) return NULL; } -static void input_poll_callback(void *userdata, s32 revents) +static void parse_user_state(struct cmc *c, struct aki_packet *packet) { - struct cmc *c = (struct cmc *)userdata; - (void)revents; - struct ncinput input; - do { - if (!cmc_ui_read_input(&c->ui, &input)) break; - if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) { - switch (input.id) { - case 'q': - aki_event_loop_break_one(&c->loop); - break; - } - } - } while (1); + u32 size = aki_packet_read_u32(packet); + for (u32 i = 0; i < size; i++) { + struct cmc_search *search = al_alloc_object(struct cmc_search); + search->id = aki_packet_read_s32(packet); + str s; + aki_packet_read_str(packet, &s); + al_str_clone(&search->module, &s); + aki_packet_read_str(packet, &s); + al_str_clone(&search->query, &s); + al_array_init(search->pages); + al_array_push(c->searches, search); + } } static void client_callback(void *userdata, u8 op, void *opaque) @@ -57,19 +44,28 @@ static void client_callback(void *userdata, u8 op, void *opaque) struct cmc *c = (struct cmc *)userdata; switch (op) { case CAMU_CLIENT_LOGIN: { + struct aki_packet *packet = (struct aki_packet *)opaque; switch (c->command) { case CLI_OPEN_UI: - if (!cmc_ui_init(&c->ui)) { + parse_user_state(c, packet); + if (!cmc_ui_init(&c->ui, &c->loop, c)) { aki_event_loop_break_one(&c->loop); } - aki_poll_init(&c->input_poll, input_poll_callback, c); - aki_poll_set(&c->input_poll, cmc_ui_get_input_fd(&c->ui), AKI_POLL_READ); - aki_poll_start(&c->input_poll, &c->loop); break; case CLI_ADD: { str *arg; al_array_foreach_ptr(c->args, i, arg) { camu_client_add_from_path(&c->client, al_str_c("default"), arg); + al_str_free(arg); + } + camu_client_disconnect(&c->client); + break; + } + case CLI_SEARCH: { + str *arg; + al_array_foreach_ptr(c->args, i, arg) { + camu_client_create_search(&c->client, al_str_c("youtube"), arg); + al_str_free(arg); } camu_client_disconnect(&c->client); break; @@ -136,11 +132,22 @@ static bool parse_cmd(s32 argc, wchar_t **argv) al_array_push(c.args, arg); } } + } else if (al_str_eq(al_str_cr(argv[1]), al_str_c("search"))) { + if (argc < 3) return false; + c.command = CLI_SEARCH; + str query; + al_str_from(&query, ""); + for (s32 i = 2; i < argc; i++) { + al_str_cat(&query, al_str_cr(argv[i])); + if (i != argc - 1) { + al_str_cat(&query, al_str_c(" ")); + } + } + al_array_push(c.args, query); } return true; } - #ifndef _WIN32 s32 main(s32 argc, char *argv[]) #else @@ -152,12 +159,11 @@ s32 wmain(s32 argc, wchar_t **argv) aki_event_loop_init(&c.loop); al_array_init(c.args); + if (!parse_cmd(argc, argv)) return EXIT_FAILURE; camu_post_cache_init(&c.cache); al_array_init(c.searches); - if (!parse_cmd(argc, argv)) return EXIT_FAILURE; - c.client.callback = client_callback; c.client.userdata = &c; if (!camu_client_login(&c.client, &c.loop, CAMU_LOCAL_TYPE, CAMU_LOCAL_ADDR, CAMU_PORT, al_str_c("andrew"))) { diff --git a/src/fruits/cmc/cmc.h b/src/fruits/cmc/cmc.h index a9c8544..353bd87 100644 --- a/src/fruits/cmc/cmc.h +++ b/src/fruits/cmc/cmc.h @@ -4,6 +4,11 @@ #include <al/array.h> #include <al/str.h> +#include "../../libclient/client.h" +#include "../../portal/src/post_cache.h" + +#include "ui/ui.h" + struct cmc_search_page { u32 num; array(str) list; @@ -11,5 +16,17 @@ struct cmc_search_page { struct cmc_search { s32 id; + str module; + str query; array(struct cmc_search_page *) pages; }; + +struct cmc { + struct aki_event_loop loop; + u8 command; + array(str) args; + struct camu_client client; + struct camu_post_cache cache; + array(struct cmc_search *) searches; + struct cmc_ui ui; +}; diff --git a/src/fruits/cmc/meson.build b/src/fruits/cmc/meson.build index 627391e..76efff5 100644 --- a/src/fruits/cmc/meson.build +++ b/src/fruits/cmc/meson.build @@ -8,7 +8,10 @@ endif use_tui = true if use_tui - cmc_src += ['ui/ui.c'] + cmc_src += [ + 'ui/ui.c', + 'ui/pane_search.c', + ] cmc_deps += [dependency('notcurses')] endif diff --git a/src/fruits/cmc/ui/pane_search.c b/src/fruits/cmc/ui/pane_search.c new file mode 100644 index 0000000..cc549e1 --- /dev/null +++ b/src/fruits/cmc/ui/pane_search.c @@ -0,0 +1,39 @@ +#include "../cmc.h" + +#include "ui.h" + +void cmc_sp_init(struct cmc_ui *ui) +{ + if (ui->c->searches.size > 0) { + ui->sp.search = al_array_at(ui->c->searches, 0); + } + cmc_sp_layout(ui, notcurses_stdplane(ui->nc)); +} + +void cmc_sp_layout(struct cmc_ui *ui, struct ncplane *parent) +{ + if (ui->sp.n) ncplane_destroy(ui->sp.n); + struct ncplane_options nopts = { + .rows = ui->term_rows, + .cols = ui->term_cols, + .x = 0, + .y = 0 + }; + ui->sp.n = ncplane_create(parent, &nopts); +} + +bool cmc_sp_handle_input(struct cmc_ui *ui, struct ncinput *input) +{ + (void)ui; + (void)input; + return true; +} + +void cmc_sp_render(struct cmc_ui *ui) +{ + if (ui->sp.search) { + char *c_str = al_str_to_c_str(&ui->sp.search->query); + ncplane_putstr_yx(ui->sp.n, 0, 0, c_str); + al_free(c_str); + } +} diff --git a/src/fruits/cmc/ui/tile.h b/src/fruits/cmc/ui/tile.h new file mode 100644 index 0000000..6f15c7d --- /dev/null +++ b/src/fruits/cmc/ui/tile.h @@ -0,0 +1,5 @@ +#pragma once + +struct cmc_tile { + +}; diff --git a/src/fruits/cmc/ui/ui.c b/src/fruits/cmc/ui/ui.c index a838f96..58d630f 100644 --- a/src/fruits/cmc/ui/ui.c +++ b/src/fruits/cmc/ui/ui.c @@ -2,28 +2,69 @@ #include "ui.h" -bool cmc_ui_init(struct cmc_ui *ui) +static s32 resize_cb(struct ncplane *p) { - al_memset(ui, 0, sizeof(struct cmc_ui)); - if (!(ui->nc = notcurses_init(NULL, stdin))) { - return false; - } - return true; + struct cmc_ui *ui = (struct cmc_ui *)ncplane_userptr(p); + (void)ui; + return 0; } -s32 cmc_ui_get_input_fd(struct cmc_ui *ui) -{ - return notcurses_inputready_fd(ui->nc); -} - -bool cmc_ui_read_input(struct cmc_ui *ui, struct ncinput *input) +static bool read_input(struct cmc_ui *ui, struct ncinput *input) { al_memset(input, 0, sizeof(struct ncinput)); u32 ret = notcurses_get_nblock(ui->nc, input); return !(ret == (u32)-1 || ret == 0); } +static void input_poll_callback(void *userdata, s32 revents) +{ + struct cmc_ui *ui = (struct cmc_ui *)userdata; + (void)revents; + struct ncinput input; + do { + if (!read_input(ui, &input)) break; + bool do_render = false; + if (input.evtype == NCTYPE_PRESS || input.evtype == NCTYPE_UNKNOWN) { + switch (input.id) { + case 'q': + aki_event_loop_break_one(ui->loop); + break; + } + } + switch (ui->pane) { + case CMC_PANE_SEARCH: + if (cmc_sp_handle_input(ui, &input)) { + cmc_sp_render(ui); + do_render |= true; + } + } + if (do_render) { + notcurses_render(ui->nc); + } + } while (1); +} + +bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c) +{ + al_memset(ui, 0, sizeof(struct cmc_ui)); + ui->c = c; + if (!(ui->nc = notcurses_init(NULL, stdin))) { + return false; + } + ui->loop = loop; + notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols); + cmc_sp_init(ui); + struct ncplane *stdplane = notcurses_stdplane(ui->nc); + ncplane_set_resizecb(stdplane, resize_cb); + ncplane_set_userptr(stdplane, ui); + aki_poll_init(&ui->input_poll, input_poll_callback, ui); + aki_poll_set(&ui->input_poll, notcurses_inputready_fd(ui->nc), AKI_POLL_READ); + aki_poll_start(&ui->input_poll, ui->loop); + ui->pane = CMC_PANE_SEARCH; + return true; +} + void cmc_ui_close(struct cmc_ui *ui) { - notcurses_stop(ui->nc); + if (ui->nc) notcurses_stop(ui->nc); } diff --git a/src/fruits/cmc/ui/ui.h b/src/fruits/cmc/ui/ui.h index 4af222e..1a49985 100644 --- a/src/fruits/cmc/ui/ui.h +++ b/src/fruits/cmc/ui/ui.h @@ -1,18 +1,33 @@ #pragma once #include <al/types.h> +#include <aki/event_loop.h> #include <notcurses/notcurses.h> +enum { + CMC_PANE_SEARCH = 0, +}; + +struct cmc; struct cmc_ui { struct notcurses *nc; + struct aki_event_loop *loop; u32 term_cols; u32 term_rows; bool pending_layout; + struct aki_poll input_poll; + u8 pane; + struct { + struct ncplane *n; + struct cmc_search *search; + } sp; // search pane. + struct cmc *c; }; -bool cmc_ui_init(struct cmc_ui *ui); - -s32 cmc_ui_get_input_fd(struct cmc_ui *ui); -bool cmc_ui_read_input(struct cmc_ui *ui, struct ncinput *input); - +bool cmc_ui_init(struct cmc_ui *ui, struct aki_event_loop *loop, struct cmc *c); void cmc_ui_close(struct cmc_ui *ui); + +void cmc_sp_init(struct cmc_ui *ui); +void cmc_sp_layout(struct cmc_ui *ui, struct ncplane *parent); +bool cmc_sp_handle_input(struct cmc_ui *ui, struct ncinput *input); +void cmc_sp_render(struct cmc_ui *ui); diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c index 47f90af..5eb9fee 100644 --- a/src/fruits/cmsrv/cmsrv.c +++ b/src/fruits/cmsrv/cmsrv.c @@ -50,7 +50,11 @@ static u8 server_line_callback(void *userdata, str *line) lia_list_clear(list); } else { struct aki_packet *packet = aki_packet_create(); - aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); + if (al_str_at(line, 0) == ';' || camu_is_url(line, 0)) { + aki_packet_write_u8(packet, CAMU_RESOURCE_SIMPLE_SEARCH); + } else { + aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); + } aki_packet_write_str(packet, line); camu_server_local_add(&s->server, packet); } @@ -98,9 +102,10 @@ static void quit_signal_callback(void *userdata) aki_event_loop_break_one(&s->loop); } -static s32 log_callback(void *userdata, char *message) +static s32 log_callback(void *userdata, char *message, const char *color) { struct cmsrv *s = (struct cmsrv *)userdata; + (void)color; cmsrv_ui_push_message(&s->ui, message); return al_strlen(message); } diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c index 7ae19d5..0c26881 100644 --- a/src/fruits/cmsrv/ui.c +++ b/src/fruits/cmsrv/ui.c @@ -8,6 +8,10 @@ static s32 resize_cb(struct ncplane *p) { struct cmsrv_ui *ui = (struct cmsrv_ui *)ncplane_userptr(p); + ncplane_erase(p); + notcurses_refresh(ui->nc, &ui->term_rows, &ui->term_cols); + notcurses_render(ui->nc); + notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols); ui->pending_layout = true; return 0; } @@ -52,6 +56,7 @@ void cmsrv_ui_push_message(struct cmsrv_ui *ui, char *message) static void layout_log(struct cmsrv_ui *ui, struct ncplane *parent) { + if (ui->log.n) ncplane_destroy(ui->log.n); struct ncplane_options nopts = { .rows = ui->term_rows / LOG_RATIO, .cols = ui->term_cols, @@ -88,6 +93,7 @@ static void render_log(struct cmsrv_ui *ui) static void layout_lists(struct cmsrv_ui *ui, struct ncplane *parent) { + if (ui->lists.n) ncplane_destroy(ui->lists.n); struct ncplane_options nopts = { .rows = ui->term_rows - (ui->term_rows / LOG_RATIO), .cols = ui->term_cols, @@ -147,6 +153,7 @@ static void render_lists(struct cmsrv_ui *ui) static void layout_nodes(struct cmsrv_ui *ui, struct ncplane *parent) { + if (ui->nodes.n) ncplane_destroy(ui->nodes.n); struct ncplane_options nopts = { .rows = ui->term_rows - (ui->term_rows / LOG_RATIO), .cols = ui->term_cols, @@ -191,10 +198,6 @@ void cmsrv_ui_render(struct cmsrv_ui *ui) erase_lists(ui); erase_nodes(ui); if (ui->pending_layout) { - ncplane_erase(stdplane); - notcurses_refresh(ui->nc, &ui->term_rows, &ui->term_cols); - notcurses_render(ui->nc); - notcurses_stddim_yx(ui->nc, &ui->term_rows, &ui->term_cols); layout_log(ui, stdplane); layout_lists(ui, stdplane); layout_nodes(ui, stdplane); diff --git a/src/liana/client.c b/src/liana/client.c index a841f7a..ad365b3 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -62,7 +62,7 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack client->mask |= 1 << index; break; case CAMU_STREAM_ATTACHMENT: - // Assume we have all the data we need in the AVStream object. + // Assume all the data we need is in the AVStream object. break; default: continue; diff --git a/src/liana/list.c b/src/liana/list.c index 23eabc3..073cdf0 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -77,11 +77,26 @@ static bool assume_ended(struct lia_list_entry *entry, u64 at) return false; } -static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry) +static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error) { - bool loaded; - list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &loaded); - if (loaded) { + u8 status; + list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &status); + if (status == LIANA_ENTRY_ERRORED) { + if (sequence >= 0) { + al_array_remove_at(list->entries, (u32)sequence); + if (sequence == list->current && (u32)list->current == list->entries.size) { + // This maps to the behavior of only skipping ahead on errors. + list->current--; + list->idle = true; + } + } + *error = true; + u8 meta = LIANA_META_FAILED; + list->callback(list->userdata, LIANA_LIST_META, entry, &meta); + return false; + } + *error = false; + if (status == LIANA_ENTRY_LOADED) { list->callback(list->userdata, LIANA_GET_DURATION, entry, &entry->duration); return true; } @@ -90,12 +105,13 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) { - if (list->current >= 0) { - struct lia_list_entry *current = al_array_at(list->entries, list->current); - if (!entry_load_and_get_duration(list, current)) { + struct lia_list_entry *current; + if (list->current >= 0 && !list->idle && !(current = al_array_at(list->entries, list->current))->ended) { + bool error; + if (!entry_load_and_get_duration(list, current, list->current, &error)) { + if (error) pump_queue(list); return false; } - sink->set = list->current; u64 now = aki_get_timestamp(); u8 pause; u64 at = LIANA_TIMESTAMP_INVALID; @@ -121,6 +137,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) .pause = pause, .ended = ended }; + sink->set = list->current; sink->callback(sink->userdata, LIANA_SINK_SET, current, list->current, &time); } else { sink->set = -1; @@ -145,8 +162,9 @@ static void handle_remove_sink(struct lia_list *list, void *userdata) static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) { if (list->idle) { - if (!entry_load_and_get_duration(list, entry)) { - return false; + bool error; + if (!entry_load_and_get_duration(list, entry, -1, &error)) { + return error; } list->current++; list->idle = false; @@ -163,7 +181,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) sink->set = list->current; sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time); } - al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&entry->name)); + u8 meta = LIANA_META_PLAYING; + list->callback(list->userdata, LIANA_LIST_META, entry, &meta); } else { /* if (list->queued == -1) { @@ -185,7 +204,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) */ entry->start = LIANA_TIMESTAMP_INVALID; //} - // meta queued + u8 meta = LIANA_META_QUEUED; + list->callback(list->userdata, LIANA_LIST_META, entry, &meta); } al_array_push(list->entries, entry); return true; @@ -202,7 +222,7 @@ static void adjust_current(struct lia_list *list, struct lia_list_entry *previou if (i == (u32)list->current) return; cmd->op = SKIPTO; cmd->sequence = i; - cmd->i = list->current; + cmd->value.i = list->current; list->current = i; break; } @@ -241,6 +261,18 @@ static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 return al_array_at(list->entries, sequence); } +static struct lia_list_entry *get_entry_from_id(struct lia_list *list, u32 id, s32 *sequence) +{ + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + if (entry->id == id) { + *sequence = (s32)i; + return entry; + } + } + return NULL; +} + // TODO: //if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) { @@ -260,7 +292,10 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) al_assert(current && !current->held); if (!target) return true; al_assert(current != target); - if (!entry_load_and_get_duration(list, target)) { + bool error; + if (!entry_load_and_get_duration(list, target, index, &error)) { + // index might point to a different entry after an error. + if (error) pump_queue(list); return false; } @@ -281,7 +316,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) // An ended entry may never have been paused, but a non-ended entry that wasn't set // cannot be unpaused. Checking assume_ended(target) should be safe here as long // as it can't go from true to false (consideration for seek?). - al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID); + //al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID); } // These are not equivalent to current/target->ended. @@ -342,7 +377,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time); } - al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(&target->name)); + u8 meta = LIANA_META_PLAYING; + list->callback(list->userdata, LIANA_LIST_META, target, &meta); return true; } @@ -353,8 +389,8 @@ static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->i = sequence + n; - return handle_skipto(list, cmd->sequence, cmd->i); + cmd->value.i = sequence + n; + return handle_skipto(list, cmd->sequence, cmd->value.i); } static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) @@ -408,19 +444,19 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) } } -// TODO: Take entry id instead of sequence for seek() and end()? - -static void handle_seek(struct lia_list *list, s32 sequence, f64 percent) +static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent) { - if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; - if (sequence < 0) return; - if (sequence != list->current) { - return; + struct lia_list_entry *entry; + if (sequence == LIANA_SEQUENCE_ANY) { + sequence = list->current; + if (sequence < 0) return; + entry = get_entry_from_sequence(list, sequence); + al_assert(entry); + } else { + entry = get_entry_from_id(list, id, &sequence); + if (!entry) return; } - struct lia_list_entry *entry = get_entry_from_sequence(list, sequence); - al_assert(entry); if (entry->duration == LIANA_TIMESTAMP_INVALID) { - // Until any kind of live resource buffering. return; } entry->ended = false; @@ -428,8 +464,9 @@ static void handle_seek(struct lia_list *list, s32 sequence, f64 percent) u64 now = aki_get_timestamp(); u64 pos = (u64)(entry->duration * percent); u64 at = now + LIANA_BASE_DELAY; + u8 pause = entry->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE; entry->offset = pos; - if (entry->paused_at == LIANA_TIMESTAMP_INVALID) { + if (pause == LIANA_PAUSE_RESUME) { entry->start = at; } @@ -440,31 +477,34 @@ static void handle_seek(struct lia_list *list, s32 sequence, f64 percent) struct lia_timing time = { .at = at, .seek_pos = pos, - .pause = LIANA_PAUSE_NONE, + .pause = pause, .ended = false }; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { + if (sequence == list->current && sink->set != sequence) { + sink->set = sequence; + sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, &time); + } sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, &time); } } -static bool handle_end(struct lia_list *list, s32 sequence) +static bool handle_end(struct lia_list *list, s32 id) { - al_assert(sequence != LIANA_SEQUENCE_ANY); + s32 sequence; + struct lia_list_entry *entry = get_entry_from_id(list, id, &sequence); + if (!entry) return true; - s32 size = (s32)list->entries.size; - struct lia_list_entry *entry = al_array_at(list->entries, sequence); - al_assert(entry); if (entry->ended) { al_log_warn("list", "Got end() from an already ended resource, ignoring."); return true; } entry->ended = true; - // TODO: Calculate duration for live resources? entry->offset = entry->duration; + s32 size = (s32)list->entries.size; if (sequence == list->current) { s32 next = sequence + 1; if (list->queued >= 0) { @@ -476,12 +516,13 @@ static bool handle_end(struct lia_list *list, s32 sequence) sink->queued = -1; } struct lia_list_entry *current = al_array_at(list->entries, list->current); - al_log_info("list", "Now playing: %ls.", AL_WSTR_PRINTF(¤t->name)); + u8 meta = LIANA_META_PLAYING; + list->callback(list->userdata, LIANA_LIST_META, current, &meta); } else if (next < size) { struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->i = next; + cmd->value.i = next; pump_queue(list); return false; } else { @@ -567,13 +608,13 @@ void pump_queue(struct lia_list *list) handle_unset(list); break; case SKIPTO: - if (!handle_skipto(list, cmd->sequence, cmd->i)) { + if (!handle_skipto(list, cmd->sequence, cmd->value.i)) { // Target entry not loaded. return; } break; case SKIP: - if (!handle_skip(list, cmd->sequence, cmd->i)) { + if (!handle_skip(list, cmd->sequence, cmd->value.i)) { // Converted to skipto and entry not loaded. return; } @@ -582,10 +623,10 @@ void pump_queue(struct lia_list *list) handle_toggle_pause(list, cmd->sequence, cmd->f); break; case SEEK: - handle_seek(list, cmd->sequence, cmd->f); + handle_seek(list, cmd->sequence, cmd->value.u, cmd->f); break; case END: - if (!handle_end(list, cmd->sequence)) { + if (!handle_end(list, cmd->value.u)) { // End was converted to a skip. return; } @@ -665,7 +706,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->i = index; + cmd->value.i = index; al_array_push(list->queue, cmd); pump_queue(list); } @@ -675,7 +716,7 @@ void lia_list_skip(struct lia_list *list, s32 sequence, s32 n) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SKIP; cmd->sequence = sequence; - cmd->i = n; + cmd->value.i = n; al_array_push(list->queue, cmd); pump_queue(list); } @@ -690,21 +731,22 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) pump_queue(list); } -void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent) +void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SEEK; cmd->sequence = sequence; + cmd->value.u = id; cmd->f = percent; al_array_push(list->queue, cmd); pump_queue(list); } -void lia_list_end(struct lia_list *list, s32 sequence) +void lia_list_end(struct lia_list *list, u32 id) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = END; - cmd->sequence = sequence; + cmd->value.u = id; al_array_push(list->queue, cmd); pump_queue(list); } diff --git a/src/liana/list.h b/src/liana/list.h index 90f8c02..8d47b72 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -26,7 +26,22 @@ enum { enum { LIANA_LOAD_ENTRY = 0, LIANA_GET_DURATION, - LIANA_UNLOAD_ENTRY + LIANA_UNLOAD_ENTRY, + LIANA_LIST_META +}; + +enum { + LIANA_ENTRY_PREPARING = 0, + LIANA_ENTRY_PREPARED, + LIANA_ENTRY_LOADING, + LIANA_ENTRY_LOADED, + LIANA_ENTRY_ERRORED +}; + +enum { + LIANA_META_PLAYING = 0, + LIANA_META_QUEUED, + LIANA_META_FAILED }; // NOTE: To handle an entry being queued right before a skip, keep a global @@ -75,7 +90,7 @@ struct lia_list_cmd { void *userdata; struct lia_list_entry *entry; s32 sequence; - s32 i; + union { s32 i; u32 u; } value; f64 f; }; @@ -106,8 +121,8 @@ void lia_list_unset(struct lia_list *list); void lia_list_skipto(struct lia_list *list, s32 sequence, s32 i); void lia_list_skip(struct lia_list *list, s32 sequence, s32 n); void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts); -void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent); -void lia_list_end(struct lia_list *list, s32 sequence); +void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent); +void lia_list_end(struct lia_list *list, u32 id); void lia_list_reverse(struct lia_list *list); void lia_list_sort(struct lia_list *list); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 435fcba..c4bf583 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -24,8 +24,8 @@ static void signal_callback(void *userdata) static void reset_metrics(struct lia_vcr *vcr) { - vcr->metric.current_frame = 0; - vcr->metric.last_report_ts = 0; + vcr->metric.current_frame = 0Lu; + vcr->metric.last_report_ts = 0Lu; } void lia_vcr_init(struct lia_vcr *vcr, struct aki_event_loop *loop, struct aki_packet_stream *data) @@ -162,10 +162,15 @@ static void update_metrics(struct lia_vcr *vcr, u32 size) } u64 diff; if ((diff = now - vcr->metric.last_report_ts) > 1000000Lu) { - f32 kbps = (vcr->metric.current_frame / 125.f) / (diff / 1000000.f); - al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps); - vcr->metric.current_frame = 0; vcr->metric.last_report_ts = now; + u64 frame = vcr->metric.current_frame; + vcr->metric.current_frame = 0Lu; + if (diff > 2500000Lu) { + al_log_debug("vcr", "Ignoring %lu bytes in metrics.", frame); + return; + } + f32 kbps = (frame / 125.f) / (diff / 1000000.f); + al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps); } } @@ -263,6 +268,7 @@ void lia_vcr_flush(struct lia_vcr *vcr) } aki_signal_stop(&vcr->signal); al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + reset_metrics(vcr); if (vcr->expand == VCR_EXPAND_COMPLETE) { vcr->mark.low = 0; vcr->expand = VCR_EXPAND_GROWN; diff --git a/src/libclient/client.c b/src/libclient/client.c index c43a7dd..4b9f46a 100644 --- a/src/libclient/client.c +++ b/src/libclient/client.c @@ -33,8 +33,7 @@ static struct aki_rpc_command commands[] = { static void idd_callback(void *userdata, struct aki_packet *packet) { struct camu_client *client = (struct camu_client *)userdata; - (void)packet; - client->callback(client->userdata, CAMU_CLIENT_LOGIN, NULL); + client->callback(client->userdata, CAMU_CLIENT_LOGIN, packet); } static void connection_callback(void *userdata, struct aki_rpc_connection *conn) diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 26c144b..bd33f46 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -270,7 +270,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_SKIP); aki_packet_write_s32(packet, get_sequence_for_command(sink)); - aki_packet_write_s32(packet, cmd->value.i); + aki_packet_write_s32(packet, (s32)cmd->value.i); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } @@ -301,7 +301,14 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_SEEK); - aki_packet_write_s32(packet, get_sequence_for_command(sink)); + struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque; + if (entry) { + aki_packet_write_s32(packet, entry->sequence); + aki_packet_write_u32(packet, entry->id); + } else { + aki_packet_write_s32(packet, LIANA_SEQUENCE_ANY); + aki_packet_write_u32(packet, 0); + } aki_packet_write_f64(packet, cmd->value.f); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -324,7 +331,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_END); - aki_packet_write_s32(packet, cmd->value.i); + aki_packet_write_u32(packet, (u32)cmd->value.u); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; } @@ -404,11 +411,12 @@ static void maybe_remove_previous(struct camu_sink *sink) } // Due to the looseness of the previous queue, we may have to explicitly remove an entry -// at a point if it becomes incorrect to attempt removing it's buffers. +// if it becomes incorrect to attempt removing it's buffers. // An obvious example of this is at the point an entry gets freed. See LIANA_CLIENT_REMOVE_BUFFERS. static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_entry *entry) { struct camu_sink_entry *rentry; + // al_array_remove_all() al_array_foreach_rev(sink->previous, i, rentry) { if (rentry == entry) { al_array_remove_at(sink->previous, i); @@ -416,7 +424,9 @@ static void maybe_remove_from_previous(struct camu_sink *sink, struct camu_sink_ } } -static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, struct camu_sink_entry *target) +// Every call to maybe_add_to_previous() must map to a remove_entry_buffers(). +static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *previous, + struct camu_sink_entry *target) { al_assert(previous != target && !previous->ended); @@ -429,10 +439,9 @@ static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry } } - // Every call to maybe_add_to_previous() should map to a remove_entry_buffers(). - // We can take a shortcut here because pushing an entry to previous is - // pointless if it's not currently added. - if (!AUDIO_ADDED_OR_EMPTY(previous) && !VIDEO_ADDED_OR_EMPTY(previous)) { + // If neither the entries audio or video buffer is ADDED, we don't care about adding it + // to previous (waiting for the next added entry to remove it). + if (!(previous->audio.state == BUFFER_ADDED || previous->video.state == BUFFER_ADDED)) { remove_entry_buffers(sink, previous); return; } @@ -490,7 +499,8 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { - al_assert(!entry->ended && entry->audio.state != BUFFER_ADDED); + al_assert(!entry->ended); + al_assert(entry->audio.state != BUFFER_ADDED && entry->audio.state != BUFFER_QUEUED); if (entry->audio.state == BUFFER_ENDED) { al_log_warn("sink", "Tried to add an ended audio buffer."); return; @@ -538,7 +548,7 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) #ifndef CAMU_SINK_NO_VIDEO void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { - al_assert(entry->video.state != BUFFER_ADDED); + al_assert(entry->video.state != BUFFER_ADDED && entry->video.state != BUFFER_QUEUED); if (entry->video.state == BUFFER_ENDED) { al_log_warn("sink", "Tried to add an ended video buffer."); return; @@ -589,6 +599,10 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target) #endif al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current)); } + if (sink->reconnecting) { + al_assert(sink->reconnecting == sink->current); + sink->reconnecting = NULL; + } } if (!target->ended) { @@ -621,6 +635,10 @@ static bool end_entry_and_advance_queue(struct camu_sink *sink, struct camu_sink { al_log_info("sink", "Entry ended."); entry->ended = true; + queue_cmd(sink, (struct camu_sink_cmd){ + .op = END, + .value.u = entry->id + }); maybe_remove_from_previous(sink, entry); if (sink->target) { switch_to(sink, sink->target); @@ -666,8 +684,7 @@ static void audio_buffer_callback(void *userdata, u8 op) if (entry->audio.state == BUFFER_ADDED) { remove_entry_audio_buffer(sink, entry); } - // This assert likely doesn't matter, but should be kept - // if it doesn't unnecessarialy trip. + // This assert likely doesn't matter due to the handling of the ENDED state. al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED); entry->audio.state = BUFFER_ENDED; bool run_queue = VIDEO_ENDED_OR_EMPTY(entry); @@ -678,12 +695,6 @@ static void audio_buffer_callback(void *userdata, u8 op) end_entry_and_advance_queue(sink, entry); } aki_mutex_unlock(&sink->mutex); - if (run_queue) { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = END, - .value.i = entry->sequence - }); - } break; } } @@ -710,7 +721,6 @@ static void video_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->video.track); bool swapped = false; - bool run_queue = false; aki_mutex_lock(&sink->mutex); al_log_info("sink", "Video EOF."); if (!ENTRY_IS_SINGLE_FRAME(entry)) { @@ -721,7 +731,6 @@ static void video_buffer_callback(void *userdata, u8 op) entry->video.state = BUFFER_ENDED; if (AUDIO_ENDED_OR_EMPTY(entry)) { swapped = end_entry_and_advance_queue(sink, entry); - run_queue = true; } } aki_mutex_unlock(&sink->mutex); @@ -731,12 +740,6 @@ static void video_buffer_callback(void *userdata, u8 op) .value.i = CAMU_SINK_VIDEO }); } - if (run_queue) { - queue_cmd(sink, (struct camu_sink_cmd){ - .op = END, - .value.i = entry->sequence - }); - } break; } } @@ -905,15 +908,22 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str bool reconnect = *(bool *)opaque; aki_mutex_lock(&sink->mutex); + // We need to call this if entry was added to previous then, // - it's being cleaned up after (ENTRY_MAX_AGE - 1) entries were added but none buffered. // - it was seeked. maybe_remove_from_previous(sink, entry); if (reconnect) { - // This may not be correct if entry->ended can ever be set anywhere - // besides end_entry_and_advance_queue(). entry->ended = false; + if (entry == sink->current) { + if (sink->target) { + switch_to(sink, sink->target); + sink->target = NULL; + } else { + sink->reconnecting = sink->current; + } + } } bool skip_audio = sink->audio.state == SINK_PAUSED; @@ -938,6 +948,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } } #endif + aki_mutex_unlock(&sink->mutex); while ( // Block until buffers are removed. @@ -970,7 +981,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str } case LIANA_CLIENT_RECONNECTED: { aki_mutex_lock(&sink->mutex); - if (entry == sink->current) { + if (entry == sink->reconnecting) { if (entry->audio.state > BUFFER_QUEUED) { add_audio_if_set_and_buffered(entry); } @@ -1250,14 +1261,11 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co u64 at = aki_packet_read_u64(packet); u8 pause = aki_packet_read_u8(packet); - aki_mutex_lock(&sink->mutex); struct camu_sink_entry *entry = get_entry_from_id(sink, id); - if (!entry) { - aki_mutex_unlock(&sink->mutex); - goto out; - } + if (!entry) goto out; al_assert(entry->sequence == sequence); + aki_mutex_lock(&sink->mutex); #ifdef CAMU_SINK_LOCAL (void)at; (void)pause; @@ -1279,9 +1287,9 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co break; } #endif -out: aki_mutex_unlock(&sink->mutex); +out: aki_packet_free(packet); return false; } @@ -1295,21 +1303,12 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con u32 id = aki_packet_read_u32(packet); s32 sequence = aki_packet_read_s32(packet); + (void)sequence; u64 at = aki_packet_read_u64(packet); u64 pos = aki_packet_read_u64(packet); - aki_mutex_lock(&sink->mutex); struct camu_sink_entry *entry = get_entry_from_id(sink, id); - if (!entry) { - aki_mutex_unlock(&sink->mutex); - goto out; - } - if (entry == sink->current && sink->target) { - switch_to(sink, sink->target); - sink->target = NULL; - } - al_assert(entry->sequence == sequence); - aki_mutex_unlock(&sink->mutex); + if (!entry) goto out; #ifdef CAMU_SINK_LOCAL at = 0; @@ -1429,13 +1428,11 @@ void camu_sink_seek(struct camu_sink *sink, f64 precent) aki_mutex_lock(&sink->mutex); 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 = precent, - .opaque = current - }); - } + queue_cmd(sink, (struct camu_sink_cmd){ + .op = SEEK, + .value.f = precent, + .opaque = current + }); } void camu_sink_reseek(struct camu_sink *sink) diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 674b684..48b2b87 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -64,7 +64,7 @@ struct camu_sink_entry { struct camu_sink_cmd { u8 op; - union { s64 i; f64 f; } value; + union { s64 i; u64 u; f64 f; } value; void *opaque; }; @@ -84,6 +84,7 @@ struct camu_sink { struct camu_sink_entry *current; struct camu_sink_entry *queued; struct camu_sink_entry *target; + struct camu_sink_entry *reconnecting; array(struct camu_sink_entry *) previous; array(struct camu_sink_entry *) entries; u16 lru; diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c index 5590563..086498e 100644 --- a/src/mixer/mixer.c +++ b/src/mixer/mixer.c @@ -136,9 +136,9 @@ static void add_buffer_internal(struct camu_mixer *mixer, struct camu_audio_buff static void remove_buffer_internal(struct camu_mixer *mixer, struct camu_audio_buffer *buf) { - struct camu_audio_buffer *cbuf; - al_array_foreach(mixer->buffers, i, cbuf) { - if (cbuf == buf) { + struct camu_audio_buffer *rbuf; + al_array_foreach(mixer->buffers, i, rbuf) { + if (rbuf == buf) { #ifdef CAMU_MIXER_THREADED al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); #endif @@ -244,6 +244,14 @@ void camu_mixer_run_queue(struct camu_mixer *mixer) run_queue_internal(mixer); aki_mutex_unlock(&mixer->mutex); } + +void camu_mixer_clear(struct camu_mixer *mixer) +{ + struct camu_audio_buffer *buf; + al_array_foreach(mixer->buffers, i, buf) { + al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); + } +} #endif void camu_mixer_pause(struct camu_mixer *mixer) @@ -279,16 +287,6 @@ void camu_mixer_resume(struct camu_mixer *mixer) #endif } -void camu_mixer_clear(struct camu_mixer *mixer) -{ -#ifdef CAMU_MIXER_THREADED - struct camu_audio_buffer *buf; - al_array_foreach(mixer->buffers, i, buf) { - al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); - } -#endif -} - void camu_mixer_close(struct camu_mixer *mixer) { al_array_free(mixer->buffers); diff --git a/src/mixer/mixer.h b/src/mixer/mixer.h index 356e2b1..a925e0c 100644 --- a/src/mixer/mixer.h +++ b/src/mixer/mixer.h @@ -39,8 +39,8 @@ void camu_mixer_add_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *b void camu_mixer_remove_buffer(struct camu_mixer *mixer, struct camu_audio_buffer *buf); #ifdef CAMU_MIXER_THREADED void camu_mixer_run_queue(struct camu_mixer *mixer); +void camu_mixer_clear(struct camu_mixer *mixer); #endif void camu_mixer_pause(struct camu_mixer *mixer); void camu_mixer_resume(struct camu_mixer *mixer); -void camu_mixer_clear(struct camu_mixer *mixer); void camu_mixer_close(struct camu_mixer *mixer); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index 2b81c97..e8f30c0 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -8,14 +8,17 @@ bool camu_python_init(void) { - PyPreConfig pre; - PyPreConfig_InitPythonConfig(&pre); - pre.utf8_mode = 1; - pre.dev_mode = 0; - Py_PreInitialize(&pre); + PyPreConfig config; + PyPreConfig_InitPythonConfig(&config); + config.dev_mode = 0; + config.utf8_mode = 1; + PyStatus status = Py_PreInitialize(&config); + if (PyStatus_Exception(status)) { + al_log_error("portal", "Preinitialization failed."); + return false; + } if (PyImport_AppendInittab("portal", PyInit_portal) == -1) { al_log_error("portal", "Could not extend in-built modules table."); - camu_python_close(); return false; } Py_Initialize(); @@ -23,10 +26,12 @@ bool camu_python_init(void) if (!module) { PyErr_Print(); al_log_error("portal", "Could not import module."); - camu_python_close(); - return false; + goto err; } return true; +err: + Py_Finalize(); + return false; } void camu_python_close(void) @@ -54,7 +59,10 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) if (bridge->quit) break; if (!have_python) { // Defer python init. - have_python = camu_python_init(); + if (!(have_python = camu_python_init())) { + al_log_error("portal", "Failed to initialize python."); + break; + } } struct camu_portal_cmd *cmd; al_array_foreach_ptr(bridge->queue, i, cmd) { @@ -107,14 +115,12 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) } } camu_queue_push(bridge->results, result); - aki_signal_send(&bridge->results_signal); - al_array_remove_at_iter(bridge->queue, i); } + bridge->queue.size = 0; + aki_signal_send(&bridge->results_signal); } while (1); aki_mutex_unlock(&bridge->mutex); - if (have_python) { - camu_python_close(); - } + if (have_python) camu_python_close(); return 0; } @@ -126,12 +132,12 @@ static void results_signal_callback(void *userdata) do { camu_queue_try_pop(bridge->results, size, result); if (size == 0) break; - result.callback(result.userdata, &result); + result.callback(bridge->userdata, result.userdata, &result); } while (1); } void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, - struct aki_event_loop *loop) + struct aki_event_loop *loop, void *userdata) { al_array_init(bridge->searches); bridge->cache = cache; @@ -142,11 +148,12 @@ void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache camu_queue_init(bridge->results); aki_signal_init(&bridge->results_signal, loop, results_signal_callback, bridge); aki_signal_start(&bridge->results_signal); + bridge->userdata = userdata; aki_thread_create(&bridge->thread, queue_thread, bridge); } void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, - void (*callback)(void *, struct camu_portal_result *), void *userdata) + void (*callback)(void *, void *, struct camu_portal_result *), void *userdata) { struct camu_portal_cmd cmd; cmd.op = CAMU_CLIENT_CREATE_SEARCH; @@ -163,7 +170,7 @@ void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, s } void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, - void (*callback)(void *, struct camu_portal_result *), void *userdata) + void (*callback)(void *, void *, struct camu_portal_result *), void *userdata) { struct camu_portal_cmd cmd; cmd.op = CAMU_CLIENT_GET_PAGE; diff --git a/src/portal/src/search.h b/src/portal/src/search.h index 85f4f6d..199f491 100644 --- a/src/portal/src/search.h +++ b/src/portal/src/search.h @@ -27,7 +27,7 @@ struct camu_portal_result { u8 op; s32 id; struct camu_result_page *page; - void (*callback)(void *, struct camu_portal_result *); + void (*callback)(void *, void *, struct camu_portal_result *); void *userdata; }; @@ -37,7 +37,7 @@ struct camu_portal_cmd { str query; s32 id; u32 num; - void (*callback)(void *, struct camu_portal_result *); + void (*callback)(void *, void *, struct camu_portal_result *); void *userdata; }; @@ -51,18 +51,19 @@ struct camu_portal_bridge { array(struct camu_portal_cmd) queue; queue(struct camu_portal_result) results; struct aki_signal results_signal; + void *userdata; }; bool camu_python_init(void); void camu_python_close(void); void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, - struct aki_event_loop *loop); + struct aki_event_loop *loop, void *userdata); void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, - void (*callback)(void *, struct camu_portal_result *), void *userdata); + void (*callback)(void *, void *, struct camu_portal_result *), void *userdata); void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, - void (*callback)(void *, struct camu_portal_result *), void *userdata); + void (*callback)(void *, void *, struct camu_portal_result *), void *userdata); void camu_portal_close(struct camu_portal_bridge *bridge); diff --git a/src/screen/screen.c b/src/screen/screen.c index d5f6c78..5b5d175 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -86,6 +86,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) if (scr->last_click_ts && aki_get_timestamp() - scr->last_click_ts <= 300000) { if (scr->flags & CAMU_SCREEN_MODIFIER) { f64 percent = scr->last_mouse_x / scr->width; + percent = CLAMP(percent, 0.0, 100.0); scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent); } else { if (scr->last_mouse_x >= scr->width / 2.f) { @@ -112,6 +113,7 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) switch (state) { case STELA_BUTTON_RELEASED: { f64 percent = scr->last_mouse_x / scr->width; + percent = CLAMP(percent, 0.0, 100.0); scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent); break; } @@ -472,6 +474,17 @@ void camu_screen_run_queue(struct camu_screen *scr) al_atomic_store(u8)(&scr->queued, 0, AL_ATOMIC_RELAXED); aki_mutex_unlock(&scr->mutex); } + +void camu_screen_clear(struct camu_screen *scr) +{ + aki_mutex_lock(&scr->mutex); + struct camu_screen_video *video; + al_array_foreach_ptr(scr->videos, i, video) { + struct camu_video_buffer *buf = video->buf; + al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); + } + aki_mutex_unlock(&scr->mutex); +} #endif void camu_screen_set_state(struct camu_screen *scr, s32 state) @@ -516,19 +529,6 @@ void camu_screen_wake(struct camu_screen *scr) #endif } -void camu_screen_clear(struct camu_screen *scr) -{ -#ifdef CAMU_SCREEN_THREADED - aki_mutex_lock(&scr->mutex); - struct camu_screen_video *video; - al_array_foreach_ptr(scr->videos, i, video) { - struct camu_video_buffer *buf = video->buf; - al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED); - } - aki_mutex_unlock(&scr->mutex); -#endif -} - void camu_screen_close(struct camu_screen *scr) { #ifdef CAMU_SCREEN_THREADED diff --git a/src/screen/screen.h b/src/screen/screen.h index 36d9075..349d41f 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -79,10 +79,10 @@ void camu_screen_add_buffer(struct camu_screen *scr, struct camu_video_buffer *b void camu_screen_remove_buffer(struct camu_screen *scr, struct camu_video_buffer *buf); #ifdef CAMU_SCREEN_THREADED void camu_screen_run_queue(struct camu_screen *scr); +void camu_screen_clear(struct camu_screen *scr); #endif void camu_screen_set_state(struct camu_screen *scr, s32 state); bool camu_screen_poll(struct camu_screen *scr, bool block); bool camu_screen_tick(struct camu_screen *scr); void camu_screen_wake(struct camu_screen *scr); -void camu_screen_clear(struct camu_screen *scr); void camu_screen_close(struct camu_screen *scr); diff --git a/src/server/common.h b/src/server/common.h index b939974..b6ee7a1 100644 --- a/src/server/common.h +++ b/src/server/common.h @@ -47,8 +47,9 @@ enum { enum { CAMU_RESOURCE_FILE = 0, + CAMU_RESOURCE_CDIO, CAMU_RESOURCE_PORTAL, - CAMU_RESOURCE_CDIO + CAMU_RESOURCE_SIMPLE_SEARCH }; AL_UNUSED_FUNCTION_PUSH diff --git a/src/server/db.c b/src/server/db.c index c9608bb..6f76864 100644 --- a/src/server/db.c +++ b/src/server/db.c @@ -15,7 +15,7 @@ static bool open_user(struct camu_server *server, struct aki_dir_entry *dir) json_error_t error; json_t *root = json_loadb(s.data, s.len, 0, &error); if (!root) { - al_log_error("server", "Failed to parse json: %.*s:%d:%d (%s).", + al_log_error("server", "Failed to parse %.*s:%d:%d (%s).", AL_STR_PRINTF(&dir->path), error.line, error.column, error.text); return false; } diff --git a/src/server/resource.h b/src/server/resource.h index 6056e8e..559e2b4 100644 --- a/src/server/resource.h +++ b/src/server/resource.h @@ -3,12 +3,6 @@ #include "../cache/entry.h" #include "../liana/server.h" -enum { - CAMU_RESOURCE_NOT_LOADED = 0, - CAMU_RESOURCE_LOADING, - CAMU_RESOURCE_LOADED -}; - struct camu_resource { u8 type; u8 load; diff --git a/src/server/server.c b/src/server/server.c index 109be14..4ec61da 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -52,13 +52,24 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s return NULL; } +static void write_user_state(struct camu_server *server, struct camu_user *user, struct aki_packet *packet) +{ + (void)user; + struct camu_search *search; + aki_packet_write_u32(packet, server->bridge.searches.size); + al_array_foreach(server->bridge.searches, i, search) { + aki_packet_write_s32(packet, search->id); + aki_packet_write_str(packet, &search->module); + aki_packet_write_str(packet, &search->query); + } +} + static void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable); static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; - (void)rpacket; u8 op = aki_packet_read_u8(packet); switch (op) { @@ -83,6 +94,7 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, client->user = user; al_array_push(server->clients, client); al_log_info("server", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->name)); + write_user_state(server, user, rpacket); break; } case CAMU_SINK: { @@ -103,33 +115,15 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn, return true; } -static void client_portal_callback(void *userdata, struct camu_portal_result *result) +static bool client_still_connected(struct camu_server *server, struct aki_rpc_connection *conn) { - struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata; - struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS); - aki_packet_write_u8(packet, result->op); - aki_packet_write_s32(packet, result->id); - switch (result->op) { - case CAMU_CLIENT_CREATE_SEARCH: { - break; - } - case CAMU_CLIENT_GET_PAGE: { - struct camu_result_page *page = result->page; - aki_packet_write_u32(packet, page->num); - aki_packet_write_u32(packet, page->posts.size); - struct camu_post *post; - al_array_foreach_ptr(page->posts, i, post) { - aki_packet_write_post(packet, post); - } - aki_packet_write_u32(packet, page->list.size); - str *unique_id; - al_array_foreach_ptr(page->list, i, unique_id) { - aki_packet_write_str(packet, unique_id); + struct camu_server_client *client; + al_array_foreach(server->clients, i, client) { + if (client->conn == conn) { + return true; } - break; } - } - aki_rpc_connection_command(conn, packet, NULL, NULL); + return false; } static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) @@ -187,11 +181,21 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve { struct lia_list *list = get_list_from_name(server, name); if (!list) return; - if (enable) { - lia_list_add_sink(list, list_sink_callback, sink); - } else { - lia_list_remove_sink(list, sink); + if (enable) lia_list_add_sink(list, list_sink_callback, sink); + else lia_list_remove_sink(list, sink); +} + +static void process_pending(struct camu_resource *resource) +{ + array(struct lia_list_entry *) pending; + // resource->pending may be edited in a list_pump() call. + al_array_clone(pending, resource->pending); + resource->pending.size = 0; + struct lia_list_entry *entry; + al_array_foreach(pending, i, entry) { + lia_list_pump(entry->list); } + al_array_free(pending); } static void node_callback(void *userdata, u8 op, u64 duration) @@ -199,27 +203,14 @@ static void node_callback(void *userdata, u8 op, u64 duration) struct camu_resource *resource = (struct camu_resource *)userdata; switch (op) { case LIANA_NODE_DURATION: - resource->load = CAMU_RESOURCE_LOADED; + resource->load = LIANA_ENTRY_LOADED; resource->duration = duration; - struct lia_list_entry *entry; - al_array_foreach(resource->pending, i, entry) { - lia_list_pump(entry->list); - } - resource->pending.size = 0; + process_pending(resource); break; } } -static void maybe_add_to_pending(struct camu_resource *resource, struct lia_list_entry *entry) -{ - struct lia_list_entry *rentry; - al_array_foreach(resource->pending, i, rentry) { - if (rentry == entry) return; - } - al_array_push(resource->pending, entry); -} - -static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *result) +static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *opaque) { struct camu_server *server = (struct camu_server *)userdata; (void)server; @@ -227,26 +218,64 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v switch (op) { case LIANA_LOAD_ENTRY: switch (resource->load) { - case CAMU_RESOURCE_NOT_LOADED: - resource->load = CAMU_RESOURCE_LOADING; + case LIANA_ENTRY_PREPARED: + resource->load = LIANA_ENTRY_LOADING; lia_node_get_duration(resource->node); // fallthrough - case CAMU_RESOURCE_LOADING: - maybe_add_to_pending(resource, entry); - *(bool *)result = false; + case LIANA_ENTRY_PREPARING: + case LIANA_ENTRY_LOADING: + al_array_push(resource->pending, entry); break; - case CAMU_RESOURCE_LOADED: - *(bool *)result = true; + case LIANA_ENTRY_LOADED: + case LIANA_ENTRY_ERRORED: break; } + *(u8 *)opaque = resource->load; break; case LIANA_GET_DURATION: { - *(u64 *)result = resource->duration; + *(u64 *)opaque = resource->duration; break; } case LIANA_UNLOAD_ENTRY: break; + case LIANA_LIST_META: + if (server->meta_callback) { + server->meta_callback(server->userdata, *(u8 *)opaque, entry->list, entry); + } + break; + } + +} + +static void client_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result) +{ + struct camu_server *server = (struct camu_server *)userdata0; + struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata1; + if (!client_still_connected(server, conn)) return; + struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS); + aki_packet_write_u8(packet, result->op); + aki_packet_write_s32(packet, result->id); + switch (result->op) { + case CAMU_CLIENT_CREATE_SEARCH: { + break; + } + case CAMU_CLIENT_GET_PAGE: { + struct camu_result_page *page = result->page; + aki_packet_write_u32(packet, page->num); + aki_packet_write_u32(packet, page->posts.size); + struct camu_post *post; + al_array_foreach_ptr(page->posts, i, post) { + aki_packet_write_post(packet, post); + } + aki_packet_write_u32(packet, page->list.size); + str *unique_id; + al_array_foreach_ptr(page->list, i, unique_id) { + aki_packet_write_str(packet, unique_id); + } + break; } + } + aki_rpc_connection_command(conn, packet, NULL, NULL); } static bool client_command_callback(void *userdata, struct aki_rpc_connection *conn, @@ -282,8 +311,8 @@ static bool client_command_callback(void *userdata, struct aki_rpc_connection *c } case CAMU_CLIENT_CREATE_SEARCH: { str module; - aki_packet_read_str(packet, &module); str query; + aki_packet_read_str(packet, &module); aki_packet_read_str(packet, &query); camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn); break; @@ -301,6 +330,55 @@ out: return true; } +static struct cch_entry *entry_from_post(struct camu_server *server, struct camu_post *post, u32 index) +{ + struct cch_entry *entry = NULL; + if (index <= post->media.size) { + struct camu_post_media *media = &al_array_at(post->media, index); + if (!al_str_is_empty(&media->url)) { + entry = cch_handler_http_create(&media->url, server->loop); + } + } + if (!entry) { + al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&post->unique_id), index); + return NULL; + } + entry->handler->maybe_spawn_worker(entry->handler, 0); + return entry; +} + +static void simple_search_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result) +{ + struct camu_server *server = (struct camu_server *)userdata0; + struct camu_resource_portal *portal = (struct camu_resource_portal *)userdata1; + struct camu_resource *resource = (struct camu_resource *)portal; + switch (result->op) { + case CAMU_CLIENT_CREATE_SEARCH: { + camu_portal_get_page(&server->bridge, result->id, 0, simple_search_portal_callback, portal); + return; + } + case CAMU_CLIENT_GET_PAGE: { + if (!result->page || result->page->list.size == 0) break; + struct camu_post *post = camu_post_cache_get(&server->cache, &al_array_at(result->page->list, 0)); + if (!post) break; + struct cch_entry *entry = entry_from_post(server, post, 0); + if (!entry) break; + portal->post = post; + resource->entry = entry; + resource->node = lia_server_create_node(&server->data.server, resource->entry); + resource->node->callback = node_callback; + resource->node->userdata = resource; + resource->type = CAMU_RESOURCE_PORTAL; + resource->load = LIANA_ENTRY_PREPARED; + process_pending(resource); + return; + } + } + // No return is the error case. + al_log_warn("server", "Server resource failed to load."); + resource->load = LIANA_ENTRY_ERRORED; +} + static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet) { u8 op = aki_packet_read_u8(packet); @@ -318,6 +396,21 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list al_wstr_from_str(&name, &path); resource = (struct camu_resource *)file; resource->type = CAMU_RESOURCE_FILE; + resource->load = LIANA_ENTRY_PREPARED; + break; + } + case CAMU_RESOURCE_CDIO: { + u32 track = aki_packet_read_u32(packet); + entry = cch_handler_cdio_create(); + if (!entry) return; + struct cch_chapter *chapter = &al_array_at(entry->chapters, track); + entry->handler->maybe_spawn_worker(entry->handler, chapter->start); + struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio); + cdio->track = track; + al_wstr_from_cstr(&name, "cdio"); + resource = (struct camu_resource *)cdio; + resource->type = CAMU_RESOURCE_CDIO; + resource->load = LIANA_ENTRY_PREPARED; break; } case CAMU_RESOURCE_PORTAL: { @@ -325,44 +418,44 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list aki_packet_read_str(packet, &unique_id); u32 index = aki_packet_read_u32(packet); struct camu_post *post = camu_post_cache_get(&server->cache, &unique_id); - entry = NULL; - if (index <= post->media.size) { - struct camu_post_media *media = &al_array_at(post->media, index); - if (!al_str_is_empty(&media->url)) { - entry = cch_handler_http_create(&media->url, server->loop); - } - } - if (!entry) { - al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&unique_id), index); - return; - } - entry->handler->maybe_spawn_worker(entry->handler, 0); + entry = entry_from_post(server, post, index); + if (!entry) return; struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal); portal->post = post; al_wstr_clone(&name, &post->title); resource = (struct camu_resource *)portal; resource->type = CAMU_RESOURCE_PORTAL; + resource->load = LIANA_ENTRY_PREPARED; break; } - case CAMU_RESOURCE_CDIO: { - u32 track = aki_packet_read_u32(packet); - entry = cch_handler_cdio_create(); - struct cch_chapter *chapter = &al_array_at(entry->chapters, track); - entry->handler->maybe_spawn_worker(entry->handler, chapter->start); - struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio); - cdio->track = track; - al_wstr_from_cstr(&name, "cdio"); - resource = (struct camu_resource *)cdio; - resource->type = CAMU_RESOURCE_CDIO; + case CAMU_RESOURCE_SIMPLE_SEARCH: { + str search; + aki_packet_read_str(packet, &search); + struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal); + portal->post = NULL; + al_wstr_from_str(&name, &search); + resource = (struct camu_resource *)portal; + resource->type = CAMU_RESOURCE_SIMPLE_SEARCH; + resource->load = LIANA_ENTRY_PREPARING; + str query; + if (al_str_at(&search, 0) == ';') { // search. + al_str_clone(&query, al_str_substr(&search, 1, search.len)); + } else { + al_str_from(&query, "link:"); + al_str_cat(&query, &search); + } + camu_portal_create_search(&server->bridge, al_str_c("youtube"), &query, simple_search_portal_callback, portal); + al_str_free(&query); break; } } al_assert(resource); - resource->load = CAMU_RESOURCE_NOT_LOADED; - resource->entry = entry; - resource->node = lia_server_create_node(&server->data.server, resource->entry); - resource->node->callback = node_callback; - resource->node->userdata = resource; + if (entry) { + resource->entry = entry; + resource->node = lia_server_create_node(&server->data.server, resource->entry); + resource->node->callback = node_callback; + resource->node->userdata = resource; + } resource->duration = LIANA_TIMESTAMP_INVALID; al_array_init(resource->pending); lia_list_add(list, resource, resource->duration, &name); @@ -412,8 +505,9 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn } case CAMU_LIST_SEEK: { s32 sequence = aki_packet_read_s32(packet); + u32 id = aki_packet_read_u32(packet); f64 percent = aki_packet_read_f64(packet); - lia_list_seek(list, sequence, percent); + lia_list_seek(list, sequence, id, percent); break; } case CAMU_LIST_UNSET: { @@ -421,8 +515,8 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn break; } case CAMU_LIST_END: { - s32 sequence = aki_packet_read_s32(packet); - lia_list_end(list, sequence); + u32 id = aki_packet_read_u32(packet); + lia_list_end(list, id); break; } } @@ -542,7 +636,7 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop #ifdef CAMU_HAVE_PORTAL camu_post_cache_init(&server->cache); - camu_portal_init(&server->bridge, &server->cache, server->loop); + camu_portal_init(&server->bridge, &server->cache, server->loop, server); #endif return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server); diff --git a/src/server/server.h b/src/server/server.h index aaabc3c..ec0f85b 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -45,6 +45,8 @@ struct camu_server { struct camu_portal_bridge bridge; struct camu_post_cache cache; #endif + void (*meta_callback)(void *, u8, struct lia_list *, struct lia_list_entry *); + void *userdata; }; bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop); |