diff options
| author | 2023-11-06 11:54:08 -0500 | |
|---|---|---|
| committer | 2023-11-06 11:54:08 -0500 | |
| commit | a48a68cfb04b2020737c0adfb0e7667451a22b5c (patch) | |
| tree | a92121897f84298c3798a4f2ba56de45cd100eab /src/tree | |
| download | camu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.tar.gz camu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.tar.bz2 camu-a48a68cfb04b2020737c0adfb0e7667451a22b5c.zip | |
Add camu
- Subprojects temporarily omitted
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/tree')
| -rw-r--r-- | src/tree/commands.h | 13 | ||||
| -rw-r--r-- | src/tree/meson.build | 3 | ||||
| -rw-r--r-- | src/tree/resource_manager.c | 75 | ||||
| -rw-r--r-- | src/tree/resource_manager.h | 27 | ||||
| -rw-r--r-- | src/tree/tree.c | 298 | ||||
| -rw-r--r-- | src/tree/tree.h | 59 |
6 files changed, 475 insertions, 0 deletions
diff --git a/src/tree/commands.h b/src/tree/commands.h new file mode 100644 index 0000000..1b1187e --- /dev/null +++ b/src/tree/commands.h @@ -0,0 +1,13 @@ +#pragma once + +enum { + TREE_NODE = 0, + TREE_CLIENT, + TREE_SINK +}; + +enum { + TREE_CMD_IDENTIFY = 0, + TREE_CMD_SEARCH, + TREE_CMD_ADD +}; diff --git a/src/tree/meson.build b/src/tree/meson.build new file mode 100644 index 0000000..c7301f9 --- /dev/null +++ b/src/tree/meson.build @@ -0,0 +1,3 @@ +tree_src = ['tree.c', 'resource_manager.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 new file mode 100644 index 0000000..0e95fbb --- /dev/null +++ b/src/tree/resource_manager.c @@ -0,0 +1,75 @@ +#include <al/random.h> +#include <aki/http.h> + +#include "resource_manager.h" + +static u8 packet_pool_callback(void *userdata, struct aki_packet *packet) +{ + (void)userdata; + struct tree_resource_request *request = (struct tree_resource_request *)packet->userdata; + aki_packet_stream_send_packet(request->stream, packet); + return AKI_PACKET_POOL_KEEP; +} + +static void packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) +{ + struct tree_resource_server *server = (struct tree_resource_server *)userdata; + struct tree_resource_request *request = al_alloc_object(struct tree_resource_request); + request->stream = stream; + request->id = aki_packet_read_u16(packet); + str unique_id; + aki_packet_read_string(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); + request->packet->userdata = request; + request->server = server; + if (!server->request_resource(server->userdata, request)) { + aki_packet_pool_return(&server->pool, request->packet); + } + aki_packet_free(packet); +} + +static void connection_closed_callback(void *userdata, struct aki_packet_stream *stream) +{ + (void)userdata; + (void)stream; +} + +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); +} + +static void flushed_callback(void *userdata) +{ + (void)userdata; +} + +static void connection_callback(void *userdata, struct aki_packet_stream *stream) +{ + stream->userdata = userdata; + stream->packet_callback = packet_callback; + stream->packet_sent_callback = packet_sent_callback; + stream->connection_closed_callback = connection_closed_callback; + stream->flushed_callback = flushed_callback; +} + +bool tree_resource_server_init(struct tree_resource_server *server, + bool (*request_resource)(void *, struct tree_resource_request *), void *userdata) +{ + server->request_resource = request_resource; + server->userdata = userdata; + return aki_packet_stream_init(&server->server, AKI_SOCKET_TCP, connection_callback, + NULL, NULL, NULL, 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_stream_listen(&server->server, loop, addr, port); +} diff --git a/src/tree/resource_manager.h b/src/tree/resource_manager.h new file mode 100644 index 0000000..a81e3a9 --- /dev/null +++ b/src/tree/resource_manager.h @@ -0,0 +1,27 @@ +#pragma once + +#include <aki/packet_stream.h> +#include <aki/packet_pool.h> +#include <aki/http.h> + +struct tree_resource_request { + u16 id; + str unique_id; + u32 index; + struct aki_packet *packet; + struct aki_packet_stream *stream; + struct aki_http_request request; + struct tree_resource_server *server; +}; + +struct tree_resource_server { + struct aki_packet_stream server; + struct aki_packet_pool pool; + bool (*request_resource)(void *, struct tree_resource_request *); + void *userdata; +}; + +bool tree_resource_server_init(struct tree_resource_server *server, + bool (*request_resource)(void *, struct tree_resource_request *), void *userdata); +void tree_resource_server_listen(struct tree_resource_server *server, + struct aki_event_loop *loop, str *addr, s32 port); diff --git a/src/tree/tree.c b/src/tree/tree.c new file mode 100644 index 0000000..0947be0 --- /dev/null +++ b/src/tree/tree.c @@ -0,0 +1,298 @@ +#include <al/log.h> +#include <aki/file.h> +#include <jansson.h> + +#include "../libclient/commands.h" + +#include "tree.h" + +#define USER_AGENT al_str_c("Mozilla/5.0 (X11; Linux x86_64; rv:96.0) Gecko/20100101 Firefox/96.0") + +static void http_callback(void *userdata, struct aki_http_request *response, bool success) +{ + struct tree_resource_request *request = (struct tree_resource_request *)userdata; + struct aki_packet *packet = request->packet; + if (!success) { + aki_packet_pool_return(&request->server->pool, packet); + return; + } + aki_packet_write_u16(packet, request->id); + aki_packet_write_buffer(packet, &response->response); + aki_packet_pool_submit(&request->server->pool, packet); +} + +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; + str *url = NULL; + if (al_str_eq(&post->author.unique_id, &request->unique_id)) { + url = &post->author.profile_image_url; + } else { + if (post->media.size <= request->index) return false; + struct sho_post_media *media = &al_array_at(post->media, request->index); + url = &media->thumbnail_url; + } + 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); + return true; +} + +static struct tree_user *get_user_by_username(struct tree_server *tree, str *username) +{ + struct tree_user *user; + al_array_foreach(tree->users, i, user) { + if (al_str_eq(&user->username, username)) { + return user; + } + } + return NULL; +} + +static struct tree_user *get_user_by_connection(struct tree_server *tree, struct aki_rpc_connection *conn) +{ + struct tree_user *user; + al_array_foreach(tree->users, i, user) { + struct tree_client *client; + al_array_foreach(user->clients, j, client) { + if (client->conn == conn) { + return user; + } + } + } + return NULL; +} + +static bool identify_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; + switch (aki_packet_read_u8(packet)) { + case TREE_NODE: { + struct tree_node *node = al_alloc_object(struct tree_node); + node->conn = conn; + al_array_push(tree->nodes, node); + break; + } + case TREE_CLIENT: { + str username; + aki_packet_read_string(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; + al_array_push(tree->clients, client); + al_array_push(user->clients, client); + } + break; + } + case TREE_SINK: { + struct tree_sink *sink = al_alloc_object(struct tree_sink); + sink->conn = conn; + al_array_push(tree->sinks, sink); + break; + } + } + aki_packet_free(packet); + return true; +} + +static struct tree_query *query_from_id(struct tree_user *user, s32 id) +{ + struct tree_query *query; + al_array_foreach_ptr(user->queries, i, query) { + if (query->search.id == id) { + return query; + } + } + al_array_push(user->queries, (struct tree_query){}); + query = &al_array_last(user->queries); + query->id = id; + sho_search_init(&query->search); + return query; +} + +static bool 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; + struct tree_user *user = get_user_by_connection(tree, conn); + if (!user) { + aki_packet_write_s32(rpacket, -1); + 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."); + 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); + 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); + } +out: + aki_packet_free(packet); + return true; +} + +static void connection_callback(void *userdata, struct aki_rpc_connection *conn) +{ + (void)userdata; + (void)conn; +} + +static void connection_closed_callback(void *userdata, struct aki_rpc_connection *conn) +{ + (void)userdata; + (void)conn; +} + +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 }, +}; + +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)) { + return false; + } + str s; + aki_file_read_as_str(&file, &s); + json_error_t error; + 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); + al_array_init(user->lists); + struct tree_list default_list; + al_str_from(&default_list.name, "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)); + al_array_push(tree->users, user); + return true; +} + +static bool open_db(struct tree_server *tree, str *path) +{ + struct aki_dir camu_db; + if (!aki_dir_open(&camu_db, path)) { + return false; + } + struct aki_dir_entry entry; + while (aki_dir_read(&camu_db, &entry)) { + if (al_str_eq(&entry.name, al_str_c("users"))) { + struct aki_dir users; + if (aki_dir_open(&users, &entry.path)) { + struct aki_dir_entry user; + while (aki_dir_read(&users, &user)) { + if (user.type == AKI_ENTRY_FILE) { + open_user(tree, &user); + } + } + aki_dir_close(&users); + } + } + aki_dir_entry_free(&entry); + } + aki_dir_close(&camu_db); + return true; +} + +static bool cap_callback(void *userdata, u8 op, str *name, str *unique_id, void *opaque) +{ + (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; +} + +s32 main(void) +{ + aki_common_init(); + + signal(SIGINT, sigint_handler); + + al_array_init(tree.nodes); + al_array_init(tree.clients); + al_array_init(tree.sinks); + al_array_init(tree.users); + + if (!open_db(&tree, al_str_c(CAMU_DB_PATH))) { + return EXIT_FAILURE; + } + + sho_post_cache_init(&tree.cache); + + bool py_init = sho_python_init(); + + aki_event_loop_init(&tree.loop); + + aki_rpc_init(&tree.server, AKI_SOCKET_TCP, connection_callback, + connection_closed_callback, &tree); + for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { + commands[i].userdata = &tree; + aki_rpc_add_command(&tree.server, &commands[i]); + } + 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); + + bmu_server_init(&tree.stream_server); + bmu_server_listen(&tree.stream_server, &tree.loop, al_str_c("0.0.0.0"), TREE_STREAM_PORT); + + aki_event_loop_run(&tree.loop); + + if (py_init) sho_python_close(); + + aki_common_close(); + + return EXIT_SUCCESS; +} diff --git a/src/tree/tree.h b/src/tree/tree.h new file mode 100644 index 0000000..06f98b1 --- /dev/null +++ b/src/tree/tree.h @@ -0,0 +1,59 @@ +#pragma once + +#include <aki/rpc2.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 "resource_manager.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_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; + array(struct tree_client *) clients; +}; + +struct tree_server { + struct aki_event_loop loop; + struct aki_rpc server; + array(struct tree_node *) nodes; + 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; +}; |