| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117 |
- // client2server - Metrics aggregation
- //
- // In-memory time-bucketed counters for events and commands. Buckets are
- // 1-minute windows covering the last 24 hours (1440 buckets).
- package main
- import (
- "sync"
- "time"
- )
- const (
- bucketMinutes = 1
- retentionBuckets = 24 * 60 // 24h
- )
- type metricsBucket struct {
- start time.Time
- eventsTotal int
- byType map[string]int
- commandsTotal int
- commandsOK int
- commandsFail int
- }
- type Metrics struct {
- mu sync.RWMutex
- buckets [retentionBuckets]*metricsBucket
- }
- var globalMetrics = &Metrics{}
- func (m *Metrics) nowBucket() *metricsBucket {
- startOfMinute := time.Now().Truncate(bucketMinutes * time.Minute)
- idx := int(startOfMinute.Unix()/60) % retentionBuckets
- b := m.buckets[idx]
- if b == nil || !b.start.Equal(startOfMinute) {
- b = &metricsBucket{
- start: startOfMinute,
- byType: make(map[string]int),
- }
- m.buckets[idx] = b
- }
- return b
- }
- func (m *Metrics) recordEvent(eventType string) {
- m.mu.Lock()
- defer m.mu.Unlock()
- b := m.nowBucket()
- b.eventsTotal++
- b.byType[eventType]++
- }
- func (m *Metrics) recordCommand(success bool) {
- m.mu.Lock()
- defer m.mu.Unlock()
- b := m.nowBucket()
- b.commandsTotal++
- if success {
- b.commandsOK++
- } else {
- b.commandsFail++
- }
- }
- // snapshot returns a list of buckets covering the last `since` minutes.
- func (m *Metrics) snapshot(since time.Duration) []map[string]interface{} {
- m.mu.RLock()
- defer m.mu.RUnlock()
- cutoff := time.Now().Add(-since).Truncate(bucketMinutes * time.Minute)
- out := []map[string]interface{}{}
- // Walk buckets; for each one, check if its start >= cutoff
- for _, b := range m.buckets {
- if b == nil {
- continue
- }
- if b.start.Before(cutoff) {
- continue
- }
- out = append(out, map[string]interface{}{
- "ts": b.start.Format(time.RFC3339),
- "events_total": b.eventsTotal,
- "commands_total": b.commandsTotal,
- "commands_ok": b.commandsOK,
- "commands_fail": b.commandsFail,
- "events_by_type": b.byType,
- })
- }
- return out
- }
- func (m *Metrics) summary() map[string]interface{} {
- m.mu.RLock()
- defer m.mu.RUnlock()
- b := m.nowBucket()
- // Aggregate events from this minute and previous 5 minutes for a more
- // meaningful "recent" number.
- total := b.eventsTotal
- for i := 1; i < 6; i++ {
- idx := (int(time.Now().Unix()/60) - i) % retentionBuckets
- if idx < 0 {
- idx += retentionBuckets
- }
- if prev := m.buckets[idx]; prev != nil {
- total += prev.eventsTotal
- }
- }
- return map[string]interface{}{
- "events_last_5min": total,
- "events_this_min": b.eventsTotal,
- "commands_this_min": b.commandsTotal,
- }
- }
|