Compare commits

...
2 Commits
Author SHA1 Message Date
Denozordec ff6efec4c5 feat(api): add endpoint to list community prefixes with pagination
CI / changes (push) Successful in 6s
CI / commitlint (push) Skipped
CI / web (push) Skipped
CI / openapi (push) Successful in 25s
CI / go (push) Successful in 1m4s
CI / bird2 (push) Successful in 14s
CI / release (push) Successful in 4m17s
Introduced a new GET endpoint `/v1/communities/{id}/prefixes` to retrieve unique prefixes associated with a community, including pagination support via cursor and limit parameters. Updated OpenAPI documentation to reflect this addition. Implemented backend logic in both PostgreSQL and in-memory storage to handle the new functionality, ensuring proper authorization checks and response formatting.
2026-07-23 11:10:53 +07:00
DenozordecandCursor 738d2e2256 feat(httpapi): add local audit log with portal dual-write
CI / changes (push) Successful in 7s
CI / commitlint (push) Skipped
CI / web (push) Skipped
CI / openapi (push) Successful in 32s
CI / go (push) Successful in 1m23s
CI / bird2 (push) Successful in 16s
CI / release (push) Successful in 4m56s
Локальный audit_log (миграции pg/sqlite), GET /v1/audit, запись на CRUD и async push в auth-portal (source_app=bgp).

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-07-21 13:24:54 +07:00
25 changed files with 1344 additions and 13 deletions
+1
View File
@@ -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 {
+1
View File
@@ -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 {
+1
View File
@@ -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 (по приоритету):
+10
View File
@@ -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`, и т.д.
+151
View File
@@ -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:
@@ -2856,6 +2922,57 @@ paths:
default:
$ref: "#/components/responses/DefaultProblem"
/v1/communities/{id}/prefixes:
parameters:
- $ref: "#/components/parameters/TenantId"
- $ref: "#/components/parameters/CommunityId"
- $ref: "#/components/parameters/Cursor"
- name: limit
in: query
schema:
type: integer
default: 500
maximum: 5000
get:
tags: [Communities]
summary: Префиксы community (latest revision per module)
description: |
Уникальные materialized-префиксы с данным community_id
из последней ревизии каждого модуля tenant.
Поле `prefixes` — плоский список для клиентов вроде EvoFirewall.
operationId: listCommunityPrefixes
responses:
"200":
description: Успешно.
content:
application/json:
schema:
type: object
required: [items, has_more]
properties:
items:
type: array
items:
type: object
required: [prefix]
properties:
prefix:
type: string
source:
type: string
prefixes:
type: array
items:
type: string
next_cursor:
type: ["string", "null"]
has_more:
type: boolean
"404":
$ref: "#/components/responses/NotFound"
default:
$ref: "#/components/responses/DefaultProblem"
/v1/communities/{id}:
parameters:
- $ref: "#/components/parameters/TenantId"
@@ -4520,6 +4637,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]
+118
View File
@@ -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
}
+63
View File
@@ -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")
}
}
+1
View File
@@ -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)
+6 -1
View File
@@ -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)
}
+171
View File
@@ -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)
},
}
}
+52
View File
@@ -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"])
}
}
+89 -7
View File
@@ -8,6 +8,7 @@ import (
"io"
"log"
"net/http"
"sort"
"strconv"
"strings"
"time"
@@ -57,6 +58,7 @@ func (s *Server) registerCRUDRoutes(m *http.ServeMux) {
m.HandleFunc("GET /communities", s.handleListComm)
m.HandleFunc("POST /communities", s.handlePostComm)
m.HandleFunc("GET /communities/{id}", s.handleGetComm)
m.HandleFunc("GET /communities/{id}/prefixes", s.handleListCommPrefixes)
m.HandleFunc("PATCH /communities/{id}", s.handlePatchComm)
m.HandleFunc("DELETE /communities/{id}", s.handleDeleteComm)
@@ -113,6 +115,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 +177,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 +197,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 +374,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 +405,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 +416,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 +480,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 +501,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 +512,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 +561,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 +582,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 +593,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 +642,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 +663,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 +674,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 +871,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 +890,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 +899,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)
}
@@ -921,6 +946,39 @@ func (s *Server) handleGetComm(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, commJSON(x))
}
func (s *Server) handleListCommPrefixes(w http.ResponseWriter, r *http.Request) {
a, ok := authFromContext(r.Context())
if !ok || !s.requirePerm(w, a, "bgp:directories:read") {
return
}
limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
if limit == 0 {
limit = 500
}
cursor := r.URL.Query().Get("cursor")
rows, next, more, err := s.store.ListCommunityPrefixes(a.TenantID, r.PathValue("id"), cursor, limit)
if err != nil {
writeStoreErr(w, err)
return
}
items := make([]map[string]any, 0, len(rows))
prefixes := make([]string, 0, len(rows))
for _, pr := range rows {
m := map[string]any{"prefix": pr.Prefix}
if pr.Source != "" {
m["source"] = pr.Source
}
items = append(items, m)
prefixes = append(prefixes, pr.Prefix)
}
writeJSON(w, http.StatusOK, map[string]any{
"items": items,
"prefixes": prefixes,
"next_cursor": strPtrOrNull(next),
"has_more": more,
})
}
func (s *Server) handlePostComm(w http.ResponseWriter, r *http.Request) {
a, ok := authFromContext(r.Context())
if !ok || !s.requirePerm(w, a, "bgp:directories:write") {
@@ -936,6 +994,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 +1013,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 +1022,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 +1050,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 +1094,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 +1115,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 +1139,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 +1175,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 +1184,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 +1256,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
+12 -5
View File
@@ -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
+85
View File
@@ -1252,6 +1252,91 @@ func (p *Postgres) GetCommunity(tenantID, id string) (*store.Community, error) {
return &c, nil
}
func (p *Postgres) ListCommunityPrefixes(tenantID, communityID, cursor string, limit int) ([]store.PrefixRow, string, bool, error) {
if _, err := p.GetCommunity(tenantID, communityID); err != nil {
return nil, "", false, err
}
if limit <= 0 {
limit = 500
}
if limit > 5000 {
limit = 5000
}
off := 0
if cursor != "" {
if n, err := strconv.Atoi(cursor); err == nil && n >= 0 {
off = n
}
}
ctx := context.Background()
useSnap := prefixSnapshotTableExists(ctx, p.pool)
var rows pgx.Rows
var err error
if useSnap {
rows, err = p.pool.Query(ctx, `
WITH latest AS (
SELECT DISTINCT ON (module_id) id, prefix_snapshot_id
FROM config_revision
WHERE tenant_id = $1::uuid AND module_id IS NOT NULL
ORDER BY module_id, created_at DESC
),
combined AS (
SELECT rmp.prefix::text AS prefix, COALESCE(rmp.source, '') AS source
FROM revision_materialized_prefix rmp
JOIN latest l ON l.id = rmp.revision_id
WHERE l.prefix_snapshot_id IS NULL AND rmp.community_id = $2::uuid
UNION
SELECT psr.prefix::text, COALESCE(psr.source, '')
FROM prefix_snapshot_row psr
JOIN latest l ON l.prefix_snapshot_id = psr.snapshot_id
WHERE l.prefix_snapshot_id IS NOT NULL AND psr.community_id = $2::uuid
)
SELECT prefix, source FROM combined
ORDER BY prefix
LIMIT $3 OFFSET $4`, tenantID, communityID, limit+1, off)
} else {
rows, err = p.pool.Query(ctx, `
WITH latest AS (
SELECT DISTINCT ON (module_id) id
FROM config_revision
WHERE tenant_id = $1::uuid AND module_id IS NOT NULL
ORDER BY module_id, created_at DESC
)
SELECT DISTINCT rmp.prefix::text, COALESCE(rmp.source, '')
FROM revision_materialized_prefix rmp
JOIN latest l ON l.id = rmp.revision_id
WHERE rmp.community_id = $2::uuid
ORDER BY 1
LIMIT $3 OFFSET $4`, tenantID, communityID, limit+1, off)
}
if err != nil {
return nil, "", false, err
}
defer rows.Close()
var all []store.PrefixRow
comm := communityID
for rows.Next() {
var pr store.PrefixRow
if err := rows.Scan(&pr.Prefix, &pr.Source); err != nil {
continue
}
pr.CommunityID = &comm
all = append(all, pr)
}
more := len(all) > limit
if more {
all = all[:limit]
}
next := ""
if more {
next = fmt.Sprintf("%d", off+limit)
}
if len(all) == 0 {
return nil, "", false, nil
}
return all, next, more, nil
}
func (p *Postgres) CreateCommunity(tenantID string, in *store.Community) (*store.Community, error) {
if in == nil {
return nil, store.ErrInvalidInput
+186
View File
@@ -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
}
+67
View File
@@ -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
}
}
+7
View File
@@ -60,6 +60,8 @@ type Backend interface {
CreateCommunity(tenantID string, in *Community) (*Community, error)
UpdateCommunity(tenantID, id string, patch *CommunityPatch) (*Community, error)
DeleteCommunity(tenantID, id string) error
// ListCommunityPrefixes returns unique prefixes tagged with community from latest revision per module.
ListCommunityPrefixes(tenantID, communityID, cursor string, limit int) (prefixes []PrefixRow, nextCursor string, hasMore bool, err error)
// ListPeers returns all BGP peers for a tenant (control plane may paginate in httpapi).
ListPeers(tenantID string) []*BGPPeer
@@ -132,6 +134,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)
+1
View File
@@ -50,6 +50,7 @@ type Memory struct {
maintenancePolicies map[string]*MaintenancePolicy
maintConfigAudit []*MaintenancePolicyConfigAudit
runtimeLogCleanupAudit []*RuntimeLogCleanupAudit
auditLog []*AuditEntry
// DemoIDs valid after SeedDemo()
demoTenantID string
+117
View File
@@ -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
}
+70
View File
@@ -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)
}
}
+70
View File
@@ -1,6 +1,8 @@
package store
import (
"sort"
"strconv"
"strings"
"time"
@@ -593,6 +595,74 @@ func (m *Memory) GetCommunity(tenantID, id string) (*Community, error) {
return c, nil
}
func (m *Memory) ListCommunityPrefixes(tenantID, communityID, cursor string, limit int) ([]PrefixRow, string, bool, error) {
if _, err := m.GetCommunity(tenantID, communityID); err != nil {
return nil, "", false, err
}
if limit <= 0 {
limit = 500
}
if limit > 5000 {
limit = 5000
}
off := 0
if cursor != "" {
if n, err := strconv.Atoi(cursor); err == nil && n >= 0 {
off = n
}
}
m.mu.RLock()
defer m.mu.RUnlock()
latestByModule := map[string]*Revision{}
for _, rev := range m.revisions {
if rev.TenantID != tenantID || strings.TrimSpace(rev.ModuleID) == "" {
continue
}
cur := latestByModule[rev.ModuleID]
if cur == nil || rev.CreatedAt.After(cur.CreatedAt) {
latestByModule[rev.ModuleID] = rev
}
}
seen := map[string]struct{}{}
var all []PrefixRow
comm := communityID
for _, rev := range latestByModule {
for _, pr := range m.revPrefixes[rev.ID] {
if pr.CommunityID == nil || *pr.CommunityID != communityID {
continue
}
pfx := strings.TrimSpace(pr.Prefix)
if pfx == "" {
continue
}
if _, ok := seen[pfx]; ok {
continue
}
seen[pfx] = struct{}{}
all = append(all, PrefixRow{Prefix: pfx, CommunityID: &comm, Source: pr.Source})
}
}
sort.Slice(all, func(i, j int) bool { return all[i].Prefix < all[j].Prefix })
if off > len(all) {
return nil, "", false, nil
}
end := off + limit
more := false
next := ""
if end < len(all) {
more = true
next = strconv.Itoa(end)
all = all[off:end]
} else {
all = all[off:]
}
if len(all) == 0 {
return nil, "", false, nil
}
return all, next, more, nil
}
func (m *Memory) CreateCommunity(tenantID string, in *Community) (*Community, error) {
if in == nil || strings.TrimSpace(in.Community) == "" {
return nil, ErrInvalidInput
@@ -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;
@@ -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);
@@ -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;
+27
View File
@@ -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);