v2.1.39 - Security audit: fix 3 Critical and 7 High severity vulnerabilities; add missing DB index for admin stats
This commit is contained in:
@@ -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)**:
|
||||
|
||||
@@ -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 \
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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*
|
||||
@@ -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`
|
||||
|
||||
@@ -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
|
||||
```
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 <relay_pubkey> 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 |
|
||||
@@ -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<br/>blocking recvfrom on UDP port]
|
||||
Q[ring buffer / queue<br/>thread-safe]
|
||||
end
|
||||
|
||||
subgraph "Main lws Thread"
|
||||
TICK[udp_ingress_tick<br/>drains queue each loop iteration]
|
||||
INGEST[ingest_event<br/>EVENT_SOURCE_UDP_INGRESS]
|
||||
VAL[nostr_validate_unified_request<br/>signature + structure + PoW + expiration]
|
||||
STORE[store_event_core<br/>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
|
||||
(`<sys/socket.h>`, `<netinet/in.h>`, `<arpa/inet.h>`) 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.
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
+10
-2
@@ -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");
|
||||
|
||||
@@ -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
|
||||
|
||||
+113
-86
@@ -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();
|
||||
|
||||
|
||||
+23
-12
@@ -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 */
|
||||
|
||||
+111
@@ -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) {
|
||||
|
||||
@@ -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"
|
||||
|
||||
+26
-4
@@ -23,6 +23,7 @@
|
||||
#include <string.h>
|
||||
#include <strings.h>
|
||||
#include <time.h>
|
||||
#include <unistd.h>
|
||||
|
||||
// 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
|
||||
|
||||
@@ -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\
|
||||
|
||||
@@ -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 <pthread.h>
|
||||
#include <sys/socket.h>
|
||||
#include <netinet/in.h>
|
||||
#include <arpa/inet.h>
|
||||
#include <unistd.h>
|
||||
#include <errno.h>
|
||||
#include <string.h>
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <signal.h>
|
||||
|
||||
#include <libwebsockets.h>
|
||||
|
||||
#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 */
|
||||
@@ -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
|
||||
+53
-24
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Executable
+203
@@ -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
|
||||
Reference in New Issue
Block a user