diff options
| author | 2024-10-21 19:22:50 -0400 | |
|---|---|---|
| committer | 2024-10-21 19:22:50 -0400 | |
| commit | 60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2 (patch) | |
| tree | 08ff2ce7975f523112e7ad2fe4f797b4fc7db5de /src/server | |
| parent | 2f9a0945bfeee3296cec3d38d094e4c49f9cb65f (diff) | |
| download | camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.gz camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.bz2 camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.zip | |
Everything before initial synced list
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/server')
| -rw-r--r-- | src/server/common.c | 10 | ||||
| -rw-r--r-- | src/server/common.h | 47 | ||||
| -rw-r--r-- | src/server/db.c | 31 | ||||
| -rw-r--r-- | src/server/list.c | 13 | ||||
| -rw-r--r-- | src/server/list.h | 11 | ||||
| -rw-r--r-- | src/server/local_compat.c | 213 | ||||
| -rw-r--r-- | src/server/local_compat.h | 87 | ||||
| -rw-r--r-- | src/server/meson.build | 10 | ||||
| -rw-r--r-- | src/server/server.c | 433 | ||||
| -rw-r--r-- | src/server/server.h | 61 | ||||
| -rw-r--r-- | src/server/user.c | 0 | ||||
| -rw-r--r-- | src/server/user.h | 8 |
12 files changed, 612 insertions, 312 deletions
diff --git a/src/server/common.c b/src/server/common.c new file mode 100644 index 0000000..765418c --- /dev/null +++ b/src/server/common.c @@ -0,0 +1,10 @@ +#include "common.h" + +//const str *CAMU_DB_PATH = al_str_c("/mnt/store/files/camu_db"); +str *CAMU_DB_PATH = al_str_c("/home/andrew/c/camu/data/camu_db_test"); + +str *CAMU_SERVER_IP = al_str_c("108.52.160.112"); +//str *CAMU_SERVER_IP = al_str_c("127.0.0.1"); + +str *CAMU_UNIX_PATH = al_str_c("/tmp/camu_sock"); +str *CAMU_UNIX_LOCAL = al_str_c("/tmp/cmv_sock"); diff --git a/src/server/common.h b/src/server/common.h index 26d0edf..c31efcf 100644 --- a/src/server/common.h +++ b/src/server/common.h @@ -2,16 +2,20 @@ #include <al/str.h> -#define CAMU_DB_PATH al_str_c("/mnt/store/files/camu_db") -//#define CAMU_DB_PATH al_str_c("/home/andrew/c/camu/data/camu_db_test") +extern str *CAMU_DB_PATH; #define CAMU_PORT 14356 #define CAMU_MULTIPLEX_RPC 0x53 -#define CAMU_MULTIPLEX_SHRUB 0x85 +#define CAMU_MULTIPLEX_LIANA 0x85 -//#define CAMU_SERVER_IP al_str_c("127.0.0.1") -//#define CAMU_SERVER_IP al_str_c("192.168.1.192") -#define CAMU_SERVER_IP al_str_c("108.52.160.112") +extern str *CAMU_SERVER_IP; +extern str *CAMU_UNIX_PATH; +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 enum { CAMU_NODE = 0, @@ -21,14 +25,33 @@ enum { enum { CAMU_SERVER_IDENTIFY = 0, + CAMU_SERVER_CLIENT_COMMAND, CAMU_SERVER_LIST_ACTION }; enum { - CAMU_ADD = 0, - CAMU_SKIP, - CAMU_SKIPTO, - CAMU_TOGGLE_PAUSE, - CAMU_SEEK, - CAMU_FINISHED + CAMU_CLIENT_CREATE_LIST = 0, + CAMU_CLIENT_ENABLE_SINK +}; + +enum { + CAMU_LIST_ADD = 0, + CAMU_LIST_SKIP, + CAMU_LIST_SKIPTO, + CAMU_LIST_SHUFFLE, + CAMU_LIST_TOGGLE_PAUSE, + CAMU_LIST_SEEK, + CAMU_LIST_UNSET, + CAMU_LIST_END }; + +AL_UNUSED_FUNCTION_PUSH + +static bool camu_is_url(str *s, u32 i) +{ + return al_str_cmp(s, al_str_c("https://"), i, 8) == 0 + || al_str_cmp(s, al_str_c("http://"), i, 7) == 0 + || al_str_cmp(s, al_str_c("cdda://"), i, 7) == 0; +} + +AL_UNUSED_FUNCTION_POP diff --git a/src/server/db.c b/src/server/db.c index 5042f2d..5c5e8fd 100644 --- a/src/server/db.c +++ b/src/server/db.c @@ -1,9 +1,11 @@ #include <al/log.h> +#include <aki/file.h> #include <jansson.h> +#include "list.h" #include "server.h" -static bool open_user(struct camu_server *srv, struct aki_dir_entry *dir) +static bool open_user(struct camu_server *server, struct aki_dir_entry *dir) { struct aki_file file; if (!aki_file_open(&file, &dir->path, 0)) { @@ -20,16 +22,14 @@ static bool open_user(struct camu_server *srv, struct aki_dir_entry *dir) } struct camu_user *user = al_alloc_object(struct camu_user); al_str_from(&user->name, json_string_value(json_object_get(root, "username"))); - al_array_init(user->lists); - struct camu_list *default_list = al_alloc_object(struct camu_list); - camu_list_init(default_list, al_str_c("default")); - al_array_push(user->lists, default_list); al_log_info("server", "Loaded user \"%.*s\"", AL_STR_PRINTF(&user->name)); - al_array_push(srv->users, user); + al_array_push(server->users, user); + json_decref(root); + al_str_free(&s); return true; } -bool camu_db_open(struct camu_server *srv, str *path) +bool camu_db_open(struct camu_server *server, str *path) { struct aki_dir camu_db; if (!aki_dir_open(&camu_db, path)) { @@ -43,8 +43,9 @@ bool camu_db_open(struct camu_server *srv, str *path) struct aki_dir_entry user; while (aki_dir_read(&users, &user)) { if (user.type == AKI_ENTRY_FILE) { - open_user(srv, &user); + open_user(server, &user); } + aki_dir_entry_free(&user); } aki_dir_close(&users); } @@ -54,3 +55,17 @@ bool camu_db_open(struct camu_server *srv, str *path) aki_dir_close(&camu_db); return true; } + +void camu_db_close(struct camu_server *server) +{ + struct camu_user *user; + al_array_foreach(server->users, i, user) { + al_str_free(&user->name); + } + al_array_free(server->users); + struct camu_list *list; + al_array_foreach(server->lists, i, list) { + camu_list_free(list); + } + al_array_free(server->lists); +} diff --git a/src/server/list.c b/src/server/list.c new file mode 100644 index 0000000..8ff123c --- /dev/null +++ b/src/server/list.c @@ -0,0 +1,13 @@ +#include "list.h" + +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 new file mode 100644 index 0000000..f3e52e6 --- /dev/null +++ b/src/server/list.h @@ -0,0 +1,11 @@ +#pragma once + +#include "../liana/list.h" + +struct camu_list { + str name; + struct lia_list impl; +}; + +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 new file mode 100644 index 0000000..9870ce5 --- /dev/null +++ b/src/server/local_compat.c @@ -0,0 +1,213 @@ +#include <al/log.h> + +#ifdef AKIYO_HAS_CURL +#include "../cache/handlers/http.h" +#endif +#include "../cache/handlers/file.h" +#include "../cache/handlers/cdio.h" + +#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; +#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); + if (!packet) { + aki_packet_cache_unlock(&compat->queue); + 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) { + 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); + } 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 + } + aki_packet_cache_unlock(&compat->queue); + 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_join(&compat->thread); + 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 index cdb88ce..dcd2fc7 100644 --- a/src/server/local_compat.h +++ b/src/server/local_compat.h @@ -1,68 +1,31 @@ #pragma once -#include "../portal/src/search.h" +#include <al/str.h> +#include <aki/packet_cache.h> +#include <aki/signal.h> -#ifdef AKIYO_HAS_CURL -#include "../cache/handlers/http.h" +#ifdef CAMU_HAVE_PORTAL +#include "../portal/src/search.h" #endif -#include "../cache/handlers/file.h" +#include "../util/queue.h" -static struct cch_entry *entry_for_external_path(struct camu_portal *portal, str *path) -{ - struct cch_entry *entry = NULL; -#ifdef AKIYO_HAS_CURL - bool is_search = al_str_at(path, 0) == ';'; - if (is_search || al_str_cmp(path, al_str_c("https://"), 0, 8) == 0 || - al_str_cmp(path, al_str_c("http://"), 0, 7) == 0) { - str module; - str query; - al_str_from(&module, ""); - al_str_from(&query, ""); - if (al_str_cmp(path, al_str_c("https://twitter.com"), 0, 19) == 0 || - al_str_cmp(path, al_str_c("https://x.com"), 0, 13) == 0) { - al_str_cat(&query, al_str_c("tweet:")); - al_str_cat(&query, path); - al_str_cat(&module, al_str_c("twitter")); - } else if (al_str_cmp(path, al_str_c("https://instagram.com"), 0, 21) == 0) { - s32 slash = al_str_rfind(path, '/'); - if (slash >= 0) { - al_str_cat(&query, al_str_substr(path, slash + 1, path->len)); - } - al_str_cat(&module, al_str_c("instagram")); - } else { - if (is_search) { - al_str_cat(&query, al_str_substr(path, 1, path->len)); - } else { - al_str_cat(&query, al_str_c("link:")); - al_str_cat(&query, path); - } - al_str_cat(&module, al_str_c("youtube")); - } - s32 id = camu_portal_create_search(portal, &module, &query); - al_str_free(&module); - al_str_free(&query); - if (id < 0) return NULL; - struct camu_search *search = camu_portal_get_search(portal, id); - if (!search || !camu_search_get_page(search, 0)) { - return NULL; - } - 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) { - entry = cch_handler_http_create(&media->url); - break; - } - } - if (entry) break; - } - camu_portal_discard_search(portal, id); - if (entry) entry->handler->maybe_spawn_worker(entry->handler, 0); - } else // { +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 - entry = cch_handler_file_create(path); - // } - return entry; -} + 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 d10b131..0617c2c 100644 --- a/src/server/meson.build +++ b/src/server/meson.build @@ -1,6 +1,10 @@ server_src = [ 'server.c', - 'db.c' + 'common.c', + 'list.c', + 'user.c', + 'db.c', + 'local_compat.c' ] -server_deps = [common_deps, portal, cache, shrub_server, list] -executable('server', sources: server_src, dependencies: server_deps) +server_deps = [common_deps, cache, liana_server] +server = declare_dependency(sources: server_src, dependencies: server_deps) diff --git a/src/server/server.c b/src/server/server.c index 4e11d3e..6c991eb 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -1,18 +1,17 @@ #include <al/log.h> +#include <al/lib.h> #include "../libsink/common.h" #include "server.h" #include "common.h" +#include "list.h" #include "db.h" -#ifdef CAMU_LOCAL_SOCKET -#include "local_compat.h" -#endif -static struct camu_user *get_user_by_username(struct camu_server *tree, str *username) +static struct camu_user *get_user_by_username(struct camu_server *server, str *username) { struct camu_user *user; - al_array_foreach(tree->users, i, user) { + al_array_foreach(server->users, i, user) { if (al_str_eq(&user->name, username)) { return user; } @@ -20,278 +19,303 @@ static struct camu_user *get_user_by_username(struct camu_server *tree, str *use return NULL; } -static void list_callback(void *userdata, u8 op, void *opaque, void *prev_opaque, s32 sequence) +static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) { - struct camu_srv_sink *sink = (struct camu_srv_sink *)userdata; - struct camu_srv_resource *resource = (struct camu_srv_resource *)opaque; - struct camu_srv_resource *prev = (struct camu_srv_resource *)prev_opaque; - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, op); - aki_packet_write_str(packet, CAMU_SERVER_IP); - aki_packet_write_s32(packet, CAMU_PORT); - aki_packet_write_u16(packet, resource->node->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, resource->node->start); - aki_packet_write_u8(packet, resource->node->paused != SHRUB_NOT_PAUSED); - if (prev) { - aki_packet_write_u64(packet, prev->node->paused_at); - } else { - aki_packet_write_u64(packet, 0); + struct camu_server_sink *sink = (struct camu_server_sink *)userdata; + switch (op) { + case CAMU_SINK_CLEAR: { + 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 CAMU_SINK_SET: + case CAMU_SINK_BUFFER: + case CAMU_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_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case CAMU_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 CAMU_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; } - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - if (op == CAMU_SINK_SET) { - al_printf("Now Playing: %.*s\n", AL_STR_PRINTF(&resource->unique_id)); } } static bool identify_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { - struct camu_server *srv = (struct camu_server *)userdata; + struct camu_server *server = (struct camu_server *)userdata; (void)rpacket; - switch (aki_packet_read_u8(packet)) { + + u8 op = aki_packet_read_u8(packet); + switch (op) { case CAMU_NODE: { - struct camu_srv_node *node = al_alloc_object(struct camu_srv_node); + struct camu_server_node *node = al_alloc_object(struct camu_server_node); node->conn = conn; - al_array_push(srv->nodes, node); + al_array_push(server->nodes, node); al_log_info("server", "New node."); break; } case CAMU_CLIENT: { + struct camu_server_client *client = al_alloc_object(struct camu_server_client); + client->conn = conn; str username; aki_packet_read_str(packet, &username); - struct camu_user *user = get_user_by_username(srv, &username); - if (user) { - struct camu_srv_client *client = al_alloc_object(struct camu_srv_client); - client->conn = conn; - client->user = user; - al_array_push(srv->clients, client); - //send_current_state(tree, user, conn); - al_log_info("server", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->name)); - } else { - al_log_info("server", "User \"%.*s\" not found.", AL_STR_PRINTF(&user->name)); + struct camu_user *user = get_user_by_username(server, &username); + if (!user) { + user = al_alloc_object(struct camu_user); + al_str_clone(&user->name, &username); + al_array_push(server->users, user); } + client->user = user; + al_array_push(server->clients, client); + al_log_info("server", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->name)); break; } case CAMU_SINK: { - struct camu_srv_sink *sink = al_alloc_object(struct camu_srv_sink); + struct camu_server_sink *sink = al_alloc_object(struct camu_server_sink); sink->conn = conn; - al_array_push(srv->sinks, sink); - struct camu_user *user = al_array_at(srv->users, 0); - struct camu_list *list = al_array_at(user->lists, 0); - camu_list_add_sink(list, list_callback, sink); + str name; + aki_packet_read_str(packet, &name); + 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, list_sink_callback, sink); al_log_info("server", "New sink."); break; } } + aki_packet_free(packet); return true; } -static void list_entry_callback(void *entry, u8 op, void *opaque) +static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct aki_rpc_connection *conn) { - struct camu_srv_resource *resource = (struct camu_srv_resource *)entry; - switch (op) { - case CAMU_LIST_ENTRY_ID: - *(str **)opaque = &resource->unique_id; - break; - case CAMU_LIST_IMPULSE: - shrb_node_set_start(resource->node, *(u64 *)opaque); - break; - case CAMU_LIST_USER_PAUSE: - shrb_node_user_pause(resource->node, *(u64 *)opaque); - break; - case CAMU_LIST_USER_RESUME: { - struct shrb_resume_req *req = (struct shrb_resume_req *)opaque; - shrb_node_user_resume(resource->node, req); - break; + struct camu_server_client *client; + al_array_foreach(server->clients, i, client) { + if (client->conn == conn) return client; } - case CAMU_LIST_PAUSE: - shrb_node_pause(resource->node); - break; - case CAMU_LIST_SEEK: - shrb_node_seek(resource->node, *(u64 *)opaque); - break; + return NULL; +} + +static struct camu_list *get_list_from_name(struct camu_server *server, str *name) +{ + struct camu_list *list; + al_array_foreach(server->lists, i, list) { + if (al_str_eq(&list->name, name)) return list; } + return NULL; } -#ifdef CAMU_LOCAL_SOCKET -static void list_add_local(struct camu_server *srv, struct camu_list *list, str *line) +static struct camu_server_sink *get_sink_from_name(struct camu_server *server, str *name) { - struct cch_entry *entry = entry_for_external_path(&srv->portal, line); - if (entry) { - struct shrb_node *node = shrb_server_create_node(&srv->resource, 0, entry); - if (node) { - struct camu_srv_resource *resource = al_alloc_object(struct camu_srv_resource); - al_array_push(srv->resources, resource); - al_str_clone(&resource->unique_id, line); - resource->entry = entry; - resource->node = node; - camu_list_add(list, resource, list_entry_callback, resource->node->duration, false); - } + struct camu_server_sink *sink; + al_array_foreach(server->sinks, i, sink) { + if (al_str_eq(&sink->name, name)) return sink; + } + return NULL; +} + +static bool client_command_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; + (void)rpacket; + + struct camu_server_client *client = get_client_by_connection(server, conn); + if (!client) goto out; + + u8 op = aki_packet_read_u8(packet); + switch (op) { + 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); + al_array_push(server->lists, list); + break; + } + case CAMU_CLIENT_ENABLE_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, list_sink_callback, sink); + break; + } } + +out: + aki_packet_free(packet); + return true; } -#endif static bool list_action_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { - struct camu_server *srv = (struct camu_server *)userdata; + struct camu_server *server = (struct camu_server *)userdata; (void)conn; (void)rpacket; - struct camu_user *user = al_array_at(srv->users, 0); - struct camu_list *list = al_array_at(user->lists, 0); + str name; + aki_packet_read_str(packet, &name); - switch (aki_packet_read_u8(packet)) { - case CAMU_ADD: { - str line; - aki_packet_read_str(packet, &line); -#ifdef CAMU_LOCAL_SOCKET - list_add_local(srv, list, &line); -#endif - break; + struct camu_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); + return false; } - case CAMU_SKIP: { + case CAMU_LIST_SKIP: { s32 sequence = aki_packet_read_s32(packet); s32 n = aki_packet_read_s32(packet); - camu_list_skip(list, sequence, n); + lia_list_skip(&list->impl, sequence, n); break; } - case CAMU_SKIPTO: { + case CAMU_LIST_SKIPTO: { + s32 sequence = aki_packet_read_s32(packet); s32 i = aki_packet_read_s32(packet); - camu_list_skipto(list, i); + lia_list_skipto(&list->impl, sequence, i); + break; + } + case CAMU_LIST_SHUFFLE: { + lia_list_shuffle(&list->impl); break; } - case CAMU_TOGGLE_PAUSE: { + case CAMU_LIST_TOGGLE_PAUSE: { s32 sequence = aki_packet_read_s32(packet); - u64 pos = aki_packet_read_u64(packet); - camu_list_toggle_pause(list, sequence, pos); + f64 pts = aki_packet_read_f64(packet); + lia_list_toggle_pause(&list->impl, sequence, pts); break; } - case CAMU_SEEK: { + case CAMU_LIST_SEEK: { s32 sequence = aki_packet_read_s32(packet); f64 percent = aki_packet_read_f64(packet); - camu_list_seek(list, sequence, percent); + lia_list_seek(&list->impl, sequence, percent); + break; + } + case CAMU_LIST_UNSET: { + lia_list_unset(&list->impl); break; } - case CAMU_FINISHED: { + case CAMU_LIST_END: { s32 sequence = aki_packet_read_s32(packet); - camu_list_finished(list, sequence); + lia_list_end(&list->impl, sequence); break; } } +out: aki_packet_free(packet); return false; } 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 } }; static void connection_callback(void *userdata, struct aki_rpc_connection *conn) { + // TODO: Cleanup zombie connections. (void)userdata; (void)conn; } -static void cleanup_node(struct camu_srv_node *node) +static void cleanup_node(struct camu_server_node *node) { al_free(node); } -static void cleanup_client(struct camu_srv_client *client) +static void cleanup_client(struct camu_server_client *client) { al_free(client); } -static void cleanup_sink(struct camu_srv_sink *sink) +static void cleanup_sink(struct camu_server_sink *sink) { + al_str_free(&sink->name); al_free(sink); } static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { - struct camu_server *srv = (struct camu_server *)userdata; + struct camu_server *server = (struct camu_server *)userdata; - struct camu_srv_node *node; - al_array_foreach(srv->nodes, i, node) { + 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_iter(srv->nodes, i); + al_array_remove_at(server->nodes, i); + al_log_info("server", "Node removed."); break; } } - struct camu_srv_client *client; - al_array_foreach(srv->clients, i, client) { + 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_iter(srv->clients, i); + al_array_remove_at(server->clients, i); + al_log_info("server", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->name)); break; } } - struct camu_srv_sink *sink; - al_array_foreach(srv->sinks, i, sink) { + 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_iter(srv->sinks, i); - struct camu_user *user = al_array_at(srv->users, 0); - struct camu_list *list = al_array_at(user->lists, 0); - camu_list_remove_sink(list, sink); + al_array_remove_at(server->sinks, i); + struct camu_list *list; + al_array_foreach(server->lists, j, list) { + lia_list_remove_sink(&list->impl, sink); + } cleanup_sink(sink); + al_log_info("server", "Sink removed."); break; } } } - -#ifdef CAMU_LOCAL_SOCKET -static u8 server_line_callback(void *userdata, str *line) -{ - struct camu_server *srv = (struct camu_server *)userdata; - struct camu_user *user = al_array_at(srv->users, 0); - struct camu_list *list = al_array_at(user->lists, 0); - if (al_str_eq(line, al_str_c(";NEXT"))) { - camu_list_skip(list, CAMU_SEQUENCE_INVALID, 1); - } else if (al_str_eq(line, al_str_c(";PREV"))) { - camu_list_skip(list, CAMU_SEQUENCE_INVALID, -1); - } else if (al_str_eq(line, al_str_c(";SHUFFLE"))) { - camu_list_shuffle(list); - } else if (al_str_eq(line, al_str_c(";SORT"))) { - camu_list_sort(list); - } else if (al_str_eq(line, al_str_c(";REVERSE"))) { - camu_list_reverse(list); - } else if (al_str_eq(line, al_str_c(";CLEAR"))) { - } else { - list_add_local(srv, list, line); - } - return AKI_LINE_PROCESSOR_CONTINUE; -} -#endif - -static void sigint_handler(s32 signum) -{ - (void)signum; - // explode. - exit(EXIT_FAILURE); -} - -static struct camu_server srv = { 0 }; - -static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *s) +static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *sock) { - struct camu_server *srv = (struct camu_server *)userdata; + struct camu_server *server = (struct camu_server *)userdata; switch (id) { case CAMU_MULTIPLEX_RPC: - aki_rpc_add_socket(&srv->server, s); + aki_rpc_add_socket(&server->server, sock); return true; - case CAMU_MULTIPLEX_SHRUB: - shrb_server_add_socket(&srv->resource, s); + case CAMU_MULTIPLEX_LIANA: + lia_server_add_socket(&server->data.server, sock); return true; default: break; @@ -299,56 +323,71 @@ static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *s) return false; } -s32 main(void) +static void local_compat_callback(void *userdata, struct camu_server_resource *resource) { - aki_common_init(); - - signal(SIGINT, sigint_handler); - - al_array_init(srv.nodes); - al_array_init(srv.clients); - al_array_init(srv.sinks); - al_array_init(srv.users); - - if (!camu_db_open(&srv, CAMU_DB_PATH)) return EXIT_FAILURE; - - camu_post_cache_init(&srv.cache); - - bool py_init = camu_python_init(); + 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); +} - camu_portal_init(&srv.portal, &srv.cache); +bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop) +{ + server->loop = loop; + server->addr = al_str_zero(); + al_array_init(server->nodes); + al_array_init(server->clients); + al_array_init(server->sinks); + al_array_init(server->users); + al_array_init(server->lists); - aki_event_loop_init(&srv.loop); + struct camu_list *list = al_alloc_object(struct camu_list); + camu_list_init(list, al_str_c("default")); + al_array_push(server->lists, list); - aki_rpc_init(&srv.server, &srv.loop, connection_callback, connection_closed_callback, &srv); + aki_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server); for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { - commands[i].userdata = &srv; - aki_rpc_add_command(&srv.server, &commands[i]); + commands[i].userdata = server; + aki_rpc_add_command(&server->server, &commands[i]); } - shrb_server_init(&srv.resource, &srv.loop); + lia_server_init(&server->data.server, server->loop); - aki_multiplex_socket_init(&srv.multi, AKI_SOCKET_TCP, multiplex_callback, &srv); - aki_multiplex_socket_listen(&srv.multi, &srv.loop, al_str_c("0.0.0.0"), CAMU_PORT); + server->compat.callback = local_compat_callback; + server->compat.userdata = server; + camu_local_compat_run(&server->compat, server); -#ifdef CAMU_LOCAL_SOCKET - srv.socket.type = AKI_SOCKET_UNIX; - aki_socket_init(&srv.socket); - aki_socket_set_blocking(&srv.socket, false); - srv.pro.callback = server_line_callback; - srv.pro.userdata = &srv; - aki_line_processor_init(&srv.pro, al_str_c("\n")); - aki_line_processor_open_socket(&srv.pro, &srv.socket); - if (aki_socket_listen(&srv.socket, al_str_c("/tmp/camu_sock"), 0)) { - aki_line_processor_run(&srv.pro, &srv.loop); - } -#endif + return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server); +} - aki_event_loop_run(&srv.loop); +bool camu_server_listen(struct camu_server *server, str *addr, u16 port) +{ + al_str_clone(&server->addr, addr); + if (server->multi.sock.type == AKI_SOCKET_TCP) addr = NULL; // any + return aki_multiplex_socket_listen(&server->multi, server->loop, addr, port); +} - if (py_init) camu_python_close(); +void camu_server_close(struct camu_server *server) +{ + camu_local_compat_stop(&server->compat); + lia_server_close(&server->data.server); + aki_multiplex_socket_close(&server->multi); +} - aki_common_close(); +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); + } + al_array_free(server->data.resources); + lia_server_free(&server->data.server); + aki_rpc_free(&server->server); + al_str_free(&server->addr); +} - return EXIT_SUCCESS; +void camu_server_local_add(struct camu_server *server, struct aki_packet *packet) +{ + camu_local_compat_send(&server->compat, packet); } diff --git a/src/server/server.h b/src/server/server.h index 9565bf6..528f08a 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -1,56 +1,57 @@ #pragma once -#define CAMU_LOCAL_SOCKET - #include <aki/multiplex.h> #include <aki/rpc2.h> -#ifdef CAMU_LOCAL_SOCKET -#include <aki/line_processor.h> -#endif -#include "../portal/src/search.h" #include "../cache/entry.h" -#include "../shrub/server.h" -#include "../list/list.h" +#include "../liana/server.h" -struct camu_srv_node { - struct aki_rpc_connection *conn; -}; +#include "user.h" +#include "local_compat.h" -struct camu_user { - str name; - array(struct camu_list *) lists; +struct camu_server_node { + struct aki_rpc_connection *conn; }; -struct camu_srv_client { +struct camu_server_client { struct aki_rpc_connection *conn; struct camu_user *user; }; -struct camu_srv_sink { +struct camu_server_sink { struct aki_rpc_connection *conn; + str name; + struct camu_server *server; }; -struct camu_srv_resource { +struct camu_server_resource { str unique_id; + struct camu_post *post; struct cch_entry *entry; - struct shrb_node *node; + struct lia_node *node; + u64 duration; }; struct camu_server { - struct aki_event_loop loop; + struct aki_event_loop *loop; + str addr; struct aki_multiplex_socket multi; struct aki_rpc server; - array(struct camu_srv_node *) nodes; - array(struct camu_srv_client *) clients; - array(struct camu_srv_sink *) sinks; + array(struct camu_server_node *) nodes; + array(struct camu_server_client *) clients; + array(struct camu_server_sink *) sinks; array(struct camu_user *) users; - struct shrb_server resource; - array(struct camu_srv_resource *) resources; - struct camu_portal portal; - struct camu_post_cache cache; -#ifdef CAMU_LOCAL_SOCKET - struct aki_socket socket; - struct aki_line_processor pro; -#endif + array(struct camu_list *) lists; + struct camu_local_compat compat; + struct { + struct lia_server server; + array(struct camu_server_resource *) resources; + } data; }; + +bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop); +bool camu_server_listen(struct camu_server *server, str *addr, u16 port); +void camu_server_close(struct camu_server *server); +void camu_server_free(struct camu_server *server); + +void camu_server_local_add(struct camu_server *server, struct aki_packet *packet); diff --git a/src/server/user.c b/src/server/user.c new file mode 100644 index 0000000..e69de29 --- /dev/null +++ b/src/server/user.c diff --git a/src/server/user.h b/src/server/user.h new file mode 100644 index 0000000..9800c45 --- /dev/null +++ b/src/server/user.h @@ -0,0 +1,8 @@ +#pragma once + +#include <al/str.h> +#include <al/array.h> + +struct camu_user { + str name; +}; |