()
| 1678 | } |
| 1679 | |
| 1680 | func (h *heartbeats) subscribe() { |
| 1681 | defer h.wg.Done() |
| 1682 | eb := backoff.NewExponentialBackOff() |
| 1683 | eb.MaxElapsedTime = 0 // retry indefinitely |
| 1684 | eb.MaxInterval = dbMaxBackoff |
| 1685 | bkoff := backoff.WithContext(eb, h.ctx) |
| 1686 | var cancel context.CancelFunc |
| 1687 | bErr := backoff.Retry(func() error { |
| 1688 | cancelFn, err := h.pubsub.SubscribeWithErr(EventHeartbeats, h.listen) |
| 1689 | if err != nil { |
| 1690 | h.logger.Warn(h.ctx, "failed to tunnel to heartbeats", slog.Error(err)) |
| 1691 | return err |
| 1692 | } |
| 1693 | cancel = cancelFn |
| 1694 | return nil |
| 1695 | }, bkoff) |
| 1696 | if bErr != nil { |
| 1697 | if h.ctx.Err() == nil { |
| 1698 | h.logger.Error(h.ctx, "code bug: retry failed before context canceled", slog.Error(bErr)) |
| 1699 | } |
| 1700 | return |
| 1701 | } |
| 1702 | go func() { |
| 1703 | // cancel subscription when context finishes |
| 1704 | <-h.ctx.Done() |
| 1705 | cancel() |
| 1706 | }() |
| 1707 | } |
| 1708 | |
| 1709 | func (h *heartbeats) listen(_ context.Context, msg []byte, err error) { |
| 1710 | if err != nil { |
no test coverage detected