Files

277 lines
8.0 KiB
Go

// Package observability provides the intentionally small health/metrics HTTP
// surface shared by the Linux server and Windows client. It contains no
// product state and never exposes command payloads or high-cardinality IDs.
package observability
import (
"context"
"fmt"
"math"
"net"
"net/http"
"sort"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
)
type Paths struct {
Liveness string
Readiness string
Metrics string
}
type Health struct {
ready atomic.Bool
dirty atomic.Bool
mu sync.Mutex
count map[string]uint64
gauge map[string]float64
hist map[string]*durationHistogram
}
type durationHistogram struct {
count uint64
sum float64
buckets [len(durationBuckets)]uint64
}
// The fixed buckets are intentionally shared by every duration metric. This
// keeps the Prometheus surface bounded and makes comparable operational
// latencies available without allowing callers to create arbitrary labels.
var durationBuckets = [...]float64{0.001, 0.005, 0.01, 0.05, 0.1, 0.25, 0.5, 1, 5, 15, 60}
func New() *Health {
return &Health{
count: make(map[string]uint64),
gauge: make(map[string]float64),
hist: make(map[string]*durationHistogram),
}
}
func (health *Health) SetReady(value bool) {
if health != nil {
health.ready.Store(value)
}
}
func (health *Health) SetDirty(value bool) {
if health != nil {
health.dirty.Store(value)
}
}
func (health *Health) Inc(name string) {
health.IncBy(name, 1)
}
// IncBy records a non-negative integral counter delta. Metric names are
// deliberately restricted to a small identifier grammar; callers cannot turn
// a client ID, request UUID, or other untrusted data into a metric name.
func (health *Health) IncBy(name string, delta uint64) {
if health == nil || delta == 0 || !validMetricName(name) {
return
}
health.mu.Lock()
health.count[name] += delta
health.mu.Unlock()
}
// SetGauge records a bounded-cardinality instantaneous value. NaN and
// infinities are discarded because they are not safe Prometheus samples.
func (health *Health) SetGauge(name string, value float64) {
if health == nil || !validMetricName(name) || math.IsNaN(value) || math.IsInf(value, 0) {
return
}
health.mu.Lock()
health.gauge[name] = value
health.mu.Unlock()
}
// ObserveDuration records a duration in a fixed-bucket histogram. Negative
// durations are invalid (normally a caller clock error) and are ignored.
func (health *Health) ObserveDuration(name string, duration time.Duration) {
if health == nil || !validMetricName(name) || duration < 0 {
return
}
seconds := duration.Seconds()
health.mu.Lock()
histogram := health.hist[name]
if histogram == nil {
histogram = &durationHistogram{}
health.hist[name] = histogram
}
histogram.count++
histogram.sum += seconds
for index, upperBound := range durationBuckets {
if seconds <= upperBound {
histogram.buckets[index]++
}
}
health.mu.Unlock()
}
func (health *Health) Snapshot() (ready, dirty bool, counters map[string]uint64) {
if health == nil {
return false, true, nil
}
health.mu.Lock()
defer health.mu.Unlock()
copyCounters := make(map[string]uint64, len(health.count))
for key, value := range health.count {
copyCounters[key] = value
}
return health.ready.Load(), health.dirty.Load(), copyCounters
}
type MetricsSnapshot struct {
Ready bool
Dirty bool
Counters map[string]uint64
Gauges map[string]float64
Histograms map[string]DurationHistogramSnapshot
}
type DurationHistogramSnapshot struct {
Count uint64
Sum float64
Buckets [len(durationBuckets)]uint64
}
// MetricsSnapshot returns a copy suitable for rendering or testing. It never
// includes application identifiers or payloads.
func (health *Health) MetricsSnapshot() MetricsSnapshot {
if health == nil {
return MetricsSnapshot{Dirty: true}
}
health.mu.Lock()
defer health.mu.Unlock()
result := MetricsSnapshot{
Ready: health.ready.Load(), Dirty: health.dirty.Load(),
Counters: make(map[string]uint64, len(health.count)),
Gauges: make(map[string]float64, len(health.gauge)),
Histograms: make(map[string]DurationHistogramSnapshot, len(health.hist)),
}
for name, value := range health.count {
result.Counters[name] = value
}
for name, value := range health.gauge {
result.Gauges[name] = value
}
for name, value := range health.hist {
result.Histograms[name] = DurationHistogramSnapshot{Count: value.count, Sum: value.sum, Buckets: value.buckets}
}
return result
}
func (health *Health) Handler(paths Paths) http.Handler {
if paths.Liveness == "" {
paths.Liveness = "/livez"
}
if paths.Readiness == "" {
paths.Readiness = "/readyz"
}
if paths.Metrics == "" {
paths.Metrics = "/metrics"
}
mux := http.NewServeMux()
mux.HandleFunc(paths.Liveness, func(response http.ResponseWriter, _ *http.Request) {
response.Header().Set("Content-Type", "text/plain; charset=utf-8")
response.WriteHeader(http.StatusOK)
_, _ = response.Write([]byte("live\n"))
})
mux.HandleFunc(paths.Readiness, func(response http.ResponseWriter, _ *http.Request) {
ready, dirty, _ := health.Snapshot()
response.Header().Set("Content-Type", "text/plain; charset=utf-8")
if !ready {
response.WriteHeader(http.StatusServiceUnavailable)
} else {
response.WriteHeader(http.StatusOK)
}
_, _ = fmt.Fprintf(response, "ready=%s dirty=%s\n", strconv.FormatBool(ready), strconv.FormatBool(dirty))
})
mux.HandleFunc(paths.Metrics, func(response http.ResponseWriter, _ *http.Request) {
snapshot := health.MetricsSnapshot()
response.Header().Set("Content-Type", "text/plain; version=0.0.4")
response.WriteHeader(http.StatusOK)
_, _ = fmt.Fprintf(response, "rvbox_health_ready %d\nrvbox_health_dirty %d\n", boolMetric(snapshot.Ready), boolMetric(snapshot.Dirty))
keys := make([]string, 0, len(snapshot.Counters))
for key := range snapshot.Counters {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
_, _ = fmt.Fprintf(response, "rvbox_%s_total %d\n", key, snapshot.Counters[key])
}
keys = keys[:0]
for key := range snapshot.Gauges {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
_, _ = fmt.Fprintf(response, "rvbox_%s %g\n", key, snapshot.Gauges[key])
}
keys = keys[:0]
for key := range snapshot.Histograms {
keys = append(keys, key)
}
sort.Strings(keys)
for _, key := range keys {
histogram := snapshot.Histograms[key]
for index, upperBound := range durationBuckets {
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_bucket{le=\"%g\"} %d\n", key, upperBound, histogram.Buckets[index])
}
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_bucket{le=\"+Inf\"} %d\n", key, histogram.Count)
_, _ = fmt.Fprintf(response, "rvbox_%s_seconds_sum %g\nrvbox_%s_seconds_count %d\n", key, histogram.Sum, key, histogram.Count)
}
})
return mux
}
// Serve starts an endpoint and closes it when ctx is cancelled. The caller
// may use the returned error to keep listener failures visible without making
// product startup depend on a slow or unavailable metrics consumer.
func (health *Health) Serve(ctx context.Context, listen string, paths Paths) error {
if strings.TrimSpace(listen) == "" {
return fmt.Errorf("observability listen address is empty")
}
server := &http.Server{Addr: listen, Handler: health.Handler(paths), ReadHeaderTimeout: 10 * time.Second}
listener, err := net.Listen("tcp", listen)
if err != nil {
return err
}
go func() {
<-ctx.Done()
_ = server.Shutdown(context.Background())
}()
err = server.Serve(listener)
if err == http.ErrServerClosed {
return nil
}
return err
}
func boolMetric(value bool) int {
if value {
return 1
}
return 0
}
func validMetricName(name string) bool {
if name == "" {
return false
}
for index, character := range name {
if index == 0 && !(character == '_' || character >= 'a' && character <= 'z' || character >= 'A' && character <= 'Z') {
return false
}
if !(character == '_' || character >= 'a' && character <= 'z' || character >= 'A' && character <= 'Z' || character >= '0' && character <= '9') {
return false
}
}
return true
}