845 lines
26 KiB
C
845 lines
26 KiB
C
#define _GNU_SOURCE
|
|
|
|
#include "thread_pool.h"
|
|
#include "debug.h"
|
|
#include "db_ops.h"
|
|
|
|
#include <sqlite3.h>
|
|
#include <pthread.h>
|
|
#include <stdint.h>
|
|
#include <stdio.h>
|
|
#include <stdlib.h>
|
|
#include <string.h>
|
|
|
|
typedef struct thread_pool_job_node {
|
|
uint64_t job_id;
|
|
thread_pool_job_t job;
|
|
struct thread_pool_job_node* next;
|
|
} thread_pool_job_node_t;
|
|
|
|
typedef struct {
|
|
thread_pool_job_node_t* head;
|
|
thread_pool_job_node_t* tail;
|
|
int size;
|
|
int max_size;
|
|
pthread_mutex_t mutex;
|
|
pthread_cond_t cond;
|
|
} thread_pool_queue_t;
|
|
|
|
typedef struct {
|
|
int running;
|
|
uint64_t next_job_id;
|
|
|
|
int reader_count;
|
|
pthread_t* readers;
|
|
pthread_t writer;
|
|
|
|
thread_pool_queue_t read_q;
|
|
thread_pool_queue_t write_q;
|
|
|
|
thread_pool_wake_loop_cb wake_loop_cb;
|
|
void* wake_loop_ctx;
|
|
|
|
char db_path[512];
|
|
|
|
pthread_mutex_t state_mutex;
|
|
} thread_pool_state_t;
|
|
|
|
static thread_pool_state_t g_pool = {0};
|
|
|
|
typedef struct {
|
|
pthread_mutex_t mutex;
|
|
pthread_cond_t cond;
|
|
int done;
|
|
thread_pool_status_t status;
|
|
void* result_data;
|
|
size_t result_size;
|
|
} thread_pool_wait_ctx_t;
|
|
|
|
static void wait_ctx_init(thread_pool_wait_ctx_t* ctx) {
|
|
memset(ctx, 0, sizeof(*ctx));
|
|
pthread_mutex_init(&ctx->mutex, NULL);
|
|
pthread_cond_init(&ctx->cond, NULL);
|
|
}
|
|
|
|
static void wait_ctx_destroy(thread_pool_wait_ctx_t* ctx) {
|
|
pthread_mutex_destroy(&ctx->mutex);
|
|
pthread_cond_destroy(&ctx->cond);
|
|
}
|
|
|
|
static void wait_ctx_result_cb(const thread_pool_result_t* result, void* user_ctx) {
|
|
thread_pool_wait_ctx_t* ctx = (thread_pool_wait_ctx_t*)user_ctx;
|
|
if (!ctx || !result) return;
|
|
|
|
pthread_mutex_lock(&ctx->mutex);
|
|
ctx->status = result->status;
|
|
ctx->result_data = result->result_data;
|
|
ctx->result_size = result->result_size;
|
|
ctx->done = 1;
|
|
pthread_cond_signal(&ctx->cond);
|
|
pthread_mutex_unlock(&ctx->mutex);
|
|
}
|
|
|
|
static void wait_ctx_wait(thread_pool_wait_ctx_t* ctx) {
|
|
pthread_mutex_lock(&ctx->mutex);
|
|
while (!ctx->done) {
|
|
pthread_cond_wait(&ctx->cond, &ctx->mutex);
|
|
}
|
|
pthread_mutex_unlock(&ctx->mutex);
|
|
}
|
|
|
|
static void free_count_payload(void* p) {
|
|
thread_pool_count_payload_t* payload = (thread_pool_count_payload_t*)p;
|
|
if (!payload) return;
|
|
free(payload->sql);
|
|
if (payload->bind_params) {
|
|
for (int i = 0; i < payload->bind_param_count; i++) {
|
|
free(payload->bind_params[i]);
|
|
}
|
|
free(payload->bind_params);
|
|
}
|
|
free(payload);
|
|
}
|
|
|
|
static void free_req_payload(void* p) {
|
|
thread_pool_req_payload_t* payload = (thread_pool_req_payload_t*)p;
|
|
if (!payload) return;
|
|
free(payload->sql);
|
|
if (payload->bind_params) {
|
|
for (int i = 0; i < payload->bind_param_count; i++) {
|
|
free(payload->bind_params[i]);
|
|
}
|
|
free(payload->bind_params);
|
|
}
|
|
free(payload);
|
|
}
|
|
|
|
static void free_store_event_payload(void* p) {
|
|
thread_pool_store_event_payload_t* payload = (thread_pool_store_event_payload_t*)p;
|
|
if (!payload) return;
|
|
free(payload->id);
|
|
free(payload->pubkey);
|
|
free(payload->event_type);
|
|
free(payload->content);
|
|
free(payload->sig);
|
|
free(payload->tags_json);
|
|
free(payload->event_json);
|
|
free(payload);
|
|
}
|
|
|
|
static void queue_init(thread_pool_queue_t* q, int max_size) {
|
|
memset(q, 0, sizeof(*q));
|
|
q->max_size = (max_size > 0) ? max_size : 4096;
|
|
pthread_mutex_init(&q->mutex, NULL);
|
|
pthread_cond_init(&q->cond, NULL);
|
|
}
|
|
|
|
static void queue_destroy(thread_pool_queue_t* q) {
|
|
pthread_mutex_lock(&q->mutex);
|
|
thread_pool_job_node_t* cur = q->head;
|
|
while (cur) {
|
|
thread_pool_job_node_t* next = cur->next;
|
|
if (cur->job.payload_free && cur->job.payload) {
|
|
cur->job.payload_free(cur->job.payload);
|
|
}
|
|
free(cur);
|
|
cur = next;
|
|
}
|
|
q->head = q->tail = NULL;
|
|
q->size = 0;
|
|
pthread_mutex_unlock(&q->mutex);
|
|
|
|
pthread_mutex_destroy(&q->mutex);
|
|
pthread_cond_destroy(&q->cond);
|
|
}
|
|
|
|
static int queue_push(thread_pool_queue_t* q, thread_pool_job_node_t* node) {
|
|
pthread_mutex_lock(&q->mutex);
|
|
if (q->size >= q->max_size) {
|
|
pthread_mutex_unlock(&q->mutex);
|
|
return 0;
|
|
}
|
|
|
|
node->next = NULL;
|
|
if (!q->tail) {
|
|
q->head = q->tail = node;
|
|
} else {
|
|
q->tail->next = node;
|
|
q->tail = node;
|
|
}
|
|
q->size++;
|
|
pthread_cond_signal(&q->cond);
|
|
pthread_mutex_unlock(&q->mutex);
|
|
return 1;
|
|
}
|
|
|
|
static thread_pool_job_node_t* queue_pop(thread_pool_queue_t* q) {
|
|
pthread_mutex_lock(&q->mutex);
|
|
while (g_pool.running && q->size == 0) {
|
|
pthread_cond_wait(&q->cond, &q->mutex);
|
|
}
|
|
|
|
if (!g_pool.running && q->size == 0) {
|
|
pthread_mutex_unlock(&q->mutex);
|
|
return NULL;
|
|
}
|
|
|
|
thread_pool_job_node_t* node = q->head;
|
|
q->head = node->next;
|
|
if (!q->head) q->tail = NULL;
|
|
q->size--;
|
|
pthread_mutex_unlock(&q->mutex);
|
|
return node;
|
|
}
|
|
|
|
static void wake_event_loop(void) {
|
|
if (g_pool.wake_loop_cb) {
|
|
g_pool.wake_loop_cb(g_pool.wake_loop_ctx);
|
|
}
|
|
}
|
|
|
|
static sqlite3* open_worker_connection(void) {
|
|
if (g_pool.db_path[0] == '\0') {
|
|
return NULL;
|
|
}
|
|
|
|
sqlite3* db = NULL;
|
|
int rc = sqlite3_open_v2(g_pool.db_path, &db, SQLITE_OPEN_READWRITE, NULL);
|
|
if (rc != SQLITE_OK) {
|
|
if (db) {
|
|
sqlite3_close(db);
|
|
}
|
|
return NULL;
|
|
}
|
|
|
|
sqlite3_exec(db, "PRAGMA journal_mode=WAL;", NULL, NULL, NULL);
|
|
sqlite3_busy_timeout(db, 5000);
|
|
|
|
return db;
|
|
}
|
|
|
|
static void complete_job_with_result(thread_pool_job_node_t* node,
|
|
thread_pool_status_t status,
|
|
const char* message,
|
|
void* result_data,
|
|
size_t result_size) {
|
|
thread_pool_result_t result;
|
|
memset(&result, 0, sizeof(result));
|
|
result.job_id = node->job_id;
|
|
result.type = node->job.type;
|
|
result.status = status;
|
|
result.session = node->job.session;
|
|
result.result_data = result_data;
|
|
result.result_size = result_size;
|
|
result.message = message;
|
|
|
|
if (node->job.result_cb) {
|
|
node->job.result_cb(&result, node->job.result_cb_ctx);
|
|
}
|
|
wake_event_loop();
|
|
|
|
if (node->job.payload_free && node->job.payload) {
|
|
node->job.payload_free(node->job.payload);
|
|
}
|
|
free(node);
|
|
}
|
|
|
|
static void execute_count_job(thread_pool_job_node_t* node) {
|
|
thread_pool_count_payload_t* payload = (thread_pool_count_payload_t*)node->job.payload;
|
|
if (!payload || !payload->sql) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"COUNT payload missing", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
int count = 0;
|
|
if (db_count_with_sql(payload->sql, (const char**)payload->bind_params,
|
|
payload->bind_param_count, &count) != 0) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"COUNT query failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
thread_pool_count_result_t* out = calloc(1, sizeof(*out));
|
|
if (!out) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"COUNT result allocation failed", NULL, 0);
|
|
return;
|
|
}
|
|
out->count = count;
|
|
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_OK,
|
|
"COUNT query completed", out, sizeof(*out));
|
|
}
|
|
|
|
static void execute_req_job(thread_pool_job_node_t* node) {
|
|
thread_pool_req_payload_t* payload = (thread_pool_req_payload_t*)node->job.payload;
|
|
if (!payload || !payload->sql) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"REQ payload missing", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
db_stmt_t* stmt = NULL;
|
|
if (db_prepare(payload->sql, &stmt) != DB_OK) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"REQ prepare failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
for (int i = 0; i < payload->bind_param_count; i++) {
|
|
const char* v = (payload->bind_params && payload->bind_params[i]) ? payload->bind_params[i] : "";
|
|
(void)db_bind_text_param(stmt, i + 1, v);
|
|
}
|
|
|
|
int cap = (payload->row_limit > 0) ? payload->row_limit : 500;
|
|
if (cap < 1) cap = 1;
|
|
|
|
thread_pool_req_result_t* out = calloc(1, sizeof(*out));
|
|
if (!out) {
|
|
db_finalize_stmt(stmt);
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"REQ result allocation failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
out->event_json_rows = calloc((size_t)cap, sizeof(char*));
|
|
if (!out->event_json_rows) {
|
|
db_finalize_stmt(stmt);
|
|
free(out);
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"REQ row array allocation failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
int rc = DB_OK;
|
|
while ((rc = db_step_stmt(stmt)) == DB_ROW) {
|
|
if (out->row_count >= cap) break;
|
|
const char* row = db_column_text_value(stmt, 0);
|
|
if (!row) continue;
|
|
|
|
out->event_json_rows[out->row_count] = strdup(row);
|
|
if (!out->event_json_rows[out->row_count]) {
|
|
break;
|
|
}
|
|
out->row_count++;
|
|
}
|
|
|
|
db_finalize_stmt(stmt);
|
|
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_OK,
|
|
"REQ query completed", out, sizeof(*out));
|
|
}
|
|
|
|
static void execute_store_event_job(thread_pool_job_node_t* node) {
|
|
thread_pool_store_event_payload_t* payload = (thread_pool_store_event_payload_t*)node->job.payload;
|
|
if (!payload || !payload->id || !payload->pubkey || !payload->content || !payload->sig) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"STORE_EVENT payload missing required fields", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
thread_pool_store_event_result_t* out = calloc(1, sizeof(*out));
|
|
if (!out) {
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"STORE_EVENT result allocation failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
int step_rc = DB_ERROR;
|
|
int extended_errcode = 0;
|
|
if (db_insert_event_with_json(payload->id,
|
|
payload->pubkey,
|
|
payload->created_at,
|
|
payload->kind,
|
|
payload->event_type ? payload->event_type : "regular",
|
|
payload->content,
|
|
payload->sig,
|
|
payload->tags_json ? payload->tags_json : "[]",
|
|
payload->event_json ? payload->event_json : "{}",
|
|
&step_rc,
|
|
&extended_errcode) != 0) {
|
|
free(out);
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_INTERNAL_ERROR,
|
|
"STORE_EVENT execution failed", NULL, 0);
|
|
return;
|
|
}
|
|
|
|
out->step_rc = step_rc;
|
|
out->extended_errcode = extended_errcode;
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_OK,
|
|
"STORE_EVENT completed", out, sizeof(*out));
|
|
}
|
|
|
|
static void execute_read_job(thread_pool_job_node_t* node) {
|
|
switch (node->job.type) {
|
|
case THREAD_POOL_JOB_REQ_QUERY:
|
|
execute_req_job(node);
|
|
break;
|
|
case THREAD_POOL_JOB_COUNT_QUERY:
|
|
execute_count_job(node);
|
|
break;
|
|
default:
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_NOT_IMPLEMENTED,
|
|
"Read job type not implemented", NULL, 0);
|
|
break;
|
|
}
|
|
}
|
|
|
|
static void execute_write_job(thread_pool_job_node_t* node) {
|
|
switch (node->job.type) {
|
|
case THREAD_POOL_JOB_STORE_EVENT:
|
|
execute_store_event_job(node);
|
|
break;
|
|
default:
|
|
complete_job_with_result(node, THREAD_POOL_STATUS_NOT_IMPLEMENTED,
|
|
"Write job type not implemented", NULL, 0);
|
|
break;
|
|
}
|
|
}
|
|
|
|
static void* reader_worker_main(void* arg) {
|
|
int reader_index = (int)(intptr_t)arg;
|
|
|
|
char thread_name[16];
|
|
snprintf(thread_name, sizeof(thread_name), "db-read-%d", reader_index);
|
|
pthread_setname_np(pthread_self(), thread_name);
|
|
|
|
sqlite3* worker_db = open_worker_connection();
|
|
if (!worker_db) {
|
|
DEBUG_ERROR("Reader worker failed to open SQLite connection");
|
|
return NULL;
|
|
}
|
|
|
|
db_set_thread_connection(worker_db);
|
|
|
|
while (g_pool.running) {
|
|
thread_pool_job_node_t* node = queue_pop(&g_pool.read_q);
|
|
if (!node) break;
|
|
|
|
execute_read_job(node);
|
|
}
|
|
|
|
db_clear_thread_connection();
|
|
sqlite3_close(worker_db);
|
|
return NULL;
|
|
}
|
|
|
|
static void* writer_worker_main(void* arg) {
|
|
(void)arg;
|
|
|
|
pthread_setname_np(pthread_self(), "db-write");
|
|
|
|
sqlite3* worker_db = open_worker_connection();
|
|
if (!worker_db) {
|
|
DEBUG_ERROR("Writer worker failed to open SQLite connection");
|
|
return NULL;
|
|
}
|
|
|
|
db_set_thread_connection(worker_db);
|
|
|
|
while (g_pool.running) {
|
|
thread_pool_job_node_t* node = queue_pop(&g_pool.write_q);
|
|
if (!node) break;
|
|
|
|
execute_write_job(node);
|
|
}
|
|
|
|
db_clear_thread_connection();
|
|
sqlite3_close(worker_db);
|
|
return NULL;
|
|
}
|
|
|
|
int thread_pool_init(const thread_pool_config_t* config) {
|
|
if (!config) return -1;
|
|
|
|
pthread_mutex_lock(&g_pool.state_mutex);
|
|
if (g_pool.running) {
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return 0;
|
|
}
|
|
|
|
g_pool.reader_count = (config->reader_threads > 0) ? config->reader_threads : 4;
|
|
g_pool.readers = calloc((size_t)g_pool.reader_count, sizeof(pthread_t));
|
|
if (!g_pool.readers) {
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return -1;
|
|
}
|
|
|
|
queue_init(&g_pool.read_q, config->max_queue_depth);
|
|
queue_init(&g_pool.write_q, config->max_queue_depth);
|
|
|
|
g_pool.next_job_id = 1;
|
|
g_pool.wake_loop_cb = config->wake_loop_cb;
|
|
g_pool.wake_loop_ctx = config->wake_loop_ctx;
|
|
|
|
const char* db_path = (config->db_path && config->db_path[0] != '\0') ? config->db_path : db_get_database_path();
|
|
if (!db_path || db_path[0] == '\0') {
|
|
free(g_pool.readers);
|
|
g_pool.readers = NULL;
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return -1;
|
|
}
|
|
strncpy(g_pool.db_path, db_path, sizeof(g_pool.db_path) - 1);
|
|
g_pool.db_path[sizeof(g_pool.db_path) - 1] = '\0';
|
|
|
|
g_pool.running = 1;
|
|
|
|
for (int i = 0; i < g_pool.reader_count; i++) {
|
|
if (pthread_create(&g_pool.readers[i], NULL, reader_worker_main, (void*)(intptr_t)(i + 1)) != 0) {
|
|
g_pool.running = 0;
|
|
pthread_cond_broadcast(&g_pool.read_q.cond);
|
|
pthread_cond_broadcast(&g_pool.write_q.cond);
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return -1;
|
|
}
|
|
}
|
|
|
|
if (pthread_create(&g_pool.writer, NULL, writer_worker_main, NULL) != 0) {
|
|
g_pool.running = 0;
|
|
pthread_cond_broadcast(&g_pool.read_q.cond);
|
|
pthread_cond_broadcast(&g_pool.write_q.cond);
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return -1;
|
|
}
|
|
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
DEBUG_LOG("Thread pool initialized: %d readers + 1 writer", g_pool.reader_count);
|
|
return 0;
|
|
}
|
|
|
|
void thread_pool_shutdown(void) {
|
|
pthread_mutex_lock(&g_pool.state_mutex);
|
|
if (!g_pool.running) {
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
return;
|
|
}
|
|
|
|
g_pool.running = 0;
|
|
pthread_cond_broadcast(&g_pool.read_q.cond);
|
|
pthread_cond_broadcast(&g_pool.write_q.cond);
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
|
|
for (int i = 0; i < g_pool.reader_count; i++) {
|
|
pthread_join(g_pool.readers[i], NULL);
|
|
}
|
|
pthread_join(g_pool.writer, NULL);
|
|
|
|
free(g_pool.readers);
|
|
g_pool.readers = NULL;
|
|
g_pool.reader_count = 0;
|
|
|
|
queue_destroy(&g_pool.read_q);
|
|
queue_destroy(&g_pool.write_q);
|
|
|
|
DEBUG_LOG("Thread pool shutdown complete");
|
|
}
|
|
|
|
int thread_pool_is_running(void) {
|
|
return g_pool.running;
|
|
}
|
|
|
|
static thread_pool_status_t submit_to_queue(thread_pool_queue_t* q, const thread_pool_job_t* job, uint64_t* out_job_id) {
|
|
if (!g_pool.running) return THREAD_POOL_STATUS_SHUTTING_DOWN;
|
|
if (!job) return THREAD_POOL_STATUS_INTERNAL_ERROR;
|
|
|
|
thread_pool_job_node_t* node = calloc(1, sizeof(*node));
|
|
if (!node) return THREAD_POOL_STATUS_INTERNAL_ERROR;
|
|
|
|
node->job = *job;
|
|
|
|
pthread_mutex_lock(&g_pool.state_mutex);
|
|
node->job_id = g_pool.next_job_id++;
|
|
pthread_mutex_unlock(&g_pool.state_mutex);
|
|
|
|
if (out_job_id) *out_job_id = node->job_id;
|
|
|
|
if (!queue_push(q, node)) {
|
|
if (job->payload_free && job->payload) {
|
|
job->payload_free(job->payload);
|
|
}
|
|
free(node);
|
|
return THREAD_POOL_STATUS_QUEUE_FULL;
|
|
}
|
|
|
|
return THREAD_POOL_STATUS_OK;
|
|
}
|
|
|
|
thread_pool_status_t thread_pool_submit_read(const thread_pool_job_t* job, uint64_t* out_job_id) {
|
|
return submit_to_queue(&g_pool.read_q, job, out_job_id);
|
|
}
|
|
|
|
thread_pool_status_t thread_pool_submit_write(const thread_pool_job_t* job, uint64_t* out_job_id) {
|
|
return submit_to_queue(&g_pool.write_q, job, out_job_id);
|
|
}
|
|
|
|
int thread_pool_execute_count_sync(const char* sql, const char** bind_params, int bind_param_count, int* out_count) {
|
|
if (!sql || !out_count) {
|
|
return -1;
|
|
}
|
|
|
|
if (!thread_pool_is_running()) {
|
|
return db_count_with_sql(sql, bind_params, bind_param_count, out_count);
|
|
}
|
|
|
|
thread_pool_count_payload_t* payload = calloc(1, sizeof(*payload));
|
|
if (!payload) return -1;
|
|
|
|
payload->sql = strdup(sql);
|
|
payload->bind_param_count = bind_param_count;
|
|
if (!payload->sql) {
|
|
free_count_payload(payload);
|
|
return -1;
|
|
}
|
|
|
|
if (bind_param_count > 0) {
|
|
payload->bind_params = calloc((size_t)bind_param_count, sizeof(char*));
|
|
if (!payload->bind_params) {
|
|
free_count_payload(payload);
|
|
return -1;
|
|
}
|
|
for (int i = 0; i < bind_param_count; i++) {
|
|
const char* v = (bind_params && bind_params[i]) ? bind_params[i] : "";
|
|
payload->bind_params[i] = strdup(v);
|
|
if (!payload->bind_params[i]) {
|
|
free_count_payload(payload);
|
|
return -1;
|
|
}
|
|
}
|
|
}
|
|
|
|
thread_pool_wait_ctx_t wait_ctx;
|
|
wait_ctx_init(&wait_ctx);
|
|
|
|
thread_pool_job_t job;
|
|
memset(&job, 0, sizeof(job));
|
|
job.type = THREAD_POOL_JOB_COUNT_QUERY;
|
|
job.payload = payload;
|
|
job.payload_free = free_count_payload;
|
|
job.result_cb = wait_ctx_result_cb;
|
|
job.result_cb_ctx = &wait_ctx;
|
|
|
|
thread_pool_status_t submit_rc = thread_pool_submit_read(&job, NULL);
|
|
if (submit_rc != THREAD_POOL_STATUS_OK) {
|
|
free_count_payload(payload);
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return -1;
|
|
}
|
|
|
|
wait_ctx_wait(&wait_ctx);
|
|
|
|
int rc = -1;
|
|
if (wait_ctx.status == THREAD_POOL_STATUS_OK && wait_ctx.result_data) {
|
|
thread_pool_count_result_t* result = (thread_pool_count_result_t*)wait_ctx.result_data;
|
|
*out_count = result->count;
|
|
free(result);
|
|
rc = 0;
|
|
}
|
|
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return rc;
|
|
}
|
|
|
|
int thread_pool_execute_req_sync(const char* sql, const char** bind_params, int bind_param_count,
|
|
int row_limit, thread_pool_req_result_t** out_result) {
|
|
if (!sql || !out_result) {
|
|
return -1;
|
|
}
|
|
|
|
if (!thread_pool_is_running()) {
|
|
db_stmt_t* stmt = NULL;
|
|
if (db_prepare(sql, &stmt) != DB_OK) {
|
|
return -1;
|
|
}
|
|
|
|
for (int i = 0; i < bind_param_count; i++) {
|
|
const char* v = (bind_params && bind_params[i]) ? bind_params[i] : "";
|
|
(void)db_bind_text_param(stmt, i + 1, v);
|
|
}
|
|
|
|
int cap = (row_limit > 0) ? row_limit : 500;
|
|
if (cap < 1) cap = 1;
|
|
|
|
thread_pool_req_result_t* result = calloc(1, sizeof(*result));
|
|
if (!result) {
|
|
db_finalize_stmt(stmt);
|
|
return -1;
|
|
}
|
|
|
|
result->event_json_rows = calloc((size_t)cap, sizeof(char*));
|
|
if (!result->event_json_rows) {
|
|
free(result);
|
|
db_finalize_stmt(stmt);
|
|
return -1;
|
|
}
|
|
|
|
int step_rc = DB_OK;
|
|
while ((step_rc = db_step_stmt(stmt)) == DB_ROW) {
|
|
if (result->row_count >= cap) break;
|
|
const char* row = db_column_text_value(stmt, 0);
|
|
if (!row) continue;
|
|
|
|
result->event_json_rows[result->row_count] = strdup(row);
|
|
if (!result->event_json_rows[result->row_count]) {
|
|
break;
|
|
}
|
|
result->row_count++;
|
|
}
|
|
|
|
db_finalize_stmt(stmt);
|
|
*out_result = result;
|
|
return 0;
|
|
}
|
|
|
|
thread_pool_req_payload_t* payload = calloc(1, sizeof(*payload));
|
|
if (!payload) return -1;
|
|
|
|
payload->sql = strdup(sql);
|
|
payload->bind_param_count = bind_param_count;
|
|
payload->row_limit = row_limit;
|
|
if (!payload->sql) {
|
|
free_req_payload(payload);
|
|
return -1;
|
|
}
|
|
|
|
if (bind_param_count > 0) {
|
|
payload->bind_params = calloc((size_t)bind_param_count, sizeof(char*));
|
|
if (!payload->bind_params) {
|
|
free_req_payload(payload);
|
|
return -1;
|
|
}
|
|
for (int i = 0; i < bind_param_count; i++) {
|
|
const char* v = (bind_params && bind_params[i]) ? bind_params[i] : "";
|
|
payload->bind_params[i] = strdup(v);
|
|
if (!payload->bind_params[i]) {
|
|
free_req_payload(payload);
|
|
return -1;
|
|
}
|
|
}
|
|
}
|
|
|
|
thread_pool_wait_ctx_t wait_ctx;
|
|
wait_ctx_init(&wait_ctx);
|
|
|
|
thread_pool_job_t job;
|
|
memset(&job, 0, sizeof(job));
|
|
job.type = THREAD_POOL_JOB_REQ_QUERY;
|
|
job.payload = payload;
|
|
job.payload_free = free_req_payload;
|
|
job.result_cb = wait_ctx_result_cb;
|
|
job.result_cb_ctx = &wait_ctx;
|
|
|
|
thread_pool_status_t submit_rc = thread_pool_submit_read(&job, NULL);
|
|
if (submit_rc != THREAD_POOL_STATUS_OK) {
|
|
free_req_payload(payload);
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return -1;
|
|
}
|
|
|
|
wait_ctx_wait(&wait_ctx);
|
|
|
|
int rc = -1;
|
|
if (wait_ctx.status == THREAD_POOL_STATUS_OK && wait_ctx.result_data) {
|
|
*out_result = (thread_pool_req_result_t*)wait_ctx.result_data;
|
|
rc = 0;
|
|
}
|
|
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return rc;
|
|
}
|
|
|
|
void thread_pool_free_req_result(thread_pool_req_result_t* result) {
|
|
if (!result) return;
|
|
if (result->event_json_rows) {
|
|
for (int i = 0; i < result->row_count; i++) {
|
|
free(result->event_json_rows[i]);
|
|
}
|
|
free(result->event_json_rows);
|
|
}
|
|
free(result);
|
|
}
|
|
|
|
int thread_pool_execute_store_event_sync(const thread_pool_store_event_payload_t* payload,
|
|
thread_pool_store_event_result_t* out_result) {
|
|
if (!payload || !out_result || !payload->id || !payload->pubkey || !payload->content || !payload->sig) {
|
|
return -1;
|
|
}
|
|
|
|
if (!thread_pool_is_running()) {
|
|
int step_rc = DB_ERROR;
|
|
int extended_errcode = 0;
|
|
if (db_insert_event_with_json(payload->id,
|
|
payload->pubkey,
|
|
payload->created_at,
|
|
payload->kind,
|
|
payload->event_type ? payload->event_type : "regular",
|
|
payload->content,
|
|
payload->sig,
|
|
payload->tags_json ? payload->tags_json : "[]",
|
|
payload->event_json ? payload->event_json : "{}",
|
|
&step_rc,
|
|
&extended_errcode) != 0) {
|
|
return -1;
|
|
}
|
|
out_result->step_rc = step_rc;
|
|
out_result->extended_errcode = extended_errcode;
|
|
return 0;
|
|
}
|
|
|
|
thread_pool_store_event_payload_t* payload_copy = calloc(1, sizeof(*payload_copy));
|
|
if (!payload_copy) return -1;
|
|
|
|
payload_copy->id = payload->id ? strdup(payload->id) : NULL;
|
|
payload_copy->pubkey = payload->pubkey ? strdup(payload->pubkey) : NULL;
|
|
payload_copy->created_at = payload->created_at;
|
|
payload_copy->kind = payload->kind;
|
|
payload_copy->event_type = payload->event_type ? strdup(payload->event_type) : strdup("regular");
|
|
payload_copy->content = payload->content ? strdup(payload->content) : NULL;
|
|
payload_copy->sig = payload->sig ? strdup(payload->sig) : NULL;
|
|
payload_copy->tags_json = payload->tags_json ? strdup(payload->tags_json) : strdup("[]");
|
|
payload_copy->event_json = payload->event_json ? strdup(payload->event_json) : strdup("{}");
|
|
|
|
if (!payload_copy->id || !payload_copy->pubkey || !payload_copy->event_type ||
|
|
!payload_copy->content || !payload_copy->sig || !payload_copy->tags_json || !payload_copy->event_json) {
|
|
free_store_event_payload(payload_copy);
|
|
return -1;
|
|
}
|
|
|
|
thread_pool_wait_ctx_t wait_ctx;
|
|
wait_ctx_init(&wait_ctx);
|
|
|
|
thread_pool_job_t job;
|
|
memset(&job, 0, sizeof(job));
|
|
job.type = THREAD_POOL_JOB_STORE_EVENT;
|
|
job.payload = payload_copy;
|
|
job.payload_free = free_store_event_payload;
|
|
job.result_cb = wait_ctx_result_cb;
|
|
job.result_cb_ctx = &wait_ctx;
|
|
|
|
thread_pool_status_t submit_rc = thread_pool_submit_write(&job, NULL);
|
|
if (submit_rc != THREAD_POOL_STATUS_OK) {
|
|
free_store_event_payload(payload_copy);
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return -1;
|
|
}
|
|
|
|
wait_ctx_wait(&wait_ctx);
|
|
|
|
int rc = -1;
|
|
if (wait_ctx.status == THREAD_POOL_STATUS_OK && wait_ctx.result_data) {
|
|
thread_pool_store_event_result_t* result = (thread_pool_store_event_result_t*)wait_ctx.result_data;
|
|
out_result->step_rc = result->step_rc;
|
|
out_result->extended_errcode = result->extended_errcode;
|
|
free(result);
|
|
rc = 0;
|
|
}
|
|
|
|
wait_ctx_destroy(&wait_ctx);
|
|
return rc;
|
|
}
|
|
|
|
__attribute__((constructor))
|
|
static void thread_pool_state_init_once(void) {
|
|
pthread_mutex_init(&g_pool.state_mutex, NULL);
|
|
}
|