From 02f3d3565602146bbbfce85b2719246f24036cb9 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Tue, 9 Apr 2024 11:24:01 -0400 Subject: Massive restructure and many changes - The server-side list concept is still a wip Signed-off-by: Andrew Opalach --- src/server/common.h | 42 +++++++ src/server/db.c | 56 +++++++++ src/server/db.h | 9 ++ src/server/local_compat.h | 68 ++++++++++ src/server/meson.build | 6 + src/server/server.c | 308 ++++++++++++++++++++++++++++++++++++++++++++++ src/server/server.h | 54 ++++++++ 7 files changed, 543 insertions(+) create mode 100644 src/server/common.h create mode 100644 src/server/db.c create mode 100644 src/server/db.h create mode 100644 src/server/local_compat.h create mode 100644 src/server/meson.build create mode 100644 src/server/server.c create mode 100644 src/server/server.h (limited to 'src/server') diff --git a/src/server/common.h b/src/server/common.h new file mode 100644 index 0000000..4d3a0ae --- /dev/null +++ b/src/server/common.h @@ -0,0 +1,42 @@ +#pragma once + +#include + +#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") + +#define CAMU_PORT 14356 +#define CAMU_RESOURCE_PORT 14357 + +#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") + + +enum { + CAMU_NODE = 0, + CAMU_CLIENT, + CAMU_SINK +}; + +enum { + CAMU_SRV_IDENTIFY = 0, + CAMU_SRV_LIST_ACTION, + CAMU_SRV_CREATE_SEARCH, + CAMU_SRV_GET_PAGE, + CAMU_SRV_CREATE_BROWSE, + CAMU_SRV_GET_PATH +}; + +enum { + CAMU_CONN_STATE +}; + +enum { + CAMU_LIST_ADD = 0, + CAMU_LIST_SKIP, + CAMU_LIST_SKIPTO, + CAMU_LIST_TOGGLE_PAUSE, + CAMU_LIST_SEEK, + CAMU_LIST_FINISHED +}; diff --git a/src/server/db.c b/src/server/db.c new file mode 100644 index 0000000..5042f2d --- /dev/null +++ b/src/server/db.c @@ -0,0 +1,56 @@ +#include +#include + +#include "server.h" + +static bool open_user(struct camu_server *srv, struct aki_dir_entry *dir) +{ + struct aki_file file; + if (!aki_file_open(&file, &dir->path, 0)) { + return false; + } + str s; + aki_file_read_as_str(&file, &s); + json_error_t error; + json_t *root = json_loadb(s.data, s.len, 0, &error); + if (!root) { + al_log_error("server", "Failed to parse json: %.*s:%d:%d (%s).", + AL_STR_PRINTF(&dir->path), error.line, error.column, error.text); + return false; + } + 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); + return true; +} + +bool camu_db_open(struct camu_server *srv, str *path) +{ + struct aki_dir camu_db; + if (!aki_dir_open(&camu_db, path)) { + return false; + } + struct aki_dir_entry entry; + while (aki_dir_read(&camu_db, &entry)) { + if (al_str_eq(&entry.name, al_str_c("users"))) { + struct aki_dir users; + if (aki_dir_open(&users, &entry.path)) { + struct aki_dir_entry user; + while (aki_dir_read(&users, &user)) { + if (user.type == AKI_ENTRY_FILE) { + open_user(srv, &user); + } + } + aki_dir_close(&users); + } + } + aki_dir_entry_free(&entry); + } + aki_dir_close(&camu_db); + return true; +} diff --git a/src/server/db.h b/src/server/db.h new file mode 100644 index 0000000..8e8b069 --- /dev/null +++ b/src/server/db.h @@ -0,0 +1,9 @@ +#pragma once + +#include + +struct camu_server; +struct camu_db { +}; + +bool camu_db_open(struct camu_server *srv, str *path); diff --git a/src/server/local_compat.h b/src/server/local_compat.h new file mode 100644 index 0000000..cdb88ce --- /dev/null +++ b/src/server/local_compat.h @@ -0,0 +1,68 @@ +#pragma once + +#include "../portal/src/search.h" + +#ifdef AKIYO_HAS_CURL +#include "../cache/handlers/http.h" +#endif +#include "../cache/handlers/file.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 // { +#endif + entry = cch_handler_file_create(path); + // } + return entry; +} diff --git a/src/server/meson.build b/src/server/meson.build new file mode 100644 index 0000000..d10b131 --- /dev/null +++ b/src/server/meson.build @@ -0,0 +1,6 @@ +server_src = [ + 'server.c', + 'db.c' +] +server_deps = [common_deps, portal, cache, shrub_server, list] +executable('server', sources: server_src, dependencies: server_deps) diff --git a/src/server/server.c b/src/server/server.c new file mode 100644 index 0000000..0fd1d71 --- /dev/null +++ b/src/server/server.c @@ -0,0 +1,308 @@ +#include + +#include "../libsink/common.h" +#include "../shrub/common.h" + +#include "server.h" +#include "common.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) +{ + struct camu_user *user; + al_array_foreach(tree->users, i, user) { + if (al_str_eq(&user->name, username)) { + return user; + } + } + return NULL; +} + +static void list_callback(void *userdata, u8 op, void *opaque, s32 sequence) +{ + struct camu_srv_sink *sink = (struct camu_srv_sink *)userdata; + struct camu_srv_resource *resource = (struct camu_srv_resource *)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_RESOURCE_PORT); + aki_packet_write_u16(packet, resource->node->id); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, resource->node->start); + 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; + (void)rpacket; + switch (aki_packet_read_u8(packet)) { + case CAMU_NODE: { + struct camu_srv_node *node = al_alloc_object(struct camu_srv_node); + node->conn = conn; + al_array_push(srv->nodes, node); + al_log_info("server", "New node."); + break; + } + case CAMU_CLIENT: { + 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)); + } + break; + } + case CAMU_SINK: { + struct camu_srv_sink *sink = al_alloc_object(struct camu_srv_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); + al_log_info("server", "New sink."); + break; + } + } + aki_packet_free(packet); + return true; +} + +static bool list_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; + (void)conn; + (void)rpacket; + + struct camu_user *user = al_array_at(srv->users, 0); + struct camu_list *list = al_array_at(user->lists, 0); + + switch (aki_packet_read_u8(packet)) { + case CAMU_LIST_ADD: { + break; + } + case CAMU_LIST_SKIP: { + s32 sequence = aki_packet_read_s32(packet); + s32 n = aki_packet_read_s32(packet); + camu_list_skip(list, sequence, n); + break; + } + case CAMU_LIST_SKIPTO: { + s32 i = aki_packet_read_s32(packet); + camu_list_skipto(list, i); + break; + } + 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); + break; + } + case CAMU_LIST_SEEK: { + s32 sequence = aki_packet_read_s32(packet); + f64 percent = aki_packet_read_f64(packet); + camu_list_seek(list, sequence, percent); + break; + } + case CAMU_LIST_FINISHED: { + s32 sequence = aki_packet_read_s32(packet); + camu_list_finished(list, sequence); + break; + } + } + + aki_packet_free(packet); + return false; +} + +static struct aki_rpc_command commands[] = { + { .op = CAMU_SRV_IDENTIFY, .callback = identify_command_callback, .userdata = NULL }, + { .op = CAMU_SRV_LIST_ACTION, .callback = list_command_callback, .userdata = NULL }, + //{ .op = CAMU_SRV_CREATE_SEARCH, .callback = create_search_command_callback, .userdata = NULL }, + //{ .op = CAMU_SRV_GET_PAGE, .callback = get_page_command_callback, .userdata = NULL }, +}; + +static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +{ + (void)userdata; + (void)conn; +} + +static void cleanup_node(struct camu_srv_node *node) +{ + al_free(node); +} + +static void cleanup_client(struct camu_srv_client *client) +{ + al_free(client); +} + +static void cleanup_sink(struct camu_srv_sink *sink) +{ + 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_srv_node *node; + al_array_foreach(srv->nodes, i, node) { + if (node->conn == conn) { + al_log_info("server", "Node removed."); + cleanup_node(node); + al_array_remove_at_iter(srv->nodes, i); + break; + } + } + + struct camu_srv_client *client; + al_array_foreach(srv->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); + break; + } + } + + struct camu_srv_sink *sink; + al_array_foreach(srv->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); + cleanup_sink(sink); + break; + } + } +} + +static void list_entry_callback(void *entry, u8 op, void *opaque) +{ + struct camu_srv_resource *resource = (struct camu_srv_resource *)entry; + switch (op) { + case CAMU_LIST_ENTRY_IMPULSE: { + struct camu_list_timing *timing = (struct camu_list_timing *)opaque; + shrb_node_set_start(resource->node, timing->start); + break; + } + case CAMU_LIST_ENTRY_TOGGLE_PAUSE: + shrb_node_toggle_pause(resource->node, *(u64 *)opaque); + break; + case CAMU_LIST_ENTRY_SEEK: + shrb_node_seek(resource->node, *(u64 *)opaque); + 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"))) { + } else if (al_str_eq(line, al_str_c(";CLEAR"))) { + } else { + 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); + } + } + } + return AKI_LINE_PROCESSOR_CONTINUE; +} +#endif + +static void sigint_handler(s32 signum) +{ + (void)signum; + // explode. + exit(EXIT_FAILURE); +} + +static struct camu_server srv = { 0 }; + +s32 main(void) +{ + 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(); + + camu_portal_init(&srv.portal, &srv.cache); + + aki_event_loop_init(&srv.loop); + + aki_rpc_init(&srv.server, AKI_SOCKET_TCP, connection_callback, connection_closed_callback, &srv); + for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { + commands[i].userdata = &srv; + aki_rpc_add_command(&srv.server, &commands[i]); + } + aki_rpc_listen(&srv.server, &srv.loop, al_str_c("0.0.0.0"), CAMU_PORT); + + shrb_server_init(&srv.resource); + shrb_server_listen(&srv.resource, &srv.loop, al_str_c("0.0.0.0"), CAMU_RESOURCE_PORT); + +#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 + + aki_event_loop_run(&srv.loop); + + if (py_init) camu_python_close(); + + aki_common_close(); + + return EXIT_SUCCESS; +} diff --git a/src/server/server.h b/src/server/server.h new file mode 100644 index 0000000..83e8523 --- /dev/null +++ b/src/server/server.h @@ -0,0 +1,54 @@ +#pragma once + +#define CAMU_LOCAL_SOCKET + +#include +#ifdef CAMU_LOCAL_SOCKET +#include +#endif + +#include "../portal/src/search.h" +#include "../cache/entry.h" +#include "../shrub/server.h" +#include "../list/list.h" + +struct camu_srv_node { + struct aki_rpc_connection *conn; +}; + +struct camu_user { + str name; + array(struct camu_list *) lists; +}; + +struct camu_srv_client { + struct aki_rpc_connection *conn; + struct camu_user *user; +}; + +struct camu_srv_sink { + struct aki_rpc_connection *conn; +}; + +struct camu_srv_resource { + str unique_id; + struct cch_entry *entry; + struct shrb_node *node; +}; + +struct camu_server { + struct aki_event_loop loop; + 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_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 +}; -- cgit v1.2.3-101-g0448