package db import ( "context" "fmt" "strings" "time" ) func (s *Store) DashboardOverview(limit int) (map[string]any, error) { return s.DashboardOverviewWindow(limit, "") } func (s *Store) DashboardOverviewWindow(limit int, since string) (map[string]any, error) { if limit <= 0 || limit > 200 { limit = 80 } var feedbackTotal, feedbackToday, sourceTotal, sourceVisible, releaseTotal, mailFailed int today := time.Now().UTC().Format("2006-01-02") + "T00:00:00Z" if err := s.queryRow(`SELECT (SELECT COUNT(*) FROM feedback_tickets), (SELECT COUNT(*) FROM feedback_tickets WHERE created_at >= ?), (SELECT COUNT(*) FROM source_endpoints), (SELECT COUNT(*) FROM source_endpoints WHERE enabled = 1 AND client_visible = 1), (SELECT COUNT(*) FROM release_notices), (SELECT COUNT(*) FROM mail_records WHERE status = 'failed')`, today).Scan( &feedbackTotal, &feedbackToday, &sourceTotal, &sourceVisible, &releaseTotal, &mailFailed, ); err != nil { return nil, err } warnings := []string{} statusCounts, err := s.groupCounts("feedback_tickets", "status") if err != nil { warnings = append(warnings, "feedback status unavailable: "+err.Error()) statusCounts = map[string]int{} } healthCounts, err := s.groupCounts("source_endpoints", "last_status") if err != nil { warnings = append(warnings, "source health unavailable: "+err.Error()) healthCounts = map[string]int{} } recentChecks, err := s.RecentSourceChecksWindow(limit, since) if err != nil { warnings = append(warnings, "source checks unavailable: "+err.Error()) recentChecks = []map[string]any{} } recentCalls, err := s.RecentSourceCallsWindow(limit, since) if err != nil { warnings = append(warnings, "client calls unavailable: "+err.Error()) recentCalls = []map[string]any{} } averageLatency, err := s.AverageSourceLatencyBucketsWindow(limit, since) if err != nil { warnings = append(warnings, "latency trend unavailable: "+err.Error()) averageLatency = []map[string]any{} } audit, err := s.ListAuditLogs(10) if err != nil { warnings = append(warnings, "audit summary unavailable: "+err.Error()) audit = []AuditLog{} } sourceRows, err := s.DashboardSourceRows() if err != nil { warnings = append(warnings, "source summary unavailable: "+err.Error()) sourceRows = []map[string]any{} } return map[string]any{ "ok": true, "kpis": map[string]any{ "feedbackTotal": feedbackTotal, "feedbackToday": feedbackToday, "sourceTotal": sourceTotal, "sourceVisible": sourceVisible, "releaseNotices": releaseTotal, "mailFailed": mailFailed, }, "feedbackStatus": statusCounts, "sourceHealth": healthCounts, "heartbeats": recentChecks, "averageLatency": averageLatency, "clientCalls": recentCalls, "database": s.Status(), "audit": audit, "sourceRows": sourceRows, "generatedAt": time.Now().UTC().Format(time.RFC3339), "warnings": warnings, }, nil } func (s *Store) AverageSourceLatencyBuckets(limit int) ([]map[string]any, error) { return s.AverageSourceLatencyBucketsWindow(limit, "") } func (s *Store) AverageSourceLatencyBucketsWindow(limit int, since string) ([]map[string]any, error) { if limit <= 0 || limit > 200 { limit = 80 } where := "" args := []any{} if strings.TrimSpace(since) != "" { where = " WHERE checked_at >= ?" args = append(args, since) } args = append(args, limit*4) rows, err := s.query(`SELECT checked_at, latency_ms, status FROM endpoint_health_checks`+where+` ORDER BY checked_at DESC, id DESC LIMIT ?`, args...) if err != nil { return nil, err } defer rows.Close() type bucket struct { label string total int count int ok int latest string } order := []string{} buckets := map[string]*bucket{} for rows.Next() { var checkedAt, status string var latency int if err := rows.Scan(&checkedAt, &latency, &status); err != nil { return nil, err } label := latencyBucketLabel(checkedAt) if label == "" { label = checkedAt } item, ok := buckets[label] if !ok { item = &bucket{label: label, latest: checkedAt} buckets[label] = item order = append(order, label) } item.total += latency item.count++ if status == "ok" || status == "redirected" { item.ok++ } if checkedAt > item.latest { item.latest = checkedAt } } if err := rows.Err(); err != nil { return nil, err } out := []map[string]any{} for i := len(order) - 1; i >= 0; i-- { item := buckets[order[i]] if item == nil || item.count == 0 { continue } out = append(out, map[string]any{ "label": item.label, "averageLatency": item.total / item.count, "avgLatencyMs": item.total / item.count, "sampleCount": item.count, "healthyCount": item.ok, "checkedAt": item.latest, }) } if len(out) > limit { out = out[len(out)-limit:] } return out, nil } func latencyBucketLabel(value string) string { if value == "" { return "" } parsed, err := time.Parse(time.RFC3339, value) if err != nil { if len(value) >= 16 { return value[:16] } return value } return parsed.UTC().Format("01-02 15:04") } func (s *Store) RecentSourceChecks(limit int) ([]map[string]any, error) { return s.RecentSourceChecksWindow(limit, "") } func (s *Store) RecentSourceChecksWindow(limit int, since string) ([]map[string]any, error) { where := "" args := []any{} if strings.TrimSpace(since) != "" { where = " WHERE h.checked_at >= ?" args = append(args, since) } args = append(args, limit) rows, err := s.query(`SELECT h.id, h.source_db_id, COALESCE(e.source_id, ''), COALESCE(e.name, ''), h.status, h.latency_ms, h.error, h.checked_at FROM endpoint_health_checks h LEFT JOIN source_endpoints e ON e.id = h.source_db_id `+where+` ORDER BY h.checked_at DESC, h.id DESC LIMIT ?`, args...) if err != nil { return nil, err } defer rows.Close() items := []map[string]any{} for rows.Next() { var id, sourceDBID int64 var sourceID, name, status, message, checkedAt string var latency int if err := rows.Scan(&id, &sourceDBID, &sourceID, &name, &status, &latency, &message, &checkedAt); err != nil { return nil, err } if sourceID == "" { sourceID = fmt.Sprintf("deleted-%d", sourceDBID) } if name == "" { name = fmt.Sprintf("已删除接口 #%d", sourceDBID) } items = append(items, map[string]any{"id": id, "sourceDbId": sourceDBID, "sourceId": sourceID, "name": name, "status": status, "latencyMs": latency, "error": message, "checkedAt": checkedAt}) } return items, rows.Err() } func (s *Store) RecentSourceCalls(limit int) ([]map[string]any, error) { return s.RecentSourceCallsWindow(limit, "") } func (s *Store) RecentSourceCallsWindow(limit int, since string) ([]map[string]any, error) { where := "" args := []any{} if strings.TrimSpace(since) != "" { where = " WHERE created_at >= ?" args = append(args, since) } args = append(args, limit) rows, err := s.query(`SELECT id, source_id, status, latency_ms, error, client, created_at FROM endpoint_call_logs`+where+` ORDER BY created_at DESC, id DESC LIMIT ?`, args...) if err != nil { return nil, err } defer rows.Close() items := []map[string]any{} for rows.Next() { var id int64 var sourceID, status, message, client, createdAt string var latency int if err := rows.Scan(&id, &sourceID, &status, &latency, &message, &client, &createdAt); err != nil { return nil, err } items = append(items, map[string]any{"id": id, "sourceId": sourceID, "status": status, "latencyMs": latency, "error": message, "client": client, "createdAt": createdAt}) } return items, rows.Err() } func (s *Store) DashboardSourceRows() ([]map[string]any, error) { rows, err := s.query(`SELECT source_id, category_id, category_name, name, enabled, client_visible, last_status, last_latency_ms, last_checked_at, last_error, consecutive_failure FROM source_endpoints ORDER BY category_id ASC, name ASC`) if err != nil { return nil, err } defer rows.Close() items := []map[string]any{} for rows.Next() { var sourceID, categoryID, categoryName, name, status, checkedAt, lastError string var enabled, visible, latency, failures int if err := rows.Scan(&sourceID, &categoryID, &categoryName, &name, &enabled, &visible, &status, &latency, &checkedAt, &lastError, &failures); err != nil { return nil, err } items = append(items, map[string]any{ "sourceId": sourceID, "categoryId": categoryID, "categoryName": categoryName, "name": name, "enabled": enabled == 1, "clientVisible": visible == 1, "status": status, "latencyMs": latency, "checkedAt": checkedAt, "healthError": lastError, "consecutiveFailure": failures, }) } return items, rows.Err() } func (s *Store) InsertAudit(log AuditLog) error { return s.InsertAuditContext(context.Background(), log) } func (s *Store) InsertAuditContext(ctx context.Context, log AuditLog) error { if log.CreatedAt == "" { log.CreatedAt = Now() } conn, d := s.active() if conn == nil { return fmt.Errorf("database is not available") } _, err := conn.ExecContext(ctx, d.rebind(`INSERT INTO audit_logs (actor, type, target, message, ip, user_agent, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)`), sanitize(log.Actor), sanitize(log.Type), sanitize(log.Target), sanitize(log.Message), sanitize(log.IP), sanitize(log.UserAgent), log.CreatedAt) return err } func (s *Store) ListAuditLogs(limit int) ([]AuditLog, error) { if limit <= 0 || limit > 200 { limit = 100 } rows, err := s.query(`SELECT id, actor, type, target, message, ip, user_agent, created_at FROM audit_logs ORDER BY id DESC LIMIT ?`, limit) if err != nil { return nil, err } defer rows.Close() return scanAuditRows(rows) } func (s *Store) ListAuditLogsPage(filters AuditFilters) (AuditPage, error) { page := filters.Page if page <= 0 { page = 1 } perPage := filters.PerPage if perPage <= 0 { perPage = 35 } if perPage > 100 { perPage = 100 } where, args := auditWhere(filters) var total int if err := s.queryRow(`SELECT COUNT(*) FROM audit_logs`+where, args...).Scan(&total); err != nil { return AuditPage{}, err } offset := (page - 1) * perPage queryArgs := append(append([]any{}, args...), perPage, offset) rows, err := s.query(`SELECT id, actor, type, target, message, ip, user_agent, created_at FROM audit_logs`+where+` ORDER BY id DESC LIMIT ? OFFSET ?`, queryArgs...) if err != nil { return AuditPage{}, err } defer rows.Close() items, err := scanAuditRows(rows) if err != nil { return AuditPage{}, err } return AuditPage{Items: items, Total: total, Page: page, PerPage: perPage}, nil } func auditWhere(filters AuditFilters) (string, []any) { clauses := []string{} args := []any{} if value := strings.TrimSpace(filters.Type); value != "" { clauses = append(clauses, "type = ?") args = append(args, sanitize(value)) } if value := strings.TrimSpace(filters.Target); value != "" { clauses = append(clauses, "target = ?") args = append(args, sanitize(value)) } if value := strings.TrimSpace(filters.Query); value != "" { clauses = append(clauses, "(actor LIKE ? OR type LIKE ? OR target LIKE ? OR message LIKE ? OR ip LIKE ?)") like := "%" + sanitize(value) + "%" args = append(args, like, like, like, like, like) } if len(clauses) == 0 { return "", args } return " WHERE " + strings.Join(clauses, " AND "), args } func (s *Store) ListAuditLogsForTarget(target string, limit int) ([]AuditLog, error) { if limit <= 0 || limit > 200 { limit = 100 } rows, err := s.query(`SELECT id, actor, type, target, message, ip, user_agent, created_at FROM audit_logs WHERE target = ? ORDER BY id DESC LIMIT ?`, target, limit) if err != nil { return nil, err } defer rows.Close() return scanAuditRows(rows) } func (s *Store) countTable(table string) (int, error) { if !validStatsTable(table) { return 0, fmt.Errorf("invalid table %q", table) } var total int err := s.queryRow(`SELECT COUNT(*) FROM ` + table).Scan(&total) return total, err } func (s *Store) countWhere(table, where string, args ...any) (int, error) { if !validStatsTable(table) { return 0, fmt.Errorf("invalid table %q", table) } var total int err := s.queryRow(`SELECT COUNT(*) FROM `+table+` WHERE `+where, args...).Scan(&total) return total, err } func (s *Store) groupCounts(table, column string) (map[string]int, error) { if !validStatsColumn(table, column) { return nil, fmt.Errorf("invalid group %s.%s", table, column) } rows, err := s.query(`SELECT ` + column + `, COUNT(*) FROM ` + table + ` GROUP BY ` + column) if err != nil { return nil, err } defer rows.Close() out := map[string]int{} for rows.Next() { var key string var count int if err := rows.Scan(&key, &count); err != nil { return nil, err } if key == "" { key = "unknown" } out[key] = count } return out, rows.Err() } func validStatsTable(table string) bool { switch table { case "feedback_tickets", "source_endpoints", "release_notices", "mail_records": return true default: return false } } func validStatsColumn(table, column string) bool { switch table + "." + column { case "feedback_tickets.status", "source_endpoints.last_status": return true default: return false } }