summaryrefslogtreecommitdiff
path: root/src/server
diff options
context:
space:
mode:
Diffstat (limited to 'src/server')
-rw-r--r--src/server/common.h3
-rw-r--r--src/server/db.c2
-rw-r--r--src/server/resource.h6
-rw-r--r--src/server/server.c264
-rw-r--r--src/server/server.h2
5 files changed, 184 insertions, 93 deletions
diff --git a/src/server/common.h b/src/server/common.h
index b939974..b6ee7a1 100644
--- a/src/server/common.h
+++ b/src/server/common.h
@@ -47,8 +47,9 @@ enum {
enum {
CAMU_RESOURCE_FILE = 0,
+ CAMU_RESOURCE_CDIO,
CAMU_RESOURCE_PORTAL,
- CAMU_RESOURCE_CDIO
+ CAMU_RESOURCE_SIMPLE_SEARCH
};
AL_UNUSED_FUNCTION_PUSH
diff --git a/src/server/db.c b/src/server/db.c
index c9608bb..6f76864 100644
--- a/src/server/db.c
+++ b/src/server/db.c
@@ -15,7 +15,7 @@ static bool open_user(struct camu_server *server, struct aki_dir_entry *dir)
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_log_error("server", "Failed to parse %.*s:%d:%d (%s).",
AL_STR_PRINTF(&dir->path), error.line, error.column, error.text);
return false;
}
diff --git a/src/server/resource.h b/src/server/resource.h
index 6056e8e..559e2b4 100644
--- a/src/server/resource.h
+++ b/src/server/resource.h
@@ -3,12 +3,6 @@
#include "../cache/entry.h"
#include "../liana/server.h"
-enum {
- CAMU_RESOURCE_NOT_LOADED = 0,
- CAMU_RESOURCE_LOADING,
- CAMU_RESOURCE_LOADED
-};
-
struct camu_resource {
u8 type;
u8 load;
diff --git a/src/server/server.c b/src/server/server.c
index 109be14..4ec61da 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -52,13 +52,24 @@ static struct camu_server_sink *get_sink_from_name(struct camu_server *server, s
return NULL;
}
+static void write_user_state(struct camu_server *server, struct camu_user *user, struct aki_packet *packet)
+{
+ (void)user;
+ struct camu_search *search;
+ aki_packet_write_u32(packet, server->bridge.searches.size);
+ al_array_foreach(server->bridge.searches, i, search) {
+ aki_packet_write_s32(packet, search->id);
+ aki_packet_write_str(packet, &search->module);
+ aki_packet_write_str(packet, &search->query);
+ }
+}
+
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;
- (void)rpacket;
u8 op = aki_packet_read_u8(packet);
switch (op) {
@@ -83,6 +94,7 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
client->user = user;
al_array_push(server->clients, client);
al_log_info("server", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->name));
+ write_user_state(server, user, rpacket);
break;
}
case CAMU_SINK: {
@@ -103,33 +115,15 @@ static bool identify_callback(void *userdata, struct aki_rpc_connection *conn,
return true;
}
-static void client_portal_callback(void *userdata, struct camu_portal_result *result)
+static bool client_still_connected(struct camu_server *server, struct aki_rpc_connection *conn)
{
- 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);
+ struct camu_server_client *client;
+ al_array_foreach(server->clients, i, client) {
+ if (client->conn == conn) {
+ return true;
}
- break;
}
- }
- aki_rpc_connection_command(conn, packet, NULL, NULL);
+ return false;
}
static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing)
@@ -187,11 +181,21 @@ void handle_toggle_sink(struct camu_server *server, str *name, struct camu_serve
{
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);
+ if (enable) lia_list_add_sink(list, list_sink_callback, sink);
+ else lia_list_remove_sink(list, sink);
+}
+
+static void process_pending(struct camu_resource *resource)
+{
+ array(struct lia_list_entry *) pending;
+ // resource->pending may be edited in a list_pump() call.
+ al_array_clone(pending, resource->pending);
+ resource->pending.size = 0;
+ struct lia_list_entry *entry;
+ al_array_foreach(pending, i, entry) {
+ lia_list_pump(entry->list);
}
+ al_array_free(pending);
}
static void node_callback(void *userdata, u8 op, u64 duration)
@@ -199,27 +203,14 @@ 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->load = LIANA_ENTRY_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;
+ process_pending(resource);
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)
+static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, void *opaque)
{
struct camu_server *server = (struct camu_server *)userdata;
(void)server;
@@ -227,26 +218,64 @@ static void list_callback(void *userdata, u8 op, struct lia_list_entry *entry, v
switch (op) {
case LIANA_LOAD_ENTRY:
switch (resource->load) {
- case CAMU_RESOURCE_NOT_LOADED:
- resource->load = CAMU_RESOURCE_LOADING;
+ case LIANA_ENTRY_PREPARED:
+ resource->load = LIANA_ENTRY_LOADING;
lia_node_get_duration(resource->node);
// fallthrough
- case CAMU_RESOURCE_LOADING:
- maybe_add_to_pending(resource, entry);
- *(bool *)result = false;
+ case LIANA_ENTRY_PREPARING:
+ case LIANA_ENTRY_LOADING:
+ al_array_push(resource->pending, entry);
break;
- case CAMU_RESOURCE_LOADED:
- *(bool *)result = true;
+ case LIANA_ENTRY_LOADED:
+ case LIANA_ENTRY_ERRORED:
break;
}
+ *(u8 *)opaque = resource->load;
break;
case LIANA_GET_DURATION: {
- *(u64 *)result = resource->duration;
+ *(u64 *)opaque = resource->duration;
break;
}
case LIANA_UNLOAD_ENTRY:
break;
+ case LIANA_LIST_META:
+ if (server->meta_callback) {
+ server->meta_callback(server->userdata, *(u8 *)opaque, entry->list, entry);
+ }
+ break;
+ }
+
+}
+
+static void client_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result)
+{
+ struct camu_server *server = (struct camu_server *)userdata0;
+ struct aki_rpc_connection *conn = (struct aki_rpc_connection *)userdata1;
+ if (!client_still_connected(server, conn)) return;
+ 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 bool client_command_callback(void *userdata, struct aki_rpc_connection *conn,
@@ -282,8 +311,8 @@ static bool client_command_callback(void *userdata, struct aki_rpc_connection *c
}
case CAMU_CLIENT_CREATE_SEARCH: {
str module;
- aki_packet_read_str(packet, &module);
str query;
+ aki_packet_read_str(packet, &module);
aki_packet_read_str(packet, &query);
camu_portal_create_search(&server->bridge, &module, &query, client_portal_callback, conn);
break;
@@ -301,6 +330,55 @@ out:
return true;
}
+static struct cch_entry *entry_from_post(struct camu_server *server, struct camu_post *post, u32 index)
+{
+ 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(&post->unique_id), index);
+ return NULL;
+ }
+ entry->handler->maybe_spawn_worker(entry->handler, 0);
+ return entry;
+}
+
+static void simple_search_portal_callback(void *userdata0, void *userdata1, struct camu_portal_result *result)
+{
+ struct camu_server *server = (struct camu_server *)userdata0;
+ struct camu_resource_portal *portal = (struct camu_resource_portal *)userdata1;
+ struct camu_resource *resource = (struct camu_resource *)portal;
+ switch (result->op) {
+ case CAMU_CLIENT_CREATE_SEARCH: {
+ camu_portal_get_page(&server->bridge, result->id, 0, simple_search_portal_callback, portal);
+ return;
+ }
+ case CAMU_CLIENT_GET_PAGE: {
+ if (!result->page || result->page->list.size == 0) break;
+ struct camu_post *post = camu_post_cache_get(&server->cache, &al_array_at(result->page->list, 0));
+ if (!post) break;
+ struct cch_entry *entry = entry_from_post(server, post, 0);
+ if (!entry) break;
+ portal->post = post;
+ resource->entry = entry;
+ resource->node = lia_server_create_node(&server->data.server, resource->entry);
+ resource->node->callback = node_callback;
+ resource->node->userdata = resource;
+ resource->type = CAMU_RESOURCE_PORTAL;
+ resource->load = LIANA_ENTRY_PREPARED;
+ process_pending(resource);
+ return;
+ }
+ }
+ // No return is the error case.
+ al_log_warn("server", "Server resource failed to load.");
+ resource->load = LIANA_ENTRY_ERRORED;
+}
+
static void handle_add_command(struct camu_server *server, struct lia_list *list, struct aki_packet *packet)
{
u8 op = aki_packet_read_u8(packet);
@@ -318,6 +396,21 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
al_wstr_from_str(&name, &path);
resource = (struct camu_resource *)file;
resource->type = CAMU_RESOURCE_FILE;
+ resource->load = LIANA_ENTRY_PREPARED;
+ break;
+ }
+ case CAMU_RESOURCE_CDIO: {
+ u32 track = aki_packet_read_u32(packet);
+ entry = cch_handler_cdio_create();
+ if (!entry) return;
+ struct cch_chapter *chapter = &al_array_at(entry->chapters, track);
+ entry->handler->maybe_spawn_worker(entry->handler, chapter->start);
+ struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio);
+ cdio->track = track;
+ al_wstr_from_cstr(&name, "cdio");
+ resource = (struct camu_resource *)cdio;
+ resource->type = CAMU_RESOURCE_CDIO;
+ resource->load = LIANA_ENTRY_PREPARED;
break;
}
case CAMU_RESOURCE_PORTAL: {
@@ -325,44 +418,44 @@ static void handle_add_command(struct camu_server *server, struct lia_list *list
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);
- 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;
- }
- entry->handler->maybe_spawn_worker(entry->handler, 0);
+ entry = entry_from_post(server, post, index);
+ if (!entry) return;
struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal);
portal->post = post;
al_wstr_clone(&name, &post->title);
resource = (struct camu_resource *)portal;
resource->type = CAMU_RESOURCE_PORTAL;
+ resource->load = LIANA_ENTRY_PREPARED;
break;
}
- case CAMU_RESOURCE_CDIO: {
- u32 track = aki_packet_read_u32(packet);
- entry = cch_handler_cdio_create();
- struct cch_chapter *chapter = &al_array_at(entry->chapters, track);
- entry->handler->maybe_spawn_worker(entry->handler, chapter->start);
- struct camu_resource_cdio *cdio = al_alloc_object(struct camu_resource_cdio);
- cdio->track = track;
- al_wstr_from_cstr(&name, "cdio");
- resource = (struct camu_resource *)cdio;
- resource->type = CAMU_RESOURCE_CDIO;
+ case CAMU_RESOURCE_SIMPLE_SEARCH: {
+ str search;
+ aki_packet_read_str(packet, &search);
+ struct camu_resource_portal *portal = al_alloc_object(struct camu_resource_portal);
+ portal->post = NULL;
+ al_wstr_from_str(&name, &search);
+ resource = (struct camu_resource *)portal;
+ resource->type = CAMU_RESOURCE_SIMPLE_SEARCH;
+ resource->load = LIANA_ENTRY_PREPARING;
+ str query;
+ if (al_str_at(&search, 0) == ';') { // search.
+ al_str_clone(&query, al_str_substr(&search, 1, search.len));
+ } else {
+ al_str_from(&query, "link:");
+ al_str_cat(&query, &search);
+ }
+ camu_portal_create_search(&server->bridge, al_str_c("youtube"), &query, simple_search_portal_callback, portal);
+ al_str_free(&query);
break;
}
}
al_assert(resource);
- resource->load = CAMU_RESOURCE_NOT_LOADED;
- resource->entry = entry;
- resource->node = lia_server_create_node(&server->data.server, resource->entry);
- resource->node->callback = node_callback;
- resource->node->userdata = resource;
+ if (entry) {
+ resource->entry = entry;
+ resource->node = lia_server_create_node(&server->data.server, resource->entry);
+ resource->node->callback = node_callback;
+ resource->node->userdata = resource;
+ }
resource->duration = LIANA_TIMESTAMP_INVALID;
al_array_init(resource->pending);
lia_list_add(list, resource, resource->duration, &name);
@@ -412,8 +505,9 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn
}
case CAMU_LIST_SEEK: {
s32 sequence = aki_packet_read_s32(packet);
+ u32 id = aki_packet_read_u32(packet);
f64 percent = aki_packet_read_f64(packet);
- lia_list_seek(list, sequence, percent);
+ lia_list_seek(list, sequence, id, percent);
break;
}
case CAMU_LIST_UNSET: {
@@ -421,8 +515,8 @@ static bool list_action_callback(void *userdata, struct aki_rpc_connection *conn
break;
}
case CAMU_LIST_END: {
- s32 sequence = aki_packet_read_s32(packet);
- lia_list_end(list, sequence);
+ u32 id = aki_packet_read_u32(packet);
+ lia_list_end(list, id);
break;
}
}
@@ -542,7 +636,7 @@ bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop
#ifdef CAMU_HAVE_PORTAL
camu_post_cache_init(&server->cache);
- camu_portal_init(&server->bridge, &server->cache, server->loop);
+ camu_portal_init(&server->bridge, &server->cache, server->loop, server);
#endif
return aki_multiplex_socket_init(&server->multi, type, multiplex_callback, server);
diff --git a/src/server/server.h b/src/server/server.h
index aaabc3c..ec0f85b 100644
--- a/src/server/server.h
+++ b/src/server/server.h
@@ -45,6 +45,8 @@ struct camu_server {
struct camu_portal_bridge bridge;
struct camu_post_cache cache;
#endif
+ void (*meta_callback)(void *, u8, struct lia_list *, struct lia_list_entry *);
+ void *userdata;
};
bool camu_server_init(struct camu_server *server, u8 type, struct aki_event_loop *loop);