From c897faeeb49f1d7a9ccb922917eabc4e5e1126bd Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 07 Jul 2026 21:14:00 +0000 Subject: [PATCH] appview,knotserver,spindle/ingester: no more 2-day past cap Lewis: May this revision serve well! --- appview/ingester.go | 5 ----- jetstream/jetstream.go | 51 ++++++++++++++++++++++++++++----------------------- jetstream/jetstream_test.go | 81 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ knotserver/ingester.go | 5 ----- spindle/ingester.go | 5 ----- 5 file(s) changed, 109 insertion(s)(+), 38 deletion(s)(-) diff --git a/appview/ingester.go b/appview/ingester.go --- a/appview/ingester.go +++ b/appview/ingester.go @@ -134,11 +134,6 @@ l.Warn("failed to ingest record, skipping", "err", err) } - lastTimeUs := e.TimeUS + 1 - if saveErr := i.Db.SaveLastTimeUs(lastTimeUs); saveErr != nil { - l.Error("failed to save cursor", "err", saveErr) - } - return nil } } diff --git a/jetstream/jetstream.go b/jetstream/jetstream.go --- a/jetstream/jetstream.go +++ b/jetstream/jetstream.go @@ -7,6 +7,7 @@ "os" "os/signal" "sync" + "sync/atomic" "syscall" "time" @@ -34,6 +35,8 @@ db DB waitForDid bool mu sync.RWMutex + + lastSeenUs atomic.Int64 cancel context.CancelFunc cancelMu sync.Mutex @@ -71,7 +74,6 @@ // since this closure references j.WantedDids; it should auto-update // existing instances of the closure when j.WantedDids is mutated return func(ctx context.Context, evt *models.Event) error { - j.mu.RLock() // empty filter => all dids allowed matches := len(j.wantedDids) == 0 @@ -82,11 +84,13 @@ } j.mu.RUnlock() + var err error if matches { - return processFunc(ctx, evt) - } else { - return nil + err = processFunc(ctx, evt) } + + j.lastSeenUs.Store(evt.TimeUS + 1) + return err } } @@ -113,7 +117,7 @@ } // StartJetstream starts the jetstream client and processes events using the provided processFunc. -// The caller is responsible for saving the last time_us to the database (just use your db.UpdateLastTimeUs). +// The client persists the last time_us cursor itself via the DB it was constructed with. func (j *JetstreamClient) StartJetstream(ctx context.Context, processFunc func(context.Context, *models.Event) error) error { logger := j.l @@ -151,7 +155,7 @@ func (j *JetstreamClient) connectAndRead(ctx context.Context) { l := log.FromContext(ctx) for { - cursor := j.getLastTimeUs(ctx) + cursor := j.resumeCursor(ctx) connCtx, cancel := context.WithCancel(ctx) j.cancelMu.Lock() @@ -185,9 +189,20 @@ case <-ctx.Done(): return case <-ticker.C: - j.db.SaveLastTimeUs(time.Now().UnixMicro()) + if seen := j.lastSeenUs.Load(); seen != 0 { + if err := j.db.SaveLastTimeUs(seen); err != nil { + log.FromContext(ctx).Error("failed to save cursor", "error", err) + } + } } } +} + +func (j *JetstreamClient) resumeCursor(ctx context.Context) *int64 { + if seen := j.lastSeenUs.Load(); seen != 0 { + return &seen + } + return j.getLastTimeUs(ctx) } func (j *JetstreamClient) getLastTimeUs(ctx context.Context) *int64 { @@ -196,18 +211,7 @@ if err != nil { l.Warn("couldn't get last time us, starting from now", "error", err) lastTimeUs = time.Now().UnixMicro() - err = j.db.SaveLastTimeUs(lastTimeUs) - if err != nil { - l.Error("failed to save last time us", "error", err) - } - } - - // If last time is older than 2 days, start from now - if time.Now().UnixMicro()-lastTimeUs > 2*24*60*60*1000*1000 { - lastTimeUs = time.Now().UnixMicro() - l.Warn("last time us is older than 2 days; discarding that and starting from now") - err = j.db.SaveLastTimeUs(lastTimeUs) - if err != nil { + if err = j.db.SaveLastTimeUs(lastTimeUs); err != nil { l.Error("failed to save last time us", "error", err) } } @@ -234,11 +238,12 @@ sig := <-sigChan j.l.Info("Received signal, initiating graceful shutdown", "signal", sig) - lastTimeUs := time.Now().UnixMicro() - if err := j.db.SaveLastTimeUs(lastTimeUs); err != nil { - j.l.Error("Failed to save last time during shutdown", "error", err) + if seen := j.lastSeenUs.Load(); seen != 0 { + if err := j.db.SaveLastTimeUs(seen); err != nil { + j.l.Error("Failed to save last time during shutdown", "error", err) + } + j.l.Info("Saved lastTimeUs before shutdown", "lastTimeUs", seen) } - j.l.Info("Saved lastTimeUs before shutdown", "lastTimeUs", lastTimeUs) j.cancelMu.Lock() if j.cancel != nil { diff --git a/jetstream/jetstream_test.go b/jetstream/jetstream_test.go new file mode 100644 --- /dev/null +++ b/jetstream/jetstream_test.go @@ -0,0 +1,81 @@ +package jetstream + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/bluesky-social/jetstream/pkg/models" +) + +type fakeCursorDB struct { + saved int64 + savedErr error +} + +func (f *fakeCursorDB) GetLastTimeUs() (int64, error) { return f.saved, f.savedErr } +func (f *fakeCursorDB) SaveLastTimeUs(int64) error { return nil } + +func cursorFor(db DB) int64 { + j := &JetstreamClient{db: db, wantedDids: make(Set[string])} + return *j.getLastTimeUs(context.Background()) +} + +const twoDaysUs = int64(2 * 24 * 60 * 60 * 1000 * 1000) + +func TestStaleCursorIsHonored(t *testing.T) { + old := time.Now().UnixMicro() - 5*twoDaysUs + if got := cursorFor(&fakeCursorDB{saved: old}); got != old { + t.Fatalf("stale cursor must be honored, not snapped forward: got %d, want %d", got, old) + } +} + +func TestZeroCursorIsHonored(t *testing.T) { + if got := cursorFor(&fakeCursorDB{saved: 0}); got != 0 { + t.Fatalf("cursor 0 must replay from the start: got %d, want 0", got) + } +} + +func TestMissingCursorStartsFromNow(t *testing.T) { + before := time.Now().UnixMicro() + got := cursorFor(&fakeCursorDB{savedErr: errors.New("no row")}) + after := time.Now().UnixMicro() + if got < before || got > after { + t.Fatalf("missing cursor must start from now: got %d, want within [%d,%d]", got, before, after) + } +} + +func TestLastSeenTracksEveryEventPreFilter(t *testing.T) { + j := &JetstreamClient{wantedDids: Set[string]{"did:plc:boltless": {}}} + wrapped := j.withDidFilter(func(context.Context, *models.Event) error { return nil }) + + _ = wrapped(context.Background(), &models.Event{Did: "did:plc:akshay", TimeUS: 100}) + if got := j.lastSeenUs.Load(); got != 101 { + t.Fatalf("filtered-out event must still advance last-seen: got %d, want 101", got) + } + + _ = wrapped(context.Background(), &models.Event{Did: "did:plc:boltless", TimeUS: 200}) + if got := j.lastSeenUs.Load(); got != 201 { + t.Fatalf("matching event must advance last-seen: got %d, want 201", got) + } +} + +func TestLastSeenAdvancesAfterProcessing(t *testing.T) { + j := &JetstreamClient{wantedDids: make(Set[string])} + + var seenDuringProcess int64 + wrapped := j.withDidFilter(func(context.Context, *models.Event) error { + seenDuringProcess = j.lastSeenUs.Load() + return nil + }) + + _ = wrapped(context.Background(), &models.Event{Did: "did:plc:boltless", TimeUS: 500}) + + if seenDuringProcess == 501 { + t.Fatal("cursor advanced before the event finished processing; a crash mid-process would skip it") + } + if got := j.lastSeenUs.Load(); got != 501 { + t.Fatalf("cursor must advance once processing returns: got %d, want 501", got) + } +} diff --git a/knotserver/ingester.go b/knotserver/ingester.go --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -139,10 +139,5 @@ h.l.Warn("failed to process event, skipping", args...) } - lastTimeUs := event.TimeUS + 1 - if saveErr := h.db.SaveLastTimeUs(lastTimeUs); saveErr != nil { - h.l.Error("failed to save cursor", "err", saveErr) - } - return nil } diff --git a/spindle/ingester.go b/spindle/ingester.go --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -37,11 +37,6 @@ s.l.Warn("failed to process message, skipping", "nsid", e.Commit.Collection, "did", e.Did, "rkey", e.Commit.RKey, "err", err) } - lastTimeUs := e.TimeUS + 1 - if saveErr := s.db.SaveLastTimeUs(lastTimeUs); saveErr != nil { - s.l.Error("failed to save cursor", "err", saveErr) - } - return nil } } -- tangled.sh