diff --git a/cmd/evobgp-all/main.go b/cmd/evobgp-all/main.go index 3e71dc0..8576144 100644 --- a/cmd/evobgp-all/main.go +++ b/cmd/evobgp-all/main.go @@ -42,6 +42,7 @@ func main() { AuthPortalURL: firstNonEmpty(os.Getenv("EVOBGP_AUTH_PORTAL_URL"), os.Getenv("AUTH_PORTAL_URL")), PortalTenantID: strings.TrimSpace(os.Getenv("EVOBGP_PORTAL_TENANT_ID")), AuthRequired: boolFromEnv("EVOBGP_AUTH_REQUIRED", "AUTH_REQUIRED"), + AuditIngestSecret: firstNonEmpty(os.Getenv("EVOBGP_AUTH_AUDIT_INGEST_SECRET"), os.Getenv("AUTH_AUDIT_INGEST_SECRET")), } srv, err := httpapi.New(opts) if err != nil { diff --git a/cmd/evobgp-api/main.go b/cmd/evobgp-api/main.go index d292931..6258950 100644 --- a/cmd/evobgp-api/main.go +++ b/cmd/evobgp-api/main.go @@ -37,6 +37,7 @@ func main() { AuthPortalURL: firstNonEmpty(os.Getenv("EVOBGP_AUTH_PORTAL_URL"), os.Getenv("AUTH_PORTAL_URL")), PortalTenantID: strings.TrimSpace(os.Getenv("EVOBGP_PORTAL_TENANT_ID")), AuthRequired: boolFromEnv("EVOBGP_AUTH_REQUIRED", "AUTH_REQUIRED"), + AuditIngestSecret: firstNonEmpty(os.Getenv("EVOBGP_AUTH_AUDIT_INGEST_SECRET"), os.Getenv("AUTH_AUDIT_INGEST_SECRET")), } srv, err := httpapi.New(opts) if err != nil { diff --git a/docs/access.md b/docs/access.md index 2877e79..73bb3eb 100644 --- a/docs/access.md +++ b/docs/access.md @@ -12,6 +12,7 @@ | `AUTH_JWT_SECRET` / `EVOBGP_AUTH_JWT_SECRET` | Тот же секрет, что `JWT_SECRET` портала (HS256) | | `AUTH_ISSUER` | Issuer JWT (как на портале) | | `AUTH_PORTAL_URL` | URL портала (также `GET /v1/auth/config`) | +| `AUTH_AUDIT_INGEST_SECRET` / `EVOBGP_AUTH_AUDIT_INGEST_SECRET` | Shared secret для push CRUD audit в auth-portal (`POST /api/v1/ingest/audit`, `source_app=bgp`) | | `EVOBGP_PORTAL_TENANT_ID` | Fallback tenant для portal JWT, если в токене нет `bgp_tenant_id` / `tenants.bgp` | Источник tenant (по приоритету): diff --git a/docs/api.md b/docs/api.md index 5080db7..108a75e 100644 --- a/docs/api.md +++ b/docs/api.md @@ -116,6 +116,16 @@ `{filename}` — только basename, паттерн `^[a-z0-9][a-z0-9_.-]*\.log$`. Очистка пишет строку в таблицу `runtime_log_cleanup_audit` (миграция `000026`). +## CRUD audit (`/v1/audit`) + +Локальный журнал изменений CRUD (modules, peers, settings, API keys, …). Миграция `000030_audit_log`. Чтение — `bgp:monitoring:read` (viewer+). + +| Метод | Путь | Роль | Назначение | +|-------|------|------|------------| +| `GET` | `/v1/audit` | viewer+ | Пагинированный audit (`cursor`, `limit`, опционально `action`, `severity`) | + +При `AUTH_PORTAL_URL` + `AUTH_AUDIT_INGEST_SECRET` каждая запись дополнительно отправляется в auth-portal (`POST /api/v1/ingest/audit`, `source_app=bgp`). + ## Соглашения из OpenAPI - Ошибки в стиле **RFC 9457** (`application/problem+json`): `type`, `title`, `status`, `detail`, и т.д. diff --git a/docs/openapi.yaml b/docs/openapi.yaml index b849fbf..56ca745 100644 --- a/docs/openapi.yaml +++ b/docs/openapi.yaml @@ -59,6 +59,8 @@ tags: description: Сессия текущего API-ключа (tenant и роль). - name: Monitoring description: Наблюдаемость PostgreSQL и корреляция (instance-level, viewer+). Maintenance — operator. + - name: Audit + description: Журнал CRUD-изменений tenant (локально + опциональный push в auth-portal). Чтение — bgp:monitoring:read. - name: Maintenance description: Политики обслуживания PostgreSQL (instance-scoped). CRUD и запуск — operator. - name: RuntimeLogs @@ -1297,6 +1299,70 @@ components: has_more: type: boolean + AuditSeverity: + type: string + enum: [info, warning, critical] + + AuditLogEntry: + type: object + required: + [id, tenant_id, event_id, source_app, action, severity, summary, created_at] + properties: + id: + $ref: "#/components/schemas/ResourceId" + tenant_id: + $ref: "#/components/schemas/ResourceId" + event_id: + type: string + description: Stable id for portal ingest deduplication (prefix bgp-). + source_app: + type: string + enum: [bgp] + action: + type: string + description: Machine action key (e.g. bgp.module.create). + severity: + $ref: "#/components/schemas/AuditSeverity" + actor_user_id: + type: ["string", "null"] + actor_email: + type: ["string", "null"] + actor_name: + type: ["string", "null"] + actor_api_key_prefix: + type: ["string", "null"] + target_type: + type: ["string", "null"] + enum: [app_resource, null] + target_id: + type: ["string", "null"] + summary: + type: string + details: + type: ["object", "null"] + additionalProperties: true + ip: + type: ["string", "null"] + created_at: + type: string + format: date-time + portal_pushed_at: + type: ["string", "null"] + format: date-time + + AuditLogList: + type: object + required: [items] + properties: + items: + type: array + items: + $ref: "#/components/schemas/AuditLogEntry" + next_cursor: + type: string + has_more: + type: boolean + RuntimeLogAutoPolicy: type: object properties: @@ -4520,6 +4586,40 @@ paths: default: $ref: "#/components/responses/DefaultProblem" + /v1/audit: + get: + tags: [Audit] + summary: Журнал CRUD audit tenant + description: | + Локальный журнал изменений (modules, peers, settings, API keys и т.д.). + При настроенных `AUTH_PORTAL_URL` + `AUTH_AUDIT_INGEST_SECRET` события также + отправляются в auth-portal ingest (`source_app=bgp`). + operationId: listAuditLog + parameters: + - $ref: "#/components/parameters/TenantId" + - $ref: "#/components/parameters/Cursor" + - $ref: "#/components/parameters/Limit" + - name: action + in: query + schema: + type: string + description: Filter by action prefix/key (exact match). + - name: severity + in: query + schema: + $ref: "#/components/schemas/AuditSeverity" + responses: + "200": + description: Успешно. + content: + application/json: + schema: + $ref: "#/components/schemas/AuditLogList" + "400": + $ref: "#/components/responses/BadRequest" + default: + $ref: "#/components/responses/DefaultProblem" + /v1/settings: get: tags: [Settings] diff --git a/internal/audit/portal.go b/internal/audit/portal.go new file mode 100644 index 0000000..57cdcb7 --- /dev/null +++ b/internal/audit/portal.go @@ -0,0 +1,118 @@ +// Package audit pushes local audit events to auth-portal ingest API. +package audit + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "log" + "net/http" + "strings" + "time" + + "evobgp/internal/httpclient" + "evobgp/internal/store" +) + +const ingestPath = "/api/v1/ingest/audit" + +// PortalPusher sends audit rows to auth-portal (best-effort, async-friendly). +type PortalPusher struct { + BaseURL string + Secret string + HTTPClient *http.Client + MarkPushed func(id string) error +} + +// PushEvent posts one audit entry to portal ingest. +func (p *PortalPusher) PushEvent(ctx context.Context, entry *store.AuditEntry) error { + if p == nil || entry == nil { + return nil + } + base := strings.TrimRight(strings.TrimSpace(p.BaseURL), "/") + secret := strings.TrimSpace(p.Secret) + if base == "" || secret == "" { + return nil + } + hc := p.HTTPClient + if hc == nil { + hc = httpclient.New(15 * time.Second) + } + body := map[string]any{ + "events": []map[string]any{p.eventPayload(entry)}, + } + raw, err := json.Marshal(body) + if err != nil { + return fmt.Errorf("audit: marshal ingest: %w", err) + } + req, err := http.NewRequestWithContext(ctx, http.MethodPost, base+ingestPath, bytes.NewReader(raw)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + req.Header.Set("Authorization", "Bearer "+secret) + resp, err := hc.Do(req) + if err != nil { + return fmt.Errorf("audit: portal ingest: %w", err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode >= 300 { + b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("audit: portal ingest %s: %s", resp.Status, strings.TrimSpace(string(b))) + } + if p.MarkPushed != nil { + if err := p.MarkPushed(entry.ID); err != nil { + log.Printf("audit: mark portal pushed id=%s: %v", entry.ID, err) + } + } + return nil +} + +func (p *PortalPusher) eventPayload(entry *store.AuditEntry) map[string]any { + ev := map[string]any{ + "event_id": entry.EventID, + "source_app": store.AuditSourceAppBGP, + "action": entry.Action, + "severity": entry.Severity, + "summary": entry.Summary, + "created_at": entry.CreatedAt.UTC().Format(time.RFC3339Nano), + } + if entry.ActorUserID != "" { + ev["actor_user_id"] = entry.ActorUserID + } else { + ev["actor_user_id"] = nil + } + if entry.ActorEmail != "" { + ev["actor_email"] = entry.ActorEmail + } else { + ev["actor_email"] = nil + } + if entry.ActorName != "" { + ev["actor_name"] = entry.ActorName + } else { + ev["actor_name"] = nil + } + if entry.TargetType != "" { + ev["target_type"] = entry.TargetType + } else { + ev["target_type"] = nil + } + if entry.TargetID != "" { + ev["target_id"] = entry.TargetID + } else { + ev["target_id"] = nil + } + if entry.Details != nil { + ev["details"] = entry.Details + } else { + ev["details"] = nil + } + if entry.IP != "" { + ev["ip"] = entry.IP + } else { + ev["ip"] = nil + } + return ev +} diff --git a/internal/audit/portal_test.go b/internal/audit/portal_test.go new file mode 100644 index 0000000..fafeb07 --- /dev/null +++ b/internal/audit/portal_test.go @@ -0,0 +1,63 @@ +package audit + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "evobgp/internal/store" +) + +func TestPortalPusherPushEvent(t *testing.T) { + var got struct { + Events []map[string]any `json:"events"` + } + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != ingestPath { + t.Fatalf("path=%s", r.URL.Path) + } + if r.Header.Get("Authorization") != "Bearer test-secret" { + t.Fatalf("auth=%q", r.Header.Get("Authorization")) + } + _ = json.NewDecoder(r.Body).Decode(&got) + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(map[string]int{"accepted": 1, "duplicates": 0}) + })) + defer srv.Close() + + marked := false + p := &PortalPusher{ + BaseURL: srv.URL, + Secret: "test-secret", + MarkPushed: func(id string) error { + marked = id == "local-id" + return nil + }, + } + entry := &store.AuditEntry{ + ID: "local-id", + EventID: "bgp-test-event", + Action: "bgp.module.create", + Severity: store.AuditSeverityInfo, + Summary: "Created module", + SourceApp: store.AuditSourceAppBGP, + CreatedAt: time.Now().UTC(), + TargetType: store.AuditTargetAppResource, + TargetID: "mod-1", + } + if err := p.PushEvent(context.Background(), entry); err != nil { + t.Fatal(err) + } + if len(got.Events) != 1 { + t.Fatalf("events=%d", len(got.Events)) + } + if got.Events[0]["source_app"] != "bgp" { + t.Fatalf("source_app=%v", got.Events[0]["source_app"]) + } + if !marked { + t.Fatal("expected mark pushed") + } +} diff --git a/internal/httpapi/routes.go b/internal/httpapi/routes.go index 8896b73..23ec82e 100644 --- a/internal/httpapi/routes.go +++ b/internal/httpapi/routes.go @@ -83,6 +83,7 @@ func (s *Server) registerV1(m *http.ServeMux) { m.HandleFunc("GET /speakers/{speaker_id}/bundle/{revision_id}", s.handleNodeBundle) m.HandleFunc("POST /nodes/enroll", s.handleNodeEnroll) s.registerCRUDRoutes(m) + s.registerAuditRoutes(m) s.registerPostgresMonitoringRoutes(m) s.registerPostgresMaintenanceRoutes(m) s.registerMaintenanceRoutes(m) diff --git a/internal/httpapi/routes_api_keys.go b/internal/httpapi/routes_api_keys.go index 170d976..734f2d0 100644 --- a/internal/httpapi/routes_api_keys.go +++ b/internal/httpapi/routes_api_keys.go @@ -131,6 +131,7 @@ func (s *Server) handlePostAPIKey(w http.ResponseWriter, r *http.Request) { } out := apiKeyJSON(&created.APIKey) out["token"] = created.Token + s.recordCRUDAudit(r, a, "bgp.api_key.create", "Created API key "+created.Name, created.ID, map[string]any{"api_key_id": created.ID, "role": created.Role}) writeJSON(w, http.StatusCreated, out) } @@ -187,6 +188,7 @@ func (s *Server) handlePatchAPIKey(w http.ResponseWriter, r *http.Request) { writeProblem(w, http.StatusInternalServerError, "Internal Server Error", "failed to reload api keys") return } + s.recordCRUDAudit(r, a, "bgp.api_key.update", "Updated API key "+k.Name, k.ID, map[string]any{"api_key_id": k.ID, "role": k.Role}) writeJSON(w, http.StatusOK, apiKeyJSON(k)) } @@ -195,7 +197,8 @@ func (s *Server) handleDeleteAPIKey(w http.ResponseWriter, r *http.Request) { if !ok || !s.requirePerm(w, a, "bgp:access:admin") { return } - if err := s.store.RevokeAPIKey(a.TenantID, r.PathValue("id")); err != nil { + keyID := r.PathValue("id") + if err := s.store.RevokeAPIKey(a.TenantID, keyID); err != nil { writeStoreErr(w, err) return } @@ -203,6 +206,7 @@ func (s *Server) handleDeleteAPIKey(w http.ResponseWriter, r *http.Request) { writeProblem(w, http.StatusInternalServerError, "Internal Server Error", "failed to reload api keys") return } + s.recordCRUDAudit(r, a, "bgp.api_key.revoke", "Revoked API key", keyID, map[string]any{"api_key_id": keyID}) w.WriteHeader(http.StatusNoContent) } @@ -222,5 +226,6 @@ func (s *Server) handleRotateAPIKey(w http.ResponseWriter, r *http.Request) { } out := apiKeyJSON(&rotated.APIKey) out["token"] = rotated.Token + s.recordCRUDAudit(r, a, "bgp.api_key.rotate", "Rotated API key "+rotated.Name, rotated.ID, map[string]any{"api_key_id": rotated.ID}) writeJSON(w, http.StatusOK, out) } diff --git a/internal/httpapi/routes_audit.go b/internal/httpapi/routes_audit.go new file mode 100644 index 0000000..b00ceca --- /dev/null +++ b/internal/httpapi/routes_audit.go @@ -0,0 +1,171 @@ +package httpapi + +import ( + "context" + "log" + "net" + "net/http" + "strings" + "time" + + "evobgp/internal/audit" + "evobgp/internal/store" +) + +func (s *Server) registerAuditRoutes(m *http.ServeMux) { + m.HandleFunc("GET /audit", s.handleListAudit) +} + +func (s *Server) handleListAudit(w http.ResponseWriter, r *http.Request) { + a, ok := authFromContext(r.Context()) + if !ok || !s.requirePerm(w, a, "bgp:monitoring:read") { + return + } + cursor := r.URL.Query().Get("cursor") + limit := parseLimitQuery(r, 20, 200) + filter := store.AuditListFilter{ + Action: strings.TrimSpace(r.URL.Query().Get("action")), + Severity: strings.TrimSpace(r.URL.Query().Get("severity")), + } + if filter.Severity != "" && !store.ValidAuditSeverity(filter.Severity) { + writeProblem(w, http.StatusBadRequest, "Bad Request", "invalid severity") + return + } + items, next, hasMore, err := s.store.ListAudit(a.TenantID, cursor, limit, filter) + if err != nil { + writeInternalError(w, "audit_list", err) + return + } + out := make([]map[string]any, 0, len(items)) + for _, row := range items { + out = append(out, auditEntryJSON(row)) + } + writeJSON(w, http.StatusOK, map[string]any{"items": out, "next_cursor": next, "has_more": hasMore}) +} + +func auditEntryJSON(row *store.AuditEntry) map[string]any { + if row == nil { + return map[string]any{} + } + m := map[string]any{ + "id": row.ID, + "tenant_id": row.TenantID, + "event_id": row.EventID, + "source_app": row.SourceApp, + "action": row.Action, + "severity": row.Severity, + "actor_user_id": strPtrOrNull(row.ActorUserID), + "actor_email": strPtrOrNull(row.ActorEmail), + "actor_name": strPtrOrNull(row.ActorName), + "actor_api_key_prefix": strPtrOrNull(row.ActorAPIKeyPrefix), + "target_type": strPtrOrNull(row.TargetType), + "target_id": strPtrOrNull(row.TargetID), + "summary": row.Summary, + "details": row.Details, + "ip": strPtrOrNull(row.IP), + "created_at": row.CreatedAt.UTC().Format(time.RFC3339Nano), + "portal_pushed_at": nil, + } + if row.PortalPushedAt != nil { + m["portal_pushed_at"] = row.PortalPushedAt.UTC().Format(time.RFC3339Nano) + } + if m["details"] == nil { + m["details"] = nil + } + return m +} + +func (s *Server) recordCRUDAudit(r *http.Request, a Auth, action, summary, targetID string, details map[string]any) { + if s == nil || s.store == nil { + return + } + in := store.AuditAppendInput{ + TenantID: a.TenantID, + Action: action, + Severity: store.AuditSeverityInfo, + TargetType: store.AuditTargetAppResource, + TargetID: targetID, + Summary: summary, + Details: details, + IP: clientIP(r), + } + fillAuditActor(&in, a) + entry, err := s.store.AppendAudit(in) + if err != nil { + log.Printf("httpapi: audit append action=%s: %v", action, err) + return + } + s.pushAuditToPortal(entry) +} + +func fillAuditActor(in *store.AuditAppendInput, a Auth) { + if in == nil { + return + } + if a.Kind == AuthKindJWT { + in.ActorUserID = strings.TrimSpace(a.UserID) + in.ActorEmail = strings.TrimSpace(a.Email) + if in.ActorEmail != "" { + in.ActorName = in.ActorEmail + } + return + } + prefix := actorPrefix(a) + in.ActorAPIKeyPrefix = prefix + if prefix != "" { + in.ActorName = "apikey:" + prefix + } +} + +func (s *Server) pushAuditToPortal(entry *store.AuditEntry) { + if s == nil || s.auditPusher == nil || entry == nil { + return + } + pusher := s.auditPusher + go func() { + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + if err := pusher.PushEvent(ctx, entry); err != nil { + log.Printf("httpapi: audit portal push event_id=%s: %v", entry.EventID, err) + } + }() +} + +func clientIP(r *http.Request) string { + if r == nil { + return "" + } + if xff := strings.TrimSpace(r.Header.Get("X-Forwarded-For")); xff != "" { + parts := strings.Split(xff, ",") + if len(parts) > 0 { + return strings.TrimSpace(parts[0]) + } + } + if xrip := strings.TrimSpace(r.Header.Get("X-Real-IP")); xrip != "" { + return xrip + } + host, _, err := net.SplitHostPort(strings.TrimSpace(r.RemoteAddr)) + if err != nil { + return strings.TrimSpace(r.RemoteAddr) + } + return host +} + +// initAuditPusher wires portal push when URL and secret are configured. +func (s *Server) initAuditPusher(portalURL, ingestSecret string) { + base := strings.TrimSpace(portalURL) + secret := strings.TrimSpace(ingestSecret) + if base == "" || secret == "" { + return + } + s.auditPusher = &audit.PortalPusher{ + BaseURL: base, + Secret: secret, + MarkPushed: func(id string) error { + if s.store == nil { + return nil + } + return s.store.MarkAuditPortalPushed(id) + }, + } +} diff --git a/internal/httpapi/routes_audit_test.go b/internal/httpapi/routes_audit_test.go new file mode 100644 index 0000000..18df07f --- /dev/null +++ b/internal/httpapi/routes_audit_test.go @@ -0,0 +1,52 @@ +package httpapi + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + + "evobgp/internal/store" +) + +func TestHandleListAudit(t *testing.T) { + mem := store.NewMemory() + mem.SeedDemo() + tenant, _, _, _, _ := mem.DemoIDs() + + srv, err := New(Options{SeedDemo: false, InsecureDev: true}) + if err != nil { + t.Fatal(err) + } + srv.store = mem + + _, err = mem.AppendAudit(store.AuditAppendInput{ + TenantID: tenant, + Action: "bgp.module.create", + Summary: "Created module demo", + TargetID: "mod-x", + }) + if err != nil { + t.Fatal(err) + } + + req := httptest.NewRequest(http.MethodGet, "/v1/audit", nil) + req.Header.Set("Authorization", "Bearer dev") + rec := httptest.NewRecorder() + srv.Handler().ServeHTTP(rec, req) + if rec.Code != http.StatusOK { + t.Fatalf("status=%d body=%s", rec.Code, rec.Body.String()) + } + var body struct { + Items []map[string]any `json:"items"` + } + if err := json.Unmarshal(rec.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if len(body.Items) != 1 { + t.Fatalf("items=%d", len(body.Items)) + } + if body.Items[0]["action"] != "bgp.module.create" { + t.Fatalf("action=%v", body.Items[0]["action"]) + } +} diff --git a/internal/httpapi/routes_crud.go b/internal/httpapi/routes_crud.go index 4026033..f233da7 100644 --- a/internal/httpapi/routes_crud.go +++ b/internal/httpapi/routes_crud.go @@ -8,6 +8,7 @@ import ( "io" "log" "net/http" + "sort" "strconv" "strings" "time" @@ -113,6 +114,7 @@ func (s *Server) handlePostModule(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.module.create", "Created module "+mod.Name, mod.ID, map[string]any{"module_id": mod.ID, "type": mod.Type, "name": mod.Name}) writeJSON(w, http.StatusCreated, moduleJSON(mod)) } @@ -174,6 +176,7 @@ func (s *Server) handlePatchModule(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.module.update", "Updated module "+mod.Name, mod.ID, map[string]any{"module_id": mod.ID, "name": mod.Name}) writeJSON(w, http.StatusOK, moduleJSON(mod)) } @@ -193,6 +196,7 @@ func (s *Server) handleDeleteModule(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.module.delete", "Deleted module", moduleID, map[string]any{"module_id": moduleID}) w.WriteHeader(http.StatusNoContent) } @@ -369,6 +373,7 @@ func (s *Server) handlePostCDNSource(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.cdn_source.create", "Created CDN source", x.ID, map[string]any{"module_id": mid, "source_id": x.ID, "url": x.URL}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "cdn_source_create") writeJSON(w, http.StatusCreated, cdnSourceJSON(x)) } @@ -399,6 +404,7 @@ func (s *Server) handlePatchCDNSource(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.cdn_source.update", "Updated CDN source", x.ID, map[string]any{"module_id": mid, "source_id": x.ID}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "cdn_source_patch") writeJSON(w, http.StatusOK, cdnSourceJSON(x)) } @@ -409,10 +415,12 @@ func (s *Server) handleDeleteCDNSource(w http.ResponseWriter, r *http.Request) { return } mid := r.PathValue("module_id") - if err := s.store.DeleteCDNSource(a.TenantID, mid, r.PathValue("source_id")); err != nil { + sourceID := r.PathValue("source_id") + if err := s.store.DeleteCDNSource(a.TenantID, mid, sourceID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.cdn_source.delete", "Deleted CDN source", sourceID, map[string]any{"module_id": mid, "source_id": sourceID}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "cdn_source_delete") w.WriteHeader(http.StatusNoContent) } @@ -471,6 +479,7 @@ func (s *Server) handlePostAS(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.as_entry.create", "Created AS entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "asn": x.ASN}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "as_entry_create") writeJSON(w, http.StatusCreated, asEntryJSON(x)) } @@ -491,6 +500,7 @@ func (s *Server) handlePatchAS(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.as_entry.update", "Updated AS entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "asn": x.ASN}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "as_entry_patch") writeJSON(w, http.StatusOK, asEntryJSON(x)) } @@ -501,10 +511,12 @@ func (s *Server) handleDeleteAS(w http.ResponseWriter, r *http.Request) { return } mid := r.PathValue("module_id") - if err := s.store.DeleteASEntry(a.TenantID, mid, r.PathValue("entry_id")); err != nil { + entryID := r.PathValue("entry_id") + if err := s.store.DeleteASEntry(a.TenantID, mid, entryID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.as_entry.delete", "Deleted AS entry", entryID, map[string]any{"module_id": mid, "entry_id": entryID}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "as_entry_delete") w.WriteHeader(http.StatusNoContent) } @@ -548,6 +560,7 @@ func (s *Server) handlePostDomain(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.domain_entry.create", "Created domain entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "fqdn": x.FQDN}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "domain_entry_create") writeJSON(w, http.StatusCreated, domainEntryJSON(x)) } @@ -568,6 +581,7 @@ func (s *Server) handlePatchDomain(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.domain_entry.update", "Updated domain entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "fqdn": x.FQDN}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "domain_entry_patch") writeJSON(w, http.StatusOK, domainEntryJSON(x)) } @@ -578,10 +592,12 @@ func (s *Server) handleDeleteDomain(w http.ResponseWriter, r *http.Request) { return } mid := r.PathValue("module_id") - if err := s.store.DeleteDomainEntry(a.TenantID, mid, r.PathValue("entry_id")); err != nil { + entryID := r.PathValue("entry_id") + if err := s.store.DeleteDomainEntry(a.TenantID, mid, entryID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.domain_entry.delete", "Deleted domain entry", entryID, map[string]any{"module_id": mid, "entry_id": entryID}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "domain_entry_delete") w.WriteHeader(http.StatusNoContent) } @@ -625,6 +641,7 @@ func (s *Server) handlePostIPRange(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.ip_range.create", "Created IP range entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "prefix": x.Prefix}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "ip_range_create") writeJSON(w, http.StatusCreated, ipRangeJSON(x)) } @@ -645,6 +662,7 @@ func (s *Server) handlePatchIPRange(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.ip_range.update", "Updated IP range entry", x.ID, map[string]any{"module_id": mid, "entry_id": x.ID, "prefix": x.Prefix}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "ip_range_patch") writeJSON(w, http.StatusOK, ipRangeJSON(x)) } @@ -655,10 +673,12 @@ func (s *Server) handleDeleteIPRange(w http.ResponseWriter, r *http.Request) { return } mid := r.PathValue("module_id") - if err := s.store.DeleteIPRangeEntry(a.TenantID, mid, r.PathValue("entry_id")); err != nil { + entryID := r.PathValue("entry_id") + if err := s.store.DeleteIPRangeEntry(a.TenantID, mid, entryID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.ip_range.delete", "Deleted IP range entry", entryID, map[string]any{"module_id": mid, "entry_id": entryID}) s.enqueueModuleRefreshIfEnabled(a.TenantID, mid, "ip_range_delete") w.WriteHeader(http.StatusNoContent) } @@ -850,6 +870,7 @@ func (s *Server) handlePostDoh(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.doh_profile.create", "Created DoH profile "+x.Name, x.ID, map[string]any{"profile_id": x.ID, "name": x.Name}) writeJSON(w, http.StatusCreated, dohJSON(x)) } @@ -868,6 +889,7 @@ func (s *Server) handlePatchDoh(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.doh_profile.update", "Updated DoH profile "+x.Name, x.ID, map[string]any{"profile_id": x.ID, "name": x.Name}) writeJSON(w, http.StatusOK, dohJSON(x)) } @@ -876,10 +898,12 @@ func (s *Server) handleDeleteDoh(w http.ResponseWriter, r *http.Request) { if !ok || !s.requirePerm(w, a, "bgp:directories:write") { return } - if err := s.store.DeleteDohProfile(a.TenantID, r.PathValue("id")); err != nil { + profileID := r.PathValue("id") + if err := s.store.DeleteDohProfile(a.TenantID, profileID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.doh_profile.delete", "Deleted DoH profile", profileID, map[string]any{"profile_id": profileID}) w.WriteHeader(http.StatusNoContent) } @@ -936,6 +960,7 @@ func (s *Server) handlePostComm(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.community.create", "Created community "+x.Community, x.ID, map[string]any{"community_id": x.ID, "community": x.Community}) writeJSON(w, http.StatusCreated, commJSON(x)) } @@ -954,6 +979,7 @@ func (s *Server) handlePatchComm(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.community.update", "Updated community "+x.Community, x.ID, map[string]any{"community_id": x.ID, "community": x.Community}) writeJSON(w, http.StatusOK, commJSON(x)) } @@ -962,10 +988,12 @@ func (s *Server) handleDeleteComm(w http.ResponseWriter, r *http.Request) { if !ok || !s.requirePerm(w, a, "bgp:directories:write") { return } - if err := s.store.DeleteCommunity(a.TenantID, r.PathValue("id")); err != nil { + commID := r.PathValue("id") + if err := s.store.DeleteCommunity(a.TenantID, commID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.community.delete", "Deleted community", commID, map[string]any{"community_id": commID}) w.WriteHeader(http.StatusNoContent) } @@ -988,6 +1016,7 @@ func (s *Server) handlePostPeer(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.peer.create", "Created BGP peer "+x.Name, x.ID, map[string]any{"peer_id": x.ID, "neighbor": x.Neighbor}) s.enqueuePeerReconcile(a.TenantID, "peer_create") writeJSON(w, http.StatusCreated, peerJSON(x)) } @@ -1031,6 +1060,7 @@ func (s *Server) handlePatchPeer(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.peer.update", "Updated BGP peer "+x.Name, x.ID, map[string]any{"peer_id": x.ID, "neighbor": x.Neighbor}) s.enqueuePeerReconcile(a.TenantID, "peer_patch") writeJSON(w, http.StatusOK, peerJSON(x)) } @@ -1051,6 +1081,7 @@ func (s *Server) handleDeletePeer(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.peer.delete", "Deleted BGP peer", peerID, map[string]any{"peer_id": peerID}) s.enqueuePeerReconcile(a.TenantID, "peer_delete") w.WriteHeader(http.StatusNoContent) } @@ -1074,6 +1105,7 @@ func (s *Server) handlePostSpeaker(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.speaker.create", "Created speaker "+x.ID, x.ID, map[string]any{"speaker_id": x.ID, "role": x.Role}) resp := speakerJSONFromStore(s.store, x) if meta := store.ParseSpeakerMeta(x.MetaJSON); meta.AgentSecret != "" { resp["agent_secret"] = meta.AgentSecret @@ -1109,6 +1141,7 @@ func (s *Server) handlePatchSpeaker(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.speaker.update", "Updated speaker "+x.ID, x.ID, map[string]any{"speaker_id": x.ID, "role": x.Role}) writeJSON(w, http.StatusOK, speakerJSONFromStore(s.store, x)) } @@ -1117,10 +1150,12 @@ func (s *Server) handleDeleteSpeaker(w http.ResponseWriter, r *http.Request) { if !ok || !s.requirePerm(w, a, "bgp:network:write") { return } - if err := s.store.DeleteSpeaker(a.TenantID, r.PathValue("speaker_id")); err != nil { + speakerID := r.PathValue("speaker_id") + if err := s.store.DeleteSpeaker(a.TenantID, speakerID); err != nil { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.speaker.delete", "Deleted speaker", speakerID, map[string]any{"speaker_id": speakerID}) w.WriteHeader(http.StatusNoContent) } @@ -1187,9 +1222,22 @@ func (s *Server) handlePatchSettings(w http.ResponseWriter, r *http.Request) { writeStoreErr(w, err) return } + s.recordCRUDAudit(r, a, "bgp.settings.update", "Updated tenant settings", a.TenantID, map[string]any{"keys": settingsAuditKeys(body)}) writeJSON(w, http.StatusOK, map[string]string{"status": "ok"}) } +func settingsAuditKeys(body map[string]any) []string { + if len(body) == 0 { + return nil + } + keys := make([]string, 0, len(body)) + for k := range body { + keys = append(keys, k) + } + sort.Strings(keys) + return keys +} + func parseRevisionRetentionMinutes(v any) (int, bool) { const minMinutes = 15 const maxMinutes = 30 * 24 * 60 diff --git a/internal/httpapi/server.go b/internal/httpapi/server.go index e36055b..a319b61 100644 --- a/internal/httpapi/server.go +++ b/internal/httpapi/server.go @@ -10,6 +10,7 @@ import ( "strings" "time" + "evobgp/internal/audit" "evobgp/internal/jobs" "evobgp/internal/maintenance" "evobgp/internal/pgmonitor" @@ -38,11 +39,13 @@ type Server struct { mux *http.ServeMux // Portal / dual-auth (JWT) configuration. - jwtSecret string - authIssuer string - authPortalURL string - portalTenantID string - authRequired bool + jwtSecret string + authIssuer string + authPortalURL string + portalTenantID string + authRequired bool + auditIngestSecret string + auditPusher *audit.PortalPusher } // Options configures the API server. @@ -63,6 +66,8 @@ type Options struct { AuthPortalURL string // AUTH_PORTAL_URL (returned by /v1/auth/config for the UI) PortalTenantID string // fallback when JWT has no bgp_tenant_id / tenants.bgp AuthRequired bool // AUTH_REQUIRED / EVOBGP_AUTH_REQUIRED (surfaced via /v1/auth/config) + // AuditIngestSecret — AUTH_AUDIT_INGEST_SECRET for portal push (optional). + AuditIngestSecret string } // New constructs Server and wiring for async jobs. @@ -123,10 +128,12 @@ func New(opts Options) (*Server, error) { authPortalURL: strings.TrimSpace(opts.AuthPortalURL), portalTenantID: strings.TrimSpace(opts.PortalTenantID), authRequired: opts.AuthRequired, + auditIngestSecret: strings.TrimSpace(opts.AuditIngestSecret), } if s.authIssuer == "" { s.authIssuer = "https://auth.shnt.top" } + s.initAuditPusher(s.authPortalURL, s.auditIngestSecret) s.mux = http.NewServeMux() s.registerRoutes() return s, nil diff --git a/internal/repository/postgres_audit.go b/internal/repository/postgres_audit.go new file mode 100644 index 0000000..68cd7ba --- /dev/null +++ b/internal/repository/postgres_audit.go @@ -0,0 +1,186 @@ +package repository + +import ( + "context" + "encoding/json" + "strconv" + "strings" + "time" + + "github.com/google/uuid" + + "evobgp/internal/store" +) + +// AppendAudit inserts a tenant-scoped audit row. +func (p *Postgres) AppendAudit(in store.AuditAppendInput) (*store.AuditEntry, error) { + if strings.TrimSpace(in.TenantID) == "" || strings.TrimSpace(in.Action) == "" || strings.TrimSpace(in.Summary) == "" { + return nil, store.ErrInvalidInput + } + sev := strings.TrimSpace(in.Severity) + if sev == "" { + sev = store.AuditSeverityInfo + } + if !store.ValidAuditSeverity(sev) { + return nil, store.ErrInvalidInput + } + ctx := context.Background() + id := uuid.NewString() + eventID := "bgp-" + uuid.NewString() + var detailJSON []byte + if in.Details != nil { + detailJSON, _ = json.Marshal(in.Details) + } + var createdAt time.Time + err := p.pool.QueryRow(ctx, ` + INSERT INTO audit_log + (id, tenant_id, event_id, source_app, action, severity, + actor_user_id, actor_email, actor_name, actor_api_key_prefix, + target_type, target_id, summary, details_json, ip, created_at) + VALUES ($1, $2, $3, 'bgp', $4, $5, $6, $7, $8, $9, $10, $11, $12, $13::jsonb, $14, now()) + RETURNING created_at`, + id, strings.TrimSpace(in.TenantID), eventID, strings.TrimSpace(in.Action), sev, + nullIfEmpty(in.ActorUserID), nullIfEmpty(in.ActorEmail), nullIfEmpty(in.ActorName), + nullIfEmpty(in.ActorAPIKeyPrefix), nullIfEmpty(in.TargetType), nullIfEmpty(in.TargetID), + strings.TrimSpace(in.Summary), nullJSONBytes(detailJSON), nullIfEmpty(in.IP), + ).Scan(&createdAt) + if err != nil { + return nil, err + } + return &store.AuditEntry{ + ID: id, + TenantID: strings.TrimSpace(in.TenantID), + EventID: eventID, + SourceApp: store.AuditSourceAppBGP, + Action: strings.TrimSpace(in.Action), + Severity: sev, + ActorUserID: strings.TrimSpace(in.ActorUserID), + ActorEmail: strings.TrimSpace(in.ActorEmail), + ActorName: strings.TrimSpace(in.ActorName), + ActorAPIKeyPrefix: strings.TrimSpace(in.ActorAPIKeyPrefix), + TargetType: strings.TrimSpace(in.TargetType), + TargetID: strings.TrimSpace(in.TargetID), + Summary: strings.TrimSpace(in.Summary), + Details: in.Details, + IP: strings.TrimSpace(in.IP), + CreatedAt: createdAt.UTC(), + }, nil +} + +// ListAudit returns paginated audit rows for a tenant. +func (p *Postgres) ListAudit(tenantID, cursor string, limit int, filter store.AuditListFilter) ([]*store.AuditEntry, string, bool, error) { + if limit <= 0 { + limit = 50 + } + off := 0 + if cursor != "" { + if n, err := strconv.Atoi(cursor); err == nil && n >= 0 { + off = n + } + } + ctx := context.Background() + args := []any{tenantID} + where := "tenant_id = $1" + argN := 2 + if a := strings.TrimSpace(filter.Action); a != "" { + where += " AND action = $" + strconv.Itoa(argN) + args = append(args, a) + argN++ + } + if s := strings.TrimSpace(filter.Severity); s != "" { + where += " AND severity = $" + strconv.Itoa(argN) + args = append(args, s) + argN++ + } + args = append(args, limit+1, off) + q := ` + SELECT id, tenant_id, event_id, source_app, action, severity, + actor_user_id, actor_email, actor_name, actor_api_key_prefix, + target_type, target_id, summary, details_json, ip, created_at, portal_pushed_at + FROM audit_log + WHERE ` + where + ` + ORDER BY created_at DESC, id DESC + LIMIT $` + strconv.Itoa(argN) + ` OFFSET $` + strconv.Itoa(argN+1) + rows, err := p.pool.Query(ctx, q, args...) + if err != nil { + return nil, "", false, err + } + defer rows.Close() + var out []*store.AuditEntry + for rows.Next() { + row, err := scanAuditEntry(rows.Scan) + if err != nil { + return nil, "", false, err + } + out = append(out, row) + } + if err := rows.Err(); err != nil { + return nil, "", false, err + } + more := len(out) > limit + if more { + out = out[:limit] + } + next := "" + if more { + next = strconv.Itoa(off + limit) + } + return out, next, more, nil +} + +// MarkAuditPortalPushed sets portal_pushed_at for a row. +func (p *Postgres) MarkAuditPortalPushed(id string) error { + ctx := context.Background() + tag, err := p.pool.Exec(ctx, `UPDATE audit_log SET portal_pushed_at = now() WHERE id = $1`, id) + if err != nil { + return err + } + if tag.RowsAffected() == 0 { + return store.ErrNotFound + } + return nil +} + +func scanAuditEntry(scan func(dest ...any) error) (*store.AuditEntry, error) { + var row store.AuditEntry + var actorUserID, actorEmail, actorName, actorPrefix, targetType, targetID, ip *string + var detailRaw []byte + var portalPushed *time.Time + if err := scan( + &row.ID, &row.TenantID, &row.EventID, &row.SourceApp, &row.Action, &row.Severity, + &actorUserID, &actorEmail, &actorName, &actorPrefix, + &targetType, &targetID, &row.Summary, &detailRaw, &ip, &row.CreatedAt, &portalPushed, + ); err != nil { + return nil, err + } + row.CreatedAt = row.CreatedAt.UTC() + if actorUserID != nil { + row.ActorUserID = *actorUserID + } + if actorEmail != nil { + row.ActorEmail = *actorEmail + } + if actorName != nil { + row.ActorName = *actorName + } + if actorPrefix != nil { + row.ActorAPIKeyPrefix = *actorPrefix + } + if targetType != nil { + row.TargetType = *targetType + } + if targetID != nil { + row.TargetID = *targetID + } + if ip != nil { + row.IP = *ip + } + if len(detailRaw) > 0 { + _ = json.Unmarshal(detailRaw, &row.Details) + } + if portalPushed != nil { + t := portalPushed.UTC() + row.PortalPushedAt = &t + } + return &row, nil +} diff --git a/internal/store/audit.go b/internal/store/audit.go new file mode 100644 index 0000000..e1c3efb --- /dev/null +++ b/internal/store/audit.go @@ -0,0 +1,67 @@ +package store + +import ( + "strings" + "time" +) + +const ( + AuditSourceAppBGP = "bgp" + AuditTargetAppResource = "app_resource" + AuditSeverityInfo = "info" + AuditSeverityWarning = "warning" + AuditSeverityCritical = "critical" +) + +// AuditEntry is a persisted CRUD / settings audit row (local + portal ingest). +type AuditEntry struct { + ID string + TenantID string + EventID string + SourceApp string + Action string + Severity string + ActorUserID string + ActorEmail string + ActorName string + ActorAPIKeyPrefix string + TargetType string + TargetID string + Summary string + Details map[string]any + IP string + CreatedAt time.Time + PortalPushedAt *time.Time +} + +// AuditAppendInput is input for AppendAudit. +type AuditAppendInput struct { + TenantID string + Action string + Severity string + ActorUserID string + ActorEmail string + ActorName string + ActorAPIKeyPrefix string + TargetType string + TargetID string + Summary string + Details map[string]any + IP string +} + +// AuditListFilter optional query filters for ListAudit. +type AuditListFilter struct { + Action string + Severity string +} + +// ValidAuditSeverity reports whether s is an allowed severity. +func ValidAuditSeverity(s string) bool { + switch strings.ToLower(strings.TrimSpace(s)) { + case AuditSeverityInfo, AuditSeverityWarning, AuditSeverityCritical: + return true + default: + return false + } +} diff --git a/internal/store/backend.go b/internal/store/backend.go index 18b2dab..bef881c 100644 --- a/internal/store/backend.go +++ b/internal/store/backend.go @@ -132,6 +132,11 @@ type Backend interface { AppendRuntimeLogCleanupAudit(tenantID, actor, filename, action string, sizeBefore int64, sizeAfter *int64, detail map[string]any) (string, error) ListRuntimeLogCleanupAudit(tenantID, cursor string, limit int) ([]*RuntimeLogCleanupAudit, string, bool, error) + // CRUD audit log (tenant-scoped; optional portal ingest push from httpapi). + AppendAudit(in AuditAppendInput) (*AuditEntry, error) + ListAudit(tenantID, cursor string, limit int, filter AuditListFilter) ([]*AuditEntry, string, bool, error) + MarkAuditPortalPushed(id string) error + // Firewall blocklist clients and policy rules. ListFirewallClients(tenantID string) ([]*FirewallClient, error) GetFirewallClient(tenantID, id string) (*FirewallClient, error) diff --git a/internal/store/memory.go b/internal/store/memory.go index 04c7db8..58803a6 100644 --- a/internal/store/memory.go +++ b/internal/store/memory.go @@ -50,6 +50,7 @@ type Memory struct { maintenancePolicies map[string]*MaintenancePolicy maintConfigAudit []*MaintenancePolicyConfigAudit runtimeLogCleanupAudit []*RuntimeLogCleanupAudit + auditLog []*AuditEntry // DemoIDs valid after SeedDemo() demoTenantID string diff --git a/internal/store/memory_audit.go b/internal/store/memory_audit.go new file mode 100644 index 0000000..653d2f1 --- /dev/null +++ b/internal/store/memory_audit.go @@ -0,0 +1,117 @@ +package store + +import ( + "sort" + "strings" + "time" + + "github.com/google/uuid" +) + +func (m *Memory) AppendAudit(in AuditAppendInput) (*AuditEntry, error) { + if strings.TrimSpace(in.TenantID) == "" || strings.TrimSpace(in.Action) == "" || strings.TrimSpace(in.Summary) == "" { + return nil, ErrInvalidInput + } + sev := strings.TrimSpace(in.Severity) + if sev == "" { + sev = AuditSeverityInfo + } + if !ValidAuditSeverity(sev) { + return nil, ErrInvalidInput + } + now := time.Now().UTC() + row := &AuditEntry{ + ID: uuid.NewString(), + TenantID: strings.TrimSpace(in.TenantID), + EventID: "bgp-" + uuid.NewString(), + SourceApp: AuditSourceAppBGP, + Action: strings.TrimSpace(in.Action), + Severity: sev, + ActorUserID: strings.TrimSpace(in.ActorUserID), + ActorEmail: strings.TrimSpace(in.ActorEmail), + ActorName: strings.TrimSpace(in.ActorName), + ActorAPIKeyPrefix: strings.TrimSpace(in.ActorAPIKeyPrefix), + TargetType: strings.TrimSpace(in.TargetType), + TargetID: strings.TrimSpace(in.TargetID), + Summary: strings.TrimSpace(in.Summary), + Details: in.Details, + IP: strings.TrimSpace(in.IP), + CreatedAt: now, + } + m.mu.Lock() + defer m.mu.Unlock() + m.auditLog = append(m.auditLog, row) + return cloneAuditEntry(row), nil +} + +func (m *Memory) ListAudit(tenantID, cursor string, limit int, filter AuditListFilter) ([]*AuditEntry, string, bool, error) { + if limit <= 0 { + limit = 50 + } + m.mu.RLock() + defer m.mu.RUnlock() + var filtered []*AuditEntry + for _, row := range m.auditLog { + if row.TenantID != tenantID { + continue + } + if a := strings.TrimSpace(filter.Action); a != "" && row.Action != a { + continue + } + if s := strings.TrimSpace(filter.Severity); s != "" && row.Severity != s { + continue + } + filtered = append(filtered, row) + } + sort.Slice(filtered, func(i, j int) bool { + if filtered[i].CreatedAt.Equal(filtered[j].CreatedAt) { + return filtered[i].ID > filtered[j].ID + } + return filtered[i].CreatedAt.After(filtered[j].CreatedAt) + }) + off := parseMaintCursor(cursor) + end := off + limit + next := "" + hasMore := false + if end > len(filtered) { + end = len(filtered) + } else if end < len(filtered) { + hasMore = true + next = formatMaintCursor(end) + } + if off >= len(filtered) { + return nil, "", false, nil + } + out := make([]*AuditEntry, end-off) + for i := off; i < end; i++ { + out[i-off] = cloneAuditEntry(filtered[i]) + } + return out, next, hasMore, nil +} + +func (m *Memory) MarkAuditPortalPushed(id string) error { + m.mu.Lock() + defer m.mu.Unlock() + for _, row := range m.auditLog { + if row.ID == id { + now := time.Now().UTC() + row.PortalPushedAt = &now + return nil + } + } + return ErrNotFound +} + +func cloneAuditEntry(row *AuditEntry) *AuditEntry { + if row == nil { + return nil + } + cp := *row + if row.Details != nil { + cp.Details = make(map[string]any, len(row.Details)) + for k, v := range row.Details { + cp.Details[k] = v + } + } + return &cp +} diff --git a/internal/store/memory_audit_test.go b/internal/store/memory_audit_test.go new file mode 100644 index 0000000..9139ed1 --- /dev/null +++ b/internal/store/memory_audit_test.go @@ -0,0 +1,70 @@ +package store + +import "testing" + +func TestMemoryAppendAndListAudit(t *testing.T) { + m := NewMemory() + tenantA := "tenant-a" + tenantB := "tenant-b" + + entry, err := m.AppendAudit(AuditAppendInput{ + TenantID: tenantA, + Action: "bgp.module.create", + Summary: "Created module test", + TargetID: "mod-1", + }) + if err != nil { + t.Fatal(err) + } + if entry == nil || entry.EventID == "" || entry.SourceApp != AuditSourceAppBGP { + t.Fatalf("unexpected entry: %+v", entry) + } + + if _, err := m.AppendAudit(AuditAppendInput{ + TenantID: tenantB, + Action: "bgp.peer.delete", + Summary: "Deleted peer", + }); err != nil { + t.Fatal(err) + } + + items, _, hasMore, err := m.ListAudit(tenantA, "", 10, AuditListFilter{}) + if err != nil { + t.Fatal(err) + } + if len(items) != 1 || hasMore { + t.Fatalf("items=%d hasMore=%v", len(items), hasMore) + } + if items[0].Action != "bgp.module.create" { + t.Fatalf("action=%s", items[0].Action) + } + + filtered, _, _, err := m.ListAudit(tenantA, "", 10, AuditListFilter{Action: "bgp.peer.delete"}) + if err != nil { + t.Fatal(err) + } + if len(filtered) != 0 { + t.Fatalf("expected empty filter result, got %d", len(filtered)) + } + + if err := m.MarkAuditPortalPushed(entry.ID); err != nil { + t.Fatal(err) + } + items2, _, _, err := m.ListAudit(tenantA, "", 10, AuditListFilter{}) + if err != nil { + t.Fatal(err) + } + if items2[0].PortalPushedAt == nil { + t.Fatal("expected portal_pushed_at") + } +} + +func TestMemoryAppendAuditValidation(t *testing.T) { + m := NewMemory() + if _, err := m.AppendAudit(AuditAppendInput{}); err != ErrInvalidInput { + t.Fatalf("err=%v", err) + } + if _, err := m.AppendAudit(AuditAppendInput{TenantID: "t", Action: "x", Summary: "s", Severity: "bad"}); err != ErrInvalidInput { + t.Fatalf("err=%v", err) + } +} diff --git a/migrations/postgres/000030_audit_log.down.sql b/migrations/postgres/000030_audit_log.down.sql new file mode 100644 index 0000000..1be3935 --- /dev/null +++ b/migrations/postgres/000030_audit_log.down.sql @@ -0,0 +1,4 @@ +DROP INDEX IF EXISTS idx_audit_log_tenant_action; +DROP INDEX IF EXISTS idx_audit_log_tenant_created; +DROP INDEX IF EXISTS idx_audit_log_event_id; +DROP TABLE IF EXISTS audit_log; diff --git a/migrations/postgres/000030_audit_log.up.sql b/migrations/postgres/000030_audit_log.up.sql new file mode 100644 index 0000000..9e0b241 --- /dev/null +++ b/migrations/postgres/000030_audit_log.up.sql @@ -0,0 +1,30 @@ +CREATE TABLE IF NOT EXISTS audit_log ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL, + event_id TEXT NOT NULL, + source_app TEXT NOT NULL DEFAULT 'bgp', + action TEXT NOT NULL, + severity TEXT NOT NULL DEFAULT 'info', + actor_user_id TEXT, + actor_email TEXT, + actor_name TEXT, + actor_api_key_prefix TEXT, + target_type TEXT, + target_id TEXT, + summary TEXT NOT NULL, + details_json JSONB, + ip TEXT, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + portal_pushed_at TIMESTAMPTZ, + CONSTRAINT audit_log_severity_chk CHECK (severity IN ('info', 'warning', 'critical')), + CONSTRAINT audit_log_source_app_chk CHECK (source_app = 'bgp'), + CONSTRAINT audit_log_summary_chk CHECK (length(trim(summary)) > 0) +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_audit_log_event_id ON audit_log (event_id); + +CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_created + ON audit_log (tenant_id, created_at DESC); + +CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_action + ON audit_log (tenant_id, action, created_at DESC); diff --git a/migrations/sqlite/000030_audit_log.down.sql b/migrations/sqlite/000030_audit_log.down.sql new file mode 100644 index 0000000..1be3935 --- /dev/null +++ b/migrations/sqlite/000030_audit_log.down.sql @@ -0,0 +1,4 @@ +DROP INDEX IF EXISTS idx_audit_log_tenant_action; +DROP INDEX IF EXISTS idx_audit_log_tenant_created; +DROP INDEX IF EXISTS idx_audit_log_event_id; +DROP TABLE IF EXISTS audit_log; diff --git a/migrations/sqlite/000030_audit_log.up.sql b/migrations/sqlite/000030_audit_log.up.sql new file mode 100644 index 0000000..4a80fb1 --- /dev/null +++ b/migrations/sqlite/000030_audit_log.up.sql @@ -0,0 +1,27 @@ +CREATE TABLE IF NOT EXISTS audit_log ( + id TEXT PRIMARY KEY, + tenant_id TEXT NOT NULL, + event_id TEXT NOT NULL, + source_app TEXT NOT NULL DEFAULT 'bgp', + action TEXT NOT NULL, + severity TEXT NOT NULL DEFAULT 'info', + actor_user_id TEXT, + actor_email TEXT, + actor_name TEXT, + actor_api_key_prefix TEXT, + target_type TEXT, + target_id TEXT, + summary TEXT NOT NULL, + details_json TEXT, + ip TEXT, + created_at TEXT NOT NULL, + portal_pushed_at TEXT +); + +CREATE UNIQUE INDEX IF NOT EXISTS idx_audit_log_event_id ON audit_log (event_id); + +CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_created + ON audit_log (tenant_id, created_at DESC); + +CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_action + ON audit_log (tenant_id, action, created_at DESC);