summaryrefslogtreecommitdiff
path: root/src/tree
diff options
context:
space:
mode:
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.c112
-rw-r--r--src/tree/list.h31
-rw-r--r--src/tree/meson.build6
-rw-r--r--src/tree/resource_manager.c6
-rw-r--r--src/tree/tree.c227
-rw-r--r--src/tree/tree.h39
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;
};