package transporthealth import ( "context" "log" "sync" "time" ) type Checker interface { Check(context.Context) error } type Monitor struct { checker Checker interval time.Duration timeout time.Duration mu sync.RWMutex ready bool checkedAt time.Time } func New(checker Checker, interval, timeout time.Duration) *Monitor { if interval < time.Minute { interval = 15 * time.Minute } if timeout <= 0 { timeout = 90 * time.Second } return &Monitor{checker: checker, interval: interval, timeout: timeout} } func (m *Monitor) Run(ctx context.Context) { m.Probe(ctx) ticker := time.NewTicker(m.interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: m.Probe(ctx) } } } func (m *Monitor) Probe(ctx context.Context) { probeCtx, cancel := context.WithTimeout(ctx, m.timeout) defer cancel() if err := m.checker.Check(probeCtx); err != nil { m.MarkFailure(err) log.Printf("transport probe failed: %v", err) return } m.MarkSuccess() log.Printf("transport probe succeeded") } func (m *Monitor) Ready() bool { m.mu.RLock() ready := m.ready checkedAt := m.checkedAt m.mu.RUnlock() if !ready || checkedAt.IsZero() { return false } return time.Since(checkedAt) <= 2*m.interval+m.timeout } func (m *Monitor) MarkSuccess() { m.mu.Lock() m.ready = true m.checkedAt = time.Now() m.mu.Unlock() } func (m *Monitor) MarkFailure(error) { m.mu.Lock() m.ready = false m.checkedAt = time.Now() m.mu.Unlock() }