package controlplane import ( "database/sql" "errors" "fmt" "strings" "time" ) func backfillAuditEventFields(db *sql.DB) error { _, err := db.Exec(`UPDATE audit_events SET operator = CASE WHEN operator = '' THEN actor ELSE operator END, result = CASE WHEN result = '' THEN outcome ELSE result END, security_detail = CASE WHEN security_detail = '' THEN detail ELSE security_detail END, resource_type = CASE WHEN resource_type = '' OR resource_type = 'system' THEN CASE WHEN action LIKE 'agent.%' THEN 'agent' WHEN action LIKE 'pairing.%' THEN 'pairing' WHEN action LIKE 'llm.%' THEN 'llm' WHEN action LIKE 'skill.%' THEN 'skill' WHEN action LIKE 'mcp.%' THEN 'mcp' WHEN action LIKE 'task.%' THEN 'task' ELSE 'system' END ELSE resource_type END, target_agent_id = CASE WHEN target_agent_id != '' THEN target_agent_id WHEN actor LIKE 'agent:%' THEN substr(actor, 7) WHEN target LIKE 'agent:%' THEN CASE WHEN instr(substr(target, 7), ':') = 0 THEN substr(target, 7) ELSE substr(substr(target, 7), 1, instr(substr(target, 7), ':') - 1) END ELSE '' END`) if err != nil { return fmt.Errorf("无法升级审计记录: %w", err) } return nil } // AuditEvent is a safe local operational record. Legacy fields remain for // compatibility; unified fields power the new server-filtered audit view. type AuditEvent struct { ID int64 `json:"id"` CreatedAt string `json:"createdAt"` Operator string `json:"operator"` TargetAgentID string `json:"targetAgentID"` ResourceType string `json:"resourceType"` Action string `json:"action"` Result string `json:"result"` SecurityDetail string `json:"securityDetail"` Actor string `json:"actor"` Target string `json:"target"` Outcome string `json:"outcome"` Detail string `json:"detail"` } type AuditFilter struct { Page int `json:"page"` PageSize int `json:"pageSize"` AgentID string `json:"agentID"` ResourceType string `json:"resourceType"` Result string `json:"result"` FromAt string `json:"fromAt"` ToAt string `json:"toAt"` } type AuditPage struct { Items []AuditEvent `json:"items"` Page int `json:"page"` PageSize int `json:"pageSize"` TotalItems int `json:"totalItems"` TotalPages int `json:"totalPages"` } var auditResourceTypes = map[string]struct{}{ "agent": {}, "pairing": {}, "llm": {}, "skill": {}, "mcp": {}, "task": {}, "template": {}, "schedule": {}, "delegation": {}, "plugin": {}, "remote-operation": {}, "system": {}, } func (c *ControlPlane) recordAudit(action, actor, target, outcome, detail string) { resourceType := auditResourceType(action) targetAgentID := auditTargetAgentID(actor, target) safeDetail := limit(sanitizeDisplayText(detail), 1000) operator := limit(sanitizeDisplayText(actor), 120) result := limit(sanitizeDisplayText(outcome), 40) _, _ = c.db.Exec(`INSERT INTO audit_events( action, actor, target, outcome, detail, created_at, operator, target_agent_id, resource_type, result, security_detail ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, limit(action, 80), operator, limit(sanitizeDisplayText(target), 160), result, safeDetail, time.Now().UTC().Format(time.RFC3339Nano), operator, targetAgentID, resourceType, result, safeDetail, ) } func auditResourceType(action string) string { prefix, _, found := strings.Cut(strings.TrimSpace(action), ".") if !found { return "system" } switch prefix { case "agent": return "agent" case "pairing": return "pairing" case "llm": return "llm" case "skill": return "skill" case "mcp": return "mcp" case "task": return "task" case "template": return "template" case "schedule": return "schedule" case "delegation": return "delegation" case "plugin": return "plugin" case "remote": return "remote-operation" default: return "system" } } func auditTargetAgentID(actor, target string) string { for _, value := range []string{target, actor} { if !strings.HasPrefix(value, "agent:") { continue } value = strings.TrimPrefix(value, "agent:") if agentID, _, _ := strings.Cut(value, ":"); agentID != "" { return limit(agentID, 80) } } return "" } func normalizeAuditFilter(filter AuditFilter) (AuditFilter, error) { filter.AgentID, filter.ResourceType = strings.TrimSpace(filter.AgentID), strings.TrimSpace(filter.ResourceType) filter.Result, filter.FromAt, filter.ToAt = strings.TrimSpace(filter.Result), strings.TrimSpace(filter.FromAt), strings.TrimSpace(filter.ToAt) if filter.Page < 1 { filter.Page = 1 } if filter.PageSize < 1 { filter.PageSize = 25 } if filter.PageSize > 100 { filter.PageSize = 100 } if filter.ResourceType != "" { if _, ok := auditResourceTypes[filter.ResourceType]; !ok { return filter, errors.New("审计资源类型无效") } } if len(filter.AgentID) > 80 || len(filter.Result) > 40 { return filter, errors.New("审计筛选条件过长") } for _, value := range []string{filter.FromAt, filter.ToAt} { if value != "" { if _, err := time.Parse(time.RFC3339, value); err != nil { return filter, errors.New("审计时间必须使用 RFC3339 格式") } } } if filter.FromAt != "" && filter.ToAt != "" && filter.FromAt > filter.ToAt { return filter, errors.New("审计起始时间不能晚于结束时间") } return filter, nil } func auditWhere(filter AuditFilter) (string, []any) { clauses, args := make([]string, 0, 5), make([]any, 0, 5) for _, criterion := range []struct{ column, value string }{{"target_agent_id", filter.AgentID}, {"resource_type", filter.ResourceType}, {"result", filter.Result}} { if criterion.value != "" { clauses, args = append(clauses, criterion.column+" = ?"), append(args, criterion.value) } } if filter.FromAt != "" { clauses, args = append(clauses, "created_at >= ?"), append(args, filter.FromAt) } if filter.ToAt != "" { clauses, args = append(clauses, "created_at <= ?"), append(args, filter.ToAt) } if len(clauses) == 0 { return "", args } return " WHERE " + strings.Join(clauses, " AND "), args } // AuditTrail provides a paged, server-filtered unified audit view. func (c *ControlPlane) AuditTrail(filter AuditFilter) (AuditPage, error) { filter, err := normalizeAuditFilter(filter) if err != nil { return AuditPage{}, err } where, args := auditWhere(filter) page := AuditPage{Items: make([]AuditEvent, 0), Page: filter.Page, PageSize: filter.PageSize} if err := c.db.QueryRow("SELECT COUNT(*) FROM audit_events"+where, args...).Scan(&page.TotalItems); err != nil { return page, fmt.Errorf("无法统计审计记录: %w", err) } page.TotalPages = (page.TotalItems + filter.PageSize - 1) / filter.PageSize if page.TotalPages > 0 && page.Page > page.TotalPages { page.Page = page.TotalPages } queryArgs := append(append([]any{}, args...), page.PageSize, (page.Page-1)*page.PageSize) rows, err := c.db.Query(`SELECT event_id, created_at, operator, target_agent_id, resource_type, action, result, security_detail, actor, target, outcome, detail FROM audit_events`+where+" ORDER BY event_id DESC LIMIT ? OFFSET ?", queryArgs...) if err != nil { return page, fmt.Errorf("无法读取审计记录: %w", err) } defer rows.Close() for rows.Next() { var event AuditEvent if err := rows.Scan(&event.ID, &event.CreatedAt, &event.Operator, &event.TargetAgentID, &event.ResourceType, &event.Action, &event.Result, &event.SecurityDetail, &event.Actor, &event.Target, &event.Outcome, &event.Detail); err != nil { return page, err } event.SecurityDetail, event.Detail = sanitizeDisplayText(event.SecurityDetail), sanitizeDisplayText(event.Detail) page.Items = append(page.Items, event) } return page, rows.Err() } // RecentAuditEvents is retained for existing callers and returns the newest 100 safe events. func (c *ControlPlane) RecentAuditEvents() ([]AuditEvent, error) { page, err := c.AuditTrail(AuditFilter{Page: 1, PageSize: 100}) return page.Items, err }