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