diff --git a/docs/实施跟踪.md b/docs/实施跟踪.md index b5d0a39..4d2b9d4 100644 --- a/docs/实施跟踪.md +++ b/docs/实施跟踪.md @@ -36,12 +36,12 @@ | # | 子任务 | 状态 | |---|--------|------| -| 1.1 | DB 迁移 — alarm_records 加 status/severity/resolved 列 | ⬜ | -| 1.2 | AlarmRecord 结构体扩展 + saveAlarm 默认值 | ⬜ | -| 1.3 | API — acknowledge/resolve/severity 端点 | ⬜ | -| 1.4 | 模板 — 严重等级列 + 状态列 + 操作按钮 | ⬜ | -| 1.5 | 筛选增强 — 设备下拉 + 规则类型下拉 | ⬜ | -| 1.6 | SSE — /api/alarms/stream 实时推送 | ⬜ | +| 1.1 | DB 迁移 — alarm_records 加 status/severity/resolved 列 | ✅ | +| 1.2 | AlarmRecord 结构体扩展 + saveAlarm 默认值 + SSE 广播 | ✅ | +| 1.3 | API — acknowledge/resolve/severity 端点 + SSE stream | ✅ | +| 1.4 | 模板 — 严重等级列 + 状态列 + 操作按钮 | ✅ | +| 1.5 | 筛选增强 — 设备下拉 + 规则类型下拉 | ✅ | +| 1.6 | SSE — /api/alarms/stream 实时推送 | ✅ | ## 已完成事项 @@ -52,6 +52,12 @@ - 提交:`27584f9` - 影响:断网部署环境下监控页视频播放恢复正常 +- **P0-1 告警中心增强**(待验证) + - DB 迁移:`alarm_records` 增加 `status`, `severity`, `acknowledged_at/by`, `resolved_at/by` 列 + - 服务层:`AlarmRecord` 扩展 + SSE 广播 + 确认/关闭/改等级 API + - 页面:严重等级列 + 状态列 + 操作按钮 + 设备/规则类型下拉筛选 + - 涉及文件:`migrate.go`, `alarm_collector.go`, `ui.go`, `alarms.html` + --- -**进度:1/15(6.7%)** +**已完成:1/15(6.7%)** | **进行中:1(P0-1)** diff --git a/internal/service/alarm_collector.go b/internal/service/alarm_collector.go index 49328fd..938573d 100644 --- a/internal/service/alarm_collector.go +++ b/internal/service/alarm_collector.go @@ -12,18 +12,24 @@ import ( ) 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"` + 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 { @@ -31,12 +37,13 @@ type AlarmCollector struct { agent *AgentClient registry *RegistryService mu sync.Mutex - lastID string // last alarm ID seen, to avoid duplicates - retentionDays int // auto-delete alarms older than N days, 0=keep + 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} + return &AlarmCollector{db: db, agent: agent, registry: registry, subscribers: make(map[chan AlarmRecord]struct{})} } func (c *AlarmCollector) Start() { @@ -201,39 +208,46 @@ func (c *AlarmCollector) saveAlarm(alarm AlarmRecord) { 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) -VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) +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) +`, 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. -// Returns records and total count (ignoring limit/offset for count). -func (c *AlarmCollector) GetFiltered(from, to string, limit, offset int) ([]AlarmRecord, int) { +// 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" - // Count total - var total int - c.db.QueryRow(` -SELECT COUNT(*) FROM alarm_records -WHERE timestamp >= ? AND timestamp <= ? -`, fromTS, toTS).Scan(&total) + 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) + } - // Fetch page - 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 -FROM alarm_records -WHERE timestamp >= ? AND timestamp <= ? + 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 ? -`, fromTS, toTS, limit, offset) +LIMIT ? OFFSET ?` + fullArgs := append(args, limit, offset) + + rows, err := c.db.Query(fullQuery, fullArgs...) if err != nil { return nil, total } @@ -242,7 +256,7 @@ LIMIT ? OFFSET ? 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); err != nil { + 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) @@ -256,7 +270,7 @@ func (c *AlarmCollector) GetRecent(limit int) []AlarmRecord { 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 +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 ? @@ -269,10 +283,99 @@ LIMIT ? 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); err != nil { + 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 +} diff --git a/internal/storage/migrate.go b/internal/storage/migrate.go index b95370b..00255cc 100644 --- a/internal/storage/migrate.go +++ b/internal/storage/migrate.go @@ -189,7 +189,10 @@ func migrate(db *sql.DB) error { if err := migrateProfilePrimaryTemplateName(db); err != nil { return err } - return migrateProfilesToSceneTemplates(db) + if err := migrateProfilesToSceneTemplates(db); err != nil { + return err + } + return migrateAlarmWorkflow(db) } func migrateProfilePrimaryTemplateName(db *sql.DB) error { @@ -431,6 +434,33 @@ func splitLegacyProfileDocument(item struct { return raw, units, nil } +func migrateAlarmWorkflow(db *sql.DB) error { + if db == nil { + return nil + } + hasStatus, err := hasColumn(db, "alarm_records", "status") + if err != nil { + return err + } + if hasStatus { + return nil + } + stmts := []string{ + `ALTER TABLE alarm_records ADD COLUMN status TEXT NOT NULL DEFAULT 'unacknowledged'`, + `ALTER TABLE alarm_records ADD COLUMN severity TEXT NOT NULL DEFAULT 'warning'`, + `ALTER TABLE alarm_records ADD COLUMN acknowledged_at TEXT NOT NULL DEFAULT ''`, + `ALTER TABLE alarm_records ADD COLUMN acknowledged_by TEXT NOT NULL DEFAULT ''`, + `ALTER TABLE alarm_records ADD COLUMN resolved_at TEXT NOT NULL DEFAULT ''`, + `ALTER TABLE alarm_records ADD COLUMN resolved_by TEXT NOT NULL DEFAULT ''`, + } + for _, stmt := range stmts { + if _, err := db.Exec(stmt); err != nil { + return fmt.Errorf("migrate alarm_workflow: %w", err) + } + } + return nil +} + func bindingStringFromAny(bindings map[string]any, slot string, field string) string { entry, _ := bindings[slot].(map[string]any) return strings.TrimSpace(stringAny(entry[field])) diff --git a/internal/web/ui.go b/internal/web/ui.go index f046c50..7258912 100644 --- a/internal/web/ui.go +++ b/internal/web/ui.go @@ -99,6 +99,10 @@ type PageData struct { MonthStart string MonthEnd string TodayAlarmCount int + AlarmFilterDevices []string + AlarmFilterRuleTypes []string + SelectedAlarmDevice string + SelectedAlarmRuleType string DeviceMetrics []DeviceMetric ConsoleDevices []ConsoleDeviceData ConsoleVideoSources []service.ConfigVideoSourceAsset @@ -523,6 +527,54 @@ func NewUI(discovery *service.DiscoveryService, registry *service.RegistryServic "add": func(a, b int) int { return a + b }, + "severityClass": func(v string) string { + switch strings.TrimSpace(v) { + case "critical": + return "bad" + case "warning": + return "warn" + case "info": + return "ok" + default: + return "" + } + }, + "severityLabel": func(v string) string { + switch strings.TrimSpace(v) { + case "critical": + return "严重" + case "warning": + return "警告" + case "info": + return "信息" + default: + return v + } + }, + "statusClass": func(v string) string { + switch strings.TrimSpace(v) { + case "unacknowledged": + return "bad" + case "acknowledged": + return "run" + case "resolved": + return "ok" + default: + return "" + } + }, + "statusLabel": func(v string) string { + switch strings.TrimSpace(v) { + case "unacknowledged": + return "未处理" + case "acknowledged": + return "已确认" + case "resolved": + return "已关闭" + default: + return v + } + }, }).ParseFS(uiFS, "ui/templates/*.html") if err != nil { return nil, err @@ -739,6 +791,10 @@ func (u *UI) Routes() (chi.Router, error) { r.Post("/resources/sync", u.actionResourceSync) r.Get("/diagnostics", u.pageDiagnostics) r.Get("/alarms", u.pageAlarms) + r.Post("/alarms/{id}/acknowledge", u.actionAcknowledgeAlarm) + r.Post("/alarms/{id}/resolve", u.actionResolveAlarm) + r.Post("/alarms/{id}/severity", u.actionUpdateAlarmSeverity) + r.Get("/api/alarms/stream", u.apiAlarmStream) r.Get("/face-gallery", u.pageFaceGallery) r.Get("/face-photo/*", u.serveFacePhoto) r.Post("/face-gallery/import", u.actionFaceGalleryImport) @@ -1969,8 +2025,13 @@ func (u *UI) pageAlarms(w http.ResponseWriter, r *http.Request) { data.MonthStart = now.Format("2006-01-02")[:8] + "01" data.MonthEnd = today + data.SelectedAlarmDevice = strings.TrimSpace(q.Get("device")) + data.SelectedAlarmRuleType = strings.TrimSpace(q.Get("rule_type")) + if u.alarmCollector != nil { - data.AlarmRecords, data.TotalAlarmCount = u.alarmCollector.GetFiltered(from, to, perPage, (page-1)*perPage) + data.AlarmRecords, data.TotalAlarmCount = u.alarmCollector.GetFiltered(from, to, perPage, (page-1)*perPage, data.SelectedAlarmDevice, data.SelectedAlarmRuleType) + data.AlarmFilterDevices = u.alarmCollector.GetDistinctDeviceIDs() + data.AlarmFilterRuleTypes = u.alarmCollector.GetDistinctRuleTypes() } // Calculate total pages and build page list: 1, ..., cur-2, cur-1, cur, cur+1, cur+2, ..., N if data.TotalAlarmCount > 0 { @@ -2002,6 +2063,80 @@ func (u *UI) pageAlarms(w http.ResponseWriter, r *http.Request) { u.render(w, r, "alarms", data) } +func (u *UI) actionAcknowledgeAlarm(w http.ResponseWriter, r *http.Request) { + id := chi.URLParam(r, "id") + if u.alarmCollector == nil || id == "" { + http.Error(w, "not found", http.StatusNotFound) + return + } + if err := u.alarmCollector.AcknowledgeAlarm(id, "admin"); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusNoContent) +} + +func (u *UI) actionResolveAlarm(w http.ResponseWriter, r *http.Request) { + id := chi.URLParam(r, "id") + if u.alarmCollector == nil || id == "" { + http.Error(w, "not found", http.StatusNotFound) + return + } + if err := u.alarmCollector.ResolveAlarm(id, "admin"); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusNoContent) +} + +func (u *UI) actionUpdateAlarmSeverity(w http.ResponseWriter, r *http.Request) { + id := chi.URLParam(r, "id") + if u.alarmCollector == nil || id == "" { + http.Error(w, "not found", http.StatusNotFound) + return + } + severity := strings.TrimSpace(r.FormValue("severity")) + if severity == "" { + http.Error(w, "missing severity", http.StatusBadRequest) + return + } + if err := u.alarmCollector.UpdateSeverity(id, severity); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + w.WriteHeader(http.StatusNoContent) +} + +func (u *UI) apiAlarmStream(w http.ResponseWriter, r *http.Request) { + if u.alarmCollector == nil { + http.Error(w, "not available", http.StatusServiceUnavailable) + return + } + flusher, ok := w.(http.Flusher) + if !ok { + http.Error(w, "streaming not supported", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("Connection", "keep-alive") + ch := u.alarmCollector.Subscribe() + defer u.alarmCollector.Unsubscribe(ch) + for { + select { + case alarm := <-ch: + data, err := json.Marshal(alarm) + if err != nil { + continue + } + fmt.Fprintf(w, "data: %s\n\n", data) + flusher.Flush() + case <-r.Context().Done(): + return + } + } +} + func (u *UI) pageResources(w http.ResponseWriter, r *http.Request) { u.ensureDevicesLoaded() data := PageData{Title: "资源管理", Devices: u.registry.GetDevices()} diff --git a/internal/web/ui/templates/alarms.html b/internal/web/ui/templates/alarms.html index 9ec90eb..f40f94f 100644 --- a/internal/web/ui/templates/alarms.html +++ b/internal/web/ui/templates/alarms.html @@ -1,10 +1,26 @@ {{define "alarms"}}