diff options
Diffstat (limited to 'src/tree')
| -rw-r--r-- | src/tree/common.h | 13 | ||||
| -rw-r--r-- | src/tree/list.c | 178 | ||||
| -rw-r--r-- | src/tree/list.h | 13 | ||||
| -rw-r--r-- | src/tree/tree.c | 97 | ||||
| -rw-r--r-- | src/tree/tree.h | 2 |
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 }; |