From 5e3641e5e692c3f2f644a4bb809c88727cb8bee9 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Fri, 8 Nov 2024 14:53:40 -0500 Subject: Command queue for portal and list, work on server Most of the server stuff can undoubtedly be simplified. I'm still working that out. Signed-off-by: Andrew Opalach --- src/portal/meson.build | 2 +- src/portal/src/packet_ext.c | 4 +- src/portal/src/search.c | 215 +++++++++++++++++++++++++++++++++----------- src/portal/src/search.h | 47 ++++++++-- 4 files changed, 207 insertions(+), 61 deletions(-) (limited to 'src/portal') diff --git a/src/portal/meson.build b/src/portal/meson.build index 7cf31a7..4e5ff56 100644 --- a/src/portal/meson.build +++ b/src/portal/meson.build @@ -8,7 +8,7 @@ portal_deps = [] portal_args = ['-DCAMU_HAVE_PORTAL'] cpy_dir = join_paths(meson.current_source_dir(), 'cpy') -run_command(join_paths(cpy_dir, 'build.sh'), cpy_dir) +run_command(join_paths(cpy_dir, 'build.sh'), cpy_dir, check: false) python3_embed = import('python').find_installation('python3.12').dependency(embed: true) portal_deps += [python3_embed] diff --git a/src/portal/src/packet_ext.c b/src/portal/src/packet_ext.c index e2909ef..833e7d9 100644 --- a/src/portal/src/packet_ext.c +++ b/src/portal/src/packet_ext.c @@ -6,7 +6,7 @@ static void aki_packet_write_optional_int(struct aki_packet *packet, optional_in AKI_PACKET_WRITE_TYPE(packet, bool, o->set); } -void aki_packet_write_camu_post(struct aki_packet *packet, struct camu_post *post) +void aki_packet_write_post(struct aki_packet *packet, struct camu_post *post) { AKI_PACKET_WRITE_TYPE(packet, u16, post->version); AKI_PACKET_WRITE_TYPE(packet, u8, post->type); @@ -53,7 +53,7 @@ static void aki_packet_read_optional_int(struct aki_packet *packet, optional_int AKI_PACKET_READ_TYPE(packet, bool, o->set); } -void aki_packet_read_camu_post(struct aki_packet *packet, struct camu_post *post) +void aki_packet_read_post(struct aki_packet *packet, struct camu_post *post) { camu_post_reset(post); AKI_PACKET_READ_TYPE(packet, u16, post->version); diff --git a/src/portal/src/search.c b/src/portal/src/search.c index bcb0643..e6022f5 100644 --- a/src/portal/src/search.c +++ b/src/portal/src/search.c @@ -1,5 +1,7 @@ #include +#include "../../server/common.h" + #include "../cpy/portal.c" #include "search.h" @@ -32,49 +34,168 @@ void camu_python_close(void) if (Py_IsInitialized()) Py_Finalize(); } -static struct camu_result_page *page_at_index(struct camu_search *search, u32 num) +static struct camu_search *get_search_by_id(struct camu_portal_bridge *bridge, s32 id) { - struct camu_result_page *page; - al_array_foreach_ptr(search->pages, i, page) { - if (page->num == num) return page; + struct camu_search *search; + al_array_foreach(bridge->searches, i, search) { + if (search->id == id) return search; } - al_array_push(search->pages, (struct camu_result_page){ 0 }); - page = &al_array_last(search->pages); - page->num = num; - al_array_init(page->posts); - al_array_init(page->list); - return page; + return NULL; } +static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) +{ + struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; + aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS); + bool have_python = false; + aki_mutex_lock(&bridge->mutex); + do { + aki_cond_wait(&bridge->cond, &bridge->mutex); + if (bridge->quit) break; + if (!have_python) { + // Defer python init. + have_python = camu_python_init(); + } + struct camu_portal_cmd *cmd; + al_array_foreach_ptr(bridge->queue, i, cmd) { + struct camu_portal_result result; + result.op = cmd->op; + result.callback = cmd->callback; + result.userdata = cmd->userdata; + switch (cmd->op) { + case CAMU_CLIENT_CREATE_SEARCH: { + s32 id = portal_bridge_search(&cmd->module, &cmd->query); + if (id >= 0) { + struct camu_search *search = al_alloc_object(struct camu_search); + search->page = 0; + al_array_init(search->pages); + search->id = id; + al_str_clone(&search->module, &cmd->module); + al_str_clone(&search->query, &cmd->query); + search->bridge = bridge; + al_array_push(bridge->searches, search); + al_log_info("portal", "New search %x (%.*s).", id, AL_STR_PRINTF(&cmd->query)); + result.id = search->id; + } else { + } + al_str_free(&cmd->module); + al_str_free(&cmd->query); + break; + } + case CAMU_CLIENT_GET_PAGE: { + struct camu_search *search = get_search_by_id(bridge, cmd->id); + if (search) { + result.id = search->id; + struct camu_result_page *page; + al_array_foreach_ptr(search->pages, j, page) { + if (page->num == cmd->num) break; + } + al_log_info("portal", "Loading page %i (%.*s).", cmd->num, AL_STR_PRINTF(&search->query)); + if (portal_bridge_get_page(search, search->id, cmd->num) == -1) { + break; + } + page = &al_array_at(search->pages, cmd->num); + if (bridge->cache) { + struct camu_post *post; + al_array_foreach_ptr(page->posts, j, post) { + camu_post_cache_push(bridge->cache, post); + } + } + result.page = page; + } + break; + } + } + camu_queue_push(bridge->results, result); + aki_signal_send(&bridge->results_signal); + al_array_remove_at_iter(bridge->queue, i); + } + } while (1); + aki_mutex_unlock(&bridge->mutex); + if (have_python) { + camu_python_close(); + } + return 0; +} -void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache) +static void results_signal_callback(void *userdata) +{ + struct camu_portal_bridge *bridge = (struct camu_portal_bridge *)userdata; + u32 size; + struct camu_portal_result result; + do { + camu_queue_try_pop(bridge->results, size, result); + if (size == 0) break; + result.callback(result.userdata, &result); + } while (1); +} + +void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, + struct aki_event_loop *loop) { - bridge->cache = cache; al_array_init(bridge->searches); + bridge->cache = cache; + bridge->quit = 0; + aki_mutex_init(&bridge->mutex); + aki_cond_init(&bridge->cond); + al_array_init(bridge->queue); + camu_queue_init(bridge->results); + aki_signal_init(&bridge->results_signal, results_signal_callback, bridge); + aki_signal_start(&bridge->results_signal, loop); + aki_thread_create(&bridge->thread, queue_thread, bridge); +} + +void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, + void (*callback)(void *, struct camu_portal_result *), void *userdata) +{ + struct camu_portal_cmd cmd; + cmd.op = CAMU_CLIENT_CREATE_SEARCH; + al_str_clone(&cmd.module, module); + al_str_clone(&cmd.query, query); + cmd.callback = callback; + cmd.userdata = userdata; + aki_mutex_lock(&bridge->mutex); + al_array_push(bridge->queue, cmd); + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); + } + aki_mutex_unlock(&bridge->mutex); } -static void camu_search_init_internal(struct camu_search *search) +void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, + void (*callback)(void *, struct camu_portal_result *), void *userdata) { - search->page = 0; - al_array_init(search->pages); + struct camu_portal_cmd cmd; + cmd.op = CAMU_CLIENT_GET_PAGE; + cmd.id = id; + cmd.num = num; + cmd.callback = callback; + cmd.userdata = userdata; + aki_mutex_lock(&bridge->mutex); + al_array_push(bridge->queue, cmd); + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); + } + aki_mutex_unlock(&bridge->mutex); } -s32 camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query) +void camu_portal_close(struct camu_portal_bridge *bridge) { - s32 id = portal_bridge_search(module, query); - if (id >= 0) { - struct camu_search *search = al_alloc_object(struct camu_search); - camu_search_init_internal(search); - search->id = id; - al_str_clone(&search->module, module); - al_str_clone(&search->query, query); - search->bridge = bridge; - al_array_push(bridge->searches, search); + aki_mutex_lock(&bridge->mutex); + bridge->quit = 1; + if (aki_cond_is_waiting(&bridge->cond)) { + aki_cond_signal(&bridge->cond); } - al_log_info("portal", "New search %x (%.*s).", id, AL_STR_PRINTF(query)); - return id; + aki_mutex_unlock(&bridge->mutex); + aki_thread_join(&bridge->thread); + aki_signal_stop(&bridge->results_signal); + camu_queue_free(bridge->results); + al_array_free(bridge->queue); + aki_cond_destroy(&bridge->cond); + aki_mutex_destroy(&bridge->mutex); } +/* struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id) { struct camu_search *search; @@ -90,42 +211,30 @@ void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id) (void)id; } -void camu_portal_close(struct camu_portal_bridge *bridge) -{ - (void)bridge; -} - -bool camu_search_get_page(struct camu_search *search, u32 num) +void camu_search_free(struct camu_search *search) { struct camu_result_page *page; al_array_foreach_ptr(search->pages, i, page) { - if (page->num == num) goto out; - } - al_log_info("portal", "Loading page %i (%.*s).", num, AL_STR_PRINTF(&search->query)); - if (portal_bridge_get_page(search, search->id, num) == -1) { - return false; - } - struct camu_portal_bridge *bridge = search->bridge; - if (bridge->cache) { - struct camu_post *post; - al_array_foreach_ptr(al_array_at(search->pages, num).posts, i, post) { - camu_post_cache_push(bridge->cache, post); - } + // TODO: Free camu_post ? + al_array_free(page->posts); + al_array_free(page->list); } -out: - search->page = num; - return true; + al_array_free(search->pages); } +*/ -void camu_search_free(struct camu_search *search) +static struct camu_result_page *page_at_index(struct camu_search *search, u32 num) { struct camu_result_page *page; al_array_foreach_ptr(search->pages, i, page) { - // TODO: Free camu_post ? - al_array_free(page->posts); - al_array_free(page->list); + if (page->num == num) return page; } - al_array_free(search->pages); + al_array_push(search->pages, (struct camu_result_page){ 0 }); + page = &al_array_last(search->pages); + page->num = num; + al_array_init(page->posts); + al_array_init(page->list); + return page; } void camu_search_add_post(struct camu_search *search, u32 num, struct camu_post *post) diff --git a/src/portal/src/search.h b/src/portal/src/search.h index 1f43ded..85f4f6d 100644 --- a/src/portal/src/search.h +++ b/src/portal/src/search.h @@ -1,8 +1,13 @@ #pragma once +#include +#include + #include "post.h" #include "post_cache.h" +#include "../../util/queue.h" + struct camu_result_page { u32 num; array(struct camu_post) posts; @@ -18,21 +23,53 @@ struct camu_search { struct camu_portal_bridge *bridge; }; +struct camu_portal_result { + u8 op; + s32 id; + struct camu_result_page *page; + void (*callback)(void *, struct camu_portal_result *); + void *userdata; +}; + +struct camu_portal_cmd { + u8 op; + str module; + str query; + s32 id; + u32 num; + void (*callback)(void *, struct camu_portal_result *); + void *userdata; +}; + struct camu_portal_bridge { array(struct camu_search *) searches; struct camu_post_cache *cache; + u8 quit; + struct aki_thread thread; + struct aki_mutex mutex; + struct aki_cond cond; + array(struct camu_portal_cmd) queue; + queue(struct camu_portal_result) results; + struct aki_signal results_signal; }; bool camu_python_init(void); void camu_python_close(void); -void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache); -s32 camu_portal_create_search(struct camu_portal_bridge *bridge, str *module_str, str *search_str); -struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id); -void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id); +void camu_portal_init(struct camu_portal_bridge *bridge, struct camu_post_cache *cache, + struct aki_event_loop *loop); + +void camu_portal_create_search(struct camu_portal_bridge *bridge, str *module, str *query, + void (*callback)(void *, struct camu_portal_result *), void *userdata); +void camu_portal_get_page(struct camu_portal_bridge *bridge, s32 id, u32 num, + void (*callback)(void *, struct camu_portal_result *), void *userdata); + void camu_portal_close(struct camu_portal_bridge *bridge); -bool camu_search_get_page(struct camu_search *search, u32 num); +/* +struct camu_search *camu_portal_get_search(struct camu_portal_bridge *bridge, s32 id); +void camu_portal_discard_search(struct camu_portal_bridge *bridge, s32 id); +*/ // Python internal. void camu_search_add_post(struct camu_search *search, u32 num, struct camu_post *post); -- cgit v1.2.3-101-g0448