diff --git a/knotclient/events.go b/knotclient/events.go index c13a65d5..808ba595 100644 --- a/knotclient/events.go +++ b/knotclient/events.go @@ -180,7 +180,7 @@ func (c *EventConsumer) worker(ctx context.Context) { } // update cursor - c.cfg.CursorStore.Set(j.source.Knot, time.Now().Unix()) + c.cfg.CursorStore.Set(j.source.Knot, time.Now().UnixNano()) 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/events.go b/knotserver/events.go index 4471e634..b6fbf60d 100644 --- a/knotserver/events.go +++ b/knotserver/events.go @@ -54,12 +54,12 @@ func (h *Handle) Events(w http.ResponseWriter, r *http.Request) { } // complete backfill first before going to live data - l.Info("going through backfill", "cursor", cursor) l.Debug("going through backfill", "cursor", cursor) if err := h.streamOps(conn, &cursor); err != nil { l.Error("failed to backfill", "err", err) return } + for { // wait for new data or timeout select {