summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
Diffstat (limited to 'src/server')
-rw-r--r--src/server/common.c10
-rw-r--r--src/server/common.h47
-rw-r--r--src/server/db.c31
-rw-r--r--src/server/list.c13
-rw-r--r--src/server/list.h11
-rw-r--r--src/server/local_compat.c213
-rw-r--r--src/server/local_compat.h87
-rw-r--r--src/server/meson.build10
-rw-r--r--src/server/server.c433
-rw-r--r--src/server/server.h61
-rw-r--r--src/server/user.c0
-rw-r--r--src/server/user.h8
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;
+};