summaryrefslogtreecommitdiff
path: root/src/tree
diff options
context:
space:
mode:
Diffstat (limited to 'src/tree')
-rw-r--r--src/tree/common.h13
-rw-r--r--src/tree/list.c178
-rw-r--r--src/tree/list.h13
-rw-r--r--src/tree/tree.c97
-rw-r--r--src/tree/tree.h2
5 files changed, 228 insertions, 75 deletions
diff --git a/src/tree/common.h b/src/tree/common.h
index 59c713d..de6258b 100644
--- a/src/tree/common.h
+++ b/src/tree/common.h
@@ -6,8 +6,8 @@
#define TREE_RESOURCE_PORT 14357
#define TREE_STREAM_PORT 14358
-//#define TREE_SERVER_IP al_str_c("127.0.0.1")
-#define TREE_SERVER_IP al_str_c("108.52.160.112")
+#define TREE_SERVER_IP al_str_c("127.0.0.1")
+//#define TREE_SERVER_IP al_str_c("108.52.160.112")
enum {
TREE_NODE = 0,
@@ -17,8 +17,9 @@ enum {
enum {
TREE_CMD_IDENTIFY = 0,
- TREE_CMD_STATUS,
- TREE_CMD_SEARCH,
- TREE_CMD_RESUME_SEARCH,
- TREE_CMD_ADD
+ TREE_CMD_UPDATE_STATE,
+ TREE_CMD_CREATE_SEARCH,
+ TREE_CMD_GET_PAGE,
+ TREE_CMD_ADD,
+ TREE_CMD_SKIP
};
diff --git a/src/tree/list.c b/src/tree/list.c
index c1630b7..948d3cc 100644
--- a/src/tree/list.c
+++ b/src/tree/list.c
@@ -1,3 +1,5 @@
+#include <al/random.h>
+
#ifdef AKIYO_HAS_CURL
#include "../cache/handlers/http.h"
#endif
@@ -14,11 +16,13 @@ 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;
+ list->backwards = false;
al_array_init(list->entries);
al_array_init(list->sinks);
list->tree = tree;
}
+/*
static void send_buffer_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry)
{
struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_BUFFER);
@@ -27,20 +31,33 @@ static void send_buffer_cmd(struct tree_list *list, struct tree_sink *sink, stru
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_sink *sink, struct tree_list_entry *entry)
{
struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_SET);
+ aki_packet_write_str(packet, TREE_SERVER_IP);
+ 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_queue_cmd(struct tree_list *list, struct tree_sink *sink, struct tree_list_entry *entry)
+{
+ struct aki_packet *packet = aki_rpc_get_packet(&list->tree->server, CAMU_SINK_CMD_QUEUE);
+ aki_packet_write_str(packet, TREE_SERVER_IP);
+ 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);
}
+*/
void tree_list_add_sink(struct tree_list *list, struct tree_sink *sink)
{
al_array_push(list->sinks, sink);
if (list->set == list->current) {
- struct tree_list_entry *entry = al_array_at(list->entries, list->current);
- send_buffer_cmd(list, sink, entry);
+ struct tree_list_entry *entry = &al_array_at(list->entries, list->current);
send_set_cmd(list, sink, entry);
}
}
@@ -62,21 +79,16 @@ bool tree_list_add(struct tree_list *list, str *unique_id, u32 index)
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);
+ struct tree_list_entry entry;
+ entry.index = index;
+ al_str_clone(&entry.unique_id, unique_id);
+ entry.post = post;
- entry->node_id = bmu_server_create_node(&tree->streams.server, entry->entry);
+ entry.entry = cch_handler_http_create(&al_array_at(entry.post->media, entry.index).url);
+ if (!entry.entry) return false;
+ entry.entry->handler->maybe_spawn_worker(entry.entry->handler, 0);
- entry->buffer_requested = false;
+ entry.node_id = bmu_server_create_node(&tree->streams.server, entry.entry);
al_array_push(list->entries, entry);
@@ -85,21 +97,79 @@ bool tree_list_add(struct tree_list *list, str *unique_id, u32 index)
return true;
}
+static struct cch_entry *entry_for_external_path(struct sho_client *client, str *path)
+{
+ struct cch_entry *entry = NULL;
+#ifdef AKIYO_HAS_CURL
+ bool is_search = al_str_at(path, 0) == ';';
+ if (is_search || al_str_cmp(path, al_str_c("https://"), 0, 8) == 0 ||
+ al_str_cmp(path, al_str_c("http://"), 0, 7) == 0) {
+ str module;
+ str query;
+ al_str_from(&module, "");
+ al_str_from(&query, "");
+ if (al_str_cmp(path, al_str_c("https://twitter.com"), 0, 19) == 0 ||
+ al_str_cmp(path, al_str_c("https://x.com"), 0, 13) == 0) {
+ al_str_cat(&query, al_str_c("tweet:"));
+ al_str_cat(&query, path);
+ al_str_cat(&module, al_str_c("twitter"));
+ } else if (al_str_cmp(path, al_str_c("https://instagram.com"), 0, 21) == 0) {
+ s32 slash = al_str_rfind(path, '/');
+ if (slash >= 0) {
+ al_str_cat(&query, al_str_substr(path, slash + 1, path->len));
+ }
+ al_str_cat(&module, al_str_c("instagram"));
+ } else {
+ if (is_search) {
+ al_str_cat(&query, al_str_substr(path, 1, path->len));
+ } else {
+ al_str_cat(&query, al_str_c("link:"));
+ al_str_cat(&query, path);
+ }
+ al_str_cat(&module, al_str_c("youtube"));
+ }
+ s32 id = sho_client_create_search(client, &module, &query);
+ al_str_free(&module);
+ al_str_free(&query);
+ if (id < 0) return NULL;
+ struct sho_search *search = sho_client_get_search(client, id);
+ if (!search || !sho_search_get_page(search, 0)) return NULL;
+ struct sho_result_page *page = &al_array_at(search->pages, 0);
+ struct sho_post *post;
+ al_array_foreach_ptr(page->posts, i, post) {
+ struct sho_post_media *media;
+ al_array_foreach_ptr(post->media, j, media) {
+ if (media->url.len > 0) {
+ entry = cch_handler_http_create(&media->url);
+ break;
+ }
+ }
+ if (entry) break;
+ }
+ sho_client_discard_search(client, id);
+ if (entry) entry->handler->maybe_spawn_worker(entry->handler, 0);
+ } else // {
+#endif
+ entry = cch_handler_file_create(path);
+ // }
+ return entry;
+}
+
bool tree_list_add_external(struct tree_list *list, str *path)
{
struct tree_server *tree = list->tree;
- struct tree_list_entry *entry = al_alloc_object(struct tree_list_entry);
+ struct tree_list_entry entry;
+ entry.index = 0;
+ al_str_clone(&entry.unique_id, path);
+ entry.post = NULL;
- entry->entry = cch_handler_file_create(path);
- if (!entry->entry) {
- al_free(entry);
- return false;
- }
+ struct tree_user *user = al_array_at(tree->users, 0);
- entry->node_id = bmu_server_create_node(&tree->streams.server, entry->entry);
+ entry.entry = entry_for_external_path(&user->search, path);
+ if (!entry.entry) return false;
- entry->buffer_requested = false;
+ entry.node_id = bmu_server_create_node(&tree->streams.server, entry.entry);
al_array_push(list->entries, entry);
@@ -115,25 +185,63 @@ void tree_list_skip(struct tree_list *list, s32 n)
return;
}
list->current += n;
+ list->backwards = n < 0;
tree_list_pump(list);
}
+static s32 sort_list_func(void *_a, void *_b)
+{
+ struct tree_list_entry *a = (struct tree_list_entry *)_a;
+ struct tree_list_entry *b = (struct tree_list_entry *)_b;
+ return al_str_cmp(&a->unique_id, &b->unique_id, 0, a->unique_id.len);
+}
+
+void tree_list_sort(struct tree_list *list)
+{
+ al_array_sort(list->entries, struct tree_list_entry, sort_list_func);
+}
+
+void tree_list_shuffle(struct tree_list *list)
+{
+ u32 size = list->entries.size;
+ if (size == 0) return;
+ for (u32 i = 0; i < size - 1; i++) {
+ u32 j = i + al_rand() / (AL_RAND_MAX / (size - i) + 1);
+ struct tree_list_entry tmp = al_array_at(list->entries, j);
+ al_array_at(list->entries, j) = al_array_at(list->entries, i);
+ al_array_at(list->entries, i) = tmp;
+ }
+ list->set = -1;
+ tree_list_pump(list);
+}
+
+void tree_list_clear(struct tree_list *list)
+{
+ //list->entries.size = 0;
+}
+
void tree_list_pump(struct tree_list *list)
{
s32 size = (s32)list->entries.size;
if (size <= list->current) return;
struct tree_sink *sink;
- if (list->set != list->current) {
- struct tree_list_entry *entry = al_array_at(list->entries, list->current);
- if (!entry->buffer_requested) {
- al_array_foreach(list->sinks, i, sink) {
- send_buffer_cmd(list, sink, entry);
- }
- entry->buffer_requested = true;
- }
- al_array_foreach(list->sinks, i, sink) {
- send_set_cmd(list, sink, entry);
- }
- list->set = list->current;
+ if (list->current + 1 < size) {
+ // This needs a massive rethinking on how to sync the lists.
+ // - Query sinks on server side skip request?
+ //struct tree_list_entry *upcoming = &al_array_at(list->entries, list->current + 1);
+ //al_array_foreach(list->sinks, i, sink) {
+ // send_queue_cmd(list, sink, upcoming);
+ //}
}
+ if (list->set == list->current) return;
+ struct tree_list_entry *entry = &al_array_at(list->entries, list->current);
+ al_array_foreach(list->sinks, i, sink) {
+ send_set_cmd(list, sink, entry);
+ }
+ list->set = list->current;
+}
+
+void tree_list_free(struct tree_list *list)
+{
+ al_str_free(&list->name);
}
diff --git a/src/tree/list.h b/src/tree/list.h
index 3defe9f..3702983 100644
--- a/src/tree/list.h
+++ b/src/tree/list.h
@@ -6,11 +6,11 @@
#include "../cache/entry.h"
struct tree_list_entry {
- struct sho_post *post;
u32 index;
- u16 node_id;
- bool buffer_requested;
+ str unique_id;
+ struct sho_post *post;
struct cch_entry *entry;
+ u16 node_id;
};
struct tree_sink;
@@ -18,7 +18,8 @@ struct tree_list {
str name;
s32 set;
s32 current;
- array(struct tree_list_entry *) entries;
+ bool backwards;
+ array(struct tree_list_entry) entries;
array(struct tree_sink *) sinks;
struct tree_server *tree;
};
@@ -29,4 +30,8 @@ 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);
bool tree_list_add_external(struct tree_list *list, str *path);
void tree_list_skip(struct tree_list *list, s32 n);
+void tree_list_sort(struct tree_list *list);
+void tree_list_shuffle(struct tree_list *list);
+void tree_list_clear(struct tree_list *list);
void tree_list_pump(struct tree_list *list);
+void tree_list_free(struct tree_list *list);
diff --git a/src/tree/tree.c b/src/tree/tree.c
index 367e0ff..8dcc501 100644
--- a/src/tree/tree.c
+++ b/src/tree/tree.c
@@ -68,17 +68,33 @@ static struct tree_user *get_user_by_connection(struct tree_server *tree, struct
return user;
}
}
+ struct tree_list *list;
+ al_array_foreach(user->lists, j, list) {
+ struct tree_sink *sink;
+ al_array_foreach(list->sinks, k, sink) {
+ if (sink->conn == conn) {
+ return user;
+ }
+ }
+ }
}
+
return NULL;
}
-static void send_status(struct tree_server *server, struct tree_user *user, struct aki_rpc_connection *conn)
+static void send_current_state(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);
+ struct aki_packet *packet = aki_rpc_get_packet(&server->server, TREE_CMD_UPDATE_STATE);
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, list->entries.size);
+ struct tree_list_entry *entry;
+ al_array_foreach_ptr(list->entries, i, entry) {
+ aki_packet_write_str(packet, &entry->unique_id);
+ }
+ aki_packet_write_s32(packet, list->current);
}
aki_packet_write_u32(packet, user->search.searches.size);
struct sho_search *search;
@@ -86,7 +102,7 @@ static void send_status(struct tree_server *server, struct tree_user *user, stru
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_packet_write_u32(packet, search->last_page);
}
aki_rpc_connection_command(conn, packet, NULL, NULL);
}
@@ -114,7 +130,7 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection
client->user = user;
al_array_push(tree->clients, client);
al_array_push(user->clients, client);
- send_status(tree, user, conn);
+ send_current_state(tree, user, conn);
al_log_info("tree", "User \"%.*s\" logged in.", AL_STR_PRINTF(&user->username));
}
break;
@@ -134,7 +150,7 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection
return true;
}
-static bool search_command_callback(void *userdata, struct aki_rpc_connection *conn,
+static bool create_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;
@@ -143,17 +159,18 @@ static bool search_command_callback(void *userdata, struct aki_rpc_connection *c
aki_packet_write_s32(rpacket, -1);
goto out;
}
+
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);
+ aki_packet_write_s32(rpacket, sho_client_create_search(&user->search, &module, &query));
+
out:
aki_packet_free(packet);
return true;
}
-static bool resume_search_command_callback(void *userdata, struct aki_rpc_connection *conn,
+static bool get_page_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;
@@ -162,28 +179,25 @@ static bool resume_search_command_callback(void *userdata, struct aki_rpc_connec
aki_packet_write_s32(rpacket, -1);
goto out;
}
+
s32 id = aki_packet_read_s32(packet);
- s32 req_page = aki_packet_read_s32(packet);
+ u32 page_request = aki_packet_read_u32(packet);
+
struct sho_search *search = sho_client_get_search(&user->search, id);
if (!search) {
aki_packet_write_s32(rpacket, -1);
goto out;
}
- 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;
- }
+
+ if (!sho_search_get_page(search, page_request)) {
+ 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_request);
+ struct sho_result_page *page = &al_array_at(search->pages, page_request);
aki_packet_write_u32(rpacket, page->posts.size);
struct sho_post *post;
al_array_foreach_ptr(page->posts, i, post) {
@@ -194,6 +208,7 @@ static bool resume_search_command_callback(void *userdata, struct aki_rpc_connec
al_array_foreach_ptr(page->list, i, unique_id) {
aki_packet_write_str(rpacket, unique_id);
}
+
out:
aki_packet_free(packet);
return true;
@@ -220,6 +235,24 @@ out:
return false;
}
+static bool skip_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;
+
+ s32 n = aki_packet_read_s32(packet);
+
+ struct tree_list *list = al_array_at(user->lists, 0);
+ tree_list_skip(list, n);
+
+out:
+ aki_packet_free(packet);
+ return false;
+}
+
static void connection_callback(void *userdata, struct aki_rpc_connection *conn)
{
(void)userdata;
@@ -282,9 +315,10 @@ static void connection_closed_callback(void *userdata, struct aki_rpc_connection
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 }
+ { .op = TREE_CMD_CREATE_SEARCH, .callback = create_search_command_callback, .userdata = NULL },
+ { .op = TREE_CMD_GET_PAGE, .callback = get_page_command_callback, .userdata = NULL },
+ { .op = TREE_CMD_ADD, .callback = add_command_callback, .userdata = NULL },
+ { .op = TREE_CMD_SKIP, .callback = skip_command_callback, .userdata = NULL }
};
static bool open_user(struct tree_server *tree, struct aki_dir_entry *dir)
@@ -347,6 +381,11 @@ static u8 line_callback(void *userdata, str *line)
} else if (al_str_eq(line, al_str_c(";PREV"))) {
tree_list_skip(list, -1);
} else if (al_str_eq(line, al_str_c(";SHUFFLE"))) {
+ tree_list_shuffle(list);
+ } else if (al_str_eq(line, al_str_c(";SORT"))) {
+ tree_list_sort(list);
+ } else if (al_str_eq(line, al_str_c(";CLEAR"))) {
+ tree_list_clear(list);
} else {
tree_list_add_external(list, line);
}
@@ -402,12 +441,12 @@ s32 main(void)
tree.socket.type = AKI_SOCKET_UNIX;
aki_socket_init(&tree.socket);
aki_socket_set_blocking(&tree.socket, false);
- tree.pro.callback = line_callback;
- tree.pro.userdata = &tree;
- aki_line_processor_init(&tree.pro, al_str_c("\n"));
- aki_line_processor_open_socket(&tree.pro, &tree.socket);
+ tree.cli.callback = line_callback;
+ tree.cli.userdata = &tree;
+ aki_line_processor_init(&tree.cli, al_str_c("\n"));
+ aki_line_processor_open_socket(&tree.cli, &tree.socket);
if (aki_socket_listen(&tree.socket, al_str_c("/tmp/tree_sock"), 0)) {
- aki_line_processor_run(&tree.pro, &tree.loop);
+ aki_line_processor_run(&tree.cli, &tree.loop);
}
#endif
diff --git a/src/tree/tree.h b/src/tree/tree.h
index 6339889..9f8a19a 100644
--- a/src/tree/tree.h
+++ b/src/tree/tree.h
@@ -53,6 +53,6 @@ struct tree_server {
} streams;
#if TREE_USE_SOCKET
struct aki_socket socket;
- struct aki_line_processor pro;
+ struct aki_line_processor cli;
#endif
};