Files
c-relay-pg/src/db_ops_postgres.c
T

2876 lines
101 KiB
C

#define _GNU_SOURCE
#include "db_ops_postgres.h"
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
struct postgres_db_stmt {
char* sql_original;
char* sql_converted;
int placeholder_count;
char** param_values; // 1-based via [index-1]
void* result; // PGresult* when libpq is enabled
int current_row;
int row_count;
int executed;
};
static __thread char g_postgres_error[512] = "postgres backend not implemented";
static char g_postgres_connection_string[512] = {0};
static void postgres_set_error_text(const char* msg) {
if (!msg) msg = "postgres backend error";
strncpy(g_postgres_error, msg, sizeof(g_postgres_error) - 1);
g_postgres_error[sizeof(g_postgres_error) - 1] = '\0';
}
static char* postgres_strdup(const char* s) {
if (!s) return NULL;
size_t n = strlen(s) + 1;
char* p = (char*)malloc(n);
if (!p) return NULL;
memcpy(p, s, n);
return p;
}
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
#include "pg_schema.h"
#include <libpq-fe.h>
#include <sys/select.h>
#include <sys/socket.h>
#include <unistd.h>
#include <errno.h>
static PGconn* g_pg_conn = NULL;
static __thread PGconn* g_thread_pg_conn = NULL;
static void postgres_set_error_from_conn(PGconn* conn, const char* fallback);
static PGconn* postgres_db_active_connection(void) {
PGconn* conn = g_thread_pg_conn ? g_thread_pg_conn : g_pg_conn;
if (!conn) return NULL;
if (PQstatus(conn) == CONNECTION_OK) {
return conn;
}
// Lightweight recovery path for transient disconnects.
PQreset(conn);
if (PQstatus(conn) == CONNECTION_OK) {
postgres_set_error_text("ok");
return conn;
}
postgres_set_error_from_conn(conn, "postgres connection lost");
return NULL;
}
static void postgres_set_error_from_conn(PGconn* conn, const char* fallback) {
const char* err = conn ? PQerrorMessage(conn) : NULL;
if (err && err[0] != '\0') {
postgres_set_error_text(err);
} else {
postgres_set_error_text(fallback ? fallback : "postgres backend error");
}
}
static void postgres_stmt_clear_result(postgres_db_stmt_t* stmt) {
if (!stmt || !stmt->result) return;
PQclear((PGresult*)stmt->result);
stmt->result = NULL;
stmt->row_count = 0;
stmt->current_row = 0;
stmt->executed = 0;
}
static char* postgres_convert_qmark_to_dollar(const char* sql, int* out_count) {
if (!sql || !out_count) return NULL;
size_t len = strlen(sql);
size_t cap = len + 32;
char* out = (char*)malloc(cap);
if (!out) return NULL;
int idx = 1;
size_t j = 0;
for (size_t i = 0; i < len; i++) {
if (sql[i] == '?') {
char buf[32];
int written = snprintf(buf, sizeof(buf), "$%d", idx++);
if (written <= 0) {
free(out);
return NULL;
}
if (j + (size_t)written + 1 > cap) {
cap = (cap * 2) + 64;
char* grown = (char*)realloc(out, cap);
if (!grown) {
free(out);
return NULL;
}
out = grown;
}
memcpy(out + j, buf, (size_t)written);
j += (size_t)written;
} else {
if (j + 2 > cap) {
cap *= 2;
char* grown = (char*)realloc(out, cap);
if (!grown) {
free(out);
return NULL;
}
out = grown;
}
out[j++] = sql[i];
}
}
out[j] = '\0';
*out_count = idx - 1;
return out;
}
int postgres_db_apply_schema(void) {
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) {
postgres_set_error_text("postgres_db_apply_schema: no active connection");
return DB_ERROR;
}
PGresult* res = PQexec(conn, EMBEDDED_PG_SCHEMA_SQL);
if (!res) {
postgres_set_error_from_conn(conn, "postgres_db_apply_schema failed");
return DB_ERROR;
}
ExecStatusType st = PQresultStatus(res);
if (st != PGRES_COMMAND_OK && st != PGRES_TUPLES_OK) {
const char* msg = PQresultErrorMessage(res);
postgres_set_error_text((msg && msg[0] != '\0') ? msg : "postgres_db_apply_schema failed");
PQclear(res);
return DB_ERROR;
}
PQclear(res);
return DB_OK;
}
int postgres_db_init(const char* connection_string) {
if (connection_string && connection_string[0] != '\0') {
strncpy(g_postgres_connection_string, connection_string, sizeof(g_postgres_connection_string) - 1);
g_postgres_connection_string[sizeof(g_postgres_connection_string) - 1] = '\0';
}
if (g_postgres_connection_string[0] == '\0') {
postgres_set_error_text("postgres_db_init: empty connection string");
return DB_MISUSE;
}
if (g_pg_conn) {
PQfinish(g_pg_conn);
g_pg_conn = NULL;
}
g_pg_conn = PQconnectdb(g_postgres_connection_string);
if (!g_pg_conn || PQstatus(g_pg_conn) != CONNECTION_OK) {
postgres_set_error_from_conn(g_pg_conn, "postgres_db_init: connection failed");
if (g_pg_conn) {
PQfinish(g_pg_conn);
g_pg_conn = NULL;
}
return DB_ERROR;
}
if (postgres_db_apply_schema() != DB_OK) {
PQfinish(g_pg_conn);
g_pg_conn = NULL;
return DB_ERROR;
}
postgres_set_error_text("ok");
return DB_OK;
}
void postgres_db_close(void) {
if (g_pg_conn) {
PQfinish(g_pg_conn);
g_pg_conn = NULL;
}
}
int postgres_db_is_available(void) {
PGconn* conn = postgres_db_active_connection();
return conn && PQstatus(conn) == CONNECTION_OK;
}
const char* postgres_db_last_error(void) {
PGconn* conn = postgres_db_active_connection();
if (conn && PQstatus(conn) != CONNECTION_OK) {
postgres_set_error_from_conn(conn, g_postgres_error);
}
return g_postgres_error;
}
const char* postgres_db_get_database_path(void) { return g_postgres_connection_string; }
int postgres_db_set_thread_connection(void* connection) {
g_thread_pg_conn = (PGconn*)connection;
return DB_OK;
}
void postgres_db_clear_thread_connection(void) {
g_thread_pg_conn = NULL;
}
int postgres_db_open_worker_connection(const char* db_path, void** out_connection) {
if (!out_connection) return DB_MISUSE;
*out_connection = NULL;
const char* effective = (db_path && db_path[0] != '\0') ? db_path : g_postgres_connection_string;
if (!effective || effective[0] == '\0') {
postgres_set_error_text("postgres_db_open_worker_connection: empty connection string");
return DB_MISUSE;
}
PGconn* conn = PQconnectdb(effective);
if (!conn || PQstatus(conn) != CONNECTION_OK) {
postgres_set_error_from_conn(conn, "postgres_db_open_worker_connection: connection failed");
if (conn) PQfinish(conn);
return DB_ERROR;
}
postgres_db_set_thread_connection(conn);
*out_connection = (void*)conn;
return DB_OK;
}
// Issue LISTEN <channel> on a worker connection. The connection must remain
// valid for the lifetime of the listen. Returns DB_OK on success, DB_ERROR on
// failure. Channel is treated as an identifier (validated alnum/underscore).
int postgres_db_worker_listen(void* connection, const char* channel) {
PGconn* conn = (PGconn*)connection;
if (!conn || PQstatus(conn) != CONNECTION_OK || !channel || !channel[0]) {
return DB_ERROR;
}
// Validate channel name to avoid SQL injection via LISTEN.
for (const char* p = channel; *p; p++) {
if (!( ((*p >= 'a' && *p <= 'z') || (*p >= 'A' && *p <= 'Z') ||
(*p >= '0' && *p <= '9') || *p == '_') )) {
postgres_set_error_text("postgres_db_worker_listen: invalid channel name");
return DB_ERROR;
}
}
char sql[128];
snprintf(sql, sizeof(sql), "LISTEN %s", channel);
PGresult* res = PQexec(conn, sql);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
postgres_set_error_from_conn(conn, "postgres_db_worker_listen: LISTEN failed");
if (res) PQclear(res);
return DB_ERROR;
}
PQclear(res);
return DB_OK;
}
// Wait up to timeout_ms for a notification on the worker connection.
// Returns 1 if a notification was consumed, 0 on timeout, -1 on error.
// This does NOT register a LISTEN; the caller must have issued LISTEN first.
int postgres_db_worker_poll_notify(void* connection, int timeout_ms) {
PGconn* conn = (PGconn*)connection;
if (!conn || PQstatus(conn) != CONNECTION_OK) {
return -1;
}
// First, drain any already-queued notifications without blocking.
PQconsumeInput(conn);
PGnotify* notify = PQnotifies(conn);
if (notify) {
PQfreemem(notify);
return 1;
}
int sock = PQsocket(conn);
if (sock < 0) {
return -1;
}
fd_set fds;
FD_ZERO(&fds);
FD_SET(sock, &fds);
struct timeval tv;
tv.tv_sec = timeout_ms / 1000;
tv.tv_usec = (timeout_ms % 1000) * 1000;
int rc = select(sock + 1, &fds, NULL, NULL, &tv);
if (rc < 0) {
if (errno == EINTR) return 0; // signal interrupted: treat as timeout-ish
return -1;
}
if (rc == 0) {
return 0; // timeout
}
// Socket readable: consume input and check for notifications.
if (PQconsumeInput(conn) == 0) {
return -1;
}
notify = PQnotifies(conn);
if (notify) {
PQfreemem(notify);
return 1;
}
// Readable but no notification (could be a keepalive or other traffic).
return 0;
}
// Like postgres_db_worker_poll_notify but also returns channel + payload.
// Caller must free *out_channel and *out_payload (both NULL on timeout/error).
int postgres_db_worker_poll_notify_with_payload(void* connection, int timeout_ms,
char** out_channel, char** out_payload) {
PGconn* conn = (PGconn*)connection;
if (out_channel) *out_channel = NULL;
if (out_payload) *out_payload = NULL;
if (!conn || PQstatus(conn) != CONNECTION_OK) {
return -1;
}
// First, drain any already-queued notifications without blocking.
PQconsumeInput(conn);
PGnotify* notify = PQnotifies(conn);
if (notify) {
if (out_channel && notify->relname) *out_channel = strdup(notify->relname);
if (out_payload && notify->extra) *out_payload = strdup(notify->extra);
PQfreemem(notify);
return 1;
}
int sock = PQsocket(conn);
if (sock < 0) {
return -1;
}
fd_set fds;
FD_ZERO(&fds);
FD_SET(sock, &fds);
struct timeval tv;
tv.tv_sec = timeout_ms / 1000;
tv.tv_usec = (timeout_ms % 1000) * 1000;
int rc = select(sock + 1, &fds, NULL, NULL, &tv);
if (rc < 0) {
if (errno == EINTR) return 0;
return -1;
}
if (rc == 0) {
return 0;
}
if (PQconsumeInput(conn) == 0) {
return -1;
}
notify = PQnotifies(conn);
if (notify) {
if (out_channel && notify->relname) *out_channel = strdup(notify->relname);
if (out_payload && notify->extra) *out_payload = strdup(notify->extra);
PQfreemem(notify);
return 1;
}
return 0;
}
void postgres_db_close_worker_connection(void* connection) {
PGconn* conn = (PGconn*)connection;
if (!conn) return;
if (g_thread_pg_conn == conn) {
postgres_db_clear_thread_connection();
}
PQfinish(conn);
}
int postgres_db_prepare(const char* sql, postgres_db_stmt_t** out_stmt) {
if (!sql || !out_stmt) return DB_MISUSE;
*out_stmt = NULL;
postgres_db_stmt_t* stmt = (postgres_db_stmt_t*)calloc(1, sizeof(postgres_db_stmt_t));
if (!stmt) return DB_ERROR;
stmt->sql_original = postgres_strdup(sql);
stmt->sql_converted = postgres_convert_qmark_to_dollar(sql, &stmt->placeholder_count);
if (!stmt->sql_original || !stmt->sql_converted) {
free(stmt->sql_original);
free(stmt->sql_converted);
free(stmt);
return DB_ERROR;
}
if (stmt->placeholder_count > 0) {
stmt->param_values = (char**)calloc((size_t)stmt->placeholder_count, sizeof(char*));
if (!stmt->param_values) {
free(stmt->sql_original);
free(stmt->sql_converted);
free(stmt);
return DB_ERROR;
}
}
*out_stmt = stmt;
return DB_OK;
}
int postgres_db_bind_text_param(postgres_db_stmt_t* stmt, int index, const char* value) {
if (!stmt || index <= 0 || index > stmt->placeholder_count) return DB_MISUSE;
int pos = index - 1;
free(stmt->param_values[pos]);
stmt->param_values[pos] = postgres_strdup(value ? value : "");
if (!stmt->param_values[pos]) return DB_ERROR;
return DB_OK;
}
int postgres_db_bind_int_param(postgres_db_stmt_t* stmt, int index, int value) {
char buf[32];
snprintf(buf, sizeof(buf), "%d", value);
return postgres_db_bind_text_param(stmt, index, buf);
}
int postgres_db_bind_int64_param(postgres_db_stmt_t* stmt, int index, long long value) {
char buf[64];
snprintf(buf, sizeof(buf), "%lld", value);
return postgres_db_bind_text_param(stmt, index, buf);
}
int postgres_db_step_stmt(postgres_db_stmt_t* stmt) {
if (!stmt) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) {
postgres_set_error_text("postgres_db_step_stmt: no active connection");
return DB_ERROR;
}
if (!stmt->executed) {
postgres_stmt_clear_result(stmt);
const char* const* params = (const char* const*)stmt->param_values;
PGresult* res = PQexecParams(conn,
stmt->sql_converted,
stmt->placeholder_count,
NULL,
params,
NULL,
NULL,
0);
if (!res) {
postgres_set_error_from_conn(conn, "postgres_db_step_stmt: query failed");
return DB_ERROR;
}
ExecStatusType st = PQresultStatus(res);
if (st == PGRES_TUPLES_OK) {
stmt->result = res;
stmt->row_count = PQntuples(res);
stmt->current_row = 0;
stmt->executed = 1;
return stmt->row_count > 0 ? DB_ROW : DB_DONE;
}
if (st == PGRES_COMMAND_OK) {
PQclear(res);
stmt->executed = 1;
stmt->row_count = 0;
stmt->current_row = 0;
return DB_DONE;
}
const char* msg = PQresultErrorMessage(res);
postgres_set_error_text((msg && msg[0] != '\0') ? msg : "postgres_db_step_stmt failed");
PQclear(res);
return DB_ERROR;
}
if (stmt->result && stmt->current_row + 1 < stmt->row_count) {
stmt->current_row++;
return DB_ROW;
}
return DB_DONE;
}
int postgres_db_reset_stmt(postgres_db_stmt_t* stmt) {
if (!stmt) return DB_MISUSE;
// Always reset execution state so command statements can be re-executed
// after each bind cycle (mirrors sqlite3_reset behavior expected by callers).
postgres_stmt_clear_result(stmt);
stmt->executed = 0;
stmt->current_row = 0;
stmt->row_count = 0;
return DB_OK;
}
const char* postgres_db_column_text_value(postgres_db_stmt_t* stmt, int col) {
if (!stmt || !stmt->result) return NULL;
if (stmt->current_row < 0 || stmt->current_row >= stmt->row_count) return NULL;
PGresult* res = (PGresult*)stmt->result;
if (col < 0 || col >= PQnfields(res)) return NULL;
if (PQgetisnull(res, stmt->current_row, col)) return NULL;
return PQgetvalue(res, stmt->current_row, col);
}
int postgres_db_column_int_value(postgres_db_stmt_t* stmt, int col) {
const char* v = postgres_db_column_text_value(stmt, col);
return v ? (int)strtol(v, NULL, 10) : 0;
}
long long postgres_db_column_int64_value(postgres_db_stmt_t* stmt, int col) {
const char* v = postgres_db_column_text_value(stmt, col);
return v ? strtoll(v, NULL, 10) : 0;
}
double postgres_db_column_double_value(postgres_db_stmt_t* stmt, int col) {
const char* v = postgres_db_column_text_value(stmt, col);
return v ? strtod(v, NULL) : 0.0;
}
void postgres_db_finalize_stmt(postgres_db_stmt_t* stmt) {
if (!stmt) return;
postgres_stmt_clear_result(stmt);
if (stmt->param_values) {
for (int i = 0; i < stmt->placeholder_count; i++) {
free(stmt->param_values[i]);
}
free(stmt->param_values);
}
free(stmt->sql_original);
free(stmt->sql_converted);
free(stmt);
}
#else
int postgres_db_apply_schema(void) {
return DB_ERROR;
}
int postgres_db_init(const char* connection_string) {
if (connection_string && connection_string[0] != '\0') {
strncpy(g_postgres_connection_string, connection_string, sizeof(g_postgres_connection_string) - 1);
g_postgres_connection_string[sizeof(g_postgres_connection_string) - 1] = '\0';
}
postgres_set_error_text("postgres backend built without libpq support");
return DB_ERROR;
}
void postgres_db_close(void) {}
int postgres_db_is_available(void) { return 0; }
const char* postgres_db_last_error(void) { return g_postgres_error; }
const char* postgres_db_get_database_path(void) { return g_postgres_connection_string; }
int postgres_db_set_thread_connection(void* connection) { (void)connection; return DB_ERROR; }
void postgres_db_clear_thread_connection(void) {}
int postgres_db_open_worker_connection(const char* db_path, void** out_connection) {
(void)db_path;
if (out_connection) *out_connection = NULL;
return DB_ERROR;
}
void postgres_db_close_worker_connection(void* connection) { (void)connection; }
int postgres_db_worker_listen(void* connection, const char* channel) {
(void)connection; (void)channel; return -1;
}
int postgres_db_worker_poll_notify(void* connection, int timeout_ms) {
(void)connection; (void)timeout_ms; return -1;
}
int postgres_db_prepare(const char* sql, postgres_db_stmt_t** out_stmt) {
(void)sql;
if (out_stmt) *out_stmt = NULL;
return DB_ERROR;
}
int postgres_db_bind_text_param(postgres_db_stmt_t* stmt, int index, const char* value) {
(void)stmt; (void)index; (void)value; return DB_ERROR;
}
int postgres_db_bind_int_param(postgres_db_stmt_t* stmt, int index, int value) {
(void)stmt; (void)index; (void)value; return DB_ERROR;
}
int postgres_db_bind_int64_param(postgres_db_stmt_t* stmt, int index, long long value) {
(void)stmt; (void)index; (void)value; return DB_ERROR;
}
int postgres_db_step_stmt(postgres_db_stmt_t* stmt) { (void)stmt; return DB_ERROR; }
int postgres_db_reset_stmt(postgres_db_stmt_t* stmt) { (void)stmt; return DB_ERROR; }
const char* postgres_db_column_text_value(postgres_db_stmt_t* stmt, int col) { (void)stmt; (void)col; return NULL; }
int postgres_db_column_int_value(postgres_db_stmt_t* stmt, int col) { (void)stmt; (void)col; return 0; }
long long postgres_db_column_int64_value(postgres_db_stmt_t* stmt, int col) { (void)stmt; (void)col; return 0; }
double postgres_db_column_double_value(postgres_db_stmt_t* stmt, int col) { (void)stmt; (void)col; return 0.0; }
void postgres_db_finalize_stmt(postgres_db_stmt_t* stmt) { (void)stmt; }
#endif
int postgres_db_log_subscription_created(const char* sub_id, const char* wsi_ptr,
const char* client_ip, const char* filter_json) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!sub_id || !wsi_ptr || !client_ip) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[4] = { sub_id, wsi_ptr, client_ip, filter_json ? filter_json : "[]" };
PGresult* res = PQexecParams(conn,
"INSERT INTO subscriptions (subscription_id, wsi_pointer, client_ip, event_type, filter_json) "
"VALUES ($1, $2, $3, 'created', $4)",
4, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)sub_id; (void)wsi_ptr; (void)client_ip; (void)filter_json;
return DB_ERROR;
#endif
}
int postgres_db_log_subscription_closed(const char* sub_id, const char* client_ip) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!sub_id) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params_insert[2] = { sub_id, client_ip ? client_ip : "unknown" };
PGresult* res = PQexecParams(conn,
"INSERT INTO subscriptions (subscription_id, wsi_pointer, client_ip, event_type) VALUES ($1, '', $2, 'closed')",
2, NULL, params_insert, NULL, NULL, 0);
if (res) PQclear(res);
const char* params_upd[1] = { sub_id };
res = PQexecParams(conn,
"UPDATE subscriptions SET ended_at = EXTRACT(EPOCH FROM NOW())::BIGINT "
"WHERE subscription_id = $1 AND event_type = 'created' AND ended_at IS NULL",
1, NULL, params_upd, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)sub_id; (void)client_ip;
return DB_ERROR;
#endif
}
int postgres_db_log_subscription_disconnected(const char* client_ip) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!client_ip) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params_upd[1] = { client_ip };
PGresult* res = PQexecParams(conn,
"UPDATE subscriptions SET ended_at = EXTRACT(EPOCH FROM NOW())::BIGINT "
"WHERE client_ip = $1 AND event_type = 'created' AND ended_at IS NULL",
1, NULL, params_upd, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
int changes = 0;
const char* tuples = PQcmdTuples(res);
if (tuples && tuples[0] != '\0') changes = atoi(tuples);
PQclear(res);
if (changes > 0) {
const char* params_ins[1] = { client_ip };
res = PQexecParams(conn,
"INSERT INTO subscriptions (subscription_id, wsi_pointer, client_ip, event_type) "
"VALUES ('disconnect', '', $1, 'disconnected')",
1, NULL, params_ins, NULL, NULL, 0);
if (res) PQclear(res);
}
return changes;
#else
(void)client_ip;
return DB_ERROR;
#endif
}
int postgres_db_update_subscription_events_sent(const char* sub_id, int events_sent) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!sub_id) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
char sent_buf[32];
snprintf(sent_buf, sizeof(sent_buf), "%d", events_sent);
const char* params[2] = { sent_buf, sub_id };
PGresult* res = PQexecParams(conn,
"UPDATE subscriptions SET events_sent = $1::INT "
"WHERE subscription_id = $2 AND event_type = 'created'",
2, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)sub_id; (void)events_sent;
return DB_ERROR;
#endif
}
int postgres_db_cleanup_orphaned_subscriptions(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn,
"UPDATE subscriptions SET ended_at = EXTRACT(EPOCH FROM NOW())::BIGINT "
"WHERE event_type = 'created' AND ended_at IS NULL");
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
int changes = 0;
const char* tuples = PQcmdTuples(res);
if (tuples && tuples[0] != '\0') changes = atoi(tuples);
PQclear(res);
return changes;
#else
return DB_ERROR;
#endif
}
int postgres_db_get_event_pubkey(const char* event_id, char* pubkey_out, size_t pubkey_out_size) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event_id || !pubkey_out || pubkey_out_size == 0) return DB_MISUSE;
pubkey_out[0] = '\0';
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[1] = { event_id };
PGresult* res = PQexecParams(conn,
"SELECT pubkey FROM events WHERE id = $1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_event_pubkey failed");
}
return DB_ERROR;
}
int found = 0;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
snprintf(pubkey_out, pubkey_out_size, "%s", PQgetvalue(res, 0, 0));
found = 1;
}
PQclear(res);
return found;
#else
(void)event_id;
if (pubkey_out && pubkey_out_size > 0) pubkey_out[0] = '\0';
return DB_ERROR;
#endif
}
int postgres_db_delete_event_by_id(const char* event_id, const char* requester_pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event_id || !requester_pubkey) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[2] = { event_id, requester_pubkey };
PGresult* res = PQexecParams(conn,
"DELETE FROM events WHERE id = $1 AND pubkey = $2",
2, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_delete_event_by_id failed");
}
return DB_ERROR;
}
const char* tuples = PQcmdTuples(res);
int deleted = (tuples && tuples[0] != '\0') ? atoi(tuples) : 0;
PQclear(res);
return deleted;
#else
(void)event_id; (void)requester_pubkey;
return DB_ERROR;
#endif
}
int postgres_db_delete_events_by_address(const char* pubkey, int kind,
const char* d_tag, long before_timestamp) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
char kind_buf[32];
char before_buf[64];
snprintf(kind_buf, sizeof(kind_buf), "%d", kind);
snprintf(before_buf, sizeof(before_buf), "%ld", before_timestamp);
PGresult* res = NULL;
if (d_tag && d_tag[0] != '\0') {
const char* params[4] = { kind_buf, pubkey, before_buf, d_tag };
res = PQexecParams(conn,
"DELETE FROM events "
"WHERE kind = $1::INT AND pubkey = $2 AND created_at < $3::BIGINT "
"AND d_tag_value = $4",
4, NULL, params, NULL, NULL, 0);
} else {
const char* params[3] = { kind_buf, pubkey, before_buf };
res = PQexecParams(conn,
"DELETE FROM events "
"WHERE kind = $1::INT AND pubkey = $2 AND created_at < $3::BIGINT",
3, NULL, params, NULL, NULL, 0);
}
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_delete_events_by_address failed");
}
return DB_ERROR;
}
const char* tuples = PQcmdTuples(res);
int deleted = (tuples && tuples[0] != '\0') ? atoi(tuples) : 0;
PQclear(res);
return deleted;
#else
(void)pubkey; (void)kind; (void)d_tag; (void)before_timestamp;
return DB_ERROR;
#endif
}
int postgres_db_is_pubkey_blacklisted(const char* pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey) return 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return 0;
const char* params[1] = { pubkey };
PGresult* res = PQexecParams(conn,
"SELECT 1 FROM auth_rules WHERE rule_type='blacklist' AND pattern_type='pubkey' AND pattern_value=$1 AND active=1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return 0;
}
int exists = PQntuples(res) > 0 ? 1 : 0;
PQclear(res);
return exists;
#else
(void)pubkey;
return 0;
#endif
}
int postgres_db_is_hash_blacklisted(const char* resource_hash) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!resource_hash) return 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return 0;
const char* params[1] = { resource_hash };
PGresult* res = PQexecParams(conn,
"SELECT 1 FROM auth_rules WHERE rule_type='blacklist' AND pattern_type='hash' AND pattern_value=$1 AND active=1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return 0;
}
int exists = PQntuples(res) > 0 ? 1 : 0;
PQclear(res);
return exists;
#else
(void)resource_hash;
return 0;
#endif
}
int postgres_db_is_pubkey_whitelisted(const char* pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey) return 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return 0;
const char* params[1] = { pubkey };
PGresult* res = PQexecParams(conn,
"SELECT 1 FROM auth_rules WHERE rule_type IN ('whitelist','wot_whitelist') AND pattern_type='pubkey' AND pattern_value=$1 AND active=1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return 0;
}
int exists = PQntuples(res) > 0 ? 1 : 0;
PQclear(res);
return exists;
#else
(void)pubkey;
return 0;
#endif
}
int postgres_db_count_active_whitelist_rules(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return 0;
PGresult* res = PQexec(conn,
"SELECT COUNT(*) FROM auth_rules WHERE rule_type IN ('whitelist','wot_whitelist') AND active=1");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return 0;
}
int count = (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) ? (int)strtol(PQgetvalue(res, 0, 0), NULL, 10) : 0;
PQclear(res);
return count;
#else
return 0;
#endif
}
int postgres_db_count_with_sql(const char* sql, const char** bind_params, int bind_param_count, int* out_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!sql || !out_count || bind_param_count < 0) return DB_MISUSE;
*out_count = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
int placeholder_count = 0;
char* converted_sql = postgres_convert_qmark_to_dollar(sql, &placeholder_count);
if (!converted_sql) {
postgres_set_error_text("postgres_db_count_with_sql: failed to convert sql placeholders");
return DB_ERROR;
}
if (placeholder_count != bind_param_count) {
free(converted_sql);
postgres_set_error_text("postgres_db_count_with_sql: bind parameter count mismatch");
return DB_MISUSE;
}
PGresult* res = PQexecParams(conn,
converted_sql,
bind_param_count,
NULL,
bind_param_count > 0 ? bind_params : NULL,
NULL,
NULL,
0);
free(converted_sql);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_count_with_sql failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
*out_count = (int)strtol(PQgetvalue(res, 0, 0), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
(void)sql; (void)bind_params; (void)bind_param_count;
if (out_count) *out_count = 0;
return DB_ERROR;
#endif
}
char* postgres_db_execute_readonly_query_json(const char* query, const char* request_id,
char* error_message, size_t error_size,
int max_rows, int timeout_ms) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!query || !request_id || !error_message || error_size == 0) return NULL;
error_message[0] = '\0';
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) {
snprintf(error_message, error_size, "No active PostgreSQL connection");
return NULL;
}
const char* p = query;
while (*p == ' ' || *p == '\t' || *p == '\r' || *p == '\n') p++;
if (strncasecmp(p, "select", 6) != 0 && strncasecmp(p, "with", 4) != 0) {
snprintf(error_message, error_size, "Only read-only SELECT/WITH queries are allowed");
return NULL;
}
if (timeout_ms > 0) {
char timeout_sql[64];
snprintf(timeout_sql, sizeof(timeout_sql), "SET statement_timeout = %d", timeout_ms);
PGresult* timeout_res = PQexec(conn, timeout_sql);
if (timeout_res) PQclear(timeout_res);
}
PGresult* res = PQexec(conn, query);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
const char* err = res ? PQresultErrorMessage(res) : PQerrorMessage(conn);
snprintf(error_message, error_size, "SQL execution failed: %s", (err && err[0] != '\0') ? err : "unknown");
if (res) PQclear(res);
if (timeout_ms > 0) {
PGresult* reset_res = PQexec(conn, "SET statement_timeout = 0");
if (reset_res) PQclear(reset_res);
}
return NULL;
}
cJSON* response = cJSON_CreateObject();
if (!response) {
PQclear(res);
snprintf(error_message, error_size, "Failed to allocate query response object");
return NULL;
}
cJSON_AddStringToObject(response, "query_type", "sql_query");
cJSON_AddStringToObject(response, "request_id", request_id);
cJSON_AddNumberToObject(response, "timestamp", (double)time(NULL));
cJSON_AddStringToObject(response, "query", query);
int col_count = PQnfields(res);
cJSON* columns = cJSON_CreateArray();
if (!columns) {
cJSON_Delete(response);
PQclear(res);
snprintf(error_message, error_size, "Failed to allocate columns array");
return NULL;
}
for (int i = 0; i < col_count; i++) {
const char* col_name = PQfname(res, i);
cJSON_AddItemToArray(columns, cJSON_CreateString(col_name ? col_name : ""));
}
cJSON_AddItemToObject(response, "columns", columns);
cJSON* rows = cJSON_CreateArray();
if (!rows) {
cJSON_Delete(response);
PQclear(res);
snprintf(error_message, error_size, "Failed to allocate rows array");
return NULL;
}
int row_limit = (max_rows > 0) ? max_rows : 1000;
int available_rows = PQntuples(res);
int row_count = (available_rows > row_limit) ? row_limit : available_rows;
for (int r = 0; r < row_count; r++) {
cJSON* row = cJSON_CreateArray();
if (!row) {
cJSON_Delete(rows);
cJSON_Delete(response);
PQclear(res);
snprintf(error_message, error_size, "Failed to allocate row array");
return NULL;
}
for (int c = 0; c < col_count; c++) {
if (PQgetisnull(res, r, c)) {
cJSON_AddItemToArray(row, cJSON_CreateNull());
} else {
cJSON_AddItemToArray(row, cJSON_CreateString(PQgetvalue(res, r, c)));
}
}
cJSON_AddItemToArray(rows, row);
}
if (available_rows > row_limit) {
cJSON_AddStringToObject(response, "warning", "Result truncated to maximum row limit");
}
cJSON_AddNumberToObject(response, "row_count", row_count);
cJSON_AddNumberToObject(response, "execution_time_ms", 0);
cJSON_AddItemToObject(response, "rows", rows);
char* out = cJSON_Print(response);
cJSON_Delete(response);
PQclear(res);
if (timeout_ms > 0) {
PGresult* reset_res = PQexec(conn, "SET statement_timeout = 0");
if (reset_res) PQclear(reset_res);
}
if (!out) {
snprintf(error_message, error_size, "Failed to generate JSON response");
return NULL;
}
return out;
#else
(void)query; (void)request_id; (void)max_rows; (void)timeout_ms;
if (error_message && error_size > 0) {
strncpy(error_message, g_postgres_error, error_size - 1);
error_message[error_size - 1] = '\0';
}
return NULL;
#endif
}
int postgres_db_get_total_event_count_ll(long long* out_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!out_count) return DB_MISUSE;
*out_count = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn, "SELECT COUNT(*) FROM events");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_total_event_count_ll failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
*out_count = strtoll(PQgetvalue(res, 0, 0), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
if (out_count) *out_count = 0;
return DB_ERROR;
#endif
}
int postgres_db_get_event_count_since(time_t cutoff, long long* out_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!out_count) return DB_MISUSE;
*out_count = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
char cutoff_buf[32];
snprintf(cutoff_buf, sizeof(cutoff_buf), "%lld", (long long)cutoff);
const char* params[1] = { cutoff_buf };
PGresult* res = PQexecParams(conn,
"SELECT COUNT(*) FROM events WHERE created_at >= $1::BIGINT",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_event_count_since failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
*out_count = strtoll(PQgetvalue(res, 0, 0), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
(void)cutoff;
if (out_count) *out_count = 0;
return DB_ERROR;
#endif
}
int postgres_db_get_storage_size_bytes(long long* out_size) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!out_size) return DB_MISUSE;
*out_size = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn, "SELECT pg_database_size(current_database())");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_storage_size_bytes failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
*out_size = strtoll(PQgetvalue(res, 0, 0), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
if (out_size) *out_size = 0;
return DB_ERROR;
#endif
}
cJSON* postgres_db_get_event_kind_distribution_rows(long long* out_total_events) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (out_total_events) *out_total_events = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
PGresult* res = PQexec(conn,
"SELECT kind, COUNT(*) AS count FROM events GROUP BY kind ORDER BY count DESC");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_event_kind_distribution_rows failed");
}
return NULL;
}
cJSON* rows = cJSON_CreateArray();
if (!rows) {
PQclear(res);
return NULL;
}
long long total = 0;
int n = PQntuples(res);
for (int i = 0; i < n; i++) {
cJSON* row = cJSON_CreateObject();
if (!row) {
cJSON_Delete(rows);
PQclear(res);
return NULL;
}
int kind = PQgetisnull(res, i, 0) ? 0 : (int)strtol(PQgetvalue(res, i, 0), NULL, 10);
long long count = PQgetisnull(res, i, 1) ? 0 : strtoll(PQgetvalue(res, i, 1), NULL, 10);
total += count;
cJSON_AddNumberToObject(row, "kind", kind);
cJSON_AddNumberToObject(row, "count", (double)count);
cJSON_AddItemToArray(rows, row);
}
PQclear(res);
if (out_total_events) *out_total_events = total;
return rows;
#else
if (out_total_events) *out_total_events = 0;
return NULL;
#endif
}
cJSON* postgres_db_get_top_pubkeys_rows(int limit) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
char limit_buf[16];
snprintf(limit_buf, sizeof(limit_buf), "%d", (limit > 0) ? limit : 10);
const char* params[1] = { limit_buf };
PGresult* res = PQexecParams(conn,
"SELECT pubkey, COUNT(*) AS count FROM events GROUP BY pubkey ORDER BY count DESC LIMIT $1::INT",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_top_pubkeys_rows failed");
}
return NULL;
}
cJSON* rows = cJSON_CreateArray();
if (!rows) {
PQclear(res);
return NULL;
}
int n = PQntuples(res);
for (int i = 0; i < n; i++) {
cJSON* row = cJSON_CreateObject();
if (!row) {
cJSON_Delete(rows);
PQclear(res);
return NULL;
}
const char* pubkey = PQgetisnull(res, i, 0) ? "" : PQgetvalue(res, i, 0);
long long count = PQgetisnull(res, i, 1) ? 0 : strtoll(PQgetvalue(res, i, 1), NULL, 10);
cJSON_AddStringToObject(row, "pubkey", pubkey);
cJSON_AddNumberToObject(row, "count", (double)count);
cJSON_AddItemToArray(rows, row);
}
PQclear(res);
return rows;
#else
(void)limit;
return NULL;
#endif
}
cJSON* postgres_db_get_subscription_details_rows(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
PGresult* res = PQexec(conn,
"SELECT subscription_id AS id, client_ip, COALESCE(filter_json, '[]') AS filter_json, "
"events_sent, created_at, "
"COALESCE(duration, CASE WHEN ended_at IS NULL THEN EXTRACT(EPOCH FROM NOW())::BIGINT - created_at "
"ELSE ended_at - created_at END) AS duration_seconds, "
"wsi_pointer "
"FROM subscriptions "
"WHERE event_type = 'created' "
"ORDER BY created_at DESC");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_subscription_details_rows failed");
}
return NULL;
}
cJSON* rows = cJSON_CreateArray();
if (!rows) {
PQclear(res);
return NULL;
}
int n = PQntuples(res);
for (int i = 0; i < n; i++) {
cJSON* row = cJSON_CreateObject();
if (!row) {
cJSON_Delete(rows);
PQclear(res);
return NULL;
}
cJSON_AddStringToObject(row, "id", PQgetisnull(res, i, 0) ? "" : PQgetvalue(res, i, 0));
cJSON_AddStringToObject(row, "client_ip", PQgetisnull(res, i, 1) ? "" : PQgetvalue(res, i, 1));
cJSON_AddStringToObject(row, "filter_json", PQgetisnull(res, i, 2) ? "[]" : PQgetvalue(res, i, 2));
cJSON_AddNumberToObject(row, "events_sent", PQgetisnull(res, i, 3) ? 0 : (double)strtoll(PQgetvalue(res, i, 3), NULL, 10));
cJSON_AddNumberToObject(row, "created_at", PQgetisnull(res, i, 4) ? 0 : (double)strtoll(PQgetvalue(res, i, 4), NULL, 10));
cJSON_AddNumberToObject(row, "duration_seconds", PQgetisnull(res, i, 5) ? 0 : (double)strtoll(PQgetvalue(res, i, 5), NULL, 10));
cJSON_AddStringToObject(row, "wsi_pointer", PQgetisnull(res, i, 6) ? "" : PQgetvalue(res, i, 6));
cJSON_AddItemToArray(rows, row);
}
PQclear(res);
return rows;
#else
return NULL;
#endif
}
cJSON* postgres_db_get_all_config_rows(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
PGresult* res = PQexec(conn, "SELECT key, value FROM config ORDER BY key");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_all_config_rows failed");
}
return NULL;
}
cJSON* rows = cJSON_CreateArray();
if (!rows) {
PQclear(res);
return NULL;
}
int n = PQntuples(res);
for (int i = 0; i < n; i++) {
cJSON* row = cJSON_CreateObject();
if (!row) {
cJSON_Delete(rows);
PQclear(res);
return NULL;
}
const char* k = PQgetisnull(res, i, 0) ? "" : PQgetvalue(res, i, 0);
const char* v = PQgetisnull(res, i, 1) ? "" : PQgetvalue(res, i, 1);
cJSON_AddStringToObject(row, "key", k);
cJSON_AddStringToObject(row, "value", v);
cJSON_AddItemToArray(rows, row);
}
PQclear(res);
return rows;
#else
return NULL;
#endif
}
char* postgres_db_get_config_value_dup(const char* key) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!key) return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
const char* params[1] = { key };
PGresult* res = PQexecParams(conn,
"SELECT value FROM config WHERE key = $1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_config_value_dup failed");
}
return NULL;
}
char* out = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
out = postgres_strdup(PQgetvalue(res, 0, 0));
}
PQclear(res);
return out;
#else
(void)key;
return NULL;
#endif
}
int postgres_db_set_config_value_full(const char* key, const char* value, const char* data_type,
const char* description, const char* category, int requires_restart) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!key || !value || !data_type) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[6];
params[0] = key;
params[1] = value;
params[2] = data_type;
params[3] = description ? description : "";
params[4] = category ? category : "general";
char rr_buf[16];
snprintf(rr_buf, sizeof(rr_buf), "%d", requires_restart);
params[5] = rr_buf;
PGresult* res = PQexecParams(conn,
"INSERT INTO config (key, value, data_type, description, category, requires_restart, created_at, updated_at) "
"VALUES ($1, $2, $3, $4, $5, $6::int, EXTRACT(EPOCH FROM NOW())::BIGINT, EXTRACT(EPOCH FROM NOW())::BIGINT) "
"ON CONFLICT (key) DO UPDATE SET "
"value = EXCLUDED.value, data_type = EXCLUDED.data_type, description = EXCLUDED.description, "
"category = EXCLUDED.category, requires_restart = EXCLUDED.requires_restart, "
"updated_at = EXTRACT(EPOCH FROM NOW())::BIGINT",
6, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_set_config_value_full failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)key; (void)value; (void)data_type; (void)description; (void)category; (void)requires_restart;
return DB_ERROR;
#endif
}
int postgres_db_update_config_value_only(const char* key, const char* value) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!key || !value) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[2] = { value, key };
PGresult* res = PQexecParams(conn,
"UPDATE config SET value = $1, updated_at = EXTRACT(EPOCH FROM NOW())::BIGINT WHERE key = $2",
2, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_update_config_value_only failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)key; (void)value;
return DB_ERROR;
#endif
}
int postgres_db_upsert_config_value(const char* key, const char* value, const char* data_type) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!key || !value || !data_type) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[3] = { key, value, data_type };
PGresult* res = PQexecParams(conn,
"INSERT INTO config (key, value, data_type, created_at, updated_at) "
"VALUES ($1, $2, $3, EXTRACT(EPOCH FROM NOW())::BIGINT, EXTRACT(EPOCH FROM NOW())::BIGINT) "
"ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, data_type = EXCLUDED.data_type, "
"updated_at = EXTRACT(EPOCH FROM NOW())::BIGINT",
3, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_upsert_config_value failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)key; (void)value; (void)data_type;
return DB_ERROR;
#endif
}
int postgres_db_store_relay_private_key_hex(const char* relay_privkey_hex) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!relay_privkey_hex) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[1] = { relay_privkey_hex };
PGresult* res = PQexecParams(conn,
"INSERT INTO relay_seckey (id, private_key_hex, created_at) "
"VALUES (1, $1, EXTRACT(EPOCH FROM NOW())::BIGINT) "
"ON CONFLICT (id) DO UPDATE SET private_key_hex = EXCLUDED.private_key_hex",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_store_relay_private_key_hex failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)relay_privkey_hex;
return DB_ERROR;
#endif
}
char* postgres_db_get_relay_private_key_hex_dup(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
PGresult* res = PQexec(conn, "SELECT private_key_hex FROM relay_seckey WHERE id = 1");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_relay_private_key_hex_dup failed");
}
return NULL;
}
char* out = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
out = postgres_strdup(PQgetvalue(res, 0, 0));
}
PQclear(res);
return out;
#else
return NULL;
#endif
}
int postgres_db_store_config_event(const cJSON* event) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const cJSON* id_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "id");
const cJSON* pubkey_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "pubkey");
const cJSON* created_at_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "created_at");
const cJSON* kind_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "kind");
const cJSON* content_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "content");
const cJSON* sig_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "sig");
const cJSON* tags_obj = cJSON_GetObjectItemCaseSensitive((cJSON*)event, "tags");
if (!cJSON_IsString(id_obj) || !cJSON_IsString(pubkey_obj) ||
!cJSON_IsNumber(created_at_obj) || !cJSON_IsNumber(kind_obj) ||
!cJSON_IsString(content_obj) || !cJSON_IsString(sig_obj)) {
return DB_MISUSE;
}
char created_buf[64];
char kind_buf[32];
snprintf(created_buf, sizeof(created_buf), "%lld", (long long)cJSON_GetNumberValue(created_at_obj));
snprintf(kind_buf, sizeof(kind_buf), "%d", (int)cJSON_GetNumberValue(kind_obj));
char* tags_json = tags_obj ? cJSON_PrintUnformatted((cJSON*)tags_obj) : NULL;
if (!tags_json) tags_json = postgres_strdup("[]");
char* event_json = cJSON_PrintUnformatted((cJSON*)event);
if (!event_json) {
free(tags_json);
return DB_ERROR;
}
const char* params[9];
params[0] = cJSON_GetStringValue(id_obj);
params[1] = cJSON_GetStringValue(pubkey_obj);
params[2] = created_buf;
params[3] = kind_buf;
params[4] = "addressable";
params[5] = cJSON_GetStringValue(content_obj);
params[6] = cJSON_GetStringValue(sig_obj);
params[7] = tags_json ? tags_json : "[]";
params[8] = event_json;
PGresult* res = PQexecParams(conn,
"INSERT INTO events (id, pubkey, created_at, kind, event_type, content, sig, tags, event_json) "
"VALUES ($1, $2, $3::BIGINT, $4::INT, $5, $6, $7, $8::jsonb, $9) "
"ON CONFLICT (id) DO UPDATE SET "
"pubkey = EXCLUDED.pubkey, created_at = EXCLUDED.created_at, kind = EXCLUDED.kind, "
"event_type = EXCLUDED.event_type, content = EXCLUDED.content, sig = EXCLUDED.sig, "
"tags = EXCLUDED.tags, event_json = EXCLUDED.event_json",
9, NULL, params, NULL, NULL, 0);
free(tags_json);
free(event_json);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_store_config_event failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)event;
return DB_ERROR;
#endif
}
int postgres_db_insert_event_with_json(const char* id, const char* pubkey, long long created_at,
int kind, const char* event_type, const char* content,
const char* sig, const char* tags_json, const char* event_json,
int* out_step_rc, int* out_extended_errcode) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!id || !pubkey || !event_type || !content || !sig || !tags_json || !event_json) {
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_MISUSE;
return DB_MISUSE;
}
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) {
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
}
char created_buf[64];
char kind_buf[32];
snprintf(created_buf, sizeof(created_buf), "%lld", created_at);
snprintf(kind_buf, sizeof(kind_buf), "%d", kind);
const char* params[9] = {
id, pubkey, created_buf, kind_buf, event_type, content, sig, tags_json, event_json
};
int use_tx = (strcmp(event_type, "replaceable") == 0 || strcmp(event_type, "addressable") == 0);
if (use_tx) {
PGresult* tx_begin = PQexec(conn, "BEGIN");
if (!tx_begin || PQresultStatus(tx_begin) != PGRES_COMMAND_OK) {
if (tx_begin) {
postgres_set_error_text(PQresultErrorMessage(tx_begin));
PQclear(tx_begin);
} else {
postgres_set_error_from_conn(conn, "postgres_db_insert_event_with_json BEGIN failed");
}
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
}
PQclear(tx_begin);
PGresult* del_res = NULL;
if (strcmp(event_type, "replaceable") == 0) {
const char* del_params[4] = { pubkey, kind_buf, created_buf, id };
del_res = PQexecParams(conn,
"DELETE FROM events "
"WHERE event_type = 'replaceable' "
" AND pubkey = $1 "
" AND kind = $2::INT "
" AND (created_at < $3::BIGINT "
" OR (created_at = $3::BIGINT AND id > $4))",
4, NULL, del_params, NULL, NULL, 0);
} else {
const char* del_params[5] = { pubkey, kind_buf, created_buf, id, tags_json };
del_res = PQexecParams(conn,
"DELETE FROM events "
"WHERE event_type = 'addressable' "
" AND pubkey = $1 "
" AND kind = $2::INT "
" AND d_tag_value = COALESCE(( "
" SELECT tag->>1 "
" FROM jsonb_array_elements($5::jsonb) AS tag "
" WHERE jsonb_typeof(tag) = 'array' "
" AND jsonb_array_length(tag) >= 2 "
" AND tag->>0 = 'd' "
" LIMIT 1 "
" ), '') "
" AND (created_at < $3::BIGINT "
" OR (created_at = $3::BIGINT AND id > $4))",
5, NULL, del_params, NULL, NULL, 0);
}
if (!del_res || PQresultStatus(del_res) != PGRES_COMMAND_OK) {
if (del_res) {
postgres_set_error_text(PQresultErrorMessage(del_res));
PQclear(del_res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_insert_event_with_json DELETE failed");
}
PGresult* tx_rollback = PQexec(conn, "ROLLBACK");
if (tx_rollback) PQclear(tx_rollback);
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
}
PQclear(del_res);
}
PGresult* res = PQexecParams(conn,
"INSERT INTO events (id, pubkey, created_at, kind, event_type, content, sig, tags, event_json) "
"VALUES ($1, $2, $3::BIGINT, $4::INT, $5, $6, $7, $8::jsonb, $9) "
"ON CONFLICT (id) DO NOTHING",
9, NULL, params, NULL, NULL, 0);
if (!res) {
postgres_set_error_from_conn(conn, "postgres_db_insert_event_with_json INSERT failed");
if (use_tx) {
PGresult* tx_rollback = PQexec(conn, "ROLLBACK");
if (tx_rollback) PQclear(tx_rollback);
}
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
}
ExecStatusType st = PQresultStatus(res);
if (st != PGRES_COMMAND_OK) {
const char* sqlstate = PQresultErrorField(res, PG_DIAG_SQLSTATE);
int is_constraint = (sqlstate && strcmp(sqlstate, "23505") == 0);
if (out_extended_errcode) {
*out_extended_errcode = is_constraint ? DB_CONSTRAINT : DB_ERROR;
}
if (out_step_rc) {
*out_step_rc = is_constraint ? DB_CONSTRAINT : DB_ERROR;
}
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
if (use_tx) {
PGresult* tx_rollback = PQexec(conn, "ROLLBACK");
if (tx_rollback) PQclear(tx_rollback);
}
return is_constraint ? DB_OK : DB_ERROR;
}
const char* tuples = PQcmdTuples(res);
int inserted = (tuples && tuples[0] != '\0') ? atoi(tuples) : 0;
PQclear(res);
if (use_tx) {
PGresult* tx_commit = PQexec(conn, "COMMIT");
if (!tx_commit || PQresultStatus(tx_commit) != PGRES_COMMAND_OK) {
if (tx_commit) {
postgres_set_error_text(PQresultErrorMessage(tx_commit));
PQclear(tx_commit);
} else {
postgres_set_error_from_conn(conn, "postgres_db_insert_event_with_json COMMIT failed");
}
PGresult* tx_rollback = PQexec(conn, "ROLLBACK");
if (tx_rollback) PQclear(tx_rollback);
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
}
PQclear(tx_commit);
}
if (out_step_rc) *out_step_rc = inserted > 0 ? DB_DONE : DB_CONSTRAINT;
if (out_extended_errcode) *out_extended_errcode = inserted > 0 ? DB_OK : DB_CONSTRAINT;
return DB_OK;
#else
(void)id; (void)pubkey; (void)created_at; (void)kind; (void)event_type;
(void)content; (void)sig; (void)tags_json; (void)event_json;
if (out_step_rc) *out_step_rc = DB_ERROR;
if (out_extended_errcode) *out_extended_errcode = DB_ERROR;
return DB_ERROR;
#endif
}
int postgres_db_get_event_time_bounds(long long* out_min_created_at, long long* out_max_created_at) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!out_min_created_at || !out_max_created_at) return DB_MISUSE;
*out_min_created_at = 0;
*out_max_created_at = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn, "SELECT COALESCE(MIN(created_at),0), COALESCE(MAX(created_at),0) FROM events");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_event_time_bounds failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0) {
if (!PQgetisnull(res, 0, 0)) *out_min_created_at = strtoll(PQgetvalue(res, 0, 0), NULL, 10);
if (!PQgetisnull(res, 0, 1)) *out_max_created_at = strtoll(PQgetvalue(res, 0, 1), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
if (out_min_created_at) *out_min_created_at = 0;
if (out_max_created_at) *out_max_created_at = 0;
return DB_ERROR;
#endif
}
int postgres_db_event_id_exists(const char* event_id, int* out_exists) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event_id || !out_exists) return DB_MISUSE;
*out_exists = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[1] = { event_id };
PGresult* res = PQexecParams(conn,
"SELECT 1 FROM events WHERE id = $1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_event_id_exists failed");
}
return DB_ERROR;
}
*out_exists = PQntuples(res) > 0 ? 1 : 0;
PQclear(res);
return DB_OK;
#else
(void)event_id;
if (out_exists) *out_exists = 0;
return DB_ERROR;
#endif
}
cJSON* postgres_db_retrieve_event_by_id(const char* event_id) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event_id) return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
const char* params[1] = { event_id };
PGresult* res = PQexecParams(conn,
"SELECT id, pubkey, created_at, kind, content, sig, tags FROM events WHERE id = $1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_retrieve_event_by_id failed");
}
return NULL;
}
if (PQntuples(res) == 0) {
PQclear(res);
return NULL;
}
cJSON* event = cJSON_CreateObject();
if (!event) {
PQclear(res);
return NULL;
}
cJSON_AddStringToObject(event, "id", PQgetisnull(res, 0, 0) ? "" : PQgetvalue(res, 0, 0));
cJSON_AddStringToObject(event, "pubkey", PQgetisnull(res, 0, 1) ? "" : PQgetvalue(res, 0, 1));
cJSON_AddNumberToObject(event, "created_at", PQgetisnull(res, 0, 2) ? 0 : (double)strtoll(PQgetvalue(res, 0, 2), NULL, 10));
cJSON_AddNumberToObject(event, "kind", PQgetisnull(res, 0, 3) ? 0 : (double)strtol(PQgetvalue(res, 0, 3), NULL, 10));
cJSON_AddStringToObject(event, "content", PQgetisnull(res, 0, 4) ? "" : PQgetvalue(res, 0, 4));
cJSON_AddStringToObject(event, "sig", PQgetisnull(res, 0, 5) ? "" : PQgetvalue(res, 0, 5));
const char* tags_text = PQgetisnull(res, 0, 6) ? "[]" : PQgetvalue(res, 0, 6);
cJSON* tags = cJSON_Parse(tags_text);
if (!tags) tags = cJSON_CreateArray();
cJSON_AddItemToObject(event, "tags", tags);
PQclear(res);
return event;
#else
(void)event_id;
return NULL;
#endif
}
char* postgres_db_get_latest_event_pubkey_for_kind_dup(int kind) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
char kind_buf[32];
snprintf(kind_buf, sizeof(kind_buf), "%d", kind);
const char* params[1] = { kind_buf };
PGresult* res = PQexecParams(conn,
"SELECT pubkey FROM events WHERE kind = $1::INT ORDER BY created_at DESC LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_latest_event_pubkey_for_kind_dup failed");
}
return NULL;
}
char* out = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
out = postgres_strdup(PQgetvalue(res, 0, 0));
}
PQclear(res);
return out;
#else
(void)kind;
return NULL;
#endif
}
int postgres_db_get_config_row_count(int* out_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!out_count) return DB_MISUSE;
*out_count = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn, "SELECT COUNT(*) FROM config");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_get_config_row_count failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
*out_count = (int)strtol(PQgetvalue(res, 0, 0), NULL, 10);
}
PQclear(res);
return DB_OK;
#else
if (out_count) *out_count = 0;
return DB_ERROR;
#endif
}
int postgres_db_store_event_tags_cjson(const char* event_id, const cJSON* tags) {
(void)event_id;
(void)tags;
// PostgreSQL stores tags in events.tags JSONB directly.
return DB_OK;
}
int postgres_db_populate_event_tags_from_existing(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn,
"INSERT INTO event_tags (event_id, tag_name, tag_value, tag_index) "
"SELECT e.id, t.tag->>0 AS tag_name, t.tag->>1 AS tag_value, (t.ord - 1)::INT AS tag_index "
"FROM events e "
"CROSS JOIN LATERAL jsonb_array_elements(e.tags) WITH ORDINALITY AS t(tag, ord) "
"WHERE jsonb_typeof(t.tag) = 'array' "
" AND jsonb_array_length(t.tag) >= 2 "
" AND COALESCE(t.tag->>0, '') <> '' "
" AND COALESCE(t.tag->>1, '') <> '' "
"ON CONFLICT (event_id, tag_name, tag_value, tag_index) DO NOTHING");
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_populate_event_tags_from_existing failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
return DB_ERROR;
#endif
}
int postgres_db_add_auth_rule(const char* rule_type, const char* pattern_type, const char* pattern_value) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!rule_type || !pattern_type || !pattern_value) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[3] = { rule_type, pattern_type, pattern_value };
PGresult* res = PQexecParams(conn,
"INSERT INTO auth_rules (rule_type, pattern_type, pattern_value, active, created_at, updated_at) "
"VALUES ($1, $2, $3, 1, EXTRACT(EPOCH FROM NOW())::BIGINT, EXTRACT(EPOCH FROM NOW())::BIGINT)",
3, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_add_auth_rule failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)rule_type; (void)pattern_type; (void)pattern_value;
return DB_ERROR;
#endif
}
int postgres_db_remove_auth_rule(const char* rule_type, const char* pattern_type, const char* pattern_value) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!rule_type || !pattern_type || !pattern_value) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[3] = { rule_type, pattern_type, pattern_value };
PGresult* res = PQexecParams(conn,
"DELETE FROM auth_rules WHERE rule_type=$1 AND pattern_type=$2 AND pattern_value=$3",
3, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_remove_auth_rule failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)rule_type; (void)pattern_type; (void)pattern_value;
return DB_ERROR;
#endif
}
int postgres_db_delete_wot_whitelist_rules(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn, "DELETE FROM auth_rules WHERE rule_type='wot_whitelist'");
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
return DB_ERROR;
#endif
}
int postgres_db_count_wot_whitelist_rules(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return 0;
PGresult* res = PQexec(conn, "SELECT COUNT(*) FROM auth_rules WHERE rule_type='wot_whitelist' AND active=1");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return 0;
}
int count = (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) ? (int)strtol(PQgetvalue(res, 0, 0), NULL, 10) : 0;
PQclear(res);
return count;
#else
return 0;
#endif
}
int postgres_db_table_exists(const char* table_name, int* out_exists) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!table_name || !out_exists) return DB_MISUSE;
*out_exists = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
const char* params[1] = { table_name };
PGresult* res = PQexecParams(conn,
"SELECT 1 FROM information_schema.tables WHERE table_schema='public' AND table_name=$1 LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_table_exists failed");
}
return DB_ERROR;
}
*out_exists = PQntuples(res) > 0 ? 1 : 0;
PQclear(res);
return DB_OK;
#else
(void)table_name;
if (out_exists) *out_exists = 0;
return DB_ERROR;
#endif
}
char* postgres_db_get_schema_version_dup(void) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
PGresult* res = PQexec(conn, "SELECT value FROM schema_info WHERE key='version'");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
}
return NULL;
}
char* out = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
out = postgres_strdup(PQgetvalue(res, 0, 0));
}
PQclear(res);
return out;
#else
return NULL;
#endif
}
int postgres_db_exec_sql(const char* sql) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!sql) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) {
postgres_set_error_text("postgres_db_exec_sql: no active connection");
return DB_ERROR;
}
PGresult* res = PQexec(conn, sql);
if (!res) {
postgres_set_error_from_conn(conn, "postgres_db_exec_sql failed");
return DB_ERROR;
}
ExecStatusType st = PQresultStatus(res);
int ok = (st == PGRES_COMMAND_OK || st == PGRES_TUPLES_OK);
if (!ok) {
const char* msg = PQresultErrorMessage(res);
postgres_set_error_text((msg && msg[0] != '\0') ? msg : "postgres_db_exec_sql failed");
}
PQclear(res);
return ok ? DB_OK : DB_ERROR;
#else
(void)sql;
return DB_ERROR;
#endif
}
int postgres_db_wal_checkpoint_passive(void) { return DB_OK; }
int postgres_db_wal_checkpoint_truncate(void) { return DB_OK; }
// ---------------------------------------------------------------------------
// Caching relay inbox helpers
//
// These functions back the c-relay-pg inbox poller that destructively
// dequeues events written into the caching_event_inbox table by the external
// caching application. They are PostgreSQL-only; the SQLite dispatch path
// returns errors/no-ops.
// ---------------------------------------------------------------------------
cJSON* postgres_db_caching_inbox_dequeue_batch(int batch_size, int* out_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (out_count) *out_count = 0;
if (batch_size <= 0) {
postgres_set_error_text("postgres_db_caching_inbox_dequeue_batch: batch_size must be positive");
return NULL;
}
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
char batch_buf[16];
snprintf(batch_buf, sizeof(batch_buf), "%d", batch_size);
const char* params[1] = { batch_buf };
PGresult* res = PQexecParams(conn,
"WITH selected AS ("
" SELECT queue_id"
" FROM caching_event_inbox"
" ORDER BY priority ASC, received_at ASC, queue_id ASC"
" LIMIT $1"
" FOR UPDATE"
") "
"DELETE FROM caching_event_inbox AS inbox "
"USING selected "
"WHERE inbox.queue_id = selected.queue_id "
"RETURNING inbox.event_id, inbox.event_json, inbox.source_relay, "
" inbox.source_class, inbox.priority, inbox.received_at",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_caching_inbox_dequeue_batch failed");
}
return NULL;
}
int n = PQntuples(res);
if (n == 0) {
PQclear(res);
if (out_count) *out_count = 0;
return NULL;
}
cJSON* rows = cJSON_CreateArray();
if (!rows) {
PQclear(res);
return NULL;
}
for (int i = 0; i < n; i++) {
cJSON* row = cJSON_CreateObject();
if (!row) {
cJSON_Delete(rows);
PQclear(res);
if (out_count) *out_count = 0;
return NULL;
}
// event_id (col 0): TEXT, never null per schema.
cJSON_AddStringToObject(row, "event_id",
PQgetisnull(res, i, 0) ? "" : PQgetvalue(res, i, 0));
// event_json (col 1): JSONB. PQgetvalue returns the canonical text form.
cJSON_AddStringToObject(row, "event_json",
PQgetisnull(res, i, 1) ? "" : PQgetvalue(res, i, 1));
// source_relay (col 2): nullable TEXT.
if (PQgetisnull(res, i, 2)) {
cJSON_AddNullToObject(row, "source_relay");
} else {
cJSON_AddStringToObject(row, "source_relay", PQgetvalue(res, i, 2));
}
// source_class (col 3): TEXT, never null per schema.
cJSON_AddStringToObject(row, "source_class",
PQgetisnull(res, i, 3) ? "" : PQgetvalue(res, i, 3));
// priority (col 4): SMALLINT.
int priority = PQgetisnull(res, i, 4) ? 0 : (int)strtol(PQgetvalue(res, i, 4), NULL, 10);
cJSON_AddNumberToObject(row, "priority", priority);
// received_at (col 5): BIGINT.
long long received_at = PQgetisnull(res, i, 5) ? 0 : strtoll(PQgetvalue(res, i, 5), NULL, 10);
cJSON_AddNumberToObject(row, "received_at", (double)received_at);
cJSON_AddItemToArray(rows, row);
}
PQclear(res);
if (out_count) *out_count = n;
return rows;
#else
if (out_count) *out_count = 0;
(void)batch_size;
return NULL;
#endif
}
int postgres_db_caching_inbox_pending_counts(int* out_live_count, int* out_backfill_count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (out_live_count) *out_live_count = 0;
if (out_backfill_count) *out_backfill_count = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn,
"SELECT "
" COALESCE(SUM(CASE WHEN priority = 0 THEN 1 ELSE 0 END), 0), "
" COALESCE(SUM(CASE WHEN priority = 1 THEN 1 ELSE 0 END), 0) "
"FROM caching_event_inbox");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_caching_inbox_pending_counts failed");
}
return DB_ERROR;
}
long long live = PQgetisnull(res, 0, 0) ? 0 : strtoll(PQgetvalue(res, 0, 0), NULL, 10);
long long backfill = PQgetisnull(res, 0, 1) ? 0 : strtoll(PQgetvalue(res, 0, 1), NULL, 10);
if (out_live_count) *out_live_count = (int)live;
if (out_backfill_count) *out_backfill_count = (int)backfill;
PQclear(res);
return DB_OK;
#else
if (out_live_count) *out_live_count = 0;
if (out_backfill_count) *out_backfill_count = 0;
return DB_ERROR;
#endif
}
int postgres_db_caching_inbox_oldest_age(int* out_oldest_age_seconds) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (out_oldest_age_seconds) *out_oldest_age_seconds = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn,
"SELECT COALESCE(EXTRACT(EPOCH FROM NOW())::BIGINT - MIN(received_at), 0) "
"FROM caching_event_inbox");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_caching_inbox_oldest_age failed");
}
return DB_ERROR;
}
long long age = PQgetisnull(res, 0, 0) ? 0 : strtoll(PQgetvalue(res, 0, 0), NULL, 10);
if (out_oldest_age_seconds) *out_oldest_age_seconds = (int)age;
PQclear(res);
return DB_OK;
#else
if (out_oldest_age_seconds) *out_oldest_age_seconds = 0;
return DB_ERROR;
#endif
}
int postgres_db_caching_inbox_insert(const char* event_id, const char* event_json,
const char* source_relay, const char* source_class,
int priority) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!event_id || !event_json || !source_class) return DB_MISUSE;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
char priority_buf[16];
snprintf(priority_buf, sizeof(priority_buf), "%d", priority);
const char* params[5] = {
event_id,
event_json,
source_relay, // may be NULL -> libpq sends a SQL NULL
source_class,
priority_buf
};
PGresult* res = PQexecParams(conn,
"INSERT INTO caching_event_inbox "
" (event_id, event_json, source_relay, source_class, priority) "
"VALUES ($1, $2::jsonb, $3, $4, $5) "
"ON CONFLICT (event_id) DO NOTHING",
5, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_caching_inbox_insert failed");
}
return DB_ERROR;
}
PQclear(res);
return DB_OK;
#else
(void)event_id; (void)event_json; (void)source_relay; (void)source_class; (void)priority;
return DB_ERROR;
#endif
}
int postgres_db_caching_backfill_author_counts(int* out_complete, int* out_total) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (out_complete) *out_complete = 0;
if (out_total) *out_total = 0;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR;
PGresult* res = PQexec(conn,
"SELECT backfill_authors_complete, backfill_authors_total "
"FROM caching_service_state WHERE id = 1");
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) {
postgres_set_error_text(PQresultErrorMessage(res));
PQclear(res);
} else {
postgres_set_error_from_conn(conn, "postgres_db_caching_backfill_author_counts failed");
}
return DB_ERROR;
}
if (PQntuples(res) > 0) {
if (out_complete) {
*out_complete = PQgetisnull(res, 0, 0) ? 0
: (int)strtol(PQgetvalue(res, 0, 0), NULL, 10);
}
if (out_total) {
*out_total = PQgetisnull(res, 0, 1) ? 0
: (int)strtol(PQgetvalue(res, 0, 1), NULL, 10);
}
}
PQclear(res);
return DB_OK;
#else
if (out_complete) *out_complete = 0;
if (out_total) *out_total = 0;
return DB_ERROR;
#endif
}
/* ------------------------------------------------------------------ */
/* Profile metadata and outbox relay helpers */
/* ------------------------------------------------------------------ */
/* Returns the most recent kind-0 metadata for a pubkey as a cJSON object
* with name/picture/about/nip05/display_name fields parsed from the event
* content JSON. Returns NULL if no kind-0 event exists or on error.
* Caller must cJSON_Delete() the result. */
cJSON* postgres_db_get_profile_metadata(const char* pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey || pubkey[0] == '\0') return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
const char* params[1] = { pubkey };
PGresult* res = PQexecParams(conn,
"SELECT content FROM events WHERE pubkey = $1 AND kind = 0 "
"ORDER BY created_at DESC LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) PQclear(res);
return NULL;
}
cJSON* result = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
const char* content = PQgetvalue(res, 0, 0);
if (content && content[0] != '\0') {
/* Parse the kind-0 content JSON and extract known fields. */
cJSON* parsed = cJSON_Parse(content);
if (parsed) {
result = cJSON_CreateObject();
if (!result) {
cJSON_Delete(parsed);
PQclear(res);
return NULL;
}
/* Extract standard NIP-01 metadata fields. */
const char* fields[] = {"name", "display_name", "picture",
"about", "nip05", "website", "lud16", "lud06"};
for (size_t i = 0; i < sizeof(fields) / sizeof(fields[0]); i++) {
cJSON* val = cJSON_GetObjectItemCaseSensitive(parsed, fields[i]);
if (val && cJSON_IsString(val) && val->valuestring[0] != '\0') {
cJSON_AddStringToObject(result, fields[i], val->valuestring);
}
}
cJSON_Delete(parsed);
}
}
}
PQclear(res);
return result;
#else
(void)pubkey;
return NULL;
#endif
}
/* Returns relay URLs from the most recent kind-10002 event for a pubkey,
* as a cJSON array of strings. Returns NULL if no kind-10002 event exists
* or on error. Caller must cJSON_Delete() the result. */
cJSON* postgres_db_get_outbox_relays(const char* pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey || pubkey[0] == '\0') return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
const char* params[1] = { pubkey };
PGresult* res = PQexecParams(conn,
"SELECT tags FROM events WHERE pubkey = $1 AND kind = 10002 "
"ORDER BY created_at DESC LIMIT 1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) PQclear(res);
return NULL;
}
cJSON* result = NULL;
if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) {
const char* tags_json = PQgetvalue(res, 0, 0);
if (tags_json && tags_json[0] != '\0') {
cJSON* tags = cJSON_Parse(tags_json);
if (tags && cJSON_IsArray(tags)) {
result = cJSON_CreateArray();
if (!result) {
cJSON_Delete(tags);
PQclear(res);
return NULL;
}
int n = cJSON_GetArraySize(tags);
for (int i = 0; i < n; i++) {
cJSON* tag = cJSON_GetArrayItem(tags, i);
if (!tag || !cJSON_IsArray(tag)) continue;
int tag_len = cJSON_GetArraySize(tag);
if (tag_len < 2) continue;
cJSON* tag_name = cJSON_GetArrayItem(tag, 0);
cJSON* tag_val = cJSON_GetArrayItem(tag, 1);
if (tag_name && cJSON_IsString(tag_name) &&
strcmp(tag_name->valuestring, "r") == 0 &&
tag_val && cJSON_IsString(tag_val) &&
tag_val->valuestring[0] != '\0') {
/* Check for a marker (3rd element) - skip "read" only
* relays, keep "write" and unmarked (both) relays. */
int is_write = 1;
if (tag_len >= 3) {
cJSON* marker = cJSON_GetArrayItem(tag, 2);
if (marker && cJSON_IsString(marker) &&
strcmp(marker->valuestring, "read") == 0) {
is_write = 0;
}
}
if (is_write) {
cJSON_AddItemToArray(result,
cJSON_CreateString(tag_val->valuestring));
}
}
}
cJSON_Delete(tags);
}
}
}
PQclear(res);
return result;
#else
(void)pubkey;
return NULL;
#endif
}
/* ---- Profile cache (profiles table) -------------------------------------- */
/* Forward declaration for config access (defined in config.c). */
const char* get_config_value(const char* key);
/* Resolve the display name per the profile_name_preference config key.
* Returns display_name if preference is "display_name" (default) and it is
* non-empty, else name. If preference is "name", returns name if non-empty,
* else display_name. Returns "" if neither is set. Never returns NULL. */
static const char* profile_resolve_best_name(const char* name,
const char* display_name) {
if (!name) name = "";
if (!display_name) display_name = "";
const char* pref = get_config_value("profile_name_preference");
if (pref && strcmp(pref, "name") == 0) {
/* Prefer name, fall back to display_name. */
if (name[0] != '\0') return name;
return display_name;
}
/* Default: prefer display_name, fall back to name. */
if (display_name[0] != '\0') return display_name;
return name;
}
/* Build a profile cJSON object from a PGresult row. The columns must be in
* the order: pubkey, name, display_name, about, picture, banner, nip05,
* website, lud16, lud06. Returns a new object (caller owns) or NULL. */
static cJSON* profile_obj_from_pgrow(PGresult* res, int row) {
cJSON* obj = cJSON_CreateObject();
if (!obj) return NULL;
const char* pubkey = PQgetvalue(res, row, 0);
const char* name = PQgetisnull(res, row, 1) ? "" : PQgetvalue(res, row, 1);
const char* display_name = PQgetisnull(res, row, 2) ? "" : PQgetvalue(res, row, 2);
cJSON_AddStringToObject(obj, "pubkey", pubkey ? pubkey : "");
cJSON_AddStringToObject(obj, "name", name);
cJSON_AddStringToObject(obj, "display_name", display_name);
cJSON_AddStringToObject(obj, "best_name",
profile_resolve_best_name(name, display_name));
/* Optional fields — only add when non-empty. */
const char* opt_fields[] = {"about", "picture", "banner", "nip05",
"website", "lud16", "lud06"};
for (size_t i = 0; i < sizeof(opt_fields) / sizeof(opt_fields[0]); i++) {
int col = (int)(i + 3);
if (!PQgetisnull(res, row, col)) {
const char* val = PQgetvalue(res, row, col);
if (val && val[0] != '\0') {
cJSON_AddStringToObject(obj, opt_fields[i], val);
}
}
}
return obj;
}
/* Returns a single cached profile from the profiles table, or NULL if no
* profile is cached or on error. Caller must cJSON_Delete(). */
cJSON* postgres_db_get_profile(const char* pubkey) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkey || pubkey[0] == '\0') return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
const char* params[1] = { pubkey };
PGresult* res = PQexecParams(conn,
"SELECT pubkey, name, display_name, about, picture, banner, "
"nip05, website, lud16, lud06 FROM profiles WHERE pubkey = $1",
1, NULL, params, NULL, NULL, 0);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) PQclear(res);
return NULL;
}
cJSON* result = NULL;
if (PQntuples(res) > 0) {
result = profile_obj_from_pgrow(res, 0);
}
PQclear(res);
return result;
#else
(void)pubkey;
return NULL;
#endif
}
/* Batch lookup: one query for many pubkeys. Returns a cJSON object keyed by
* pubkey hex -> profile object. Pubkeys with no cached profile are absent.
* Caller must cJSON_Delete(). Returns NULL on error. */
cJSON* postgres_db_get_profiles(const char** pubkeys, int count) {
#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ)
if (!pubkeys || count <= 0) return NULL;
PGconn* conn = postgres_db_active_connection();
if (!conn || PQstatus(conn) != CONNECTION_OK) return NULL;
/* Build a PostgreSQL array literal: {pk1,pk2,...}
* Each pubkey is 64 hex chars (no quotes needed), but we validate to
* prevent SQL injection since we're building the literal manually. */
size_t buf_size = (size_t)count * 66 + 4;
char* array_lit = (char*)malloc(buf_size);
if (!array_lit) return NULL;
size_t pos = 0;
array_lit[pos++] = '{';
for (int i = 0; i < count; i++) {
if (!pubkeys[i] || strlen(pubkeys[i]) != 64) {
/* Skip invalid pubkeys rather than including them. */
continue;
}
/* Validate hex characters. */
int valid = 1;
for (int j = 0; j < 64; j++) {
char c = pubkeys[i][j];
if (!((c >= '0' && c <= '9') || (c >= 'a' && c <= 'f') ||
(c >= 'A' && c <= 'F'))) {
valid = 0;
break;
}
}
if (!valid) continue;
if (pos > 1) array_lit[pos++] = ',';
memcpy(array_lit + pos, pubkeys[i], 64);
pos += 64;
}
array_lit[pos++] = '}';
array_lit[pos] = '\0';
/* If no valid pubkeys, return an empty object. */
if (pos <= 2) {
free(array_lit);
cJSON* empty = cJSON_CreateObject();
return empty;
}
const char* params[1] = { array_lit };
PGresult* res = PQexecParams(conn,
"SELECT pubkey, name, display_name, about, picture, banner, "
"nip05, website, lud16, lud06 FROM profiles "
"WHERE pubkey = ANY($1::text[])",
1, NULL, params, NULL, NULL, 0);
free(array_lit);
if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) {
if (res) PQclear(res);
return NULL;
}
cJSON* result = cJSON_CreateObject();
if (!result) {
PQclear(res);
return NULL;
}
int n = PQntuples(res);
for (int i = 0; i < n; i++) {
cJSON* obj = profile_obj_from_pgrow(res, i);
if (obj) {
const char* pk = PQgetvalue(res, i, 0);
if (pk && pk[0] != '\0') {
cJSON_AddItemToObject(result, pk, obj);
} else {
cJSON_Delete(obj);
}
}
}
PQclear(res);
return result;
#else
(void)pubkeys;
(void)count;
return NULL;
#endif
}