307 lines
12 KiB
Markdown
307 lines
12 KiB
Markdown
# Forward Catch-Up Plan: Bridging the Gap When Caching Resumes
|
|
|
|
## Problem
|
|
|
|
When caching is turned off (or the caching service stops), events posted by
|
|
followed authors during the downtime are missed. The current backfill walks
|
|
**backward** from `until_cursor` toward the beginning of time — it does not
|
|
cover events **newer** than the cursor. The live subscriber uses
|
|
`since = time(NULL)`, so it only catches events from the moment it connects.
|
|
Events in the gap between "last event we know about" and "caching resumed"
|
|
are lost.
|
|
|
|
## Schema Change: `last_event_at` column
|
|
|
|
Add a `last_event_at BIGINT NOT NULL DEFAULT 0` column to
|
|
`caching_followed_pubkeys`. This records the `created_at` of the most
|
|
recent event we've ever seen for this author — not a proxy like
|
|
`updated_at` (which records when we last *touched* the row, which could
|
|
be a zero-event progress write or a follow-graph refresh).
|
|
|
|
`since = last_event_at + 1, until = now` is exact: it catches every
|
|
event the relay doesn't yet have, with no assumptions about when the
|
|
last event occurred relative to when we last checked.
|
|
|
|
### Population
|
|
|
|
- **Backfill**: when publishing a page of events, compute
|
|
`max(created_at)` across the page (we already compute
|
|
`min(created_at)` for the cursor advance — add a parallel max).
|
|
Update `last_event_at = GREATEST(last_event_at, page_max_created_at)`.
|
|
- **Live subscriber**: when an event is received and published, update
|
|
`last_event_at = GREATEST(last_event_at, event.created_at)`.
|
|
- **One-time seed from `events` table**: on schema upgrade, set
|
|
`last_event_at = COALESCE((SELECT MAX(created_at) FROM events WHERE
|
|
pubkey = caching_followed_pubkeys.pubkey), 0)` for all existing rows.
|
|
This seeds the column from the relay's own data.
|
|
|
|
### Migration
|
|
|
|
Added to `src/pg_schema.sql` and `src/pg_schema.h`:
|
|
|
|
```sql
|
|
ALTER TABLE caching_followed_pubkeys
|
|
ADD COLUMN IF NOT EXISTS last_event_at BIGINT NOT NULL DEFAULT 0;
|
|
|
|
-- One-time seed from the events table.
|
|
UPDATE caching_followed_pubkeys fp
|
|
SET last_event_at = COALESCE(
|
|
(SELECT MAX(e.created_at) FROM events e WHERE e.pubkey = fp.pubkey),
|
|
0
|
|
)
|
|
WHERE fp.last_event_at = 0;
|
|
```
|
|
|
|
No schema version bump needed — `ALTER TABLE ADD COLUMN IF NOT EXISTS`
|
|
is idempotent and the seed `UPDATE` is guarded by `WHERE last_event_at = 0`.
|
|
|
|
## Current Architecture
|
|
|
|
### Tables
|
|
|
|
**`caching_followed_pubkeys`** (per-author state):
|
|
| Column | Purpose |
|
|
|---|---|
|
|
| `pubkey` | PK |
|
|
| `until_cursor` | Unix timestamp; backfill queries `until = this`, walks backward |
|
|
| `backfill_complete` | TRUE when drained to the beginning of time |
|
|
| `events_fetched` | Cumulative count |
|
|
| `last_seen` | Updated on follow-graph refresh |
|
|
| `updated_at` | Updated on every backfill progress write and completion mark |
|
|
| `last_event_at` | **NEW**: `created_at` of the most recent event we've seen for this author |
|
|
|
|
**`caching_backfill_relay_progress`** (per-author-per-relay state):
|
|
| Column | Purpose |
|
|
|---|---|
|
|
| `author_pubkey, relay_url` | Composite PK |
|
|
| `until_cursor` | Per-relay backward-walk cursor |
|
|
| `complete` | TRUE when this relay is drained for this author |
|
|
| `updated_at` | Updated on every relay progress write |
|
|
|
|
### Backfill flow ([`caching/src/backfill.c`](../caching/src/backfill.c:256))
|
|
|
|
1. Pick next incomplete author (round-robin)
|
|
2. For each incomplete relay for that author:
|
|
- Query `authors=[pk], until=until_cursor, limit=page_size`
|
|
- Publish events to the relay via the sink
|
|
- Advance `until_cursor` to `oldest_event_created_at - 1`
|
|
- If < page_size events + EOSE: mark relay `complete = TRUE`
|
|
3. If all relays complete: mark author `backfill_complete = TRUE`
|
|
|
|
### Restart behavior
|
|
|
|
- **Normal restart** (no `--restart`): cursor and completion state preserved.
|
|
Backfill resumes the backward walk from stored cursors. **Gap not covered.**
|
|
- **`--restart`**: all cursors reset to 0, all completion flags cleared. Full
|
|
re-drain from `now` backward. **Re-fetches everything** but covers the gap
|
|
incidentally (since it starts from `now`).
|
|
|
|
### Live subscriber ([`caching/src/live_subscriber.c`](../caching/src/live_subscriber.c:86))
|
|
|
|
Uses `since = time(NULL)` — only catches events from the moment it connects.
|
|
No gap-bridging.
|
|
|
|
## The Gap
|
|
|
|
```
|
|
Time ──────────────────────────────────────────────────────►
|
|
│ │ │
|
|
last_event_at caching now
|
|
(most recent event turned (caching
|
|
we know about) off resumed)
|
|
│
|
|
└── events posted here are missed
|
|
```
|
|
|
|
For **completed authors** (`backfill_complete = TRUE`): the backward walk is
|
|
done. `last_event_at` tells us the most recent event we have. Events posted
|
|
after `last_event_at` are in the gap.
|
|
|
|
For **incomplete authors**: the backward walk is still in progress. The cursor
|
|
is somewhere in the past, walking backward. Events newer than the cursor are
|
|
not fetched by backfill. The live subscriber covers events from `now` forward.
|
|
The gap is between `last_event_at` and `now`.
|
|
|
|
## Solution: Forward Catch-Up Phase
|
|
|
|
Add a **forward catch-up** phase that runs once when the caching service
|
|
starts (or when backfill is re-enabled), before the normal backward-drain
|
|
backfill loop begins.
|
|
|
|
### Logic
|
|
|
|
For **every** followed author (both complete and incomplete):
|
|
1. Read `last_event_at` from `caching_followed_pubkeys`.
|
|
2. If `last_event_at = 0`, skip (no events known yet — the backward drain
|
|
will handle it).
|
|
3. Query one or two outbox relays:
|
|
`authors=[pk], since=last_event_at + 1, until=now, limit=page_size`
|
|
4. Publish all returned events to the relay via the sink.
|
|
5. Update `last_event_at` to the max `created_at` seen (or `now` if no
|
|
events were returned, to avoid re-querying the same empty window next
|
|
time).
|
|
|
|
This is safe for both complete and incomplete authors:
|
|
- **Complete authors**: the backward drain is done, so the forward catch-up
|
|
is the only thing needed.
|
|
- **Incomplete authors**: the backward drain walks *below* `until_cursor`,
|
|
so events above `until_cursor` up to `last_event_at` were already fetched
|
|
during the initial drain (when `until_cursor` started at `now`). The
|
|
forward catch-up fills from `last_event_at + 1` to `now` — the gap that
|
|
formed while caching was off.
|
|
|
|
### When to run
|
|
|
|
- **On caching service startup** (not `--restart`, which does a full reset).
|
|
- The caching service process starts when `caching_enabled` is turned on,
|
|
so this covers the "caching was off, now it's on" case.
|
|
|
|
### Implementation
|
|
|
|
#### 1. Schema: add `last_event_at` column
|
|
|
|
In `src/pg_schema.sql` and `src/pg_schema.h`:
|
|
|
|
```sql
|
|
ALTER TABLE caching_followed_pubkeys
|
|
ADD COLUMN IF NOT EXISTS last_event_at BIGINT NOT NULL DEFAULT 0;
|
|
|
|
UPDATE caching_followed_pubkeys fp
|
|
SET last_event_at = COALESCE(
|
|
(SELECT MAX(e.created_at) FROM events e WHERE e.pubkey = fp.pubkey),
|
|
0
|
|
)
|
|
WHERE fp.last_event_at = 0;
|
|
```
|
|
|
|
#### 2. New function: `pg_inbox_get_authors_for_catchup()`
|
|
|
|
In `caching/src/pg_inbox.c`:
|
|
|
|
```c
|
|
/* Returns a cJSON array of {pubkey, last_event_at} objects for all
|
|
* followed authors where last_event_at > 0. Caller must cJSON_Delete().
|
|
* Returns NULL on error. */
|
|
cJSON* pg_inbox_get_authors_for_catchup(void);
|
|
```
|
|
|
|
SQL:
|
|
```sql
|
|
SELECT pubkey, last_event_at
|
|
FROM caching_followed_pubkeys
|
|
WHERE last_event_at > 0
|
|
ORDER BY last_event_at ASC
|
|
```
|
|
|
|
#### 3. New function: `pg_inbox_update_last_event_at()`
|
|
|
|
```c
|
|
/* Update last_event_at for a pubkey to the max of current and new value. */
|
|
int pg_inbox_update_last_event_at(const char *pk, long event_created_at);
|
|
```
|
|
|
|
SQL:
|
|
```sql
|
|
UPDATE caching_followed_pubkeys
|
|
SET last_event_at = GREATEST(last_event_at, $2::BIGINT)
|
|
WHERE pubkey = $1
|
|
```
|
|
|
|
#### 4. New function: `cr_forward_catchup()`
|
|
|
|
In a new file `caching/src/forward_catchup.c`:
|
|
|
|
```c
|
|
/* Run forward catch-up for all followed authors.
|
|
* For each author with last_event_at > 0, query events from
|
|
* last_event_at + 1 to now and publish them to the sink.
|
|
* Returns 0 on success, -1 on error. */
|
|
int cr_forward_catchup(cr_config_t *cfg,
|
|
nostr_relay_pool_t *upstream,
|
|
cr_sink_t *sink);
|
|
```
|
|
|
|
Flow:
|
|
1. Call `pg_inbox_get_authors_for_catchup()` to get the list.
|
|
2. For each author:
|
|
a. Get the author's outbox relays from `caching_backfill_relay_progress`
|
|
(any relay, since we just need one good source).
|
|
b. Query `authors=[pk], since=last_event_at + 1, until=now,
|
|
limit=page_size` on one relay.
|
|
c. Publish all returned events to the sink.
|
|
d. If events were returned, update `last_event_at` to the max
|
|
`created_at` in the batch. If no events, update `last_event_at`
|
|
to `now` (so we don't re-query the same empty window).
|
|
3. Log: "forward catch-up: N authors checked, M events published".
|
|
|
|
#### 5. Update backfill to maintain `last_event_at`
|
|
|
|
In `caching/src/backfill.c`, in the page-publishing loop (around line 375):
|
|
- Add a `find_newest_created_at()` helper (parallel to the existing
|
|
`find_oldest_created_at()`).
|
|
- After publishing a page, call
|
|
`pg_inbox_update_last_event_at(pk, newest_created_at)`.
|
|
|
|
#### 6. Update live subscriber to maintain `last_event_at`
|
|
|
|
In `caching/src/live_subscriber.c`, in the event-received callback:
|
|
- Extract `created_at` from the event.
|
|
- Call `pg_inbox_update_last_event_at(pubkey, created_at)`.
|
|
|
|
#### 7. Call from `main.c`
|
|
|
|
In `caching/src/main.c`, after relay discovery and followed-set sync,
|
|
before the main loop:
|
|
|
|
```c
|
|
/* Forward catch-up: bridge the gap for all followed authors. */
|
|
if (cfg->backfill.enabled && pg_conn && !restart) {
|
|
DEBUG_INFO("forward catch-up: checking for missed events");
|
|
cr_forward_catchup(&cfg, upstream, &sink);
|
|
}
|
|
```
|
|
|
|
This runs once at startup. It's not throttled — it's a one-time pass.
|
|
|
|
### Edge cases
|
|
|
|
- **`--restart` flag**: full reset already starts from `now`, so forward
|
|
catch-up is skipped. The `last_event_at` seed from the `events` table
|
|
will set it to the most recent known event, and the backward drain from
|
|
`now` will cover everything.
|
|
- **`last_event_at = 0`**: author has no known events. Skip — the backward
|
|
drain handles it.
|
|
- **Very large gap** (caching off for months): the forward catch-up query
|
|
may return many events. Use `limit = page_size` and paginate if needed.
|
|
The relay's own dedup (unique index on event ID) handles duplicates.
|
|
- **Relay doesn't support `since`**: most Nostr relays support `since`/
|
|
`until` (NIP-01). If ignored, the relay returns all events — dedup
|
|
handles it.
|
|
|
|
### Files to change
|
|
|
|
| File | Change |
|
|
|---|---|
|
|
| `src/pg_schema.sql` | `ALTER TABLE` add `last_event_at` + seed from `events` |
|
|
| `src/pg_schema.h` | Mirror the above |
|
|
| `caching/src/forward_catchup.c` | New: `cr_forward_catchup()` |
|
|
| `caching/src/forward_catchup.h` | New: declaration |
|
|
| `caching/src/pg_inbox.c` | New: `pg_inbox_get_authors_for_catchup()`, `pg_inbox_update_last_event_at()` |
|
|
| `caching/src/pg_inbox.h` | New: declarations |
|
|
| `caching/src/backfill.c` | Add `find_newest_created_at()`, call `pg_inbox_update_last_event_at()` after each page |
|
|
| `caching/src/live_subscriber.c` | Call `pg_inbox_update_last_event_at()` on event receipt |
|
|
| `caching/src/main.c` | Call `cr_forward_catchup()` at startup |
|
|
| `caching/Makefile` | Add `forward_catchup.c` to sources |
|
|
|
|
### Sequencing
|
|
|
|
```mermaid
|
|
graph TD
|
|
A[Caching service starts] --> B{Is --restart?}
|
|
B -- Yes --> C[Reset all progress, full re-drain from now]
|
|
B -- No --> D[Forward catch-up: all authors with last_event_at > 0]
|
|
D --> E[Normal backward-drain backfill loop]
|
|
E --> F[Live subscriber: since = now, ongoing]
|
|
F --> G[Live subscriber updates last_event_at on each event]
|
|
E --> H[Backfill updates last_event_at on each page]
|