// 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 }