From 81f2dc52dc42d01c89dff45a5407ec40b8863052 Mon Sep 17 00:00:00 2001 From: Your Name Date: Tue, 19 May 2026 02:31:19 +0530 Subject: feat: local Nostr relay with relay selection, sync, and integration tests Local Nostr relay (NIP-01) on port 4869 with LittleFS 4MB storage. All events published locally first, then synced to public relays via REQ-diff. Relay selection via NIP-11 HTTP probing with NIP-77 scoring and auto-failover. Components: - wisp_relay: 16-file local relay (ws_server, storage_engine, sub_manager, broadcaster, relay_validator, router, handlers, rate_limiter, nip11, deletion, flash_monitor, relay_types) - esp_littlefs: LittleFS VFS integration (git submodule) - negentropy: for future NIP-77 binary sync (git submodule) New source files: - local_relay.c/h: thin wrapper for relay init/start/publish - relay_selector.c/h: NIP-11 probe + scoring + auto-failover - sync_manager.c/h: REQ-diff sync (primary 30min, fallback 6h) Bug fixes: - config.c: use-after-free (cJSON_Delete before seed_relays/sync parsing) - local_relay: moved init to app_main for boot-time start (not gated on STA IP) Flash layout: 4MB LittleFS partition at 0x500000 for relay_store Test results (Board B, live hardware): - Smoke: ping + HTTP 4869 + NIP-11: PASS - NIP-11 info document: 10/11 PASS - WS pub/sub (connect, REQ/EOSE, EVENT/OK, CLOSE, concurrent): 6/6 PASS - Unit tests (relay_validator + relay_selector): 13/13 PASS Hardware test make targets in physical-router-test-automation/: - make relay-build, relay-flash-b, relay-test-smoke/nip11/pubsub/sync/full --- components/wisp_relay/CMakeLists.txt | 16 ++ components/wisp_relay/broadcaster.c | 33 +++ components/wisp_relay/broadcaster.h | 11 + components/wisp_relay/deletion.c | 190 ++++++++++++++ components/wisp_relay/deletion.h | 11 + components/wisp_relay/flash_monitor.c | 30 +++ components/wisp_relay/flash_monitor.h | 16 ++ components/wisp_relay/handlers.c | 203 +++++++++++++++ components/wisp_relay/handlers.h | 10 + components/wisp_relay/idf_component.yml | 1 + components/wisp_relay/nip11_relay.c | 53 ++++ components/wisp_relay/nip11_relay.h | 9 + components/wisp_relay/rate_limiter.c | 98 ++++++++ components/wisp_relay/rate_limiter.h | 40 +++ components/wisp_relay/relay_core.h | 27 ++ components/wisp_relay/relay_types.c | 21 ++ components/wisp_relay/relay_types.h | 43 ++++ components/wisp_relay/relay_validator.c | 176 +++++++++++++ components/wisp_relay/relay_validator.h | 45 ++++ components/wisp_relay/router.c | 140 +++++++++++ components/wisp_relay/router.h | 19 ++ components/wisp_relay/storage_engine.c | 402 ++++++++++++++++++++++++++++++ components/wisp_relay/storage_engine.h | 88 +++++++ components/wisp_relay/sub_manager.c | 272 ++++++++++++++++++++ components/wisp_relay/sub_manager.h | 92 +++++++ components/wisp_relay/ws_server.c | 426 ++++++++++++++++++++++++++++++++ components/wisp_relay/ws_server.h | 41 +++ 27 files changed, 2513 insertions(+) create mode 100644 components/wisp_relay/CMakeLists.txt create mode 100644 components/wisp_relay/broadcaster.c create mode 100644 components/wisp_relay/broadcaster.h create mode 100644 components/wisp_relay/deletion.c create mode 100644 components/wisp_relay/deletion.h create mode 100644 components/wisp_relay/flash_monitor.c create mode 100644 components/wisp_relay/flash_monitor.h create mode 100644 components/wisp_relay/handlers.c create mode 100644 components/wisp_relay/handlers.h create mode 100644 components/wisp_relay/idf_component.yml create mode 100644 components/wisp_relay/nip11_relay.c create mode 100644 components/wisp_relay/nip11_relay.h create mode 100644 components/wisp_relay/rate_limiter.c create mode 100644 components/wisp_relay/rate_limiter.h create mode 100644 components/wisp_relay/relay_core.h create mode 100644 components/wisp_relay/relay_types.c create mode 100644 components/wisp_relay/relay_types.h create mode 100644 components/wisp_relay/relay_validator.c create mode 100644 components/wisp_relay/relay_validator.h create mode 100644 components/wisp_relay/router.c create mode 100644 components/wisp_relay/router.h create mode 100644 components/wisp_relay/storage_engine.c create mode 100644 components/wisp_relay/storage_engine.h create mode 100644 components/wisp_relay/sub_manager.c create mode 100644 components/wisp_relay/sub_manager.h create mode 100644 components/wisp_relay/ws_server.c create mode 100644 components/wisp_relay/ws_server.h (limited to 'components/wisp_relay') diff --git a/components/wisp_relay/CMakeLists.txt b/components/wisp_relay/CMakeLists.txt new file mode 100644 index 0000000..5da9a9c --- /dev/null +++ b/components/wisp_relay/CMakeLists.txt @@ -0,0 +1,16 @@ +idf_component_register( + SRCS "ws_server.c" + "storage_engine.c" + "sub_manager.c" + "broadcaster.c" + "rate_limiter.c" + "nip11_relay.c" + "deletion.c" + "flash_monitor.c" + "relay_validator.c" + "router.c" + "handlers.c" + "relay_types.c" + INCLUDE_DIRS "." + REQUIRES esp_http_server esp_timer nvs_flash log json esp_littlefs mbedtls secp256k1 +) diff --git a/components/wisp_relay/broadcaster.c b/components/wisp_relay/broadcaster.c new file mode 100644 index 0000000..738cbdb --- /dev/null +++ b/components/wisp_relay/broadcaster.c @@ -0,0 +1,33 @@ +#include "broadcaster.h" +#include "relay_core.h" +#include "router.h" +#include "sub_manager.h" +#include "esp_log.h" + +static const char *TAG = "broadcaster"; + +void broadcaster_fanout_json(relay_ctx_t *ctx, const char *event_json, + size_t event_len, int event_kind, + const char *event_pubkey_hex, + uint64_t event_created_at) +{ + if (!ctx || !ctx->sub_manager) return; + + sub_match_result_t matches; + sub_manager_match_json(ctx->sub_manager, event_json, event_len, event_kind, + event_pubkey_hex, event_created_at, &matches); + + if (matches.count == 0) { + ESP_LOGD(TAG, "No subscribers for event kind=%d", event_kind); + return; + } + + ESP_LOGD(TAG, "Broadcasting event kind=%d to %d subscriptions", + event_kind, matches.count); + + for (uint8_t i = 0; i < matches.count; i++) { + sub_match_entry_t *entry = &matches.matches[i]; + router_send_event(ctx, entry->conn_fd, entry->sub_id, + event_json, event_len); + } +} diff --git a/components/wisp_relay/broadcaster.h b/components/wisp_relay/broadcaster.h new file mode 100644 index 0000000..0b29f71 --- /dev/null +++ b/components/wisp_relay/broadcaster.h @@ -0,0 +1,11 @@ +#ifndef BROADCASTER_H +#define BROADCASTER_H + +#include "relay_core.h" + +void broadcaster_fanout_json(relay_ctx_t *ctx, const char *event_json, + size_t event_len, int event_kind, + const char *event_pubkey_hex, + uint64_t event_created_at); + +#endif diff --git a/components/wisp_relay/deletion.c b/components/wisp_relay/deletion.c new file mode 100644 index 0000000..7ad3c22 --- /dev/null +++ b/components/wisp_relay/deletion.c @@ -0,0 +1,190 @@ +#include "deletion.h" +#include "relay_types.h" +#include "cJSON.h" +#include "esp_log.h" +#include +#include +#include + +static const char *TAG = "deletion"; + +static int extract_event_id_field(const char *event_json, size_t len, + uint8_t id_out[32]) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, len); + if (!obj) return -1; + cJSON *id_item = cJSON_GetObjectItem(obj, "id"); + if (!id_item || !cJSON_IsString(id_item) || strlen(id_item->valuestring) != 64) { + cJSON_Delete(obj); + return -1; + } + int ret = relay_hex_to_bytes(id_item->valuestring, 64, id_out, 32); + cJSON_Delete(obj); + return ret; +} + +static char *extract_pubkey_hex(const char *event_json, size_t len) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, len); + if (!obj) return NULL; + cJSON *pk = cJSON_GetObjectItem(obj, "pubkey"); + char *result = NULL; + if (pk && cJSON_IsString(pk)) result = strdup(pk->valuestring); + cJSON_Delete(obj); + return result; +} + +static int delete_by_e_tags(storage_engine_t *storage, const char *event_json, + size_t len, const char *deleter_pubkey) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, len); + if (!obj) return 0; + + cJSON *tags = cJSON_GetObjectItem(obj, "tags"); + if (!tags || !cJSON_IsArray(tags)) { cJSON_Delete(obj); return 0; } + + int deleted = 0; + int array_size = cJSON_GetArraySize(tags); + + for (int i = 0; i < array_size; i++) { + cJSON *tag = cJSON_GetArrayItem(tags, i); + if (!cJSON_IsArray(tag)) continue; + cJSON *tag_name = cJSON_GetArrayItem(tag, 0); + if (!tag_name || !cJSON_IsString(tag_name)) continue; + if (strcmp(tag_name->valuestring, "e") != 0) continue; + + cJSON *tag_val = cJSON_GetArrayItem(tag, 1); + if (!tag_val || !cJSON_IsString(tag_val)) continue; + + uint8_t event_id[32]; + if (relay_hex_to_bytes(tag_val->valuestring, 64, event_id, 32) != 0) continue; + + storage_error_t err = storage_delete_event(storage, event_id); + if (err == STORAGE_OK) { + deleted++; + ESP_LOGI(TAG, "Deleted event: %.16s...", tag_val->valuestring); + } + } + + cJSON_Delete(obj); + return deleted; +} + +static int delete_by_a_tags(storage_engine_t *storage, const char *event_json, + size_t len, const char *deleter_pubkey, + uint64_t created_at) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, len); + if (!obj) return 0; + + cJSON *tags = cJSON_GetObjectItem(obj, "tags"); + if (!tags || !cJSON_IsArray(tags)) { cJSON_Delete(obj); return 0; } + + int deleted = 0; + int array_size = cJSON_GetArraySize(tags); + + for (int i = 0; i < array_size; i++) { + cJSON *tag = cJSON_GetArrayItem(tags, i); + if (!cJSON_IsArray(tag)) continue; + cJSON *tag_name = cJSON_GetArrayItem(tag, 0); + if (!tag_name || !cJSON_IsString(tag_name)) continue; + if (strcmp(tag_name->valuestring, "a") != 0) continue; + + cJSON *tag_val = cJSON_GetArrayItem(tag, 1); + if (!tag_val || !cJSON_IsString(tag_val)) continue; + + int32_t kind; + char pubkey[65] = {0}; + char d_tag[256] = ""; + if (sscanf(tag_val->valuestring, "%" SCNd32 ":%64[^:]:%255s", + &kind, pubkey, d_tag) < 2) + continue; + + if (strcmp(pubkey, deleter_pubkey) != 0) continue; + + char **results = NULL; + uint16_t count = 0; + storage_query_events_json(storage, kind, pubkey, 100, &results, &count); + for (uint16_t e = 0; e < count; e++) { + if (storage_delete_event(storage, (const uint8_t *)results[e]) == STORAGE_OK) { + deleted++; + } + } + storage_free_query_results(results, count); + } + + cJSON_Delete(obj); + return deleted; +} + +static int delete_by_k_tags(storage_engine_t *storage, const char *event_json, + size_t len, const char *deleter_pubkey, + uint64_t created_at) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, len); + if (!obj) return 0; + + cJSON *tags = cJSON_GetObjectItem(obj, "tags"); + if (!tags || !cJSON_IsArray(tags)) { cJSON_Delete(obj); return 0; } + + int deleted = 0; + int array_size = cJSON_GetArraySize(tags); + + for (int i = 0; i < array_size; i++) { + cJSON *tag = cJSON_GetArrayItem(tags, i); + if (!cJSON_IsArray(tag)) continue; + cJSON *tag_name = cJSON_GetArrayItem(tag, 0); + if (!tag_name || !cJSON_IsString(tag_name)) continue; + if (strcmp(tag_name->valuestring, "k") != 0) continue; + + cJSON *tag_val = cJSON_GetArrayItem(tag, 1); + if (!tag_val || !cJSON_IsString(tag_val)) continue; + + int kind = atoi(tag_val->valuestring); + + char **results = NULL; + uint16_t count = 0; + storage_query_events_json(storage, kind, deleter_pubkey, 500, &results, &count); + for (uint16_t e = 0; e < count; e++) { + uint8_t eid[32]; + if (extract_event_id_field(results[e], strlen(results[e]), eid) == 0) { + storage_delete_event(storage, eid); + deleted++; + } + } + storage_free_query_results(results, count); + } + + cJSON_Delete(obj); + return deleted; +} + +int deletion_process_json(storage_engine_t *storage, const char *event_json, + size_t event_len) +{ + if (!storage || !event_json) return 0; + + cJSON *obj = cJSON_ParseWithLength(event_json, event_len); + if (!obj) return 0; + cJSON *kind_item = cJSON_GetObjectItem(obj, "kind"); + int kind = kind_item ? kind_item->valueint : 0; + cJSON *pk_item = cJSON_GetObjectItem(obj, "pubkey"); + const char *pubkey = pk_item ? pk_item->valuestring : ""; + cJSON *ca_item = cJSON_GetObjectItem(obj, "created_at"); + uint64_t created_at = ca_item ? (uint64_t)ca_item->valuedouble : 0; + cJSON_Delete(obj); + + if (kind != NOSTR_KIND_DELETION) return 0; + + char *deleter_pk = strdup(pubkey); + if (!deleter_pk) return 0; + + int deleted = 0; + deleted += delete_by_e_tags(storage, event_json, event_len, deleter_pk); + deleted += delete_by_a_tags(storage, event_json, event_len, deleter_pk, created_at); + deleted += delete_by_k_tags(storage, event_json, event_len, deleter_pk, created_at); + + free(deleter_pk); + ESP_LOGI(TAG, "Deletion processed: %d events removed", deleted); + return deleted; +} diff --git a/components/wisp_relay/deletion.h b/components/wisp_relay/deletion.h new file mode 100644 index 0000000..b494a8e --- /dev/null +++ b/components/wisp_relay/deletion.h @@ -0,0 +1,11 @@ +#ifndef DELETION_H +#define DELETION_H + +#include "storage_engine.h" + +#define NOSTR_KIND_DELETION 5 + +int deletion_process_json(storage_engine_t *storage, const char *event_json, + size_t event_len); + +#endif diff --git a/components/wisp_relay/flash_monitor.c b/components/wisp_relay/flash_monitor.c new file mode 100644 index 0000000..ceb8c3b --- /dev/null +++ b/components/wisp_relay/flash_monitor.c @@ -0,0 +1,30 @@ +#include "flash_monitor.h" +#include "esp_littlefs.h" +#include "esp_log.h" +#include + +static const char *TAG = "flash_monitor"; + +void flash_get_health(const char *partition_label, flash_health_t *health) +{ + memset(health, 0, sizeof(flash_health_t)); + + esp_err_t ret = esp_littlefs_info(partition_label, + &health->total_bytes, + &health->used_bytes); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to get LittleFS info: %s", esp_err_to_name(ret)); + return; + } + + if (health->total_bytes == 0) { + health->free_bytes = 0; + health->usage_percent = 0.0f; + } else { + health->free_bytes = health->total_bytes - health->used_bytes; + health->usage_percent = (float)health->used_bytes / health->total_bytes * 100.0f; + } + + ESP_LOGD(TAG, "Flash: %.1f%% used (%zu/%zu bytes)", + health->usage_percent, health->used_bytes, health->total_bytes); +} diff --git a/components/wisp_relay/flash_monitor.h b/components/wisp_relay/flash_monitor.h new file mode 100644 index 0000000..86f1b53 --- /dev/null +++ b/components/wisp_relay/flash_monitor.h @@ -0,0 +1,16 @@ +#ifndef FLASH_MONITOR_H +#define FLASH_MONITOR_H + +#include +#include + +typedef struct { + size_t total_bytes; + size_t used_bytes; + size_t free_bytes; + float usage_percent; +} flash_health_t; + +void flash_get_health(const char *partition_label, flash_health_t *health); + +#endif diff --git a/components/wisp_relay/handlers.c b/components/wisp_relay/handlers.c new file mode 100644 index 0000000..2164725 --- /dev/null +++ b/components/wisp_relay/handlers.c @@ -0,0 +1,203 @@ +#include "handlers.h" +#include "router.h" +#include "storage_engine.h" +#include "sub_manager.h" +#include "relay_validator.h" +#include "broadcaster.h" +#include "deletion.h" +#include "rate_limiter.h" +#include "relay_types.h" +#include "cJSON.h" +#include "esp_log.h" +#include + +static const char *TAG = "handlers"; + +int handle_event(relay_ctx_t *ctx, int conn_fd, const char *event_json, size_t event_len) +{ + if (!ctx || !event_json) return -1; + + if (ctx->rate_limiter) { + if (!rate_limiter_check(ctx->rate_limiter, conn_fd, RATE_TYPE_EVENT)) { + router_send_ok(ctx, conn_fd, "", false, "rate limited"); + return -1; + } + } + + cJSON *obj = cJSON_ParseWithLength(event_json, event_len); + if (!obj) { + router_send_ok(ctx, conn_fd, "", false, "invalid JSON"); + return -1; + } + + cJSON *id_item = cJSON_GetObjectItem(obj, "id"); + cJSON *pubkey_item = cJSON_GetObjectItem(obj, "pubkey"); + cJSON *kind_item = cJSON_GetObjectItem(obj, "kind"); + cJSON *ca_item = cJSON_GetObjectItem(obj, "created_at"); + + if (!id_item || !pubkey_item || !kind_item || !ca_item) { + cJSON_Delete(obj); + router_send_ok(ctx, conn_fd, "", false, "missing required fields"); + return -1; + } + + const char *id_hex = id_item->valuestring; + const char *pubkey_hex = pubkey_item->valuestring; + int kind = kind_item->valueint; + uint64_t created_at = (uint64_t)ca_item->valuedouble; + + if (ctx->config.max_future_sec > 0) { + int64_t now = (int64_t)(xTaskGetTickCount() / configTICK_RATE_HZ); + if ((int64_t)created_at > now + ctx->config.max_future_sec) { + cJSON_Delete(obj); + router_send_ok(ctx, conn_fd, id_hex, false, "created_at too far in future"); + return -1; + } + } + + uint8_t event_id[32]; + if (relay_hex_to_bytes(id_hex, 64, event_id, 32) != 0) { + cJSON_Delete(obj); + router_send_ok(ctx, conn_fd, "", false, "invalid event id"); + return -1; + } + + if (storage_event_exists(ctx->storage, event_id)) { + cJSON_Delete(obj); + router_send_ok(ctx, conn_fd, id_hex, true, "duplicate"); + return 0; + } + + if (!relay_validator_verify_event(event_json, event_len)) { + cJSON_Delete(obj); + router_send_ok(ctx, conn_fd, id_hex, false, "invalid signature"); + return -1; + } + + cJSON_Delete(obj); + + storage_error_t err = storage_save_event_json(ctx->storage, event_json, event_len); + if (err != STORAGE_OK) { + const char *msg = (err == STORAGE_ERR_FULL) ? "relay full" : + (err == STORAGE_ERR_DUPLICATE) ? "duplicate" : "storage error"; + router_send_ok(ctx, conn_fd, id_hex, false, msg); + return -1; + } + + router_send_ok(ctx, conn_fd, id_hex, true, ""); + + if (kind == NOSTR_KIND_DELETION) { + deletion_process_json(ctx->storage, event_json, event_len); + } + + broadcaster_fanout_json(ctx, event_json, event_len, kind, pubkey_hex, created_at); + + return 0; +} + +static void parse_filter_json(const char *json, sub_filter_t *filter) +{ + memset(filter, 0, sizeof(sub_filter_t)); + cJSON *obj = cJSON_Parse(json); + if (!obj) return; + + cJSON *arr; + + arr = cJSON_GetObjectItem(obj, "ids"); + if (arr && cJSON_IsArray(arr)) { + filter->ids_count = cJSON_GetArraySize(arr); + if (filter->ids_count > SUB_MAX_FILTER_IDS) filter->ids_count = SUB_MAX_FILTER_IDS; + for (size_t i = 0; i < filter->ids_count; i++) + filter->ids[i] = strdup(cJSON_GetArrayItem(arr, i)->valuestring); + } + + arr = cJSON_GetObjectItem(obj, "authors"); + if (arr && cJSON_IsArray(arr)) { + filter->authors_count = cJSON_GetArraySize(arr); + if (filter->authors_count > SUB_MAX_FILTER_AUTHORS) filter->authors_count = SUB_MAX_FILTER_AUTHORS; + for (size_t i = 0; i < filter->authors_count; i++) + filter->authors[i] = strdup(cJSON_GetArrayItem(arr, i)->valuestring); + } + + arr = cJSON_GetObjectItem(obj, "kinds"); + if (arr && cJSON_IsArray(arr)) { + filter->kinds_count = cJSON_GetArraySize(arr); + if (filter->kinds_count > SUB_MAX_FILTER_KINDS) filter->kinds_count = SUB_MAX_FILTER_KINDS; + for (size_t i = 0; i < filter->kinds_count; i++) + filter->kinds[i] = cJSON_GetArrayItem(arr, i)->valueint; + } + + arr = cJSON_GetObjectItem(obj, "#e"); + if (arr && cJSON_IsArray(arr)) { + filter->e_tags_count = cJSON_GetArraySize(arr); + if (filter->e_tags_count > SUB_MAX_FILTER_ETAGS) filter->e_tags_count = SUB_MAX_FILTER_ETAGS; + for (size_t i = 0; i < filter->e_tags_count; i++) + filter->e_tags[i] = strdup(cJSON_GetArrayItem(arr, i)->valuestring); + } + + arr = cJSON_GetObjectItem(obj, "#p"); + if (arr && cJSON_IsArray(arr)) { + filter->p_tags_count = cJSON_GetArraySize(arr); + if (filter->p_tags_count > SUB_MAX_FILTER_PTAGS) filter->p_tags_count = SUB_MAX_FILTER_PTAGS; + for (size_t i = 0; i < filter->p_tags_count; i++) + filter->p_tags[i] = strdup(cJSON_GetArrayItem(arr, i)->valuestring); + } + + cJSON *since = cJSON_GetObjectItem(obj, "since"); + if (since) filter->since = (int64_t)since->valuedouble; + cJSON *until = cJSON_GetObjectItem(obj, "until"); + if (until) filter->until = (int64_t)until->valuedouble; + cJSON *limit = cJSON_GetObjectItem(obj, "limit"); + if (limit) filter->limit = limit->valueint; + + cJSON_Delete(obj); +} + +void handle_req(relay_ctx_t *ctx, int conn_fd, const char *sub_id, const char *filters_json) +{ + if (!ctx || !sub_id) return; + + if (ctx->rate_limiter) { + if (!rate_limiter_check(ctx->rate_limiter, conn_fd, RATE_TYPE_REQ)) { + router_send_closed(ctx, conn_fd, sub_id, "rate limited"); + return; + } + } + + sub_filter_t filter; + parse_filter_json(filters_json, &filter); + + int query_kind = -1; + const char *query_author = NULL; + int query_limit = filter.limit > 0 ? filter.limit : 100; + + if (filter.kinds_count > 0) query_kind = filter.kinds[0]; + if (filter.authors_count > 0) query_author = filter.authors[0]; + + char **results = NULL; + uint16_t count = 0; + storage_query_events_json(ctx->storage, query_kind, query_author, + query_limit, &results, &count); + + for (uint16_t i = 0; i < count; i++) { + router_send_event(ctx, conn_fd, sub_id, results[i], strlen(results[i])); + } + storage_free_query_results(results, count); + + router_send_eose(ctx, conn_fd, sub_id); + + sub_manager_add(ctx->sub_manager, conn_fd, sub_id, &filter, 1); + + sub_filter_t *f = &filter; + for (size_t i = 0; i < f->ids_count; i++) free(f->ids[i]); + for (size_t i = 0; i < f->authors_count; i++) free(f->authors[i]); + for (size_t i = 0; i < f->e_tags_count; i++) free(f->e_tags[i]); + for (size_t i = 0; i < f->p_tags_count; i++) free(f->p_tags[i]); +} + +int handle_close(relay_ctx_t *ctx, int conn_fd, const char *sub_id) +{ + if (!ctx || !sub_id) return -1; + sub_manager_remove(ctx->sub_manager, conn_fd, sub_id); + return 0; +} diff --git a/components/wisp_relay/handlers.h b/components/wisp_relay/handlers.h new file mode 100644 index 0000000..91621bf --- /dev/null +++ b/components/wisp_relay/handlers.h @@ -0,0 +1,10 @@ +#ifndef HANDLERS_H +#define HANDLERS_H + +#include "relay_core.h" + +int handle_event(relay_ctx_t *ctx, int conn_fd, const char *event_json, size_t event_len); +void handle_req(relay_ctx_t *ctx, int conn_fd, const char *sub_id, const char *filters_json); +int handle_close(relay_ctx_t *ctx, int conn_fd, const char *sub_id); + +#endif diff --git a/components/wisp_relay/idf_component.yml b/components/wisp_relay/idf_component.yml new file mode 100644 index 0000000..c093387 --- /dev/null +++ b/components/wisp_relay/idf_component.yml @@ -0,0 +1 @@ +dependencies: {} diff --git a/components/wisp_relay/nip11_relay.c b/components/wisp_relay/nip11_relay.c new file mode 100644 index 0000000..4e1df37 --- /dev/null +++ b/components/wisp_relay/nip11_relay.c @@ -0,0 +1,53 @@ +#include "nip11_relay.h" +#include + +static const char *NIP11_JSON = +"{" + "\"name\":\"TollGate Relay\"," + "\"description\":\"Local Nostr relay with 21-day TTL and negentropy sync\"," + "\"pubkey\":\"\"," + "\"contact\":\"\"," + "\"supported_nips\":[1,9,11,20,40,77]," + "\"software\":\"https://github.com/nicobao/esp32-tollgate\"," + "\"version\":\"1.0.0\"," + "\"limitation\":{" + "\"max_message_length\":65536," + "\"max_subscriptions\":8," + "\"max_filters\":4," + "\"max_limit\":500," + "\"max_subid_length\":64," + "\"max_event_tags\":100," + "\"max_content_length\":32768," + "\"min_pow_difficulty\":0," + "\"auth_required\":false," + "\"payment_required\":false" + "}," + "\"retention\":[{\"kinds\":[0,1,2,3,4,5,6,7],\"time\":1814400}]," + "\"relay_countries\":[\"DE\"]" +"}"; + +esp_err_t relay_nip11_handler(httpd_req_t *req) +{ + char accept[64] = ""; + httpd_req_get_hdr_value_str(req, "Accept", accept, sizeof(accept)); + + if (strstr(accept, "application/nostr+json")) { + httpd_resp_set_type(req, "application/nostr+json"); + } else { + httpd_resp_set_type(req, "application/json"); + } + + httpd_resp_set_hdr(req, "Access-Control-Allow-Origin", "*"); + httpd_resp_set_hdr(req, "Access-Control-Allow-Headers", "Content-Type, Accept"); + httpd_resp_set_hdr(req, "Access-Control-Allow-Methods", "GET, OPTIONS"); + return httpd_resp_send(req, NIP11_JSON, strlen(NIP11_JSON)); +} + +esp_err_t relay_nip11_options_handler(httpd_req_t *req) +{ + httpd_resp_set_hdr(req, "Access-Control-Allow-Origin", "*"); + httpd_resp_set_hdr(req, "Access-Control-Allow-Headers", "Content-Type, Accept"); + httpd_resp_set_hdr(req, "Access-Control-Allow-Methods", "GET, OPTIONS"); + httpd_resp_set_status(req, "204 No Content"); + return httpd_resp_send(req, NULL, 0); +} diff --git a/components/wisp_relay/nip11_relay.h b/components/wisp_relay/nip11_relay.h new file mode 100644 index 0000000..84f7971 --- /dev/null +++ b/components/wisp_relay/nip11_relay.h @@ -0,0 +1,9 @@ +#ifndef NIP11_RELAY_H +#define NIP11_RELAY_H + +#include "esp_http_server.h" + +esp_err_t relay_nip11_handler(httpd_req_t *req); +esp_err_t relay_nip11_options_handler(httpd_req_t *req); + +#endif diff --git a/components/wisp_relay/rate_limiter.c b/components/wisp_relay/rate_limiter.c new file mode 100644 index 0000000..7734e03 --- /dev/null +++ b/components/wisp_relay/rate_limiter.c @@ -0,0 +1,98 @@ +#include "rate_limiter.h" +#include "esp_timer.h" +#include "esp_log.h" +#include + +static const char *TAG = "rate_limiter"; + +void rate_limiter_init(rate_limiter_t *rl, const rate_config_t *config) +{ + memset(rl, 0, sizeof(rate_limiter_t)); + rl->lock = xSemaphoreCreateMutex(); + if (config) { + memcpy(&rl->config, config, sizeof(rate_config_t)); + } else { + rl->config.events_per_minute = 30; + rl->config.reqs_per_minute = 60; + } +} + +void rate_limiter_destroy(rate_limiter_t *rl) +{ + if (!rl) return; + if (rl->lock) { + vSemaphoreDelete(rl->lock); + rl->lock = NULL; + } +} + +static rate_bucket_t* get_bucket(rate_limiter_t *rl, int fd) +{ + for (int i = 0; i < RATE_LIMITER_MAX_BUCKETS; i++) { + if (rl->buckets[i].active && rl->buckets[i].fd == fd) { + return &rl->buckets[i]; + } + } + for (int i = 0; i < RATE_LIMITER_MAX_BUCKETS; i++) { + if (!rl->buckets[i].active) { + rl->buckets[i].fd = fd; + rl->buckets[i].active = true; + rl->buckets[i].event_count = 0; + rl->buckets[i].req_count = 0; + rl->buckets[i].window_start = esp_timer_get_time() / 1000000; + return &rl->buckets[i]; + } + } + return NULL; +} + +bool rate_limiter_check(rate_limiter_t *rl, int fd, rate_type_t type) +{ + xSemaphoreTake(rl->lock, portMAX_DELAY); + + rate_bucket_t *bucket = get_bucket(rl, fd); + if (!bucket) { + xSemaphoreGive(rl->lock); + return false; + } + + uint32_t now = esp_timer_get_time() / 1000000; + + if (now - bucket->window_start >= 60) { + bucket->event_count = 0; + bucket->req_count = 0; + bucket->window_start = now; + } + + bool allowed = true; + if (type == RATE_TYPE_EVENT) { + if (bucket->event_count >= rl->config.events_per_minute) { + ESP_LOGW(TAG, "Rate limited: fd=%d events=%d", fd, bucket->event_count); + allowed = false; + } else { + bucket->event_count++; + } + } else { + if (bucket->req_count >= rl->config.reqs_per_minute) { + ESP_LOGW(TAG, "Rate limited: fd=%d reqs=%d", fd, bucket->req_count); + allowed = false; + } else { + bucket->req_count++; + } + } + + xSemaphoreGive(rl->lock); + return allowed; +} + +void rate_limiter_reset(rate_limiter_t *rl, int fd) +{ + xSemaphoreTake(rl->lock, portMAX_DELAY); + for (int i = 0; i < RATE_LIMITER_MAX_BUCKETS; i++) { + if (rl->buckets[i].active && rl->buckets[i].fd == fd) { + rl->buckets[i].active = false; + break; + } + } + xSemaphoreGive(rl->lock); +} diff --git a/components/wisp_relay/rate_limiter.h b/components/wisp_relay/rate_limiter.h new file mode 100644 index 0000000..655ddf2 --- /dev/null +++ b/components/wisp_relay/rate_limiter.h @@ -0,0 +1,40 @@ +#ifndef RATE_LIMITER_H +#define RATE_LIMITER_H + +#include +#include +#include "freertos/FreeRTOS.h" +#include "freertos/semphr.h" + +#define RATE_LIMITER_MAX_BUCKETS 16 + +typedef enum { + RATE_TYPE_EVENT, + RATE_TYPE_REQ, +} rate_type_t; + +typedef struct { + uint16_t events_per_minute; + uint16_t reqs_per_minute; +} rate_config_t; + +typedef struct { + int fd; + uint16_t event_count; + uint16_t req_count; + uint32_t window_start; + bool active; +} rate_bucket_t; + +typedef struct rate_limiter { + rate_config_t config; + rate_bucket_t buckets[RATE_LIMITER_MAX_BUCKETS]; + SemaphoreHandle_t lock; +} rate_limiter_t; + +void rate_limiter_init(rate_limiter_t *rl, const rate_config_t *config); +void rate_limiter_destroy(rate_limiter_t *rl); +bool rate_limiter_check(rate_limiter_t *rl, int fd, rate_type_t type); +void rate_limiter_reset(rate_limiter_t *rl, int fd); + +#endif diff --git a/components/wisp_relay/relay_core.h b/components/wisp_relay/relay_core.h new file mode 100644 index 0000000..d8e7096 --- /dev/null +++ b/components/wisp_relay/relay_core.h @@ -0,0 +1,27 @@ +#ifndef RELAY_CORE_H +#define RELAY_CORE_H + +#include + +#include "ws_server.h" + +typedef struct sub_manager sub_manager_t; +typedef struct storage_engine storage_engine_t; +typedef struct rate_limiter rate_limiter_t; + +typedef struct relay_ctx { + ws_server_t ws_server; + sub_manager_t *sub_manager; + storage_engine_t *storage; + rate_limiter_t *rate_limiter; + + struct { + uint16_t port; + uint32_t max_event_age_sec; + uint8_t max_subs_per_conn; + uint8_t max_filters_per_sub; + int64_t max_future_sec; + } config; +} relay_ctx_t; + +#endif diff --git a/components/wisp_relay/relay_types.c b/components/wisp_relay/relay_types.c new file mode 100644 index 0000000..9833885 --- /dev/null +++ b/components/wisp_relay/relay_types.c @@ -0,0 +1,21 @@ +#include "relay_types.h" +#include +#include + +int relay_hex_to_bytes(const char *hex, size_t hex_len, uint8_t *out, size_t out_len) +{ + if (hex_len != out_len * 2) return -1; + for (size_t i = 0; i < out_len; i++) { + unsigned int byte; + if (sscanf(hex + i * 2, "%02x", &byte) != 1) return -1; + out[i] = (uint8_t)byte; + } + return 0; +} + +void relay_bytes_to_hex(const uint8_t *bytes, size_t len, char *hex) +{ + for (size_t i = 0; i < len; i++) + sprintf(hex + i * 2, "%02x", bytes[i]); + hex[len * 2] = '\0'; +} diff --git a/components/wisp_relay/relay_types.h b/components/wisp_relay/relay_types.h new file mode 100644 index 0000000..343e51b --- /dev/null +++ b/components/wisp_relay/relay_types.h @@ -0,0 +1,43 @@ +#ifndef RELAY_TYPES_H +#define RELAY_TYPES_H + +#include +#include +#include + +#define RELAY_MAX_EVENT_SIZE 8192 +#define RELAY_ID_SIZE 32 +#define RELAY_SIG_SIZE 64 +#define RELAY_MAX_TAGS 100 +#define RELAY_MAX_TAG_VALUES 10 + +typedef struct relay_event { + uint8_t id[RELAY_ID_SIZE]; + uint8_t pubkey[RELAY_ID_SIZE]; + uint64_t created_at; + int kind; + uint8_t sig[RELAY_SIG_SIZE]; + char content[RELAY_MAX_EVENT_SIZE]; + size_t content_len; +} relay_event_t; + +typedef struct { + char **ids; + size_t ids_count; + char **authors; + size_t authors_count; + int32_t *kinds; + size_t kinds_count; + char **e_tags; + size_t e_tags_count; + char **p_tags; + size_t p_tags_count; + int64_t since; + int64_t until; + int limit; +} relay_filter_t; + +int relay_hex_to_bytes(const char *hex, size_t hex_len, uint8_t *out, size_t out_len); +void relay_bytes_to_hex(const uint8_t *bytes, size_t len, char *hex); + +#endif diff --git a/components/wisp_relay/relay_validator.c b/components/wisp_relay/relay_validator.c new file mode 100644 index 0000000..eb40d22 --- /dev/null +++ b/components/wisp_relay/relay_validator.c @@ -0,0 +1,176 @@ +#include "relay_validator.h" +#include "relay_types.h" +#include "esp_log.h" +#include "mbedtls/sha256.h" +#include "secp256k1.h" +#include "secp256k1_extrakeys.h" +#include "secp256k1_schnorrsig.h" +#include "cJSON.h" +#include "freertos/FreeRTOS.h" +#include "freertos/task.h" +#include +#include +#include +#include + +static const char *TAG = "relay_validator"; + +static int hex_to_bytes(const char *hex, size_t hex_len, uint8_t *out, size_t out_len) +{ + if (hex_len != out_len * 2) return -1; + for (size_t i = 0; i < out_len; i++) { + unsigned int byte; + if (sscanf(hex + i * 2, "%02x", &byte) != 1) return -1; + out[i] = (uint8_t)byte; + } + return 0; +} + +static char *serialize_event_for_id(const char *event_json, size_t event_len) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, event_len); + if (!obj) return NULL; + + cJSON *serial = cJSON_CreateArray(); + cJSON_AddItemToArray(serial, cJSON_CreateNumber(0)); + cJSON_AddItemToArray(serial, cJSON_CreateString( + cJSON_GetObjectItem(obj, "pubkey")->valuestring)); + cJSON_AddItemToArray(serial, cJSON_CreateNumber( + cJSON_GetObjectItem(obj, "created_at")->valuedouble)); + cJSON_AddItemToArray(serial, cJSON_CreateNumber( + cJSON_GetObjectItem(obj, "kind")->valueint)); + cJSON *tags = cJSON_GetObjectItem(obj, "tags"); + cJSON_AddItemToArray(serial, cJSON_Duplicate(tags, 1)); + cJSON_AddItemToArray(serial, cJSON_CreateString( + cJSON_GetObjectItem(obj, "content")->valuestring)); + + char *result = cJSON_PrintUnformatted(serial); + cJSON_Delete(serial); + cJSON_Delete(obj); + return result; +} + +static bool verify_event_id(const char *event_json, size_t event_len, + const uint8_t expected_id[32]) +{ + char *serialized = serialize_event_for_id(event_json, event_len); + if (!serialized) return false; + + uint8_t hash[32]; + mbedtls_sha256((const unsigned char *)serialized, strlen(serialized), hash, 0); + free(serialized); + + return memcmp(hash, expected_id, 32) == 0; +} + +static bool verify_schnorr_sig(const uint8_t pubkey[32], const uint8_t msg[32], + const uint8_t sig[64]) +{ + secp256k1_context *ctx = secp256k1_context_create(SECP256K1_CONTEXT_VERIFY); + if (!ctx) return false; + + secp256k1_xonly_pubkey xonly_pub; + if (!secp256k1_xonly_pubkey_parse(ctx, &xonly_pub, pubkey)) { + secp256k1_context_destroy(ctx); + return false; + } + + bool valid = secp256k1_schnorrsig_verify(ctx, sig, msg, 32, &xonly_pub); + secp256k1_context_destroy(ctx); + return valid; +} + +bool relay_validator_verify_event(const char *event_json, size_t event_len) +{ + cJSON *obj = cJSON_ParseWithLength(event_json, event_len); + if (!obj) { + ESP_LOGD(TAG, "Invalid JSON"); + return false; + } + + cJSON *id_item = cJSON_GetObjectItem(obj, "id"); + cJSON *pk_item = cJSON_GetObjectItem(obj, "pubkey"); + cJSON *sig_item = cJSON_GetObjectItem(obj, "sig"); + + if (!id_item || !pk_item || !sig_item) { + cJSON_Delete(obj); + ESP_LOGD(TAG, "Missing required fields"); + return false; + } + + const char *id_hex = id_item->valuestring; + const char *pk_hex = pk_item->valuestring; + const char *sig_hex = sig_item->valuestring; + + if (strlen(id_hex) != 64 || strlen(pk_hex) != 64 || strlen(sig_hex) != 128) { + cJSON_Delete(obj); + ESP_LOGD(TAG, "Invalid field lengths"); + return false; + } + + uint8_t event_id[32], pubkey[32], sig[64]; + if (hex_to_bytes(id_hex, 64, event_id, 32) != 0 || + hex_to_bytes(pk_hex, 64, pubkey, 32) != 0 || + hex_to_bytes(sig_hex, 128, sig, 64) != 0) { + cJSON_Delete(obj); + ESP_LOGD(TAG, "Invalid hex encoding"); + return false; + } + + cJSON_Delete(obj); + + if (!verify_event_id(event_json, event_len, event_id)) { + ESP_LOGD(TAG, "Event ID mismatch"); + return false; + } + + if (!verify_schnorr_sig(pubkey, event_id, sig)) { + ESP_LOGD(TAG, "Invalid signature"); + return false; + } + + return true; +} + +validation_result_t relay_validator_check(const uint8_t *id, + const uint8_t *pubkey, + uint64_t created_at, + int kind, + const char *content, + size_t content_len, + const char *tags_json, + const uint8_t *sig, + const validator_config_t *config) +{ + (void)content; (void)content_len; (void)tags_json; + + if (config) { + if (config->max_future_sec > 0) { + int64_t now = (int64_t)(xTaskGetTickCount() / configTICK_RATE_HZ); + if ((int64_t)created_at > now + config->max_future_sec) + return VALIDATION_ERR_FUTURE; + } + } + + if (!verify_schnorr_sig(pubkey, id, sig)) + return VALIDATION_ERR_SIG; + + return VALIDATION_OK; +} + +const char *relay_validator_result_string(validation_result_t result) +{ + switch (result) { + case VALIDATION_OK: return "ok"; + case VALIDATION_ERR_SCHEMA: return "invalid: schema"; + case VALIDATION_ERR_ID: return "invalid: event id"; + case VALIDATION_ERR_SIG: return "invalid: signature"; + case VALIDATION_ERR_EXPIRED: return "invalid: expired"; + case VALIDATION_ERR_FUTURE: return "invalid: future"; + case VALIDATION_ERR_DUPLICATE: return "duplicate"; + case VALIDATION_ERR_POW: return "pow: insufficient"; + case VALIDATION_ERR_BLOCKED: return "blocked"; + case VALIDATION_ERR_TOO_OLD: return "invalid: too old"; + default: return "error: unknown"; + } +} diff --git a/components/wisp_relay/relay_validator.h b/components/wisp_relay/relay_validator.h new file mode 100644 index 0000000..c07308f --- /dev/null +++ b/components/wisp_relay/relay_validator.h @@ -0,0 +1,45 @@ +#ifndef RELAY_VALIDATOR_H +#define RELAY_VALIDATOR_H + +#include +#include +#include + +typedef enum { + VALIDATION_OK = 0, + VALIDATION_ERR_SCHEMA, + VALIDATION_ERR_ID, + VALIDATION_ERR_SIG, + VALIDATION_ERR_EXPIRED, + VALIDATION_ERR_FUTURE, + VALIDATION_ERR_DUPLICATE, + VALIDATION_ERR_POW, + VALIDATION_ERR_BLOCKED, + VALIDATION_ERR_TOO_OLD, +} validation_result_t; + +typedef struct { + uint32_t max_event_age_sec; + int64_t max_future_sec; + uint8_t min_pow_difficulty; + bool check_duplicates; +} validator_config_t; + +typedef struct relay_event relay_event_t; +typedef struct storage_engine storage_engine_t; + +validation_result_t relay_validator_check(const uint8_t *id, + const uint8_t *pubkey, + uint64_t created_at, + int kind, + const char *content, + size_t content_len, + const char *tags_json, + const uint8_t *sig, + const validator_config_t *config); + +bool relay_validator_verify_event(const char *event_json, size_t event_len); + +const char *relay_validator_result_string(validation_result_t result); + +#endif diff --git a/components/wisp_relay/router.c b/components/wisp_relay/router.c new file mode 100644 index 0000000..05aa7d4 --- /dev/null +++ b/components/wisp_relay/router.c @@ -0,0 +1,140 @@ +#include "router.h" +#include "ws_server.h" +#include "handlers.h" +#include "sub_manager.h" +#include "cJSON.h" +#include "esp_log.h" +#include + +static const char *TAG = "router"; + +esp_err_t router_send_notice(relay_ctx_t *ctx, int conn_fd, const char *message) +{ + cJSON *arr = cJSON_CreateArray(); + cJSON_AddItemToArray(arr, cJSON_CreateString("NOTICE")); + cJSON_AddItemToArray(arr, cJSON_CreateString(message)); + char *json = cJSON_PrintUnformatted(arr); + cJSON_Delete(arr); + esp_err_t ret = ws_server_send(&ctx->ws_server, conn_fd, json, strlen(json)); + cJSON_free(json); + return ret; +} + +esp_err_t router_send_ok(relay_ctx_t *ctx, int conn_fd, const char *event_id_hex, + bool accepted, const char *message) +{ + cJSON *arr = cJSON_CreateArray(); + cJSON_AddItemToArray(arr, cJSON_CreateString("OK")); + cJSON_AddItemToArray(arr, cJSON_CreateString(event_id_hex)); + cJSON_AddItemToArray(arr, cJSON_CreateBool(accepted)); + cJSON_AddItemToArray(arr, cJSON_CreateString(message ? message : "")); + char *json = cJSON_PrintUnformatted(arr); + cJSON_Delete(arr); + esp_err_t ret = ws_server_send(&ctx->ws_server, conn_fd, json, strlen(json)); + cJSON_free(json); + return ret; +} + +esp_err_t router_send_eose(relay_ctx_t *ctx, int conn_fd, const char *sub_id) +{ + cJSON *arr = cJSON_CreateArray(); + cJSON_AddItemToArray(arr, cJSON_CreateString("EOSE")); + cJSON_AddItemToArray(arr, cJSON_CreateString(sub_id)); + char *json = cJSON_PrintUnformatted(arr); + cJSON_Delete(arr); + esp_err_t ret = ws_server_send(&ctx->ws_server, conn_fd, json, strlen(json)); + cJSON_free(json); + return ret; +} + +esp_err_t router_send_closed(relay_ctx_t *ctx, int conn_fd, const char *sub_id, + const char *message) +{ + cJSON *arr = cJSON_CreateArray(); + cJSON_AddItemToArray(arr, cJSON_CreateString("CLOSED")); + cJSON_AddItemToArray(arr, cJSON_CreateString(sub_id)); + cJSON_AddItemToArray(arr, cJSON_CreateString(message ? message : "")); + char *json = cJSON_PrintUnformatted(arr); + cJSON_Delete(arr); + esp_err_t ret = ws_server_send(&ctx->ws_server, conn_fd, json, strlen(json)); + cJSON_free(json); + return ret; +} + +esp_err_t router_send_event(relay_ctx_t *ctx, int conn_fd, const char *sub_id, + const char *event_json, size_t event_len) +{ + size_t buf_size = event_len + strlen(sub_id) + 32; + char *buf = malloc(buf_size); + if (!buf) return ESP_ERR_NO_MEM; + int n = snprintf(buf, buf_size, "[\"EVENT\",\"%s\",%.*s]", sub_id, (int)event_len, event_json); + esp_err_t ret = ws_server_send(&ctx->ws_server, conn_fd, buf, n); + free(buf); + return ret; +} + +static void on_ws_message(int fd, const char *data, size_t len) +{ + extern relay_ctx_t g_relay_ctx; + router_dispatch(&g_relay_ctx, fd, data, len); +} + +static void on_ws_disconnect(int fd) +{ + extern relay_ctx_t g_relay_ctx; + if (g_relay_ctx.sub_manager) { + sub_manager_remove_all(g_relay_ctx.sub_manager, fd); + } +} + +void router_dispatch(relay_ctx_t *ctx, int conn_fd, const char *data, size_t len) +{ + cJSON *arr = cJSON_ParseWithLength(data, len); + if (!arr || !cJSON_IsArray(arr)) { + router_send_notice(ctx, conn_fd, "invalid JSON"); + if (arr) cJSON_Delete(arr); + return; + } + + int array_size = cJSON_GetArraySize(arr); + if (array_size < 2) { + router_send_notice(ctx, conn_fd, "array too short"); + cJSON_Delete(arr); + return; + } + + cJSON *cmd = cJSON_GetArrayItem(arr, 0); + if (!cmd || !cJSON_IsString(cmd)) { + router_send_notice(ctx, conn_fd, "invalid command"); + cJSON_Delete(arr); + return; + } + + const char *cmd_str = cmd->valuestring; + + if (strcmp(cmd_str, "EVENT") == 0 && array_size >= 2) { + cJSON *event_obj = cJSON_GetArrayItem(arr, 1); + if (event_obj) { + char *event_json = cJSON_PrintUnformatted(event_obj); + handle_event(ctx, conn_fd, event_json, strlen(event_json)); + cJSON_free(event_json); + } + } else if (strcmp(cmd_str, "REQ") == 0 && array_size >= 3) { + cJSON *sub_id_item = cJSON_GetArrayItem(arr, 1); + if (sub_id_item && cJSON_IsString(sub_id_item)) { + cJSON *filter_obj = cJSON_GetArrayItem(arr, 2); + char *filter_json = filter_obj ? cJSON_PrintUnformatted(filter_obj) : strdup("{}"); + handle_req(ctx, conn_fd, sub_id_item->valuestring, filter_json); + free(filter_json); + } + } else if (strcmp(cmd_str, "CLOSE") == 0 && array_size >= 2) { + cJSON *sub_id_item = cJSON_GetArrayItem(arr, 1); + if (sub_id_item && cJSON_IsString(sub_id_item)) { + handle_close(ctx, conn_fd, sub_id_item->valuestring); + } + } else { + router_send_notice(ctx, conn_fd, "unknown command"); + } + + cJSON_Delete(arr); +} diff --git a/components/wisp_relay/router.h b/components/wisp_relay/router.h new file mode 100644 index 0000000..9afd46e --- /dev/null +++ b/components/wisp_relay/router.h @@ -0,0 +1,19 @@ +#ifndef ROUTER_H +#define ROUTER_H + +#include "relay_core.h" +#include +#include + +esp_err_t router_send_notice(relay_ctx_t *ctx, int conn_fd, const char *message); +esp_err_t router_send_ok(relay_ctx_t *ctx, int conn_fd, const char *event_id_hex, + bool accepted, const char *message); +esp_err_t router_send_eose(relay_ctx_t *ctx, int conn_fd, const char *sub_id); +esp_err_t router_send_closed(relay_ctx_t *ctx, int conn_fd, const char *sub_id, + const char *message); +esp_err_t router_send_event(relay_ctx_t *ctx, int conn_fd, const char *sub_id, + const char *event_json, size_t event_len); + +void router_dispatch(relay_ctx_t *ctx, int conn_fd, const char *data, size_t len); + +#endif diff --git a/components/wisp_relay/storage_engine.c b/components/wisp_relay/storage_engine.c new file mode 100644 index 0000000..d26705b --- /dev/null +++ b/components/wisp_relay/storage_engine.c @@ -0,0 +1,402 @@ +#include "storage_engine.h" +#include "esp_littlefs.h" +#include "esp_log.h" +#include "nvs_flash.h" +#include "nvs.h" +#include +#include +#include +#include +#include +#include + +static const char *TAG = "storage"; + +#define INDEX_NVS_NAMESPACE "nostr_idx" +#define EVENTS_DIR "/littlefs/events" + +static void get_event_path(const uint8_t event_id[32], uint32_t file_index, + char *path, size_t len) +{ + char id_hex[33]; + for (int i = 0; i < 16; i++) sprintf(id_hex + i * 2, "%02x", event_id[i]); + snprintf(path, len, EVENTS_DIR "/%02x/%s_%08" PRIx32 ".json", + event_id[0], id_hex, file_index); +} + +static int save_index_to_nvs(storage_engine_t *engine) +{ + nvs_handle_t nvs; + esp_err_t err = nvs_open(INDEX_NVS_NAMESPACE, NVS_READWRITE, &nvs); + if (err != ESP_OK) return STORAGE_ERR_IO; + + nvs_set_u16(nvs, "count", engine->index_count); + nvs_set_u32(nvs, "next_idx", engine->next_file_index); + + const uint16_t chunk_size = 50; + for (uint16_t i = 0; i < engine->index_count; i += chunk_size) { + char key[16]; + snprintf(key, sizeof(key), "idx_%u", i / chunk_size); + uint16_t entries = engine->index_count - i; + if (entries > chunk_size) entries = chunk_size; + nvs_set_blob(nvs, key, &engine->index[i], entries * sizeof(storage_index_entry_t)); + } + nvs_commit(nvs); + nvs_close(nvs); + return STORAGE_OK; +} + +static int load_index_from_nvs(storage_engine_t *engine) +{ + nvs_handle_t nvs; + esp_err_t err = nvs_open(INDEX_NVS_NAMESPACE, NVS_READONLY, &nvs); + if (err == ESP_ERR_NVS_NOT_FOUND) return STORAGE_OK; + if (err != ESP_OK) return STORAGE_ERR_IO; + + err = nvs_get_u16(nvs, "count", &engine->index_count); + if (err != ESP_OK) { nvs_close(nvs); return STORAGE_ERR_IO; } + if (engine->index_count > engine->max_index_entries) engine->index_count = engine->max_index_entries; + + err = nvs_get_u32(nvs, "next_idx", &engine->next_file_index); + if (err != ESP_OK) { nvs_close(nvs); return STORAGE_ERR_IO; } + + const uint16_t chunk_size = 50; + for (uint16_t i = 0; i < engine->index_count; i += chunk_size) { + char key[16]; + snprintf(key, sizeof(key), "idx_%u", i / chunk_size); + uint16_t entries = engine->index_count - i; + if (entries > chunk_size) entries = chunk_size; + size_t len = entries * sizeof(storage_index_entry_t); + nvs_get_blob(nvs, key, &engine->index[i], &len); + } + nvs_close(nvs); + return STORAGE_OK; +} + +static storage_index_entry_t *find_index_entry(storage_engine_t *engine, + const uint8_t event_id[32]) +{ + for (uint16_t i = 0; i < engine->index_count; i++) { + if (memcmp(engine->index[i].event_id, event_id, 32) == 0 && + !(engine->index[i].flags & STORAGE_FLAG_DELETED)) { + return &engine->index[i]; + } + } + return NULL; +} + +static void parse_event_meta(const char *json, size_t len, + uint8_t *id_out, uint8_t *pubkey_out, + uint64_t *created_at_out, int *kind_out) +{ + extern int relay_hex_to_bytes(const char *hex, size_t hex_len, uint8_t *out, size_t out_len); + extern void relay_bytes_to_hex(const uint8_t *bytes, size_t len, char *hex); + + id_out[0] = 0; pubkey_out[0] = 0; *created_at_out = 0; *kind_out = 0; + + const char *p; + p = strstr(json, "\"id\":\""); + if (p) relay_hex_to_bytes(p + 6, 64, id_out, 32); + p = strstr(json, "\"pubkey\":\""); + if (p) relay_hex_to_bytes(p + 10, 64, pubkey_out, 32); + p = strstr(json, "\"created_at\":"); + if (p) *created_at_out = strtoull(p + 13, NULL, 10); + p = strstr(json, "\"kind\":"); + if (p) *kind_out = atoi(p + 7); +} + +esp_err_t storage_init(storage_engine_t *engine, uint32_t default_ttl_sec) +{ + memset(engine, 0, sizeof(storage_engine_t)); + engine->default_ttl_sec = default_ttl_sec; + strcpy(engine->mount_point, "/littlefs"); + + engine->lock = xSemaphoreCreateMutex(); + if (!engine->lock) return ESP_ERR_NO_MEM; + + engine->max_index_entries = STORAGE_INDEX_ENTRIES; + engine->index = heap_caps_calloc(engine->max_index_entries, + sizeof(storage_index_entry_t), + MALLOC_CAP_SPIRAM | MALLOC_CAP_8BIT); + if (!engine->index) { + engine->max_index_entries = 1000; + engine->index = calloc(engine->max_index_entries, sizeof(storage_index_entry_t)); + if (!engine->index) { vSemaphoreDelete(engine->lock); return ESP_ERR_NO_MEM; } + } + + esp_vfs_littlefs_conf_t conf = { + .base_path = "/littlefs", + .partition_label = STORAGE_PARTITION_LABEL, + .format_if_mount_failed = true, + .dont_mount = false, + }; + + esp_err_t ret = esp_vfs_littlefs_register(&conf); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to mount LittleFS: %s", esp_err_to_name(ret)); + free(engine->index); + vSemaphoreDelete(engine->lock); + return ret; + } + + mkdir(EVENTS_DIR, 0755); + for (int i = 0; i < 256; i++) { + char subdir[64]; + snprintf(subdir, sizeof(subdir), EVENTS_DIR "/%02x", i); + mkdir(subdir, 0755); + } + + int load_err = load_index_from_nvs(engine); + if (load_err != STORAGE_OK) { + ESP_LOGW(TAG, "Failed to load index, starting fresh"); + engine->index_count = 0; + engine->next_file_index = 0; + } + + engine->initialized = true; + + size_t total, used; + esp_littlefs_info(STORAGE_PARTITION_LABEL, &total, &used); + ESP_LOGI(TAG, "Storage initialized: %" PRIu16 " events, %zu/%zu bytes used", + engine->index_count, used, total); + return ESP_OK; +} + +void storage_destroy(storage_engine_t *engine) +{ + if (!engine->initialized) return; + if (engine->cleanup_task) { + engine->cleanup_stop = true; + while (engine->cleanup_task != NULL) vTaskDelay(pdMS_TO_TICKS(100)); + } + save_index_to_nvs(engine); + esp_vfs_littlefs_unregister(STORAGE_PARTITION_LABEL); + if (engine->index) { free(engine->index); engine->index = NULL; } + if (engine->lock) { vSemaphoreDelete(engine->lock); engine->lock = NULL; } + engine->initialized = false; +} + +storage_error_t storage_save_event_json(storage_engine_t *engine, + const char *event_json, + size_t event_json_len) +{ + if (!engine->initialized) return STORAGE_ERR_NOT_INITIALIZED; + + uint8_t id[32] = {0}, pubkey[32] = {0}; + uint64_t created_at = 0; + int kind = 0; + parse_event_meta(event_json, event_json_len, id, pubkey, &created_at, &kind); + + xSemaphoreTake(engine->lock, portMAX_DELAY); + + if (find_index_entry(engine, id)) { + xSemaphoreGive(engine->lock); + return STORAGE_ERR_DUPLICATE; + } + if (engine->index_count >= engine->max_index_entries) { + xSemaphoreGive(engine->lock); + return STORAGE_ERR_FULL; + } + + char path[128]; + get_event_path(id, engine->next_file_index, path, sizeof(path)); + FILE *f = fopen(path, "wb"); + if (!f) { + char dir[64]; + snprintf(dir, sizeof(dir), EVENTS_DIR "/%02x", id[0]); + mkdir(dir, 0755); + f = fopen(path, "wb"); + } + if (!f) { xSemaphoreGive(engine->lock); return STORAGE_ERR_IO; } + + fwrite(event_json, 1, event_json_len, f); + fclose(f); + + storage_index_entry_t *entry = &engine->index[engine->index_count]; + memcpy(entry->event_id, id, 32); + entry->created_at = (uint32_t)created_at; + entry->kind = kind; + memcpy(entry->pubkey_prefix, pubkey, 4); + entry->file_index = engine->next_file_index; + entry->flags = 0; + entry->expires_at = (uint32_t)time(NULL) + engine->default_ttl_sec; + + engine->index_count++; + engine->next_file_index++; + if (engine->index_count % 10 == 0) save_index_to_nvs(engine); + + xSemaphoreGive(engine->lock); + return STORAGE_OK; +} + +bool storage_event_exists(storage_engine_t *engine, const uint8_t event_id[32]) +{ + if (!engine->initialized) return false; + xSemaphoreTake(engine->lock, portMAX_DELAY); + bool exists = (find_index_entry(engine, event_id) != NULL); + xSemaphoreGive(engine->lock); + return exists; +} + +storage_error_t storage_query_events_json(storage_engine_t *engine, + int kind, + const char *author_hex, + int limit, + char ***results, + uint16_t *count) +{ + if (!engine->initialized) return STORAGE_ERR_NOT_INITIALIZED; + *results = NULL; + *count = 0; + if (limit > 500) limit = 500; + if (limit <= 0) limit = 100; + + char **out = calloc(limit, sizeof(char *)); + if (!out) return STORAGE_ERR_NO_MEM; + + xSemaphoreTake(engine->lock, portMAX_DELAY); + uint32_t now = (uint32_t)time(NULL); + uint16_t found = 0; + + uint8_t author_prefix[4] = {0}; + int have_author = 0; + if (author_hex && strlen(author_hex) >= 8) { + extern int relay_hex_to_bytes(const char *, size_t, uint8_t *, size_t); + relay_hex_to_bytes(author_hex, 8, author_prefix, 4); + have_author = 1; + } + + for (int i = engine->index_count - 1; i >= 0 && found < limit; i--) { + storage_index_entry_t *e = &engine->index[i]; + if (e->flags & STORAGE_FLAG_DELETED) continue; + if (e->expires_at > 0 && e->expires_at < now) continue; + if (kind > 0 && e->kind != kind) continue; + if (have_author && memcmp(e->pubkey_prefix, author_prefix, 4) != 0) continue; + + char path[128]; + get_event_path(e->event_id, e->file_index, path, sizeof(path)); + FILE *f = fopen(path, "rb"); + if (!f) continue; + fseek(f, 0, SEEK_END); + long sz = ftell(f); + fseek(f, 0, SEEK_SET); + if (sz <= 0 || sz > STORAGE_MAX_EVENT_SIZE) { fclose(f); continue; } + char *buf = malloc(sz + 1); + fread(buf, 1, sz, f); + buf[sz] = '\0'; + fclose(f); + out[found++] = buf; + } + + xSemaphoreGive(engine->lock); + *results = out; + *count = found; + return STORAGE_OK; +} + +void storage_free_query_results(char **results, uint16_t count) +{ + if (!results) return; + for (uint16_t i = 0; i < count; i++) free(results[i]); + free(results); +} + +storage_error_t storage_delete_event(storage_engine_t *engine, const uint8_t event_id[32]) +{ + if (!engine->initialized) return STORAGE_ERR_NOT_INITIALIZED; + xSemaphoreTake(engine->lock, portMAX_DELAY); + storage_index_entry_t *e = find_index_entry(engine, event_id); + if (!e) { xSemaphoreGive(engine->lock); return STORAGE_ERR_NOT_FOUND; } + char path[128]; + get_event_path(e->event_id, e->file_index, path, sizeof(path)); + unlink(path); + e->flags |= STORAGE_FLAG_DELETED; + save_index_to_nvs(engine); + xSemaphoreGive(engine->lock); + return STORAGE_OK; +} + +int storage_purge_expired(storage_engine_t *engine) +{ + if (!engine->initialized) return 0; + xSemaphoreTake(engine->lock, portMAX_DELAY); + uint32_t now = (uint32_t)time(NULL); + int purged = 0; + for (uint16_t i = 0; i < engine->index_count; i++) { + if (engine->index[i].flags & STORAGE_FLAG_DELETED) continue; + if (engine->index[i].expires_at > 0 && engine->index[i].expires_at < now) { + char path[128]; + get_event_path(engine->index[i].event_id, engine->index[i].file_index, path, sizeof(path)); + unlink(path); + engine->index[i].flags |= STORAGE_FLAG_DELETED; + purged++; + } + } + if (purged > 0) { save_index_to_nvs(engine); ESP_LOGI(TAG, "Purged %d expired events", purged); } + xSemaphoreGive(engine->lock); + return purged; +} + +int storage_compact_index(storage_engine_t *engine) +{ + if (!engine->initialized) return 0; + xSemaphoreTake(engine->lock, portMAX_DELAY); + uint16_t write_idx = 0; + int compacted = 0; + for (uint16_t read_idx = 0; read_idx < engine->index_count; read_idx++) { + if (!(engine->index[read_idx].flags & STORAGE_FLAG_DELETED)) { + if (write_idx != read_idx) + memcpy(&engine->index[write_idx], &engine->index[read_idx], sizeof(storage_index_entry_t)); + write_idx++; + } else { + compacted++; + } + } + if (compacted > 0) { + engine->index_count = write_idx; + save_index_to_nvs(engine); + ESP_LOGI(TAG, "Compacted: removed %d, %" PRIu16 " remaining", compacted, engine->index_count); + } + xSemaphoreGive(engine->lock); + return compacted; +} + +void storage_get_stats(storage_engine_t *engine, storage_stats_t *stats) +{ + memset(stats, 0, sizeof(storage_stats_t)); + if (!engine->initialized) return; + xSemaphoreTake(engine->lock, portMAX_DELAY); + uint32_t now = (uint32_t)time(NULL); + for (uint16_t i = 0; i < engine->index_count; i++) { + if (engine->index[i].flags & STORAGE_FLAG_DELETED) continue; + if (engine->index[i].expires_at > 0 && engine->index[i].expires_at < now) continue; + stats->total_events++; + } + size_t total, used; + esp_littlefs_info(STORAGE_PARTITION_LABEL, &total, &used); + stats->total_bytes = total; + stats->free_bytes = total - used; + xSemaphoreGive(engine->lock); +} + +static void storage_cleanup_task(void *arg) +{ + storage_engine_t *engine = (storage_engine_t *)arg; + int cycles = 0; + while (!engine->cleanup_stop) { + for (int i = 0; i < 60 && !engine->cleanup_stop; i++) vTaskDelay(pdMS_TO_TICKS(1000)); + if (engine->cleanup_stop) break; + storage_purge_expired(engine); + if (++cycles >= 10) { storage_compact_index(engine); cycles = 0; } + } + engine->cleanup_task = NULL; + vTaskDelete(NULL); +} + +esp_err_t storage_start_cleanup_task(storage_engine_t *engine) +{ + engine->cleanup_stop = false; + BaseType_t ret = xTaskCreate(storage_cleanup_task, "relay_cleanup", 4096, engine, 2, &engine->cleanup_task); + if (ret != pdPASS) { engine->cleanup_task = NULL; return ESP_ERR_NO_MEM; } + return ESP_OK; +} diff --git a/components/wisp_relay/storage_engine.h b/components/wisp_relay/storage_engine.h new file mode 100644 index 0000000..4e17113 --- /dev/null +++ b/components/wisp_relay/storage_engine.h @@ -0,0 +1,88 @@ +#ifndef STORAGE_ENGINE_H +#define STORAGE_ENGINE_H + +#include +#include +#include "esp_err.h" +#include "freertos/FreeRTOS.h" +#include "freertos/semphr.h" +#include "freertos/task.h" + +#define STORAGE_MAX_EVENTS 5000 +#define STORAGE_MAX_EVENT_SIZE 8192 +#define STORAGE_INDEX_ENTRIES 5000 +#define STORAGE_PARTITION_LABEL "relay_store" + +typedef enum { + STORAGE_OK = 0, + STORAGE_ERR_NOT_INITIALIZED, + STORAGE_ERR_FULL, + STORAGE_ERR_DUPLICATE, + STORAGE_ERR_NOT_FOUND, + STORAGE_ERR_IO, + STORAGE_ERR_NO_MEM, + STORAGE_ERR_SERIALIZE +} storage_error_t; + +#define STORAGE_FLAG_DELETED 0x01 + +typedef struct __attribute__((packed)) { + uint8_t event_id[32]; + uint32_t created_at; + uint32_t expires_at; + uint32_t file_index; + uint16_t kind; + uint8_t pubkey_prefix[4]; + uint8_t flags; + uint8_t reserved; +} storage_index_entry_t; + +typedef struct { + uint32_t total_events; + uint32_t total_bytes; + uint32_t free_bytes; + uint32_t oldest_event_ts; + uint32_t newest_event_ts; +} storage_stats_t; + +typedef struct storage_engine { + storage_index_entry_t *index; + uint16_t index_count; + uint16_t max_index_entries; + uint32_t next_file_index; + SemaphoreHandle_t lock; + TaskHandle_t cleanup_task; + bool initialized; + bool cleanup_stop; + char mount_point[16]; + uint32_t default_ttl_sec; +} storage_engine_t; + +esp_err_t storage_init(storage_engine_t *engine, uint32_t default_ttl_sec); +void storage_destroy(storage_engine_t *engine); + +storage_error_t storage_save_event_json(storage_engine_t *engine, + const char *event_json, + size_t event_json_len); + +storage_error_t storage_query_events_json(storage_engine_t *engine, + int kind, + const char *author_hex, + int limit, + char ***results, + uint16_t *count); + +void storage_free_query_results(char **results, uint16_t count); + +bool storage_event_exists(storage_engine_t *engine, const uint8_t event_id[32]); + +storage_error_t storage_delete_event(storage_engine_t *engine, const uint8_t event_id[32]); + +int storage_purge_expired(storage_engine_t *engine); +int storage_compact_index(storage_engine_t *engine); + +void storage_get_stats(storage_engine_t *engine, storage_stats_t *stats); + +esp_err_t storage_start_cleanup_task(storage_engine_t *engine); + +#endif diff --git a/components/wisp_relay/sub_manager.c b/components/wisp_relay/sub_manager.c new file mode 100644 index 0000000..a1da2e3 --- /dev/null +++ b/components/wisp_relay/sub_manager.c @@ -0,0 +1,272 @@ +#include "sub_manager.h" +#include "relay_types.h" +#include "esp_log.h" +#include +#include + +static const char *TAG = "sub_mgr"; + +static void filter_clear(sub_filter_t *f) +{ + for (size_t i = 0; i < f->ids_count; i++) free(f->ids[i]); + for (size_t i = 0; i < f->authors_count; i++) free(f->authors[i]); + for (size_t i = 0; i < f->e_tags_count; i++) free(f->e_tags[i]); + for (size_t i = 0; i < f->p_tags_count; i++) free(f->p_tags[i]); + memset(f, 0, sizeof(sub_filter_t)); +} + +static bool filter_copy(sub_filter_t *dst, const sub_filter_t *src) +{ + memset(dst, 0, sizeof(sub_filter_t)); + + size_t ids_count = src->ids_count > SUB_MAX_FILTER_IDS ? SUB_MAX_FILTER_IDS : src->ids_count; + for (size_t i = 0; i < ids_count; i++) { + dst->ids[i] = strdup(src->ids[i]); + if (!dst->ids[i]) goto fail; + } + dst->ids_count = ids_count; + + size_t authors_count = src->authors_count > SUB_MAX_FILTER_AUTHORS ? SUB_MAX_FILTER_AUTHORS : src->authors_count; + for (size_t i = 0; i < authors_count; i++) { + dst->authors[i] = strdup(src->authors[i]); + if (!dst->authors[i]) goto fail; + } + dst->authors_count = authors_count; + + size_t kinds_count = src->kinds_count > SUB_MAX_FILTER_KINDS ? SUB_MAX_FILTER_KINDS : src->kinds_count; + memcpy(dst->kinds, src->kinds, kinds_count * sizeof(int32_t)); + dst->kinds_count = kinds_count; + + size_t e_tags_count = src->e_tags_count > SUB_MAX_FILTER_ETAGS ? SUB_MAX_FILTER_ETAGS : src->e_tags_count; + for (size_t i = 0; i < e_tags_count; i++) { + dst->e_tags[i] = strdup(src->e_tags[i]); + if (!dst->e_tags[i]) goto fail; + } + dst->e_tags_count = e_tags_count; + + size_t p_tags_count = src->p_tags_count > SUB_MAX_FILTER_PTAGS ? SUB_MAX_FILTER_PTAGS : src->p_tags_count; + for (size_t i = 0; i < p_tags_count; i++) { + dst->p_tags[i] = strdup(src->p_tags[i]); + if (!dst->p_tags[i]) goto fail; + } + dst->p_tags_count = p_tags_count; + + dst->since = src->since; + dst->until = src->until; + dst->limit = src->limit; + return true; + +fail: + filter_clear(dst); + return false; +} + +static void clear_subscription(subscription_t *sub) +{ + for (uint8_t i = 0; i < sub->filter_count; i++) { + filter_clear(&sub->filters[i]); + } + memset(sub, 0, sizeof(subscription_t)); +} + +esp_err_t sub_manager_init(sub_manager_t *mgr) +{ + memset(mgr, 0, sizeof(sub_manager_t)); + mgr->lock = xSemaphoreCreateMutex(); + if (!mgr->lock) return ESP_ERR_NO_MEM; + ESP_LOGI(TAG, "Initialized (max=%d, per_conn=%d)", SUB_MAX_TOTAL, SUB_MAX_PER_CONN); + return ESP_OK; +} + +void sub_manager_destroy(sub_manager_t *mgr) +{ + if (!mgr) return; + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (mgr->subs[i].active) clear_subscription(&mgr->subs[i]); + } + if (mgr->lock) { vSemaphoreDelete(mgr->lock); mgr->lock = NULL; } +} + +static subscription_t *find_sub(sub_manager_t *mgr, int conn_fd, const char *sub_id) +{ + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (mgr->subs[i].active && mgr->subs[i].conn_fd == conn_fd && + strcmp(mgr->subs[i].sub_id, sub_id) == 0) + return &mgr->subs[i]; + } + return NULL; +} + +static subscription_t *find_free_slot(sub_manager_t *mgr) +{ + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (!mgr->subs[i].active) return &mgr->subs[i]; + } + return NULL; +} + +static bool hex_prefix_match(const char *prefix, size_t prefix_len, + const char *full, size_t full_len) +{ + if (prefix_len == 0) return true; + if (prefix_len > full_len) return false; + return memcmp(prefix, full, prefix_len) == 0; +} + +static bool filter_matches_event(const sub_filter_t *f, int event_kind, + const char *pubkey_hex, uint64_t created_at) +{ + if (f->kinds_count > 0) { + bool found = false; + for (size_t i = 0; i < f->kinds_count; i++) { + if (f->kinds[i] == event_kind) { found = true; break; } + } + if (!found) return false; + } + + if (f->authors_count > 0) { + bool found = false; + for (size_t i = 0; i < f->authors_count; i++) { + if (hex_prefix_match(f->authors[i], strlen(f->authors[i]), + pubkey_hex, strlen(pubkey_hex))) { + found = true; break; + } + } + if (!found) return false; + } + + if (f->since > 0 && (int64_t)created_at < f->since) return false; + if (f->until > 0 && (int64_t)created_at > f->until) return false; + + return true; +} + +void sub_manager_match_json(sub_manager_t *mgr, const char *event_json, + size_t event_len, int event_kind, + const char *event_pubkey_hex, + uint64_t event_created_at, + sub_match_result_t *result) +{ + result->count = 0; + (void)event_json; + (void)event_len; + + xSemaphoreTake(mgr->lock, portMAX_DELAY); + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + subscription_t *sub = &mgr->subs[i]; + if (!sub->active) continue; + + bool matched = false; + for (uint8_t f = 0; f < sub->filter_count; f++) { + if (filter_matches_event(&sub->filters[f], event_kind, + event_pubkey_hex, event_created_at)) { + matched = true; + break; + } + } + if (matched) { + sub_match_entry_t *entry = &result->matches[result->count++]; + entry->conn_fd = sub->conn_fd; + memcpy(entry->sub_id, sub->sub_id, sizeof(entry->sub_id)); + } + } + xSemaphoreGive(mgr->lock); +} + +sub_error_t sub_manager_add(sub_manager_t *mgr, int conn_fd, + const char *sub_id, + const sub_filter_t *filters, + size_t filter_count) +{ + if (filter_count > SUB_MAX_FILTERS) filter_count = SUB_MAX_FILTERS; + + xSemaphoreTake(mgr->lock, portMAX_DELAY); + + subscription_t *existing = find_sub(mgr, conn_fd, sub_id); + if (existing) { + for (uint8_t i = 0; i < existing->filter_count; i++) + filter_clear(&existing->filters[i]); + existing->events_sent = 0; + for (size_t i = 0; i < filter_count; i++) { + if (!filter_copy(&existing->filters[i], &filters[i])) { + existing->filter_count = (uint8_t)i; + xSemaphoreGive(mgr->lock); + return SUB_ERR_MEMORY; + } + } + existing->filter_count = (uint8_t)filter_count; + xSemaphoreGive(mgr->lock); + return SUB_OK; + } + + uint8_t conn_count = 0; + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (mgr->subs[i].active && mgr->subs[i].conn_fd == conn_fd) conn_count++; + } + if (conn_count >= SUB_MAX_PER_CONN) { + xSemaphoreGive(mgr->lock); + return SUB_ERR_TOO_MANY_FILTERS; + } + + subscription_t *slot = find_free_slot(mgr); + if (!slot) { xSemaphoreGive(mgr->lock); return SUB_ERR_MEMORY; } + + memset(slot, 0, sizeof(subscription_t)); + strncpy(slot->sub_id, sub_id, SUB_MAX_ID_LEN); + slot->sub_id[SUB_MAX_ID_LEN] = '\0'; + slot->conn_fd = conn_fd; + + for (size_t i = 0; i < filter_count; i++) { + if (!filter_copy(&slot->filters[i], &filters[i])) { + slot->filter_count = (uint8_t)i; + clear_subscription(slot); + xSemaphoreGive(mgr->lock); + return SUB_ERR_MEMORY; + } + } + slot->filter_count = (uint8_t)filter_count; + slot->active = true; + mgr->active_count++; + + ESP_LOGI(TAG, "Added sub=%s fd=%d filters=%zu total=%d", + sub_id, conn_fd, filter_count, mgr->active_count); + xSemaphoreGive(mgr->lock); + return SUB_OK; +} + +sub_error_t sub_manager_remove(sub_manager_t *mgr, int conn_fd, const char *sub_id) +{ + xSemaphoreTake(mgr->lock, portMAX_DELAY); + subscription_t *sub = find_sub(mgr, conn_fd, sub_id); + if (!sub) { xSemaphoreGive(mgr->lock); return SUB_ERR_NOT_FOUND; } + clear_subscription(sub); + mgr->active_count--; + xSemaphoreGive(mgr->lock); + return SUB_OK; +} + +void sub_manager_remove_all(sub_manager_t *mgr, int conn_fd) +{ + xSemaphoreTake(mgr->lock, portMAX_DELAY); + int removed = 0; + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (mgr->subs[i].active && mgr->subs[i].conn_fd == conn_fd) { + clear_subscription(&mgr->subs[i]); + mgr->active_count--; + removed++; + } + } + if (removed > 0) ESP_LOGI(TAG, "Removed %d subs for fd=%d", removed, conn_fd); + xSemaphoreGive(mgr->lock); +} + +uint8_t sub_manager_count(sub_manager_t *mgr, int conn_fd) +{ + uint8_t count = 0; + xSemaphoreTake(mgr->lock, portMAX_DELAY); + for (int i = 0; i < SUB_MAX_TOTAL; i++) { + if (mgr->subs[i].active && mgr->subs[i].conn_fd == conn_fd) count++; + } + xSemaphoreGive(mgr->lock); + return count; +} diff --git a/components/wisp_relay/sub_manager.h b/components/wisp_relay/sub_manager.h new file mode 100644 index 0000000..64afb04 --- /dev/null +++ b/components/wisp_relay/sub_manager.h @@ -0,0 +1,92 @@ +#ifndef SUB_MANAGER_H +#define SUB_MANAGER_H + +#include +#include +#include "esp_err.h" +#include "freertos/FreeRTOS.h" +#include "freertos/semphr.h" +#include "relay_types.h" + +#define SUB_MAX_TOTAL 64 +#define SUB_MAX_PER_CONN 8 +#define SUB_MAX_FILTERS 4 +#define SUB_MAX_ID_LEN 64 + +#define SUB_MAX_FILTER_IDS 20 +#define SUB_MAX_FILTER_AUTHORS 20 +#define SUB_MAX_FILTER_KINDS 20 +#define SUB_MAX_FILTER_ETAGS 20 +#define SUB_MAX_FILTER_PTAGS 20 + +typedef enum { + SUB_OK = 0, + SUB_ERR_INVALID, + SUB_ERR_TOO_MANY_FILTERS, + SUB_ERR_MEMORY, + SUB_ERR_NOT_FOUND, +} sub_error_t; + +typedef struct { + char *ids[SUB_MAX_FILTER_IDS]; + size_t ids_count; + char *authors[SUB_MAX_FILTER_AUTHORS]; + size_t authors_count; + int32_t kinds[SUB_MAX_FILTER_KINDS]; + size_t kinds_count; + char *e_tags[SUB_MAX_FILTER_ETAGS]; + size_t e_tags_count; + char *p_tags[SUB_MAX_FILTER_PTAGS]; + size_t p_tags_count; + int64_t since; + int64_t until; + int limit; +} sub_filter_t; + +typedef struct { + char sub_id[SUB_MAX_ID_LEN + 1]; + int conn_fd; + sub_filter_t filters[SUB_MAX_FILTERS]; + uint8_t filter_count; + uint16_t events_sent; + bool active; +} subscription_t; + +typedef struct sub_manager { + subscription_t subs[SUB_MAX_TOTAL]; + SemaphoreHandle_t lock; + uint16_t active_count; +} sub_manager_t; + +typedef struct { + int conn_fd; + char sub_id[SUB_MAX_ID_LEN + 1]; +} sub_match_entry_t; + +typedef struct { + sub_match_entry_t matches[SUB_MAX_TOTAL]; + uint8_t count; +} sub_match_result_t; + +esp_err_t sub_manager_init(sub_manager_t *mgr); +void sub_manager_destroy(sub_manager_t *mgr); + +sub_error_t sub_manager_add(sub_manager_t *mgr, int conn_fd, + const char *sub_id, + const sub_filter_t *filters, + size_t filter_count); + +sub_error_t sub_manager_remove(sub_manager_t *mgr, int conn_fd, + const char *sub_id); + +void sub_manager_remove_all(sub_manager_t *mgr, int conn_fd); + +void sub_manager_match_json(sub_manager_t *mgr, const char *event_json, + size_t event_len, int event_kind, + const char *event_pubkey_hex, + uint64_t event_created_at, + sub_match_result_t *result); + +uint8_t sub_manager_count(sub_manager_t *mgr, int conn_fd); + +#endif diff --git a/components/wisp_relay/ws_server.c b/components/wisp_relay/ws_server.c new file mode 100644 index 0000000..a973ca6 --- /dev/null +++ b/components/wisp_relay/ws_server.c @@ -0,0 +1,426 @@ +#include "ws_server.h" +#include "nip11_relay.h" +#include "esp_log.h" +#include "esp_timer.h" +#include +#include +#include +#include +#include +#include +#include + +static const char *TAG = "ws_server"; +static ws_message_cb_t g_message_callback = NULL; +static ws_disconnect_cb_t g_disconnect_callback = NULL; +static ws_server_t *g_server = NULL; +static __thread httpd_req_t *g_current_req = NULL; + +static ws_connection_t* find_free_slot(ws_server_t *server) +{ + for (int i = 0; i < WS_MAX_CONNECTIONS; i++) { + if (!server->connections[i].active) { + return &server->connections[i]; + } + } + return NULL; +} + +static ws_connection_t* find_connection_by_fd(ws_server_t *server, int fd) +{ + for (int i = 0; i < WS_MAX_CONNECTIONS; i++) { + if (server->connections[i].active && server->connections[i].fd == fd) { + return &server->connections[i]; + } + } + return NULL; +} + +static void update_connection_activity(ws_server_t *server, int fd) +{ + xSemaphoreTake(server->lock, portMAX_DELAY); + ws_connection_t *conn = find_connection_by_fd(server, fd); + if (conn) { + conn->last_activity = esp_timer_get_time() / 1000000; + } + xSemaphoreGive(server->lock); +} + +static void set_unknown_ip(char *ip_buf, size_t buf_len) +{ + if (buf_len == 0) { + return; + } + strncpy(ip_buf, "unknown", buf_len - 1); + ip_buf[buf_len - 1] = '\0'; +} + +static void get_client_ip(int fd, char *ip_buf, size_t buf_len) +{ + if (buf_len == 0) { + return; + } + + struct sockaddr_storage addr; + socklen_t addr_len = sizeof(addr); + + if (getpeername(fd, (struct sockaddr *)&addr, &addr_len) != 0) { + set_unknown_ip(ip_buf, buf_len); + return; + } + + const char *result = NULL; + if (addr.ss_family == AF_INET) { + struct sockaddr_in *addr_in = (struct sockaddr_in *)&addr; + result = inet_ntop(AF_INET, &addr_in->sin_addr, ip_buf, buf_len); + } + if (!result) { + set_unknown_ip(ip_buf, buf_len); + } +} + +static esp_err_t on_open(httpd_handle_t hd, int sockfd) +{ + if (!g_server) return ESP_FAIL; + + xSemaphoreTake(g_server->lock, portMAX_DELAY); + + if (g_server->connection_count >= WS_MAX_CONNECTIONS) { + xSemaphoreGive(g_server->lock); + ESP_LOGW(TAG, "Connection rejected - max connections reached"); + return ESP_FAIL; + } + + ws_connection_t *conn = find_free_slot(g_server); + if (!conn) { + xSemaphoreGive(g_server->lock); + ESP_LOGE(TAG, "No free slot despite connection_count < WS_MAX_CONNECTIONS (fd=%d)", sockfd); + return ESP_FAIL; + } + + struct linger so_linger = { .l_onoff = 1, .l_linger = 0 }; + setsockopt(sockfd, SOL_SOCKET, SO_LINGER, &so_linger, sizeof(so_linger)); + + int nodelay = 1; + setsockopt(sockfd, IPPROTO_TCP, TCP_NODELAY, &nodelay, sizeof(nodelay)); + + conn->fd = sockfd; + conn->active = true; + conn->connected_at = esp_timer_get_time() / 1000000; + conn->last_activity = conn->connected_at; + get_client_ip(sockfd, conn->remote_ip, sizeof(conn->remote_ip)); + g_server->connection_count++; + ESP_LOGI(TAG, "New connection from %s (fd=%d, total=%d)", + conn->remote_ip, sockfd, g_server->connection_count); + + xSemaphoreGive(g_server->lock); + return ESP_OK; +} + +static void on_close(httpd_handle_t hd, int sockfd) +{ + if (!g_server) return; + + if (g_disconnect_callback) { + g_disconnect_callback(sockfd); + } + + xSemaphoreTake(g_server->lock, portMAX_DELAY); + + ws_connection_t *conn = find_connection_by_fd(g_server, sockfd); + if (conn) { + ESP_LOGI(TAG, "Connection closed (fd=%d, ip=%s)", sockfd, conn->remote_ip); + memset(conn, 0, sizeof(ws_connection_t)); + g_server->connection_count--; + } + + xSemaphoreGive(g_server->lock); +} + +void ws_server_set_disconnect_cb(ws_disconnect_cb_t cb) +{ + g_disconnect_callback = cb; +} + +static esp_err_t ws_handler(httpd_req_t *req) +{ + if (req->method == HTTP_GET) { + char upgrade[16] = {0}; + if (httpd_req_get_hdr_value_str(req, "Upgrade", upgrade, sizeof(upgrade)) != ESP_OK || + strcasecmp(upgrade, "websocket") != 0) { + return relay_nip11_handler(req); + } + ESP_LOGD(TAG, "WebSocket handshake completed"); + return ESP_OK; + } + + httpd_ws_frame_t ws_pkt; + memset(&ws_pkt, 0, sizeof(httpd_ws_frame_t)); + ws_pkt.type = HTTPD_WS_TYPE_TEXT; + + esp_err_t ret = httpd_ws_recv_frame(req, &ws_pkt, 0); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to get frame len: %d", ret); + return ret; + } + + if (ws_pkt.len == 0) { + return ESP_OK; + } + + if (ws_pkt.len > WS_MAX_FRAME_SIZE) { + ESP_LOGW(TAG, "Frame too large: %zu bytes", ws_pkt.len); + return ESP_FAIL; + } + + ws_pkt.payload = malloc(ws_pkt.len + 1); + if (!ws_pkt.payload) { + ESP_LOGE(TAG, "Failed to allocate %zu bytes", ws_pkt.len); + return ESP_ERR_NO_MEM; + } + + ret = httpd_ws_recv_frame(req, &ws_pkt, ws_pkt.len); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to receive frame: %d", ret); + free(ws_pkt.payload); + return ret; + } + + ((char *)ws_pkt.payload)[ws_pkt.len] = '\0'; + + int fd = httpd_req_to_sockfd(req); + if (g_server) { + update_connection_activity(g_server, fd); + } + + switch (ws_pkt.type) { + case HTTPD_WS_TYPE_TEXT: + ESP_LOGD(TAG, "Received %zu bytes from fd=%d", ws_pkt.len, fd); + if (g_message_callback) { + g_current_req = req; + g_message_callback(fd, (char *)ws_pkt.payload, ws_pkt.len); + g_current_req = NULL; + } + break; + + case HTTPD_WS_TYPE_PING: + ws_pkt.type = HTTPD_WS_TYPE_PONG; + ret = httpd_ws_send_frame(req, &ws_pkt); + if (ret != ESP_OK) { + ESP_LOGW(TAG, "Failed to send PONG to fd=%d: %d", fd, ret); + free(ws_pkt.payload); + return ret; + } + break; + + case HTTPD_WS_TYPE_CLOSE: { + ESP_LOGD(TAG, "Received CLOSE frame from fd=%d", fd); + free(ws_pkt.payload); + httpd_ws_frame_t close_pkt = { + .type = HTTPD_WS_TYPE_CLOSE, + .payload = NULL, + .len = 0, + }; + httpd_ws_send_frame(req, &close_pkt); + return ESP_FAIL; + } + + default: + break; + } + + free(ws_pkt.payload); + return ESP_OK; +} + +typedef struct { + httpd_handle_t hd; + int fd; + char *data; + size_t len; +} async_send_arg_t; + +static void ws_async_send(void *arg) +{ + async_send_arg_t *a = (async_send_arg_t *)arg; + + httpd_ws_frame_t ws_pkt = { + .type = HTTPD_WS_TYPE_TEXT, + .payload = (uint8_t *)a->data, + .len = a->len, + }; + + esp_err_t ret = httpd_ws_send_frame_async(a->hd, a->fd, &ws_pkt); + if (ret != ESP_OK) { + ESP_LOGW(TAG, "Async send failed to fd=%d: %d", a->fd, ret); + } + + free(a->data); + free(a); +} + +static void cleanup_server_init(ws_server_t *server, bool stop_httpd) +{ + g_server = NULL; + g_message_callback = NULL; + if (stop_httpd && server->server) { + httpd_stop(server->server); + server->server = NULL; + } + if (server->lock) { + vSemaphoreDelete(server->lock); + server->lock = NULL; + } +} + +esp_err_t ws_server_init(ws_server_t *server, uint16_t port, ws_message_cb_t on_message) +{ + if (server->server != NULL) { + ESP_LOGE(TAG, "Server already initialized, call ws_server_stop first"); + return ESP_ERR_INVALID_STATE; + } + + memset(server, 0, sizeof(ws_server_t)); + server->lock = xSemaphoreCreateMutex(); + if (!server->lock) { + return ESP_ERR_NO_MEM; + } + + g_server = server; + g_message_callback = on_message; + + httpd_config_t config = HTTPD_DEFAULT_CONFIG(); + config.server_port = port; + config.ctrl_port = port + 1; + config.max_open_sockets = WS_MAX_CONNECTIONS; + config.backlog_conn = WS_MAX_CONNECTIONS; + config.lru_purge_enable = true; + config.recv_wait_timeout = 3; + config.send_wait_timeout = 3; + config.keep_alive_enable = true; + config.keep_alive_idle = 5; + config.keep_alive_interval = 1; + config.keep_alive_count = 3; + config.stack_size = 12288; + config.open_fn = on_open; + config.close_fn = on_close; + + esp_err_t ret = httpd_start(&server->server, &config); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to start server: %d", ret); + cleanup_server_init(server, false); + return ret; + } + + httpd_uri_t ws_uri = { + .uri = "/", + .method = HTTP_GET, + .handler = ws_handler, + .user_ctx = NULL, + .is_websocket = true, + .handle_ws_control_frames = true, + }; + + ret = httpd_register_uri_handler(server->server, &ws_uri); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to register WS handler: %d", ret); + cleanup_server_init(server, true); + return ret; + } + + httpd_uri_t options_uri = { + .uri = "/", + .method = HTTP_OPTIONS, + .handler = relay_nip11_options_handler, + .user_ctx = NULL, + }; + + ret = httpd_register_uri_handler(server->server, &options_uri); + if (ret != ESP_OK) { + ESP_LOGE(TAG, "Failed to register OPTIONS handler: %d", ret); + } + + ESP_LOGI(TAG, "WebSocket server started on port %d", port); + return ESP_OK; +} + +void ws_server_stop(ws_server_t *server) +{ + g_server = NULL; + g_message_callback = NULL; + g_disconnect_callback = NULL; + + if (server->server) { + httpd_stop(server->server); + server->server = NULL; + } + if (server->lock) { + vSemaphoreDelete(server->lock); + server->lock = NULL; + } + memset(server->connections, 0, sizeof(server->connections)); + server->connection_count = 0; +} + +bool ws_server_is_running(ws_server_t *server) +{ + return server && server->server != NULL; +} + +esp_err_t ws_server_send(ws_server_t *server, int fd, const char *data, size_t len) +{ + if (!server->server) return ESP_ERR_INVALID_STATE; + + if (g_current_req && httpd_req_to_sockfd(g_current_req) == fd) { + httpd_ws_frame_t ws_pkt = { + .type = HTTPD_WS_TYPE_TEXT, + .payload = (uint8_t *)data, + .len = len, + }; + return httpd_ws_send_frame(g_current_req, &ws_pkt); + } + + async_send_arg_t *arg = malloc(sizeof(async_send_arg_t)); + if (!arg) return ESP_ERR_NO_MEM; + + arg->data = malloc(len); + if (!arg->data) { + free(arg); + return ESP_ERR_NO_MEM; + } + + memcpy(arg->data, data, len); + arg->hd = server->server; + arg->fd = fd; + arg->len = len; + + esp_err_t ret = httpd_queue_work(server->server, ws_async_send, arg); + if (ret != ESP_OK) { + free(arg->data); + free(arg); + return ret; + } + return ESP_OK; +} + +esp_err_t ws_server_broadcast(ws_server_t *server, const char *data, size_t len) +{ + xSemaphoreTake(server->lock, portMAX_DELAY); + + for (int i = 0; i < WS_MAX_CONNECTIONS; i++) { + if (server->connections[i].active) { + ws_server_send(server, server->connections[i].fd, data, len); + } + } + + xSemaphoreGive(server->lock); + return ESP_OK; +} + +void ws_server_close_connection(ws_server_t *server, int fd) +{ + if (!server || !server->server) { + return; + } + httpd_sess_trigger_close(server->server, fd); +} diff --git a/components/wisp_relay/ws_server.h b/components/wisp_relay/ws_server.h new file mode 100644 index 0000000..4fe616e --- /dev/null +++ b/components/wisp_relay/ws_server.h @@ -0,0 +1,41 @@ +#ifndef WS_SERVER_H +#define WS_SERVER_H + +#include +#include +#include +#include "esp_http_server.h" +#include "freertos/FreeRTOS.h" +#include "freertos/semphr.h" + +#define WS_MAX_CONNECTIONS 8 +#define WS_MAX_FRAME_SIZE 65536 +#define WS_IP_ADDR_MAX_LEN 48 + +typedef struct { + int fd; + bool active; + uint32_t connected_at; + uint32_t last_activity; + char remote_ip[WS_IP_ADDR_MAX_LEN]; +} ws_connection_t; + +typedef struct { + httpd_handle_t server; + ws_connection_t connections[WS_MAX_CONNECTIONS]; + SemaphoreHandle_t lock; + uint8_t connection_count; +} ws_server_t; + +typedef void (*ws_message_cb_t)(int fd, const char *data, size_t len); +typedef void (*ws_disconnect_cb_t)(int fd); + +esp_err_t ws_server_init(ws_server_t *server, uint16_t port, ws_message_cb_t on_message); +void ws_server_set_disconnect_cb(ws_disconnect_cb_t cb); +void ws_server_stop(ws_server_t *server); +bool ws_server_is_running(ws_server_t *server); +esp_err_t ws_server_send(ws_server_t *server, int fd, const char *data, size_t len); +esp_err_t ws_server_broadcast(ws_server_t *server, const char *data, size_t len); +void ws_server_close_connection(ws_server_t *server, int fd); + +#endif -- cgit v1.2.3