summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
Diffstat (limited to 'src/server')
-rw-r--r--src/server/common.h10
-rw-r--r--src/server/db.c5
-rw-r--r--src/server/list.c71
-rw-r--r--src/server/list.h16
-rw-r--r--src/server/local_compat.c219
-rw-r--r--src/server/local_compat.h31
-rw-r--r--src/server/meson.build2
-rw-r--r--src/server/resource.h24
-rw-r--r--src/server/server.c314
-rw-r--r--src/server/server.h15
10 files changed, 306 insertions, 401 deletions
diff --git a/src/server/common.h b/src/server/common.h
index c31efcf..bde8106 100644
--- a/src/server/common.h
+++ b/src/server/common.h
@@ -14,6 +14,7 @@ extern str *CAMU_UNIX_LOCAL;
//#define CAMU_LOCAL_TYPE AKI_SOCKET_TCP
//#define CAMU_LOCAL_ADDR CAMU_SERVER_IP
+
#define CAMU_LOCAL_TYPE AKI_SOCKET_UNIX
#define CAMU_LOCAL_ADDR CAMU_UNIX_PATH
@@ -31,7 +32,9 @@ enum {
enum {
CAMU_CLIENT_CREATE_LIST = 0,
- CAMU_CLIENT_ENABLE_SINK
+ CAMU_CLIENT_TOGGLE_SINK,
+ CAMU_CLIENT_CREATE_SEARCH,
+ CAMU_CLIENT_GET_PAGE
};
enum {
@@ -45,6 +48,11 @@ enum {
CAMU_LIST_END
};
+enum {
+ CAMU_RESOURCE_FILE = 0,
+ CAMU_RESOURCE_PORTAL
+};
+
AL_UNUSED_FUNCTION_PUSH
static bool camu_is_url(str *s, u32 i)
diff --git a/src/server/db.c b/src/server/db.c
index 5c5e8fd..c9608bb 100644
--- a/src/server/db.c
+++ b/src/server/db.c
@@ -2,7 +2,6 @@
#include <aki/file.h>
#include <jansson.h>
-#include "list.h"
#include "server.h"
static bool open_user(struct camu_server *server, struct aki_dir_entry *dir)
@@ -63,9 +62,9 @@ void camu_db_close(struct camu_server *server)
al_str_free(&user->name);
}
al_array_free(server->users);
- struct camu_list *list;
+ struct lia_list *list;
al_array_foreach(server->lists, i, list) {
- camu_list_free(list);
+ lia_list_free(list);
}
al_array_free(server->lists);
}
diff --git a/src/server/list.c b/src/server/list.c
deleted file mode 100644
index e3565d0..0000000
--- a/src/server/list.c
+++ /dev/null
@@ -1,71 +0,0 @@
-#include "../libsink/common.h"
-#include "../server/common.h"
-
-#include "list.h"
-#include "server.h"
-
-void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry)
-{
- struct camu_server *server = (struct camu_server *)userdata;
- (void)server;
- (void)op;
- (void)entry;
-}
-
-void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing)
-{
- struct camu_server_sink *sink = (struct camu_server_sink *)userdata;
- switch (op) {
- case LIANA_SINK_SET:
- case LIANA_SINK_BUFFER:
- case LIANA_SINK_BUFFER_AND_QUEUE: {
- struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque;
- struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET);
- aki_packet_write_u8(packet, op);
- aki_packet_write_str(packet, &sink->server->addr);
- aki_packet_write_u16(packet, CAMU_PORT);
- aki_packet_write_u16(packet, resource->node->id);
- aki_packet_write_s32(packet, sequence);
- aki_packet_write_u64(packet, timing->at);
- aki_packet_write_u64(packet, timing->seek_pos);
- aki_packet_write_u8(packet, timing->pause);
- aki_packet_write_bool(packet, timing->ended);
- aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
- break;
- }
- case LIANA_SINK_UNSET: {
- struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET);
- aki_packet_write_u8(packet, op);
- aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
- break;
- }
- case LIANA_SINK_PAUSE: {
- struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE);
- aki_packet_write_s32(packet, sequence);
- aki_packet_write_u64(packet, timing->at);
- aki_packet_write_u8(packet, timing->pause);
- aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
- break;
- }
- case LIANA_SINK_SEEK: {
- struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK);
- aki_packet_write_s32(packet, sequence);
- aki_packet_write_u64(packet, timing->at);
- aki_packet_write_u64(packet, timing->seek_pos);
- aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
- break;
- }
- }
-}
-
-void camu_list_init(struct camu_list *list, str *name)
-{
- al_str_clone(&list->name, name);
- lia_list_init(&list->impl);
-}
-
-void camu_list_free(struct camu_list *list)
-{
- lia_list_free(&list->impl);
- al_str_free(&list->name);
-}
diff --git a/src/server/list.h b/src/server/list.h
deleted file mode 100644
index e472024..0000000
--- a/src/server/list.h
+++ /dev/null
@@ -1,16 +0,0 @@
-#pragma once
-
-#include "../liana/list.h"
-
-#include "resource.h"
-
-struct camu_list {
- str name;
- struct lia_list impl;
-};
-
-void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry);
-void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing);
-
-void camu_list_init(struct camu_list *list, str *name);
-void camu_list_free(struct camu_list *list);
diff --git a/src/server/local_compat.c b/src/server/local_compat.c
deleted file mode 100644
index 39698c5..0000000
--- a/src/server/local_compat.c
+++ /dev/null
@@ -1,219 +0,0 @@
-#include <al/log.h>
-
-#ifdef AKIYO_HAS_CURL
-#include "../cache/handlers/http.h"
-#endif
-#include "../cache/handlers/file.h"
-#ifdef LIANA_HAVE_CDIO
-#include "../cache/handlers/cdio.h"
-#endif
-
-#include "local_compat.h"
-#include "common.h"
-#include "server.h"
-
-#ifdef CAMU_HAVE_PORTAL
-static bool uri_for_local(str *local, bool search, struct camu_portal_bridge *bridge, str *uri, struct camu_post **selected)
-{
- if (search || camu_is_url(local, 0)) {
- str module;
- str query;
- al_str_from(&module, "");
- al_str_from(&query, "");
- if (al_str_cmp(local, al_str_c("https://twitter.com"), 0, 19) == 0 ||
- al_str_cmp(local, al_str_c("https://x.com"), 0, 13) == 0) {
- al_str_cat(&query, al_str_c("tweet:"));
- al_str_cat(&query, local);
- al_str_cat(&module, al_str_c("twitter"));
- } else if (al_str_cmp(local, al_str_c("https://instagram.com"), 0, 21) == 0) {
- s32 slash = al_str_rfind(local, '/');
- if (slash >= 0) {
- al_str_cat(&query, al_str_substr(local, slash + 1, local->len));
- }
- al_str_cat(&module, al_str_c("instagram"));
- } else {
- if (search) {
- al_str_cat(&query, al_str_substr(local, 1, local->len));
- } else {
- al_str_cat(&query, al_str_c("link:"));
- al_str_cat(&query, local);
- }
- al_str_cat(&module, al_str_c("youtube"));
- }
- s32 id = camu_portal_create_search(bridge, &module, &query);
- al_str_free(&module);
- al_str_free(&query);
- if (id < 0) return false;
- struct camu_search *search = camu_portal_get_search(bridge, id);
- if (!search || !camu_search_get_page(search, 0)) {
- al_log_info("local_compat", "Search failed.");
- return false;
- }
- struct camu_result_page *page = &al_array_at(search->pages, 0);
- struct camu_post *post;
- al_array_foreach_ptr(page->posts, i, post) {
- struct camu_post_media *media;
- al_array_foreach_ptr(post->media, j, media) {
- if (media->url.len > 0) {
- al_str_clone(uri, &media->url);
- *selected = post;
- break;
- }
- }
- if (*selected) break;
- }
- camu_portal_discard_search(bridge, id);
- }
- return *selected != NULL;
-}
-#endif
-
-static void worker_signal_callback(void *userdata)
-{
- struct camu_local_compat *compat = (struct camu_local_compat *)userdata;
- u32 size;
- struct camu_server_resource *resource;
- do {
- camu_queue_try_pop(compat->pending, size, resource);
- if (size == 0) break;
- struct cch_handler *handler = resource->entry->handler;
- if (al_str_eq(&handler->liana, al_str_c("codec"))) {
- handler->maybe_spawn_worker(handler, 0);
- }
- } while (1);
-}
-
-static aki_thread_result AKI_THREADCALL queue_thread(void *userdata)
-{
- struct camu_local_compat *compat = (struct camu_local_compat *)userdata;
- aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS);
-#ifdef CAMU_HAVE_PORTAL
- bool have_python = false;
-#endif
- u32 count;
- while (aki_packet_cache_wait(&compat->queue, &count)) {
- struct aki_packet *packet = aki_packet_cache_pop(&compat->queue);
- aki_packet_cache_unlock(&compat->queue);
- if (!packet) {
- break;
- }
- str local;
- aki_packet_read_str(packet, &local);
- struct cch_entry *entry = NULL;
- struct camu_post *post = NULL;
- if (al_str_cmp(&local, al_str_c("cdda://"), 0, 7) == 0) {
-#ifdef LIANA_HAVE_CDIO
- entry = cch_handler_cdio_create();
- struct cch_chapter *chapter = &al_array_at(entry->chapters, 0);
- if (local.len > 7) {
- s64 index = al_str_to_long(al_str_substr(&local, 7, local.len), 10);
- if (index != INT64_MIN && index != INT64_MAX && index > 0 && index <= entry->chapters.size) {
- chapter = &al_array_at(entry->chapters, index - 1);
- }
- }
- entry->chapter = chapter;
- entry->handler->maybe_spawn_worker(entry->handler, chapter->start);
-#endif
- } else {
-#ifdef CAMU_HAVE_PORTAL
-#ifndef AKIYO_HAS_CURL
-#error "Curl required to use portal"
-#endif
- bool is_search = al_str_at(&local, 0) == ';';
- if (is_search || camu_is_url(&local, 0)) {
- if (!have_python) {
- // Defer python init.
- have_python = camu_python_init();
- }
- if (have_python) {
- str uri;
- if (uri_for_local(&local, is_search, &compat->bridge, &uri, &post)) {
- entry = cch_handler_http_create(&uri, compat->server->loop);
- }
- }
- } else {
- entry = cch_handler_file_create(&local);
- }
-#else
-#ifdef AKIYO_HAS_CURL
- if (camu_is_url(&local, 0)) {
- entry = cch_handler_http_create(&local, compat->server->loop);
- } else {
-#endif
- entry = cch_handler_file_create(&local);
-#ifdef AKIYO_HAS_CURL
- }
-#endif
-#endif
- }
- if (entry) {
- struct camu_server_resource *resource = al_alloc_object(struct camu_server_resource);
- al_str_clone(&resource->unique_id, &local);
- resource->post = post;
- resource->entry = entry;
- resource->node = lia_server_create_node(&compat->server->data.server, entry);
- camu_queue_push(compat->pending, resource);
- aki_signal_send(&compat->worker_signal);
- resource->duration = lia_node_get_duration(resource->node);
- al_array_push(compat->server->data.resources, resource);
- camu_queue_push(compat->results, resource);
- aki_signal_send(&compat->result_signal);
- } else {
- al_log_info("local_compat", "No resource could be created for: %.*s.", AL_STR_PRINTF(&local));
- }
- aki_packet_free(packet);
- }
-#ifdef CAMU_HAVE_PORTAL
- if (have_python) camu_python_close();
-#endif
- return 0;
-}
-
-static void result_signal_callback(void *userdata)
-{
- struct camu_local_compat *compat = (struct camu_local_compat *)userdata;
- u32 size;
- struct camu_server_resource *resource;
- do {
- camu_queue_try_pop(compat->results, size, resource);
- if (size == 0) break;
- compat->callback(compat->userdata, resource);
- } while (1);
-}
-
-void camu_local_compat_run(struct camu_local_compat *compat, struct camu_server *server)
-{
- compat->server = server;
-#ifdef CAMU_HAVE_PORTAL
- camu_post_cache_init(&compat->cache);
- camu_portal_init(&compat->bridge, &compat->cache);
-#endif
- // Size 0 to flush on the first packet.
- aki_packet_cache_init(&compat->queue, 0);
- aki_signal_init(&compat->worker_signal, worker_signal_callback, compat);
- aki_signal_start(&compat->worker_signal, compat->server->loop);
- camu_queue_init(compat->pending);
- aki_signal_init(&compat->result_signal, result_signal_callback, compat);
- aki_signal_start(&compat->result_signal, compat->server->loop);
- camu_queue_init(compat->results);
- aki_thread_create(&compat->thread, queue_thread, compat);
-}
-
-void camu_local_compat_send(struct camu_local_compat *compat, struct aki_packet *packet)
-{
- aki_packet_cache_send_packet(&compat->queue, packet);
-}
-
-void camu_local_compat_stop(struct camu_local_compat *compat)
-{
- aki_packet_cache_disable(&compat->queue);
- aki_thread_cancel(&compat->thread);
- aki_thread_join(&compat->thread);
- aki_signal_stop(&compat->worker_signal);
- aki_signal_stop(&compat->result_signal);
- camu_queue_free(compat->results);
- struct aki_packet *packet;
- while ((packet = aki_packet_cache_pop(&compat->queue))) {
- aki_packet_free(packet);
- }
-}
diff --git a/src/server/local_compat.h b/src/server/local_compat.h
deleted file mode 100644
index dcd2fc7..0000000
--- a/src/server/local_compat.h
+++ /dev/null
@@ -1,31 +0,0 @@
-#pragma once
-
-#include <al/str.h>
-#include <aki/packet_cache.h>
-#include <aki/signal.h>
-
-#ifdef CAMU_HAVE_PORTAL
-#include "../portal/src/search.h"
-#endif
-#include "../util/queue.h"
-
-struct camu_server_resource;
-struct camu_local_compat {
- struct aki_thread thread;
-#ifdef CAMU_HAVE_PORTAL
- struct camu_portal_bridge bridge;
- struct camu_post_cache cache;
-#endif
- struct aki_packet_cache queue;
- struct aki_signal worker_signal;
- queue(struct camu_server_resource *) pending;
- struct aki_signal result_signal;
- queue(struct camu_server_resource *) results;
- struct camu_server *server;
- void (*callback)(void *, struct camu_server_resource *);
- void *userdata;
-};
-
-void camu_local_compat_run(struct camu_local_compat *compat, struct camu_server *server);
-void camu_local_compat_send(struct camu_local_compat *compat, struct aki_packet *packet);
-void camu_local_compat_stop(struct camu_local_compat *compat);
diff --git a/src/server/meson.build b/src/server/meson.build
index 0617c2c..88f80c0 100644
--- a/src/server/meson.build
+++ b/src/server/meson.build
@@ -1,10 +1,8 @@
server_src = [
'server.c',
'common.c',
- 'list.c',
'user.c',
'db.c',
- 'local_compat.c'
]
server_deps = [common_deps, cache, liana_server]
server = declare_dependency(sources: server_src, dependencies: server_deps)
diff --git a/src/server/resource.h b/src/server/resource.h
index c1cf619..87250e9 100644
--- a/src/server/resource.h
+++ b/src/server/resource.h
@@ -2,12 +2,28 @@
#include "../cache/entry.h"
#include "../liana/server.h"
-#include "../portal/src/post.h"
-struct camu_server_resource {
- str unique_id;
- struct camu_post *post;
+enum {
+ CAMU_RESOURCE_NOT_LOADED = 0,
+ CAMU_RESOURCE_LOADING,
+ CAMU_RESOURCE_LOADED
+};
+
+struct camu_resource {
+ u8 type;
+ u8 load;
struct cch_entry *entry;
struct lia_node *node;
u64 duration;
+ array(struct lia_list_entry *) pending;
+};
+
+struct camu_resource_file {
+ struct camu_resource r;
+ str path;
+};
+
+struct camu_resource_portal {
+ struct camu_resource r;
+ struct camu_post *post;
};
diff --git a/src/server/server.c b/src/server/server.c
index d7f8bc9..d3050ca 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -1,9 +1,16 @@
#include <al/log.h>
#include <al/lib.h>
+#include "../cache/handlers/file.h"
+#include "../cache/handlers/http.h"
+#include "../libclient/common.h"
+#include "../libsink/common.h"
+#ifdef CAMU_HAVE_PORTAL
+#include "../portal/src/packet_ext.h"
+#endif
+
#include "server.h"
#include "common.h"
-#include "list.h"
#include "db.h"
static struct camu_user *get_user_by_username(struct camu_server *server, str *username)
@@ -26,9 +33,9 @@ static struct camu_server_client *get_client_by_connection(struct camu_server *s
return NULL;
}
-static struct camu_list *get_list_from_name(struct camu_server *server, str *name)
+static struct lia_list *get_list_from_name(struct camu_server *server, str *name)
{
- struct camu_list *list;
+ struct lia_list *list;
al_array_foreach(server->lists, i, list) {
if (al_str_eq(&list->name, name)) return list;
}
@@ -44,7 +51,9 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s
return NULL;
}
-static bool identify_command_callback(void *userdata, struct aki_rpc_connection *conn,
+static void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable);
+
+static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
@@ -83,8 +92,7 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection
al_str_clone(&sink->name, &name);
sink->server = server;
al_array_push(server->sinks, sink);
- struct camu_list *list = al_array_at(server->lists, 0);
- lia_list_add_sink(&list->impl, camu_list_sink_callback, sink);
+ handle_toggle_sink(server, al_str_c("default"), sink, true);
al_log_info("server", "New sink.");
break;
}
@@ -94,7 +102,93 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection
return true;
}
-static bool client_command_command_callback(void *userdata, struct aki_rpc_connection *conn,
+static void client_portal_callback(void *userdata, struct camu_portal_result *result)
+{
+ struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata;
+ struct aki_packet *packet = aki_rpc_get_packet(conn->rpc, CAMU_CLIENT_RESULTS);
+ aki_packet_write_u8(packet, result->op);
+ aki_packet_write_s32(packet, result->id);
+ switch (result->op) {
+ case CAMU_CLIENT_CREATE_SEARCH: {
+ break;
+ }
+ case CAMU_CLIENT_GET_PAGE: {
+ struct camu_result_page *page = result->page;
+ aki_packet_write_u32(packet, page->num);
+ aki_packet_write_u32(packet, page->posts.size);
+ struct camu_post *post;
+ al_array_foreach_ptr(page->posts, i, post) {
+ aki_packet_write_post(packet, post);
+ }
+ aki_packet_write_u32(packet, page->list.size);
+ str *unique_id;
+ al_array_foreach_ptr(page->list, i, unique_id) {
+ aki_packet_write_str(packet, unique_id);
+ }
+ break;
+ }
+ }
+ aki_rpc_connection_command(conn, packet, NULL, NULL);
+}
+
+static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing)
+{
+ struct camu_server_sink *sink = (struct camu_server_sink *)userdata;
+ switch (op) {
+ case LIANA_SINK_SET:
+ case LIANA_SINK_BUFFER:
+ case LIANA_SINK_BUFFER_AND_QUEUE: {
+ struct camu_resource *resource = (struct camu_resource *)entry->opaque;
+ struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET);
+ aki_packet_write_u8(packet, op);
+ aki_packet_write_str(packet, &sink->server->addr);
+ aki_packet_write_u16(packet, CAMU_PORT);
+ aki_packet_write_u16(packet, resource->node->id);
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_u64(packet, timing->at);
+ aki_packet_write_u64(packet, timing->seek_pos);
+ aki_packet_write_u8(packet, timing->pause);
+ aki_packet_write_bool(packet, timing->ended);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
+ case LIANA_SINK_UNSET: {
+ struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET);
+ aki_packet_write_u8(packet, op);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
+ case LIANA_SINK_PAUSE: {
+ struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE);
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_u64(packet, timing->at);
+ aki_packet_write_u8(packet, timing->pause);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
+ case LIANA_SINK_SEEK: {
+ struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK);
+ aki_packet_write_s32(packet, sequence);
+ aki_packet_write_u64(packet, timing->at);
+ aki_packet_write_u64(packet, timing->seek_pos);
+ aki_rpc_connection_command(sink->conn, packet, NULL, NULL);
+ break;
+ }
+ }
+}
+
+void handle_toggle_sink(struct camu_server *server, str *name, struct camu_server_sink *sink, bool enable)
+{
+ struct lia_list *list = get_list_from_name(server, name);
+ if (!list) return;
+ if (enable) {
+ lia_list_add_sink(list, list_sink_callback, sink);
+ } else {
+ lia_list_remove_sink(list, sink);
+ }
+}
+
+static bool client_command_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
@@ -108,18 +202,33 @@ static bool client_command_command_callback(void *userdata, struct aki_rpc_conne
case CAMU_CLIENT_CREATE_LIST: {
str name;
aki_packet_read_str(packet, &name);
- struct camu_list *list = al_alloc_object(struct camu_list);
- camu_list_init(list, &name);
+ struct lia_list *list = al_alloc_object(struct lia_list);
+ lia_list_init(list, &name);
al_array_push(server->lists, list);
break;
}
- case CAMU_CLIENT_ENABLE_SINK: {
+ case CAMU_CLIENT_TOGGLE_SINK: {
str name;
aki_packet_read_str(packet, &name);
- struct camu_list *list = get_list_from_name(server, &name);
- aki_packet_read_str(packet, &name);
struct camu_server_sink *sink = get_sink_from_name(server, &name);
- lia_list_add_sink(&list->impl, camu_list_sink_callback, sink);
+ if (!sink) goto out;
+ aki_packet_read_str(packet, &name); // list name.
+ bool enable = aki_packet_read_bool(packet);
+ handle_toggle_sink(server, &name, sink, enable);
+ break;
+ }
+ case CAMU_CLIENT_CREATE_SEARCH: {
+ str module;
+ aki_packet_read_str(packet, &module);
+ str query;
+ aki_packet_read_str(packet, &query);
+ camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn);
+ break;
+ }
+ case CAMU_CLIENT_GET_PAGE: {
+ s32 id = aki_packet_read_s32(packet);
+ u32 num = aki_packet_read_u32(packet);
+ camu_portal_get_page(&server->bridge, id, num, client_portal_callback, conn);
break;
}
}
@@ -129,7 +238,116 @@ out:
return true;
}
-static bool list_action_command_callback(void *userdata, struct aki_rpc_connection *conn,
+static void node_callback(void *userdata, u8 op, u64 duration)
+{
+ struct camu_resource *resource = (struct camu_resource *)userdata;
+ switch (op) {
+ case LIANA_NODE_DURATION:
+ resource->load = CAMU_RESOURCE_LOADED;
+ resource->duration = duration;
+ struct lia_list_entry *entry;
+ al_array_foreach(resource->pending, i, entry) {
+ lia_list_pump(entry->list);
+ }
+ resource->pending.size = 0;
+ break;
+ }
+}
+
+static void maybe_add_to_pending(struct camu_resource *resource, struct lia_list_entry *entry)
+{
+ struct lia_list_entry *rentry;
+ al_array_foreach(resource->pending, i, rentry) {
+ if (rentry == entry) return;
+ }
+ al_array_push(resource->pending, entry);
+}
+
+static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *result)
+{
+ struct camu_server *server = (struct camu_server *)userdata;
+ (void)server;
+ struct camu_resource *resource = (struct camu_resource *)entry->opaque;
+ switch (op) {
+ case LIANA_LOAD_ENTRY:
+ switch (resource->load) {
+ case CAMU_RESOURCE_NOT_LOADED:
+ resource->load = CAMU_RESOURCE_LOADING;
+ lia_node_get_duration(resource->node);
+ // fallthrough
+ case CAMU_RESOURCE_LOADING:
+ maybe_add_to_pending(resource, entry);
+ *(bool *)result = false;
+ break;
+ case CAMU_RESOURCE_LOADED:
+ *(bool *)result = true;
+ break;
+ }
+ break;
+ case LIANA_GET_DURATION: {
+ *(u64 *)result = resource->duration;
+ break;
+ }
+ case LIANA_UNLOAD_ENTRY:
+ break;
+ }
+}
+
+static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet)
+{
+ u8 op = aki_packet_read_u8(packet);
+ switch (op) {
+ case CAMU_RESOURCE_FILE: {
+ str path;
+ aki_packet_read_str(packet, &path);
+ struct camu_resource_file *resource = al_alloc_object(struct camu_resource_file);
+ resource->r.type = CAMU_RESOURCE_FILE;
+ resource->r.load = CAMU_RESOURCE_NOT_LOADED;
+ al_str_clone(&resource->path, &path);
+ resource->r.entry = cch_handler_file_create(&path);
+ resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry);
+ resource->r.node->callback = node_callback;
+ resource->r.node->userdata = (struct camu_resource *)resource;
+ resource->r.duration = LIANA_TIMESTAMP_INVALID;
+ al_array_init(resource->r.pending);
+ lia_list_add(list, resource, resource->r.duration, &path);
+ break;
+ }
+ case CAMU_RESOURCE_PORTAL: {
+ str unique_id;
+ aki_packet_read_str(packet, &unique_id);
+ u32 index = aki_packet_read_u32(packet);
+ struct camu_post *post = camu_post_cache_get(&server->cache, &unique_id);
+ struct cch_entry *entry = NULL;
+ if (index <= post->media.size) {
+ struct camu_post_media *media = &al_array_at(post->media, index);
+ if (!al_str_is_empty(&media->url)) {
+ entry = cch_handler_http_create(&media->url, server->loop);
+ }
+ }
+ if (!entry) {
+ al_log_warn("server", "Failed to load resource %.*s %u.", AL_STR_PRINTF(&unique_id), index);
+ return;
+ }
+ struct camu_resource_portal *resource = al_alloc_object(struct camu_resource_portal);
+ resource->r.type = CAMU_RESOURCE_PORTAL;
+ resource->r.load = CAMU_RESOURCE_NOT_LOADED;
+ resource->post = post;
+ resource->r.entry = entry;
+ struct cch_handler *handler = resource->r.entry->handler;
+ handler->maybe_spawn_worker(handler, 0);
+ resource->r.node = lia_server_create_node(&server->data.server, resource->r.entry);
+ resource->r.node->callback = node_callback;
+ resource->r.node->userdata = (struct camu_resource *)resource;
+ resource->r.duration = LIANA_TIMESTAMP_INVALID;
+ al_array_init(resource->r.pending);
+ lia_list_add(list, resource, resource->r.duration, al_str_c("dfdd"));
+ break;
+ }
+ }
+}
+
+static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn,
struct aki_packet *packet, struct aki_packet *rpacket)
{
struct camu_server *server = (struct camu_server *)userdata;
@@ -139,50 +357,50 @@ static bool list_action_command_callback(void *userdata, struct aki_rpc_connecti
str name;
aki_packet_read_str(packet, &name);
- struct camu_list *list = get_list_from_name(server, &name);
+ struct lia_list *list = get_list_from_name(server, &name);
if (!list) goto out;
u8 op = aki_packet_read_u8(packet);
switch (op) {
case CAMU_LIST_ADD: {
- camu_local_compat_send(&server->compat, packet);
+ handle_add_command(server, list, packet);
return false;
}
case CAMU_LIST_SKIP: {
s32 sequence = aki_packet_read_s32(packet);
s32 n = aki_packet_read_s32(packet);
- lia_list_skip(&list->impl, sequence, n);
+ lia_list_skip(list, sequence, n);
break;
}
case CAMU_LIST_SKIPTO: {
s32 sequence = aki_packet_read_s32(packet);
s32 i = aki_packet_read_s32(packet);
- lia_list_skipto(&list->impl, sequence, i);
+ lia_list_skipto(list, sequence, i);
break;
}
case CAMU_LIST_SHUFFLE: {
- lia_list_shuffle(&list->impl);
+ lia_list_shuffle(list);
break;
}
case CAMU_LIST_TOGGLE_PAUSE: {
s32 sequence = aki_packet_read_s32(packet);
f64 pts = aki_packet_read_f64(packet);
- lia_list_toggle_pause(&list->impl, sequence, pts);
+ lia_list_toggle_pause(list, sequence, pts);
break;
}
case CAMU_LIST_SEEK: {
s32 sequence = aki_packet_read_s32(packet);
f64 percent = aki_packet_read_f64(packet);
- lia_list_seek(&list->impl, sequence, percent);
+ lia_list_seek(list, sequence, percent);
break;
}
case CAMU_LIST_UNSET: {
- lia_list_unset(&list->impl);
+ lia_list_unset(list);
break;
}
case CAMU_LIST_END: {
s32 sequence = aki_packet_read_s32(packet);
- lia_list_end(&list->impl, sequence);
+ lia_list_end(list, sequence);
break;
}
}
@@ -193,9 +411,9 @@ out:
}
static struct aki_rpc_command commands[] = {
- { .op = CAMU_SERVER_IDENTIFY, .callback = identify_command_callback, .userdata = NULL },
- { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_command_callback, .userdata = NULL },
- { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_command_callback, .userdata = NULL }
+ { .op = CAMU_SERVER_IDENTIFY, .callback = identify_callback, .userdata = NULL },
+ { .op = CAMU_SERVER_CLIENT_COMMAND, .callback = client_command_callback, .userdata = NULL },
+ { .op = CAMU_SERVER_LIST_ACTION, .callback = list_action_callback, .userdata = NULL }
};
static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
@@ -228,9 +446,9 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
struct camu_server_node *node;
al_array_foreach(server->nodes, i, node) {
if (node->conn == conn) {
+ al_log_info("server", "Node removed.");
cleanup_node(node);
al_array_remove_at(server->nodes, i);
- al_log_info("server", "Node removed.");
break;
}
}
@@ -238,9 +456,9 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
struct camu_server_client *client;
al_array_foreach(server->clients, i, client) {
if (client->conn == conn) {
+ al_log_info("server", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->name));
cleanup_client(client);
al_array_remove_at(server->clients, i);
- al_log_info("server", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->name));
break;
}
}
@@ -248,13 +466,13 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
struct camu_server_sink *sink;
al_array_foreach(server->sinks, i, sink) {
if (sink->conn == conn) {
+ al_log_info("server", "Sink removed.");
al_array_remove_at(server->sinks, i);
- struct camu_list *list;
+ struct lia_list *list;
al_array_foreach(server->lists, j, list) {
- lia_list_remove_sink(&list->impl, sink);
+ lia_list_remove_sink(list, sink);
}
cleanup_sink(sink);
- al_log_info("server", "Sink removed.");
break;
}
}
@@ -276,13 +494,6 @@ static bool multiplex_callback(void *userdata, u8 id, struct aki_socket *sock)
return false;
}
-static void local_compat_callback(void *userdata, struct camu_server_resource *resource)
-{
- struct camu_server *server = (struct camu_server *)userdata;
- struct camu_list *list = al_array_last(server->lists);
- if (resource) lia_list_add(&list->impl, resource, resource->duration, &resource->unique_id);
-}
-
bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop)
{
server->loop = loop;
@@ -293,8 +504,10 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop
al_array_init(server->users);
al_array_init(server->lists);
- struct camu_list *list = al_alloc_object(struct camu_list);
- camu_list_init(list, al_str_c("default"));
+ struct lia_list *list = al_alloc_object(struct lia_list);
+ lia_list_init(list, al_str_c("default"));
+ list->callback = list_callback;
+ list->userdata = server;
al_array_push(server->lists, list);
aki_rpc_init(&server->server, server->loop, connection_callback, connection_closed_callback, server);
@@ -305,9 +518,10 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop
lia_server_init(&server->data.server, server->loop);
- server->compat.callback = local_compat_callback;
- server->compat.userdata = server;
- camu_local_compat_run(&server->compat, server);
+#ifdef CAMU_HAVE_PORTAL
+ camu_post_cache_init(&server->cache);
+ camu_portal_init(&server->bridge, &server->cache, server->loop);
+#endif
return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server);
}
@@ -321,7 +535,9 @@ bool camu_server_listen(struct camu_server *server, str *addr, u16 port)
void camu_server_close(struct camu_server *server)
{
- camu_local_compat_stop(&server->compat);
+#ifdef CAMU_HAVE_PORTAL
+ camu_portal_close(&server->bridge);
+#endif
lia_server_close(&server->data.server);
aki_multiplex_socket_close(&server->multi);
}
@@ -330,10 +546,10 @@ void camu_server_free(struct camu_server *server)
{
// This needs to happen in flight.
// Too many things can block us before we get here.
- struct camu_server_resource *resource;
- al_array_foreach(server->data.resources, i, resource) {
- cch_entry_free(&resource->entry);
- }
+ //struct camu_server_resource *resource;
+ //al_array_foreach(server->data.resources, i, resource) {
+ // cch_entry_free(&resource->entry);
+ //}
al_array_free(server->data.resources);
lia_server_free(&server->data.server);
aki_rpc_free(&server->server);
@@ -342,5 +558,5 @@ void camu_server_free(struct camu_server *server)
void camu_server_local_add(struct camu_server *server, struct aki_packet *packet)
{
- camu_local_compat_send(&server->compat, packet);
+ handle_add_command(server, al_array_last(server->lists), packet);
}
diff --git a/src/server/server.h b/src/server/server.h
index 6246f79..aaabc3c 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -3,12 +3,14 @@
#include <aki/multiplex.h>
#include <aki/rpc2.h>
-#include "../cache/entry.h"
#include "../liana/server.h"
+#include "../liana/list.h"
+#ifdef CAMU_HAVE_PORTAL
+#include "../portal/src/search.h"
+#endif
#include "user.h"
#include "resource.h"
-#include "local_compat.h"
struct camu_server_node {
struct aki_rpc_connection *conn;
@@ -34,12 +36,15 @@ struct camu_server {
array(struct camu_server_client *) clients;
array(struct camu_server_sink *) sinks;
array(struct camu_user *) users;
- array(struct camu_list *) lists;
- struct camu_local_compat compat;
+ array(struct lia_list *) lists;
struct {
struct lia_server server;
- array(struct camu_server_resource *) resources;
+ array(struct camu_resource *) resources;
} data;
+#ifdef CAMU_HAVE_PORTAL
+ struct camu_portal_bridge bridge;
+ struct camu_post_cache cache;
+#endif
};
bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop);