diff --git a/relay.pid b/relay.pid index 7206daf..0c20008 100644 --- a/relay.pid +++ b/relay.pid @@ -1 +1 @@ -4055658 +66905 diff --git a/src/db_ops.c b/src/db_ops.c index 9e42415..c71ddc6 100644 --- a/src/db_ops.c +++ b/src/db_ops.c @@ -70,6 +70,7 @@ int db_update_subscription_events_sent(const char* sub_id, int events_sent) { return postgres_db_update_subscription_events_sent(sub_id, events_sent); } int db_cleanup_orphaned_subscriptions(void) { return postgres_db_cleanup_orphaned_subscriptions(); } +int db_prune_ended_subscriptions(int max_batches) { return postgres_db_prune_ended_subscriptions(max_batches); } int db_get_event_pubkey(const char* event_id, char* pubkey_out, size_t pubkey_out_size) { return postgres_db_get_event_pubkey(event_id, pubkey_out, pubkey_out_size); diff --git a/src/db_ops.h b/src/db_ops.h index 0046d01..cc70920 100644 --- a/src/db_ops.h +++ b/src/db_ops.h @@ -65,6 +65,8 @@ int db_log_subscription_closed(const char* sub_id, const char* client_ip); int db_log_subscription_disconnected(const char* client_ip); int db_update_subscription_events_sent(const char* sub_id, int events_sent); int db_cleanup_orphaned_subscriptions(void); +// Retention: prune ended subscription rows in batches (never touches active rows) +int db_prune_ended_subscriptions(int max_batches); // NIP-09 event deletion helpers int db_get_event_pubkey(const char* event_id, char* pubkey_out, size_t pubkey_out_size); diff --git a/src/db_ops_postgres.c b/src/db_ops_postgres.c index bd2b20c..2c99e76 100644 --- a/src/db_ops_postgres.c +++ b/src/db_ops_postgres.c @@ -761,6 +761,61 @@ int postgres_db_update_subscription_events_sent(const char* sub_id, int events_s #endif } +// Prune ended subscription rows (event_type != 'created' OR ended_at IS NOT NULL). +// Subscriptions are ephemeral: clients always reconnect with fresh REQs, so +// ended rows have no value and only create scan burden for the admin stats +// queries. Runs in batches to avoid long locks; returns total rows deleted +// (>= 0) or DB_ERROR. Never touches active rows (event_type = 'created' AND +// ended_at IS NULL). +int postgres_db_prune_ended_subscriptions(int max_batches) { +#if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ) + PGconn* conn = postgres_db_active_connection(); + if (!conn || PQstatus(conn) != CONNECTION_OK) return DB_ERROR; + if (max_batches <= 0) return 0; + + int total_deleted = 0; + for (int batch = 0; batch < max_batches; batch++) { + PGresult* res = PQexecParams(conn, + "WITH del AS ( " + " DELETE FROM subscriptions " + " WHERE id IN ( " + " SELECT id FROM subscriptions " + " WHERE event_type <> 'created' OR ended_at IS NOT NULL " + " LIMIT 50000 " + " ) " + " RETURNING 1 " + ") " + "SELECT count(*) FROM del", + 0, NULL, NULL, NULL, NULL, 0); + if (!res || PQresultStatus(res) != PGRES_TUPLES_OK) { + if (res) { + postgres_set_error_text(PQresultErrorMessage(res)); + PQclear(res); + } else { + postgres_set_error_from_conn(conn, "postgres_db_prune_ended_subscriptions failed"); + } + return (total_deleted > 0) ? total_deleted : DB_ERROR; + } + + int deleted = 0; + if (PQntuples(res) > 0 && !PQgetisnull(res, 0, 0)) { + deleted = atoi(PQgetvalue(res, 0, 0)); + } + PQclear(res); + + total_deleted += deleted; + if (deleted == 0) { + break; // nothing left to prune + } + } + + return total_deleted; +#else + (void)max_batches; + return DB_ERROR; +#endif +} + int postgres_db_cleanup_orphaned_subscriptions(void) { #if defined(DB_BACKEND_POSTGRES) && defined(HAVE_LIBPQ) PGconn* conn = postgres_db_active_connection(); diff --git a/src/db_ops_postgres.h b/src/db_ops_postgres.h index 46acc0f..f4756e1 100644 --- a/src/db_ops_postgres.h +++ b/src/db_ops_postgres.h @@ -42,6 +42,7 @@ int postgres_db_log_subscription_closed(const char* sub_id, const char* client_i int postgres_db_log_subscription_disconnected(const char* client_ip); int postgres_db_update_subscription_events_sent(const char* sub_id, int events_sent); int postgres_db_cleanup_orphaned_subscriptions(void); +int postgres_db_prune_ended_subscriptions(int max_batches); int postgres_db_get_event_pubkey(const char* event_id, char* pubkey_out, size_t pubkey_out_size); int postgres_db_delete_event_by_id(const char* event_id, const char* requester_pubkey); diff --git a/src/main.h b/src/main.h index 7f770cf..9add3a6 100644 --- a/src/main.h +++ b/src/main.h @@ -13,8 +13,8 @@ // Using CRELAY_ prefix to avoid conflicts with nostr_core_lib VERSION macros #define CRELAY_VERSION_MAJOR 2 #define CRELAY_VERSION_MINOR 1 -#define CRELAY_VERSION_PATCH 42 -#define CRELAY_VERSION "v2.1.42" +#define CRELAY_VERSION_PATCH 43 +#define CRELAY_VERSION "v2.1.43" // Relay metadata (authoritative source for NIP-11 information) #define RELAY_NAME "C-Relay-PG" diff --git a/src/subscriptions.c b/src/subscriptions.c index 618f8bc..63aa52a 100644 --- a/src/subscriptions.c +++ b/src/subscriptions.c @@ -1274,6 +1274,18 @@ void cleanup_all_subscriptions_on_startup(void) { } else { DEBUG_LOG("Startup cleanup: no orphaned subscriptions found"); } + + // Retention: prune ended subscription rows left behind by previous runs. + // Subscriptions are ephemeral (clients reconnect with fresh REQs), so + // ended rows have no value and only burden the admin stats queries. + // Capped at 40 batches (2M rows) per startup; the periodic timer + // finishes the rest if the backlog is larger. + int pruned = db_prune_ended_subscriptions(40); + if (pruned > 0) { + DEBUG_LOG("Startup retention: pruned %d ended subscription rows", pruned); + } else if (pruned < 0) { + DEBUG_WARN("Startup retention: prune failed"); + } } diff --git a/src/websockets.c b/src/websockets.c index c36d8d5..03c589f 100644 --- a/src/websockets.c +++ b/src/websockets.c @@ -3505,6 +3505,8 @@ int start_websocket_relay(int port_override, int strict_port) { // Static variable for connection age check timing static time_t last_connection_age_check = 0; + // Static variable for subscription retention pruning (every 10 minutes) + static time_t last_subscription_prune = 0; if (start_async_event_worker() != 0) { DEBUG_WARN("Async event worker failed to start; EVENT path will remain synchronous"); @@ -3591,6 +3593,22 @@ int start_websocket_relay(int port_override, int strict_port) { ip_ban_cleanup(); ip_ban_log_stats(); } + + // Subscription retention (every 10 minutes): prune ended rows. + // Subscriptions are ephemeral — clients reconnect with fresh REQs — + // so ended rows are deleted instead of accumulating forever (this + // table previously grew to 39M rows / 24 GB and crushed the admin + // stats queries). Small cap keeps each run fast; steady state has + // only a few hundred ended rows per interval. + if (current_time - last_subscription_prune >= 600) { + last_subscription_prune = current_time; + int pruned = db_prune_ended_subscriptions(4); + if (pruned > 0) { + DEBUG_LOG("Subscription retention: pruned %d ended rows", pruned); + } else if (pruned < 0) { + DEBUG_WARN("Subscription retention: prune failed"); + } + } } }