diff --git a/server/event_emmiter.go b/server/event_emmiter.go index fe3bfba..c04052b 100644 --- a/server/event_emmiter.go +++ b/server/event_emmiter.go @@ -16,12 +16,14 @@ func (s *Server) emmitEvents(ctx context.Context) error { logger := s.logger.With("component", "event-emmiter") ident := "self" - var since *int64 - // TODO: track since + + day := time.Hour * 24 + + since := time.Now().Add(-day).UnixMilli() evts, evtManCancel, err := s.evtman.Subscribe(ctx, ident, func(evt *events.XRPCStreamEvent) bool { return true - }, since) + }, &since) if err != nil { return err } diff --git a/server/persist.go b/server/persist.go index 700061f..4204ea3 100644 --- a/server/persist.go +++ b/server/persist.go @@ -60,7 +60,7 @@ func (p *DbPersister) Persist(ctx context.Context, e *events.XRPCStreamEvent) er rec := &models.EventRecord{} if err := p.Db.Order("seq desc").Limit(1).First(rec).Error; err != nil { slog.Error("fetching most recent event record", "error", err) - rec.Seq = time.Now().Unix() + rec.Seq = time.Now().UnixMilli() } // if the error on inserting the event record is a constraint error, it means that @@ -140,6 +140,8 @@ 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)) + if len(records) == 0 { return nil }