diff --git a/server/event_emmiter.go b/server/event_emmiter.go index c04052b..175d108 100644 --- a/server/event_emmiter.go +++ b/server/event_emmiter.go @@ -3,11 +3,13 @@ package server import ( "bytes" "context" + "log/slog" "net/http" "time" "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/lex/util" + "github.com/haileyok/cocoon/models" ) func (s *Server) emmitEvents(ctx context.Context) error { @@ -17,9 +19,15 @@ func (s *Server) emmitEvents(ctx context.Context) error { logger := s.logger.With("component", "event-emmiter") ident := "self" - day := time.Hour * 24 + // get the most recent event + rec := &models.EventRecord{} + err := s.db.Raw(ctx, "SELECT * from event_records ORDER BY seq DESC LIMIT 1", nil).Scan(&rec).Error + if err != nil { + slog.Error("fetching most recent event record", "error", err) + rec.Seq = time.Now().UnixMilli() + } - since := time.Now().Add(-day).UnixMilli() + since := rec.Seq evts, evtManCancel, err := s.evtman.Subscribe(ctx, ident, func(evt *events.XRPCStreamEvent) bool { return true @@ -72,11 +80,17 @@ func (s *Server) emmitEvents(ctx context.Context) error { } // TODO: use a HTTP client here not the default - _, err := http.Post(s.config.SubscribeReposServiceURL, "", buf) + resp, err := http.Post(s.config.SubscribeReposServiceURL, "", buf) if err != nil { logger.Error("posting to web server", "error", err) return } + if resp.StatusCode == http.StatusAccepted { + logger.Info("posted event to subscribe repos") + return + } + + logger.Error("posting event to subscribe repos", "status", resp.StatusCode) }() } diff --git a/server/persist.go b/server/persist.go index 4204ea3..c880746 100644 --- a/server/persist.go +++ b/server/persist.go @@ -116,6 +116,7 @@ func (p *DbPersister) Persist(ctx context.Context, e *events.XRPCStreamEvent) er }, retry.RetryIf(retryIfFunc)) if err != nil { + slog.Error(err.Error()) return err } @@ -140,7 +141,7 @@ func (p *DbPersister) Playback(ctx context.Context, since int64, cb func(*events return fmt.Errorf("failed to query events: %w", err) } - slog.Info("playback", "len of records", len(records)) + slog.Info("playback", "len of records", len(records), "cursor", cursor) if len(records) == 0 { return nil