summaryrefslogtreecommitdiff
path: root/src/portal
diff options
context:
space:
mode:
Diffstat (limited to 'src/portal')
-rw-r--r--src/portal/meson.build2
-rw-r--r--src/portal/src/packet_ext.c4
-rw-r--r--src/portal/src/search.c215
-rw-r--r--src/portal/src/search.h47
4 files changed, 207 insertions, 61 deletions
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 <al/log.h>
+#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 <aki/thread.h>
+#include <aki/signal.h>
+
#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);