From 491e97e2755a7bf755b2d9b919d8053cd3946aad Mon Sep 17 00:00:00 2001 From: Laan Tungir Date: Thu, 13 Aug 2026 09:49:48 -0400 Subject: [PATCH] v2.1.39 - Security audit: fix 3 Critical and 7 High severity vulnerabilities; add missing DB index for admin stats --- API.md | 9 + Dockerfile.alpine-musl | 2 +- Makefile | 2 +- README.md | 1 + audits/security_remediation.md | 83 +++ docs/configuration_guide.md | 43 ++ docs/deployment_guide.md | 28 +- docs/user_guide.md | 28 + plans/status_publish_implementation_plan.md | 162 +++++ plans/udp_ingress_integration_plan.md | 693 ++++++++++++++++++++ relay.pid | 2 +- src/api.c | 116 +++- src/db_ops_postgres.c | 12 +- src/default_config_event.h | 8 +- src/main.c | 199 +++--- src/main.h | 35 +- src/nip042.c | 111 ++++ src/pg_schema.h | 2 + src/request_validator.c | 30 +- src/sql_schema.h | 3 + src/udp_ingress.c | 358 ++++++++++ src/udp_ingress.h | 25 + src/websockets.c | 77 ++- tests/udp_ingress_test.sh | 203 ++++++ 24 files changed, 2062 insertions(+), 170 deletions(-) create mode 100644 audits/security_remediation.md create mode 100644 plans/status_publish_implementation_plan.md create mode 100644 plans/udp_ingress_integration_plan.md create mode 100644 src/udp_ingress.c create mode 100644 src/udp_ingress.h create mode 100755 tests/udp_ingress_test.sh diff --git a/API.md b/API.md index c711820..4a34b4e 100644 --- a/API.md +++ b/API.md @@ -559,6 +559,15 @@ ORDER BY count DESC |-----|------|---------|-------------| | `kind_24567_reporting_throttle_sec` | integer | 5 | Monitoring event throttle | +### UDP Nostr Ingress + +| Key | Type | Default | Description | +|-----|------|---------|-------------| +| `udp_ingress_enabled` | boolean | false | Enable UDP event listener (no-handshake, no-response) | +| `udp_ingress_port` | integer | 443 | UDP port (TCP/UDP namespaces are separate) | +| `udp_ingress_bind_addr` | string | "0.0.0.0" | UDP bind address | +| `udp_ingress_max_datagram_size` | integer | 1472 | Max datagram size (single-packet guarantee) | + ### Dynamic vs Restart-Required **Dynamic (No Restart)**: diff --git a/Dockerfile.alpine-musl b/Dockerfile.alpine-musl index f6dc72d..afdf73e 100644 --- a/Dockerfile.alpine-musl +++ b/Dockerfile.alpine-musl @@ -125,7 +125,7 @@ RUN if [ "$DEBUG_BUILD" = "true" ]; then \ src/main.c src/config.c src/dm_admin.c src/request_validator.c \ src/nip009.c src/nip011.c src/nip013.c src/nip040.c src/nip042.c \ src/websockets.c src/subscriptions.c src/api.c src/embedded_web_content.c src/ip_ban.c \ - src/caching_inbox_poller.c src/caching_service_launcher.c \ + src/caching_inbox_poller.c src/caching_service_launcher.c src/udp_ingress.c \ src/db_ops.c src/db_ops_postgres.c src/thread_pool.c \ -o /build/c_relay_pg_static \ c_utils_lib/libc_utils.a \ diff --git a/Makefile b/Makefile index 489811c..9efa718 100644 --- a/Makefile +++ b/Makefile @@ -9,7 +9,7 @@ LIBS = -lwebsockets -lz -ldl -lpthread -lm -L/usr/local/lib -lsecp256k1 -lssl -l BUILD_DIR = build # Source files -MAIN_SRC = src/main.c src/config.c src/dm_admin.c src/request_validator.c src/nip009.c src/nip011.c src/nip013.c src/nip040.c src/nip042.c src/websockets.c src/subscriptions.c src/api.c src/embedded_web_content.c src/ip_ban.c src/thread_pool.c src/caching_inbox_poller.c src/caching_service_launcher.c +MAIN_SRC = src/main.c src/config.c src/dm_admin.c src/request_validator.c src/nip009.c src/nip011.c src/nip013.c src/nip040.c src/nip042.c src/websockets.c src/subscriptions.c src/api.c src/embedded_web_content.c src/ip_ban.c src/thread_pool.c src/caching_inbox_poller.c src/caching_service_launcher.c src/udp_ingress.c DB_OPS_SRC = src/db_ops.c src/db_ops_postgres.c NOSTR_CORE_LIB = nostr_core_lib/libnostr_core_x64.a diff --git a/README.md b/README.md index 7cbc0f4..980fe49 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,7 @@ Control your relay by sending direct messages from any Nostr client: - **Static binary available** - Single file deployment with zero dependencies - **Efficient memory management** - Optimized for long-running operation - **WebSocket native** - Built on libwebsockets for optimal protocol support +- **UDP Nostr Ingress** - Receive events as single UDP datagrams (â‰Ī1472 bytes) with no handshake, no connection state, and no response — the most censorship-resistant way to publish to a Nostr relay. Send an event with a single line: `nak event -k 1 -c "Hello!" | nc -u -w1 laantungir.net 443` ## 📋 Supported NIPs diff --git a/audits/security_remediation.md b/audits/security_remediation.md new file mode 100644 index 0000000..bcae4f3 --- /dev/null +++ b/audits/security_remediation.md @@ -0,0 +1,83 @@ +# C-Relay-PG Security Remediation Summary + +**Date:** 2026-08-13 +**Status:** 10 of 29 findings remediated (3 Critical, 7 High) + +--- + +## ✅ Fixed and Deployed + +### Critical (3) + +| ID | Title | Fix | File(s) | +|----|-------|-----|---------| +| C-01 | SQL Injection via NIP-50 Search Filter | Replaced manual quote escaping with parameterized queries (`?` placeholders + bind parameters) | [`src/websockets.c`](src/websockets.c), [`src/main.c`](src/main.c) | +| C-02 | Weak RNG for NIP-42 Challenges | Replaced `srand(time(NULL))` + `rand()` with `/dev/urandom` for cryptographically secure random bytes | [`src/request_validator.c`](src/request_validator.c) | +| C-03 | Command Injection via popen() | Documented that all `popen()` calls use fixed compile-time strings (no user input); added `execve()`-based alternative for PHP report generation | [`src/api.c`](src/api.c) | + +### High (7) + +| ID | Title | Fix | File(s) | +|----|-------|-----|---------| +| H-01 | SQL Injection via LISTEN Channel | Added `PQescapeIdentifier()` for defense-in-depth alongside existing character validation | [`src/db_ops_postgres.c`](src/db_ops_postgres.c) | +| H-02 | Use-After-Free in Async Event Completion | Added `current_pss` null check before `send_ok_response()` | [`src/websockets.c`](src/websockets.c) | +| H-03 | Buffer Overflow in Config Parsing | Replaced all `strcpy()` calls with `snprintf()` using fixed buffer size (128) | [`src/api.c`](src/api.c) | +| H-04 | SQL Injection via Direct Query Execution | Added whitelist-based validation (SELECT/WITH only) with improved blacklist and word-boundary matching | [`src/api.c`](src/api.c) | +| H-05 | Missing Rate Limiting on Auth Endpoints | Added per-IP rate limiting (10 attempts per 60-second window) with mutex-protected tracking | [`src/nip042.c`](src/nip042.c) | +| H-06 | Unbounded Memory Growth in Message Reassembly | Added 10MB maximum message size limit with early rejection before allocation | [`src/websockets.c`](src/websockets.c) | +| H-07 | Race Condition in Connection Tracking | Added `pthread_mutex_lock(&pss->session_lock)` when reading `session_active`/`connection_established` | [`src/websockets.c`](src/websockets.c) | + +### Additional Fixes + +| Fix | Description | File(s) | +|-----|-------------|---------| +| Missing DB index | Added `idx_subscriptions_active_lookup` partial index for admin stats page queries | [`src/pg_schema.h`](src/pg_schema.h), [`src/sql_schema.h`](src/sql_schema.h) | + +--- + +## 📋 Remaining Findings + +### Medium (11) + +| ID | Title | File | Description | +|----|-------|------|-------------| +| M-01 | Information Disclosure via SQL Error Messages | `src/db_ops_postgres.c` | SQL error details leaked to clients in error responses | +| M-02 | Predictable Config Change IDs | `src/api.c:2457-2467` | Config change IDs generated with predictable pattern | +| M-03 | Missing CSRF Protection in Admin API | `admin/api/*.php` | No CSRF tokens on state-changing admin endpoints | +| M-04 | Weak Session Timeout Handling | `src/nip042.c:48` | NIP-42 challenges have 10-minute expiration (too long) | +| M-05 | Unvalidated UDP Source Address | `src/udp_ingress.c:132` | UDP datagrams accepted from any source without validation | +| M-06 | Missing Input Validation on Admin Commands | `src/dm_admin.c` | Admin commands not validated against exact whitelist | +| M-07 | Potential Integer Overflow in Query Limits | `src/websockets.c:2120` | Integer overflow possible in query limit calculations | +| M-08 | Unrestricted File Read via Embedded Files | `src/api.c:1199` | No path traversal protection for embedded file serving | +| M-09 | Unvalidated Relay URLs in Caching Service | `caching/src/relay_discovery.c:103` | Relay URLs from NIP-65 events accepted without scheme/format validation | +| M-10 | Config File Written with World-Readable Permissions | `caching/src/config.c:265` | Config file written with mode 0644 instead of 0640/0600 | +| M-11 | Use of eval() in Restart Script | `make_and_restart_relay.sh:651` | Shell script uses `eval` with user-controlled arguments | + +### Low/Info (8) + +| ID | Title | File | Description | +|----|-------|------|-------------| +| L-01 | Hardcoded Default Database Credentials | `admin/lib/config.php:24-28` | Default DB credentials (`crelay`/`crelay`) in config file | +| L-02 | Missing Security Headers in HTTP Responses | `src/nip011.c` | No CSP, HSTS, X-Frame-Options headers on HTTP responses | +| L-03 | Debug Information Leakage in Production | `src/main.c`, `src/config.c` | Debug logging may expose sensitive info in production | +| L-04 | Inconsistent Error Handling | Multiple files | Mix of error return codes, NULL returns, and output params | +| L-05 | Missing Input Sanitization for Log Messages | Multiple files | User input logged without sanitization (log injection) | +| L-06 | Potential String Truncation in strncpy Calls | `caching/src/*.c` | Multiple `strncpy` calls without explicit null termination | +| L-07 | No Binary Integrity Verification in Deployment | `deploy_lt.sh` | No checksum/signature verification on deployed binaries | +| L-08 | Hardcoded Database Credentials in Systemd Service | `systemd/c-relay-pg-local.service:13` | DB password visible in service file and process list | + +--- + +## Summary + +| Severity | Total | Fixed | Remaining | +|----------|-------|-------|-----------| +| Critical | 3 | 3 | 0 | +| High | 7 | 7 | 0 | +| Medium | 11 | 0 | 11 | +| Low/Info | 8 | 0 | 8 | +| **Total** | **29** | **10** | **19** | + +--- + +*Generated: 2026-08-13* diff --git a/docs/configuration_guide.md b/docs/configuration_guide.md index f793409..9ba3668 100644 --- a/docs/configuration_guide.md +++ b/docs/configuration_guide.md @@ -175,6 +175,49 @@ Configuration events follow the standard Nostr event format with kind 33334: - **Impact**: Allows some flexibility in expiration timing - **Example**: `"600"` (10 minute grace period) +### UDP Nostr Ingress (No-Handshake Event Reception) + +The relay can receive Nostr events as single UDP datagrams — no TCP connection, +no WebSocket handshake, no TLS negotiation. The event carries its own Schnorr +signature, which replaces the handshake. This provides censorship resistance: +a passive observer cannot distinguish a Nostr event from random UDP noise, and +active probing cannot detect the relay (invalid events are silently dropped). + +#### `udp_ingress_enabled` +- **Description**: Enable the UDP event listener +- **Default**: `"false"` +- **Values**: `"true"` or `"false"` +- **Impact**: When enabled, the relay binds a UDP socket and accepts Nostr events + as single datagrams. Events are validated through the same pipeline as WebSocket + events (signature, structure, expiration, PoW) and broadcast to subscribers. +- **Note**: This is a dynamic config key — the listener starts/stops without + restart when the value changes (via `udp_ingress_tick()` in the main loop). + +#### `udp_ingress_port` +- **Description**: UDP port for the event listener +- **Default**: `"443"` +- **Range**: `1` to `65535` +- **Impact**: Ports below 1024 require `CAP_NET_BIND_SERVICE` capability (see + deployment guide). TCP and UDP port namespaces are separate — UDP 443 can + coexist with an HTTPS server on TCP 443. +- **Recommendation**: Use port 443 to blend with QUIC/HTTP3 traffic for maximum + censorship resistance. Use port 8888 for local testing. + +#### `udp_ingress_bind_addr` +- **Description**: IP address to bind the UDP socket to +- **Default**: `"0.0.0.0"` +- **Format**: IPv4 address string +- **Example**: `"127.0.0.1"` (listen only on localhost) + +#### `udp_ingress_max_datagram_size` +- **Description**: Maximum allowed UDP datagram size in bytes +- **Default**: `"1472"` +- **Range**: `512` to `65535` +- **Impact**: 1472 bytes is the single-packet guarantee (1500 Ethernet MTU − 20 + IP header − 8 UDP header). Datagrams larger than this are silently dropped. + This ensures the event fits in a single unfragmented IP packet, which is the + foundation of the censorship-resistance property. + ### NIP-59 Gift Wrap Timestamp Configuration #### `nip59_timestamp_max_delay_sec` diff --git a/docs/deployment_guide.md b/docs/deployment_guide.md index 13cdbe7..bcae83d 100644 --- a/docs/deployment_guide.md +++ b/docs/deployment_guide.md @@ -90,6 +90,24 @@ sudo cp systemd/c-relay-pg.service /etc/systemd/system/ sudo systemctl daemon-reload ``` +### UDP Nostr Ingress (Privileged Port Binding) + +If UDP ingress is enabled on a port below 1024 (e.g., the default port 443), +the relay needs the `CAP_NET_BIND_SERVICE` capability to bind the port as a +non-root user. Add this to the `[Service]` section of the systemd unit file: + +```ini +AmbientCapabilities=CAP_NET_BIND_SERVICE +``` + +This grants the `c-relay-pg` service user permission to bind privileged ports +without running as root. The capability is automatically inherited by the +forked relay process and does not require any changes to the binary. + +**Note**: TCP and UDP port namespaces are separate. The UDP listener on port +443 coexists with an HTTPS server (nginx, caddy) on TCP port 443 without +conflict. This is the same mechanism used by QUIC/HTTP3. + ### Service Management #### Start and Enable Service @@ -119,9 +137,12 @@ sudo journalctl -u c-relay-pg --no-pager | grep "Admin Private Key" #### UFW (Ubuntu) ```bash -# Allow relay port +# Allow relay WebSocket port sudo ufw allow 8888/tcp +# Allow UDP Nostr ingress (if enabled) +sudo ufw allow 443/udp + # Allow SSH (ensure you don't lock yourself out) sudo ufw allow 22/tcp @@ -131,9 +152,12 @@ sudo ufw enable #### iptables ```bash -# Allow relay port +# Allow relay WebSocket port sudo iptables -A INPUT -p tcp --dport 8888 -j ACCEPT +# Allow UDP Nostr ingress (if enabled) +sudo iptables -A INPUT -p udp --dport 443 -j ACCEPT + # Save rules (Ubuntu/Debian) sudo iptables-save > /etc/iptables/rules.v4 ``` diff --git a/docs/user_guide.md b/docs/user_guide.md index a7089a7..4629076 100644 --- a/docs/user_guide.md +++ b/docs/user_guide.md @@ -46,6 +46,26 @@ Your relay is now available at: - **WebSocket**: `ws://localhost:8888` - **NIP-11 Info**: `http://localhost:8888` (with `Accept: application/nostr+json` header) - **Web Admin Interface**: `http://localhost:8888/api/` (serves embedded admin interface) +- **UDP Nostr Ingress**: `udp://localhost:443` (send events as single UDP datagrams, no handshake) + +#### Sending Events via UDP + +If UDP ingress is enabled (`udp_ingress_enabled=true`), you can send Nostr events +as single UDP datagrams — no TCP connection, no WebSocket handshake, no TLS: + +```bash +# Simplest: nak + nc (zero custom code, nak uses a default key) +nak event -k 1 -c "Hello via UDP!" | nc -u -w1 laantungir.net 443 + +# Using nak + Python (4 lines, explicit key) +nak event -k 1 -c "Hello via UDP!" --sec $(nak key generate) \ + | python3 -c "import socket, sys; sock=socket.socket(socket.AF_INET,socket.SOCK_DGRAM); sock.sendto(sys.stdin.read().encode(),('laantungir.net',443))" +``` + +The event is received as a single UDP datagram (â‰Ī1472 bytes), signature-validated, +stored, and broadcast to all WebSocket subscribers. Invalid events are silently +dropped — the relay never responds to UDP, making it indistinguishable from a +non-Nostr service to active probing. ## Installation @@ -213,6 +233,14 @@ Send this to your relay via WebSocket, and changes are applied immediately. | `nip40_expiration_filter` | Filter expired events | "true" | "true", "false" | | `nip40_expiration_grace_period` | Grace period (seconds) | "300" | 0-86400 | +#### UDP Nostr Ingress +| Parameter | Description | Default | Options | +|-----------|-------------|---------|---------| +| `udp_ingress_enabled` | Enable UDP event listener | "false" | "true", "false" | +| `udp_ingress_port` | UDP port for event reception | "443" | 1-65535 | +| `udp_ingress_bind_addr` | UDP bind address | "0.0.0.0" | IPv4 address | +| `udp_ingress_max_datagram_size` | Max datagram size (bytes) | "1472" | 512-65535 | + ## Web Admin Interface The relay includes a built-in web-based administration interface that provides a user-friendly way to manage your relay without command-line tools. diff --git a/plans/status_publish_implementation_plan.md b/plans/status_publish_implementation_plan.md new file mode 100644 index 0000000..23f8ec1 --- /dev/null +++ b/plans/status_publish_implementation_plan.md @@ -0,0 +1,162 @@ +# Implementation Plan: Republish Kind 1 Status Events to External Relays + +## Overview + +The main relay (`c_relay_pg`) already does all the work via +[`generate_and_post_status_event()`](src/api.c:976): +- Generates stats text +- Gets the relay private key +- Creates and signs a kind 1 event +- Stores it in the `events` table +- Broadcasts to local WebSocket clients + +The caching service just needs to **detect new kind 1 events by the relay's +pubkey** and republish them to upstream relays. No key management, no stats +generation, no event signing — just a simple poll-and-republish. + +## Architecture + +```mermaid +flowchart TB + subgraph MainRelay[c_relay_pg - already done] + Gen[generate_and_post_status_event] + Store[Store kind 1 event in events table] + LocalBC[Broadcast to local clients] + Gen --> Store --> LocalBC + end + + subgraph CachingService[caching_relay - new] + Poll[Poll events table for new kind 1\nby relay pubkey] + Republish[Republish to upstream relays\nvia nostr_relay_pool_publish_async] + Poll --> Republish + end + + subgraph PostgreSQL + Events[events table] + Config[config table\nrelay_pubkey] + end + + subgraph UpstreamRelays + R1[Relay 1] + R2[Relay 2] + R3[Relay N] + end + + Store --> Events + Config --> Poll + Events --> Poll + Republish --> R1 + Republish --> R2 + Republish --> R3 +``` + +## Implementation Steps + +### Step 1: Add `pg_inbox_get_latest_status_event()` to pg_inbox + +**Files**: [`caching/src/pg_inbox.h`](caching/src/pg_inbox.h), +[`caching/src/pg_inbox.c`](caching/src/pg_inbox.c) + +A function that queries the `events` table for the most recent kind 1 event +by the relay's pubkey. The relay pubkey is read from the `config` table +(same as the main relay does at [`src/api.c:899`](src/api.c:899)). + +```c +/* Get the latest kind 1 status event by the relay pubkey. + * Reads relay_pubkey from config, then queries events table. + * On success: fills out_event_json with a malloc'd string (caller frees) + * and out_event_id with the 64-char hex id. Returns 0 on success, + * -1 on error, 1 if no status event found. */ +int pg_inbox_get_latest_status_event(char **out_event_json, + char *out_event_id, int id_len); +``` + +SQL: +```sql +SELECT e.event_json, e.id +FROM events e +WHERE e.pubkey = $1 AND e.kind = 1 +ORDER BY e.created_at DESC +LIMIT 1 +``` + +### Step 2: Add status republish tick to main loop + +**File**: [`caching/src/main.c`](caching/src/main.c) + +Add a periodic task in the main loop that: + +1. Checks if `kind_1_status_posts_hours` config is enabled (via + `pg_inbox_get_config_value()`) +2. Polls for the latest kind 1 event by the relay pubkey +3. If the event ID is different from the last one we republished: + - Parses the event JSON with `cJSON_Parse()` + - Gets the list of connected upstream relay URLs + - Calls `nostr_relay_pool_publish_async(upstream, urls, count, event, callback, NULL)` + - Records the event ID as "last republished" +4. Runs on a reasonable poll interval (e.g., every 60 seconds) — the actual + publish frequency is controlled by the main relay's + `kind_1_status_posts_hours` config, which determines how often new kind 1 + events appear in the database + +A `time_t last_status_check = 0` and `char last_republished_id[65] = {0}` +track state. + +The publish callback is a simple logging function (same pattern as +[`sink_publish_cb()`](caching/src/relay_sink.c:18)). + +**Event ownership**: `nostr_relay_pool_publish_async()` takes a `cJSON* event`. +We need to verify whether the pool takes ownership or we need to keep it +alive. The safest approach: pass the parsed event and free it in the callback, +or make a copy. We'll check the nostr_core_lib implementation during coding. + +### Step 3: Build the caching service + +```bash +cd caching && make +``` + +### Step 4: Build and start the test relay on port 8888 + +```bash +./make_and_restart_relay.sh +``` + +Production on port 7777 (database `crelay_prod`) is untouched. + +### Step 5: Enable status publishing in the test database + +```bash +PGPASSWORD=crelay psql -h localhost -U crelay -d crelay -c \ + "INSERT INTO config (key, value, type, description, category, editable) \ + VALUES ('kind_1_status_posts_hours', '1', 'int', \ + 'Hours between kind 1 status posts', 'relay', 1) \ + ON CONFLICT (key) DO UPDATE SET value='1';" +``` + +### Step 6: Run the caching service and verify + +```bash +./build/caching_relay -p "host=127.0.0.1 dbname=crelay user=crelay password=crelay" -d 4 +``` + +Check `caching_relay.log` for republish log messages. + +### Step 7: Verify published events + +```bash +# Get relay pubkey +PGPASSWORD=crelay psql -h localhost -U crelay -d crelay -c \ + "SELECT value FROM config WHERE key='relay_pubkey';" + +# Query an upstream relay +nak req -k 1 -a wss://relay.damus.io +``` + +## Files to Modify + +| File | Change | +|------|--------| +| [`caching/src/pg_inbox.h`](caching/src/pg_inbox.h) | Add `pg_inbox_get_latest_status_event()` declaration | +| [`caching/src/pg_inbox.c`](caching/src/pg_inbox.c) | Implement `pg_inbox_get_latest_status_event()` | +| [`caching/src/main.c`](caching/src/main.c) | Add status republish tick + publish callback | diff --git a/plans/udp_ingress_integration_plan.md b/plans/udp_ingress_integration_plan.md new file mode 100644 index 0000000..e37f744 --- /dev/null +++ b/plans/udp_ingress_integration_plan.md @@ -0,0 +1,693 @@ +# UDP Nostr Ingress Integration Plan + +## Overview + +Integrate the [`udp_nostr`](../../lt/udp_nostr) concept into c-relay-pg: a direct UDP +listener that receives signed Nostr events as single datagrams (â‰Ī1472 bytes, JSON +format), validates them through the existing authoritative pipeline, stores them, and +broadcasts to WebSocket subscribers — all with no handshake, no response, and no +connection state. + +The design is structured so a future **authoritative DNS relay** ingress (events +received via DNS queries on port 53) can reuse the same completion queue, config keys, +and `ingest_event()` source path with minimal additional work. + +--- + +## Background & Key Constraints + +From the [`udp_nostr`](../../lt/udp_nostr) project: + +- **No handshake**: A single UDP datagram carries a self-validating Nostr event. The + signature replaces the handshake. The relay never responds to invalid events + (silent drop — no protocol fingerprint for active probing). +- **Single-packet guarantee**: 1500 Ethernet MTU − 20 IP − 8 UDP = **1472 bytes** + max payload. Events larger than this are silently dropped (no fragmentation). +- **Stateless**: No connection state, no session, no OK response. The relay receives, + validates, stores, and forgets the sender. +- **Fire-and-forget**: UDP provides no delivery guarantee. This is acceptable for the + ingress path — the sender does not wait for an acknowledgment. + +From c-relay-pg: + +- The main thread runs the libwebsockets service loop (`lws_service()` with 1000ms + timeout) in [`start_websocket_relay()`](src/websockets.c:3432). +- [`ingest_event()`](src/main.c:1281) is the shared entry point for non-client event + ingestion. It performs signature/structure/expiration/PoW validation, stores through + the writer thread pool, and pushes a completion record to a queue drained on the lws + thread by [`process_inbox_event_completions()`](src/main.c:1241). +- The `EVENT_SOURCE_CACHING_INBOX` source mode already implements the exact semantics + UDP needs: no `wsi`/`pss`, no OK response, no admin command execution (kind 23456 + stored as data), still validates and broadcasts. +- The caching inbox poller ([`src/caching_inbox_poller.c`](src/caching_inbox_poller.c)) + runs on the main lws thread and calls `ingest_event()` directly — it does NOT use a + separate thread. The UDP listener CANNOT do this because `recvfrom()` is blocking. + +--- + +## Architecture + +```mermaid +flowchart TD + subgraph "UDP Sender (any client)" + S[ nak event ... | python udp_nostr_send.py ] + end + + subgraph "c-relay-pg process" + subgraph "UDP Listener Thread" + UDPL[udp_ingress_thread
blocking recvfrom on UDP port] + Q[ring buffer / queue
thread-safe] + end + + subgraph "Main lws Thread" + TICK[udp_ingress_tick
drains queue each loop iteration] + INGEST[ingest_event
EVENT_SOURCE_UDP_INGRESS] + VAL[nostr_validate_unified_request
signature + structure + PoW + expiration] + STORE[store_event_core
via writer thread pool] + COMP[external_ingress_completion_queue] + BCAST[broadcast_event_to_subscriptions] + end + + subgraph "Writer Thread Pool" + WP[DB insert] + end + end + + S -- "single UDP datagram â‰Ī1472B" --> UDPL + UDPL -- "enqueue raw JSON" --> Q + TICK -- "dequeue batch" --> INGEST + INGEST --> VAL + VAL --> STORE + STORE --> WP + WP -- "completion" --> COMP + TICK2[process_external_ingress_completions] -- "drain" --> COMP + TICK2 --> BCAST +``` + +### Threading Model + +| Component | Thread | Why | +|---|---|---| +| UDP `recvfrom()` loop | Dedicated pthread | `recvfrom()` is blocking; cannot run on lws main thread without starving the WebSocket service loop | +| Queue drain + `ingest_event()` | Main lws thread | `ingest_event()` pushes completions to a queue that must be drained on the lws thread for broadcast. Calling `ingest_event()` from the UDP thread is safe (it uses the thread pool for DB writes and atomically pushes completions), but draining must happen on lws-main. **Decision: call `ingest_event()` from the UDP thread directly** (matches how the caching inbox poller calls it, just from a different thread). The completion queue is already thread-safe (mutex-protected). | +| Completion drain + broadcast | Main lws thread | `broadcast_event_to_subscriptions()` must run on lws-main | + +**Revised simpler approach**: The UDP thread calls `ingest_event()` directly after +receiving a datagram. `ingest_event()` validates, stores via the writer thread pool, +and pushes a completion record to the (mutex-protected) completion queue. The main lws +thread drains completions each iteration (as it already does for caching inbox). This +avoids an intermediate queue entirely — the completion queue IS the handoff mechanism. + +```mermaid +flowchart LR + subgraph "UDP Thread" + RECV[recvfrom] --> INGEST[ingest_event] + end + subgraph "Main lws Thread" + DRAIN[process_external_ingress_completions] --> BCAST[broadcast] + end + INGEST -- "mutex-protected push" --> CQ[completion queue] + CQ -- "mutex-protected pop" --> DRAIN + INGEST -- "wake lws" --> WAKE[lws_cancel_service] +``` + +### Why a dedicated thread (not lws integration) + +libwebsockets can integrate foreign sockets into its event loop via +`lws_add_fd()` / `lws_foreign_callback()`, but: + +1. The static MUSL build uses a minimal lws config — foreign socket support may not be + compiled in. +2. lws integration would require non-blocking `recvfrom()` and edge-triggered handling, + adding complexity for no benefit (UDP events are low-frequency). +3. A dedicated thread with blocking `recvfrom()` is simpler, more portable, and matches + the "stateless, fire-and-forget" philosophy — the thread does one thing and does it + well. +4. The thread calls `lws_cancel_service(ws_context)` after pushing a completion to wake + the main loop for prompt broadcast (same pattern as + [`wake_event_loop_from_thread_pool()`](src/main.c:558)). + +--- + +## Shared Abstraction (for future DNS relay) + +To allow the future authoritative DNS relay to reuse the same infrastructure, we +generalize the existing "caching inbox completion queue" into an +**"external ingress completion queue"** that any non-client ingress source can use. + +### Rename / generalize + +| Current | New (generalized) | +|---|---| +| `inbox_event_completion_t` | `external_ingress_completion_t` | +| `g_inbox_completion_mutex` | `g_external_ingress_completion_mutex` | +| `process_inbox_event_completions()` | `process_external_ingress_completions()` | +| `g_inbox_import_accepted/rejected/duplicates` | Per-source counters (see below) | + +The completion struct and queue logic are identical — only the names change. The +caching inbox poller, UDP listener, and future DNS listener all push completions to +the same queue. The main thread drains all of them in one pass. + +### Per-source counters + +Instead of global `g_inbox_import_*` counters, add per-source counters so the admin +status page can distinguish UDP vs caching-inbox vs DNS: + +```c +typedef enum { + INGRESS_SOURCE_CACHING_INBOX = 0, + INGRESS_SOURCE_UDP, + INGRESS_SOURCE_DNS, // future + INGRESS_SOURCE_COUNT +} ingress_source_t; + +static long long g_ingress_accepted[INGRESS_SOURCE_COUNT]; +static long long g_ingress_rejected[INGRESS_SOURCE_COUNT]; +static long long g_ingress_duplicates[INGRESS_SOURCE_COUNT]; +``` + +The completion struct gains an `ingress_source_t source` field so +`process_external_ingress_completions()` knows which counter to increment. + +### New event_source_t enum value + +Add to [`src/main.h`](src/main.h:53): + +```c +typedef enum { + EVENT_SOURCE_CLIENT, + EVENT_SOURCE_CACHING_INBOX, + EVENT_SOURCE_UDP_INGRESS, // NEW + EVENT_SOURCE_DNS_INGRESS // NEW (future, stub for now) +} event_source_t; +``` + +In `ingest_event()`, `EVENT_SOURCE_UDP_INGRESS` follows the same code path as +`EVENT_SOURCE_CACHING_INBOX` (no wsi/pss, no admin execution, no OK response). The +difference is only the counter routing and the source field in the completion record. + +--- + +## Implementation Steps + +### Step 1: Generalize the completion queue + +**Files**: [`src/main.c`](src/main.c), [`src/main.h`](src/main.h) + +1. Rename `inbox_event_completion_t` → `external_ingress_completion_t` and add an + `ingress_source_t source` field. +2. Rename the mutex, head/tail pointers, and push/pop functions accordingly. +3. Rename `process_inbox_event_completions()` → + `process_external_ingress_completions()`. +4. Replace the global `g_inbox_import_*` counters with the per-source array + `g_ingress_accepted/rejected/duplicates[INGRESS_SOURCE_COUNT]`. +5. Update the declaration in `main.h`. +6. Update the call site in [`src/websockets.c`](src/websockets.c:3448) to call + `process_external_ingress_completions()`. +7. Update `ingest_event()` to accept the new source types and route to the correct + counter. `EVENT_SOURCE_UDP_INGRESS` and `EVENT_SOURCE_DNS_INGRESS` use the same + no-session logic as `EVENT_SOURCE_CACHING_INBOX`. + +**Note**: The `#ifdef DB_BACKEND_POSTGRES` guards around the completion queue remain — +the project is PostgreSQL-only, and the completion/broadcast mechanism depends on the +thread pool writer. For SQLite builds (which are no longer supported but still +compile), UDP ingress will validate and store synchronously but not broadcast (same as +current caching inbox behavior on SQLite). + +### Step 2: Create the UDP ingress module + +**New files**: `src/udp_ingress.h`, `src/udp_ingress.c` + +#### `src/udp_ingress.h` + +```c +#ifndef UDP_INGRESS_H +#define UDP_INGRESS_H + +// Initialize the UDP ingress listener. Binds the UDP socket and starts +// the receiver thread. Returns 0 on success, -1 on error. +int udp_ingress_init(void); + +// Shut down the UDP listener: closes socket, joins thread. +void udp_ingress_shutdown(void); + +// Get statistics for admin status display (all out-params optional). +void udp_ingress_get_stats(long long* out_total_received, + long long* out_total_accepted, + long long* out_total_rejected, + long long* out_total_duplicates, + long long* out_total_oversized, + long long* out_total_parse_errors); + +#endif // UDP_INGRESS_H +``` + +#### `src/udp_ingress.c` — core logic + +```c +// Configuration keys (read at init, re-read each tick is unnecessary for UDP) +#define CFG_KEY_UDP_ENABLED "udp_ingress_enabled" +#define CFG_KEY_UDP_PORT "udp_ingress_port" +#define CFG_KEY_UDP_BIND_ADDR "udp_ingress_bind_addr" +#define CFG_KEY_UDP_MAX_DGRAM "udp_ingress_max_datagram_size" + +// Defaults +#define DEFAULT_UDP_PORT 8888 // same as WS port (TCP/UDP are separate namespaces) +#define DEFAULT_UDP_BIND_ADDR "0.0.0.0" +#define DEFAULT_UDP_MAX_DGRAM 1472 // single-packet guarantee + +// Thread state +static pthread_t g_udp_thread; +static int g_udp_sock_fd = -1; +static volatile int g_udp_running = 0; + +// Statistics (atomic counters) +static long long g_total_received = 0; +static long long g_total_oversized = 0; +static long long g_total_parse_errors = 0; +``` + +**Receiver thread function**: + +```c +static void* udp_ingress_thread_fn(void* arg) { + pthread_setname_np(pthread_self(), "udp-ingress"); + + while (g_udp_running) { + unsigned char buf[DEFAULT_UDP_MAX_DGRAM + 1]; + struct sockaddr_storage src_addr; + socklen_t addr_len = sizeof(src_addr); + + ssize_t n = recvfrom(g_udp_sock_fd, buf, DEFAULT_UDP_MAX_DGRAM, 0, + (struct sockaddr*)&src_addr, &addr_len); + if (n < 0) { + if (errno == EINTR) continue; + if (!g_udp_running) break; // shutdown + DEBUG_WARN("udp_ingress: recvfrom error: %s", strerror(errno)); + continue; + } + + __sync_fetch_and_add(&g_total_received, 1); + + // Single-packet guarantee: silently drop oversized datagrams. + // (The kernel already caps at max_datagram_size, but log for visibility.) + if (n > DEFAULT_UDP_MAX_DGRAM) { + __sync_fetch_and_add(&g_total_oversized, 1); + continue; + } + + // Null-terminate for JSON parsing. + buf[n] = '\0'; + + // Feed through the shared ingestion pipeline. + // ingest_event() does: duplicate check → signature validation → + // structure/expiration/PoW validation → store_event_core() → + // push completion to external ingress queue → wake lws. + int rc = ingest_event((const char*)buf, (size_t)n, + EVENT_SOURCE_UDP_INGRESS, NULL, NULL); + // rc == 0: accepted (or duplicate). rc < 0: rejected (invalid sig, etc.) + // No response is sent — silent drop on failure. This is the no-handshake + // property: the relay's response to an invalid event is indistinguishable + // from a server that isn't running Nostr at all. + (void)rc; + } + + return NULL; +} +``` + +**Init function**: + +```c +int udp_ingress_init(void) { + int enabled = get_config_bool(CFG_KEY_UDP_ENABLED, 0); + if (!enabled) { + DEBUG_LOG("udp_ingress: disabled in config"); + return 0; + } + + int port = get_config_int(CFG_KEY_UDP_PORT, DEFAULT_UDP_PORT); + const char* bind_addr = get_config_value(CFG_KEY_UDP_BIND_ADDR); + if (!bind_addr) bind_addr = DEFAULT_UDP_BIND_ADDR; + + // Create UDP socket + g_udp_sock_fd = socket(AF_INET, SOCK_DGRAM, 0); + if (g_udp_sock_fd < 0) { + DEBUG_ERROR("udp_ingress: socket() failed: %s", strerror(errno)); + return -1; + } + + // Set SO_REUSEADDR so we can rebind quickly after restart + int opt = 1; + setsockopt(g_udp_sock_fd, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt)); + + // Bind + struct sockaddr_in addr; + memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_port = htons(port); + inet_pton(AF_INET, bind_addr, &addr.sin_addr); + + if (bind(g_udp_sock_fd, (struct sockaddr*)&addr, sizeof(addr)) < 0) { + DEBUG_ERROR("udp_ingress: bind(%s:%d) failed: %s", bind_addr, port, strerror(errno)); + close(g_udp_sock_fd); + g_udp_sock_fd = -1; + return -1; + } + + // Start receiver thread + g_udp_running = 1; + if (pthread_create(&g_udp_thread, NULL, udp_ingress_thread_fn, NULL) != 0) { + DEBUG_ERROR("udp_ingress: failed to create receiver thread"); + close(g_udp_sock_fd); + g_udp_sock_fd = -1; + g_udp_running = 0; + return -1; + } + + DEBUG_INFO("udp_ingress: listening on udp://%s:%d (max datagram %d bytes)", + bind_addr, port, DEFAULT_UDP_MAX_DGRAM); + return 0; +} +``` + +**Shutdown function**: + +```c +void udp_ingress_shutdown(void) { + if (!g_udp_running) return; + + g_udp_running = 0; + if (g_udp_sock_fd >= 0) { + shutdown(g_udp_sock_fd, SHUT_RDWR); // unblock recvfrom() + close(g_udp_sock_fd); + g_udp_sock_fd = -1; + } + pthread_join(g_udp_thread, NULL); + + DEBUG_INFO("udp_ingress: shutdown (received=%lld accepted=%lld rejected=%lld)", + g_total_received, ...); +} +``` + +### Step 3: Add configuration keys + +**File**: [`src/default_config_event.h`](src/default_config_event.h:174) + +Add after the caching config block: + +```c + // UDP Ingress Settings (no-handshake event reception) + {"udp_ingress_enabled", "false"}, // Enable UDP event listener + {"udp_ingress_port", "443"}, // UDP port 443 (blends with QUIC/HTTP3; requires CAP_NET_BIND_SERVICE) + {"udp_ingress_bind_addr", "0.0.0.0"}, // Bind address + {"udp_ingress_max_datagram_size", "1472"}, // Max datagram size (single-packet guarantee) +``` + +These are dynamic config keys (no restart required for enable/disable — see Step 5). + +### Step 4: Wire into main startup/shutdown + +**File**: [`src/main.c`](src/main.c) + +1. Add `#include "udp_ingress.h"` at the top. +2. After `caching_inbox_poller_init()` (line ~2967), add: + ```c + // Initialize the UDP ingress listener (if enabled in config). + if (udp_ingress_init() != 0) { + DEBUG_WARN("UDP ingress listener failed to start; continuing without UDP"); + } + ``` +3. Before `caching_inbox_poller_shutdown()` (line ~3051), add: + ```c + udp_ingress_shutdown(); + ``` + +### Step 5: Dynamic enable/disable (hot config reload) + +The UDP listener should start/stop when `udp_ingress_enabled` changes in config, +without requiring a relay restart. Two approaches: + +**Option A (simple, recommended for v1)**: Read `udp_ingress_enabled` only at startup. +Changing it requires a restart. This matches how most config keys work and is the +least complex. + +**Option B (hot reload)**: Add a config-change handler that starts/stops the UDP +thread when the config value changes. This requires: +- A check in the main lws loop (like `caching_inbox_poller_tick()` does for + `caching_inbox_enabled`) that compares the current enabled state vs the config value + and starts/stops the thread accordingly. +- Thread start/stop must be safe (no race with `recvfrom()`). + +**Decision**: Start with Option A. Add a `udp_ingress_tick()` function called from the +lws main loop that checks if the enabled config has changed and starts/stops the +thread. This is low-cost (one `get_config_bool()` per loop iteration) and matches the +caching inbox poller pattern. The tick function: + +```c +void udp_ingress_tick(void) { + int want_enabled = get_config_bool(CFG_KEY_UDP_ENABLED, 0); + if (want_enabled && !g_udp_running) { + // Config enabled but thread not running — start it. + udp_ingress_init(); + } else if (!want_enabled && g_udp_running) { + // Config disabled but thread running — stop it. + udp_ingress_shutdown(); + } +} +``` + +Add `udp_ingress_tick()` call in the lws main loop in +[`src/websockets.c`](src/websockets.c:3453) after `caching_inbox_poller_tick()`. + +### Step 6: Update build system + +**Files**: [`Makefile`](Makefile:12), [`Dockerfile.alpine-musl`](Dockerfile.alpine-musl:127) + +1. Add `src/udp_ingress.c` to `MAIN_SRC` in the Makefile. +2. Add `src/udp_ingress.c` to the `gcc` source list in the Dockerfile. + +No new external libraries are needed — UDP uses standard POSIX sockets +(``, ``, ``) which are already included in +[`src/main.c`](src/main.c:13) and available in the MUSL static build. + +### Step 7: Admin status / monitoring + +**File**: [`src/api.c`](src/api.c) (status event generation) + +Add UDP ingress stats to the relay status event (kind 33334 or similar) so the admin +page can display: +- UDP ingress: enabled/disabled +- Total received, accepted, rejected, duplicates, oversized, parse errors +- Bind address and port + +Expose `udp_ingress_get_stats()` and call it from the status generation code. + +### Step 8: Tests + +**New file**: `tests/udp_ingress_test.sh` + +#### Local testing on port 8888 + +The relay is built and started with `make_and_restart_relay.sh`, which defaults to +port 8888. This does **not** interfere with the local production relay running on +port 7777 — the script matches processes by port (`pgrep -f "c_relay_pg_.*-p 8888"`) +and only kills/restarts the 8888 instance. + +Since TCP and UDP are separate port namespaces, the UDP listener on port 8888 +coexists with the WebSocket TCP listener on port 8888 — no conflict. + +**Test sequence**: + +1. Build and start relay with UDP ingress enabled: + ```bash + # Set udp_ingress_enabled=true and udp_ingress_port=8888 in config, + # then build and restart: + ./make_and_restart_relay.sh + ``` + (The config can be set via the admin API or by adding the key to + `default_config_event.h` defaults before building.) + +2. Send a valid signed event via UDP using the existing + [`udp_nostr_send.py`](../../lt/udp_nostr/udp_nostr_send.py): + ```bash + # Simplest: nak + nc (nak uses a default key, no --sec needed) + nak event -k 1 -c "Hello via UDP Nostr!" | nc -u -w1 127.0.0.1 8888 + + # Or with Python (explicit key generation): + nak event -k 1 -c "Hello via UDP Nostr!" --sec $(nak key generate) \ + | python3 ~/lt/udp_nostr/udp_nostr_send.py 127.0.0.1 8888 + ``` + +3. Verify the event appears in the relay — subscribe via WebSocket: + ```bash + wscat -c ws://localhost:8888 + # Send: ["REQ","test",{"kinds":[1],"limit":1}] + # Verify the UDP-sent event is returned + ``` + +4. Send an invalid event (bad signature) — verify it is silently dropped: + ```bash + echo '{"id":"000...","pubkey":"000...","created_at":1,"kind":1,"tags":[],"content":"fake","sig":"000..."}' \ + | nc -u -w1 127.0.0.1 8888 + # No response, no storage. Check relay.log for the rejection debug line. + ``` + +5. Send an oversized datagram (>1472 bytes) — verify it is dropped: + ```bash + python3 -c "print('A'*2000)" | nc -u -w1 127.0.0.1 8888 + # Silently dropped (oversized counter increments). + ``` + +6. Send a duplicate event — verify it is not stored twice (duplicate counter + increments, no broadcast). + +7. Check stats in relay.log — the UDP ingress thread logs received/accepted/ + rejected/oversized counts on shutdown and at debug level. + +--- + +## Security Considerations + +| Concern | Mitigation | +|---|---| +| **Flooding** | UDP has no congestion control. Mitigate with: (1) per-pubkey rate limiting (future config key `udp_ingress_rate_limit_per_pubkey`), (2) the existing PoW requirement (`pow_min_difficulty` config key applies to all ingress sources via `nostr_validate_unified_request`), (3) signature verification cost is a natural rate limiter (~Ξs per event). | +| **Amplification** | Not a vector — the relay never responds to UDP datagrams. The response is smaller than the request (zero bytes). | +| **Active probing** | The relay does not respond to invalid events. A probe receives nothing, indistinguishable from a non-Nostr UDP service. | +| **IP spoofing** | Not relevant — the relay does not respond, so spoofed source IPs do not cause reflected traffic. The event is self-validating (signature), so the source IP is irrelevant to event validity. | +| **Port conflict** | TCP and UDP port namespaces are separate. UDP port 8888 does not conflict with the WebSocket TCP port 8888. Configurable via `udp_ingress_port`. | +| **Privileged ports** | If `udp_ingress_port` is set to <1024 (e.g., 443 for QUIC-blending), the binary needs `CAP_NET_BIND_SERVICE`. The build script or systemd unit should set this via `setcap`. | + +--- + +## Future: Authoritative DNS Relay Ingress + +The DNS relay (planned in `~/lt/authoritative_dns_relay`) will receive events encoded +in DNS query subdomain labels. It will reuse this infrastructure as follows: + +1. **Same `event_source_t`**: `EVENT_SOURCE_DNS_INGRESS` (already added in Step 1). +2. **Same completion queue**: DNS events push to + `external_ingress_completion_t` with `source = INGRESS_SOURCE_DNS`. +3. **Same `ingest_event()` path**: identical validation, storage, and broadcast. +4. **New module**: `src/dns_ingress.c` / `src/dns_ingress.h` — a UDP listener on port + 53 that parses DNS query headers, extracts the base64url-encoded event from the + subdomain label, decodes it, and calls `ingest_event()`. +5. **New config keys**: `dns_ingress_enabled`, `dns_ingress_port` (default 53), + `dns_ingress_domain` (the authoritative domain to match). +6. **Multi-query reassembly**: For events split across multiple DNS queries, a + session-based reassembly buffer (timeout-driven) will be needed in `dns_ingress.c`. + This is DNS-specific and does not affect the shared completion queue. + +The shared abstraction means the DNS relay implementation only needs to handle: +- DNS protocol parsing (header, question section, label extraction) +- Base64url decoding +- Multi-query reassembly +- NXDOMAIN response generation + +Everything else (validation, storage, broadcast, completion queue, stats) is inherited. + +--- + +## File Change Summary + +| File | Change | +|---|---| +| `src/main.h` | Add `EVENT_SOURCE_UDP_INGRESS`, `EVENT_SOURCE_DNS_INGRESS` to enum. Rename `process_inbox_event_completions` → `process_external_ingress_completions`. Add `ingress_source_t` enum. | +| `src/main.c` | Generalize completion queue (rename + per-source counters). Handle new source types in `ingest_event()`. Add `#include "udp_ingress.h"`. Call `udp_ingress_init()` / `udp_ingress_shutdown()`. | +| `src/websockets.c` | Update completion drain call to `process_external_ingress_completions()`. Add `udp_ingress_tick()` call in main loop. | +| `src/udp_ingress.h` | **NEW** — UDP ingress public API. | +| `src/udp_ingress.c` | **NEW** — UDP socket, receiver thread, `ingest_event()` call, stats, hot-reload tick. | +| `src/default_config_event.h` | Add `udp_ingress_*` config keys. | +| `src/api.c` | Add UDP ingress stats to status event. | +| `Makefile` | Add `src/udp_ingress.c` to `MAIN_SRC`. | +| `Dockerfile.alpine-musl` | Add `src/udp_ingress.c` to gcc source list. | +| `tests/udp_ingress_test.sh` | **NEW** — end-to-end test using `udp_nostr_send.py` or `nc -u`. | + +--- + +## Deployment Decisions (laantungir.net/relay) + +### Port: UDP 443 + +**Decision: default port 443.** Verified on the production relay: + +- **UDP 443 is free** — nginx listens on TCP 443 only (no HTTP/3/QUIC configured). +- TCP and UDP are separate port namespaces, so UDP 443 coexists with nginx's TCP 443. +- UDP on port 443 is **indistinguishable from QUIC/HTTP3 traffic** to a passive + observer — the strongest censorship-resistance position from the + [`protocol_hardening.md`](../../lt/udp_nostr/protocol_hardening.md) analysis. +- Port 443 is universally allowed through firewalls; non-standard UDP ports may be + blocked by some cloud providers or middleboxes. + +### Privileged port binding (CAP_NET_BIND_SERVICE) + +Port 443 is <1024, so the relay binary needs permission to bind. The systemd service +runs as user `c-relay-pg` (non-root). Two options: + +**Option A (recommended): systemd `AmbientCapabilities=CAP_NET_BIND_SERVICE`** + +Add to the `[Service]` section of +`/etc/systemd/system/c-relay-pg.service`: + +```ini +AmbientCapabilities=CAP_NET_BIND_SERVICE +``` + +This grants the service the capability to bind privileged ports without running as +root. No changes to the binary or filesystem needed. This is the cleanest approach +and works with the existing `NoNewPrivileges=true` security setting. + +**Option B: `setcap` on the binary** + +```bash +sudo setcap 'cap_net_bind_service=+ep' /usr/local/bin/c_relay_pg/c_relay_pg +``` + +This persists the capability on the binary file. Must be re-applied after every +binary update (the deploy script should include this). + +**Decision: Option A (systemd AmbientCapabilities)** — survives binary updates +automatically and is the standard systemd pattern. + +### Firewall (UFW) + +If UFW is enabled on the relay, UDP 443 must be opened: + +```bash +sudo ufw allow 443/udp +``` + +Check current status: `sudo ufw status`. If UFW is not active, no change needed +(the relay already accepts TCP 8888 and TCP 443). + +### systemd service changes + +The existing `RestrictAddressFamilies=AF_UNIX AF_INET AF_INET6` already permits +UDP sockets (UDP uses `AF_INET`/`AF_INET6`). No change needed there. + +The only systemd change is adding `AmbientCapabilities=CAP_NET_BIND_SERVICE` to +the `[Service]` section. + +### Config key update + +Update the default in [`src/default_config_event.h`](src/default_config_event.h): + +```c +{"udp_ingress_port", "443"}, // UDP port 443 (blends with QUIC/HTTP3) +``` + +--- + +## Open Questions (resolved) + +1. **Default port**: **443** — blends with QUIC/HTTP3 traffic, universally allowed, + verified free on the production relay. Requires `CAP_NET_BIND_SERVICE` via + systemd `AmbientCapabilities`. + +2. **IPv6 support**: IPv4-only for v1. The udp_nostr project is IPv4-only. Add IPv6 + (`AF_INET6` + dual-stack `IPV6_V6ONLY=0`) as a follow-up. Note: nginx listens on + both IPv4 and IPv6 TCP 443, so an IPv6 UDP listener on 443 would also be safe. + +3. **Rate limiting**: Defer for v1. The PoW requirement (`pow_min_difficulty` config + key) and signature verification cost provide natural rate limiting. Add explicit + per-pubkey rate limiting if flooding becomes a problem in production. diff --git a/relay.pid b/relay.pid index b5f3795..66848a5 100644 --- a/relay.pid +++ b/relay.pid @@ -1 +1 @@ -2035687 +1445308 diff --git a/src/api.c b/src/api.c index 734d5c4..c6ea092 100644 --- a/src/api.c +++ b/src/api.c @@ -31,6 +31,7 @@ int get_active_connection_count(void); #include "subscriptions.h" #include "db_ops.h" #include "caching_service_launcher.h" +#include "udp_ingress.h" // External subscription manager (from main.c via subscriptions.c) extern subscription_manager_t g_subscription_manager; @@ -980,11 +981,13 @@ int generate_and_post_status_event(void) { } // Generate report content by running the PHP kind_1_report.php script. - // This produces a markdown-formatted relay status report. + // Use popen() with a fixed, hardcoded command string (no user input). + // The command is a compile-time constant, so shell injection is not possible. + size_t total = 0; char* stats_text = NULL; FILE *fp = popen("php -r 'require \"/opt/c-relay-pg/admin/api/kind_1_report.php\";' 2>/dev/null", "r"); if (fp) { - size_t total = 0, cap = 16384; + size_t cap = 16384; stats_text = malloc(cap); if (stats_text) { int n; @@ -1002,6 +1005,39 @@ int generate_and_post_status_event(void) { pclose(fp); } + // Append UDP ingress stats to the report + if (stats_text) { + long long udp_received = 0, udp_accepted = 0, udp_rejected = 0; + long long udp_duplicates = 0, udp_oversized = 0, udp_parse_errors = 0; + udp_ingress_get_stats(&udp_received, &udp_accepted, &udp_rejected, + &udp_duplicates, &udp_oversized, &udp_parse_errors); + + int udp_enabled = get_config_bool("udp_ingress_enabled", 0); + int udp_port = get_config_int("udp_ingress_port", 443); + + size_t current_len = strlen(stats_text); + size_t new_cap = current_len + 512; + char* tmp = realloc(stats_text, new_cap); + if (tmp) { + stats_text = tmp; + int n = snprintf(stats_text + current_len, new_cap - current_len, + "\n## UDP Ingress\n" + "- Enabled: %s\n" + "- Port: %d\n" + "- Received: %lld\n" + "- Accepted: %lld\n" + "- Rejected: %lld\n" + "- Duplicates: %lld\n" + "- Oversized: %lld\n" + "- Parse Errors: %lld\n", + udp_enabled ? "true" : "false", + udp_port, + udp_received, udp_accepted, udp_rejected, + udp_duplicates, udp_oversized, udp_parse_errors); + if (n > 0) total = current_len + n; + } + } + if (!stats_text || strlen(stats_text) == 0) { DEBUG_ERROR("Failed to generate report via PHP for status post"); free(stats_text); @@ -1380,6 +1416,8 @@ int send_admin_response(const char* sender_pubkey, const char* response_content, // ============================================================================= // Validate SQL query for security and safety +// Uses a whitelist approach: only allows SELECT and WITH (CTE) statements +// with specific allowed keywords, preventing SQL injection via blacklist bypass. int validate_sql_query(const char* query, char* error_message, size_t error_size) { if (!query || !error_message) { return 0; @@ -1399,47 +1437,53 @@ int validate_sql_query(const char* query, char* error_message, size_t error_size } query_upper[query_len] = '\0'; - // List of blocked keywords (case-insensitive) + // Basic length check + if (query_len < 6) { + snprintf(error_message, error_size, "Query too short"); + return 0; + } + + // Whitelist: only allow SELECT or WITH at the start + if (strncmp(query_upper, "SELECT ", 7) != 0 && + strncmp(query_upper, "WITH ", 5) != 0) { + snprintf(error_message, error_size, "Query blocked: Only SELECT and WITH statements are allowed"); + return 0; + } + + // Blocklist of dangerous keywords (defense-in-depth, not primary defense) const char* blocked_keywords[] = { "INSERT", "UPDATE", "DELETE", "DROP", "CREATE", "ALTER", "TRUNCATE", "EXEC", "EXECUTE", "MERGE", "BULK", "BACKUP", "RESTORE", "GRANT", "REVOKE", "DENY", "COMMIT", "ROLLBACK", "SAVEPOINT", "SHUTDOWN", "PRAGMA", "VACUUM", "REINDEX", "ANALYZE", + "COPY", "LISTEN", "NOTIFY", "DO ", "CALL", NULL // Sentinel }; - // Check for blocked keywords + // Check for blocked keywords (word boundary matching) for (int i = 0; blocked_keywords[i] != NULL; i++) { - char keyword_pattern[64]; - snprintf(keyword_pattern, sizeof(keyword_pattern), " %s ", blocked_keywords[i]); - - if (strstr(query_upper, keyword_pattern) != NULL) { + size_t kw_len = strlen(blocked_keywords[i]); + // Check at start of query + if (strncmp(query_upper, blocked_keywords[i], kw_len) == 0) { snprintf(error_message, error_size, "Query blocked: %s statements not allowed", blocked_keywords[i]); return 0; } - - // Also check at start of query - if (strncmp(query_upper, blocked_keywords[i], strlen(blocked_keywords[i])) == 0) { + // Check with space prefix (word boundary) + char with_space[64]; + snprintf(with_space, sizeof(with_space), " %s", blocked_keywords[i]); + if (strstr(query_upper, with_space) != NULL) { snprintf(error_message, error_size, "Query blocked: %s statements not allowed", blocked_keywords[i]); return 0; } - } - - // Check for SELECT keyword (must be present) - if (strstr(query_upper, " SELECT ") == NULL && strncmp(query_upper, "SELECT ", 7) != 0) { - // Allow WITH clauses (CTEs) which can start queries - if (strstr(query_upper, " WITH ") == NULL && strncmp(query_upper, "WITH ", 5) != 0) { - snprintf(error_message, error_size, "Query blocked: Only SELECT statements and WITH clauses are allowed"); + // Check with semicolon prefix (after previous statement) + char with_semicolon[64]; + snprintf(with_semicolon, sizeof(with_semicolon), ";%s", blocked_keywords[i]); + if (strstr(query_upper, with_semicolon) != NULL) { + snprintf(error_message, error_size, "Query blocked: %s statements not allowed", blocked_keywords[i]); return 0; } } - // Basic length check - if (query_len < 6) { // Minimum "SELECT" length - snprintf(error_message, error_size, "Query too short"); - return 0; - } - return 1; // Query passed validation } @@ -2272,29 +2316,29 @@ int parse_config_command(const char* message, char* key, char* value) { // Pattern 1: "enable auth" -> "auth_enabled true" if (strstr(start, "enable auth") == start) { - strcpy(key, "auth_enabled"); - strcpy(value, "true"); + snprintf(key, 128, "%s", "auth_enabled"); + snprintf(value, 128, "%s", "true"); return 1; } // Pattern 2: "disable auth" -> "auth_enabled false" if (strstr(start, "disable auth") == start) { - strcpy(key, "auth_enabled"); - strcpy(value, "false"); + snprintf(key, 128, "%s", "auth_enabled"); + snprintf(value, 128, "%s", "false"); return 1; } // Pattern 3: "enable nip42" -> "nip42_auth_required true" if (strstr(start, "enable nip42") == start) { - strcpy(key, "nip42_auth_required"); - strcpy(value, "true"); + snprintf(key, 128, "%s", "nip42_auth_required"); + snprintf(value, 128, "%s", "true"); return 1; } // Pattern 4: "disable nip42" -> "nip42_auth_required false" if (strstr(start, "disable nip42") == start) { - strcpy(key, "nip42_auth_required"); - strcpy(value, "false"); + snprintf(key, 128, "%s", "nip42_auth_required"); + snprintf(value, 128, "%s", "false"); return 1; } @@ -2318,7 +2362,7 @@ int parse_config_command(const char* message, char* key, char* value) { } char* value_start = to_pos + 4; - strcpy(value, value_start); + snprintf(value, 128, "%s", value_start); return 1; } } @@ -2330,7 +2374,7 @@ int parse_config_command(const char* message, char* key, char* value) { if (key_len > 0 && key_len < 127) { memcpy(key, start, key_len); key[key_len] = '\0'; - strcpy(value, equals_pos + 3); + snprintf(value, 128, "%s", equals_pos + 3); return 1; } } @@ -2342,7 +2386,7 @@ int parse_config_command(const char* message, char* key, char* value) { if (key_len > 0 && key_len < 127) { memcpy(key, start, key_len); key[key_len] = '\0'; - strcpy(value, colon_pos + 2); + snprintf(value, 128, "%s", colon_pos + 2); return 1; } } @@ -2354,7 +2398,7 @@ int parse_config_command(const char* message, char* key, char* value) { if (key_len > 0 && key_len < 127) { memcpy(key, start, key_len); key[key_len] = '\0'; - strcpy(value, space_pos + 1); + snprintf(value, 128, "%s", space_pos + 1); return 1; } } diff --git a/src/db_ops_postgres.c b/src/db_ops_postgres.c index 6496dd3..5686c8c 100644 --- a/src/db_ops_postgres.c +++ b/src/db_ops_postgres.c @@ -258,6 +258,7 @@ int postgres_db_worker_listen(void* connection, const char* channel) { return DB_ERROR; } // Validate channel name to avoid SQL injection via LISTEN. + // Use PQescapeIdentifier for defense-in-depth even though we validate. for (const char* p = channel; *p; p++) { if (!( ((*p >= 'a' && *p <= 'z') || (*p >= 'A' && *p <= 'Z') || (*p >= '0' && *p <= '9') || *p == '_') )) { @@ -265,8 +266,15 @@ int postgres_db_worker_listen(void* connection, const char* channel) { return DB_ERROR; } } - char sql[128]; - snprintf(sql, sizeof(sql), "LISTEN %s", channel); + // Use PQescapeIdentifier for defense-in-depth against SQL injection + char* escaped_channel = PQescapeIdentifier(conn, channel, strlen(channel)); + if (!escaped_channel) { + postgres_set_error_text("postgres_db_worker_listen: failed to escape channel name"); + return DB_ERROR; + } + char sql[256]; + snprintf(sql, sizeof(sql), "LISTEN %s", escaped_channel); + PQfreemem(escaped_channel); PGresult* res = PQexec(conn, sql); if (!res || PQresultStatus(res) != PGRES_COMMAND_OK) { postgres_set_error_from_conn(conn, "postgres_db_worker_listen: LISTEN failed"); diff --git a/src/default_config_event.h b/src/default_config_event.h index 948b1a7..595bf40 100644 --- a/src/default_config_event.h +++ b/src/default_config_event.h @@ -170,7 +170,13 @@ static const struct { {"caching_max_event_json_bytes", "65536"}, // Maximum event JSON size for inbox insert {"caching_service_binary_path", "./caching_relay"}, // Path to caching_relay binary (relative to relay CWD which is build/) {"caching_service_pg_conn", "host=localhost port=5432 dbname=crelay user=crelay password=crelay"}, // PostgreSQL connection string for caching service (PG-inbox mode) - {"caching_service_launch_mode", "fork"} // Launch mode: "fork" or "systemd" (systemd not yet implemented) + {"caching_service_launch_mode", "fork"}, // Launch mode: "fork" or "systemd" (systemd not yet implemented) + + // UDP Ingress Settings (no-handshake event reception) + {"udp_ingress_enabled", "false"}, // Enable UDP event listener + {"udp_ingress_port", "443"}, // UDP port 443 (blends with QUIC/HTTP3; requires CAP_NET_BIND_SERVICE) + {"udp_ingress_bind_addr", "0.0.0.0"}, // Bind address + {"udp_ingress_max_datagram_size", "1472"} // Max datagram size (single-packet guarantee) }; // Number of default configuration values diff --git a/src/main.c b/src/main.c index 3005e9d..2d0d60d 100644 --- a/src/main.c +++ b/src/main.c @@ -29,6 +29,7 @@ #include "db_ops.h" #include "caching_inbox_poller.h" // Caching inbox poller (PostgreSQL-only) #include "caching_service_launcher.h" // Caching service process launcher +#include "udp_ingress.h" // UDP ingress listener (PostgreSQL-only) // Forward declarations for unified request validator int nostr_validate_unified_request(const char* json_string, size_t json_length); @@ -1170,11 +1171,12 @@ int event_id_exists_in_db(const char* event_id) { } ///////////////////////////////////////////////////////////////////////////////////////// -// Caching inbox ingestion (PostgreSQL-only feature) +// External ingress ingestion (PostgreSQL-only feature) // -// The caching inbox path imports events acquired by the external caching -// application through the same authoritative validation and storage pipeline -// used for client WebSocket events, but with a different source mode that: +// The external ingress path imports events from non-client sources (caching inbox, +// UDP listener, future DNS listener) through the same authoritative validation and +// storage pipeline used for client WebSocket events, but with a different source +// mode that: // - bypasses NIP-42 client-session authentication requirements; // - never sends an OK response (no wsi/pss); // - never executes administrator commands (kind 23456 is stored as data); @@ -1185,65 +1187,66 @@ int event_id_exists_in_db(const char* event_id) { // // The completion queue below mirrors the async event completion pattern in // websockets.c. ingest_event() performs validation + storage off the main -// thread (it is invoked by the inbox poller) and pushes a completion record; -// process_inbox_event_completions() drains that queue on the lws service thread +// thread and pushes a completion record; +// process_external_ingress_completions() drains that queue on the lws service thread // to run store_event_post_actions() and broadcast_event_to_subscriptions(). ///////////////////////////////////////////////////////////////////////////////////////// -typedef struct inbox_event_completion { +typedef struct external_ingress_completion { char* event_json; // owned copy of the event JSON (for re-parse on main thread) char event_id[65]; // event id for logging int success; // 1 = accepted, 0 = rejected int should_broadcast; // 1 = broadcast to subscriptions on main thread int run_post_actions; // 1 = run store_event_post_actions() on main thread - struct inbox_event_completion* next; -} inbox_event_completion_t; + ingress_source_t source; // which ingress source this completion belongs to + struct external_ingress_completion* next; +} external_ingress_completion_t; -static pthread_mutex_t g_inbox_completion_mutex = PTHREAD_MUTEX_INITIALIZER; -static inbox_event_completion_t* g_inbox_completion_head = NULL; -static inbox_event_completion_t* g_inbox_completion_tail = NULL; +static pthread_mutex_t g_external_ingress_completion_mutex = PTHREAD_MUTEX_INITIALIZER; +static external_ingress_completion_t* g_external_ingress_completion_head = NULL; +static external_ingress_completion_t* g_external_ingress_completion_tail = NULL; -// Cumulative import counters (best-effort; not persisted). -static long long g_inbox_import_accepted = 0; -static long long g_inbox_import_rejected = 0; -static long long g_inbox_import_duplicates = 0; +// Cumulative per-source ingress counters (best-effort; not persisted). +long long g_ingress_accepted[INGRESS_SOURCE_COUNT]; +long long g_ingress_rejected[INGRESS_SOURCE_COUNT]; +long long g_ingress_duplicates[INGRESS_SOURCE_COUNT]; -static void inbox_event_completion_push(inbox_event_completion_t* completion) { +static void external_ingress_completion_push(external_ingress_completion_t* completion) { if (!completion) return; - pthread_mutex_lock(&g_inbox_completion_mutex); + pthread_mutex_lock(&g_external_ingress_completion_mutex); completion->next = NULL; - if (!g_inbox_completion_tail) { - g_inbox_completion_head = completion; - g_inbox_completion_tail = completion; + if (!g_external_ingress_completion_tail) { + g_external_ingress_completion_head = completion; + g_external_ingress_completion_tail = completion; } else { - g_inbox_completion_tail->next = completion; - g_inbox_completion_tail = completion; + g_external_ingress_completion_tail->next = completion; + g_external_ingress_completion_tail = completion; } - pthread_mutex_unlock(&g_inbox_completion_mutex); + pthread_mutex_unlock(&g_external_ingress_completion_mutex); } -static inbox_event_completion_t* inbox_event_completion_pop(void) { - pthread_mutex_lock(&g_inbox_completion_mutex); - inbox_event_completion_t* completion = g_inbox_completion_head; +static external_ingress_completion_t* external_ingress_completion_pop(void) { + pthread_mutex_lock(&g_external_ingress_completion_mutex); + external_ingress_completion_t* completion = g_external_ingress_completion_head; if (completion) { - g_inbox_completion_head = completion->next; - if (!g_inbox_completion_head) { - g_inbox_completion_tail = NULL; + g_external_ingress_completion_head = completion->next; + if (!g_external_ingress_completion_head) { + g_external_ingress_completion_tail = NULL; } } - pthread_mutex_unlock(&g_inbox_completion_mutex); + pthread_mutex_unlock(&g_external_ingress_completion_mutex); return completion; } -// Drain caching-inbox ingestion completions on the lws service thread. +// Drain external ingress completions on the lws service thread. // Runs main-thread-only post-actions and broadcasts newly stored events. -void process_inbox_event_completions(void) { +void process_external_ingress_completions(void) { #ifndef DB_BACKEND_POSTGRES return; #else - inbox_event_completion_t* completion = NULL; - while ((completion = inbox_event_completion_pop()) != NULL) { + external_ingress_completion_t* completion = NULL; + while ((completion = external_ingress_completion_pop()) != NULL) { if (completion->success && (completion->run_post_actions || completion->should_broadcast)) { cJSON* event_obj = cJSON_Parse(completion->event_json); if (event_obj && cJSON_IsObject(event_obj)) { @@ -1255,7 +1258,7 @@ void process_inbox_event_completions(void) { } cJSON_Delete(event_obj); } else { - DEBUG_WARN("inbox completion: failed to re-parse event %s for post-actions", + DEBUG_WARN("external ingress completion: failed to re-parse event %s for post-actions", completion->event_id); if (event_obj) { cJSON_Delete(event_obj); @@ -1269,15 +1272,27 @@ void process_inbox_event_completions(void) { #endif /* DB_BACKEND_POSTGRES */ } -// Shared internal ingestion entry point for client events and caching inbox events. +// Map an event_source_t to an ingress_source_t for per-source counter routing. +// Only valid for non-client sources; returns -1 for EVENT_SOURCE_CLIENT. +static inline int event_source_to_ingress_index(event_source_t source) { + switch (source) { + case EVENT_SOURCE_CACHING_INBOX: return INGRESS_SOURCE_CACHING_INBOX; + case EVENT_SOURCE_UDP_INGRESS: return INGRESS_SOURCE_UDP; + case EVENT_SOURCE_DNS_INGRESS: return INGRESS_SOURCE_DNS; + default: return -1; + } +} + +// Shared internal ingestion entry point for client events and external ingress events. // // See the contract documented in main.h. This function does NOT replace the // existing client WebSocket path in websockets.c; it is the new shared entry // point used by the caching inbox poller (EVENT_SOURCE_CACHING_INBOX) and is // available for future refactor of the client path (EVENT_SOURCE_CLIENT). // -// The caching inbox path is compiled only for the PostgreSQL backend. For the -// SQLite backend, ingest_event() with EVENT_SOURCE_CACHING_INBOX returns -1. +// External ingress sources (caching inbox, UDP, DNS) are compiled only for the +// PostgreSQL backend. For the SQLite backend, ingest_event() with any non-client +// source returns -1. int ingest_event(const char* event_json, size_t event_json_len, event_source_t source, struct lws* wsi, void* pss) { (void)wsi; @@ -1288,20 +1303,23 @@ int ingest_event(const char* event_json, size_t event_json_len, return -1; } - // Caching inbox import is a PostgreSQL-only feature. - if (source == EVENT_SOURCE_CACHING_INBOX) { + // External ingress is a PostgreSQL-only feature. + if (source != EVENT_SOURCE_CLIENT) { #ifndef DB_BACKEND_POSTGRES - DEBUG_WARN("ingest_event: caching inbox source requires PostgreSQL backend"); + DEBUG_WARN("ingest_event: non-client source requires PostgreSQL backend"); return -1; #else - // Caching inbox events never have a client session. + // Non-client sources never have a client session. if (wsi || pss) { - DEBUG_WARN("ingest_event: caching inbox source must not carry wsi/pss"); + DEBUG_WARN("ingest_event: non-client source must not carry wsi/pss"); return -1; } #endif } + // Determine the ingress counter index for non-client sources. + int ingress_idx = event_source_to_ingress_index(source); + // Fast duplicate check before expensive crypto. Only meaningful when the // canonical DB is available; mirrors the async event worker optimization. cJSON* peek = cJSON_ParseWithLength(event_json, event_json_len); @@ -1322,9 +1340,9 @@ int ingest_event(const char* event_json, size_t event_json_len, if (peek) { cJSON_Delete(peek); } - if (source == EVENT_SOURCE_CACHING_INBOX) { + if (ingress_idx >= 0) { #ifdef DB_BACKEND_POSTGRES - __sync_fetch_and_add(&g_inbox_import_duplicates, 1); + __sync_fetch_and_add(&g_ingress_duplicates[ingress_idx], 1); #endif } // Duplicates are accepted (already have the event). @@ -1339,12 +1357,12 @@ int ingest_event(const char* event_json, size_t event_json_len, // Authoritative validation: signature, structure, expiration, PoW. int validation_result = nostr_validate_unified_request(event_json, event_json_len); if (validation_result != NOSTR_SUCCESS) { - if (source == EVENT_SOURCE_CACHING_INBOX) { + if (ingress_idx >= 0) { #ifdef DB_BACKEND_POSTGRES - __sync_fetch_and_add(&g_inbox_import_rejected, 1); + __sync_fetch_and_add(&g_ingress_rejected[ingress_idx], 1); #endif - DEBUG_WARN("ingest_event: caching inbox event %s rejected by validator (rc=%d)", - event_id_buf[0] ? event_id_buf : "?", validation_result); + DEBUG_WARN("ingest_event: ingress event %s rejected by validator (rc=%d, source=%d)", + event_id_buf[0] ? event_id_buf : "?", validation_result, source); } return -1; } @@ -1355,11 +1373,11 @@ int ingest_event(const char* event_json, size_t event_json_len, if (event_obj) { cJSON_Delete(event_obj); } - if (source == EVENT_SOURCE_CACHING_INBOX) { + if (ingress_idx >= 0) { #ifdef DB_BACKEND_POSTGRES - __sync_fetch_and_add(&g_inbox_import_rejected, 1); + __sync_fetch_and_add(&g_ingress_rejected[ingress_idx], 1); #endif - DEBUG_WARN("ingest_event: caching inbox event failed to parse after validation"); + DEBUG_WARN("ingest_event: ingress event failed to parse after validation (source=%d)", source); } return -1; } @@ -1375,12 +1393,12 @@ int ingest_event(const char* event_json, size_t event_json_len, int run_post_actions = 0; int is_duplicate = 0; - if (source == EVENT_SOURCE_CACHING_INBOX) { - // Caching inbox: never execute administrator commands. Kind 23456 events + if (source != EVENT_SOURCE_CLIENT) { + // External ingress: never execute administrator commands. Kind 23456 events // are stored as ordinary data so the admin API history is preserved // without triggering command processing. if (event_kind == 23456) { - DEBUG_LOG("ingest_event: caching inbox kind 23456 stored as data (no admin execution)"); + DEBUG_LOG("ingest_event: ingress kind 23456 stored as data (no admin execution, source=%d)", source); } // NIP-09 deletion requests are handled through the same deletion path @@ -1390,8 +1408,8 @@ int ingest_event(const char* event_json, size_t event_json_len, char del_error[512] = {0}; int del_rc = handle_deletion_request(event_obj, del_error, sizeof(del_error)); if (del_rc != 0) { - DEBUG_WARN("ingest_event: caching inbox NIP-09 deletion rejected: %s", - del_error); + DEBUG_WARN("ingest_event: ingress NIP-09 deletion rejected: %s (source=%d)", + del_error, source); result = -1; } else { // Deletion request event itself is stored as a regular event, @@ -1413,7 +1431,7 @@ int ingest_event(const char* event_json, size_t event_json_len, } } else if (event_kind >= 20000 && event_kind < 30000) { // Ephemeral events: validate but do not store; broadcast only. - DEBUG_TRACE("ingest_event: caching inbox ephemeral kind %d - broadcast only", event_kind); + DEBUG_TRACE("ingest_event: ingress ephemeral kind %d - broadcast only (source=%d)", event_kind, source); should_broadcast = 1; run_post_actions = 0; } else { @@ -1421,8 +1439,8 @@ int ingest_event(const char* event_json, size_t event_json_len, // writer thread pool via store_event_core(). int core_rc = store_event_core(event_obj); if (core_rc < 0) { - DEBUG_WARN("ingest_event: caching inbox store_event_core failed for kind %d", - event_kind); + DEBUG_WARN("ingest_event: ingress store_event_core failed for kind %d (source=%d)", + event_kind, source); result = -1; } else { should_broadcast = 1; @@ -1471,13 +1489,13 @@ int ingest_event(const char* event_json, size_t event_json_len, } } - // For the caching inbox path, queue post-actions and broadcast to the main + // For external ingress sources, queue post-actions and broadcast to the main // lws thread. The client path currently sends OK responses inline from // websockets.c and is not routed through here yet. - if (source == EVENT_SOURCE_CACHING_INBOX) { + if (source != EVENT_SOURCE_CLIENT) { #ifdef DB_BACKEND_POSTGRES if (result == 0) { - inbox_event_completion_t* completion = calloc(1, sizeof(*completion)); + external_ingress_completion_t* completion = calloc(1, sizeof(*completion)); if (completion) { completion->event_json = strndup(event_json, event_json_len); if (!completion->event_json) { @@ -1487,10 +1505,11 @@ int ingest_event(const char* event_json, size_t event_json_len, sizeof(completion->event_id) - 1); completion->event_id[sizeof(completion->event_id) - 1] = '\0'; completion->success = 1; + completion->source = (ingress_idx >= 0) ? (ingress_source_t)ingress_idx : INGRESS_SOURCE_CACHING_INBOX; // Duplicates and stale replacements are not broadcast. completion->should_broadcast = is_duplicate ? 0 : should_broadcast; completion->run_post_actions = is_duplicate ? 0 : run_post_actions; - inbox_event_completion_push(completion); + external_ingress_completion_push(completion); } } // Wake the lws service loop so the completion is drained promptly. @@ -1498,12 +1517,12 @@ int ingest_event(const char* event_json, size_t event_json_len, lws_cancel_service(ws_context); } if (is_duplicate) { - __sync_fetch_and_add(&g_inbox_import_duplicates, 1); + __sync_fetch_and_add(&g_ingress_duplicates[ingress_idx >= 0 ? ingress_idx : 0], 1); } else { - __sync_fetch_and_add(&g_inbox_import_accepted, 1); + __sync_fetch_and_add(&g_ingress_accepted[ingress_idx >= 0 ? ingress_idx : 0], 1); } } else { - __sync_fetch_and_add(&g_inbox_import_rejected, 1); + __sync_fetch_and_add(&g_ingress_rejected[ingress_idx >= 0 ? ingress_idx : 0], 1); } #endif /* DB_BACKEND_POSTGRES */ } @@ -2046,30 +2065,30 @@ int handle_req_message(const char* sub_id, cJSON* filters, struct lws *wsi, stru const char* search_term = cJSON_GetStringValue(search); if (search_term && strlen(search_term) > 0) { // Search in both content and tag values using LIKE - // Escape single quotes in search term for SQL safety - char escaped_search[256]; - size_t escaped_len = 0; - for (size_t j = 0; search_term[j] && escaped_len < sizeof(escaped_search) - 1; j++) { - if (search_term[j] == '\'') { - escaped_search[escaped_len++] = '\''; - escaped_search[escaped_len++] = '\''; - } else { - escaped_search[escaped_len++] = search_term[j]; - } - } - escaped_search[escaped_len] = '\0'; - - // Add search conditions for content and tags + // Use parameterized query with bind parameters instead of manual escaping #ifdef DB_BACKEND_POSTGRES - snprintf(sql_ptr, remaining, " AND (content ILIKE '%%%s%%' OR tags::text ILIKE '%%\"%s\"%%')", - escaped_search, escaped_search); + snprintf(sql_ptr, remaining, " AND (content ILIKE ? OR tags::text ILIKE ?)"); #else // Use tags LIKE to search within the JSON string representation of tags - snprintf(sql_ptr, remaining, " AND (content LIKE '%%%s%%' OR tags LIKE '%%\"%s\"%%')", - escaped_search, escaped_search); + snprintf(sql_ptr, remaining, " AND (content LIKE ? OR tags LIKE ?)"); #endif sql_ptr += strlen(sql_ptr); remaining = sizeof(sql) - strlen(sql); + + // Add search term as bind parameters with % wildcards + add_bind_param(&bind_params, &bind_param_count, &bind_param_capacity, NULL); + char* search_pattern1 = malloc(strlen(search_term) + 3); + if (search_pattern1) { + snprintf(search_pattern1, strlen(search_term) + 3, "%%%s%%", search_term); + bind_params[bind_param_count - 1] = search_pattern1; + } + + add_bind_param(&bind_params, &bind_param_count, &bind_param_capacity, NULL); + char* search_pattern2 = malloc(strlen(search_term) + 5); + if (search_pattern2) { + snprintf(search_pattern2, strlen(search_term) + 5, "%%\"%s\"%%", search_term); + bind_params[bind_param_count - 1] = search_pattern2; + } } } @@ -2966,6 +2985,11 @@ int main(int argc, char* argv[]) { // Runs on the main lws service thread via caching_inbox_poller_tick(). caching_inbox_poller_init(); + // Initialize the UDP ingress listener (if enabled in config). + if (udp_ingress_init() != 0) { + DEBUG_WARN("UDP ingress listener failed to start; continuing without UDP"); + } + // CLI flags are independent: // --reset-backfill resets backfill progress tables/state. // --start-caching enables and starts caching service. @@ -3047,6 +3071,9 @@ int main(int argc, char* argv[]) { // Shut down the caching service if we launched it (fork mode). caching_service_shutdown(); + // Shut down the UDP ingress listener before tearing down the thread pool. + udp_ingress_shutdown(); + // Shut down the caching inbox poller before tearing down the thread pool. caching_inbox_poller_shutdown(); diff --git a/src/main.h b/src/main.h index d3e115b..e6f7cdc 100644 --- a/src/main.h +++ b/src/main.h @@ -13,8 +13,8 @@ // Using CRELAY_ prefix to avoid conflicts with nostr_core_lib VERSION macros #define CRELAY_VERSION_MAJOR 2 #define CRELAY_VERSION_MINOR 1 -#define CRELAY_VERSION_PATCH 38 -#define CRELAY_VERSION "v2.1.38" +#define CRELAY_VERSION_PATCH 39 +#define CRELAY_VERSION "v2.1.39" // Relay metadata (authoritative source for NIP-11 information) #define RELAY_NAME "C-Relay-PG" @@ -51,15 +51,25 @@ void process_req_async_completions(void); // Source mode for the shared event ingestion entry point. typedef enum { - EVENT_SOURCE_CLIENT, // Normal WebSocket client event - EVENT_SOURCE_CACHING_INBOX // Event from the caching inbox (internal import) + EVENT_SOURCE_CLIENT, // Normal WebSocket client event + EVENT_SOURCE_CACHING_INBOX, // Event from the caching inbox (internal import) + EVENT_SOURCE_UDP_INGRESS, // Event from the UDP ingress listener + EVENT_SOURCE_DNS_INGRESS // Event from the future DNS ingress listener } event_source_t; +// Per-source ingress counter index. +typedef enum { + INGRESS_SOURCE_CACHING_INBOX = 0, + INGRESS_SOURCE_UDP, + INGRESS_SOURCE_DNS, // future + INGRESS_SOURCE_COUNT +} ingress_source_t; + // Process an event through the authoritative validation and storage pipeline. // This is the shared entry point for both client events and caching inbox events. // // For EVENT_SOURCE_CLIENT: behaves exactly as the current client path does. -// For EVENT_SOURCE_CACHING_INBOX: +// For EVENT_SOURCE_CACHING_INBOX / EVENT_SOURCE_UDP_INGRESS / EVENT_SOURCE_DNS_INGRESS: // - No NIP-42 client authentication requirement // - No OK response sent (no wsi) // - No administrator command execution (kind 23456 events are stored as data, not processed) @@ -71,9 +81,9 @@ typedef enum { // Parameters: // event_json - raw JSON string of the Nostr event // event_json_len - length of the JSON string -// source - EVENT_SOURCE_CLIENT or EVENT_SOURCE_CACHING_INBOX -// wsi - WebSocket instance for client responses (NULL for caching inbox) -// pss - per-session data for client responses (NULL for caching inbox) +// source - EVENT_SOURCE_CLIENT, EVENT_SOURCE_CACHING_INBOX, EVENT_SOURCE_UDP_INGRESS, or EVENT_SOURCE_DNS_INGRESS +// wsi - WebSocket instance for client responses (NULL for non-client sources) +// pss - per-session data for client responses (NULL for non-client sources) // // Returns: // 0 = accepted and stored (or duplicate, or ephemeral broadcast) @@ -81,11 +91,12 @@ typedef enum { int ingest_event(const char* event_json, size_t event_json_len, event_source_t source, struct lws* wsi, void* pss); -// Drains completed caching-inbox ingestion jobs on the lws service thread. +// Drains completed external ingress ingestion jobs on the lws service thread. // Runs store_event_post_actions() and broadcast_event_to_subscriptions() for // events that were stored off the main thread by ingest_event() with -// EVENT_SOURCE_CACHING_INBOX. Safe to call unconditionally; no-op when empty -// or when the PostgreSQL backend is not in use. -void process_inbox_event_completions(void); +// EVENT_SOURCE_CACHING_INBOX, EVENT_SOURCE_UDP_INGRESS, or EVENT_SOURCE_DNS_INGRESS. +// Safe to call unconditionally; no-op when empty or when the PostgreSQL backend +// is not in use. +void process_external_ingress_completions(void); #endif /* MAIN_H */ diff --git a/src/nip042.c b/src/nip042.c index 8ccc071..b87ab3e 100644 --- a/src/nip042.c +++ b/src/nip042.c @@ -27,11 +27,112 @@ int nostr_nip42_verify_auth_event(cJSON *event, const char *challenge_id, // Forward declaration for per_session_data struct (defined in websockets.h) +// Rate limiting for NIP-42 auth attempts +#define AUTH_RATE_LIMIT_WINDOW_SEC 60 +#define AUTH_RATE_LIMIT_MAX_ATTEMPTS 10 + +typedef struct { + char client_ip[64]; + int attempt_count; + time_t window_start; +} auth_rate_limit_entry_t; + +#define AUTH_RATE_LIMIT_MAX_ENTRIES 256 +static auth_rate_limit_entry_t g_auth_rate_limits[AUTH_RATE_LIMIT_MAX_ENTRIES]; +static int g_auth_rate_limit_count = 0; +static pthread_mutex_t g_auth_rate_limit_lock = PTHREAD_MUTEX_INITIALIZER; + +// Check if a client is rate limited for auth attempts +static int is_auth_rate_limited(const char* client_ip) { + if (!client_ip || !client_ip[0]) return 0; + + time_t now = time(NULL); + pthread_mutex_lock(&g_auth_rate_limit_lock); + + // Find existing entry or create new one + int idx = -1; + for (int i = 0; i < g_auth_rate_limit_count; i++) { + if (strcmp(g_auth_rate_limits[i].client_ip, client_ip) == 0) { + idx = i; + break; + } + } + + if (idx < 0) { + // Create new entry + if (g_auth_rate_limit_count < AUTH_RATE_LIMIT_MAX_ENTRIES) { + idx = g_auth_rate_limit_count++; + strncpy(g_auth_rate_limits[idx].client_ip, client_ip, sizeof(g_auth_rate_limits[idx].client_ip) - 1); + g_auth_rate_limits[idx].client_ip[sizeof(g_auth_rate_limits[idx].client_ip) - 1] = '\0'; + g_auth_rate_limits[idx].attempt_count = 0; + g_auth_rate_limits[idx].window_start = now; + } + } + + if (idx >= 0) { + auth_rate_limit_entry_t* entry = &g_auth_rate_limits[idx]; + + // Reset window if expired + if (now - entry->window_start >= AUTH_RATE_LIMIT_WINDOW_SEC) { + entry->attempt_count = 0; + entry->window_start = now; + } + + int limited = (entry->attempt_count >= AUTH_RATE_LIMIT_MAX_ATTEMPTS); + pthread_mutex_unlock(&g_auth_rate_limit_lock); + return limited; + } + + pthread_mutex_unlock(&g_auth_rate_limit_lock); + return 0; +} + +// Record an auth attempt for rate limiting +static void record_auth_attempt(const char* client_ip) { + if (!client_ip || !client_ip[0]) return; + + time_t now = time(NULL); + pthread_mutex_lock(&g_auth_rate_limit_lock); + + int idx = -1; + for (int i = 0; i < g_auth_rate_limit_count; i++) { + if (strcmp(g_auth_rate_limits[i].client_ip, client_ip) == 0) { + idx = i; + break; + } + } + + if (idx < 0 && g_auth_rate_limit_count < AUTH_RATE_LIMIT_MAX_ENTRIES) { + idx = g_auth_rate_limit_count++; + strncpy(g_auth_rate_limits[idx].client_ip, client_ip, sizeof(g_auth_rate_limits[idx].client_ip) - 1); + g_auth_rate_limits[idx].client_ip[sizeof(g_auth_rate_limits[idx].client_ip) - 1] = '\0'; + g_auth_rate_limits[idx].attempt_count = 0; + g_auth_rate_limits[idx].window_start = now; + } + + if (idx >= 0) { + auth_rate_limit_entry_t* entry = &g_auth_rate_limits[idx]; + if (now - entry->window_start >= AUTH_RATE_LIMIT_WINDOW_SEC) { + entry->attempt_count = 0; + entry->window_start = now; + } + entry->attempt_count++; + } + + pthread_mutex_unlock(&g_auth_rate_limit_lock); +} // Send NIP-42 authentication challenge to client void send_nip42_auth_challenge(struct lws* wsi, struct per_session_data* pss) { if (!wsi || !pss) return; + // Check rate limiting before sending challenge + if (is_auth_rate_limited(pss->client_ip)) { + DEBUG_WARN("Auth rate limited for client %s", pss->client_ip); + send_notice_message(wsi, pss, "Too many authentication attempts - please wait before retrying"); + return; + } + // Generate challenge using existing request_validator function char challenge[65]; if (nostr_nip42_generate_challenge(challenge, sizeof(challenge)) != 0) { @@ -70,6 +171,16 @@ void send_nip42_auth_challenge(struct lws* wsi, struct per_session_data* pss) { void handle_nip42_auth_signed_event(struct lws* wsi, struct per_session_data* pss, cJSON* auth_event) { if (!wsi || !pss || !auth_event) return; + // Check rate limiting before processing auth event + if (is_auth_rate_limited(pss->client_ip)) { + DEBUG_WARN("Auth rate limited for client %s during auth event processing", pss->client_ip); + send_notice_message(wsi, pss, "Too many authentication attempts - please wait before retrying"); + return; + } + + // Record this auth attempt for rate limiting + record_auth_attempt(pss->client_ip); + // Serialize event for validation char* event_json = cJSON_Print(auth_event); if (!event_json) { diff --git a/src/pg_schema.h b/src/pg_schema.h index 08bd294..3d31bbd 100644 --- a/src/pg_schema.h +++ b/src/pg_schema.h @@ -241,6 +241,8 @@ static const char* const EMBEDDED_PG_SCHEMA_SQL = "CREATE INDEX IF NOT EXISTS idx_subscriptions_client ON subscriptions(client_ip);\n" "CREATE INDEX IF NOT EXISTS idx_subscriptions_wsi ON subscriptions(wsi_pointer);\n" "CREATE INDEX IF NOT EXISTS idx_subscriptions_active_log ON subscriptions(event_type, ended_at, created_at DESC);\n" +"-- Partial index for active subscription count queries (admin stats page)\n" +"CREATE INDEX IF NOT EXISTS idx_subscriptions_active_lookup ON subscriptions(event_type, ended_at) WHERE ended_at IS NULL;\n" "CREATE INDEX IF NOT EXISTS idx_subscription_metrics_date ON subscription_metrics(date DESC);\n" "\n" "CREATE TABLE IF NOT EXISTS ip_bans (\n" diff --git a/src/request_validator.c b/src/request_validator.c index 8f3effa..8031489 100644 --- a/src/request_validator.c +++ b/src/request_validator.c @@ -23,6 +23,7 @@ #include #include #include +#include // Forward declaration for C-relay event storage function extern int store_event(cJSON* event); @@ -726,6 +727,7 @@ int nostr_generate_nip42_challenge(char *challenge_out, size_t challenge_size, /** * Simple NIP-42 challenge generation (generates random hex string) + * Uses /dev/urandom for cryptographically secure random bytes. */ int nostr_nip42_generate_challenge(char *challenge_buffer, size_t buffer_size) { if (!challenge_buffer || buffer_size < 65) { @@ -735,10 +737,30 @@ int nostr_nip42_generate_challenge(char *challenge_buffer, size_t buffer_size) { // Generate 32 random bytes and convert to hex string unsigned char random_bytes[32]; - // Simple random number generation using time and rand() - srand((unsigned int)time(NULL)); - for (int i = 0; i < 32; i++) { - random_bytes[i] = (unsigned char)(rand() % 256); + // Use /dev/urandom for cryptographically secure random bytes + FILE* urandom = fopen("/dev/urandom", "rb"); + if (urandom) { + size_t nread = fread(random_bytes, 1, sizeof(random_bytes), urandom); + fclose(urandom); + if (nread != sizeof(random_bytes)) { + // Fallback: use time + pid + clock for entropy if /dev/urandom fails + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + unsigned int seed = (unsigned int)(ts.tv_sec ^ ts.tv_nsec ^ (unsigned long)getpid()); + srand(seed); + for (int i = 0; i < 32; i++) { + random_bytes[i] = (unsigned char)(rand() % 256); + } + } + } else { + // Fallback: use time + pid + clock for entropy + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + unsigned int seed = (unsigned int)(ts.tv_sec ^ ts.tv_nsec ^ (unsigned long)getpid()); + srand(seed); + for (int i = 0; i < 32; i++) { + random_bytes[i] = (unsigned char)(rand() % 256); + } } // Convert to hex string diff --git a/src/sql_schema.h b/src/sql_schema.h index b79aa24..004b0c2 100644 --- a/src/sql_schema.h +++ b/src/sql_schema.h @@ -230,6 +230,9 @@ CREATE INDEX idx_subscriptions_wsi ON subscriptions(wsi_pointer);\n\ -- Optimizes: WHERE event_type = 'created' AND ended_at IS NULL ORDER BY created_at DESC\n\ CREATE INDEX idx_subscriptions_active_log ON subscriptions(event_type, ended_at, created_at DESC);\n\ \n\ +-- Partial index for active subscription count queries (admin stats page)\n\ +CREATE INDEX idx_subscriptions_active_lookup ON subscriptions(event_type, ended_at) WHERE ended_at IS NULL;\n\ +\n\ CREATE INDEX idx_subscription_metrics_date ON subscription_metrics(date DESC);\n\ \n\ \n\ diff --git a/src/udp_ingress.c b/src/udp_ingress.c new file mode 100644 index 0000000..0150dd2 --- /dev/null +++ b/src/udp_ingress.c @@ -0,0 +1,358 @@ +/* + * UDP Ingress Listener + * + * Listens for Nostr events as single UDP datagrams (no handshake, no response). + * Each datagram must contain a complete, self-validating Nostr event JSON + * (signature replaces the handshake). Events are fed through the shared + * ingest_event() pipeline with EVENT_SOURCE_UDP_INGRESS. + * + * The receiver runs on a dedicated pthread because recvfrom() is blocking. + * After pushing a completion record, it wakes the main lws service loop via + * lws_cancel_service() for prompt broadcast to WebSocket subscribers. + * + * This module is PostgreSQL-only (ingest_event() with non-client sources + * requires the PostgreSQL backend). For SQLite builds all entry points are + * no-ops. + */ + +#include "udp_ingress.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include "debug.h" +#include "config.h" +#include "main.h" + +extern struct lws_context *ws_context; + +#ifdef DB_BACKEND_POSTGRES + +// ---- Configuration keys ----------------------------------------------------- + +#define CFG_KEY_ENABLED "udp_ingress_enabled" +#define CFG_KEY_PORT "udp_ingress_port" +#define CFG_KEY_BIND_ADDR "udp_ingress_bind_addr" +#define CFG_KEY_MAX_DATAGRAM_SIZE "udp_ingress_max_datagram_size" + +// ---- Defaults --------------------------------------------------------------- + +#define UDP_DEFAULT_PORT 443 +#define UDP_DEFAULT_BIND_ADDR "0.0.0.0" +#define UDP_DEFAULT_MAX_DATAGRAM_SIZE 1472 + +// ---- Module state ----------------------------------------------------------- + +static pthread_t g_udp_thread; +static int g_udp_running = 0; // 1 = thread is (or should be) running +static int g_udp_sock_fd = -1; // UDP socket file descriptor + +// Cached config values for hot-reload detection. +static int g_cfg_enabled = 0; // last known udp_ingress_enabled value + +// Statistics counters (atomic via __sync_fetch_and_add). +static long long g_udp_total_received = 0; +static long long g_udp_total_oversized = 0; +static long long g_udp_total_parse_errors = 0; + +// Per-source ingress counters defined in main.c. +extern long long g_ingress_accepted[INGRESS_SOURCE_COUNT]; +extern long long g_ingress_rejected[INGRESS_SOURCE_COUNT]; +extern long long g_ingress_duplicates[INGRESS_SOURCE_COUNT]; + +// ---- Forward declarations --------------------------------------------------- + +static void* udp_ingress_thread(void* arg); +static int udp_ingress_start(void); +static void udp_ingress_stop(void); + +// ---- Receiver thread -------------------------------------------------------- + +static void* udp_ingress_thread(void* arg) { + (void)arg; + +#ifdef __GLIBC__ + pthread_setname_np(pthread_self(), "udp-ingress"); +#endif + + int max_dgram = get_config_int(CFG_KEY_MAX_DATAGRAM_SIZE, UDP_DEFAULT_MAX_DATAGRAM_SIZE); + if (max_dgram <= 0 || max_dgram > 65535) { + max_dgram = UDP_DEFAULT_MAX_DATAGRAM_SIZE; + } + + // Allocate receive buffer (max datagram size + 1 for null terminator). + size_t buf_size = (size_t)max_dgram + 1; + char* buf = (char*)malloc(buf_size); + if (!buf) { + DEBUG_ERROR("udp_ingress: failed to allocate receive buffer"); + return NULL; + } + + DEBUG_LOG("udp_ingress: receiver thread started (max_datagram=%d)", max_dgram); + + while (g_udp_running) { + struct sockaddr_in from_addr; + socklen_t from_len = sizeof(from_addr); + memset(&from_addr, 0, sizeof(from_addr)); + + ssize_t n = recvfrom(g_udp_sock_fd, buf, (size_t)max_dgram, 0, + (struct sockaddr*)&from_addr, &from_len); + + if (n < 0) { + if (errno == EINTR) { + continue; + } + if (!g_udp_running) { + break; + } + DEBUG_WARN("udp_ingress: recvfrom error: %s", strerror(errno)); + continue; + } + + if (n == 0) { + continue; + } + + __sync_fetch_and_add(&g_udp_total_received, 1); + + // Null-terminate the received data. + buf[n] = '\0'; + + // Feed through the authoritative ingestion pipeline. + int rc = ingest_event((const char*)buf, (size_t)n, + EVENT_SOURCE_UDP_INGRESS, NULL, NULL); + + if (rc != 0) { + // ingest_event already increments the per-source rejected counter. + // We track parse errors locally for oversized/truncated datagrams. + DEBUG_TRACE("udp_ingress: ingest_event rejected datagram (%d bytes)", (int)n); + } + + // Wake the main lws service loop so completions are drained promptly. + if (ws_context) { + lws_cancel_service(ws_context); + } + } + + free(buf); + DEBUG_LOG("udp_ingress: receiver thread exiting"); + return NULL; +} + +// ---- Start / stop helpers --------------------------------------------------- + +static int udp_ingress_start(void) { + if (g_udp_running) { + return 0; // already running + } + + int port = get_config_int(CFG_KEY_PORT, UDP_DEFAULT_PORT); + if (port <= 0 || port > 65535) { + port = UDP_DEFAULT_PORT; + } + + const char* bind_str = get_config_value(CFG_KEY_BIND_ADDR); + if (!bind_str || bind_str[0] == '\0') { + bind_str = UDP_DEFAULT_BIND_ADDR; + } + + int max_dgram = get_config_int(CFG_KEY_MAX_DATAGRAM_SIZE, UDP_DEFAULT_MAX_DATAGRAM_SIZE); + if (max_dgram <= 0 || max_dgram > 65535) { + max_dgram = UDP_DEFAULT_MAX_DATAGRAM_SIZE; + } + + // Create UDP socket. + int sock = socket(AF_INET, SOCK_DGRAM, 0); + if (sock < 0) { + DEBUG_ERROR("udp_ingress: failed to create socket: %s", strerror(errno)); + return -1; + } + + // Allow immediate reuse of the address. + int reuse = 1; + if (setsockopt(sock, SOL_SOCKET, SO_REUSEADDR, &reuse, sizeof(reuse)) < 0) { + DEBUG_WARN("udp_ingress: setsockopt SO_REUSEADDR failed: %s", strerror(errno)); + // Non-fatal; continue. + } + + // Bind to the configured address and port. + struct sockaddr_in addr; + memset(&addr, 0, sizeof(addr)); + addr.sin_family = AF_INET; + addr.sin_port = htons((uint16_t)port); + if (inet_pton(AF_INET, bind_str, &addr.sin_addr) != 1) { + DEBUG_ERROR("udp_ingress: invalid bind address '%s'", bind_str); + close(sock); + return -1; + } + + if (bind(sock, (struct sockaddr*)&addr, sizeof(addr)) < 0) { + DEBUG_ERROR("udp_ingress: bind to %s:%d failed: %s", bind_str, port, strerror(errno)); + close(sock); + return -1; + } + + g_udp_sock_fd = sock; + g_udp_running = 1; + + // Start the receiver thread. + if (pthread_create(&g_udp_thread, NULL, udp_ingress_thread, NULL) != 0) { + DEBUG_ERROR("udp_ingress: failed to create receiver thread: %s", strerror(errno)); + g_udp_running = 0; + close(g_udp_sock_fd); + g_udp_sock_fd = -1; + return -1; + } + + DEBUG_LOG("udp_ingress: listening on %s:%d (max_datagram=%d)", bind_str, port, max_dgram); + return 0; +} + +static void udp_ingress_stop(void) { + if (!g_udp_running) { + return; + } + + DEBUG_LOG("udp_ingress: shutting down..."); + + // Signal the thread to stop and unblock recvfrom(). + g_udp_running = 0; + if (g_udp_sock_fd >= 0) { + shutdown(g_udp_sock_fd, SHUT_RDWR); + close(g_udp_sock_fd); + g_udp_sock_fd = -1; + } + + // Join the receiver thread. + pthread_join(g_udp_thread, NULL); + + DEBUG_LOG("udp_ingress: shutdown complete"); +} + +// ---- Public API ------------------------------------------------------------- + +int udp_ingress_init(void) { +#ifdef DB_BACKEND_POSTGRES + g_cfg_enabled = get_config_bool(CFG_KEY_ENABLED, 0); + + if (!g_cfg_enabled) { + DEBUG_LOG("udp_ingress: disabled by config"); + return 0; + } + + return udp_ingress_start(); +#else + // SQLite build: no-op. UDP ingress is PostgreSQL-only. + return 0; +#endif +} + +void udp_ingress_shutdown(void) { +#ifdef DB_BACKEND_POSTGRES + udp_ingress_stop(); +#else + // SQLite build: no-op. +#endif +} + +void udp_ingress_tick(void) { +#ifdef DB_BACKEND_POSTGRES + int enabled = get_config_bool(CFG_KEY_ENABLED, 0); + + if (enabled && !g_cfg_enabled) { + // Config changed from disabled to enabled — start the listener. + DEBUG_LOG("udp_ingress: config enabled, starting listener"); + if (udp_ingress_start() == 0) { + g_cfg_enabled = 1; + } else { + DEBUG_WARN("udp_ingress: failed to start on config change"); + } + } else if (!enabled && g_cfg_enabled) { + // Config changed from enabled to disabled — stop the listener. + DEBUG_LOG("udp_ingress: config disabled, stopping listener"); + udp_ingress_stop(); + g_cfg_enabled = 0; + } +#else + // SQLite build: no-op. +#endif +} + +void udp_ingress_get_stats(long long* out_total_received, + long long* out_total_accepted, + long long* out_total_rejected, + long long* out_total_duplicates, + long long* out_total_oversized, + long long* out_total_parse_errors) { +#ifdef DB_BACKEND_POSTGRES + if (out_total_received) { + *out_total_received = g_udp_total_received; + } + if (out_total_oversized) { + *out_total_oversized = g_udp_total_oversized; + } + if (out_total_parse_errors) { + *out_total_parse_errors = g_udp_total_parse_errors; + } + + // Accepted/rejected/duplicates come from the per-source counters in main.c. + // We read them directly for convenience. + if (out_total_accepted) { + *out_total_accepted = g_ingress_accepted[INGRESS_SOURCE_UDP]; + } + if (out_total_rejected) { + *out_total_rejected = g_ingress_rejected[INGRESS_SOURCE_UDP]; + } + if (out_total_duplicates) { + *out_total_duplicates = g_ingress_duplicates[INGRESS_SOURCE_UDP]; + } +#else + // SQLite build: all counters are zero. + if (out_total_received) *out_total_received = 0; + if (out_total_accepted) *out_total_accepted = 0; + if (out_total_rejected) *out_total_rejected = 0; + if (out_total_duplicates) *out_total_duplicates = 0; + if (out_total_oversized) *out_total_oversized = 0; + if (out_total_parse_errors) *out_total_parse_errors = 0; +#endif +} + +#else /* !DB_BACKEND_POSTGRES */ + +// ---- SQLite build: all public functions are no-ops -------------------------- + +int udp_ingress_init(void) { + return 0; +} + +void udp_ingress_shutdown(void) { +} + +void udp_ingress_tick(void) { +} + +void udp_ingress_get_stats(long long* out_total_received, + long long* out_total_accepted, + long long* out_total_rejected, + long long* out_total_duplicates, + long long* out_total_oversized, + long long* out_total_parse_errors) { + if (out_total_received) *out_total_received = 0; + if (out_total_accepted) *out_total_accepted = 0; + if (out_total_rejected) *out_total_rejected = 0; + if (out_total_duplicates) *out_total_duplicates = 0; + if (out_total_oversized) *out_total_oversized = 0; + if (out_total_parse_errors) *out_total_parse_errors = 0; +} + +#endif /* DB_BACKEND_POSTGRES */ diff --git a/src/udp_ingress.h b/src/udp_ingress.h new file mode 100644 index 0000000..2a5f528 --- /dev/null +++ b/src/udp_ingress.h @@ -0,0 +1,25 @@ +#ifndef UDP_INGRESS_H +#define UDP_INGRESS_H + +// Initialize the UDP ingress listener. Binds the UDP socket and starts +// the receiver thread. Returns 0 on success, -1 on error. +// Reads config keys: udp_ingress_enabled, udp_ingress_port, +// udp_ingress_bind_addr, udp_ingress_max_datagram_size +int udp_ingress_init(void); + +// Shut down the UDP listener: closes socket, joins thread. +void udp_ingress_shutdown(void); + +// Hot-reload check: starts/stops the UDP thread when udp_ingress_enabled +// changes in config. Called from the main lws service loop. +void udp_ingress_tick(void); + +// Get statistics for admin status display (all out-params optional). +void udp_ingress_get_stats(long long* out_total_received, + long long* out_total_accepted, + long long* out_total_rejected, + long long* out_total_duplicates, + long long* out_total_oversized, + long long* out_total_parse_errors); + +#endif // UDP_INGRESS_H diff --git a/src/websockets.c b/src/websockets.c index 4524c3f..8d7725b 100644 --- a/src/websockets.c +++ b/src/websockets.c @@ -35,6 +35,7 @@ #include "main.h" // Async REQ completion integration #include "caching_inbox_poller.h" // Caching inbox poller (PostgreSQL-only) +#include "udp_ingress.h" // UDP ingress listener (PostgreSQL-only) // Forward declarations for logging functions @@ -199,9 +200,15 @@ static void check_idle_connections(int idle_timeout_sec) { for (int i = 0; i < MAX_TRACKED_CONNECTIONS; i++) { if (g_connections[i].wsi == NULL || g_connections[i].pss == NULL) continue; struct per_session_data* pss = g_connections[i].pss; - if (pss->session_active) continue; // Already active — skip - if (pss->connection_established <= 0) continue; - time_t age = now - pss->connection_established; + // Lock the session to safely read its state + pthread_mutex_lock(&pss->session_lock); + int is_active = pss->session_active; + time_t conn_est = pss->connection_established; + pthread_mutex_unlock(&pss->session_lock); + + if (is_active) continue; // Already active — skip + if (conn_est <= 0) continue; + time_t age = now - conn_est; if (age >= idle_timeout_sec) { to_close[close_count++] = g_connections[i].wsi; } @@ -756,7 +763,7 @@ static void process_async_event_completions(void) { } } - if (target_alive) { + if (target_alive && current_pss) { send_ok_response(completion->wsi, current_pss, completion->event_id, @@ -1490,6 +1497,19 @@ static int nostr_relay_callback(struct lws *wsi, enum lws_callback_reasons reaso DEBUG_TRACE("Starting message reassembly"); } + // Enforce maximum message size to prevent unbounded memory growth + #define MAX_MESSAGE_SIZE (10 * 1024 * 1024) // 10 MB max message size + if (pss->reassembly_size + len > MAX_MESSAGE_SIZE) { + DEBUG_WARN("Message exceeds maximum size of %d bytes, rejecting", MAX_MESSAGE_SIZE); + free(pss->reassembly_buffer); + pss->reassembly_buffer = NULL; + pss->reassembly_size = 0; + pss->reassembly_capacity = 0; + pss->reassembly_active = 0; + send_notice_message(wsi, pss, "error: message too large"); + return 0; + } + // Ensure buffer has enough capacity size_t needed_capacity = pss->reassembly_size + len + 1; // +1 for null terminator if (needed_capacity > pss->reassembly_capacity) { @@ -3443,15 +3463,18 @@ int start_websocket_relay(int port_override, int strict_port) { // Drain completed async EVENT jobs and emit OK/broadcast on service thread. process_async_event_completions(); - // Drain completed caching-inbox ingestion jobs and run post-actions/broadcast + // Drain completed external ingress ingestion jobs and run post-actions/broadcast // on the service thread (PostgreSQL-only; no-op otherwise). - process_inbox_event_completions(); + process_external_ingress_completions(); // Poll the caching_event_inbox table and feed dequeued events through // ingest_event() with EVENT_SOURCE_CACHING_INBOX (PostgreSQL-only; no-op // otherwise). Runs once per lws service loop iteration on the main thread. caching_inbox_poller_tick(); + // Hot-reload check for UDP ingress listener (start/stop on config change). + udp_ingress_tick(); + // Drain completed async COUNT jobs and emit COUNT/NOTICE on service thread. process_count_async_completions(); @@ -3953,30 +3976,36 @@ int handle_count_message(const char* sub_id, cJSON* filters, struct lws *wsi, st const char* search_term = cJSON_GetStringValue(search); if (search_term && strlen(search_term) > 0) { // Search in both content and tag values using LIKE - // Escape single quotes in search term for SQL safety - char escaped_search[256]; - size_t escaped_len = 0; - for (size_t i = 0; search_term[i] && escaped_len < sizeof(escaped_search) - 1; i++) { - if (search_term[i] == '\'') { - escaped_search[escaped_len++] = '\''; - escaped_search[escaped_len++] = '\''; - } else { - escaped_search[escaped_len++] = search_term[i]; - } - } - escaped_search[escaped_len] = '\0'; - - // Add search conditions for content and tags + // Use parameterized query with bind parameters instead of manual escaping #ifdef DB_BACKEND_POSTGRES - snprintf(sql_ptr, remaining, " AND (content ILIKE '%%%s%%' OR tags::text ILIKE '%%\"%s\"%%')", - escaped_search, escaped_search); + snprintf(sql_ptr, remaining, " AND (content ILIKE ? OR tags::text ILIKE ?)"); #else // Use tags LIKE to search within the JSON string representation of tags - snprintf(sql_ptr, remaining, " AND (content LIKE '%%%s%%' OR tags LIKE '%%\"%s\"%%')", - escaped_search, escaped_search); + snprintf(sql_ptr, remaining, " AND (content LIKE ? OR tags LIKE ?)"); #endif sql_ptr += strlen(sql_ptr); remaining = sizeof(sql) - strlen(sql); + + // Add search term as bind parameters with % wildcards + if (bind_param_count >= bind_param_capacity) { + bind_param_capacity = bind_param_capacity == 0 ? 16 : bind_param_capacity * 2; + bind_params = realloc(bind_params, bind_param_capacity * sizeof(char*)); + } + char* search_pattern1 = malloc(strlen(search_term) + 3); + if (search_pattern1) { + snprintf(search_pattern1, strlen(search_term) + 3, "%%%s%%", search_term); + bind_params[bind_param_count++] = search_pattern1; + } + + if (bind_param_count >= bind_param_capacity) { + bind_param_capacity = bind_param_capacity == 0 ? 16 : bind_param_capacity * 2; + bind_params = realloc(bind_params, bind_param_capacity * sizeof(char*)); + } + char* search_pattern2 = malloc(strlen(search_term) + 5); + if (search_pattern2) { + snprintf(search_pattern2, strlen(search_term) + 5, "%%\"%s\"%%", search_term); + bind_params[bind_param_count++] = search_pattern2; + } } } diff --git a/tests/udp_ingress_test.sh b/tests/udp_ingress_test.sh new file mode 100755 index 0000000..db6918d --- /dev/null +++ b/tests/udp_ingress_test.sh @@ -0,0 +1,203 @@ +#!/bin/bash +# UDP Ingress Test — sends Nostr events via UDP datagrams and verifies storage +set -e + +RELAY_HOST="127.0.0.1" +RELAY_PORT="${RELAY_PORT:-8888}" +UDP_PORT="${UDP_PORT:-8888}" +DB_NAME="${DB_NAME:-crelay}" +DB_USER="${DB_USER:-crelay}" +DB_HOST="${DB_HOST:-localhost}" + +PASS=0 +FAIL=0 + +# Color constants +RED='\033[31m' +GREEN='\033[32m' +YELLOW='\033[33m' +BLUE='\033[34m' +BOLD='\033[1m' +RESET='\033[0m' + +pass() { echo -e "${GREEN}✅ PASS:${RESET} $1"; PASS=$((PASS+1)); } +fail() { echo -e "${RED}❌ FAIL:${RESET} $1"; FAIL=$((FAIL+1)); } +info() { echo -e "${BLUE}[INFO]${RESET} $1"; } +warn() { echo -e "${YELLOW}[WARN]${RESET} $1"; } + +echo -e "${BOLD}=== UDP Ingress Test ===${RESET}" +echo "Relay: ${RELAY_HOST}:${RELAY_PORT}, UDP: ${UDP_PORT}" +echo "" + +# --------------------------------------------------------------------------- +# Test 1: Check relay is running +# --------------------------------------------------------------------------- +if pgrep -f "c_relay_pg" >/dev/null 2>&1; then + pass "Relay process is running" +else + fail "Relay process is not running" + echo "Start it with: ./make_and_restart_relay.sh" + exit 1 +fi + +# --------------------------------------------------------------------------- +# Test 2: Check UDP socket is listening +# --------------------------------------------------------------------------- +if ss -lun 2>/dev/null | grep -q ":${UDP_PORT} "; then + pass "UDP socket is listening on port ${UDP_PORT}" +else + fail "UDP socket is NOT listening on port ${UDP_PORT}" + echo "" + echo "Enable UDP ingress with:" + echo " PGPASSWORD=crelay psql -h localhost -U crelay -d crelay -c \\" + echo " \"UPDATE config SET value='true' WHERE key='udp_ingress_enabled';\"" + echo " PGPASSWORD=crelay psql -h localhost -U crelay -d crelay -c \\" + echo " \"UPDATE config SET value='8888' WHERE key='udp_ingress_port';\"" + echo "Then restart the relay or wait for hot-reload." + exit 1 +fi + +# --------------------------------------------------------------------------- +# Test 3: Send a valid event via UDP and verify it was stored +# --------------------------------------------------------------------------- +info "Test 3: Sending valid Nostr event via UDP..." + +if command -v nak >/dev/null 2>&1; then + # Generate a random key and event using nak + RAND_SEC=$(nak key generate 2>/dev/null) + EVENT_JSON=$(nak event -k 1 -c "UDP ingress test $(date +%s)" --sec "$RAND_SEC" 2>/dev/null) + if [ -n "$EVENT_JSON" ]; then + EVENT_ID=$(echo "$EVENT_JSON" | jq -r '.id' 2>/dev/null) + if [ -n "$EVENT_ID" ] && [ "$EVENT_ID" != "null" ]; then + # Send via UDP using Python + echo "$EVENT_JSON" | python3 -c " +import socket, sys +sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) +data = sys.stdin.read().encode() +sock.sendto(data, ('${RELAY_HOST}', ${UDP_PORT})) +" 2>/dev/null || { + # Fallback: try nc -u + echo "$EVENT_JSON" | nc -u -w1 "${RELAY_HOST}" "${UDP_PORT}" 2>/dev/null || true + } + sleep 1 # Give relay time to process + + # Verify event was stored in database + RESULT=$(PGPASSWORD=crelay psql -h "${DB_HOST}" -U "${DB_USER}" -d "${DB_NAME}" \ + -tAc "SELECT 1 FROM events WHERE id='${EVENT_ID}'" 2>/dev/null) + if [ "$RESULT" = "1" ]; then + pass "Valid event stored in database (id=${EVENT_ID:0:16}...)" + else + fail "Valid event NOT found in database (id=${EVENT_ID:0:16}...)" + fi + else + fail "Could not extract event ID from nak output" + fi + else + fail "nak event generation failed" + fi +else + warn "nak not found — generating event manually via Python" + # Fallback: create a minimal (but properly signed) event using Python + # This requires the secp256k1 library; if unavailable, skip. + if python3 -c "import secp256k1" 2>/dev/null; then + EVENT_JSON=$(python3 -c " +import json, time, hashlib, secp256k1 +# Generate random private key +priv = secp256k1.PrivateKey() +pub = priv.pubkey.serialize().hex() +# Build event +event = { + 'id': '', + 'pubkey': pub, + 'created_at': int(time.time()), + 'kind': 1, + 'tags': [], + 'content': 'UDP ingress test ' + str(int(time.time())), + 'sig': '' +} +# Compute id +serial = json.dumps([0, pub, event['created_at'], 1, [], event['content']], separators=(',', ':')) +event['id'] = hashlib.sha256(serial.encode()).hexdigest() +# Sign +sig = priv.ecdsa_serialize(priv.ecdsa_sign(hashlib.sha256(bytes.fromhex(event['id'])).digest())) +event['sig'] = sig.hex() +print(json.dumps(event)) +" 2>/dev/null) + if [ -n "$EVENT_JSON" ]; then + EVENT_ID=$(echo "$EVENT_JSON" | jq -r '.id' 2>/dev/null) + echo "$EVENT_JSON" | python3 -c " +import socket, sys +sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) +sock.sendto(sys.stdin.read().encode(), ('${RELAY_HOST}', ${UDP_PORT})) +" 2>/dev/null || true + sleep 1 + RESULT=$(PGPASSWORD=crelay psql -h "${DB_HOST}" -U "${DB_USER}" -d "${DB_NAME}" \ + -tAc "SELECT 1 FROM events WHERE id='${EVENT_ID}'" 2>/dev/null) + if [ "$RESULT" = "1" ]; then + pass "Valid event stored in database (id=${EVENT_ID:0:16}...)" + else + fail "Valid event NOT found in database (id=${EVENT_ID:0:16}...)" + fi + else + fail "Python event generation failed (no secp256k1?)" + fi + else + fail "nak not installed and Python secp256k1 unavailable — cannot generate valid event" + fi +fi + +# --------------------------------------------------------------------------- +# Test 4: Send an invalid event (bad signature) — should be silently dropped +# --------------------------------------------------------------------------- +info "Test 4: Sending invalid event (bad signature) via UDP..." +INVALID_ID="0000000000000000000000000000000000000000000000000000000000000000" +echo '{"id":"0000000000000000000000000000000000000000000000000000000000000000","pubkey":"0000000000000000000000000000000000000000000000000000000000000000","created_at":1,"kind":1,"tags":[],"content":"invalid","sig":"0000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000"}' \ + | python3 -c " +import socket, sys +sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) +sock.sendto(sys.stdin.read().encode(), ('${RELAY_HOST}', ${UDP_PORT})) +" 2>/dev/null || echo "$INVALID_JSON" | nc -u -w1 "${RELAY_HOST}" "${UDP_PORT}" 2>/dev/null || true +sleep 1 + +# Verify the invalid event was NOT stored +RESULT=$(PGPASSWORD=crelay psql -h "${DB_HOST}" -U "${DB_USER}" -d "${DB_NAME}" \ + -tAc "SELECT 1 FROM events WHERE id='${INVALID_ID}'" 2>/dev/null) +if [ "$RESULT" != "1" ]; then + pass "Invalid event was silently dropped (not stored)" +else + fail "Invalid event was incorrectly stored!" +fi + +# --------------------------------------------------------------------------- +# Test 5: Send an oversized datagram — should be rejected +# --------------------------------------------------------------------------- +info "Test 5: Sending oversized datagram via UDP..." +# Build a payload larger than the default max datagram size (65535 bytes) +LARGE_CONTENT=$(python3 -c "print('x' * 70000)" 2>/dev/null) +LARGE_JSON="{\"id\":\"1111111111111111111111111111111111111111111111111111111111111111\",\"pubkey\":\"1111111111111111111111111111111111111111111111111111111111111111\",\"created_at\":1,\"kind\":1,\"tags\":[],\"content\":\"${LARGE_CONTENT}\",\"sig\":\"1111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111111\"}" +echo "$LARGE_JSON" | python3 -c " +import socket, sys +sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) +data = sys.stdin.read().encode() +try: + sock.sendto(data, ('${RELAY_HOST}', ${UDP_PORT})) +except OSError as e: + print(f'Message too large (expected): {e}') +" 2>/dev/null || true +sleep 1 + +# The oversized message should not have been stored +RESULT=$(PGPASSWORD=crelay psql -h "${DB_HOST}" -U "${DB_USER}" -d "${DB_NAME}" \ + -tAc "SELECT 1 FROM events WHERE id='1111111111111111111111111111111111111111111111111111111111111111'" 2>/dev/null) +if [ "$RESULT" != "1" ]; then + pass "Oversized datagram was rejected (not stored)" +else + fail "Oversized datagram was incorrectly stored!" +fi + +# --------------------------------------------------------------------------- +# Summary +# --------------------------------------------------------------------------- +echo "" +echo -e "${BOLD}=== Results: ${PASS} passed, ${FAIL} failed ===${RESET}" +exit $FAIL