Files
kindexr/internal/nostr/reader.go
T
enki 1933306392 cleanup: pre-phase-2 fixes
- fix two stale nzbstr comment refs
- migration 002: drop api_keys.visibility and curation_set
- remove Visibility from APIKey struct, GetAPIKey, CreateAPIKey, CLI
- omit torznab size/seeders/peers attrs when data is absent
- reset relay backoff on successful connection (>= 30s uptime)
- use last_event from relays table as since on reconnect
- fix TestSearchWithResults to actually query the test server DB
2026-05-16 23:21:22 -07:00

141 lines
3.5 KiB
Go

package nostr
import (
"context"
"log/slog"
"sync/atomic"
"time"
nostrlib "github.com/nbd-wtf/go-nostr"
"git.utn.lol/enki/kindexr/internal/config"
"git.utn.lol/enki/kindexr/internal/db"
)
// Reader subscribes to Nostr relays and indexes kind 2003 events.
type Reader struct {
cfg *config.Config
db *db.DB
connected atomic.Int32
}
// New creates a Reader.
func New(cfg *config.Config, database *db.DB) *Reader {
return &Reader{cfg: cfg, db: database}
}
// ConnectedCount returns the number of currently connected relays.
func (rd *Reader) ConnectedCount() int {
return int(rd.connected.Load())
}
// Start launches one goroutine per configured relay. It returns immediately;
// relay connections run in the background until ctx is cancelled.
func (rd *Reader) Start(ctx context.Context) {
for _, relayURL := range rd.cfg.Relays {
// Seed the relays table so we can track sync state.
_ = rd.db.UpsertRelay(ctx, relayURL)
go rd.connectLoop(ctx, relayURL)
}
}
// connectLoop keeps re-connecting to a single relay with exponential backoff.
// Backoff resets to the base delay if the connection stayed up long enough to
// be considered successful (30 seconds), so a briefly-dropped live relay
// reconnects quickly rather than waiting minutes.
func (rd *Reader) connectLoop(ctx context.Context, relayURL string) {
const baseBackoff = 5 * time.Second
const maxBackoff = 5 * time.Minute
const successThreshold = 30 * time.Second
backoff := baseBackoff
for {
if ctx.Err() != nil {
return
}
start := time.Now()
if err := rd.runRelay(ctx, relayURL); err != nil {
slog.Warn("relay disconnected", "url", relayURL, "err", err, "retry_in", backoff)
}
if time.Since(start) >= successThreshold {
backoff = baseBackoff
}
select {
case <-ctx.Done():
return
case <-time.After(backoff):
if backoff < maxBackoff {
backoff *= 2
}
}
}
}
// runRelay connects, subscribes, and streams events until the connection drops
// or ctx is cancelled.
func (rd *Reader) runRelay(ctx context.Context, relayURL string) error {
relay, err := nostrlib.RelayConnect(ctx, relayURL)
if err != nil {
return err
}
defer relay.Close()
rd.connected.Add(1)
defer rd.connected.Add(-1)
_ = rd.db.UpdateRelaySync(ctx, relayURL, time.Now().Unix())
slog.Info("relay connected", "url", relayURL)
backfillSince := time.Now().Add(-time.Duration(rd.cfg.BackfillDays) * 24 * time.Hour).Unix()
lastEvent, _ := rd.db.GetRelayLastEvent(ctx, relayURL)
if lastEvent > backfillSince {
backfillSince = lastEvent
}
since := nostrlib.Timestamp(backfillSince)
sub, err := relay.Subscribe(ctx, nostrlib.Filters{{
Kinds: []int{nostrlib.KindTorrent},
Since: &since,
}})
if err != nil {
return err
}
for {
select {
case <-ctx.Done():
return nil
case event, ok := <-sub.Events:
if !ok {
return nil
}
rd.handleEvent(ctx, relayURL, event)
case <-sub.EndOfStoredEvents:
slog.Info("relay EOSE", "url", relayURL)
}
}
}
// handleEvent parses and stores a single event.
func (rd *Reader) handleEvent(ctx context.Context, relayURL string, event *nostrlib.Event) {
rec, err := Parse(event)
if err != nil {
slog.Debug("skipping event", "id", event.ID, "reason", err)
return
}
if err := rd.db.InsertTorrent(ctx, *rec); err != nil {
slog.Error("db insert failed", "id", event.ID, "err", err)
return
}
_ = rd.db.UpdateRelayLastEvent(ctx, relayURL, int64(event.CreatedAt))
slog.Debug("indexed event", "id", event.ID, "title", rec.Title)
}