summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-04-09 11:24:01 -0400
committerAndrew Opalach <andrew@akon.city> 2024-04-09 11:24:01 -0400
commit02f3d3565602146bbbfce85b2719246f24036cb9 (patch)
treec6588ffe297b777e36260effa4fa42b958ca6ba3 /src/server
parentbbf3314165182e402ff25acccddc004a87f81ef0 (diff)
downloadcamu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.gz
camu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.bz2
camu-02f3d3565602146bbbfce85b2719246f24036cb9.zip
Massive restructure and many changes
- The server-side list concept is still a wip Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/server')
-rw-r--r--src/server/common.h42
-rw-r--r--src/server/db.c56
-rw-r--r--src/server/db.h9
-rw-r--r--src/server/local_compat.h68
-rw-r--r--src/server/meson.build6
-rw-r--r--src/server/server.c308
-rw-r--r--src/server/server.h54
7 files changed, 543 insertions, 0 deletions
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 <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")
+
+#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 <al/log.h>
+#include <jansson.h>
+
+#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 <al/str.h>
+
+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 <al/log.h>
+
+#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 <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"
+
+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
+};