safesight-control/internal/service/alarm_collector.go

382 lines
12 KiB
Go

package service
import (
"database/sql"
"encoding/json"
"fmt"
"log"
"sync"
"time"
"safesight-control/internal/models"
)
type AlarmRecord struct {
ID string `json:"id"`
DeviceID string `json:"device_id"`
Channel string `json:"channel"`
Timestamp string `json:"timestamp"`
RuleName string `json:"rule_name"`
RuleType string `json:"rule_type"`
ObjectLabel string `json:"object_label"`
Confidence float64 `json:"confidence"`
SnapshotURL string `json:"snapshot_url"`
ClipURL string `json:"clip_url"`
DurationMs int64 `json:"duration_ms"`
CollectedAt string `json:"collected_at"`
Status string `json:"status"`
Severity string `json:"severity"`
AcknowledgedAt string `json:"acknowledged_at,omitempty"`
AcknowledgedBy string `json:"acknowledged_by,omitempty"`
ResolvedAt string `json:"resolved_at,omitempty"`
ResolvedBy string `json:"resolved_by,omitempty"`
}
type AlarmCollector struct {
db *sql.DB
agent *AgentClient
registry *RegistryService
mu sync.Mutex
lastID string
retentionDays int
subscribers map[chan AlarmRecord]struct{}
}
func NewAlarmCollector(db *sql.DB, agent *AgentClient, registry *RegistryService) *AlarmCollector {
return &AlarmCollector{db: db, agent: agent, registry: registry, subscribers: make(map[chan AlarmRecord]struct{})}
}
func (c *AlarmCollector) Start() {
go c.poll()
go c.cleanupLoop()
}
// SetRetention configures alarm retention in days. 0 = keep forever.
func (c *AlarmCollector) SetRetention(days int) {
c.retentionDays = days
}
func (c *AlarmCollector) cleanupLoop() {
c.cleanupOldAlarms()
ticker := time.NewTicker(24 * time.Hour)
defer ticker.Stop()
for range ticker.C {
c.cleanupOldAlarms()
}
}
func (c *AlarmCollector) cleanupOldAlarms() {
if c.db == nil || c.retentionDays <= 0 {
return
}
cutoff := time.Now().AddDate(0, 0, -c.retentionDays).Format(time.RFC3339)
_, err := c.db.Exec(`DELETE FROM alarm_records WHERE collected_at < ?`, cutoff)
if err != nil {
log.Printf("alarm cleanup: %v", err)
}
}
func (c *AlarmCollector) poll() {
// Initial delay to let the system settle
time.Sleep(5 * time.Second)
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
for range ticker.C {
c.collectFromDevices()
}
}
func (c *AlarmCollector) collectFromDevices() {
if c.agent == nil || c.registry == nil {
return
}
devices := c.registry.GetDevices()
for _, dev := range devices {
if dev == nil || !dev.Online {
continue
}
alarms, err := c.fetchDeviceAlarms(dev)
if err != nil {
continue
}
for _, alarm := range alarms {
alarm.DeviceID = dev.DeviceID
alarm.CollectedAt = time.Now().Format(time.RFC3339)
c.saveAlarm(alarm)
}
}
}
func (c *AlarmCollector) fetchDeviceAlarms(dev *models.Device) ([]AlarmRecord, error) {
c.mu.Lock()
lastID := c.lastID
c.mu.Unlock()
url := "/v1/alarms/recent?limit=100"
body, status, err := c.agent.Do("GET", dev.IP, dev.AgentPort, url, nil)
if err != nil {
return nil, err
}
if status != 200 {
return nil, nil // agent may not support alarms yet
}
var resp struct {
Alarms []map[string]any `json:"alarms"`
}
if err := json.Unmarshal(body, &resp); err != nil {
return nil, err
}
newLastID := lastID
alarms := make([]AlarmRecord, 0)
for _, a := range resp.Alarms {
id, _ := a["id"].(string)
if id == "" {
continue
}
if id == lastID {
break // reached previously seen alarms
}
// Extract fields from whatever format media-server sends
channel, _ := a["channel"].(string)
if channel == "" {
channel, _ = a["node_id"].(string) // media-server uses node_id
}
ruleName, _ := a["rule_name"].(string)
ruleType, _ := a["rule_type"].(string)
objectLabel, _ := a["object_label"].(string)
snapshotURL, _ := a["snapshot_url"].(string)
clipURL, _ := a["clip_url"].(string)
// Timestamp can be string (RFC3339) or number (unix millis)
var ts string
switch v := a["timestamp"].(type) {
case string:
ts = v
case float64:
ts = time.UnixMilli(int64(v)).Format(time.RFC3339)
}
confidence, _ := a["confidence"].(float64)
if confidence == 0 {
if score, ok := a["score"].(float64); ok {
confidence = score
}
}
// Extract from nested detections if top-level fields are empty
if dets, ok := a["detections"].([]any); ok && len(dets) > 0 {
if objectLabel == "" {
objectLabel = fmt.Sprintf("%d 个检测目标", len(dets))
}
if confidence == 0 {
if d0, ok := dets[0].(map[string]any); ok {
if s, ok := d0["score"].(float64); ok {
confidence = s
}
}
}
}
durationMs, _ := a["duration_ms"].(float64)
alarms = append(alarms, AlarmRecord{
ID: id,
Timestamp: ts,
Channel: channel,
RuleName: ruleName,
RuleType: ruleType,
ObjectLabel: objectLabel,
Confidence: confidence,
SnapshotURL: snapshotURL,
ClipURL: clipURL,
DurationMs: int64(durationMs),
})
if newLastID == "" {
newLastID = id
}
}
if newLastID != "" {
c.mu.Lock()
c.lastID = newLastID
c.mu.Unlock()
}
return alarms, nil
}
func (c *AlarmCollector) saveAlarm(alarm AlarmRecord) {
if c.db == nil {
return
}
_, err := c.db.Exec(`
INSERT INTO alarm_records(id, device_id, channel, timestamp, rule_name, rule_type, object_label, confidence, snapshot_url, clip_url, duration_ms, collected_at, status, severity)
VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO NOTHING
`, alarm.ID, alarm.DeviceID, alarm.Channel, alarm.Timestamp, alarm.RuleName, alarm.RuleType, alarm.ObjectLabel, alarm.Confidence, alarm.SnapshotURL, alarm.ClipURL, alarm.DurationMs, alarm.CollectedAt, "unacknowledged", "warning")
if err != nil {
log.Printf("alarm collector: save error: %v", err)
}
c.broadcast(alarm)
}
// GetFiltered returns alarm records within a date range, with pagination and optional filters.
// deviceID and ruleType filter values, empty string means no filter.
func (c *AlarmCollector) GetFiltered(from, to string, limit, offset int, deviceID, ruleType string) ([]AlarmRecord, int) {
if c.db == nil {
return nil, 0
}
fromTS := from + "T00:00:00"
toTS := to + "T23:59:59"
cond := "WHERE timestamp >= ? AND timestamp <= ?"
args := []any{fromTS, toTS}
if deviceID != "" {
cond += " AND device_id = ?"
args = append(args, deviceID)
}
if ruleType != "" {
cond += " AND rule_type = ?"
args = append(args, ruleType)
}
var total int
c.db.QueryRow(`SELECT COUNT(*) FROM alarm_records `+cond, args...).Scan(&total)
fullQuery := `SELECT id, device_id, channel, timestamp, rule_name, rule_type, object_label, confidence, snapshot_url, clip_url, duration_ms, collected_at, status, severity, acknowledged_at, acknowledged_by, resolved_at, resolved_by
FROM alarm_records ` + cond + `
ORDER BY timestamp DESC
LIMIT ? OFFSET ?`
fullArgs := append(args, limit, offset)
rows, err := c.db.Query(fullQuery, fullArgs...)
if err != nil {
return nil, total
}
defer rows.Close()
var alarms []AlarmRecord
for rows.Next() {
var a AlarmRecord
if err := rows.Scan(&a.ID, &a.DeviceID, &a.Channel, &a.Timestamp, &a.RuleName, &a.RuleType, &a.ObjectLabel, &a.Confidence, &a.SnapshotURL, &a.ClipURL, &a.DurationMs, &a.CollectedAt, &a.Status, &a.Severity, &a.AcknowledgedAt, &a.AcknowledgedBy, &a.ResolvedAt, &a.ResolvedBy); err != nil {
continue
}
alarms = append(alarms, a)
}
return alarms, total
}
// GetRecent returns the most recent N alarm records, newest first.
func (c *AlarmCollector) GetRecent(limit int) []AlarmRecord {
if c.db == nil || limit <= 0 {
return nil
}
rows, err := c.db.Query(`
SELECT id, device_id, channel, timestamp, rule_name, rule_type, object_label, confidence, snapshot_url, clip_url, duration_ms, collected_at, status, severity, acknowledged_at, acknowledged_by, resolved_at, resolved_by
FROM alarm_records
ORDER BY timestamp DESC
LIMIT ?
`, limit)
if err != nil {
return nil
}
defer rows.Close()
var alarms []AlarmRecord
for rows.Next() {
var a AlarmRecord
if err := rows.Scan(&a.ID, &a.DeviceID, &a.Channel, &a.Timestamp, &a.RuleName, &a.RuleType, &a.ObjectLabel, &a.Confidence, &a.SnapshotURL, &a.ClipURL, &a.DurationMs, &a.CollectedAt, &a.Status, &a.Severity, &a.AcknowledgedAt, &a.AcknowledgedBy, &a.ResolvedAt, &a.ResolvedBy); err != nil {
continue
}
alarms = append(alarms, a)
}
return alarms
}
func (c *AlarmCollector) Subscribe() chan AlarmRecord {
ch := make(chan AlarmRecord, 32)
c.mu.Lock()
c.subscribers[ch] = struct{}{}
c.mu.Unlock()
return ch
}
func (c *AlarmCollector) Unsubscribe(ch chan AlarmRecord) {
c.mu.Lock()
delete(c.subscribers, ch)
c.mu.Unlock()
}
func (c *AlarmCollector) broadcast(alarm AlarmRecord) {
c.mu.Lock()
defer c.mu.Unlock()
for ch := range c.subscribers {
select {
case ch <- alarm:
default:
}
}
}
func (c *AlarmCollector) AcknowledgeAlarm(id, actor string) error {
if c.db == nil {
return fmt.Errorf("db not available")
}
_, err := c.db.Exec(`UPDATE alarm_records SET status='acknowledged', acknowledged_at=?, acknowledged_by=? WHERE id=?`,
time.Now().Format(time.RFC3339), actor, id)
return err
}
func (c *AlarmCollector) ResolveAlarm(id, actor string) error {
if c.db == nil {
return fmt.Errorf("db not available")
}
_, err := c.db.Exec(`UPDATE alarm_records SET status='resolved', resolved_at=?, resolved_by=? WHERE id=?`,
time.Now().Format(time.RFC3339), actor, id)
return err
}
func (c *AlarmCollector) UpdateSeverity(id, severity string) error {
if c.db == nil {
return fmt.Errorf("db not available")
}
_, err := c.db.Exec(`UPDATE alarm_records SET severity=? WHERE id=?`, severity, id)
return err
}
func (c *AlarmCollector) GetDistinctRuleTypes() []string {
if c.db == nil {
return nil
}
rows, err := c.db.Query(`SELECT DISTINCT rule_type FROM alarm_records WHERE rule_type != '' ORDER BY rule_type`)
if err != nil {
return nil
}
defer rows.Close()
var types []string
for rows.Next() {
var t string
if err := rows.Scan(&t); err == nil {
types = append(types, t)
}
}
return types
}
func (c *AlarmCollector) GetDistinctDeviceIDs() []string {
if c.db == nil {
return nil
}
rows, err := c.db.Query(`SELECT DISTINCT device_id FROM alarm_records ORDER BY device_id`)
if err != nil {
return nil
}
defer rows.Close()
var ids []string
for rows.Next() {
var id string
if err := rows.Scan(&id); err == nil {
ids = append(ids, id)
}
}
return ids
}