From c184a759339a9a5d934b27e73252ed6d762586d1 Mon Sep 17 00:00:00 2001 From: Lewis Date: Sat, 4 Apr 2026 00:16:50 +0300 Subject: [PATCH] appview/ingester: harden event processing and cursor This is more of a practical proposal than solving a direct problem. We fire-and-forget our cursor updates even if an event fails. Instead, we should have it in the same thread as the record ingestion and not be silent about errors. Also for verification I think a little retry-backoff would be in order. Lewis: May this revision serve well! --- appview/ingester.go | 33 +++++++++++++++++++++------------ appview/state/state.go | 1 + eventconsumer/consumer.go | 8 ++++---- knotserver/ingester.go | 15 ++++++--------- spindle/ingester.go | 17 +++++++---------- 5 files changed, 39 insertions(+), 35 deletions(-) diff --git a/appview/ingester.go b/appview/ingester.go index 300c032b..e8d85ae4 100644 --- a/appview/ingester.go +++ b/appview/ingester.go @@ -12,6 +12,7 @@ import ( "time" + "github.com/avast/retry-go/v4" "github.com/bluesky-social/indigo/atproto/syntax" jmodels "github.com/bluesky-social/jetstream/pkg/models" "github.com/go-git/go-git/v5/plumbing" @@ -41,13 +42,6 @@ type processFunc func(ctx context.Context, e *jmodels.Event) error func (i *Ingester) Ingest() processFunc { return func(ctx context.Context, e *jmodels.Event) error { var err error - defer func() { - eventTime := e.TimeUS - lastTimeUs := eventTime + 1 - if err := i.Db.SaveLastTimeUs(lastTimeUs); err != nil { - err = fmt.Errorf("(deferred) failed to save last time us: %w", err) - } - }() l := i.Logger.With("kind", e.Kind) switch e.Kind { @@ -92,7 +86,12 @@ func (i *Ingester) Ingest() processFunc { } if err != nil { - l.Warn("refused to ingest record", "err", err) + 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 @@ -560,9 +559,13 @@ func (i *Ingester) ingestSpindle(ctx context.Context, e *jmodels.Event) error { return err } - err = serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) + err = retry.Do( + func() error { return serververify.RunVerification(ctx, instance, did, i.Config.Core.Dev) }, + retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second), + retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), + ) if err != nil { - l.Error("failed to add spindle to db", "err", err, "instance", instance) + l.Error("failed to verify spindle after retries", "err", err, "instance", instance) return err } @@ -778,9 +781,15 @@ func (i *Ingester) ingestKnot(e *jmodels.Event) error { return err } - err = serververify.RunVerification(context.Background(), domain, did, i.Config.Core.Dev) + err = retry.Do( + func() error { + return serververify.RunVerification(context.Background(), domain, did, i.Config.Core.Dev) + }, + retry.Attempts(5), retry.Delay(5*time.Second), retry.MaxDelay(80*time.Second), + retry.DelayType(retry.BackOffDelay), retry.LastErrorOnly(true), + ) if err != nil { - l.Error("failed to verify knot", "err", err, "domain", domain) + l.Error("failed to verify knot after retries", "err", err, "domain", domain) return err } diff --git a/appview/state/state.go b/appview/state/state.go index 0a236388..378f3fc6 100644 --- a/appview/state/state.go +++ b/appview/state/state.go @@ -123,6 +123,7 @@ func Make(ctx context.Context, config *config.Config) (*State, error) { tangled.KnotMemberNSID, tangled.SpindleMemberNSID, tangled.SpindleNSID, + tangled.KnotNSID, tangled.StringNSID, tangled.RepoIssueNSID, tangled.RepoIssueCommentNSID, diff --git a/eventconsumer/consumer.go b/eventconsumer/consumer.go index 328eeef8..25b22734 100644 --- a/eventconsumer/consumer.go +++ b/eventconsumer/consumer.go @@ -159,15 +159,15 @@ func (c *Consumer) worker(ctx context.Context) { return } + if err := c.cfg.ProcessFunc(ctx, j.source, msg); err != nil { + c.logger.Error("error processing message", "source", j.source, "err", err) + } + cursorVal := msg.Created if cursorVal == 0 { cursorVal = time.Now().UnixNano() } c.cfg.CursorStore.Set(j.source.Key(), cursorVal) - - if err := c.cfg.ProcessFunc(ctx, j.source, msg); err != nil { - c.logger.Error("error processing message", "source", j.source, "err", err) - } } } } diff --git a/knotserver/ingester.go b/knotserver/ingester.go index ec4cccab..ac450a6e 100644 --- a/knotserver/ingester.go +++ b/knotserver/ingester.go @@ -363,14 +363,6 @@ func (h *Knot) processMessages(ctx context.Context, event *models.Event) error { } var err error - defer func() { - eventTime := event.TimeUS - lastTimeUs := eventTime + 1 - if err := h.db.SaveLastTimeUs(lastTimeUs); err != nil { - err = fmt.Errorf("(deferred) failed to save last time us: %w", err) - } - }() - switch event.Commit.Collection { case tangled.PublicKeyNSID: err = h.processPublicKey(ctx, event) @@ -383,7 +375,12 @@ func (h *Knot) processMessages(ctx context.Context, event *models.Event) error { } if err != nil { - h.l.Debug("failed to process event", "nsid", event.Commit.Collection, "err", err) + h.l.Warn("failed to process event, skipping", "nsid", event.Commit.Collection, "err", err) + } + + 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 index 00ff1df6..43324aab 100644 --- a/spindle/ingester.go +++ b/spindle/ingester.go @@ -24,19 +24,11 @@ type Ingester func(ctx context.Context, e *models.Event) error func (s *Spindle) ingest() Ingester { return func(ctx context.Context, e *models.Event) error { - var err error - defer func() { - eventTime := e.TimeUS - lastTimeUs := eventTime + 1 - if err := s.db.SaveLastTimeUs(lastTimeUs); err != nil { - err = fmt.Errorf("(deferred) failed to save last time us: %w", err) - } - }() - if e.Kind != models.EventKindCommit { return nil } + var err error switch e.Commit.Collection { case tangled.SpindleMemberNSID: err = s.ingestMember(ctx, e) @@ -47,7 +39,12 @@ func (s *Spindle) ingest() Ingester { } if err != nil { - s.l.Debug("failed to process message", "nsid", e.Commit.Collection, "err", err) + s.l.Warn("failed to process message, skipping", "nsid", e.Commit.Collection, "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 -- 2.51.2