package observability import ( "fmt" "io" "sort" "sync" "time" ) // Metrics is a bounded in-process collector for API request health. Operation // names are normalized to a fixed vocabulary before storage. type Metrics struct { mu sync.Mutex counts map[metricKey]uint64 sums map[metricKey]time.Duration buckets map[metricKey][]uint64 } type metricKey struct{ operation, status string } // apiLatencyBucketsSeconds is deliberately fixed and small. It is wide enough // to query the documented 250 ms API SLO while keeping the exporter bounded. var apiLatencyBucketsSeconds = []float64{0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2, 5, 10} func NewMetrics() *Metrics { return &Metrics{counts: make(map[metricKey]uint64), sums: make(map[metricKey]time.Duration), buckets: make(map[metricKey][]uint64)} } func (m *Metrics) ObserveAPI(operation string, statusCode int, duration time.Duration) { if m == nil { return } if duration < 0 { duration = 0 } key := metricKey{normalizeOperation(operation), statusClass(statusCode)} m.mu.Lock() m.counts[key]++ m.sums[key] += duration bucketCounts := m.buckets[key] if bucketCounts == nil { bucketCounts = make([]uint64, len(apiLatencyBucketsSeconds)) m.buckets[key] = bucketCounts } seconds := duration.Seconds() for index, upperBound := range apiLatencyBucketsSeconds { if seconds <= upperBound { bucketCounts[index]++ } } m.mu.Unlock() } func (m *Metrics) WritePrometheus(w io.Writer) error { if m == nil { return nil } m.mu.Lock() keys := make([]metricKey, 0, len(m.counts)) for key := range m.counts { keys = append(keys, key) } sort.Slice(keys, func(i, j int) bool { if keys[i].operation != keys[j].operation { return keys[i].operation < keys[j].operation } return keys[i].status < keys[j].status }) counts := make(map[metricKey]uint64, len(keys)) sums := make(map[metricKey]time.Duration, len(keys)) buckets := make(map[metricKey][]uint64, len(keys)) for _, key := range keys { counts[key], sums[key] = m.counts[key], m.sums[key] buckets[key] = append([]uint64(nil), m.buckets[key]...) } m.mu.Unlock() if _, err := io.WriteString(w, "# TYPE cosmic_clash_api_requests_total counter\n# TYPE cosmic_clash_api_latency_seconds histogram\n"); err != nil { return err } for _, key := range keys { labels := fmt.Sprintf(`operation="%s",status="%s"`, key.operation, key.status) for index, upperBound := range apiLatencyBucketsSeconds { if _, err := fmt.Fprintf(w, "cosmic_clash_api_latency_seconds_bucket{%s,le=\"%g\"} %d\n", labels, upperBound, buckets[key][index]); err != nil { return err } } if _, err := fmt.Fprintf(w, "cosmic_clash_api_latency_seconds_bucket{%s,le=\"+Inf\"} %d\ncosmic_clash_api_requests_total{%s} %d\ncosmic_clash_api_latency_seconds_count{%s} %d\ncosmic_clash_api_latency_seconds_sum{%s} %.9f\n", labels, counts[key], labels, counts[key], labels, counts[key], labels, sums[key].Seconds()); err != nil { return err } } return nil } func normalizeOperation(operation string) string { for _, allowed := range []string{"queue", "proposal", "assignment", "profile", "ranked_profile", "server", "events", "session", "probe"} { if operation == allowed { return allowed } } return "other" } func statusClass(code int) string { if code < 100 || code > 599 { return "unknown" } return fmt.Sprintf("%dxx", code/100) }