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