package main import ( "database/sql" "time" ) type AlertEvent struct { EventID int64 `json:"eventId"` ID string `json:"id"` Severity string `json:"severity"` Source string `json:"source"` Title string `json:"title"` Message string `json:"message"` FirstSeen int64 `json:"firstSeen"` LastSeen int64 `json:"lastSeen"` ResolvedAt *int64 `json:"resolvedAt,omitempty"` AcknowledgedAt *int64 `json:"acknowledgedAt,omitempty"` } type RecoveryEvent struct { EventID int64 ID, Severity, Source, Title, Message string } func (s *Store) SyncAlerts(alerts []Alert) ([]Alert, error) { tx, err := s.db.Begin() if err != nil { return alerts, err } defer tx.Rollback() now := time.Now().Unix() rows, err := tx.Query(`SELECT id,alert_key FROM alert_events WHERE resolved_at IS NULL`) if err != nil { return alerts, err } active := map[string]int64{} for rows.Next() { var id int64 var key string if err := rows.Scan(&id, &key); err != nil { rows.Close() return alerts, err } active[key] = id } rows.Close() seen := map[string]bool{} for _, alert := range alerts { seen[alert.ID] = true if id, ok := active[alert.ID]; ok { _, err = tx.Exec(`UPDATE alert_events SET severity=?,source=?,title=?,message=?,last_seen=? WHERE id=?`, alert.Severity, alert.Source, alert.Title, alert.Message, now, id) } else { result, insertErr := tx.Exec(`INSERT INTO alert_events(alert_key,severity,source,title,message,first_seen,last_seen) VALUES(?,?,?,?,?,?,?)`, alert.ID, alert.Severity, alert.Source, alert.Title, alert.Message, now, now) if insertErr != nil { return alerts, insertErr } active[alert.ID], _ = result.LastInsertId() } if err != nil { return alerts, err } } for key, id := range active { if !seen[key] { if _, err = tx.Exec(`UPDATE alert_events SET resolved_at=? WHERE id=?`, now, id); err != nil { return alerts, err } } } if _, err = tx.Exec(`DELETE FROM alert_events WHERE (resolved_at IS NOT NULL OR acknowledged_at IS NOT NULL) AND last_seen < ?`, now-int64(30*24*time.Hour/time.Second)); err != nil { return alerts, err } if err = tx.Commit(); err != nil { return alerts, err } rows, err = s.db.Query(`SELECT id,alert_key,severity,source,title,message FROM alert_events WHERE resolved_at IS NULL AND acknowledged_at IS NULL ORDER BY CASE severity WHEN 'critical' THEN 0 WHEN 'warning' THEN 1 ELSE 2 END,first_seen DESC`) if err != nil { return alerts, err } defer rows.Close() var visible []Alert for rows.Next() { var a Alert if err = rows.Scan(&a.EventID, &a.ID, &a.Severity, &a.Source, &a.Title, &a.Message); err != nil { return alerts, err } visible = append(visible, a) } return visible, rows.Err() } func (s *Store) AcknowledgeAlert(eventID int64) error { _, err := s.db.Exec(`UPDATE alert_events SET acknowledged_at=? WHERE id=? AND acknowledged_at IS NULL`, time.Now().Unix(), eventID) return err } func (s *Store) AlertHistory(limit int) ([]AlertEvent, error) { if limit < 1 || limit > 500 { limit = 200 } rows, err := s.db.Query(`SELECT id,alert_key,severity,source,title,message,first_seen,last_seen,resolved_at,acknowledged_at FROM alert_events ORDER BY first_seen DESC LIMIT ?`, limit) if err != nil { return nil, err } defer rows.Close() var result []AlertEvent for rows.Next() { var e AlertEvent var resolved, acked sql.NullInt64 if err = rows.Scan(&e.EventID, &e.ID, &e.Severity, &e.Source, &e.Title, &e.Message, &e.FirstSeen, &e.LastSeen, &resolved, &acked); err != nil { return nil, err } if resolved.Valid { e.ResolvedAt = &resolved.Int64 } if acked.Valid { e.AcknowledgedAt = &acked.Int64 } result = append(result, e) } return result, rows.Err() } func (s *Store) PendingEmailAlerts() ([]Alert, error) { rows, err := s.db.Query(`SELECT id,alert_key,severity,source,title,message FROM alert_events WHERE resolved_at IS NULL AND acknowledged_at IS NULL AND email_sent_at IS NULL ORDER BY first_seen`) if err != nil { return nil, err } defer rows.Close() var result []Alert for rows.Next() { var a Alert if err = rows.Scan(&a.EventID, &a.ID, &a.Severity, &a.Source, &a.Title, &a.Message); err != nil { return nil, err } result = append(result, a) } return result, rows.Err() } func (s *Store) MarkAlertsEmailed(alerts []Alert) error { tx, err := s.db.Begin() if err != nil { return err } defer tx.Rollback() now := time.Now().Unix() for _, a := range alerts { if _, err = tx.Exec(`UPDATE alert_events SET email_sent_at=? WHERE id=?`, now, a.EventID); err != nil { return err } } return tx.Commit() } func (s *Store) PendingImmediateAlerts() ([]Alert, error) { rows, err := s.db.Query(`SELECT id,alert_key,severity,source,title,message FROM alert_events WHERE resolved_at IS NULL AND acknowledged_at IS NULL AND email_sent_at IS NULL AND (severity='critical' OR alert_key='ups-battery') ORDER BY first_seen`) if err != nil { return nil, err } defer rows.Close() var out []Alert for rows.Next() { var a Alert if err = rows.Scan(&a.EventID, &a.ID, &a.Severity, &a.Source, &a.Title, &a.Message); err != nil { return nil, err } out = append(out, a) } return out, rows.Err() } func (s *Store) PendingRecoveries() ([]RecoveryEvent, error) { rows, err := s.db.Query(`SELECT id,alert_key,severity,source,title,message FROM alert_events WHERE resolved_at IS NOT NULL AND email_sent_at IS NOT NULL AND recovery_sent_at IS NULL ORDER BY resolved_at`) if err != nil { return nil, err } defer rows.Close() var out []RecoveryEvent for rows.Next() { var v RecoveryEvent if err = rows.Scan(&v.EventID, &v.ID, &v.Severity, &v.Source, &v.Title, &v.Message); err != nil { return nil, err } out = append(out, v) } return out, rows.Err() } func (s *Store) MarkRecoveriesEmailed(items []RecoveryEvent) error { tx, err := s.db.Begin() if err != nil { return err } defer tx.Rollback() now := time.Now().Unix() for _, v := range items { if _, err = tx.Exec(`UPDATE alert_events SET recovery_sent_at=? WHERE id=?`, now, v.EventID); err != nil { return err } } return tx.Commit() }