diff options
Diffstat (limited to 'src/server')
| -rw-r--r-- | src/server/common.h | 3 | ||||
| -rw-r--r-- | src/server/db.c | 2 | ||||
| -rw-r--r-- | src/server/resource.h | 6 | ||||
| -rw-r--r-- | src/server/server.c | 264 | ||||
| -rw-r--r-- | src/server/server.h | 2 |
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); |