Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d8f477c6cf | ||
|
|
040eeadb13 |
+29
-10
@@ -881,6 +881,22 @@ int store_event(cJSON* event) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
// Fast duplicate check: returns 1 if event ID already exists in DB, 0 if not.
|
||||
// Uses the primary key index — single B-tree lookup, ~10μs.
|
||||
// Call this BEFORE signature verification to skip expensive crypto on duplicates.
|
||||
int event_id_exists_in_db(const char* event_id) {
|
||||
if (!g_db || !event_id || strlen(event_id) != 64) return 0;
|
||||
|
||||
sqlite3_stmt* stmt;
|
||||
const char* sql = "SELECT 1 FROM events WHERE id=? LIMIT 1";
|
||||
if (sqlite3_prepare_v2(g_db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0;
|
||||
|
||||
sqlite3_bind_text(stmt, 1, event_id, 64, SQLITE_STATIC);
|
||||
int exists = (sqlite3_step(stmt) == SQLITE_ROW) ? 1 : 0;
|
||||
sqlite3_finalize(stmt);
|
||||
return exists;
|
||||
}
|
||||
|
||||
// Populate event_tags from existing events (run once at startup)
|
||||
int populate_event_tags_from_existing(void) {
|
||||
if (!g_db) return -1;
|
||||
@@ -1552,21 +1568,24 @@ int handle_req_message(const char* sub_id, cJSON* filters, struct lws *wsi, stru
|
||||
}
|
||||
}
|
||||
|
||||
// Build EVENT message using string concatenation (much faster than cJSON operations)
|
||||
// Build EVENT message using zero-copy path: allocate with LWS_PRE prefix,
|
||||
// write directly, transfer ownership to queue — no memcpy.
|
||||
// Format: ["EVENT","<sub_id>",<event_json>]
|
||||
size_t sub_id_len = strlen(sub_id);
|
||||
size_t event_json_len = strlen(event_json_str);
|
||||
size_t msg_len = 10 + sub_id_len + 3 + event_json_len + 1; // ["EVENT",""] + sub_id + "," + event_json + ]
|
||||
|
||||
char* msg_str = malloc(msg_len + 1);
|
||||
if (msg_str) {
|
||||
snprintf(msg_str, msg_len + 1, "[\"EVENT\",\"%s\",%s]", sub_id, event_json_str);
|
||||
|
||||
// Use proper message queue system instead of direct lws_write
|
||||
if (queue_message(wsi, pss, msg_str, strlen(msg_str), LWS_WRITE_TEXT) != 0) {
|
||||
size_t msg_len = 10 + sub_id_len + 3 + event_json_len + 1;
|
||||
|
||||
unsigned char* buf = malloc(LWS_PRE + msg_len + 1);
|
||||
if (buf) {
|
||||
char* msg_ptr = (char*)(buf + LWS_PRE);
|
||||
snprintf(msg_ptr, msg_len + 1, "[\"EVENT\",\"%s\",%s]", sub_id, event_json_str);
|
||||
size_t actual_len = strlen(msg_ptr);
|
||||
|
||||
// queue_message_take_ownership takes buf ownership — no memcpy, no free needed here
|
||||
if (queue_message_take_ownership(wsi, pss, buf, actual_len, LWS_WRITE_TEXT) != 0) {
|
||||
DEBUG_ERROR("Failed to queue EVENT message for sub=%s", sub_id);
|
||||
// buf already freed by queue_message_take_ownership on failure
|
||||
}
|
||||
free(msg_str);
|
||||
}
|
||||
|
||||
cJSON_Delete(event);
|
||||
|
||||
+2
-2
@@ -13,8 +13,8 @@
|
||||
// Using CRELAY_ prefix to avoid conflicts with nostr_core_lib VERSION macros
|
||||
#define CRELAY_VERSION_MAJOR 1
|
||||
#define CRELAY_VERSION_MINOR 2
|
||||
#define CRELAY_VERSION_PATCH 9
|
||||
#define CRELAY_VERSION "v1.2.9"
|
||||
#define CRELAY_VERSION_PATCH 11
|
||||
#define CRELAY_VERSION "v1.2.11"
|
||||
|
||||
// Relay metadata (authoritative source for NIP-11 information)
|
||||
#define RELAY_NAME "C-Relay"
|
||||
|
||||
+10
-8
@@ -876,7 +876,6 @@ int broadcast_event_to_subscriptions(cJSON* event) {
|
||||
// Serialize event once per subscription using pre-serialized event_json if available,
|
||||
// otherwise fall back to cJSON serialization.
|
||||
// Format: ["EVENT","<sub_id>",<event_json>]
|
||||
cJSON* event_id_obj = cJSON_GetObjectItemCaseSensitive(event, "id");
|
||||
const char* event_json_str = NULL;
|
||||
char* event_json_allocated = NULL;
|
||||
|
||||
@@ -895,16 +894,19 @@ int broadcast_event_to_subscriptions(cJSON* event) {
|
||||
size_t event_json_len = strlen(event_json_str);
|
||||
// ["EVENT","<sub_id>",<event_json>]
|
||||
size_t msg_len = 10 + sub_id_len + 3 + event_json_len + 1;
|
||||
char* msg_str = malloc(msg_len + 1);
|
||||
if (msg_str) {
|
||||
snprintf(msg_str, msg_len + 1, "[\"EVENT\",\"%s\",%s]", current_temp->id, event_json_str);
|
||||
size_t actual_len = strlen(msg_str);
|
||||
|
||||
// Zero-copy: allocate with LWS_PRE prefix, write directly, transfer ownership to queue
|
||||
unsigned char* buf = malloc(LWS_PRE + msg_len + 1);
|
||||
if (buf) {
|
||||
char* msg_ptr = (char*)(buf + LWS_PRE);
|
||||
snprintf(msg_ptr, msg_len + 1, "[\"EVENT\",\"%s\",%s]", current_temp->id, event_json_str);
|
||||
size_t actual_len = strlen(msg_ptr);
|
||||
|
||||
DEBUG_TRACE("WS_FRAME_SEND: type=EVENT sub=%s len=%zu", current_temp->id, actual_len);
|
||||
|
||||
// Queue message for proper libwebsockets pattern
|
||||
struct per_session_data* pss = (struct per_session_data*)lws_wsi_user(current_temp->wsi);
|
||||
if (queue_message(current_temp->wsi, pss, msg_str, actual_len, LWS_WRITE_TEXT) == 0) {
|
||||
// queue_message_take_ownership takes buf ownership — no memcpy, no free needed here
|
||||
if (queue_message_take_ownership(current_temp->wsi, pss, buf, actual_len, LWS_WRITE_TEXT) == 0) {
|
||||
broadcasts++;
|
||||
|
||||
// Update events sent counter for this subscription
|
||||
@@ -922,8 +924,8 @@ int broadcast_event_to_subscriptions(cJSON* event) {
|
||||
pthread_mutex_unlock(&g_subscription_manager.subscriptions_lock);
|
||||
} else {
|
||||
DEBUG_ERROR("Failed to queue EVENT message for sub=%s", current_temp->id);
|
||||
// buf already freed by queue_message_take_ownership on failure
|
||||
}
|
||||
free(msg_str);
|
||||
}
|
||||
}
|
||||
if (event_json_allocated) {
|
||||
|
||||
@@ -64,6 +64,7 @@ int remove_subscription_from_manager(const char* sub_id, struct lws* wsi);
|
||||
// Forward declarations for event handling
|
||||
int handle_event_message(cJSON* event, char* error_message, size_t error_size);
|
||||
int nostr_validate_unified_request(const char* json_string, size_t json_length);
|
||||
int event_id_exists_in_db(const char* event_id);
|
||||
|
||||
// Forward declarations for admin event processing
|
||||
int process_admin_event_in_config(cJSON* event, char* error_message, size_t error_size, struct lws* wsi);
|
||||
@@ -194,6 +195,66 @@ int queue_message(struct lws* wsi, struct per_session_data* pss, const char* mes
|
||||
return 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Zero-copy variant of queue_message. The caller allocates a buffer of
|
||||
* (LWS_PRE + length) bytes, writes the message at (buf + LWS_PRE), then
|
||||
* passes ownership to the queue. The queue will free buf when done.
|
||||
* No memcpy is performed — eliminates one copy per queued message.
|
||||
*
|
||||
* @param wsi WebSocket instance
|
||||
* @param pss Per-session data containing message queue
|
||||
* @param buf Pre-allocated buffer of size (LWS_PRE + length); ownership transferred
|
||||
* @param length Length of message (NOT including LWS_PRE)
|
||||
* @param type LWS_WRITE_* type
|
||||
* @return 0 on success, -1 on error (buf is freed on error)
|
||||
*/
|
||||
int queue_message_take_ownership(struct lws* wsi, struct per_session_data* pss, unsigned char* buf, size_t length, enum lws_write_protocol type) {
|
||||
if (!wsi || !pss || !buf || length == 0) {
|
||||
DEBUG_ERROR("queue_message_take_ownership: invalid parameters");
|
||||
free(buf);
|
||||
return -1;
|
||||
}
|
||||
|
||||
// Drop message if queue is full
|
||||
if (pss->message_queue_count >= MAX_MESSAGE_QUEUE_SIZE) {
|
||||
DEBUG_WARN("queue_message_take_ownership: queue full (%d), dropping message",
|
||||
pss->message_queue_count);
|
||||
free(buf);
|
||||
return -1;
|
||||
}
|
||||
|
||||
struct message_queue_node* node = malloc(sizeof(struct message_queue_node));
|
||||
if (!node) {
|
||||
DEBUG_ERROR("queue_message_take_ownership: failed to allocate queue node");
|
||||
free(buf);
|
||||
return -1;
|
||||
}
|
||||
|
||||
node->data = buf; // buf already has LWS_PRE prefix — no copy needed
|
||||
node->length = length;
|
||||
node->type = type;
|
||||
node->next = NULL;
|
||||
|
||||
pthread_mutex_lock(&pss->session_lock);
|
||||
if (!pss->message_queue_head) {
|
||||
pss->message_queue_head = node;
|
||||
pss->message_queue_tail = node;
|
||||
} else {
|
||||
pss->message_queue_tail->next = node;
|
||||
pss->message_queue_tail = node;
|
||||
}
|
||||
pss->message_queue_count++;
|
||||
pthread_mutex_unlock(&pss->session_lock);
|
||||
|
||||
if (!pss->writeable_requested) {
|
||||
pss->writeable_requested = 1;
|
||||
lws_callback_on_writable(wsi);
|
||||
}
|
||||
|
||||
DEBUG_TRACE("Queued message (zero-copy): len=%zu, queue_count=%d", length, pss->message_queue_count);
|
||||
return 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Process message queue when the socket becomes writeable.
|
||||
* This function is called from LWS_CALLBACK_SERVER_WRITEABLE.
|
||||
@@ -609,6 +670,27 @@ static int nostr_relay_callback(struct lws *wsi, enum lws_callback_reasons reaso
|
||||
return 0;
|
||||
}
|
||||
|
||||
// Early duplicate check: skip expensive crypto for events already in DB.
|
||||
// The event id is a 64-char hex string — look it up via primary key index (~10μs).
|
||||
// This must happen AFTER we have the id string but BEFORE signature verification.
|
||||
cJSON* id_obj_dup = cJSON_GetObjectItemCaseSensitive(event, "id");
|
||||
if (id_obj_dup && cJSON_IsString(id_obj_dup)) {
|
||||
const char* event_id_str = cJSON_GetStringValue(id_obj_dup);
|
||||
if (event_id_str && event_id_exists_in_db(event_id_str)) {
|
||||
// Already have this event — send OK true and skip crypto
|
||||
DEBUG_TRACE("Duplicate event %s — skipping signature verification", event_id_str);
|
||||
char ok_msg[128];
|
||||
snprintf(ok_msg, sizeof(ok_msg),
|
||||
"[\"OK\",\"%s\",true,\"duplicate: already have this event\"]",
|
||||
event_id_str);
|
||||
size_t ok_len = strlen(ok_msg);
|
||||
queue_message(wsi, pss, ok_msg, ok_len, LWS_WRITE_TEXT);
|
||||
free(event_json_str);
|
||||
cJSON_Delete(json);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
// Call unified validator with JSON string
|
||||
size_t event_json_len = strlen(event_json_str);
|
||||
int validation_result = nostr_validate_unified_request(event_json_str, event_json_len);
|
||||
@@ -1344,6 +1426,25 @@ static int nostr_relay_callback(struct lws *wsi, enum lws_callback_reasons reaso
|
||||
return 0;
|
||||
}
|
||||
|
||||
// Early duplicate check: skip expensive crypto for events already in DB.
|
||||
cJSON* id_obj_dup2 = cJSON_GetObjectItemCaseSensitive(event, "id");
|
||||
if (id_obj_dup2 && cJSON_IsString(id_obj_dup2)) {
|
||||
const char* event_id_str2 = cJSON_GetStringValue(id_obj_dup2);
|
||||
if (event_id_str2 && event_id_exists_in_db(event_id_str2)) {
|
||||
DEBUG_TRACE("Duplicate event %s — skipping signature verification", event_id_str2);
|
||||
char ok_msg2[128];
|
||||
snprintf(ok_msg2, sizeof(ok_msg2),
|
||||
"[\"OK\",\"%s\",true,\"duplicate: already have this event\"]",
|
||||
event_id_str2);
|
||||
size_t ok_len2 = strlen(ok_msg2);
|
||||
queue_message(wsi, pss, ok_msg2, ok_len2, LWS_WRITE_TEXT);
|
||||
free(event_json_str);
|
||||
cJSON_Delete(json);
|
||||
free(message);
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
// Call unified validator with JSON string
|
||||
size_t event_json_len = strlen(event_json_str);
|
||||
int validation_result = nostr_validate_unified_request(event_json_str, event_json_len);
|
||||
|
||||
@@ -107,6 +107,10 @@ int start_websocket_relay(int port_override, int strict_port);
|
||||
int queue_message(struct lws* wsi, struct per_session_data* pss, const char* message, size_t length, enum lws_write_protocol type);
|
||||
int process_message_queue(struct lws* wsi, struct per_session_data* pss);
|
||||
|
||||
// Zero-copy variant: caller allocates (LWS_PRE + length) bytes, writes message at buf+LWS_PRE,
|
||||
// then passes ownership to the queue. The queue will free buf when done. No memcpy performed.
|
||||
int queue_message_take_ownership(struct lws* wsi, struct per_session_data* pss, unsigned char* buf, size_t length, enum lws_write_protocol type);
|
||||
|
||||
// Auth rules checking function from request_validator.c
|
||||
int check_database_auth_rules(const char *pubkey, const char *operation, const char *resource_hash);
|
||||
|
||||
|
||||
Reference in New Issue
Block a user