diff options
| author | 2024-11-08 14:53:40 -0500 | |
|---|---|---|
| committer | 2024-11-08 14:53:40 -0500 | |
| commit | 5e3641e5e692c3f2f644a4bb809c88727cb8bee9 (patch) | |
| tree | fbac2007fdcebda13406350fdcef977910427078 /src | |
| parent | f56abfafcd4fa722b807278b138f805112cd953e (diff) | |
| download | camu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.tar.gz camu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.tar.bz2 camu-5e3641e5e692c3f2f644a4bb809c88727cb8bee9.zip | |
Command queue for portal and list, work on server
Most of the server stuff can undoubtedly be simplified. I'm still
working that out.
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src')
40 files changed, 1242 insertions, 695 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 113ea1d..8d33f79 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -143,7 +143,7 @@ static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame s32 sample_count = frame->nb_samples; AVStream *stream = buf->stream->av.stream; f64 pts = frame->best_effort_timestamp * av_q2d(stream->time_base); - f64 duration = camu_audio_format_samples_to_sec(&buf->fmt.in, frame->nb_samples); + f64 duration = camu_audio_format_samples_to_sec(&buf->fmt.in, sample_count); if (pts + duration >= camu_clock_get_base_pts(buf->clock)) { if (buf->pts == -1.0) buf->pts = pts; u8 **data = frame->data; @@ -163,15 +163,20 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra al_assert(al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED) == FLOWING); switch (frame->mode) { case CAMU_NORMAL: { - u8 *store[AV_NUM_DATA_POINTERS] = { 0 }; - store[0] = frame->data; - u8 **data = store; s32 sample_count = frame->audio.sample_count; - if (buf->fmt.resampler_needed) { - sample_count = buf->resamp->convert(buf->resamp, (const u8 **)store, sample_count); - data = buf->resamp->get_data(buf->resamp); + f64 pts = frame->pts; + f64 duration = camu_audio_format_samples_to_sec(&buf->fmt.in, sample_count); + if (pts + duration >= camu_clock_get_base_pts(buf->clock)) { + if (buf->pts == -1.0) buf->pts = pts; + u8 *store[AV_NUM_DATA_POINTERS] = { 0 }; + store[0] = frame->data; + u8 **data = store; + if (buf->fmt.resampler_needed) { + sample_count = buf->resamp->convert(buf->resamp, (const u8 **)store, sample_count); + data = buf->resamp->get_data(buf->resamp); + } + push_internal(buf, data[0], sample_count); } - push_internal(buf, data[0], sample_count); break; } #ifdef CAMU_HAVE_FFMPEG diff --git a/src/buffer/video.c b/src/buffer/video.c index 0ffc00d..91c3687 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -20,7 +20,7 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl buf->clock = clock; buf->latency = 0.0; buf->pts = -1.0; - // Defaulting single_frame to true should simplify non-configured buffers in sink. + // Defaulting single_frame to true can simplify non-configured buffers in sink. buf->single_frame = true; buf->queue = renderer->create_queue(renderer); buf->buffered = false; @@ -61,6 +61,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_code camu_video_format_copy(&buf->fmt.in, fmt); buf->single_frame = duration == 0 || frame_rate.den == 0; + if (buf->single_frame) frame_rate = (AVRational){ 0, 1 }; buf->avg_frame_duration = buf->single_frame ? 0.0 : av_q2d(av_inv_q(frame_rate)); const char *format_name = av_get_pix_fmt_name(buf->fmt.in.format); diff --git a/src/fruits/cmc/cli.h b/src/fruits/cmc/cli.h new file mode 100644 index 0000000..4373054 --- /dev/null +++ b/src/fruits/cmc/cli.h @@ -0,0 +1,4 @@ +#pragma once + +struct cmc_cli { +}; diff --git a/src/fruits/cmc/cmc.c b/src/fruits/cmc/cmc.c new file mode 100644 index 0000000..ad7c75e --- /dev/null +++ b/src/fruits/cmc/cmc.c @@ -0,0 +1,98 @@ +#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 "cmc.h" + +struct cmc { + struct aki_event_loop loop; + struct camu_client client; + struct camu_post_cache cache; + array(struct cmc_search *) searches; +}; + +static struct cmc_search *get_search_by_id(struct cmc *c, s32 id) +{ + struct cmc_search *search; + al_array_foreach(c->searches, i, search) { + if (search->id == id) return search; + } + return NULL; +} + +static void client_callback(void *userdata, u8 op, void *opaque) +{ + struct cmc *c = (struct cmc *)userdata; + switch (op) { + case CAMU_CLIENT_LOGIN: { + camu_client_create_search(&c->client, al_str_c("youtube"), al_str_c("")); + break; + } + case CAMU_CLIENT_SEARCH_CREATED: { + struct cmc_search *search = al_alloc_object(struct cmc_search); + search->id = *(s32 *)opaque; + al_array_init(search->pages); + al_array_push(c->searches, search); + camu_client_get_page(&c->client, search->id, 0); + break; + } + case CAMU_CLIENT_PAGE_RESULTS: { + struct aki_packet *packet = (struct aki_packet *)opaque; + s32 id = aki_packet_read_s32(packet); + struct cmc_search *search = get_search_by_id(c, id); + if (!search) return; + struct cmc_search_page *page = al_alloc_object(struct cmc_search_page); + page->num = aki_packet_read_u32(packet); + u32 posts = aki_packet_read_u32(packet); + for (u32 i = 0; i < posts; i++) { + struct camu_post post; + aki_packet_read_post(packet, &post); + camu_post_cache_push(&c->cache, &post); + } + u32 ids = aki_packet_read_u32(packet); + for (u32 i = 0; i < ids; i++) { + str unique_id; + aki_packet_read_str(packet, &unique_id); + str s; + al_str_clone(&s, &unique_id); + al_array_push(page->list, s); + } + al_array_push(search->pages, page); + camu_client_add(&c->client, al_str_c("default"), &al_array_at(page->list, 0), 0); + break; + } + } +} + +static struct cmc c = { 0 }; + +#ifndef _WIN32 +s32 main(s32 argc, char *argv[]) +#else +s32 wmain(s32 argc, wchar_t **argv) +#endif +{ + (void)argc; + (void)argv; + if (!aki_common_init()) return EXIT_FAILURE; + + aki_event_loop_init(&c.loop); + + camu_post_cache_init(&c.cache); + + al_array_init(c.searches); + + 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"))) { + return EXIT_FAILURE; + } + + aki_event_loop_run(&c.loop); + + return EXIT_SUCCESS; +} diff --git a/src/fruits/cmc/cmc.h b/src/fruits/cmc/cmc.h new file mode 100644 index 0000000..a9c8544 --- /dev/null +++ b/src/fruits/cmc/cmc.h @@ -0,0 +1,15 @@ +#pragma once + +#include <al/types.h> +#include <al/array.h> +#include <al/str.h> + +struct cmc_search_page { + u32 num; + array(str) list; +}; + +struct cmc_search { + s32 id; + array(struct cmc_search_page *) pages; +}; diff --git a/src/fruits/cmc/cmc.py b/src/fruits/cmc/cmc.py deleted file mode 100644 index 5c0a5b6..0000000 --- a/src/fruits/cmc/cmc.py +++ /dev/null @@ -1,17 +0,0 @@ -import cffi -import locale -from notcurses import notcurses - -ffi = cffi.FFI() - -print(ffi.string(notcurses.lib.notcurses_version())) - -class CMC: - def __init__(self): - pass - -c = CMC() - -locale.setlocale(locale.LC_ALL, "") -nc = notcurses.Notcurses() -nc.render() diff --git a/src/fruits/cmc/meson.build b/src/fruits/cmc/meson.build new file mode 100644 index 0000000..4cbe553 --- /dev/null +++ b/src/fruits/cmc/meson.build @@ -0,0 +1,19 @@ +cmc_src = ['cmc.c'] +cmc_deps = [libclient] +cmc_args = [] + +if get_option('portal').enabled() + cmc_deps += [portal] +endif + +use_tui = true +if use_tui + #cmc_src += ['ui.c'] + cmc_deps += [dependency('notcurses')] +endif + +if is_windows and meson.is_cross_build() + cmc_args += ['-static', '-static-libgcc', '-static-libstdc++', '-municode', '-mwindows'] +endif + +executable('cmc', cmc_src, dependencies: cmc_deps, link_args: cmc_args) diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c index 08f67dd..c6f4546 100644 --- a/src/fruits/cmsrv/cmsrv.c +++ b/src/fruits/cmsrv/cmsrv.c @@ -8,13 +8,9 @@ #include <aki/timer.h> #include "../../server/server.h" -#include "../../server/list.h" #include "../../server/common.h" #include "../../cache/handlers/cdio.h" #include "../../codec/ffmpeg/common.h" -#ifdef CAMU_LOCAL_SOCKET -#include "../../server/local_compat.h" -#endif #include "ui.h" @@ -26,7 +22,6 @@ struct cmsrv { struct { struct aki_socket sock; struct aki_line_processor cli; - struct camu_local_compat compat; } local; #endif struct cmsrv_ui ui; @@ -38,23 +33,24 @@ struct cmsrv { static u8 server_line_callback(void *userdata, str *line) { struct cmsrv *s = (struct cmsrv *)userdata; - struct camu_list *list = al_array_at(s->server.lists, 0); + struct lia_list *list = al_array_at(s->server.lists, 0); if (al_str_eq(line, al_str_c(";PAUSE"))) { - lia_list_toggle_pause(&list->impl, LIANA_SEQUENCE_ANY, -1.0); + lia_list_toggle_pause(list, LIANA_SEQUENCE_ANY, -1.0); } else if (al_str_eq(line, al_str_c(";NEXT"))) { - lia_list_skip(&list->impl, LIANA_SEQUENCE_ANY, 1); + lia_list_skip(list, LIANA_SEQUENCE_ANY, 1); } else if (al_str_eq(line, al_str_c(";PREV"))) { - lia_list_skip(&list->impl, LIANA_SEQUENCE_ANY, -1); + lia_list_skip(list, LIANA_SEQUENCE_ANY, -1); } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) { - lia_list_shuffle(&list->impl); + lia_list_shuffle(list); } else if (al_str_eq(line, al_str_c(";SORT"))) { - lia_list_sort(&list->impl); + lia_list_sort(list); } else if (al_str_eq(line, al_str_c(";REVERSE"))) { - lia_list_reverse(&list->impl); + lia_list_reverse(list); } else if (al_str_eq(line, al_str_c(";CLEAR"))) { - lia_list_clear(&list->impl); + lia_list_clear(list); } else { struct aki_packet *packet = aki_packet_create(); + aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); aki_packet_write_str(packet, line); camu_server_local_add(&s->server, packet); } @@ -62,27 +58,6 @@ static u8 server_line_callback(void *userdata, str *line) } #endif -static void meta_callback(void *userdata, u8 op, struct lia_list_entry *entry) -{ - struct cmsrv *s = (struct cmsrv *)userdata; - (void)s; - switch (op) { - case LIANA_META_PLAYING: { - struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; -#ifdef CAMU_HAVE_PORTAL - if (resource->post) { - al_log_info("server", "now playing: %.*ls.", AL_WSTR_PRINTF(&resource->post->title)); - } else { -#endif - al_log_info("server", "now playing: %.*s.", AL_STR_PRINTF(&resource->unique_id)); -#ifdef CAMU_HAVE_PORTAL - } -#endif - break; - } - } -} - static void render_timer_callback(void *userdata, struct aki_timer *timer) { struct cmsrv *s = (struct cmsrv *)userdata; @@ -161,9 +136,6 @@ s32 wmain(s32 argc, wchar_t **argv) if (!camu_server_init(&s.server, CAMU_LOCAL_TYPE, &s.loop)) return EXIT_FAILURE; camu_server_listen(&s.server, CAMU_LOCAL_ADDR, CAMU_PORT); - struct camu_list *list = al_array_at(s.server.lists, 0); - list->impl.callback = meta_callback; - list->impl.userdata = &s; #ifdef CAMU_LOCAL_SOCKET s.local.sock.type = AKI_SOCKET_UNIX; diff --git a/src/fruits/cmsrv/ui.c b/src/fruits/cmsrv/ui.c index 48f50ca..c5cceef 100644 --- a/src/fruits/cmsrv/ui.c +++ b/src/fruits/cmsrv/ui.c @@ -1,7 +1,5 @@ #include <al/lib.h> -#include "../../server/list.h" - #include "ui.h" #define LOG_RATIO 1.3 @@ -112,22 +110,20 @@ static void render_lists(struct cmsrv_ui *ui) u32 current_line = 0; s32 entries_per_list = 10; - struct camu_list *list; + struct lia_list *list; al_array_foreach(ui->server->lists, i, list) { if (current_line++ >= max_height) break; char *c_str = al_str_to_c_str(&list->name); ncplane_putnstr_yx(p, i, 0, max_width, c_str); al_free(c_str); - struct lia_list *impl = &list->impl; - s32 current = (s32)impl->current; - s32 index = AL_MAX(current - (entries_per_list / 2), 0); - s32 size = (s32)impl->entries.size; + s32 index = AL_MAX(list->current - (entries_per_list / 2), 0); + s32 size = (s32)list->entries.size; s32 end = AL_MIN(index + entries_per_list, size); for (s32 j = index; j < end; j++) { - struct lia_list_entry *entry = al_array_at(impl->entries, j); + struct lia_list_entry *entry = al_array_at(list->entries, j); c_str = al_str_to_c_str(&entry->name); u32 y = i + (j - index) + 1; - if (j == current) { + if (j == list->current) { ncplane_putchar_yx(p, y, 1, '>'); ncplane_putnstr_yx(p, y, 3, max_width - 3, c_str); } else { diff --git a/src/fruits/cmv/cmv.c b/src/fruits/cmv/cmv.c index 8eaa8b7..eb96623 100644 --- a/src/fruits/cmv/cmv.c +++ b/src/fruits/cmv/cmv.c @@ -7,7 +7,6 @@ #include "../../codec/ffmpeg/common.h" #ifndef CAMU_SINK_ONLY #include "../../server/server.h" -#include "../../server/list.h" #else #include "../../server/common.c" #endif @@ -21,19 +20,6 @@ struct cmv { }; #ifndef CAMU_SINK_ONLY -static void meta_callback(void *userdata, u8 op, struct lia_list_entry *entry) -{ - struct cmv *c = (struct cmv *)userdata; - (void)c; - switch (op) { - case LIANA_META_PLAYING: { - struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; - al_log_info("server", "now playing: %.*s.", AL_STR_PRINTF(&resource->unique_id)); - break; - } - } -} - static void exit_callback(void *userdata, struct camu_desktop *desktop) { struct cmv *c = (struct cmv *)userdata; @@ -87,9 +73,6 @@ s32 wmain(s32 argc, wchar_t **argv) type = AKI_SOCKET_UNIX; addr = CAMU_UNIX_LOCAL; if (!camu_server_init(&c.server, type, &c.loop)) failure(); if (!camu_server_listen(&c.server, addr, CAMU_PORT)) failure(); - struct camu_list *list = al_array_at(c.server.lists, 0); - list->impl.callback = meta_callback; - list->impl.userdata = &c; } else { type = CAMU_LOCAL_TYPE; addr = CAMU_LOCAL_ADDR; @@ -108,7 +91,13 @@ s32 wmain(s32 argc, wchar_t **argv) al_wstr_to_str(al_wstr_cr(argv[i]), &arg); #endif struct aki_packet *packet = aki_packet_create(); - aki_packet_write_str(packet, &arg); + if (al_str_at(&arg, 0) == ';') { + aki_packet_write_u8(packet, CAMU_RESOURCE_PORTAL); + aki_packet_write_str(packet, al_str_substr(&arg, 1, arg.len)); + } else { + aki_packet_write_u8(packet, CAMU_RESOURCE_FILE); + aki_packet_write_str(packet, &arg); + } camu_server_local_add(&c.server, packet); } } diff --git a/src/liana/client.c b/src/liana/client.c index 55a925a..527376a 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -60,6 +60,7 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack break; case AVMEDIA_TYPE_VIDEO: client->mask |= 1 << index; + //continue; break; case AVMEDIA_TYPE_SUBTITLE: default: @@ -106,7 +107,7 @@ static void info_packet_callback(void *userdata, struct aki_packet_stream *strea struct lia_client *client = (struct lia_client *)userdata; client->connection_id = aki_packet_read_u16(packet); parse_info_packet(client, packet); - al_assert(client->mask != 0); + //al_assert(client->mask != 0); aki_packet_free(packet); stream->packet_callback = data_packet_callback; struct aki_packet *rpacket = aki_packet_create(); diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index 118f6b5..d271829 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -45,14 +45,20 @@ static void cdio_server_subscribe(struct lia_server_handler *handler, s32 mask) static u64 cdio_server_get_duration(struct lia_server_handler *handler) { - (void)handler; - return 0; + struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; + struct cch_chapter *first = &al_array_at(cdio->handle->entry->chapters, 0); + struct cch_chapter *last = &al_array_last(cdio->handle->entry->chapters); + f64 seconds = camu_audio_format_bytes_to_sec(&cdio->fmt, (last->end - first->start) * CDIO_CD_FRAMESIZE_RAW); + al_log_info("cdio", "Length: %.2fs.", seconds); + return (u64)(seconds * 1000000.0); } static bool cdio_server_seek(struct lia_server_handler *handler, u64 pos) { struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; - return cch_handle_seek(cdio->handle, pos, SEEK_SET); + (void)cdio; + (void)pos; + return false; } static void cdio_server_step(struct lia_server_handler *handler) diff --git a/src/liana/list.c b/src/liana/list.c index 970d247..fc34aca 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -6,16 +6,6 @@ #include "list.h" #include "list_cmp.h" -void lia_list_init(struct lia_list *list) -{ - list->current = -1; - list->previous = -1; - list->queued = -1; - list->idle = true; - al_array_init(list->entries); - al_array_init(list->sinks); -} - /* static void buffer_ahead(struct lia_list *list) { @@ -33,13 +23,33 @@ static void buffer_ahead(struct lia_list *list) } */ -static void unset_all(struct lia_list *list) +enum { + ADD_SINK = 0, + REMOVE_SINK, + ADD, + UNSET, + SKIPTO, + SKIP, + TOGGLE_PAUSE, + SEEK, + END, + REVERSE, + SORT, + SHUFFLE, + CLEAR +}; + +void lia_list_init(struct lia_list *list, str *name) { - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - sink->set = -1; - sink->queued = -1; - } + al_str_clone(&list->name, name); + list->current = -1; + list->previous = -1; + list->queued = -1; + list->idle = true; + al_array_init(list->entries); + al_array_init(list->sinks); + al_array_init(list->queue); + list->cmd = NULL; } static bool assume_ended(struct lia_list_entry *entry, u64 at) @@ -57,15 +67,25 @@ static bool assume_ended(struct lia_list_entry *entry, u64 at) return false; } -void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata) +static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry) +{ + bool loaded; + list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &loaded); + if (loaded) { + list->callback(list->userdata, LIANA_GET_DURATION, entry, &entry->duration); + return true; + } + return false; +} + +static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) { - struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink); - sink->callback = callback; - sink->userdata = userdata; - al_array_push(list->sinks, sink); if (list->current >= 0) { - sink->set = list->current; struct lia_list_entry *current = al_array_at(list->entries, list->current); + if (!entry_load_and_get_duration(list, current)) { + return false; + } + sink->set = list->current; u64 now = aki_get_timestamp(); u8 pause; u64 at = LIANA_TIMESTAMP_INVALID; @@ -92,9 +112,11 @@ void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struc sink->set = -1; } sink->queued = -1; + al_array_push(list->sinks, sink); + return true; } -void lia_list_remove_sink(struct lia_list *list, void *userdata) +static void handle_remove_sink(struct lia_list *list, void *userdata) { struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { @@ -106,18 +128,12 @@ void lia_list_remove_sink(struct lia_list *list, void *userdata) } } -void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) +static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) { - struct lia_list_entry *entry = al_alloc_object(struct lia_list_entry); - entry->opaque = opaque; - entry->paused_at = LIANA_TIMESTAMP_INVALID; - entry->held = false; - entry->offset = 0; - entry->ended = false; - entry->duration = duration; - al_str_clone(&entry->name, name); - al_array_push(list->entries, entry); if (list->idle) { + if (!entry_load_and_get_duration(list, entry)) { + return false; + } list->current++; list->idle = false; entry->start = aki_get_timestamp() + LIANA_BASE_DELAY; @@ -133,10 +149,11 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) sink->set = list->current; sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time); } - if (list->callback) list->callback(list->userdata, LIANA_META_PLAYING, entry); + // meta playing } else { /* if (list->queued == -1) { + // TODODODO: this is based on addeding entry to list->entries BEFORE this point. struct lia_list_entry *current = al_array_at(list->entries, list->current); list->queued = list->current + 1; entry->start = current->start + (current->duration - current->offset); @@ -154,11 +171,22 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) */ entry->start = LIANA_TIMESTAMP_INVALID; //} - if (list->callback) list->callback(list->userdata, LIANA_META_QUEUED, entry); + // meta queued } + al_array_push(list->entries, entry); + return true; } -void lia_list_unset(struct lia_list *list) +static void unset_all(struct lia_list *list) +{ + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->set = -1; + sink->queued = -1; + } +} + +static void handle_unset(struct lia_list *list) { unset_all(list); list->current = list->entries.size - 1; @@ -178,39 +206,40 @@ static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 return al_array_at(list->entries, sequence); } +// TODO: +// - Think about what is means for an entry to be done. never put into a pause state? +// - clock_end()?? +// - Sink needs to handle case where entry gets queued but the list already expects it to be playing +// - It's possible to know if sink->current is done during a set command, synchronously. +// So, check that when queueing an entry. +// - Can clock be ended during a queue command in any other case? +// - In the simplest case of our only operation being skip, how could client's become desynced? +// - Then with toggle pause +// - Is it safe to assert paused state on the client. +// - Do queued +// - Do seek +//if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) { -void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) +static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) { - if (index == list->current) return; + if (index == list->current) return true; if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; if (sequence == list->previous) { // This can happen but almost certainly won't be expected behavior. - return; + return true; } struct lia_list_entry *current = get_entry_from_sequence(list, sequence); struct lia_list_entry *target = get_entry_from_sequence(list, index); - if (!current || !target) return; + if (!current || !target) return true; + if (!entry_load_and_get_duration(list, target)) { + return false; + } u64 now = aki_get_timestamp(); u64 at = now + LIANA_BASE_DELAY; u8 pause; - // TODO: - // - Think about what is means for an entry to be done. never put into a pause state? - // - clock_end()?? - // - Sink needs to handle case where entry gets queued but the list already expects it to be playing - // - It's possible to know if sink->current is done during a set command, synchronously. - // So, check that when queueing an entry. - // - Can clock be ended during a queue command in any other case? - // - In the simplest case of our only operation being skip, how could client's become desynced? - // - Then with toggle pause - // - Is it safe to assert paused state on the client. - // - Do queued - // - Do seek - - //if (current->start != LIANA_TIMESTAMP_INVALID && current->start > ts - LIANA_BASE_PING) { - // This should only happen if `start` has never been set. if (target->start == LIANA_TIMESTAMP_INVALID && target->paused_at == LIANA_TIMESTAMP_INVALID) { target->start = at; @@ -264,16 +293,18 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time); } - list->callback(list->userdata, LIANA_META_PLAYING, target); + // meta playing + + return true; } -void lia_list_skip(struct lia_list *list, s32 sequence, s32 n) +static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) { if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; - lia_list_skipto(list, sequence, sequence + n); + return handle_skipto(list, sequence, sequence + n); } -void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) +static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) { if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; struct lia_list_entry *current = get_entry_from_sequence(list, sequence); @@ -284,7 +315,9 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) u64 at; switch (pause) { case LIANA_PAUSE_PAUSE: - al_assert(pts != -1.0); + // This assert should exist but the correct behavior for this is unfinished. + //al_assert(pts != -1.0); + (void)pts; current->paused_at = now + LIANA_PAUSE_DELAY; current->offset += current->paused_at - current->start; current->start = LIANA_TIMESTAMP_INVALID; @@ -310,9 +343,10 @@ 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) +static void handle_seek(struct lia_list *list, s32 sequence, f64 percent) { if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + if (sequence < 0) return; list->idle = false; struct lia_list_entry *current = get_entry_from_sequence(list, sequence); u64 pos = (u64)(current->duration * percent); @@ -332,7 +366,7 @@ void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent) } } -void lia_list_end(struct lia_list *list, s32 sequence) +static void handle_end(struct lia_list *list, s32 sequence) { al_assert(sequence != LIANA_SEQUENCE_ANY); if (sequence != list->current) return; @@ -353,11 +387,13 @@ void lia_list_end(struct lia_list *list, s32 sequence) al_array_foreach(list->sinks, i, sink) { sink->queued = -1; } - if (list->callback) { - list->callback(list->userdata, LIANA_META_PLAYING, al_array_at(list->entries, list->current)); - } + // meta playing } else if (next < size) { - lia_list_skipto(list, sequence, next); + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = SKIPTO; + cmd->sequence = sequence; + cmd->i = next; + al_array_push(list->queue, cmd); } else { list->idle = true; struct lia_list_sink *sink; @@ -367,7 +403,7 @@ void lia_list_end(struct lia_list *list, s32 sequence) } } -void lia_list_reverse(struct lia_list *list) +static void handle_reverse(struct lia_list *list) { u32 size = list->entries.size; for (u32 i = 0; i < size; i++) { @@ -378,13 +414,13 @@ void lia_list_reverse(struct lia_list *list) unset_all(list); } -void lia_list_sort(struct lia_list *list) +static void handle_sort(struct lia_list *list) { al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare); unset_all(list); } -void lia_list_shuffle(struct lia_list *list) +static void handle_shuffle(struct lia_list *list) { u32 size = list->entries.size; if (size == 0) return; @@ -395,7 +431,7 @@ void lia_list_shuffle(struct lia_list *list) unset_all(list); } -void lia_list_clear(struct lia_list *list) +static void handle_clear(struct lia_list *list) { struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { @@ -409,6 +445,200 @@ void lia_list_clear(struct lia_list *list) list->idle = true; } +static void pump_queue(struct lia_list *list) +{ + if (!list->cmd) { + if (list->queue.size == 0) return; + al_array_pop_at(list->queue, 0, list->cmd); + } + struct lia_list_cmd *cmd = list->cmd; + switch (cmd->op) { + case ADD_SINK: + if (!handle_add_sink(list, cmd->sink)) { + return; + } + break; + case REMOVE_SINK: + handle_remove_sink(list, cmd->userdata); + break; + case ADD: + if (!handle_add(list, cmd->entry)) { + return; + } + break; + case UNSET: + handle_unset(list); + break; + case SKIPTO: + if (!handle_skipto(list, cmd->sequence, cmd->i)) { + return; + } + break; + case SKIP: + if (!handle_skip(list, cmd->sequence, cmd->i)) { + return; + } + break; + case TOGGLE_PAUSE: + handle_toggle_pause(list, cmd->sequence, cmd->f); + break; + case SEEK: + handle_seek(list, cmd->sequence, cmd->f); + break; + case END: + handle_end(list, cmd->sequence); + break; + case REVERSE: + handle_reverse(list); + break; + case SORT: + handle_sort(list); + break; + case SHUFFLE: + handle_shuffle(list); + break; + case CLEAR: + handle_clear(list); + break; + } + al_free(cmd); + list->cmd = NULL; + pump_queue(list); +} + +void lia_list_pump(struct lia_list *list) +{ + pump_queue(list); +} + +void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata) +{ + struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink); + sink->callback = callback; + sink->userdata = userdata; + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = ADD_SINK; + cmd->sink = sink; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_remove_sink(struct lia_list *list, void *userdata) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = REMOVE_SINK; + cmd->userdata = userdata; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) +{ + struct lia_list_entry *entry = al_alloc_object(struct lia_list_entry); + entry->opaque = opaque; + entry->paused_at = LIANA_TIMESTAMP_INVALID; + entry->held = false; + entry->offset = 0; + entry->ended = false; + entry->duration = duration; + al_str_clone(&entry->name, name); + entry->list = list; + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = ADD; + cmd->entry = entry; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_unset(struct lia_list *list) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = UNSET; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +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; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +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; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = TOGGLE_PAUSE; + cmd->sequence = sequence; + cmd->f = pts; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = SEEK; + cmd->sequence = sequence; + cmd->f = percent; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_end(struct lia_list *list, s32 sequence) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = END; + cmd->sequence = sequence; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_reverse(struct lia_list *list) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = REVERSE; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_sort(struct lia_list *list) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = SORT; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_shuffle(struct lia_list *list) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = SHUFFLE; + al_array_push(list->queue, cmd); + pump_queue(list); +} + +void lia_list_clear(struct lia_list *list) +{ + struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); + cmd->op = CLEAR; + al_array_push(list->queue, cmd); + pump_queue(list); +} + void lia_list_free(struct lia_list *list) { struct lia_list_entry *entry; @@ -422,4 +652,5 @@ void lia_list_free(struct lia_list *list) al_free(sink); } al_array_free(list->sinks); + al_str_free(&list->name); } diff --git a/src/liana/list.h b/src/liana/list.h index 5e3eaec..270ef98 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -25,14 +25,10 @@ enum { enum { LIANA_LOAD_ENTRY = 0, + LIANA_GET_DURATION, LIANA_UNLOAD_ENTRY }; -enum { - LIANA_META_PLAYING = 0, - LIANA_META_QUEUED -}; - // NOTE: To handle an entry being queued right before a skip, keep a global // "max time until all sinks buffered" and used that instead of LIANA_PAUSE_DELAY (if greater). @@ -61,6 +57,7 @@ struct lia_list_entry { bool ended; u64 duration; str name; + struct lia_list *list; }; struct lia_list_sink { @@ -70,18 +67,33 @@ struct lia_list_sink { void *userdata; }; +struct lia_list_cmd { + u8 op; + struct lia_list_sink *sink; + void *userdata; + struct lia_list_entry *entry; + s32 sequence; + s32 i; + f64 f; +}; + struct lia_list { + str name; s32 current; s32 previous; s32 queued; bool idle; array(struct lia_list_entry *) entries; array(struct lia_list_sink *) sinks; - void (*callback)(void *, u8, struct lia_list_entry *); + array(struct lia_list_cmd *) queue; + struct lia_list_cmd *cmd; + void (*callback)(void *, u8, struct lia_list_entry *, void *); void *userdata; }; -void lia_list_init(struct lia_list *list); +void lia_list_init(struct lia_list *list, str *name); + +void lia_list_pump(struct lia_list *list); void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata); void lia_list_remove_sink(struct lia_list *list, void *userdata); diff --git a/src/liana/server.c b/src/liana/server.c index 3eb303e..ccca0fd 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -289,23 +289,35 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en return node; } -u64 lia_node_get_duration(struct lia_node *node) +static aki_thread_result AKI_THREADCALL init_duration_thread(void *userdata) { - (void)node; -#if 0 - return 0; -#else - struct cch_handle handle; - cch_entry_get_handle(node->entry, &handle); - struct lia_server_handler *handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); - if (!handler->init(handler, &handle)) { - return LIANA_TIMESTAMP_INVALID; + struct lia_node *node = (struct lia_node *)userdata; + if (!node->handler->init(node->handler, &node->handle)) { + node->errored = true; + } else { + node->duration = node->handler->get_duration(node->handler); } - u64 duration = handler->get_duration(handler); - handler->free(&handler); - cch_entry_return_handle(node->entry, &handle); - return duration; -#endif + aki_signal_send(&node->signal); + return 0; +} + +static void duration_signal_callback(void *userdata) +{ + struct lia_node *node = (struct lia_node *)userdata; + aki_signal_stop(&node->signal); + aki_thread_join(&node->thread); + node->handler->free(&node->handler); + cch_entry_return_handle(node->entry, &node->handle); + node->callback(node->userdata, LIANA_NODE_DURATION, node->duration); +} + +void lia_node_get_duration(struct lia_node *node) +{ + aki_signal_init(&node->signal, duration_signal_callback, node); + aki_signal_start(&node->signal, node->server->loop); + cch_entry_get_handle(node->entry, &node->handle); + node->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); + aki_thread_create(&node->thread, init_duration_thread, node); } void lia_server_close(struct lia_server *server) diff --git a/src/liana/server.h b/src/liana/server.h index 3c4c294..2ed7c12 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -21,11 +21,25 @@ struct lia_node_connection { struct lia_node *node; }; +enum { + LIANA_NODE_DURATION = 0 +}; + struct lia_node { u16 id; struct cch_entry *entry; array(struct lia_node_connection *) connections; struct lia_server *server; + u64 duration; + // Temporary copy from node_connection. We need to + // figure out a "connection pool" structure. + struct lia_server_handler *handler; + bool errored; + struct cch_handle handle; + struct aki_thread thread; + struct aki_signal signal; + void (*callback)(void *, u8, u64); + void *userdata; }; struct lia_server { @@ -37,6 +51,6 @@ struct lia_server { bool lia_server_init(struct lia_server *server, struct aki_event_loop *loop); void lia_server_add_socket(struct lia_server *server, struct aki_socket *sock); struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry); -u64 lia_node_get_duration(struct lia_node *node); +void lia_node_get_duration(struct lia_node *node); void lia_server_close(struct lia_server *server); void lia_server_free(struct lia_server *server); diff --git a/src/libclient/client.c b/src/libclient/client.c index 58264e0..0776842 100644 --- a/src/libclient/client.c +++ b/src/libclient/client.c @@ -1,8 +1,41 @@ #include "client.h" +#include "common.h" #include "../server/common.h" -static struct aki_rpc_command commands[] = { }; +static bool results_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) +{ + struct camu_client *client = (struct camu_client *)userdata; + (void)conn; + (void)rpacket; + + u8 op = aki_packet_read_u8(packet); + switch (op) { + case CAMU_CLIENT_CREATE_SEARCH: { + s32 id = aki_packet_read_s32(packet); + client->callback(client->userdata, CAMU_CLIENT_SEARCH_CREATED, &id); + break; + } + case CAMU_CLIENT_GET_PAGE: + client->callback(client->userdata, CAMU_CLIENT_PAGE_RESULTS, packet); + break; + } + + aki_packet_free(packet); + return false; +} + +static struct aki_rpc_command commands[] = { + { .op = CAMU_CLIENT_RESULTS, .callback = results_callback, .userdata = NULL } +}; + +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); +} static void connection_callback(void *userdata, struct aki_rpc_connection *conn) { @@ -11,7 +44,7 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn) struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_IDENTIFY); aki_packet_write_u8(packet, CAMU_CLIENT); aki_packet_write_str(packet, &client->username); - aki_rpc_connection_command(client->conn, packet, NULL, NULL); + aki_rpc_connection_command(client->conn, packet, idd_callback, client); } static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) @@ -47,12 +80,42 @@ void camu_client_create_list(struct camu_client *client, str *name, aki_rpc_connection_command(client->conn, packet, callback, userdata); } -void camu_client_enable_sink(struct camu_client *client, str *list, str *sink, +void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable, void (*callback)(void *, struct aki_packet *), void *userdata) { struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); - aki_packet_write_u8(packet, CAMU_CLIENT_ENABLE_SINK); - aki_packet_write_str(packet, list); + aki_packet_write_u8(packet, CAMU_CLIENT_TOGGLE_SINK); aki_packet_write_str(packet, sink); + aki_packet_write_str(packet, list); + aki_packet_write_bool(packet, enable); aki_rpc_connection_command(client->conn, packet, callback, userdata); } + +void camu_client_create_search(struct camu_client *client, str *module, str *query) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + aki_packet_write_u8(packet, CAMU_CLIENT_CREATE_SEARCH); + aki_packet_write_str(packet, module); + aki_packet_write_str(packet, query); + aki_rpc_connection_command(client->conn, packet, NULL, NULL); +} + +void camu_client_get_page(struct camu_client *client, s32 id, u32 num) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_CLIENT_COMMAND); + aki_packet_write_u8(packet, CAMU_CLIENT_GET_PAGE); + aki_packet_write_s32(packet, id); + aki_packet_write_u32(packet, num); + aki_rpc_connection_command(client->conn, packet, NULL, NULL); +} + +void camu_client_add(struct camu_client *client, str *list, str *unique_id, u32 index) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, CAMU_SERVER_LIST_ACTION); + aki_packet_write_str(packet, list); + aki_packet_write_u8(packet, CAMU_LIST_ADD); + aki_packet_write_u8(packet, CAMU_RESOURCE_PORTAL); + aki_packet_write_str(packet, unique_id); + aki_packet_write_u32(packet, index); + aki_rpc_connection_command(client->conn, packet, NULL, NULL); +} diff --git a/src/libclient/client.h b/src/libclient/client.h index a56be26..02a6639 100644 --- a/src/libclient/client.h +++ b/src/libclient/client.h @@ -2,11 +2,19 @@ #include <aki/rpc2.h> +enum { + CAMU_CLIENT_LOGIN = 0, + CAMU_CLIENT_SEARCH_CREATED, + CAMU_CLIENT_PAGE_RESULTS +}; + struct camu_client { struct aki_event_loop *loop; str username; struct aki_rpc client; struct aki_rpc_connection *conn; + void (*callback)(void *, u8, void *); + void *userdata; }; bool camu_client_login(struct camu_client *client, struct aki_event_loop *loop, @@ -14,5 +22,10 @@ bool camu_client_login(struct camu_client *client, struct aki_event_loop *loop, void camu_client_create_list(struct camu_client *client, str *name, void (*callback)(void *, struct aki_packet *), void *userdata); -void camu_client_enable_sink(struct camu_client *client, str *list, str *sink, +void camu_client_toggle_sink(struct camu_client *client, str *sink, str *list, bool enable, void (*callback)(void *, struct aki_packet *), void *userdata); + +void camu_client_create_search(struct camu_client *client, str *module, str *query); +void camu_client_get_page(struct camu_client *client, s32 id, u32 num); + +void camu_client_add(struct camu_client *client, str *list, str *unique_id, u32 index); diff --git a/src/libclient/common.h b/src/libclient/common.h new file mode 100644 index 0000000..5b8dcbf --- /dev/null +++ b/src/libclient/common.h @@ -0,0 +1,5 @@ +#pragma once + +enum { + CAMU_CLIENT_RESULTS = 0 +}; diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 2c8790b..899023d 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -91,6 +91,8 @@ static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_e entry->audio.state = BUFFER_SET_OR_BUFFERED; } else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) { entry->audio.state = BUFFER_CONFIGURED; + } else if (entry->audio.state == BUFFER_QUEUED) { + entry->audio.state = BUFFER_INIT; } } @@ -102,6 +104,9 @@ static void remove_entry_video_buffer(struct camu_sink *sink, struct camu_sink_e entry->video.state = BUFFER_SET_OR_BUFFERED; } else if (entry->video.state == BUFFER_SET_OR_BUFFERED) { entry->video.state = BUFFER_CONFIGURED; + } else if (entry->video.state == BUFFER_QUEUED) { + // This can be hit when skipping through entries very fast. + entry->video.state = BUFFER_INIT; } } #endif @@ -276,7 +281,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) break; } case CLOSE: { - aki_signal_stop(&sink->signal); + aki_signal_stop(&sink->queue_signal); sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL); return; } @@ -298,7 +303,7 @@ static void queue_signal_callback(void *userdata) static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd) { camu_queue_push(sink->queue, cmd); - aki_signal_send(&sink->signal); + aki_signal_send(&sink->queue_signal); } static void maybe_remove_previous(struct camu_sink *sink) @@ -310,13 +315,32 @@ static void maybe_remove_previous(struct camu_sink *sink) sink->previous.size = 0; } +static void maybe_add_to_previous(struct camu_sink *sink, struct camu_sink_entry *entry, struct camu_sink_entry *current) +{ + al_assert(entry != current); + struct camu_sink_entry *rentry; + al_array_foreach_rev(sink->previous, i, rentry) { + if (rentry == current) { + // If the entry we are about to add is in previous, + // remove it immediately. + remove_entry_buffers(sink, current); + al_array_remove_at(sink->previous, i); + } + } + // Don't accept duplicates. + al_array_foreach_rev(sink->previous, i, rentry) { + if (rentry == entry) return; + } + al_array_push(sink->previous, entry); +} + void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { - u8 state = entry->audio.state; - al_assert(state != BUFFER_ADDED); - if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } else if (state == BUFFER_SET_OR_BUFFERED) { + al_assert(entry->audio.state != BUFFER_ADDED); + if (entry->audio.state == BUFFER_CONFIGURED) { + entry->audio.state = BUFFER_SET_OR_BUFFERED; + } else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) { + entry->audio.state = BUFFER_ADDED; entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); // It's possible for this entry's video buffer to have been added and removed by EOF // before this point. This needs to be a consideration for keeping sync. @@ -324,24 +348,24 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) if (VIDEO_READY_OR_EMPTY(entry)) { maybe_remove_previous(entry->sink); } +#ifndef CAMU_SINK_LOCAL camu_audio_buffer_unpause(&entry->audio.buf); +#endif queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO }); - state = BUFFER_ADDED; } - entry->audio.state = state; } #ifndef CAMU_SINK_NO_VIDEO void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { - u8 state = entry->video.state; - al_assert(state != BUFFER_ADDED); - if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; - } else if (state == BUFFER_SET_OR_BUFFERED) { + al_assert(entry->video.state != BUFFER_ADDED); + if (entry->video.state == BUFFER_CONFIGURED) { + entry->video.state = BUFFER_SET_OR_BUFFERED; + } else if (entry->video.state == BUFFER_SET_OR_BUFFERED) { + entry->video.state = BUFFER_ADDED; entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); if (AUDIO_READY_OR_EMPTY(entry)) { maybe_remove_previous(entry->sink); @@ -351,9 +375,7 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry) .op = single_frame ? STOP : START, .value.i = CAMU_SINK_VIDEO }); - state = BUFFER_ADDED; } - entry->video.state = state; } #endif @@ -416,7 +438,7 @@ static void video_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_BUFFERED: { bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); aki_mutex_lock(&sink->mutex); - if (single_frame || !entry->ended) { + if (!entry->ended || single_frame) { add_video_if_set_and_buffered(entry); } aki_mutex_unlock(&sink->mutex); @@ -477,16 +499,21 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent frames -= sink->video.renderer->get_latency(sink->video.renderer); camu_video_buffer_set_latency(&entry->video.buf, -frames); } +#else + (void)sink; + (void)entry; #endif #else // To sync clients with differing audio latencies our only option is to factor the mixer // latency directly into the audio buffer. f64 audio = camu_mixer_get_latency(sink->audio.mixer); +#ifndef CAMU_SINK_NO_VIDEO if (!BUFFER_EMPTY(&entry->video)) { s32 frames = audio / entry->video.buf.avg_frame_duration; frames += sink->video.renderer->get_latency(sink->video.renderer); camu_video_buffer_set_latency(&entry->video.buf, frames); } +#endif camu_audio_buffer_set_latency(&entry->audio.buf, audio); #endif } @@ -650,10 +677,10 @@ bool camu_sink_init(struct camu_sink *sink, struct aki_event_loop *loop, ) { 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); + aki_signal_init(&sink->queue_signal, queue_signal_callback, sink); + aki_signal_start(&sink->queue_signal, sink->loop); + camu_queue_init(sink->queue); sink->queued = NULL; sink->current = NULL; al_array_init(sink->previous); @@ -720,7 +747,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *entry) } } else { if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } } set_or_queue_entry(entry); @@ -832,8 +859,6 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn entry->ended = ended; - if (entry == sink->current) goto out; - if (op == LIANA_SINK_BUFFER) { goto out; } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { @@ -844,10 +869,14 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn #ifdef CAMU_SINK_LOCAL (void)at; (void)pause; - if (sink->current && !camu_clock_is_paused(&sink->current->clock)) { - camu_clock_pause(&sink->current->clock, 0); + if (sink->current) { + if (!camu_clock_is_paused(&sink->current->clock)) { + camu_clock_pause(&sink->current->clock, 0); + } + maybe_add_to_previous(sink, sink->current, entry); } - switch_to(sink, entry); + set_or_queue_entry(entry); + sink->current = entry; // This will resume a user paused stream. camu_clock_resume(&entry->clock, 0); #else @@ -858,7 +887,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn sink->target = NULL; } if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } set_or_queue_entry(entry); sink->current = entry; @@ -871,7 +900,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (entry == sink->current) { camu_audio_buffer_unpause(&entry->audio.buf); } else if (sink->current) { - al_array_push(sink->previous, sink->current); + maybe_add_to_previous(sink, sink->current, entry); } camu_clock_resume(&entry->clock, at); set_or_queue_entry(entry); @@ -1028,6 +1057,14 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn) aki_rpc_connection_command(sink->conn, packet, idd_callback, sink); } +static void reconnect_timer_callback(void *userdata, struct aki_timer *timer) +{ + struct camu_sink *sink = (struct camu_sink *)userdata; + (void)timer; + aki_rpc_reconnect(&sink->client, &sink->addr, sink->port); + aki_timer_stop(&sink->reconnect_timer); +} + static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { struct camu_sink *sink = (struct camu_sink *)userdata; @@ -1035,12 +1072,15 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection al_assert(sink->conn == conn); sink->conn = NULL; } + aki_timer_again(&sink->reconnect_timer); } bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str *name) { al_str_clone(&sink->name, name); sink->type = type; + aki_timer_init(&sink->reconnect_timer, sink->loop, reconnect_timer_callback, sink); + aki_timer_set_repeat(&sink->reconnect_timer, AKI_TS_FROM_USEC(1000000)); aki_rpc_init(&sink->client, sink->loop, connection_callback, connection_closed_callback, sink); for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { commands[i].userdata = sink; @@ -1050,7 +1090,9 @@ bool camu_sink_connect(struct camu_sink *sink, u8 type, str *addr, u16 port, str if (!aki_rpc_prepare_client(&sink->client, sink->type, CAMU_MULTIPLEX_RPC)) { return false; } - aki_rpc_connect(&sink->client, addr, port); + al_str_clone(&sink->addr, addr); + sink->port = port; + aki_rpc_connect(&sink->client, &sink->addr, sink->port); return true; } @@ -1142,6 +1184,8 @@ void camu_sink_stop(struct camu_sink *sink) void camu_sink_close(struct camu_sink *sink) { + aki_timer_stop(&sink->reconnect_timer); + aki_timer_disable(&sink->reconnect_timer); if (sink->conn) aki_rpc_conn_disconnect(sink->conn); struct camu_sink_entry *entry; al_array_foreach_rev(sink->entries, i, entry) { @@ -1160,4 +1204,6 @@ void camu_sink_free(struct camu_sink *sink) aki_rpc_free(&sink->client); camu_queue_free(sink->queue); aki_mutex_destroy(&sink->mutex); + al_str_free(&sink->addr); + al_str_free(&sink->name); } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index b830230..0b55ad6 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -4,6 +4,7 @@ #include <al/array.h> #include <aki/rpc2.h> #include <aki/signal.h> +#include <aki/timer.h> #include "../util/queue.h" @@ -68,11 +69,14 @@ struct camu_sink { struct aki_event_loop *loop; str name; u8 type; + str addr; + u16 port; struct aki_rpc client; struct aki_rpc_connection *conn; - struct aki_signal signal; - queue(struct camu_sink_cmd) queue; struct aki_mutex mutex; + struct aki_timer reconnect_timer; + struct aki_signal queue_signal; + queue(struct camu_sink_cmd) queue; str default_list; struct camu_sink_entry *current; struct camu_sink_entry *queued; diff --git a/src/mixer/audio_miniaudio.c b/src/mixer/audio_miniaudio.c index 72e1cd8..2a93491 100644 --- a/src/mixer/audio_miniaudio.c +++ b/src/mixer/audio_miniaudio.c @@ -54,8 +54,10 @@ static ma_allocation_callbacks alloc_callbacks = { static void miniaudio_log_callback(void *userdata, u32 level, const char *message) { (void)userdata; - if (level < MA_LOG_LEVEL_WARNING) { - al_log_debug("audio_miniaudio", message); + if (level == MA_LOG_LEVEL_ERROR) { + al_log_error("audio_miniaudio", message); + } else if (level == MA_LOG_LEVEL_WARNING) { + al_log_warn("audio_miniaudio", message); } else { al_log_debug("audio_miniaudio", message); } diff --git a/src/mixer/mixer.h b/src/mixer/mixer.h index 729657f..9a15751 100644 --- a/src/mixer/mixer.h +++ b/src/mixer/mixer.h @@ -1,7 +1,5 @@ #pragma once -#define CAMU_MIXER_THREADED - #include <al/types.h> #include <al/array.h> #ifdef CAMU_MIXER_THREADED diff --git a/src/portal/meson.build b/src/portal/meson.build index 7cf31a7..4e5ff56 100644 --- a/src/portal/meson.build +++ b/src/portal/meson.build @@ -8,7 +8,7 @@ portal_deps = [] portal_args = ['-DCAMU_HAVE_PORTAL'] cpy_dir = join_paths(meson.current_source_dir(), 'cpy') -run_command(join_paths(cpy_dir, 'build.sh'), cpy_dir) +run_command(join_paths(cpy_dir, 'build.sh'), cpy_dir, check: false) python3_embed = import('python').find_installation('python3.12').dependency(embed: true) portal_deps += [python3_embed] diff --git a/src/portal/src/packet_ext.c b/src/portal/src/packet_ext.c index e2909ef..833e7d9 100644 --- a/src/portal/src/packet_ext.c +++ b/src/portal/src/packet_ext.c @@ -6,7 +6,7 @@ static void aki_packet_write_optional_int(struct aki_packet *packet, optional_in AKI_PACKET_WRITE_TYPE(packet, bool, o->set); } -void aki_packet_write_camu_post(struct aki_packet *packet, struct camu_post *post) +void aki_packet_write_post(struct aki_packet *packet, struct camu_post *post) { AKI_PACKET_WRITE_TYPE(packet, u16, post->version); AKI_PACKET_WRITE_TYPE(packet, u8, post->type); @@ -53,7 +53,7 @@ static void aki_packet_read_optional_int(struct aki_packet *packet, optional_int AKI_PACKET_READ_TYPE(packet, bool, o->set); } -void aki_packet_read_camu_post(struct aki_packet *packet, struct camu_post *post) +void aki_packet_read_post(struct aki_packet *packet, struct camu_post *post) { camu_post_reset(post); AKI_PACKET_READ_TYPE(packet, u16, post->version); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index bcb0643..e6022f5 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -1,5 +1,7 @@ #include <al/log.h> +#include "../../server/common.h" + #include "../cpy/portal.c" #include "search.h" @@ -32,49 +34,168 @@ void camu_python_close(void) if (Py_IsInitialized()) Py_Finalize(); } -static struct camu_result_page *page_at_index(struct camu_search *search, u32 num) +static struct camu_search *get_search_by_id(struct camu_portal_bridge *bridge, s32 id) { - struct camu_result_page *page; - al_array_foreach_ptr(search->pages, i, page) { - if (page->num == num) return page; + struct camu_search *search; + al_array_foreach(bridge->searches, i, search) { + if (search->id == id) return search; } - al_array_push(search->pages, (struct camu_result_page){ 0 }); - page = &al_array_last(search->pages); - page->num = num; - al_array_init(page->posts); - al_array_init(page->list); - return page; + return NULL; } +static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) +{ + struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; + aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS); + bool have_python = false; + aki_mutex_lock(&bridge->mutex); + do { + aki_cond_wait(&bridge->cond, &bridge->mutex); + if (bridge->quit) break; + if (!have_python) { + // Defer python init. + have_python = camu_python_init(); + } + struct camu_portal_cmd *cmd; + al_array_foreach_ptr(bridge->queue, i, cmd) { + struct camu_portal_result result; + result.op = cmd->op; + result.callback = cmd->callback; + result.userdata = cmd->userdata; + switch (cmd->op) { + case CAMU_CLIENT_CREATE_SEARCH: { + s32 id = portal_bridge_search(&cmd->module, &cmd->query); + if (id >= 0) { + struct camu_search *search = al_alloc_object(struct camu_search); + search->page = 0; + al_array_init(search->pages); + search->id = id; + al_str_clone(&search->module, &cmd->module); + al_str_clone(&search->query, &cmd->query); + search->bridge = bridge; + al_array_push(bridge->searches, search); + al_log_info("portal", "New search %x (%.*s).", id, AL_STR_PRINTF(&cmd->query)); + result.id = search->id; + } else { + } + al_str_free(&cmd->module); + al_str_free(&cmd->query); + break; + } + case CAMU_CLIENT_GET_PAGE: { + struct camu_search *search = get_search_by_id(bridge, cmd->id); + if (search) { + result.id = search->id; + struct camu_result_page *page; + al_array_foreach_ptr(search->pages, j, page) { + if (page->num == cmd->num) break; + } + al_log_info("portal", "Loading page %i (%.*s).", cmd->num, AL_STR_PRINTF(&search->query)); + if (portal_bridge_get_page(search, search->id, cmd->num) == -1) { + break; + } + page = &al_array_at(search->pages, cmd->num); + if (bridge->cache) { + struct camu_post *post; + al_array_foreach_ptr(page->posts, j, post) { + camu_post_cache_push(bridge->cache, post); + } + } + result.page = page; + } + break; + } + } + camu_queue_push(bridge->results, result); + aki_signal_send(&bridge->results_signal); + al_array_remove_at_iter(bridge->queue, i); + } + } while (1); + aki_mutex_unlock(&bridge->mutex); + if (have_python) { + camu_python_close(); + } + return 0; +} -void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache) +static void results_signal_callback(void *userdata) +{ + struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; + u32 size; + struct camu_portal_result result; + do { + camu_queue_try_pop(bridge->results, size, result); + if (size == 0) break; + result.callback(result.userdata, &result); + } while (1); +} + +void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, + struct aki_event_loop *loop) { - bridge->cache = cache; al_array_init(bridge->searches); + bridge->cache = cache; + bridge->quit = 0; + aki_mutex_init(&bridge->mutex); + aki_cond_init(&bridge->cond); + al_array_init(bridge->queue); + camu_queue_init(bridge->results); + aki_signal_init(&bridge->results_signal, results_signal_callback, bridge); + aki_signal_start(&bridge->results_signal, loop); + 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) +{ + struct camu_portal_cmd cmd; + cmd.op = CAMU_CLIENT_CREATE_SEARCH; + al_str_clone(&cmd.module, module); + al_str_clone(&cmd.query, query); + cmd.callback = callback; + cmd.userdata = userdata; + aki_mutex_lock(&bridge->mutex); + al_array_push(bridge->queue, cmd); + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); + } + aki_mutex_unlock(&bridge->mutex); } -static void camu_search_init_internal(struct camu_search *search) +void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, + void (*callback)(void *, struct camu_portal_result *), void *userdata) { - search->page = 0; - al_array_init(search->pages); + struct camu_portal_cmd cmd; + cmd.op = CAMU_CLIENT_GET_PAGE; + cmd.id = id; + cmd.num = num; + cmd.callback = callback; + cmd.userdata = userdata; + aki_mutex_lock(&bridge->mutex); + al_array_push(bridge->queue, cmd); + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); + } + aki_mutex_unlock(&bridge->mutex); } -s32 camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query) +void camu_portal_close(struct camu_portal_bridge *bridge) { - s32 id = portal_bridge_search(module, query); - if (id >= 0) { - struct camu_search *search = al_alloc_object(struct camu_search); - camu_search_init_internal(search); - search->id = id; - al_str_clone(&search->module, module); - al_str_clone(&search->query, query); - search->bridge = bridge; - al_array_push(bridge->searches, search); + aki_mutex_lock(&bridge->mutex); + bridge->quit = 1; + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); } - al_log_info("portal", "New search %x (%.*s).", id, AL_STR_PRINTF(query)); - return id; + aki_mutex_unlock(&bridge->mutex); + aki_thread_join(&bridge->thread); + aki_signal_stop(&bridge->results_signal); + camu_queue_free(bridge->results); + al_array_free(bridge->queue); + aki_cond_destroy(&bridge->cond); + aki_mutex_destroy(&bridge->mutex); } +/* struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id) { struct camu_search *search; @@ -90,42 +211,30 @@ void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id) (void)id; } -void camu_portal_close(struct camu_portal_bridge *bridge) -{ - (void)bridge; -} - -bool camu_search_get_page(struct camu_search *search, u32 num) +void camu_search_free(struct camu_search *search) { struct camu_result_page *page; al_array_foreach_ptr(search->pages, i, page) { - if (page->num == num) goto out; - } - al_log_info("portal", "Loading page %i (%.*s).", num, AL_STR_PRINTF(&search->query)); - if (portal_bridge_get_page(search, search->id, num) == -1) { - return false; - } - struct camu_portal_bridge *bridge = search->bridge; - if (bridge->cache) { - struct camu_post *post; - al_array_foreach_ptr(al_array_at(search->pages, num).posts, i, post) { - camu_post_cache_push(bridge->cache, post); - } + // TODO: Free camu_post ? + al_array_free(page->posts); + al_array_free(page->list); } -out: - search->page = num; - return true; + al_array_free(search->pages); } +*/ -void camu_search_free(struct camu_search *search) +static struct camu_result_page *page_at_index(struct camu_search *search, u32 num) { struct camu_result_page *page; al_array_foreach_ptr(search->pages, i, page) { - // TODO: Free camu_post ? - al_array_free(page->posts); - al_array_free(page->list); + if (page->num == num) return page; } - al_array_free(search->pages); + al_array_push(search->pages, (struct camu_result_page){ 0 }); + page = &al_array_last(search->pages); + page->num = num; + al_array_init(page->posts); + al_array_init(page->list); + return page; } void camu_search_add_post(struct camu_search *search, u32 num, struct camu_post *post) diff --git a/src/portal/src/search.h b/src/portal/src/search.h index 1f43ded..85f4f6d 100644 --- a/src/portal/src/search.h +++ b/src/portal/src/search.h @@ -1,8 +1,13 @@ #pragma once +#include <aki/thread.h> +#include <aki/signal.h> + #include "post.h" #include "post_cache.h" +#include "../../util/queue.h" + struct camu_result_page { u32 num; array(struct camu_post) posts; @@ -18,21 +23,53 @@ struct camu_search { struct camu_portal_bridge *bridge; }; +struct camu_portal_result { + u8 op; + s32 id; + struct camu_result_page *page; + void (*callback)(void *, struct camu_portal_result *); + void *userdata; +}; + +struct camu_portal_cmd { + u8 op; + str module; + str query; + s32 id; + u32 num; + void (*callback)(void *, struct camu_portal_result *); + void *userdata; +}; + struct camu_portal_bridge { array(struct camu_search *) searches; struct camu_post_cache *cache; + u8 quit; + struct aki_thread thread; + struct aki_mutex mutex; + struct aki_cond cond; + array(struct camu_portal_cmd) queue; + queue(struct camu_portal_result) results; + struct aki_signal results_signal; }; bool camu_python_init(void); void camu_python_close(void); -void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache); -s32 camu_portal_create_search(struct camu_portal_bridge *bridge, str *module_str, str *search_str); -struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id); -void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id); +void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, + struct aki_event_loop *loop); + +void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, + void (*callback)(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 camu_portal_close(struct camu_portal_bridge *bridge); -bool camu_search_get_page(struct camu_search *search, u32 num); +/* +struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id); +void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id); +*/ // Python internal. void camu_search_add_post(struct camu_search *search, u32 num, struct camu_post *post); diff --git a/src/screen/screen.c b/src/screen/screen.c index 7f1ad30..f949c1d 100644 --- a/src/screen/screen.c +++ b/src/screen/screen.c @@ -109,8 +109,13 @@ static bool mouse_button_callback(void *userdata, u8 state, u8 button) } break; case STELA_MOUSE3: { - f64 percent = scr->last_mouse_x / scr->width; - scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent); + switch (state) { + case STELA_BUTTON_RELEASED: { + f64 percent = scr->last_mouse_x / scr->width; + scr->callback(scr->userdata, CAMU_SCREEN_SEEK, &percent); + break; + } + } break; } default: @@ -157,11 +162,13 @@ static bool key_callback(void *userdata, u8 state, u8 button) case 0x31: // n case 0x20: // d case 0x6a: // right arrow + case 0x4d: // right arrow (wine?) scr->callback(scr->userdata, CAMU_SCREEN_NEXT, NULL); break; case 0x30: // b case 0x1e: // a case 0x69: // left arrow + case 0x4b: // left arrow (wine?) scr->callback(scr->userdata, CAMU_SCREEN_PREVIOUS, NULL); break; case 0x39: // spacebar diff --git a/src/screen/screen.h b/src/screen/screen.h index 5ed33b3..a10796a 100644 --- a/src/screen/screen.h +++ b/src/screen/screen.h @@ -1,7 +1,5 @@ #pragma once -#define CAMU_SCREEN_THREADED - #include <al/array.h> #ifdef CAMU_SCREEN_THREADED #include <al/atomic.h> diff --git a/src/server/common.h b/src/server/common.h index c31efcf..bde8106 100644 --- a/src/server/common.h +++ b/src/server/common.h @@ -14,6 +14,7 @@ extern str *CAMU_UNIX_LOCAL; //#define CAMU_LOCAL_TYPE AKI_SOCKET_TCP //#define CAMU_LOCAL_ADDR CAMU_SERVER_IP + #define CAMU_LOCAL_TYPE AKI_SOCKET_UNIX #define CAMU_LOCAL_ADDR CAMU_UNIX_PATH @@ -31,7 +32,9 @@ enum { enum { CAMU_CLIENT_CREATE_LIST = 0, - CAMU_CLIENT_ENABLE_SINK + CAMU_CLIENT_TOGGLE_SINK, + CAMU_CLIENT_CREATE_SEARCH, + CAMU_CLIENT_GET_PAGE }; enum { @@ -45,6 +48,11 @@ enum { CAMU_LIST_END }; +enum { + CAMU_RESOURCE_FILE = 0, + CAMU_RESOURCE_PORTAL +}; + AL_UNUSED_FUNCTION_PUSH static bool camu_is_url(str *s, u32 i) diff --git a/src/server/db.c b/src/server/db.c index 5c5e8fd..c9608bb 100644 --- a/src/server/db.c +++ b/src/server/db.c @@ -2,7 +2,6 @@ #include <aki/file.h> #include <jansson.h> -#include "list.h" #include "server.h" static bool open_user(struct camu_server *server, struct aki_dir_entry *dir) @@ -63,9 +62,9 @@ void camu_db_close(struct camu_server *server) al_str_free(&user->name); } al_array_free(server->users); - struct camu_list *list; + struct lia_list *list; al_array_foreach(server->lists, i, list) { - camu_list_free(list); + lia_list_free(list); } al_array_free(server->lists); } diff --git a/src/server/list.c b/src/server/list.c deleted file mode 100644 index e3565d0..0000000 --- a/src/server/list.c +++ /dev/null @@ -1,71 +0,0 @@ -#include "../libsink/common.h" -#include "../server/common.h" - -#include "list.h" -#include "server.h" - -void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry) -{ - struct camu_server *server = (struct camu_server *)userdata; - (void)server; - (void)op; - (void)entry; -} - -void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) -{ - struct camu_server_sink *sink = (struct camu_server_sink *)userdata; - switch (op) { - case LIANA_SINK_SET: - case LIANA_SINK_BUFFER: - case LIANA_SINK_BUFFER_AND_QUEUE: { - struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_packet_write_str(packet, &sink->server->addr); - aki_packet_write_u16(packet, CAMU_PORT); - aki_packet_write_u16(packet, resource->node->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u64(packet, timing->seek_pos); - aki_packet_write_u8(packet, timing->pause); - aki_packet_write_bool(packet, timing->ended); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case LIANA_SINK_UNSET: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case LIANA_SINK_PAUSE: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u8(packet, timing->pause); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case LIANA_SINK_SEEK: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u64(packet, timing->seek_pos); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - } -} - -void camu_list_init(struct camu_list *list, str *name) -{ - al_str_clone(&list->name, name); - lia_list_init(&list->impl); -} - -void camu_list_free(struct camu_list *list) -{ - lia_list_free(&list->impl); - al_str_free(&list->name); -} diff --git a/src/server/list.h b/src/server/list.h deleted file mode 100644 index e472024..0000000 --- a/src/server/list.h +++ /dev/null @@ -1,16 +0,0 @@ -#pragma once - -#include "../liana/list.h" - -#include "resource.h" - -struct camu_list { - str name; - struct lia_list impl; -}; - -void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry); -void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing); - -void camu_list_init(struct camu_list *list, str *name); -void camu_list_free(struct camu_list *list); diff --git a/src/server/local_compat.c b/src/server/local_compat.c deleted file mode 100644 index 39698c5..0000000 --- a/src/server/local_compat.c +++ /dev/null @@ -1,219 +0,0 @@ -#include <al/log.h> - -#ifdef AKIYO_HAS_CURL -#include "../cache/handlers/http.h" -#endif -#include "../cache/handlers/file.h" -#ifdef LIANA_HAVE_CDIO -#include "../cache/handlers/cdio.h" -#endif - -#include "local_compat.h" -#include "common.h" -#include "server.h" - -#ifdef CAMU_HAVE_PORTAL -static bool uri_for_local(str *local, bool search, struct camu_portal_bridge *bridge, str *uri, struct camu_post **selected) -{ - if (search || camu_is_url(local, 0)) { - str module; - str query; - al_str_from(&module, ""); - al_str_from(&query, ""); - if (al_str_cmp(local, al_str_c("https://twitter.com"), 0, 19) == 0 || - al_str_cmp(local, al_str_c("https://x.com"), 0, 13) == 0) { - al_str_cat(&query, al_str_c("tweet:")); - al_str_cat(&query, local); - al_str_cat(&module, al_str_c("twitter")); - } else if (al_str_cmp(local, al_str_c("https://instagram.com"), 0, 21) == 0) { - s32 slash = al_str_rfind(local, '/'); - if (slash >= 0) { - al_str_cat(&query, al_str_substr(local, slash + 1, local->len)); - } - al_str_cat(&module, al_str_c("instagram")); - } else { - if (search) { - al_str_cat(&query, al_str_substr(local, 1, local->len)); - } else { - al_str_cat(&query, al_str_c("link:")); - al_str_cat(&query, local); - } - al_str_cat(&module, al_str_c("youtube")); - } - s32 id = camu_portal_create_search(bridge, &module, &query); - al_str_free(&module); - al_str_free(&query); - if (id < 0) return false; - struct camu_search *search = camu_portal_get_search(bridge, id); - if (!search || !camu_search_get_page(search, 0)) { - al_log_info("local_compat", "Search failed."); - return false; - } - struct camu_result_page *page = &al_array_at(search->pages, 0); - struct camu_post *post; - al_array_foreach_ptr(page->posts, i, post) { - struct camu_post_media *media; - al_array_foreach_ptr(post->media, j, media) { - if (media->url.len > 0) { - al_str_clone(uri, &media->url); - *selected = post; - break; - } - } - if (*selected) break; - } - camu_portal_discard_search(bridge, id); - } - return *selected != NULL; -} -#endif - -static void worker_signal_callback(void *userdata) -{ - struct camu_local_compat *compat = (struct camu_local_compat *)userdata; - u32 size; - struct camu_server_resource *resource; - do { - camu_queue_try_pop(compat->pending, size, resource); - if (size == 0) break; - struct cch_handler *handler = resource->entry->handler; - if (al_str_eq(&handler->liana, al_str_c("codec"))) { - handler->maybe_spawn_worker(handler, 0); - } - } while (1); -} - -static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) -{ - struct camu_local_compat *compat = (struct camu_local_compat *)userdata; - aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS); -#ifdef CAMU_HAVE_PORTAL - bool have_python = false; -#endif - u32 count; - while (aki_packet_cache_wait(&compat->queue, &count)) { - struct aki_packet *packet = aki_packet_cache_pop(&compat->queue); - aki_packet_cache_unlock(&compat->queue); - if (!packet) { - break; - } - str local; - aki_packet_read_str(packet, &local); - struct cch_entry *entry = NULL; - struct camu_post *post = NULL; - if (al_str_cmp(&local, al_str_c("cdda://"), 0, 7) == 0) { -#ifdef LIANA_HAVE_CDIO - entry = cch_handler_cdio_create(); - struct cch_chapter *chapter = &al_array_at(entry->chapters, 0); - if (local.len > 7) { - s64 index = al_str_to_long(al_str_substr(&local, 7, local.len), 10); - if (index != INT64_MIN && index != INT64_MAX && index > 0 && index <= entry->chapters.size) { - chapter = &al_array_at(entry->chapters, index - 1); - } - } - entry->chapter = chapter; - entry->handler->maybe_spawn_worker(entry->handler, chapter->start); -#endif - } else { -#ifdef CAMU_HAVE_PORTAL -#ifndef AKIYO_HAS_CURL -#error "Curl required to use portal" -#endif - bool is_search = al_str_at(&local, 0) == ';'; - if (is_search || camu_is_url(&local, 0)) { - if (!have_python) { - // Defer python init. - have_python = camu_python_init(); - } - if (have_python) { - str uri; - if (uri_for_local(&local, is_search, &compat->bridge, &uri, &post)) { - entry = cch_handler_http_create(&uri, compat->server->loop); - } - } - } else { - entry = cch_handler_file_create(&local); - } -#else -#ifdef AKIYO_HAS_CURL - if (camu_is_url(&local, 0)) { - entry = cch_handler_http_create(&local, compat->server->loop); - } else { -#endif - entry = cch_handler_file_create(&local); -#ifdef AKIYO_HAS_CURL - } -#endif -#endif - } - if (entry) { - struct camu_server_resource *resource = al_alloc_object(struct camu_server_resource); - al_str_clone(&resource->unique_id, &local); - resource->post = post; - resource->entry = entry; - resource->node = lia_server_create_node(&compat->server->data.server, entry); - camu_queue_push(compat->pending, resource); - aki_signal_send(&compat->worker_signal); - resource->duration = lia_node_get_duration(resource->node); - al_array_push(compat->server->data.resources, resource); - camu_queue_push(compat->results, resource); - aki_signal_send(&compat->result_signal); - } else { - al_log_info("local_compat", "No resource could be created for: %.*s.", AL_STR_PRINTF(&local)); - } - aki_packet_free(packet); - } -#ifdef CAMU_HAVE_PORTAL - if (have_python) camu_python_close(); -#endif - return 0; -} - -static void result_signal_callback(void *userdata) -{ - struct camu_local_compat *compat = (struct camu_local_compat *)userdata; - u32 size; - struct camu_server_resource *resource; - do { - camu_queue_try_pop(compat->results, size, resource); - if (size == 0) break; - compat->callback(compat->userdata, resource); - } while (1); -} - -void camu_local_compat_run(struct camu_local_compat *compat, struct camu_server *server) -{ - compat->server = server; -#ifdef CAMU_HAVE_PORTAL - camu_post_cache_init(&compat->cache); - camu_portal_init(&compat->bridge, &compat->cache); -#endif - // Size 0 to flush on the first packet. - aki_packet_cache_init(&compat->queue, 0); - aki_signal_init(&compat->worker_signal, worker_signal_callback, compat); - aki_signal_start(&compat->worker_signal, compat->server->loop); - camu_queue_init(compat->pending); - aki_signal_init(&compat->result_signal, result_signal_callback, compat); - aki_signal_start(&compat->result_signal, compat->server->loop); - camu_queue_init(compat->results); - aki_thread_create(&compat->thread, queue_thread, compat); -} - -void camu_local_compat_send(struct camu_local_compat *compat, struct aki_packet *packet) -{ - aki_packet_cache_send_packet(&compat->queue, packet); -} - -void camu_local_compat_stop(struct camu_local_compat *compat) -{ - aki_packet_cache_disable(&compat->queue); - aki_thread_cancel(&compat->thread); - aki_thread_join(&compat->thread); - aki_signal_stop(&compat->worker_signal); - aki_signal_stop(&compat->result_signal); - camu_queue_free(compat->results); - struct aki_packet *packet; - while ((packet = aki_packet_cache_pop(&compat->queue))) { - aki_packet_free(packet); - } -} diff --git a/src/server/local_compat.h b/src/server/local_compat.h deleted file mode 100644 index dcd2fc7..0000000 --- a/src/server/local_compat.h +++ /dev/null @@ -1,31 +0,0 @@ -#pragma once - -#include <al/str.h> -#include <aki/packet_cache.h> -#include <aki/signal.h> - -#ifdef CAMU_HAVE_PORTAL -#include "../portal/src/search.h" -#endif -#include "../util/queue.h" - -struct camu_server_resource; -struct camu_local_compat { - struct aki_thread thread; -#ifdef CAMU_HAVE_PORTAL - struct camu_portal_bridge bridge; - struct camu_post_cache cache; -#endif - struct aki_packet_cache queue; - struct aki_signal worker_signal; - queue(struct camu_server_resource *) pending; - struct aki_signal result_signal; - queue(struct camu_server_resource *) results; - struct camu_server *server; - void (*callback)(void *, struct camu_server_resource *); - void *userdata; -}; - -void camu_local_compat_run(struct camu_local_compat *compat, struct camu_server *server); -void camu_local_compat_send(struct camu_local_compat *compat, struct aki_packet *packet); -void camu_local_compat_stop(struct camu_local_compat *compat); diff --git a/src/server/meson.build b/src/server/meson.build index 0617c2c..88f80c0 100644 --- a/src/server/meson.build +++ b/src/server/meson.build @@ -1,10 +1,8 @@ server_src = [ 'server.c', 'common.c', - 'list.c', 'user.c', 'db.c', - 'local_compat.c' ] server_deps = [common_deps, cache, liana_server] server = declare_dependency(sources: server_src, dependencies: server_deps) diff --git a/src/server/resource.h b/src/server/resource.h index c1cf619..87250e9 100644 --- a/src/server/resource.h +++ b/src/server/resource.h @@ -2,12 +2,28 @@ #include "../cache/entry.h" #include "../liana/server.h" -#include "../portal/src/post.h" -struct camu_server_resource { - str unique_id; - struct camu_post *post; +enum { + CAMU_RESOURCE_NOT_LOADED = 0, + CAMU_RESOURCE_LOADING, + CAMU_RESOURCE_LOADED +}; + +struct camu_resource { + u8 type; + u8 load; struct cch_entry *entry; struct lia_node *node; u64 duration; + array(struct lia_list_entry *) pending; +}; + +struct camu_resource_file { + struct camu_resource r; + str path; +}; + +struct camu_resource_portal { + struct camu_resource r; + struct camu_post *post; }; diff --git a/src/server/server.c b/src/server/server.c index d7f8bc9..d3050ca 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -1,9 +1,16 @@ #include <al/log.h> #include <al/lib.h> +#include "../cache/handlers/file.h" +#include "../cache/handlers/http.h" +#include "../libclient/common.h" +#include "../libsink/common.h" +#ifdef CAMU_HAVE_PORTAL +#include "../portal/src/packet_ext.h" +#endif + #include "server.h" #include "common.h" -#include "list.h" #include "db.h" static struct camu_user *get_user_by_username(struct camu_server *server, str *username) @@ -26,9 +33,9 @@ static struct camu_server_client *get_client_by_connection(struct camu_server *s return NULL; } -static struct camu_list *get_list_from_name(struct camu_server *server, str *name) +static struct lia_list *get_list_from_name(struct camu_server *server, str *name) { - struct camu_list *list; + struct lia_list *list; al_array_foreach(server->lists, i, list) { if (al_str_eq(&list->name, name)) return list; } @@ -44,7 +51,9 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s return NULL; } -static bool identify_command_callback(void *userdata, struct aki_rpc_connection *conn, +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; @@ -83,8 +92,7 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection al_str_clone(&sink->name, &name); sink->server = server; al_array_push(server->sinks, sink); - struct camu_list *list = al_array_at(server->lists, 0); - lia_list_add_sink(&list->impl, camu_list_sink_callback, sink); + handle_toggle_sink(server, al_str_c("default"), sink, true); al_log_info("server", "New sink."); break; } @@ -94,7 +102,93 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection return true; } -static bool client_command_command_callback(void *userdata, struct aki_rpc_connection *conn, +static void client_portal_callback(void *userdata, struct camu_portal_result *result) +{ + 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); + } + break; + } + } + aki_rpc_connection_command(conn, packet, NULL, NULL); +} + +static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) +{ + struct camu_server_sink *sink = (struct camu_server_sink *)userdata; + switch (op) { + case LIANA_SINK_SET: + case LIANA_SINK_BUFFER: + case LIANA_SINK_BUFFER_AND_QUEUE: { + struct camu_resource *resource = (struct camu_resource *)entry->opaque; + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + aki_packet_write_u8(packet, op); + aki_packet_write_str(packet, &sink->server->addr); + aki_packet_write_u16(packet, CAMU_PORT); + aki_packet_write_u16(packet, resource->node->id); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u64(packet, timing->seek_pos); + aki_packet_write_u8(packet, timing->pause); + aki_packet_write_bool(packet, timing->ended); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_UNSET: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + aki_packet_write_u8(packet, op); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_PAUSE: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u8(packet, timing->pause); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_SEEK: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u64(packet, timing->seek_pos); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + } +} + +void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable) +{ + 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); + } +} + +static bool client_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; @@ -108,18 +202,33 @@ static bool client_command_command_callback(void *userdata, struct aki_rpc_conne case CAMU_CLIENT_CREATE_LIST: { str name; aki_packet_read_str(packet, &name); - struct camu_list *list = al_alloc_object(struct camu_list); - camu_list_init(list, &name); + struct lia_list *list = al_alloc_object(struct lia_list); + lia_list_init(list, &name); al_array_push(server->lists, list); break; } - case CAMU_CLIENT_ENABLE_SINK: { + case CAMU_CLIENT_TOGGLE_SINK: { str name; aki_packet_read_str(packet, &name); - struct camu_list *list = get_list_from_name(server, &name); - aki_packet_read_str(packet, &name); struct camu_server_sink *sink = get_sink_from_name(server, &name); - lia_list_add_sink(&list->impl, camu_list_sink_callback, sink); + if (!sink) goto out; + aki_packet_read_str(packet, &name); // list name. + bool enable = aki_packet_read_bool(packet); + handle_toggle_sink(server, &name, sink, enable); + break; + } + case CAMU_CLIENT_CREATE_SEARCH: { + str module; + aki_packet_read_str(packet, &module); + str query; + aki_packet_read_str(packet, &query); + camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn); + break; + } + case CAMU_CLIENT_GET_PAGE: { + s32 id = aki_packet_read_s32(packet); + u32 num = aki_packet_read_u32(packet); + camu_portal_get_page(&server->bridge, id, num, client_portal_callback, conn); break; } } @@ -129,7 +238,116 @@ out: return true; } -static bool list_action_command_callback(void *userdata, struct aki_rpc_connection *conn, +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->duration = duration; + struct lia_list_entry *entry; + al_array_foreach(resource->pending, i, entry) { + lia_list_pump(entry->list); + } + resource->pending.size = 0; + 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) +{ + struct camu_server *server = (struct camu_server *)userdata; + (void)server; + struct camu_resource *resource = (struct camu_resource *)entry->opaque; + switch (op) { + case LIANA_LOAD_ENTRY: + switch (resource->load) { + case CAMU_RESOURCE_NOT_LOADED: + resource->load = CAMU_RESOURCE_LOADING; + lia_node_get_duration(resource->node); + // fallthrough + case CAMU_RESOURCE_LOADING: + maybe_add_to_pending(resource, entry); + *(bool *)result = false; + break; + case CAMU_RESOURCE_LOADED: + *(bool *)result = true; + break; + } + break; + case LIANA_GET_DURATION: { + *(u64 *)result = resource->duration; + break; + } + case LIANA_UNLOAD_ENTRY: + break; + } +} + +static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet) +{ + u8 op = aki_packet_read_u8(packet); + switch (op) { + case CAMU_RESOURCE_FILE: { + str path; + aki_packet_read_str(packet, &path); + struct camu_resource_file *resource = al_alloc_object(struct camu_resource_file); + resource->r.type = CAMU_RESOURCE_FILE; + resource->r.load = CAMU_RESOURCE_NOT_LOADED; + al_str_clone(&resource->path, &path); + resource->r.entry = cch_handler_file_create(&path); + resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry); + resource->r.node->callback = node_callback; + resource->r.node->userdata = (struct camu_resource *)resource; + resource->r.duration = LIANA_TIMESTAMP_INVALID; + al_array_init(resource->r.pending); + lia_list_add(list, resource, resource->r.duration, &path); + break; + } + case CAMU_RESOURCE_PORTAL: { + str unique_id; + 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); + 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(&unique_id), index); + return; + } + struct camu_resource_portal *resource = al_alloc_object(struct camu_resource_portal); + resource->r.type = CAMU_RESOURCE_PORTAL; + resource->r.load = CAMU_RESOURCE_NOT_LOADED; + resource->post = post; + resource->r.entry = entry; + struct cch_handler *handler = resource->r.entry->handler; + handler->maybe_spawn_worker(handler, 0); + resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry); + resource->r.node->callback = node_callback; + resource->r.node->userdata = (struct camu_resource *)resource; + resource->r.duration = LIANA_TIMESTAMP_INVALID; + al_array_init(resource->r.pending); + lia_list_add(list, resource, resource->r.duration, al_str_c("dfdd")); + break; + } + } +} + +static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { struct camu_server *server = (struct camu_server *)userdata; @@ -139,50 +357,50 @@ static bool list_action_command_callback(void *userdata, struct aki_rpc_connecti str name; aki_packet_read_str(packet, &name); - struct camu_list *list = get_list_from_name(server, &name); + struct lia_list *list = get_list_from_name(server, &name); if (!list) goto out; u8 op = aki_packet_read_u8(packet); switch (op) { case CAMU_LIST_ADD: { - camu_local_compat_send(&server->compat, packet); + handle_add_command(server, list, packet); return false; } case CAMU_LIST_SKIP: { s32 sequence = aki_packet_read_s32(packet); s32 n = aki_packet_read_s32(packet); - lia_list_skip(&list->impl, sequence, n); + lia_list_skip(list, sequence, n); break; } case CAMU_LIST_SKIPTO: { s32 sequence = aki_packet_read_s32(packet); s32 i = aki_packet_read_s32(packet); - lia_list_skipto(&list->impl, sequence, i); + lia_list_skipto(list, sequence, i); break; } case CAMU_LIST_SHUFFLE: { - lia_list_shuffle(&list->impl); + lia_list_shuffle(list); break; } case CAMU_LIST_TOGGLE_PAUSE: { s32 sequence = aki_packet_read_s32(packet); f64 pts = aki_packet_read_f64(packet); - lia_list_toggle_pause(&list->impl, sequence, pts); + lia_list_toggle_pause(list, sequence, pts); break; } case CAMU_LIST_SEEK: { s32 sequence = aki_packet_read_s32(packet); f64 percent = aki_packet_read_f64(packet); - lia_list_seek(&list->impl, sequence, percent); + lia_list_seek(list, sequence, percent); break; } case CAMU_LIST_UNSET: { - lia_list_unset(&list->impl); + lia_list_unset(list); break; } case CAMU_LIST_END: { s32 sequence = aki_packet_read_s32(packet); - lia_list_end(&list->impl, sequence); + lia_list_end(list, sequence); break; } } @@ -193,9 +411,9 @@ out: } static struct aki_rpc_command commands[] = { - { .op = CAMU_SERVER_IDENTIFY, .callback = identify_command_callback, .userdata = NULL }, - { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_command_callback, .userdata = NULL }, - { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_command_callback, .userdata = NULL } + { .op = CAMU_SERVER_IDENTIFY, .callback = identify_callback, .userdata = NULL }, + { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_callback, .userdata = NULL }, + { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL } }; static void connection_callback(void *userdata, struct aki_rpc_connection *conn) @@ -228,9 +446,9 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection struct camu_server_node *node; al_array_foreach(server->nodes, i, node) { if (node->conn == conn) { + al_log_info("server", "Node removed."); cleanup_node(node); al_array_remove_at(server->nodes, i); - al_log_info("server", "Node removed."); break; } } @@ -238,9 +456,9 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection struct camu_server_client *client; al_array_foreach(server->clients, i, client) { if (client->conn == conn) { + al_log_info("server", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->name)); cleanup_client(client); al_array_remove_at(server->clients, i); - al_log_info("server", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->name)); break; } } @@ -248,13 +466,13 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection struct camu_server_sink *sink; al_array_foreach(server->sinks, i, sink) { if (sink->conn == conn) { + al_log_info("server", "Sink removed."); al_array_remove_at(server->sinks, i); - struct camu_list *list; + struct lia_list *list; al_array_foreach(server->lists, j, list) { - lia_list_remove_sink(&list->impl, sink); + lia_list_remove_sink(list, sink); } cleanup_sink(sink); - al_log_info("server", "Sink removed."); break; } } @@ -276,13 +494,6 @@ static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *sock) return false; } -static void local_compat_callback(void *userdata, struct camu_server_resource *resource) -{ - struct camu_server *server = (struct camu_server *)userdata; - struct camu_list *list = al_array_last(server->lists); - if (resource) lia_list_add(&list->impl, resource, resource->duration, &resource->unique_id); -} - bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop) { server->loop = loop; @@ -293,8 +504,10 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop al_array_init(server->users); al_array_init(server->lists); - struct camu_list *list = al_alloc_object(struct camu_list); - camu_list_init(list, al_str_c("default")); + struct lia_list *list = al_alloc_object(struct lia_list); + lia_list_init(list, al_str_c("default")); + list->callback = list_callback; + list->userdata = server; al_array_push(server->lists, list); aki_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); @@ -305,9 +518,10 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop lia_server_init(&server->data.server, server->loop); - server->compat.callback = local_compat_callback; - server->compat.userdata = server; - camu_local_compat_run(&server->compat, server); +#ifdef CAMU_HAVE_PORTAL + camu_post_cache_init(&server->cache); + camu_portal_init(&server->bridge, &server->cache, server->loop); +#endif return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server); } @@ -321,7 +535,9 @@ bool camu_server_listen(struct camu_server *server, str *addr, u16 port) void camu_server_close(struct camu_server *server) { - camu_local_compat_stop(&server->compat); +#ifdef CAMU_HAVE_PORTAL + camu_portal_close(&server->bridge); +#endif lia_server_close(&server->data.server); aki_multiplex_socket_close(&server->multi); } @@ -330,10 +546,10 @@ void camu_server_free(struct camu_server *server) { // This needs to happen in flight. // Too many things can block us before we get here. - struct camu_server_resource *resource; - al_array_foreach(server->data.resources, i, resource) { - cch_entry_free(&resource->entry); - } + //struct camu_server_resource *resource; + //al_array_foreach(server->data.resources, i, resource) { + // cch_entry_free(&resource->entry); + //} al_array_free(server->data.resources); lia_server_free(&server->data.server); aki_rpc_free(&server->server); @@ -342,5 +558,5 @@ void camu_server_free(struct camu_server *server) void camu_server_local_add(struct camu_server *server, struct aki_packet *packet) { - camu_local_compat_send(&server->compat, packet); + handle_add_command(server, al_array_last(server->lists), packet); } diff --git a/src/server/server.h b/src/server/server.h index 6246f79..aaabc3c 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -3,12 +3,14 @@ #include <aki/multiplex.h> #include <aki/rpc2.h> -#include "../cache/entry.h" #include "../liana/server.h" +#include "../liana/list.h" +#ifdef CAMU_HAVE_PORTAL +#include "../portal/src/search.h" +#endif #include "user.h" #include "resource.h" -#include "local_compat.h" struct camu_server_node { struct aki_rpc_connection *conn; @@ -34,12 +36,15 @@ struct camu_server { array(struct camu_server_client *) clients; array(struct camu_server_sink *) sinks; array(struct camu_user *) users; - array(struct camu_list *) lists; - struct camu_local_compat compat; + array(struct lia_list *) lists; struct { struct lia_server server; - array(struct camu_server_resource *) resources; + array(struct camu_resource *) resources; } data; +#ifdef CAMU_HAVE_PORTAL + struct camu_portal_bridge bridge; + struct camu_post_cache cache; +#endif }; bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop); diff --git a/src/sink/meson.build b/src/sink/meson.build index df80934..bd7b3f0 100644 --- a/src/sink/meson.build +++ b/src/sink/meson.build @@ -1,3 +1,5 @@ desktop_src = ['desktop.c'] desktop_deps = [common_deps, buffer, render, screen, mixer, libsink] -desktop = declare_dependency(sources: desktop_src, dependencies: desktop_deps) +desktop_args = ['-DCAMU_MIXER_THREADED', '-DCAMU_SCREEN_THREADED'] +desktop = declare_dependency(sources: desktop_src, dependencies: desktop_deps, + compile_args: desktop_args) |