Files
c-relay-pg/plans/forward_catchup_plan.md
T

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]