summaryrefslogtreecommitdiff
path: root/src/tree
diff options
context:
space:
mode:
Diffstat (limited to 'src/tree')
-rw-r--r--src/tree/commands.h13
-rw-r--r--src/tree/meson.build3
-rw-r--r--src/tree/resource_manager.c75
-rw-r--r--src/tree/resource_manager.h27
-rw-r--r--src/tree/tree.c298
-rw-r--r--src/tree/tree.h59
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;
+};