661 lines
27 KiB
C
661 lines
27 KiB
C
/*
|
|
* caching_relay - a daemon that caches Nostr events from followed people
|
|
* into a local relay.
|
|
*
|
|
* Architecture: see plans/plan.md. Two nostr_relay_pool_t instances:
|
|
* - upstream_pool: query/subscribe to upstream relays
|
|
* - sink_pool (in cr_sink_t): publish-only to the local relay
|
|
*
|
|
* Build: see Makefile. Statically linked C99 binary, c-relay style.
|
|
*/
|
|
#define _GNU_SOURCE
|
|
#include "main.h"
|
|
#include "debug.h"
|
|
#include "config.h"
|
|
#include "state.h"
|
|
#include "follow_graph.h"
|
|
#include "relay_sink.h"
|
|
#include "live_subscriber.h"
|
|
#include "backfill.h"
|
|
#include "forward_catchup.h"
|
|
#include "relay_discovery.h"
|
|
#include "pg_inbox.h"
|
|
#include "pg_config.h"
|
|
|
|
#include "../nostr_core_lib/nostr_core/nostr_core.h"
|
|
#include "../nostr_core_lib/nostr_core/nostr_log.h"
|
|
|
|
#include <stdio.h>
|
|
#include <stdlib.h>
|
|
#include <string.h>
|
|
#include <signal.h>
|
|
#include <time.h>
|
|
#include <unistd.h>
|
|
#include <getopt.h>
|
|
|
|
/* Forward nostr_core_lib's internal logging to our debug system. */
|
|
static void nostr_log_forwarder(int level, const char *component,
|
|
const char *message, void *user_data) {
|
|
(void)user_data;
|
|
/* Map nostr log levels to our debug levels (they're the same 1-5). */
|
|
if (level >= 5) {
|
|
DEBUG_TRACE("[nostr:%s] %s", component ? component : "?", message ? message : "");
|
|
} else if (level >= 4) {
|
|
DEBUG_LOG("[nostr:%s] %s", component ? component : "?", message ? message : "");
|
|
} else if (level >= 3) {
|
|
DEBUG_INFO("[nostr:%s] %s", component ? component : "?", message ? message : "");
|
|
} else if (level >= 2) {
|
|
DEBUG_WARN("[nostr:%s] %s", component ? component : "?", message ? message : "");
|
|
} else {
|
|
DEBUG_ERROR("[nostr:%s] %s", component ? component : "?", message ? message : "");
|
|
}
|
|
}
|
|
|
|
/* Initialize per-relay backfill progress rows for every followed pubkey.
|
|
* For each author, merges its discovered outbox relays (from relay_map)
|
|
* with the bootstrap (upstream) relays from cfg, then inserts one
|
|
* caching_backfill_relay_progress row per relay (ON CONFLICT DO NOTHING).
|
|
* Safe to call repeatedly - existing rows are left untouched. */
|
|
static void init_relay_progress_for_all(const cr_config_t *cfg,
|
|
const cr_relay_map_t *relay_map,
|
|
const cr_pubkey_set_t *followed) {
|
|
if (!followed || followed->count <= 0) return;
|
|
|
|
/* Bootstrap relay URLs from cfg->upstream_relays. */
|
|
const char **bootstrap = NULL;
|
|
int bootstrap_n = cfg->upstream_count;
|
|
if (bootstrap_n > 0) {
|
|
bootstrap = malloc(bootstrap_n * sizeof(char *));
|
|
if (!bootstrap) return;
|
|
for (int i = 0; i < bootstrap_n; i++)
|
|
bootstrap[i] = cfg->upstream_relays[i];
|
|
}
|
|
|
|
for (int i = 0; i < followed->count; i++) {
|
|
const char *pk = followed->items[i];
|
|
if (!pk || !*pk) continue;
|
|
|
|
/* Collect this author's outbox relays. */
|
|
const cr_outbox_entry_t *oe = NULL;
|
|
int outbox_n = 0;
|
|
if (relay_map) {
|
|
oe = cr_relay_map_get_outbox(relay_map, pk);
|
|
if (oe) outbox_n = oe->relay_count;
|
|
}
|
|
|
|
int total = outbox_n + bootstrap_n;
|
|
if (total <= 0) continue;
|
|
|
|
const char **urls = malloc(total * sizeof(char *));
|
|
if (!urls) continue;
|
|
int n = 0;
|
|
if (oe) {
|
|
for (int j = 0; j < oe->relay_count && n < total; j++)
|
|
urls[n++] = oe->relays[j];
|
|
}
|
|
for (int j = 0; j < bootstrap_n && n < total; j++) {
|
|
int dup = 0;
|
|
for (int k = 0; k < n; k++) {
|
|
if (strcmp(urls[k], bootstrap[j]) == 0) { dup = 1; break; }
|
|
}
|
|
if (!dup) urls[n++] = bootstrap[j];
|
|
}
|
|
|
|
if (n > 0) {
|
|
pg_inbox_init_relay_progress_for_author(pk, urls, n);
|
|
}
|
|
free(urls);
|
|
}
|
|
|
|
free(bootstrap);
|
|
}
|
|
|
|
volatile sig_atomic_t g_shutdown = 0;
|
|
static volatile sig_atomic_t g_reload = 0;
|
|
|
|
static void on_signal(int sig) {
|
|
if (sig == SIGINT || sig == SIGTERM) g_shutdown = 1;
|
|
else if (sig == SIGHUP) g_reload = 1;
|
|
}
|
|
|
|
static void usage(const char *prog) {
|
|
fprintf(stderr,
|
|
"caching_relay %s - cache Nostr events from followed people into a local relay\n"
|
|
"\n"
|
|
"Usage: %s [-c <config.jsonc>] [-p <pg-conn>] [options]\n"
|
|
"\n"
|
|
"Options:\n"
|
|
" -c, --config <file> Path to .jsonc config file (default: ./caching_relay_config.jsonc)\n"
|
|
" -p, --pg-conn <str> PostgreSQL connection string (libpq format). When provided,\n"
|
|
" config is read from the c-relay-pg config table and fetched\n"
|
|
" events are inserted into the caching_event_inbox table instead\n"
|
|
" of being published via WebSocket to a local relay.\n"
|
|
" -d, --debug <level> Log level 0-5 (0=none, 1=error, 2=warn, 3=info, 4=debug, 5=trace). Default 3.\n"
|
|
" -r, --restart Reset state to first-time startup (ignore local relay cache,\n"
|
|
" re-discover all kind-10002 from bootstrap relays, reset backfill)\n"
|
|
" -h, --help Show this help\n"
|
|
"\n"
|
|
"Without -p, the daemon uses the .jsonc config file and publishes to a local relay\n"
|
|
"via WebSocket (legacy mode). With -p, it uses PostgreSQL for config and inbox.\n"
|
|
"The config file is also the persistent state store in legacy mode; the daemon\n"
|
|
"rewrites it as backfill progresses. See plans/plan.md and caching_relay_config.jsonc.\n",
|
|
CR_VERSION, prog);
|
|
}
|
|
|
|
int main(int argc, char **argv) {
|
|
const char *config_path = NULL;
|
|
const char *pg_conn = NULL;
|
|
int log_level = DEBUG_LEVEL_INFO;
|
|
int restart = 0;
|
|
|
|
static struct option longopts[] = {
|
|
{"config", required_argument, 0, 'c'},
|
|
{"pg-conn", required_argument, 0, 'p'},
|
|
{"debug", required_argument, 0, 'd'},
|
|
{"restart", no_argument, 0, 'r'},
|
|
{"help", no_argument, 0, 'h'},
|
|
{0, 0, 0, 0}
|
|
};
|
|
int opt;
|
|
while ((opt = getopt_long(argc, argv, "c:p:d:rh", longopts, NULL)) != -1) {
|
|
switch (opt) {
|
|
case 'c': config_path = optarg; break;
|
|
case 'p': pg_conn = optarg; break;
|
|
case 'd': log_level = atoi(optarg); break;
|
|
case 'r': restart = 1; break;
|
|
case 'h': usage(argv[0]); return 0;
|
|
default: usage(argv[0]); return 1;
|
|
}
|
|
}
|
|
/* In PostgreSQL mode, the .jsonc config file is optional. */
|
|
if (!config_path && !pg_conn) {
|
|
/* Default: look for caching_relay.jsonc in the current directory. */
|
|
config_path = "caching_relay_config.jsonc";
|
|
if (access(config_path, F_OK) != 0) {
|
|
fprintf(stderr, "ERROR: no config specified and default '%s' not found\n\n", config_path);
|
|
usage(argv[0]);
|
|
return 1;
|
|
}
|
|
}
|
|
if (log_level < 0) log_level = 0;
|
|
if (log_level > 5) log_level = 5;
|
|
|
|
debug_init(log_level);
|
|
|
|
/* Enable nostr_core_lib internal logging and forward to our debug system.
|
|
* This shows the actual WebSocket messages (REQ filters, EVENT responses)
|
|
* at debug level 5 (trace). */
|
|
nostr_set_log_callback(nostr_log_forwarder, NULL);
|
|
nostr_set_log_level((nostr_log_level_t)log_level);
|
|
|
|
DEBUG_INFO("caching_relay %s starting (config=%s, pg-conn=%s, loglevel=%d%s)",
|
|
CR_VERSION,
|
|
config_path ? config_path : "(none)",
|
|
pg_conn ? "yes" : "no",
|
|
log_level, restart ? ", RESTART" : "");
|
|
|
|
/* Signals. */
|
|
struct sigaction sa;
|
|
memset(&sa, 0, sizeof(sa));
|
|
sa.sa_handler = on_signal;
|
|
sigaction(SIGINT, &sa, NULL);
|
|
sigaction(SIGTERM, &sa, NULL);
|
|
sigaction(SIGHUP, &sa, NULL);
|
|
/* Ignore SIGPIPE - relay pool handles its own socket errors. */
|
|
signal(SIGPIPE, SIG_IGN);
|
|
|
|
/* PostgreSQL inbox mode: connect and load config from the config table. */
|
|
long config_generation = -1;
|
|
if (pg_conn) {
|
|
if (pg_inbox_init(pg_conn) != 0) {
|
|
DEBUG_ERROR("failed to connect to PostgreSQL");
|
|
return 1;
|
|
}
|
|
config_generation = pg_inbox_get_config_generation();
|
|
if (config_generation < 0) {
|
|
DEBUG_WARN("could not read caching_config_generation (defaulting to 0)");
|
|
config_generation = 0;
|
|
}
|
|
DEBUG_INFO("config generation: %ld", config_generation);
|
|
}
|
|
|
|
/* Load config. */
|
|
cr_config_t cfg;
|
|
if (pg_conn) {
|
|
if (pg_config_load(&cfg) != 0) {
|
|
DEBUG_ERROR("failed to load config from PostgreSQL");
|
|
pg_inbox_shutdown();
|
|
return 1;
|
|
}
|
|
} else {
|
|
if (cr_config_load(&cfg, config_path) != 0) return 1;
|
|
}
|
|
|
|
/* --restart: reset state to first-time startup. */
|
|
if (restart) {
|
|
DEBUG_INFO("RESTART: resetting state to first-time startup");
|
|
cfg.state.backfilled_until = 0;
|
|
if (pg_conn) {
|
|
/* Reset per-author backfill progress (caching_followed_pubkeys)
|
|
* so the next run starts fresh. The followed set itself is
|
|
* preserved; only cursor/completion state is reset. */
|
|
pg_inbox_reset_backfill_progress();
|
|
} else {
|
|
cr_config_save_state(&cfg);
|
|
}
|
|
}
|
|
|
|
/* Init crypto. */
|
|
if (nostr_crypto_init() != NOSTR_SUCCESS) {
|
|
DEBUG_ERROR("nostr_crypto_init failed");
|
|
return 1;
|
|
}
|
|
|
|
/* In-memory state. */
|
|
cr_seen_ring_t seen;
|
|
cr_seen_ring_init(&seen);
|
|
cr_pubkey_set_t followed;
|
|
cr_pubkey_set_init(&followed);
|
|
|
|
/* Upstream pool. */
|
|
nostr_relay_pool_t *upstream = nostr_relay_pool_create(nostr_pool_reconnect_config_default());
|
|
if (!upstream) {
|
|
DEBUG_ERROR("failed to create upstream pool");
|
|
return 1;
|
|
}
|
|
for (int i = 0; i < cfg.upstream_count; i++) {
|
|
if (nostr_relay_pool_add_relay(upstream, cfg.upstream_relays[i]) != NOSTR_SUCCESS) {
|
|
DEBUG_WARN("failed to add upstream relay %s (continuing)", cfg.upstream_relays[i]);
|
|
} else {
|
|
DEBUG_INFO("upstream: added %s", cfg.upstream_relays[i]);
|
|
}
|
|
}
|
|
|
|
/* Sink: WebSocket pool (legacy) or PostgreSQL inbox. */
|
|
cr_sink_t sink;
|
|
if (pg_conn) {
|
|
if (cr_sink_init_pg(&sink, &seen) != 0) {
|
|
DEBUG_ERROR("failed to init sink (pg)");
|
|
return 1;
|
|
}
|
|
} else {
|
|
if (cr_sink_init(&sink, cfg.local_relay, &seen) != 0) {
|
|
DEBUG_ERROR("failed to init sink");
|
|
return 1;
|
|
}
|
|
}
|
|
|
|
/* Give the sink pool a moment to connect to the local relay before we
|
|
* start publishing (relay discovery caches kind-10002 events to local).
|
|
* No-op in PostgreSQL inbox mode (cr_sink_pump does nothing). */
|
|
if (!pg_conn) {
|
|
DEBUG_INFO("waiting for sink connection to establish...");
|
|
for (int i = 0; i < 30 && !g_shutdown; i++) {
|
|
cr_sink_pump(&sink, 100);
|
|
}
|
|
}
|
|
|
|
/* Resolve follow graph. */
|
|
if (g_shutdown) goto shutdown;
|
|
if (cr_follow_resolve(&cfg, upstream, &followed) != 0) {
|
|
DEBUG_ERROR("follow graph resolution failed");
|
|
return 1;
|
|
}
|
|
|
|
/* Sync the followed set to the DB so the backfill tick can iterate
|
|
* over it (per-author until-cursor drain). Also run the one-time
|
|
* migration from the old caching_backfill_progress table. */
|
|
if (pg_conn) {
|
|
pg_inbox_migrate_backfill_progress();
|
|
const char **root_pks = malloc(cfg.root_npub_count * sizeof(char *));
|
|
for (int i = 0; i < cfg.root_npub_count; i++)
|
|
root_pks[i] = cfg.root_hex[i];
|
|
pg_inbox_sync_followed_pubkeys((const char **)followed.items,
|
|
followed.count, root_pks,
|
|
cfg.root_npub_count);
|
|
free(root_pks);
|
|
}
|
|
|
|
if (g_shutdown) goto shutdown;
|
|
|
|
/* Discover outbox relays (NIP-65) and compute minimum covering set. */
|
|
int is_first_time = (cfg.state.backfilled_until == 0);
|
|
cr_relay_map_t relay_map;
|
|
if (cr_relay_discovery_run(&relay_map, &cfg, upstream, &sink, &followed,
|
|
is_first_time) != 0) {
|
|
DEBUG_WARN("relay discovery failed, continuing with bootstrap relays only");
|
|
memset(&relay_map, 0, sizeof(relay_map));
|
|
}
|
|
|
|
/* Add discovered outbox relays to the upstream pool. */
|
|
for (int i = 0; i < relay_map.selected_count; i++) {
|
|
/* Check if already in the pool (bootstrap relays may already be there). */
|
|
char **listed = NULL;
|
|
nostr_pool_relay_status_t *statuses = NULL;
|
|
int n = nostr_relay_pool_list_relays(upstream, &listed, &statuses);
|
|
int already = 0;
|
|
for (int j = 0; j < n; j++) {
|
|
if (strcmp(listed[j], relay_map.selected_relays[i]) == 0) {
|
|
already = 1; break;
|
|
}
|
|
}
|
|
free(listed);
|
|
free(statuses);
|
|
if (!already) {
|
|
if (nostr_relay_pool_add_relay(upstream, relay_map.selected_relays[i]) == NOSTR_SUCCESS) {
|
|
DEBUG_INFO("upstream: added outbox relay %s", relay_map.selected_relays[i]);
|
|
}
|
|
}
|
|
}
|
|
|
|
/* Log final upstream pool. */
|
|
{
|
|
char **listed = NULL;
|
|
nostr_pool_relay_status_t *statuses = NULL;
|
|
int n = nostr_relay_pool_list_relays(upstream, &listed, &statuses);
|
|
DEBUG_INFO("upstream_pool: %d relays connected:", n);
|
|
for (int j = 0; j < n; j++) {
|
|
DEBUG_INFO(" %s", listed[j]);
|
|
}
|
|
free(listed);
|
|
free(statuses);
|
|
}
|
|
|
|
/* Create per-relay backfill progress rows for every followed pubkey
|
|
* (outbox relays + bootstrap relays). Existing rows are left untouched
|
|
* so cursor/completion state is preserved across restarts. */
|
|
if (pg_conn) {
|
|
init_relay_progress_for_all(&cfg, &relay_map, &followed);
|
|
}
|
|
|
|
/* Forward catch-up: bridge the gap for events posted while caching
|
|
* was off. Runs once at startup (not on --restart, which does a full
|
|
* re-drain from now). For each followed author with last_event_at > 0,
|
|
* queries since = last_event_at + 1, until = now. */
|
|
if (pg_conn && !restart && cfg.backfill.enabled) {
|
|
DEBUG_INFO("forward catch-up: bridging gap for followed authors");
|
|
cr_forward_catchup(&cfg, upstream, &sink);
|
|
}
|
|
|
|
/* Open live subscription. */
|
|
cr_live_t live;
|
|
if (cfg.live.enabled) {
|
|
if (cr_live_open(&live, &cfg, upstream, &followed, &sink) != 0) {
|
|
DEBUG_WARN("live subscription failed to open (will retry on resubscribe)");
|
|
}
|
|
} else {
|
|
memset(&live, 0, sizeof(live));
|
|
}
|
|
|
|
/* Init backfill. */
|
|
cr_backfill_t bf;
|
|
cr_backfill_init(&bf, &cfg);
|
|
|
|
time_t last_follow_refresh = time(NULL);
|
|
time_t last_state_save = time(NULL);
|
|
time_t last_status_heartbeat = 0;
|
|
|
|
/* Initial status heartbeat in PostgreSQL mode. */
|
|
if (pg_conn) {
|
|
int bf_complete = 0, bf_total = 0;
|
|
pg_inbox_count_backfill_progress(&bf_complete, &bf_total);
|
|
pg_inbox_update_status("starting", config_generation, (long)time(NULL),
|
|
followed.count, relay_map.selected_count, 0,
|
|
bf_complete, bf_total, 0, 0, NULL, 0);
|
|
}
|
|
|
|
DEBUG_INFO("entering main loop");
|
|
|
|
while (!g_shutdown) {
|
|
if (g_reload) {
|
|
g_reload = 0;
|
|
DEBUG_INFO("SIGHUP: reloading config (state preserved)");
|
|
cr_config_t newcfg;
|
|
int reload_ok = 0;
|
|
if (pg_conn) {
|
|
if (pg_config_load(&newcfg) == 0 &&
|
|
cr_follow_decode_roots(&newcfg) == 0) reload_ok = 1;
|
|
} else if (config_path) {
|
|
if (cr_config_load(&newcfg, config_path) == 0) reload_ok = 1;
|
|
}
|
|
if (reload_ok) {
|
|
/* Preserve runtime state across reload. */
|
|
newcfg.state = cfg.state;
|
|
cr_config_free(&cfg);
|
|
cfg = newcfg;
|
|
DEBUG_INFO("config reloaded");
|
|
/* Force an immediate follow-graph refresh so any new
|
|
* root npubs are picked up right away. */
|
|
last_follow_refresh = 0;
|
|
} else {
|
|
DEBUG_ERROR("config reload failed, keeping old config");
|
|
}
|
|
}
|
|
|
|
/* PostgreSQL: check for config generation change and reload. */
|
|
if (pg_conn) {
|
|
int changed = pg_config_generation_changed(config_generation);
|
|
if (changed == 1) {
|
|
long new_gen = pg_inbox_get_config_generation();
|
|
DEBUG_INFO("config generation changed (%ld -> %ld), reloading",
|
|
config_generation, new_gen);
|
|
cr_config_t newcfg;
|
|
if (pg_config_load(&newcfg) == 0) {
|
|
/* Decode root npubs to hex — pg_config_load fills
|
|
* root_npubs[] but NOT root_hex[] / root_hex_ready.
|
|
* Without this, cr_follow_is_root() returns 0 for
|
|
* everything and the new root's follows are never
|
|
* resolved. */
|
|
if (cr_follow_decode_roots(&newcfg) != 0) {
|
|
DEBUG_ERROR("config reload: failed to decode root npubs, keeping old config");
|
|
cr_config_free(&newcfg);
|
|
} else {
|
|
newcfg.state = cfg.state;
|
|
cr_config_free(&cfg);
|
|
cfg = newcfg;
|
|
config_generation = new_gen;
|
|
DEBUG_INFO("config reloaded from PostgreSQL");
|
|
/* Force an immediate follow-graph refresh so the
|
|
* new root npub's follows are picked up right
|
|
* away (instead of waiting up to
|
|
* follow_graph_refresh_seconds). */
|
|
last_follow_refresh = 0;
|
|
/* Resubscribe live with new config. */
|
|
if (cfg.live.enabled) {
|
|
cr_live_resubscribe(&live, &cfg, upstream, &followed, &sink);
|
|
}
|
|
}
|
|
} else {
|
|
DEBUG_ERROR("PostgreSQL config reload failed, keeping old config");
|
|
}
|
|
} else if (changed < 0) {
|
|
DEBUG_WARN("config generation check failed");
|
|
}
|
|
}
|
|
|
|
/* Pump upstream pool (drives live subscription callbacks). */
|
|
nostr_relay_pool_run(upstream, 100);
|
|
|
|
/* Pump sink pool (flush publish callbacks). */
|
|
cr_sink_pump(&sink, 50);
|
|
|
|
/* Backfill tick. */
|
|
int brc = cr_backfill_tick(&bf, &cfg, upstream, &followed, &sink, &relay_map);
|
|
(void)brc;
|
|
|
|
/* Immediate follow-graph refresh when a root npub publishes a new
|
|
* kind-3 contact list (detected by the live subscriber). */
|
|
if (cfg.live.enabled && live.follow_graph_changed) {
|
|
live.follow_graph_changed = 0;
|
|
DEBUG_INFO("live: admin kind-3 detected, refreshing follow graph immediately");
|
|
cr_pubkey_set_t new_followed;
|
|
cr_pubkey_set_init(&new_followed);
|
|
if (cr_follow_resolve(&cfg, upstream, &new_followed) == 0) {
|
|
int changed = (new_followed.count != followed.count);
|
|
if (!changed) {
|
|
for (int i = 0; i < followed.count; i++) {
|
|
if (!cr_pubkey_set_contains(&new_followed, followed.items[i])) {
|
|
changed = 1; break;
|
|
}
|
|
}
|
|
}
|
|
cr_pubkey_set_free(&followed);
|
|
followed = new_followed;
|
|
if (pg_conn) {
|
|
const char **root_pks = malloc(cfg.root_npub_count * sizeof(char *));
|
|
for (int i = 0; i < cfg.root_npub_count; i++)
|
|
root_pks[i] = cfg.root_hex[i];
|
|
pg_inbox_sync_followed_pubkeys((const char **)followed.items,
|
|
followed.count, root_pks,
|
|
cfg.root_npub_count);
|
|
free(root_pks);
|
|
/* Create relay progress rows for newly added follows.
|
|
* Existing rows are preserved (ON CONFLICT DO NOTHING). */
|
|
init_relay_progress_for_all(&cfg, &relay_map, &followed);
|
|
}
|
|
if (changed && cfg.live.enabled) {
|
|
cr_live_resubscribe(&live, &cfg, upstream, &followed, &sink);
|
|
}
|
|
/* Reset the backfill in_progress flag in case we had entered
|
|
* steady-state; new follows need to be drained. */
|
|
bf.in_progress = 1;
|
|
last_follow_refresh = time(NULL);
|
|
} else {
|
|
DEBUG_WARN("immediate follow graph refresh failed");
|
|
cr_pubkey_set_free(&new_followed);
|
|
}
|
|
}
|
|
|
|
/* Periodic follow-graph refresh. */
|
|
time_t now = time(NULL);
|
|
if (cfg.follow_graph_refresh_seconds > 0 &&
|
|
(now - last_follow_refresh) >= cfg.follow_graph_refresh_seconds) {
|
|
DEBUG_INFO("refreshing follow graph");
|
|
cr_pubkey_set_t new_followed;
|
|
cr_pubkey_set_init(&new_followed);
|
|
if (cr_follow_resolve(&cfg, upstream, &new_followed) == 0) {
|
|
/* If the set changed, resubscribe live. */
|
|
int changed = (new_followed.count != followed.count);
|
|
if (!changed) {
|
|
for (int i = 0; i < followed.count; i++) {
|
|
if (!cr_pubkey_set_contains(&new_followed, followed.items[i])) {
|
|
changed = 1; break;
|
|
}
|
|
}
|
|
}
|
|
cr_pubkey_set_free(&followed);
|
|
followed = new_followed;
|
|
/* Sync the refreshed followed set to the DB so newly added
|
|
* follows get backfilled and dropped follows stop. Existing
|
|
* cursor/completion state is preserved (upsert only touches
|
|
* is_root and last_seen). */
|
|
if (pg_conn) {
|
|
const char **root_pks = malloc(cfg.root_npub_count * sizeof(char *));
|
|
for (int i = 0; i < cfg.root_npub_count; i++)
|
|
root_pks[i] = cfg.root_hex[i];
|
|
pg_inbox_sync_followed_pubkeys((const char **)followed.items,
|
|
followed.count, root_pks,
|
|
cfg.root_npub_count);
|
|
free(root_pks);
|
|
/* Create relay progress rows for newly added follows.
|
|
* Existing rows are preserved (ON CONFLICT DO NOTHING). */
|
|
init_relay_progress_for_all(&cfg, &relay_map, &followed);
|
|
}
|
|
if (changed && cfg.live.enabled) {
|
|
cr_live_resubscribe(&live, &cfg, upstream, &followed, &sink);
|
|
}
|
|
/* Reset the backfill in_progress flag in case we had entered
|
|
* steady-state; new follows need to be drained. The backfill
|
|
* tick will quickly re-enter steady-state if there are no
|
|
* incomplete authors, so this is cheap. */
|
|
bf.in_progress = 1;
|
|
} else {
|
|
DEBUG_WARN("follow graph refresh failed");
|
|
cr_pubkey_set_free(&new_followed);
|
|
}
|
|
last_follow_refresh = now;
|
|
}
|
|
|
|
/* Periodic live resubscribe. */
|
|
if (cfg.live.enabled && cfg.live.resubscribe_interval_seconds > 0 &&
|
|
(now - live.last_resubscribe) >= cfg.live.resubscribe_interval_seconds) {
|
|
cr_live_resubscribe(&live, &cfg, upstream, &followed, &sink);
|
|
}
|
|
|
|
/* Periodic state save (in case backfill didn't just save). */
|
|
if ((now - last_state_save) >= 60) {
|
|
if (!pg_conn) cr_config_save_state(&cfg);
|
|
last_state_save = now;
|
|
}
|
|
|
|
/* PostgreSQL: periodic status heartbeat. */
|
|
if (pg_conn && (now - last_status_heartbeat) >= 15) {
|
|
long events_fetched = live.events_received + bf.events_total;
|
|
long inbox_inserts = sink.published_ok;
|
|
int connected = 0;
|
|
char **listed = NULL;
|
|
nostr_pool_relay_status_t *statuses = NULL;
|
|
int n = nostr_relay_pool_list_relays(upstream, &listed, &statuses);
|
|
if (n > 0) {
|
|
/* Build per-relay status arrays for the upstream_relays table. */
|
|
const char **urls = malloc((size_t)n * sizeof(char *));
|
|
int *codes = malloc((size_t)n * sizeof(int));
|
|
const char **texts = malloc((size_t)n * sizeof(char *));
|
|
if (urls && codes && texts) {
|
|
for (int j = 0; j < n; j++) {
|
|
urls[j] = listed[j];
|
|
codes[j] = (int)statuses[j];
|
|
if (statuses[j] == NOSTR_POOL_RELAY_CONNECTED) {
|
|
texts[j] = "connected";
|
|
connected++;
|
|
} else if (statuses[j] == NOSTR_POOL_RELAY_CONNECTING) {
|
|
texts[j] = "connecting";
|
|
} else if (statuses[j] == NOSTR_POOL_RELAY_DISCONNECTED) {
|
|
texts[j] = "disconnected";
|
|
} else {
|
|
texts[j] = "error";
|
|
}
|
|
}
|
|
pg_inbox_update_upstream_relays(urls, codes, texts, n);
|
|
}
|
|
free(urls);
|
|
free(codes);
|
|
free(texts);
|
|
}
|
|
free(listed);
|
|
free(statuses);
|
|
int bf_complete = 0, bf_total = 0;
|
|
pg_inbox_count_backfill_progress(&bf_complete, &bf_total);
|
|
pg_inbox_update_status("running", config_generation, (long)now,
|
|
followed.count, relay_map.selected_count,
|
|
connected, bf_complete, bf_total,
|
|
events_fetched, inbox_inserts, NULL, 0);
|
|
last_status_heartbeat = now;
|
|
}
|
|
}
|
|
|
|
/* Graceful shutdown. */
|
|
shutdown:
|
|
DEBUG_INFO("shutting down...");
|
|
cr_live_close(&live);
|
|
if (!pg_conn) cr_config_save_state(&cfg);
|
|
if (pg_conn) {
|
|
int bf_complete = 0, bf_total = 0;
|
|
pg_inbox_count_backfill_progress(&bf_complete, &bf_total);
|
|
pg_inbox_update_status("stopped", config_generation, (long)time(NULL),
|
|
followed.count, relay_map.selected_count, 0,
|
|
bf_complete, bf_total,
|
|
live.events_received + bf.events_total,
|
|
sink.published_ok, NULL, 0);
|
|
}
|
|
cr_sink_destroy(&sink);
|
|
nostr_relay_pool_destroy(upstream);
|
|
cr_pubkey_set_free(&followed);
|
|
cr_relay_map_free(&relay_map);
|
|
cr_config_free(&cfg);
|
|
if (pg_conn) pg_inbox_shutdown();
|
|
nostr_crypto_cleanup();
|
|
DEBUG_INFO("clean exit");
|
|
return 0;
|
|
}
|