201 lines
7.5 KiB
C
201 lines
7.5 KiB
C
/*
|
|
* caching_relay - forward catch-up implementation.
|
|
*
|
|
* See forward_catchup.h for design.
|
|
*/
|
|
#include "forward_catchup.h"
|
|
#include "pg_inbox.h"
|
|
#include "config.h"
|
|
#include "follow_graph.h"
|
|
#include "relay_sink.h"
|
|
#include "debug.h"
|
|
|
|
#include <cjson/cJSON.h>
|
|
#include <time.h>
|
|
#include <string.h>
|
|
#include <stdlib.h>
|
|
|
|
/* Find the newest (maximum) created_at among an array of event JSON objects.
|
|
* Returns 1 and sets *out_newest, or 0 if no events / no valid created_at. */
|
|
static int find_newest_created_at(cJSON **events, int count, long *out_newest) {
|
|
long newest = 0;
|
|
int found = 0;
|
|
for (int i = 0; i < count; i++) {
|
|
cJSON *ca = cJSON_GetObjectItem(events[i], "created_at");
|
|
if (!ca || !cJSON_IsNumber(ca)) continue;
|
|
long ts = (long)ca->valuedouble;
|
|
if (!found || ts > newest) {
|
|
newest = ts;
|
|
found = 1;
|
|
}
|
|
}
|
|
if (found && out_newest) *out_newest = newest;
|
|
return found;
|
|
}
|
|
|
|
/* Build the kinds array for a given pubkey (admin vs regular).
|
|
* Returns NULL if no kinds filter should be applied (admin_all_kinds). */
|
|
static cJSON *build_kinds(cr_config_t *cfg, const char *pk) {
|
|
int is_admin = cr_follow_is_root(cfg, pk);
|
|
if (is_admin && cfg->admin_all_kinds) {
|
|
return NULL;
|
|
}
|
|
cJSON *kinds = cJSON_CreateArray();
|
|
if (is_admin && cfg->admin_kind_count > 0) {
|
|
for (int i = 0; i < cfg->admin_kind_count; i++)
|
|
cJSON_AddItemToArray(kinds, cJSON_CreateNumber(cfg->admin_kinds[i]));
|
|
} else {
|
|
for (int i = 0; i < cfg->kind_count; i++)
|
|
cJSON_AddItemToArray(kinds, cJSON_CreateNumber(cfg->kinds[i]));
|
|
}
|
|
return kinds;
|
|
}
|
|
|
|
/* Get relay URLs for an author from the relay progress table.
|
|
* Returns a cJSON array of relay_url strings. Caller must cJSON_Delete().
|
|
* Returns NULL on error or if no relays found. */
|
|
static cJSON *get_author_relays(const char *pk) {
|
|
/* Use pg_inbox_get_incomplete_relays for incomplete authors.
|
|
* For completed authors, we need all relays — query directly. */
|
|
cJSON *incomplete = pg_inbox_get_incomplete_relays(pk);
|
|
if (incomplete && cJSON_GetArraySize(incomplete) > 0) {
|
|
/* Extract just the relay_url strings. */
|
|
cJSON *urls = cJSON_CreateArray();
|
|
int n = cJSON_GetArraySize(incomplete);
|
|
for (int i = 0; i < n; i++) {
|
|
cJSON *entry = cJSON_GetArrayItem(incomplete, i);
|
|
cJSON *url = cJSON_GetObjectItemCaseSensitive(entry, "relay_url");
|
|
if (url && cJSON_IsString(url) && url->valuestring[0]) {
|
|
cJSON_AddItemToArray(urls, cJSON_CreateString(url->valuestring));
|
|
}
|
|
}
|
|
cJSON_Delete(incomplete);
|
|
if (cJSON_GetArraySize(urls) > 0) return urls;
|
|
cJSON_Delete(urls);
|
|
}
|
|
if (incomplete) cJSON_Delete(incomplete);
|
|
|
|
/* For completed authors, fall back to the upstream pool's connected
|
|
* relays. The caller will use the pool directly. */
|
|
return NULL;
|
|
}
|
|
|
|
int cr_forward_catchup(cr_config_t *cfg,
|
|
nostr_relay_pool_t *upstream,
|
|
cr_sink_t *sink) {
|
|
if (!cfg || !upstream || !sink) return -1;
|
|
|
|
cJSON *authors = pg_inbox_get_authors_for_catchup();
|
|
if (!authors) {
|
|
DEBUG_INFO("forward_catchup: no authors with last_event_at > 0");
|
|
return 0;
|
|
}
|
|
|
|
int author_count = cJSON_GetArraySize(authors);
|
|
if (author_count == 0) {
|
|
cJSON_Delete(authors);
|
|
DEBUG_INFO("forward_catchup: no authors need catch-up");
|
|
return 0;
|
|
}
|
|
|
|
DEBUG_INFO("forward_catchup: checking %d authors for missed events", author_count);
|
|
|
|
long now = (long)time(NULL);
|
|
int page_size = cfg->backfill.events_per_tick;
|
|
if (page_size < 1) page_size = 500;
|
|
|
|
cr_sink_set_source_class(sink, CR_SINK_CLASS_BACKFILL);
|
|
|
|
int total_events_published = 0;
|
|
int authors_caught_up = 0;
|
|
|
|
for (int i = 0; i < author_count; i++) {
|
|
cJSON *entry = cJSON_GetArrayItem(authors, i);
|
|
if (!entry) continue;
|
|
|
|
cJSON *pk_node = cJSON_GetObjectItemCaseSensitive(entry, "pubkey");
|
|
cJSON *lea_node = cJSON_GetObjectItemCaseSensitive(entry, "last_event_at");
|
|
if (!pk_node || !cJSON_IsString(pk_node) || !lea_node || !cJSON_IsNumber(lea_node))
|
|
continue;
|
|
|
|
const char *pk = pk_node->valuestring;
|
|
long last_event_at = (long)lea_node->valuedouble;
|
|
|
|
if (last_event_at <= 0 || last_event_at >= now) continue;
|
|
|
|
/* Build filter: since = last_event_at + 1, until = now. */
|
|
cJSON *filter = cJSON_CreateObject();
|
|
cJSON *authors_arr = cJSON_CreateArray();
|
|
cJSON_AddItemToArray(authors_arr, cJSON_CreateString(pk));
|
|
cJSON_AddItemToObject(filter, "authors", authors_arr);
|
|
cJSON *kinds = build_kinds(cfg, pk);
|
|
if (kinds) cJSON_AddItemToObject(filter, "kinds", kinds);
|
|
cJSON_AddItemToObject(filter, "since", cJSON_CreateNumber((double)(last_event_at + 1)));
|
|
cJSON_AddItemToObject(filter, "until", cJSON_CreateNumber((double)now));
|
|
cJSON_AddItemToObject(filter, "limit", cJSON_CreateNumber((double)page_size));
|
|
|
|
/* Try to get author-specific relays first; fall back to upstream pool. */
|
|
cJSON *author_relays = get_author_relays(pk);
|
|
|
|
int ev_count = 0;
|
|
cJSON **events = NULL;
|
|
|
|
if (author_relays && cJSON_GetArraySize(author_relays) > 0) {
|
|
/* Query using the author's outbox relays. */
|
|
int n_relays = cJSON_GetArraySize(author_relays);
|
|
const char **urls = malloc(sizeof(char *) * n_relays);
|
|
for (int r = 0; r < n_relays; r++) {
|
|
cJSON *u = cJSON_GetArrayItem(author_relays, r);
|
|
urls[r] = cJSON_IsString(u) ? u->valuestring : "";
|
|
}
|
|
events = synchronous_query_relays_with_progress(
|
|
urls, n_relays, filter, RELAY_QUERY_ALL_RESULTS,
|
|
&ev_count, 30, NULL, NULL, 0, NULL);
|
|
free(urls);
|
|
} else {
|
|
/* Fall back: query all connected upstream relays.
|
|
* Use the pool's relay list. We pass NULL for urls to let
|
|
* the pool use all connected relays. */
|
|
events = synchronous_query_relays_with_progress(
|
|
NULL, 0, filter, RELAY_QUERY_ALL_RESULTS,
|
|
&ev_count, 30, NULL, NULL, 0, NULL);
|
|
}
|
|
cJSON_Delete(author_relays);
|
|
cJSON_Delete(filter);
|
|
|
|
if (events && ev_count > 0) {
|
|
/* Publish all events. */
|
|
long newest = 0;
|
|
find_newest_created_at(events, ev_count, &newest);
|
|
|
|
for (int k = 0; k < ev_count; k++) {
|
|
cr_sink_publish(sink, events[k]);
|
|
cJSON_Delete(events[k]);
|
|
}
|
|
free(events);
|
|
total_events_published += ev_count;
|
|
authors_caught_up++;
|
|
|
|
/* Update last_event_at to the newest event we saw. */
|
|
if (newest > 0) {
|
|
pg_inbox_update_last_event_at(pk, newest);
|
|
}
|
|
|
|
DEBUG_LOG("forward_catchup: %s -> %d events (since=%ld, until=%ld)",
|
|
pk, ev_count, last_event_at + 1, now);
|
|
} else {
|
|
/* No events in the gap — update last_event_at to now so we
|
|
* don't re-query the same empty window next time. */
|
|
pg_inbox_update_last_event_at(pk, now);
|
|
if (events) free(events);
|
|
}
|
|
}
|
|
|
|
cJSON_Delete(authors);
|
|
|
|
DEBUG_INFO("forward_catchup: %d authors checked, %d had events, %d total events published",
|
|
author_count, authors_caught_up, total_events_published);
|
|
|
|
return 0;
|
|
}
|