metrics.go 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117
  1. // client2server - Metrics aggregation
  2. //
  3. // In-memory time-bucketed counters for events and commands. Buckets are
  4. // 1-minute windows covering the last 24 hours (1440 buckets).
  5. package main
  6. import (
  7. "sync"
  8. "time"
  9. )
  10. const (
  11. bucketMinutes = 1
  12. retentionBuckets = 24 * 60 // 24h
  13. )
  14. type metricsBucket struct {
  15. start time.Time
  16. eventsTotal int
  17. byType map[string]int
  18. commandsTotal int
  19. commandsOK int
  20. commandsFail int
  21. }
  22. type Metrics struct {
  23. mu sync.RWMutex
  24. buckets [retentionBuckets]*metricsBucket
  25. }
  26. var globalMetrics = &Metrics{}
  27. func (m *Metrics) nowBucket() *metricsBucket {
  28. startOfMinute := time.Now().Truncate(bucketMinutes * time.Minute)
  29. idx := int(startOfMinute.Unix()/60) % retentionBuckets
  30. b := m.buckets[idx]
  31. if b == nil || !b.start.Equal(startOfMinute) {
  32. b = &metricsBucket{
  33. start: startOfMinute,
  34. byType: make(map[string]int),
  35. }
  36. m.buckets[idx] = b
  37. }
  38. return b
  39. }
  40. func (m *Metrics) recordEvent(eventType string) {
  41. m.mu.Lock()
  42. defer m.mu.Unlock()
  43. b := m.nowBucket()
  44. b.eventsTotal++
  45. b.byType[eventType]++
  46. }
  47. func (m *Metrics) recordCommand(success bool) {
  48. m.mu.Lock()
  49. defer m.mu.Unlock()
  50. b := m.nowBucket()
  51. b.commandsTotal++
  52. if success {
  53. b.commandsOK++
  54. } else {
  55. b.commandsFail++
  56. }
  57. }
  58. // snapshot returns a list of buckets covering the last `since` minutes.
  59. func (m *Metrics) snapshot(since time.Duration) []map[string]interface{} {
  60. m.mu.RLock()
  61. defer m.mu.RUnlock()
  62. cutoff := time.Now().Add(-since).Truncate(bucketMinutes * time.Minute)
  63. out := []map[string]interface{}{}
  64. // Walk buckets; for each one, check if its start >= cutoff
  65. for _, b := range m.buckets {
  66. if b == nil {
  67. continue
  68. }
  69. if b.start.Before(cutoff) {
  70. continue
  71. }
  72. out = append(out, map[string]interface{}{
  73. "ts": b.start.Format(time.RFC3339),
  74. "events_total": b.eventsTotal,
  75. "commands_total": b.commandsTotal,
  76. "commands_ok": b.commandsOK,
  77. "commands_fail": b.commandsFail,
  78. "events_by_type": b.byType,
  79. })
  80. }
  81. return out
  82. }
  83. func (m *Metrics) summary() map[string]interface{} {
  84. m.mu.RLock()
  85. defer m.mu.RUnlock()
  86. b := m.nowBucket()
  87. // Aggregate events from this minute and previous 5 minutes for a more
  88. // meaningful "recent" number.
  89. total := b.eventsTotal
  90. for i := 1; i < 6; i++ {
  91. idx := (int(time.Now().Unix()/60) - i) % retentionBuckets
  92. if idx < 0 {
  93. idx += retentionBuckets
  94. }
  95. if prev := m.buckets[idx]; prev != nil {
  96. total += prev.eventsTotal
  97. }
  98. }
  99. return map[string]interface{}{
  100. "events_last_5min": total,
  101. "events_this_min": b.eventsTotal,
  102. "commands_this_min": b.commandsTotal,
  103. }
  104. }