diff options
| author | 2024-11-08 14:53:40 -0500 | |
|---|---|---|
| committer | 2024-11-08 14:53:40 -0500 | |
| commit | 5e3641e5e692c3f2f644a4bb809c88727cb8bee9 (patch) | |
| tree | fbac2007fdcebda13406350fdcef977910427078 /src/server | |
| 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/server')
| -rw-r--r-- | src/server/common.h | 10 | ||||
| -rw-r--r-- | src/server/db.c | 5 | ||||
| -rw-r--r-- | src/server/list.c | 71 | ||||
| -rw-r--r-- | src/server/list.h | 16 | ||||
| -rw-r--r-- | src/server/local_compat.c | 219 | ||||
| -rw-r--r-- | src/server/local_compat.h | 31 | ||||
| -rw-r--r-- | src/server/meson.build | 2 | ||||
| -rw-r--r-- | src/server/resource.h | 24 | ||||
| -rw-r--r-- | src/server/server.c | 314 | ||||
| -rw-r--r-- | src/server/server.h | 15 |
10 files changed, 306 insertions, 401 deletions
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); |