summaryrefslogtreecommitdiffstats
path: root/internal/transporthealth/monitor.go
blob: 48369ee55b3b352c7eae3cd2d6fa5dc3509b5815 (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
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()
}