From 130a0edc0405e53b45f8b2f8da0b94356480a644 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 1 Jan 2024 11:23:22 -0500 Subject: wip Signed-off-by: Andrew Opalach --- src/libclient/client.c | 164 +++++++++++++++++++++++++++++++--------- src/libclient/client.h | 41 ++++++++-- src/libclient/meson.build | 2 +- src/libclient/resource_client.c | 33 ++++---- src/libclient/resource_client.h | 7 +- src/libclient/search.c | 31 -------- src/libclient/search.h | 28 ------- 7 files changed, 182 insertions(+), 124 deletions(-) delete mode 100644 src/libclient/search.c delete mode 100644 src/libclient/search.h (limited to 'src/libclient') diff --git a/src/libclient/client.c b/src/libclient/client.c index 60998c6..e42f662 100644 --- a/src/libclient/client.c +++ b/src/libclient/client.c @@ -1,6 +1,3 @@ -#include "../tree/commands.h" -#include "../shoki/src/packet_ext.h" - #include "client.h" static void identifed_callback(void *userdata, struct aki_packet *packet) @@ -16,7 +13,7 @@ static void connection_callback(void *userdata, struct aki_rpc_connection *conn) client->conn = conn; struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_IDENTIFY); aki_packet_write_u8(packet, TREE_CLIENT); - aki_packet_write_string(packet, &client->username); + aki_packet_write_str(packet, &client->username); aki_rpc_connection_command(client->conn, packet, identifed_callback, client); } @@ -33,66 +30,161 @@ static u8 packet_pool_callback(void *userdata, struct aki_packet *packet) return AKI_PACKET_POOL_NOP; } -bool camu_client_init(struct camu_client *client) -{ - aki_event_loop_init(&client->loop); - aki_packet_pool_init_ex(&client->pool, 0, &client->loop, - AKI_PACKET_POOL_MODE_PASSTHROUGH, packet_pool_callback, client); - return aki_rpc_init(&client->client, AKI_SOCKET_TCP, connection_callback, - connection_closed_callback, client); -} - -static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) +static bool status_command_callback(void *userdata, struct aki_rpc_connection *conn, + struct aki_packet *packet, struct aki_packet *rpacket) { struct camu_client *client = (struct camu_client *)userdata; - aki_event_loop_run(&client->loop); - return 0; + (void)conn; + (void)rpacket; + + str s; + u32 size = aki_packet_read_u32(packet); + for (u32 i = 0; i < size; i++) { + struct camu_list *list = al_alloc_object(struct camu_list); + aki_packet_read_str(packet, &s); + al_str_clone(&list->name, &s); + al_array_push(client->state.lists, list); + } + + size = aki_packet_read_u32(packet); + for (u32 i = 0; i < size; i++) { + struct camu_search *search = al_alloc_object(struct camu_search); + search->id = aki_packet_read_s32(packet); + aki_packet_read_str(packet, &s); + al_str_clone(&search->provider, &s); + aki_packet_read_str(packet, &s); + al_str_clone(&search->query, &s); + search->page = aki_packet_read_s32(packet); + al_array_push(client->state.searches, search); + } + + client->callback(client->userdata, CAMU_CLIENT_STATUS_UPDATED, NULL); + + return false; } -bool camu_client_login(struct camu_client *client, str *username, str *addr, s32 port, +static struct aki_rpc_command commands[] = { + { .op = TREE_CMD_STATUS, .callback = status_command_callback, .userdata = NULL } +}; + +bool camu_client_init(struct camu_client *client, struct aki_event_loop *loop, void (*callback)(void *, u8, void *), void *userdata) { - al_str_clone(&client->username, username); + client->loop = loop; client->callback = callback; client->userdata = userdata; - if (!aki_rpc_connect(&client->client, &client->loop, addr, port)) { + client->username = AL_STR_EMPTY; + al_array_init(client->state.lists); + aki_packet_pool_init_ex(&client->pool, 0, client->loop, + AKI_PACKET_POOL_MODE_PASSTHROUGH, packet_pool_callback, client); + if (!aki_rpc_init(&client->client, AKI_SOCKET_TCP, connection_callback, + connection_closed_callback, client)) { + return false; + } + for (u32 i = 0; i < AL_ARRAY_SIZE(commands); i++) { + commands[i].userdata = client; + aki_rpc_add_command(&client->client, &commands[i]); + } + sho_post_cache_init(&client->cache); + return true; +} + +bool camu_client_login(struct camu_client *client, str *username, str *addr, s32 port) +{ + al_str_clone(&client->username, username); + if (!aki_rpc_connect(&client->client, client->loop, addr, port)) { return false; } - aki_thread_create(&client->thread, event_loop_thread, client); return true; } -static void search_callback(void *userdata, struct aki_packet *packet) +static void create_search_callback(void *userdata, struct aki_packet *packet) { struct camu_client *client = (struct camu_client *)userdata; - struct camu_search_results *results = al_alloc_object(struct camu_search_results); - al_array_init(results->results); - s32 search_id = aki_packet_read_s32(packet); - if (search_id == -1) return; - results->page = aki_packet_read_s32(packet); + s32 id = aki_packet_read_s32(packet); + client->callback(client->userdata, CAMU_CLIENT_SEARCH_CREATED, &id); +} + +void camu_client_create_search(struct camu_client *client, str *provider, str *query) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_SEARCH); + aki_packet_write_str(packet, provider); + aki_packet_write_str(packet, query); + packet->userdata = create_search_callback; + aki_packet_pool_submit(&client->pool, packet); +} + +static struct camu_search *search_from_id(struct camu_client *client, s32 id) +{ + struct camu_search *search; + al_array_foreach(client->state.searches, i, search) { + if (search->id == id) return search; + } + return NULL; +} + +static void results_callback(void *userdata, struct aki_packet *packet) +{ + struct camu_client *client = (struct camu_client *)userdata; + s32 ok = aki_packet_read_s32(packet); + if (ok != 0) return; + s32 id = aki_packet_read_s32(packet); + struct camu_search *search = search_from_id(client, id); + if (!search) return; + s32 page = aki_packet_read_s32(packet); + while (search->pages.size <= (u32)page) { + struct camu_search_results result; + result.page = page; + al_array_init(result.unique_ids); + al_array_push(search->pages, result); + } u32 size = aki_packet_read_u32(packet); for (u32 i = 0; i < size; i++) { struct sho_post post; aki_packet_read_sho_post(packet, &post); - sho_post_cache_push(&client->search.cache, &post); + sho_post_cache_push(&client->cache, &post); } + struct camu_search_results *result = &al_array_at(search->pages, page); + str s, unique_id; size = aki_packet_read_u32(packet); for (u32 i = 0; i < size; i++) { - str s, unique_id; - aki_packet_read_string(packet, &s); + aki_packet_read_str(packet, &s); al_str_clone(&unique_id, &s); - al_array_push(results->results, unique_id); + al_array_push(result->unique_ids, unique_id); } - client->callback(client->userdata, CAMU_CLIENT_RESULTS, results); - aki_packet_free(packet); + client->callback(client->userdata, CAMU_CLIENT_RESULTS, result); +} + +void camu_client_resume_search(struct camu_client *client, struct camu_search *search) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_RESUME_SEARCH); + aki_packet_write_s32(packet, search->id); + aki_packet_write_s32(packet, search->page); + packet->userdata = results_callback; + aki_packet_pool_submit(&client->pool, packet); } void camu_client_more_results(struct camu_client *client, struct camu_search *search) { - struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_SEARCH); + struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_RESUME_SEARCH); aki_packet_write_s32(packet, search->id); - aki_packet_write_string(packet, &search->provider); - aki_packet_write_string(packet, &search->query); - packet->userdata = search_callback; + aki_packet_write_s32(packet, -1); + packet->userdata = results_callback; aki_packet_pool_submit(&client->pool, packet); } + +void camu_client_add(struct camu_client *client, str *unique_id, u32 index) +{ + struct aki_packet *packet = aki_rpc_get_packet(&client->client, TREE_CMD_ADD); + aki_packet_write_str(packet, unique_id); + aki_packet_write_u32(packet, index); + packet->userdata = NULL; + aki_packet_pool_submit(&client->pool, packet); +} + +void camu_client_close(struct camu_client *client) +{ + aki_packet_pool_free(&client->pool); + al_str_free(&client->username); + al_array_free(client->state.lists); +} diff --git a/src/libclient/client.h b/src/libclient/client.h index 989c839..cdcaffe 100644 --- a/src/libclient/client.h +++ b/src/libclient/client.h @@ -3,29 +3,56 @@ #include #include #include +#include +#include -#include "search.h" +#include "../tree/common.h" enum { CAMU_CLIENT_CONNECTED = 0, - // Search client. + CAMU_CLIENT_STATUS_UPDATED, + CAMU_CLIENT_SEARCH_CREATED, CAMU_CLIENT_RESULTS }; +struct camu_list { + str name; + array(str) entries; +}; + +struct camu_search_results { + s32 page; + array(str) unique_ids; +}; + +struct camu_search { + s32 id; + str provider; + str query; + s32 page; + array(struct camu_search_results) pages; +}; + struct camu_client { - struct aki_event_loop loop; + struct aki_event_loop *loop; struct aki_rpc client; str username; + struct { + array(struct camu_list *) lists; + array(struct camu_search *) searches; + } state; + struct sho_post_cache cache; struct aki_rpc_connection *conn; - struct camu_search_client search; struct aki_packet_pool pool; - struct aki_thread thread; void (*callback)(void *, u8, void *); void *userdata; }; -bool camu_client_init(struct camu_client *client); -bool camu_client_login(struct camu_client *client, str *username, str *addr, s32 port, +bool camu_client_init(struct camu_client *client, struct aki_event_loop *loop, void (*callback)(void *, u8, void *), void *userdata); +bool camu_client_login(struct camu_client *client, str *username, str *addr, s32 port); +void camu_client_create_search(struct camu_client *client, str *provider, str *query); +void camu_client_resume_search(struct camu_client *client, struct camu_search *search); void camu_client_more_results(struct camu_client *client, struct camu_search *search); +void camu_client_add(struct camu_client *client, str *unique_id, u32 index); void camu_client_close(struct camu_client *client); diff --git a/src/libclient/meson.build b/src/libclient/meson.build index 99eabb2..a2f23b0 100644 --- a/src/libclient/meson.build +++ b/src/libclient/meson.build @@ -1,3 +1,3 @@ -libclient_src = ['client.c', 'search.c', 'resource_client.c'] +libclient_src = ['client.c', 'resource_client.c'] libclient_deps = [shoki] libclient = declare_dependency(sources: libclient_src, dependencies: libclient_deps) diff --git a/src/libclient/resource_client.c b/src/libclient/resource_client.c index 1069755..77c19c7 100644 --- a/src/libclient/resource_client.c +++ b/src/libclient/resource_client.c @@ -28,6 +28,7 @@ static void packet_callback(void *userdata, struct aki_packet_stream *stream, aki_buffer_init(&buffer); aki_packet_read_buffer(packet, &buffer); client->callback(client->userdata, id, &buffer); + aki_packet_free(packet); } static void packet_sent_callback(void *userdata, struct aki_packet *packet) @@ -42,41 +43,37 @@ static void signal_callback(void *userdata) u32 size; struct camu_resource_request request; do { - camu_queue_size(client->queue, size); - if (!size) break; - camu_queue_pop(client->queue, request); + camu_queue_try_pop(client->queue, size, request); + if (size == 0) break; struct aki_packet *packet = aki_packet_create(); aki_packet_write_u16(packet, request.id); - aki_packet_write_string(packet, &request.unique_id); + aki_packet_write_str(packet, &request.unique_id); aki_packet_write_u32(packet, request.index); aki_packet_stream_send_packet(&client->client, packet); } while (1); } -static aki_thread_result AKI_THREADCALL event_loop_thread(void *userdata) -{ - struct camu_resource_client *client = (struct camu_resource_client *)userdata; - aki_event_loop_run(&client->loop); - return 0; -} - -void camu_resource_client_run(struct camu_resource_client *client, - void (*callback)(void *, u16, struct aki_buffer *), void *userdata) +void camu_resource_client_init(struct camu_resource_client *client, + struct aki_event_loop *loop, void (*callback)(void *, u16, struct aki_buffer *), void *userdata) { + client->loop = loop; client->callback = callback; client->userdata = userdata; camu_queue_init(client->queue); - aki_event_loop_init(&client->loop); aki_signal_init(&client->signal, signal_callback, client); - aki_signal_start(&client->signal, &client->loop); aki_packet_stream_init(&client->client, AKI_SOCKET_TCP, connection_callback, connection_closed_callback, packet_callback, packet_sent_callback, client); - aki_packet_stream_connect(&client->client, &client->loop, al_str_c("127.0.0.1"), TREE_RESOURCE_PORT); +} + + +void camu_resource_client_connect(struct camu_resource_client *client, str *addr, s32 port) +{ + aki_signal_start(&client->signal, client->loop); + aki_packet_stream_connect(&client->client, client->loop, addr, port); client->connected = false; while (!client->connected) { - aki_event_loop_run_once(&client->loop); + aki_event_loop_run_once(client->loop); } - aki_thread_create(&client->thread, event_loop_thread, client); } void camu_resource_request(struct camu_resource_client *client, diff --git a/src/libclient/resource_client.h b/src/libclient/resource_client.h index 4cc1bbf..17ac31e 100644 --- a/src/libclient/resource_client.h +++ b/src/libclient/resource_client.h @@ -13,7 +13,7 @@ struct camu_resource_request { }; struct camu_resource_client { - struct aki_event_loop loop; + struct aki_event_loop *loop; struct aki_packet_stream client; bool connected; struct aki_thread thread; @@ -23,7 +23,8 @@ struct camu_resource_client { void *userdata; }; -void camu_resource_client_run(struct camu_resource_client *client, - void (*callback)(void *, u16, struct aki_buffer *), void *userdata); +void camu_resource_client_init(struct camu_resource_client *client, + struct aki_event_loop *loop, void (*callback)(void *, u16, struct aki_buffer *), void *userdata); +void camu_resource_client_connect(struct camu_resource_client *client, str *addr, s32 port); void camu_resource_request(struct camu_resource_client *client, str *unique_id, u32 index, u16 id); diff --git a/src/libclient/search.c b/src/libclient/search.c deleted file mode 100644 index f74f83d..0000000 --- a/src/libclient/search.c +++ /dev/null @@ -1,31 +0,0 @@ -#include "search.h" - -void camu_search_client_init(struct camu_search_client *client) -{ - client->id = 0; - sho_post_cache_init(&client->cache); -} - -void camu_search_set(struct camu_search_client *client, struct camu_search *search, - str *provider, str *query) -{ - search->id = client->id++; - al_str_clone(&search->provider, provider); - al_str_clone(&search->query, query); -} - -void camu_search_free(struct camu_search *search) -{ - al_str_free(&search->provider); - al_str_free(&search->query); -} - -void camu_search_results_free(struct camu_search_results *results) -{ - str *unique_id; - al_array_foreach_ptr(results->results, i, unique_id) { - al_str_free(unique_id); - } - al_array_free(results->results); - al_free(results); -} diff --git a/src/libclient/search.h b/src/libclient/search.h deleted file mode 100644 index b0d8554..0000000 --- a/src/libclient/search.h +++ /dev/null @@ -1,28 +0,0 @@ -#pragma once - -#include -#include - -#include "../shoki/src/post_cache.h" - -struct camu_search { - s32 id; - str provider; - str query; -}; - -struct camu_search_results { - u32 page; - array(str) results; -}; - -struct camu_search_client { - s32 id; - struct sho_post_cache cache; -}; - -void camu_search_client_init(struct camu_search_client *client); -void camu_search_set(struct camu_search_client *client, struct camu_search *search, - str *provider, str *query); -void camu_search_free(struct camu_search *search); -void camu_search_results_free(struct camu_search_results *results); -- cgit v1.2.3-101-g0448