diff options
Diffstat (limited to 'src/tree')
| -rw-r--r-- | src/tree/common.h (renamed from src/tree/commands.h) | 6 | ||||
| -rw-r--r-- | src/tree/list.c | 112 | ||||
| -rw-r--r-- | src/tree/list.h | 31 | ||||
| -rw-r--r-- | src/tree/meson.build | 6 | ||||
| -rw-r--r-- | src/tree/resource_manager.c | 6 | ||||
| -rw-r--r-- | src/tree/tree.c | 227 | ||||
| -rw-r--r-- | src/tree/tree.h | 39 |
7 files changed, 332 insertions, 95 deletions
diff --git a/src/tree/commands.h b/src/tree/common.h index 1b1187e..4e3ec3a 100644 --- a/src/tree/commands.h +++ b/src/tree/common.h @@ -1,5 +1,9 @@ #pragma once +#define TREE_PORT 4356 +#define TREE_RESOURCE_PORT 4357 +#define TREE_STREAM_PORT 4358 + enum { TREE_NODE = 0, TREE_CLIENT, @@ -8,6 +12,8 @@ enum { enum { TREE_CMD_IDENTIFY = 0, + TREE_CMD_STATUS, TREE_CMD_SEARCH, + TREE_CMD_RESUME_SEARCH, TREE_CMD_ADD }; diff --git a/src/tree/list.c b/src/tree/list.c new file mode 100644 index 0000000..016c0de --- /dev/null +++ b/src/tree/list.c @@ -0,0 +1,112 @@ +#ifdef AKIYO_HAS_CURL +#include "../cache/handlers/http.h" +#endif + +#include "../libsink/common.h" + +#include "list.h" +#include "tree.h" + +void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name) +{ + al_str_clone(&list->name, name); + list->set = -1; + list->current = 0; + al_array_init(list->entries); + al_array_init(list->sinks); + list->tree = tree; +} + +void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink) +{ + al_array_push(list->sinks, sink); +} + +void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink) +{ + struct tree_sink *rsink; + al_array_foreach(list->sinks, i, rsink) { + if (rsink == sink) { + al_array_remove_at_iter(list->sinks, i); + break; + } + } +} + +bool tree_list_add(struct tree_list *list, str *unique_id, u32 index) +{ + struct tree_server *tree = list->tree; + struct sho_post *post = sho_post_cache_get(&tree->resources.cache, unique_id); + if (!post) return false; + + struct tree_list_entry *entry = al_alloc_object(struct tree_list_entry); + + entry->post = post; + entry->index = index; + + entry->entry = cch_handler_http_create(&al_array_at(entry->post->media, entry->index).url); + if (!entry->entry) { + al_free(entry); + return false; + } + entry->entry->handler->maybe_spawn_worker(entry->entry->handler, 0); + + entry->node_id = bmu_server_create_node(&tree->streams.server, entry->entry); + + entry->buffer_requested = false; + + al_array_push(list->entries, entry); + + tree_list_pump(list); + + return true; +} + +void tree_list_skip(struct tree_list *list, s32 n) +{ + s32 size = (s32)list->entries.size; + if (list->current + n < 0 || list->current + n >= size) { + return; + } + list->current += n; + tree_list_pump(list); +} + +static void send_buffer_cmd(struct tree_list *list, struct tree_list_entry *entry) +{ + struct aki_packet *packet; + struct tree_sink *sink; + al_array_foreach(list->sinks, i, sink) { + packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER); + aki_packet_write_str(packet, al_str_c("127.0.0.1")); + aki_packet_write_s32(packet, TREE_STREAM_PORT); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + } +} + +static void send_set_cmd(struct tree_list *list, struct tree_list_entry *entry) +{ + struct aki_packet *packet; + struct tree_sink *sink; + al_array_foreach(list->sinks, i, sink) { + packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET); + aki_packet_write_u16(packet, entry->node_id); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + } +} + +void tree_list_pump(struct tree_list *list) +{ + s32 size = (s32)list->entries.size; + if (size <= list->current) return; + if (list->set != list->current) { + struct tree_list_entry *entry = al_array_at(list->entries, list->current); + if (!entry->buffer_requested) { + send_buffer_cmd(list, entry); + entry->buffer_requested = true; + } + send_set_cmd(list, entry); + list->set = list->current; + } +} diff --git a/src/tree/list.h b/src/tree/list.h new file mode 100644 index 0000000..ac7f8ee --- /dev/null +++ b/src/tree/list.h @@ -0,0 +1,31 @@ +#pragma once + +#include <al/str.h> +#include <al/array.h> + +#include "../cache/entry.h" + +struct tree_list_entry { + struct sho_post *post; + u32 index; + u16 node_id; + bool buffer_requested; + struct cch_entry *entry; +}; + +struct tree_sink; +struct tree_list { + str name; + s32 set; + s32 current; + array(struct tree_list_entry *) entries; + array(struct tree_sink *) sinks; + struct tree_server *tree; +}; + +void tree_list_init(struct tree_list *list, struct tree_server *tree, str *name); +void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink); +void tree_list_remove_sink(struct tree_list *list, struct tree_sink *sink); +bool tree_list_add(struct tree_list *list, str *unique_id, u32 index); +void tree_list_skip(struct tree_list *list, s32 n); +void tree_list_pump(struct tree_list *list); diff --git a/src/tree/meson.build b/src/tree/meson.build index c7301f9..e9b8670 100644 --- a/src/tree/meson.build +++ b/src/tree/meson.build @@ -1,3 +1,7 @@ -tree_src = ['tree.c', 'resource_manager.c'] +tree_src = [ + 'tree.c', + 'resource_manager.c', + 'list.c' +] tree_deps = [common_deps, shoki, cache, cap, bimu, av] executable('tree', sources: tree_src, dependencies: tree_deps) diff --git a/src/tree/resource_manager.c b/src/tree/resource_manager.c index 0e95fbb..dd24310 100644 --- a/src/tree/resource_manager.c +++ b/src/tree/resource_manager.c @@ -18,7 +18,7 @@ static void packet_callback(void *userdata, struct aki_packet_stream *stream, st request->stream = stream; request->id = aki_packet_read_u16(packet); str unique_id; - aki_packet_read_string(packet, &unique_id); + aki_packet_read_str(packet, &unique_id); al_str_clone(&request->unique_id, &unique_id); request->index = aki_packet_read_u32(packet); request->packet = aki_packet_pool_get(&server->pool); @@ -40,8 +40,8 @@ static void packet_sent_callback(void *userdata, struct aki_packet *packet) { struct tree_resource_server *server = (struct tree_resource_server *)userdata; struct tree_resource_request *request = (struct tree_resource_request *)packet->userdata; - aki_packet_pool_return(&server->pool, packet); aki_http_request_close(&request->request); + aki_packet_pool_return(&server->pool, packet); } static void flushed_callback(void *userdata) @@ -70,6 +70,6 @@ bool tree_resource_server_init(struct tree_resource_server *server, void tree_resource_server_listen(struct tree_resource_server *server, struct aki_event_loop *loop, str *addr, s32 port) { - aki_packet_pool_init(&server->pool, 25, loop, packet_pool_callback, server); + aki_packet_pool_init(&server->pool, 50, loop, packet_pool_callback, server); aki_packet_stream_listen(&server->server, loop, addr, port); } diff --git a/src/tree/tree.c b/src/tree/tree.c index 0947be0..014397f 100644 --- a/src/tree/tree.c +++ b/src/tree/tree.c @@ -24,8 +24,12 @@ static void http_callback(void *userdata, struct aki_http_request *response, boo static bool request_resource(void *userdata, struct tree_resource_request *request) { struct tree_server *tree = (struct tree_server *)userdata; - struct sho_post *post = sho_post_cache_get(&tree->cache, &request->unique_id); - if (!post) return false; + struct sho_post *post = sho_post_cache_get(&tree->resources.cache, &request->unique_id); + struct aki_packet *packet = request->packet; + if (!post) { + aki_packet_pool_return(&request->server->pool, packet); + return false; + } str *url = NULL; if (al_str_eq(&post->author.unique_id, &request->unique_id)) { url = &post->author.profile_image_url; @@ -37,7 +41,10 @@ static bool request_resource(void *userdata, struct tree_resource_request *reque aki_http_request_init(&request->request); aki_http_set_url(&request->request.http, url); aki_http_set_user_agent(&request->request.http, USER_AGENT); - aki_http_request(&request->request, AKI_HTTP_GET, &tree->loop, http_callback, request); + if (!aki_http_request(&request->request, AKI_HTTP_GET, &tree->loop, http_callback, request)) { + aki_packet_pool_return(&request->server->pool, packet); + return false; + } return true; } @@ -66,6 +73,25 @@ static struct tree_user *get_user_by_connection(struct tree_server *tree, struct return NULL; } +static void send_status(struct tree_server *server, struct tree_user *user, struct aki_rpc_connection *conn) +{ + struct aki_packet *packet = aki_rpc_get_packet(&server->server, TREE_CMD_STATUS); + aki_packet_write_u32(packet, user->lists.size); + struct tree_list *list; + al_array_foreach(user->lists, i, list) { + aki_packet_write_str(packet, &list->name); + } + aki_packet_write_u32(packet, user->search.searches.size); + struct sho_search *search; + al_array_foreach(user->search.searches, i, search) { + aki_packet_write_s32(packet, search->id); + aki_packet_write_str(packet, &search->module); + aki_packet_write_str(packet, &search->query); + aki_packet_write_s32(packet, search->page); + } + aki_rpc_connection_command(conn, packet, NULL, NULL); +} + static bool identify_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { @@ -76,18 +102,21 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection struct tree_node *node = al_alloc_object(struct tree_node); node->conn = conn; al_array_push(tree->nodes, node); + al_log_info("tree", "New node."); break; } case TREE_CLIENT: { str username; - aki_packet_read_string(packet, &username); + aki_packet_read_str(packet, &username); struct tree_user *user = get_user_by_username(tree, &username); if (user) { - al_log_info("tree", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->username)); struct tree_client *client = al_alloc_object(struct tree_client); client->conn = conn; + client->user = user; al_array_push(tree->clients, client); al_array_push(user->clients, client); + send_status(tree, user, conn); + al_log_info("tree", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->username)); } break; } @@ -95,6 +124,10 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection struct tree_sink *sink = al_alloc_object(struct tree_sink); sink->conn = conn; al_array_push(tree->sinks, sink); + struct tree_user *user = al_array_at(tree->users, 0); + struct tree_list *list = al_array_at(user->lists, 0); + tree_list_add_sink(list, sink); + al_log_info("tree", "New sink."); break; } } @@ -102,22 +135,26 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection return true; } -static struct tree_query *query_from_id(struct tree_user *user, s32 id) +static bool search_command_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) { - struct tree_query *query; - al_array_foreach_ptr(user->queries, i, query) { - if (query->search.id == id) { - return query; - } + struct tree_server *tree = (struct tree_server *)userdata; + struct tree_user *user = get_user_by_connection(tree, conn); + if (!user) { + aki_packet_write_s32(rpacket, -1); + goto out; } - al_array_push(user->queries, (struct tree_query){}); - query = &al_array_last(user->queries); - query->id = id; - sho_search_init(&query->search); - return query; + str module, query; + aki_packet_read_str(packet, &module); + aki_packet_read_str(packet, &query); + s32 id = sho_client_create_search(&user->search, &module, &query); + aki_packet_write_s32(rpacket, id); +out: + aki_packet_free(packet); + return true; } -static bool search_command_callback(void *userdata, struct aki_rpc_connection *conn, +static bool resume_search_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { struct tree_server *tree = (struct tree_server *)userdata; @@ -127,66 +164,134 @@ static bool search_command_callback(void *userdata, struct aki_rpc_connection *c goto out; } s32 id = aki_packet_read_s32(packet); - struct tree_query *query = query_from_id(user, id); - str provider, query_str; - aki_packet_read_string(packet, &provider); - aki_packet_read_string(packet, &query_str); - s32 page_num = sho_search_more_results(&query->search, &provider, &query_str); - if (page_num < 0) { - al_log_error("tree", "Search failed."); + s32 req_page = aki_packet_read_s32(packet); + struct sho_search *search = sho_client_get_search(&user->search, id); + if (!search) { aki_packet_write_s32(rpacket, -1); goto out; } - struct sho_result_page *page = &al_array_at(query->search.result.pages, page_num); - aki_packet_write_s32(rpacket, query->search.id); - aki_packet_write_s32(rpacket, page_num); + if (req_page >= 0) { + if (!sho_search_from_page(search, req_page)) { + aki_packet_write_s32(rpacket, -1); + goto out; + } + } else { + if (!sho_search_more_results(search)) { + aki_packet_write_s32(rpacket, -1); + goto out; + } + } + aki_packet_write_s32(rpacket, 0); + aki_packet_write_s32(rpacket, search->id); + aki_packet_write_s32(rpacket, search->page); + struct sho_result_page *page = &al_array_at(search->pages, search->page); aki_packet_write_u32(rpacket, page->posts.size); struct sho_post *post; al_array_foreach_ptr(page->posts, i, post) { aki_packet_write_sho_post(rpacket, post); - sho_post_cache_push(&tree->cache, post); } aki_packet_write_u32(rpacket, page->list.size); str *unique_id; al_array_foreach_ptr(page->list, i, unique_id) { - aki_packet_write_string(rpacket, unique_id); + aki_packet_write_str(rpacket, unique_id); } out: aki_packet_free(packet); return true; } +static bool add_command_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) +{ + struct tree_server *tree = (struct tree_server *)userdata; + (void)rpacket; + struct tree_user *user = get_user_by_connection(tree, conn); + if (!user) goto out; + + str unique_id; + aki_packet_read_str(packet, &unique_id); + u32 index = aki_packet_read_u32(packet); + + struct tree_list *list = al_array_at(user->lists, 0); + tree_list_add(list, &unique_id, index); + tree_list_skip(list, 1); + +out: + aki_packet_free(packet); + return false; +} + static void connection_callback(void *userdata, struct aki_rpc_connection *conn) { (void)userdata; (void)conn; } +static void cleanup_node(struct tree_node *node) +{ + al_free(node); +} + +static void cleanup_client(struct tree_client *client) +{ + al_free(client); +} + +static void cleanup_sink(struct tree_sink *sink) +{ + al_free(sink); +} + static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) { - (void)userdata; - (void)conn; + struct tree_server *tree = (struct tree_server *)userdata; + + struct tree_node *node; + al_array_foreach(tree->nodes, i, node) { + if (node->conn == conn) { + al_log_info("tree", "Node removed."); + cleanup_node(node); + al_array_remove_at_iter(tree->nodes, i); + break; + } + } + + struct tree_client *client; + al_array_foreach(tree->clients, i, client) { + if (client->conn == conn) { + al_log_info("tree", "User \"%.*s\" logged out.", AL_STR_PRINTF(&client->user->username)); + cleanup_client(client); + al_array_remove_at_iter(tree->clients, i); + break; + } + } + + struct tree_sink *sink; + al_array_foreach(tree->sinks, i, sink) { + if (sink->conn == conn) { + al_log_info("tree", "Sink removed."); + cleanup_sink(sink); + al_array_remove_at_iter(tree->sinks, i); + struct tree_user *user = al_array_at(tree->users, 0); + struct tree_list *list = al_array_at(user->lists, 0); + tree_list_remove_sink(list, sink); + break; + } + } } static struct aki_rpc_command commands[] = { { .op = TREE_CMD_IDENTIFY, .callback = identify_command_callback, .userdata = NULL }, // Client commands. { .op = TREE_CMD_SEARCH, .callback = search_command_callback, .userdata = NULL }, + { .op = TREE_CMD_RESUME_SEARCH, .callback = resume_search_command_callback, .userdata = NULL }, + { .op = TREE_CMD_ADD, .callback = add_command_callback, .userdata = NULL } }; -static struct tree_server tree; - -void sigint_handler(s32 signum) -{ - (void)signum; - // explode. - exit(EXIT_SUCCESS); -} - static bool open_user(struct tree_server *tree, struct aki_dir_entry *dir) { struct aki_file file; - if (!aki_file_open(&file, &dir->path, false)) { + if (!aki_file_open(&file, &dir->path, 0)) { return false; } str s; @@ -195,10 +300,10 @@ static bool open_user(struct tree_server *tree, struct aki_dir_entry *dir) json_t *json = json_loadb(s.data, s.len, 0, &error); struct tree_user *user = al_alloc_object(struct tree_user); al_str_from(&user->username, json_string_value(json_object_get(json, "username"))); - al_array_init(user->queries); + sho_client_init(&user->search, &tree->resources.cache); al_array_init(user->lists); - struct tree_list default_list; - al_str_from(&default_list.name, "default"); + struct tree_list *default_list = al_alloc_object(struct tree_list); + tree_list_init(default_list, tree, al_str_c("default")); al_array_push(user->lists, default_list); al_array_init(user->clients); al_log_info("tree", "Loaded user \"%.*s\"", AL_STR_PRINTF(&user->username)); @@ -232,25 +337,15 @@ static bool open_db(struct tree_server *tree, str *path) return true; } -static bool cap_callback(void *userdata, u8 op, str *name, str *unique_id, void *opaque) +static void sigint_handler(s32 signum) { - (void)userdata; - (void)name; - (void)unique_id; - (void)opaque; - switch (op) { - case CAP_BUFFER: - return true; - case CAP_SWAP: - break; - case CAP_SET: - break; - case CAP_UNLOAD: - break; - } - return true; + (void)signum; + // explode. + exit(EXIT_SUCCESS); } +static struct tree_server tree = { 0 }; + s32 main(void) { aki_common_init(); @@ -266,7 +361,7 @@ s32 main(void) return EXIT_FAILURE; } - sho_post_cache_init(&tree.cache); + sho_post_cache_init(&tree.resources.cache); bool py_init = sho_python_init(); @@ -280,13 +375,11 @@ s32 main(void) } aki_rpc_listen(&tree.server, &tree.loop, al_str_c("0.0.0.0"), TREE_PORT); - tree_resource_server_init(&tree.resource_server, request_resource, &tree); - tree_resource_server_listen(&tree.resource_server, &tree.loop, al_str_c("0.0.0.0"), TREE_RESOURCE_PORT); - - cap_init(&tree.cap, cap_callback, &tree); + tree_resource_server_init(&tree.resources.server, request_resource, &tree); + tree_resource_server_listen(&tree.resources.server, &tree.loop, al_str_c("0.0.0.0"), TREE_RESOURCE_PORT); - bmu_server_init(&tree.stream_server); - bmu_server_listen(&tree.stream_server, &tree.loop, al_str_c("0.0.0.0"), TREE_STREAM_PORT); + bmu_server_init(&tree.streams.server); + bmu_server_listen(&tree.streams.server, &tree.loop, al_str_c("0.0.0.0"), TREE_STREAM_PORT); aki_event_loop_run(&tree.loop); diff --git a/src/tree/tree.h b/src/tree/tree.h index 06f98b1..478306a 100644 --- a/src/tree/tree.h +++ b/src/tree/tree.h @@ -1,47 +1,35 @@ #pragma once #include <aki/rpc2.h> +#include <sho/post.h> +#include <sho/post_cache.h> +#include <sho/search.h> -#include "../shoki/src/post_cache.h" -#include "../shoki/src/search.h" -#include "../shoki/src/packet_ext.h" #include "../bimu/server.h" -#include "../fruits/cap/cap.h" -#include "commands.h" +#include "common.h" #include "resource_manager.h" +#include "list.h" #define CAMU_DB_PATH "/home/andrew/c/camu/data/camu_db_test" -#define TREE_PORT 4356 -#define TREE_RESOURCE_PORT 4357 -#define TREE_STREAM_PORT 4358 - struct tree_node { struct aki_rpc_connection *conn; }; struct tree_client { struct aki_rpc_connection *conn; + struct tree_user *user; }; struct tree_sink { struct aki_rpc_connection *conn; }; -struct tree_query { - s32 id; - struct sho_search search; -}; - -struct tree_list { - str name; -}; - struct tree_user { str username; - array(struct tree_query) queries; - array(struct tree_list) lists; + struct sho_client search; + array(struct tree_list *) lists; array(struct tree_client *) clients; }; @@ -52,8 +40,11 @@ struct tree_server { array(struct tree_client *) clients; array(struct tree_sink *) sinks; array(struct tree_user *) users; - struct tree_resource_server resource_server; - struct sho_post_cache cache; - struct cap_runner cap; - struct bmu_server stream_server; + struct { + struct sho_post_cache cache; + struct tree_resource_server server; + } resources; + struct { + struct bmu_server server; + } streams; }; |