mirror of
https://github.com/taylanbakircioglu/haproxy-openmanager.git
synced 2026-09-24 11:26:29 +00:00
02b1cb2bca
Closes #13, Closes #14. This release squashes the v1.4.0 → v1.5.0 development line. v1.4.0 shipped the ACME stability & enterprise audit (Issues #10/#11/#12). v1.5.0 builds on that foundation with two co-equal headline features plus a 22-round audit campaign hardening the prior configuration surface. License remains MIT for v1.5.0 (relicense to AGPL-3.0 lands in v1.5.2). ------------------------------------------------------------------ HEADLINE FEATURE A — ACME Diagnostic Panel (Issue #13) ------------------------------------------------------------------ A live pre-flight + post-failure diagnostic surface for every ACME order, reachable from the ACME Automation page. The panel exists to make ACME failures legible to operators who do NOT have shell access to the API host. Endpoints (`backend/routers/acme_diagnostics.py`): POST /api/letsencrypt/orders/{order_id}/diagnostics Run the full 5-check suite (DNS / port-80 / routing / account / agents) and humanize the order's `error_detail` (>=11 RFC-8555 problem types, backwards compatible with legacy plain-string failures). POST /api/letsencrypt/orders/{order_id}/diagnostics/ {check_id}/rerun Re-run a single check in place — used by the "Re-run" button on every row of the modal's pre-flight table. GET /api/letsencrypt/orders/{order_id}/events Merged event timeline combining the typed `acme_order_events` rows with correlated `user_activity_logs` entries (resource_type = 'letsencrypt_order' AND resource_id = order_id). The diagnostic modal auto-tails this timeline every 5 seconds while open. Service-level checks (`backend/services/acme_diagnostics.py`): * DNS resolution via stdlib socket.gethostbyname_ex through run_in_executor (intentionally avoiding an aiodns runtime dep for v1.5.0). * Port-80 HEAD probe, target locked to the order's domains, success on HTTP 200 OR 404, warns on egress timeout (corp egress policies routinely blackhole outbound 80 — fail-hard would be too noisy). * SSRF guard: probe refuses non-public IPs and surfaces the skip in the diagnostic result; IPv4-mapped IPv6 normalisation closes the `::ffff:169.254.169.254` cloud-metadata vector. * HAProxy routing presence check: matches the order's cluster_ids to a port-80 HTTP frontend. * ACME account validity check against `letsencrypt_accounts`. * Agent presence check (>=1 active agent in target cluster). * Every sub-check wrapped in a wall-clock timeout to bound impact on the API event loop. RBAC: ssl.read for run, ssl.read for events. Per-user 5/min rate limit on both run and rerun, backed by the (user_id, action, created_at DESC) composite index. Frontend (`frontend/src/components/ACMEAutomation.js`): * "Diagnose" button on every order row + the existing "stuck order" warning row. * Modal with two tabs: - Pre-flight Checks (Antd Table with status pills + Re-run buttons + humanized error banner) - Event Log (Antd Timeline with auto-tail polling, scroll- to-bottom, pause-on-hover) * Correlation IDs surfaced in error banners and individual check fail details for backend-log lookup. ------------------------------------------------------------------ HEADLINE FEATURE B — Site Setup Wizard (Issue #14) ------------------------------------------------------------------ A single guided flow that creates a Backend + Servers + HTTP Frontend (and optional HTTPS Frontend) in one atomic transaction. Endpoints (`backend/routers/site_wizard.py`): POST /api/site-wizard/preview — diff-preview the changeset POST /api/site-wizard/create — atomic execute POST /api/site-wizard/reject — clean rollback (including any wizard_staged ACME orders) GET /api/site-wizard/drafts — draft persistence PUT /api/site-wizard/drafts/{id} — save/update DELETE /api/site-wizard/drafts/{id} Feature surface: * One screen captures both backend (mode + servers) AND frontend (http + optional https + SSL mode) inputs. * SSL modes: ACME (new order, HTTP-01 only for v1.5.0), Upload (existing PEM), Existing (link to a stored cert), or None. * ACME-staged path: wizard_staged_until watermark on the `letsencrypt_orders` row defers finalisation until agent confirmation; per-mode reject cleanly cancels and rolls back the staged order. * Live diff preview against the cluster's current generated config (renderer-evolution noise stripped — track-sc<N> dedup, per-server cookie strip, defaults-cookie inheritance, listen-block flattening). * Draft persistence with PEM stripped at save time (private keys never round-trip through the drafts table). * Per-cluster multi-tenancy: drafts and wizard_staged orders are isolated to the creating user's cluster scope. Frontend (`frontend/src/components/SiteWizard.js`): * 4-step Antd Steps flow: Backend → Frontend → SSL → Review. * Render the live diff preview inline before commit. * Antd Form-level validation mirrors backend Pydantic validators (numeric bounds, HAProxy reserved keywords, ALPN consistency, IPv6 scope-id, domain regex, server name dedup). ------------------------------------------------------------------ AUDIT CAMPAIGN — Rounds 1 → 22 (Bulgu #1 → #82) ------------------------------------------------------------------ v1.5.0 includes 22 adversarial review passes. Each round produced its own commit set in the corporate development line; this squash collapses those into the v1.5.0 release artefact. Highlights: Round 1-4 Site Wizard core: dry-run parity, single-line value injection guard, ACL -f pattern-file block, SSL parity, timeout regex, form-state pin. Round 5-7 defaults-cookie inheritance, server-named-cookie guard, fe/be mode mismatch, duplicate server names, health_check_uri + server_address validators. Round 8-10 cookie_name / cookie_options newline-injection guard, dry-run parity (round 9), TCP-mode HTTP-only feature blockers. Round 11 SSL name path traversal + health-check >= 1. Round 12-13 SSL & ACME deep dive (Bulgu #23-#32). Round 14 single-line value injection (Bulgu #33). Round 15-17 ACME multi-tenant UX, numeric bounds, HAProxy reserved keywords, ALPN/TLS consistency, all-backup, multi-domain & multi-user enterprise edges, drain/HSTS/post-completion (Bulgu #34-#53). Round 18-21 concurrency, agent state, TCP-mode HTTP-only, list size caps, IPv6 scope-id, preview account validation, TCP backend + balance uri reject (Bulgu #54-#61). Round 22 FE error visibility + 3x stale-data lockouts, referential integrity + cascade safety, authentication & authorization, multi-cluster isolation, apply_pending_changes concurrency, script injection + bulk import multi-tenancy, prefix-stripped signature comparison (Bulgu #62-#82). ------------------------------------------------------------------ NO CORPORATE-SPECIFIC ARTIFACTS ------------------------------------------------------------------ This squash deliberately sanitises corporate hostnames, container registry references, and TLS secret names into generic placeholders (`your-registry.example.com/your-org`, `haproxy-openmanager*.example.com`, `wildcard-tls`, `taylanbakircioglu/haproxy-openmanager-*`) so the public artefact contains no internal infrastructure detail. Pilot / development history that retained those values stays in the corporate fork and is NOT part of this commit.
851 lines
37 KiB
Python
851 lines
37 KiB
Python
from fastapi import APIRouter, HTTPException, Request, Header
|
|
from typing import Optional
|
|
import logging
|
|
import time
|
|
import hashlib
|
|
import json
|
|
|
|
from models import WAFRule, WAFRuleUpdate
|
|
from database.connection import get_database_connection, close_database_connection
|
|
from utils.activity_log import log_user_activity
|
|
from services.haproxy_config import generate_haproxy_config_for_cluster
|
|
from .waf_helpers import assign_frontends_and_get_clusters, create_pending_configs_for_clusters, log_activity
|
|
|
|
|
|
router = APIRouter(prefix="/api/waf", tags=["waf"])
|
|
logger = logging.getLogger(__name__)
|
|
|
|
async def validate_user_cluster_access(user_id: int, cluster_id: int, conn):
|
|
"""Validate that user has access to the specified cluster"""
|
|
# Check if cluster exists
|
|
cluster_exists = await conn.fetchval("""
|
|
SELECT id FROM haproxy_clusters WHERE id = $1
|
|
""", cluster_id)
|
|
|
|
if not cluster_exists:
|
|
raise HTTPException(
|
|
status_code=404,
|
|
detail="Cluster not found"
|
|
)
|
|
|
|
# Check if user is admin - admins can access everything
|
|
is_admin = await conn.fetchval("""
|
|
SELECT is_admin FROM users WHERE id = $1
|
|
""", user_id)
|
|
|
|
if is_admin:
|
|
logger.info(f"Admin user {user_id} granted access to cluster {cluster_id}")
|
|
return True
|
|
|
|
# Check if user_pool_access table exists (for backward compatibility)
|
|
table_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.tables
|
|
WHERE table_name = 'user_pool_access'
|
|
)
|
|
""")
|
|
|
|
if not table_exists:
|
|
# Fallback to basic validation if table doesn't exist yet
|
|
logger.warning("user_pool_access table not found, using basic cluster validation")
|
|
return True
|
|
|
|
# Check if expires_at column exists (for backward compatibility)
|
|
expires_at_exists = await conn.fetchval("""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.columns
|
|
WHERE table_name = 'user_pool_access' AND column_name = 'expires_at'
|
|
)
|
|
""")
|
|
|
|
# Proper user-pool access validation
|
|
if expires_at_exists:
|
|
user_access = await conn.fetchrow("""
|
|
SELECT upa.access_level, hc.id, hc.name
|
|
FROM haproxy_clusters hc
|
|
JOIN haproxy_cluster_pools hcp ON hc.pool_id = hcp.id
|
|
JOIN user_pool_access upa ON hcp.id = upa.pool_id
|
|
WHERE upa.user_id = $1 AND hc.id = $2 AND upa.is_active = TRUE
|
|
AND (upa.expires_at IS NULL OR upa.expires_at > CURRENT_TIMESTAMP)
|
|
""", user_id, cluster_id)
|
|
else:
|
|
# Fallback query without expires_at column
|
|
user_access = await conn.fetchrow("""
|
|
SELECT upa.access_level, hc.id, hc.name
|
|
FROM haproxy_clusters hc
|
|
JOIN haproxy_cluster_pools hcp ON hc.pool_id = hcp.id
|
|
JOIN user_pool_access upa ON hcp.id = upa.pool_id
|
|
WHERE upa.user_id = $1 AND hc.id = $2 AND upa.is_active = TRUE
|
|
""", user_id, cluster_id)
|
|
|
|
if not user_access:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="You don't have access to this cluster. Please contact your administrator."
|
|
)
|
|
|
|
logger.info(f"User {user_id} granted {user_access['access_level']} access to cluster {cluster_id}")
|
|
return True
|
|
|
|
@router.get("/stats")
|
|
async def get_waf_stats(cluster_id: Optional[int] = None, authorization: str = Header(None)):
|
|
"""Get WAF statistics and overview, optionally filtered by cluster"""
|
|
try:
|
|
# Get current user for cluster validation
|
|
from auth_middleware import get_current_user_from_token
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Validate cluster access if cluster_id provided
|
|
if cluster_id:
|
|
await validate_user_cluster_access(current_user['id'], cluster_id, conn)
|
|
|
|
# Build cluster filter condition
|
|
cluster_filter = ""
|
|
cluster_params = []
|
|
if cluster_id:
|
|
cluster_filter = " WHERE cluster_id = $1"
|
|
cluster_params = [cluster_id]
|
|
|
|
# Get total WAF rules count (cluster-filtered)
|
|
total_rules = await conn.fetchval(f"SELECT COUNT(*) FROM waf_rules{cluster_filter}", *cluster_params)
|
|
|
|
# Get rules by action type
|
|
action_stats = await conn.fetch("""
|
|
SELECT action, COUNT(*) as count
|
|
FROM waf_rules
|
|
GROUP BY action
|
|
ORDER BY count DESC
|
|
""")
|
|
|
|
# Get rules by rule type
|
|
type_stats = await conn.fetch("""
|
|
SELECT rule_type, COUNT(*) as count
|
|
FROM waf_rules
|
|
GROUP BY rule_type
|
|
ORDER BY count DESC
|
|
""")
|
|
|
|
# Get recently created rules (last 30 days)
|
|
recent_rules = await conn.fetchval("""
|
|
SELECT COUNT(*) FROM waf_rules
|
|
WHERE created_at >= CURRENT_TIMESTAMP - INTERVAL '30 days'
|
|
""")
|
|
|
|
# Get rules with frontend assignments (with schema safety)
|
|
try:
|
|
assigned_rules = await conn.fetchval("""
|
|
SELECT COUNT(DISTINCT waf_rule_id) FROM frontend_waf_rules
|
|
""")
|
|
|
|
# Get top frontend assignments
|
|
top_frontends = await conn.fetch("""
|
|
SELECT f.name, f.id, COUNT(*) as rule_count
|
|
FROM frontend_waf_rules fwr
|
|
JOIN frontends f ON fwr.frontend_id = f.id
|
|
GROUP BY f.id, f.name
|
|
ORDER BY rule_count DESC
|
|
LIMIT 5
|
|
""")
|
|
except Exception as relation_error:
|
|
logger.warning(f"frontend_waf_rules table doesn't exist: {relation_error}")
|
|
assigned_rules = 0
|
|
top_frontends = []
|
|
|
|
unassigned_rules = total_rules - (assigned_rules or 0)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
return {
|
|
"overview": {
|
|
"total_rules": total_rules or 0,
|
|
"assigned_rules": assigned_rules or 0,
|
|
"unassigned_rules": unassigned_rules,
|
|
"recent_rules": recent_rules or 0
|
|
},
|
|
"action_distribution": [
|
|
{"action": row["action"], "count": row["count"]}
|
|
for row in action_stats
|
|
],
|
|
"type_distribution": [
|
|
{"type": row["rule_type"], "count": row["count"]}
|
|
for row in type_stats
|
|
],
|
|
"top_frontends": [
|
|
{"frontend_id": row["id"], "frontend_name": row["name"], "rule_count": row["rule_count"]}
|
|
for row in top_frontends
|
|
]
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error fetching WAF stats: {e}")
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
@router.get("/rules", summary="Get WAF Rules", response_description="List of WAF rules")
|
|
async def get_waf_rules(cluster_id: Optional[int] = None):
|
|
"""
|
|
# Get WAF Rules
|
|
|
|
Retrieve all Web Application Firewall (WAF) rules for protecting web applications.
|
|
|
|
## Query Parameters
|
|
- **cluster_id** (optional): Filter rules by cluster ID
|
|
|
|
## Example Request
|
|
```bash
|
|
curl -X GET "{BASE_URL}/api/waf/rules?cluster_id=1" \\
|
|
-H "Authorization: Bearer eyJhbGciOiJIUz..."
|
|
```
|
|
|
|
## Example Response
|
|
```json
|
|
[
|
|
{
|
|
"id": 1,
|
|
"name": "Block SQL Injection",
|
|
"rule_type": "deny",
|
|
"pattern": ".*(\\'|\\\")(union|select|insert|drop).*",
|
|
"action": "deny",
|
|
"priority": 100,
|
|
"enabled": true,
|
|
"cluster_ids": [1, 2],
|
|
"frontend_names": ["web-frontend"]
|
|
}
|
|
]
|
|
```
|
|
"""
|
|
try:
|
|
conn = await get_database_connection()
|
|
|
|
# Detect if consolidated JSON config column exists (backward compatibility)
|
|
config_column_exists = await conn.fetchval(
|
|
"""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.columns
|
|
WHERE table_name = 'waf_rules' AND column_name = 'config'
|
|
)
|
|
"""
|
|
)
|
|
# Detect if last_config_status column exists (schema-safety)
|
|
status_column_exists = await conn.fetchval(
|
|
"""
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM information_schema.columns
|
|
WHERE table_name = 'waf_rules' AND column_name = 'last_config_status'
|
|
)
|
|
"""
|
|
)
|
|
|
|
# Build base query dynamically to avoid referencing non-existent column
|
|
config_select = "w.config" if config_column_exists else "'{}'::jsonb AS config"
|
|
status_select = "w.last_config_status" if status_column_exists else "'APPLIED'::text AS last_config_status"
|
|
|
|
# Base query to fetch WAF rules with their configuration and frontend assignments
|
|
base_query = f"""
|
|
SELECT w.id, w.name, w.rule_type, {config_select}, w.action,
|
|
w.priority, w.description, w.is_active, w.created_at, w.updated_at,
|
|
{status_select},
|
|
COALESCE(ARRAY_AGG(DISTINCT f.name) FILTER (WHERE f.name IS NOT NULL), ARRAY[]::VARCHAR[]) as frontend_names,
|
|
COALESCE(ARRAY_AGG(DISTINCT f.id) FILTER (WHERE f.id IS NOT NULL), ARRAY[]::INTEGER[]) as frontend_ids
|
|
FROM waf_rules w
|
|
LEFT JOIN frontend_waf_rules fwr ON w.id = fwr.waf_rule_id
|
|
LEFT JOIN frontends f ON fwr.frontend_id = f.id
|
|
"""
|
|
|
|
if cluster_id:
|
|
# Filter by cluster - include WAF rules that belong to this cluster directly or via frontend
|
|
waf_rules = await conn.fetch(
|
|
f"""{base_query}
|
|
WHERE w.cluster_id = $1 OR f.cluster_id = $1 OR (w.cluster_id IS NULL AND w.id NOT IN (SELECT waf_rule_id FROM frontend_waf_rules))
|
|
GROUP BY w.id
|
|
ORDER BY w.priority, w.name
|
|
""",
|
|
cluster_id
|
|
)
|
|
else:
|
|
# Fetch all WAF rules if no cluster_id is specified
|
|
waf_rules = await conn.fetch(
|
|
f"""{base_query}
|
|
GROUP BY w.id
|
|
ORDER BY w.priority, w.name
|
|
"""
|
|
)
|
|
|
|
# Check for pending configurations for the given cluster
|
|
pending_waf_ids = set()
|
|
if waf_rules and cluster_id:
|
|
try:
|
|
pending_configs = await conn.fetch("""
|
|
SELECT DISTINCT
|
|
(regexp_matches(version_name, '^waf-([0-9]+)-'))[1]::int as waf_id
|
|
FROM config_versions
|
|
WHERE cluster_id = $1 AND status = 'PENDING'
|
|
AND version_name ~ '^waf-[0-9]+-'
|
|
""", cluster_id)
|
|
pending_waf_ids = {pc["waf_id"] for pc in pending_configs if pc["waf_id"]}
|
|
except Exception as e:
|
|
logger.warning(f"WAF API: Failed to check pending configs: {e}")
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Format the response with debug logging
|
|
formatted_rules = []
|
|
for rule in waf_rules:
|
|
# Parse config if it's a string
|
|
config = rule["config"] or {}
|
|
if isinstance(config, str):
|
|
try:
|
|
import json
|
|
config = json.loads(config)
|
|
except:
|
|
config = {}
|
|
|
|
# DEBUG: Log config data for troubleshooting edit form issues
|
|
logger.info(f"WAF Rule {rule['id']} ({rule['name']}) - Config: {config}, Frontend IDs: {rule['frontend_ids']}")
|
|
|
|
formatted_rule = {
|
|
"id": rule["id"],
|
|
"name": rule["name"],
|
|
"rule_type": rule["rule_type"],
|
|
"config": config,
|
|
"action": rule["action"],
|
|
"priority": rule["priority"],
|
|
"description": rule["description"],
|
|
"is_active": rule["is_active"],
|
|
"last_config_status": rule.get("last_config_status", 'APPLIED'),
|
|
"created_at": rule["created_at"].isoformat().replace('+00:00', 'Z') if rule["created_at"] else None,
|
|
"updated_at": rule["updated_at"].isoformat().replace('+00:00', 'Z') if rule["updated_at"] else None,
|
|
"frontend_names": rule["frontend_names"],
|
|
"frontend_ids": rule["frontend_ids"],
|
|
"frontend_count": len(rule["frontend_ids"]),
|
|
"has_pending_config": rule["id"] in pending_waf_ids
|
|
}
|
|
formatted_rules.append(formatted_rule)
|
|
|
|
return {"rules": formatted_rules}
|
|
except Exception as e:
|
|
logger.error(f"Error getting WAF rules: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
@router.post("/rules", summary="Create WAF Rule", response_description="WAF rule created successfully")
|
|
async def create_waf_rule(waf_rule_data: dict, cluster_id: Optional[int] = None, request: Request = None, authorization: str = Header(None)):
|
|
"""
|
|
# Create WAF Rule
|
|
|
|
Create a new Web Application Firewall rule to protect against attacks.
|
|
|
|
## Request Body
|
|
- **name**: Rule name (required)
|
|
- **rule_type**: Rule type (deny, allow, rate_limit)
|
|
- **pattern**: Regex pattern to match (required)
|
|
- **action**: Action to take (deny, allow)
|
|
- **priority**: Rule priority (higher = applied first)
|
|
- **enabled**: Enable rule (default: true)
|
|
- **cluster_ids**: List of cluster IDs to apply rule
|
|
|
|
## Example Request
|
|
```bash
|
|
curl -X POST "{BASE_URL}/api/waf/rules" \\
|
|
-H "Authorization: Bearer eyJhbGciOiJIUz..." \\
|
|
-H "Content-Type: application/json" \\
|
|
-d '{
|
|
"name": "Block SQL Injection",
|
|
"rule_type": "deny",
|
|
"pattern": ".*(union|select|insert|drop).*",
|
|
"action": "deny",
|
|
"priority": 100,
|
|
"enabled": true,
|
|
"cluster_ids": [1]
|
|
}'
|
|
```
|
|
"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token, check_user_permission
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Check permission for WAF create
|
|
has_permission = await check_user_permission(current_user["id"], "waf", "create")
|
|
if not has_permission:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Insufficient permissions: waf.create required"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Validate cluster access if cluster_id provided
|
|
if cluster_id:
|
|
await validate_user_cluster_access(current_user['id'], cluster_id, conn)
|
|
|
|
# Prepare the config dictionary (support nested payload from UI)
|
|
if isinstance(waf_rule_data.get('config'), dict):
|
|
config = waf_rule_data['config']
|
|
else:
|
|
# Backward compatibility: flatten extra fields as config
|
|
config = {
|
|
k: v for k, v in waf_rule_data.items() if k not in [
|
|
'name', 'rule_type', 'action', 'priority', 'is_active',
|
|
'description', 'frontend_ids'
|
|
]
|
|
}
|
|
|
|
# Create a WAFRule object for validation
|
|
waf_rule = WAFRule(
|
|
name=waf_rule_data.get('name'),
|
|
rule_type=waf_rule_data.get('rule_type'),
|
|
action=waf_rule_data.get('action'),
|
|
priority=waf_rule_data.get('priority'),
|
|
is_active=waf_rule_data.get('is_active', True),
|
|
config=config,
|
|
description=waf_rule_data.get('description'),
|
|
frontend_ids=waf_rule_data.get('frontend_ids', [])
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Check if rule name exists for this cluster
|
|
if cluster_id:
|
|
existing = await conn.fetchrow("SELECT id FROM waf_rules WHERE name = $1 AND cluster_id = $2", waf_rule.name, cluster_id)
|
|
else:
|
|
existing = await conn.fetchrow("SELECT id FROM waf_rules WHERE name = $1 AND cluster_id IS NULL", waf_rule.name)
|
|
|
|
if existing:
|
|
await close_database_connection(conn)
|
|
cluster_info = f" in cluster {cluster_id}" if cluster_id else ""
|
|
raise HTTPException(status_code=400, detail=f"WAF rule '{waf_rule.name}' already exists{cluster_info}")
|
|
|
|
# CRITICAL VALIDATION: At least one frontend must be selected
|
|
# Use case: Prevent unintentional application to ALL frontends
|
|
# Risk: WAF rule without frontend = applies to ALL frontends in cluster (dangerous!)
|
|
if not waf_rule.frontend_ids or len(waf_rule.frontend_ids) == 0:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="At least one frontend must be selected for WAF rule. Please select target frontend(s) where this WAF rule should be applied."
|
|
)
|
|
|
|
async with conn.transaction():
|
|
rule_id = await conn.fetchval("""
|
|
INSERT INTO waf_rules (name, rule_type, config, action, priority, description, enabled, is_active, cluster_id)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) RETURNING id
|
|
""", waf_rule.name, waf_rule.rule_type, json.dumps(waf_rule.config),
|
|
waf_rule.action, waf_rule.priority, waf_rule.description, waf_rule.is_active, waf_rule.is_active, cluster_id)
|
|
|
|
# CRITICAL FIX: Mark as PENDING *before* generating config
|
|
await conn.execute("UPDATE waf_rules SET last_config_status = 'PENDING' WHERE id = $1", rule_id)
|
|
|
|
frontend_assignments, cluster_ids = await assign_frontends_and_get_clusters(conn, rule_id, waf_rule.frontend_ids, cluster_id)
|
|
|
|
sync_results = await create_pending_configs_for_clusters(conn, cluster_ids, "create", rule_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
if current_user and current_user.get('id'):
|
|
await log_activity(current_user, 'create', rule_id, waf_rule, frontend_assignments, cluster_ids, sync_results, request)
|
|
|
|
return {
|
|
"message": f"WAF rule '{waf_rule.name}' created successfully. Go to Apply Changes to activate.",
|
|
"id": rule_id,
|
|
"waf_rule": waf_rule.dict(),
|
|
"frontend_assignments": frontend_assignments,
|
|
"sync_results": sync_results
|
|
}
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Error creating WAF rule: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
@router.get("/rules/{waf_rule_id}/config-versions")
|
|
async def get_waf_rule_config_versions(waf_rule_id: int, authorization: str = Header(None)):
|
|
"""Get config version history for a specific WAF rule"""
|
|
try:
|
|
# Verify user authentication
|
|
from auth_middleware import get_current_user_from_token
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
# Get WAF rule info first
|
|
waf_info = await conn.fetchrow("""
|
|
SELECT w.id, w.name,
|
|
ARRAY_AGG(DISTINCT c.name) as cluster_names,
|
|
ARRAY_AGG(DISTINCT c.id) as cluster_ids
|
|
FROM waf_rules w
|
|
LEFT JOIN frontend_waf_rules fwr ON w.id = fwr.waf_rule_id
|
|
LEFT JOIN frontends f ON fwr.frontend_id = f.id
|
|
LEFT JOIN haproxy_clusters c ON f.cluster_id = c.id
|
|
WHERE w.id = $1
|
|
GROUP BY w.id, w.name
|
|
""", waf_rule_id)
|
|
|
|
if not waf_info:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="WAF rule not found")
|
|
|
|
# Get all APPLIED config versions that are related to this WAF rule across all clusters
|
|
cluster_ids = [cid for cid in waf_info['cluster_ids'] if cid] if waf_info['cluster_ids'] else []
|
|
|
|
if cluster_ids:
|
|
versions = await conn.fetch("""
|
|
SELECT cv.id, cv.version_name, cv.description, cv.status, cv.is_active,
|
|
cv.created_at, cv.file_size, cv.checksum, cv.cluster_id,
|
|
u.username as created_by_username,
|
|
c.name as cluster_name
|
|
FROM config_versions cv
|
|
LEFT JOIN users u ON cv.created_by = u.id
|
|
LEFT JOIN haproxy_clusters c ON cv.cluster_id = c.id
|
|
WHERE cv.cluster_id = ANY($1) AND cv.status = 'APPLIED'
|
|
AND cv.version_name ~ $2
|
|
ORDER BY cv.created_at DESC
|
|
""", cluster_ids, f'^waf-{waf_rule_id}-')
|
|
else:
|
|
versions = []
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Format the response
|
|
formatted_versions = []
|
|
for version in versions:
|
|
formatted_versions.append({
|
|
"id": version["id"],
|
|
"version_name": version["version_name"],
|
|
"description": version["description"] or "WAF rule configuration update",
|
|
"type": "WAF Rule",
|
|
"status": version["status"],
|
|
"is_active": version["is_active"],
|
|
"created_at": version["created_at"].isoformat().replace('+00:00', 'Z') if version["created_at"] else None,
|
|
"created_by": version["created_by_username"] or "System",
|
|
"cluster_name": version["cluster_name"],
|
|
"cluster_id": version["cluster_id"],
|
|
"file_size": version["file_size"],
|
|
"checksum": version["checksum"][:8] + "..." if version["checksum"] else "No checksum"
|
|
})
|
|
|
|
return {
|
|
"versions": formatted_versions,
|
|
"entity_info": {
|
|
"entityName": waf_info["name"],
|
|
"clusterNames": [name for name in waf_info["cluster_names"] if name] if waf_info["cluster_names"] else ["No Cluster"],
|
|
"clusterIds": cluster_ids
|
|
}
|
|
}
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error fetching WAF rule config versions: {e}")
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
@router.put("/rules/{rule_id}")
|
|
async def update_waf_rule(rule_id: int, waf_rule_data: dict, request: Request, authorization: str = Header(None)):
|
|
"""Update existing WAF rule with multiple frontend assignments"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token, check_user_permission
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Check permission for WAF update
|
|
has_permission = await check_user_permission(current_user["id"], "waf", "update")
|
|
if not has_permission:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Insufficient permissions: waf.update required"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
existing_rule = await conn.fetchrow("SELECT * FROM waf_rules WHERE id = $1", rule_id)
|
|
if not existing_rule:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="WAF rule not found")
|
|
|
|
# Bulgu #79 — WAF rules carry an optional cluster_id;
|
|
# validate that the operator can touch this cluster
|
|
# before mutating the row. WAF rules with cluster_id IS
|
|
# NULL are "global" and gated only by `waf.update`.
|
|
rule_cluster_id = existing_rule.get('cluster_id') if hasattr(existing_rule, 'get') else existing_rule['cluster_id']
|
|
if rule_cluster_id:
|
|
await validate_user_cluster_access(current_user['id'], rule_cluster_id, conn)
|
|
|
|
# Prepare the config dictionary for the update
|
|
import json
|
|
config = existing_rule['config']
|
|
if isinstance(config, str):
|
|
config = json.loads(config) if config else {}
|
|
elif config is None:
|
|
config = {}
|
|
|
|
# DEBUG: Log incoming data to trace config loss issue
|
|
import logging
|
|
logger = logging.getLogger(__name__)
|
|
logger.info(f"WAF Update Debug - Rule ID: {rule_id}")
|
|
logger.info(f"WAF Update Debug - Incoming data: {waf_rule_data}")
|
|
logger.info(f"WAF Update Debug - Existing config: {existing_rule['config']}")
|
|
logger.info(f"WAF Update Debug - Parsed config: {config}")
|
|
|
|
# CRITICAL FIX: Merge form fields into config object
|
|
# Frontend sends WAF-specific fields directly, not wrapped in 'config'
|
|
config_fields = [
|
|
# Header Filter fields
|
|
'header_name', 'header_condition', 'header_value',
|
|
# IP Filter fields
|
|
'ip_addresses', 'ip_action',
|
|
# Rate Limit fields
|
|
'rate_limit_requests', 'rate_limit_window',
|
|
# Request Filter fields
|
|
'http_method', 'path_pattern',
|
|
# Geo Block fields
|
|
'countries', 'geo_action',
|
|
# Size Limit fields
|
|
'max_request_size', 'max_header_size',
|
|
# Advanced/General fields (frontend field names)
|
|
'log_message', 'custom_condition', 'redirect_url',
|
|
# Legacy/alternative field names (backward compatibility)
|
|
'custom_log_message', 'custom_haproxy_condition', 'rate_limit', 'time_window', 'method', 'request_size_limit', 'geo_countries'
|
|
]
|
|
|
|
# Update config with any form fields that exist in the incoming data
|
|
for field in config_fields:
|
|
if field in waf_rule_data:
|
|
config[field] = waf_rule_data[field]
|
|
logger.info(f"WAF Update Debug - Updated config.{field} = {waf_rule_data[field]}")
|
|
|
|
# Also handle explicit config updates if provided
|
|
if 'config' in waf_rule_data:
|
|
logger.info(f"WAF Update Debug - New config from request: {waf_rule_data['config']}")
|
|
config.update(waf_rule_data['config'])
|
|
else:
|
|
logger.info("WAF Update Debug - No 'config' field in incoming data, using form fields")
|
|
|
|
logger.info(f"WAF Update Debug - Final merged config: {config}")
|
|
|
|
# Create a WAFRuleUpdate object for validation
|
|
waf_rule = WAFRuleUpdate(
|
|
name=waf_rule_data.get('name', existing_rule['name']),
|
|
rule_type=waf_rule_data.get('rule_type', existing_rule['rule_type']),
|
|
action=waf_rule_data.get('action', existing_rule['action']),
|
|
priority=waf_rule_data.get('priority', existing_rule['priority']),
|
|
is_active=waf_rule_data.get('is_active', existing_rule['is_active']),
|
|
config=config,
|
|
description=waf_rule_data.get('description', existing_rule['description']),
|
|
frontend_ids=waf_rule_data.get('frontend_ids')
|
|
)
|
|
|
|
if waf_rule.name != existing_rule["name"]:
|
|
# Check for name conflict within the same cluster
|
|
cluster_id = existing_rule.get("cluster_id")
|
|
if cluster_id:
|
|
name_exists = await conn.fetchrow("SELECT id FROM waf_rules WHERE name = $1 AND cluster_id = $2 AND id != $3", waf_rule.name, cluster_id, rule_id)
|
|
else:
|
|
name_exists = await conn.fetchrow("SELECT id FROM waf_rules WHERE name = $1 AND cluster_id IS NULL AND id != $2", waf_rule.name, rule_id)
|
|
|
|
if name_exists:
|
|
await close_database_connection(conn)
|
|
cluster_info = f" in cluster {cluster_id}" if cluster_id else ""
|
|
raise HTTPException(status_code=400, detail=f"WAF rule name '{waf_rule.name}' already exists{cluster_info}")
|
|
|
|
async with conn.transaction():
|
|
await conn.execute("""
|
|
UPDATE waf_rules SET
|
|
name = $1, rule_type = $2, action = $3, priority = $4,
|
|
description = $5, is_active = $6, updated_at = CURRENT_TIMESTAMP,
|
|
config = $7
|
|
WHERE id = $8
|
|
""", waf_rule.name, waf_rule.rule_type, waf_rule.action,
|
|
waf_rule.priority, waf_rule.description, waf_rule.is_active,
|
|
json.dumps(waf_rule.config), rule_id)
|
|
|
|
# CRITICAL FIX: Mark as PENDING *before* generating config
|
|
await conn.execute("UPDATE waf_rules SET last_config_status = 'PENDING' WHERE id = $1", rule_id)
|
|
|
|
# PHASE 2: Create entity snapshot for rollback
|
|
from utils.entity_snapshot import save_entity_snapshot
|
|
|
|
new_values = {
|
|
"name": waf_rule.name,
|
|
"rule_type": waf_rule.rule_type,
|
|
"action": waf_rule.action,
|
|
"priority": waf_rule.priority,
|
|
"description": waf_rule.description,
|
|
"is_active": waf_rule.is_active,
|
|
"config": json.dumps(waf_rule.config)
|
|
}
|
|
|
|
entity_snapshot_metadata = await save_entity_snapshot(
|
|
conn=conn,
|
|
entity_type="waf_rule",
|
|
entity_id=rule_id,
|
|
old_values=existing_rule,
|
|
new_values=new_values,
|
|
operation="UPDATE"
|
|
)
|
|
|
|
# Update frontend assignments if provided
|
|
frontend_assignments = []
|
|
cluster_ids = set()
|
|
|
|
# CRITICAL FIX: Check if frontend_ids was explicitly sent in the request payload
|
|
# We need to distinguish between:
|
|
# - frontend_ids not sent at all (None) → keep existing associations
|
|
# - frontend_ids sent as empty array ([]) → clear all associations
|
|
# - frontend_ids sent with values → update associations
|
|
frontend_ids_in_payload = 'frontend_ids' in waf_rule_data
|
|
|
|
if frontend_ids_in_payload:
|
|
# CRITICAL VALIDATION: If frontend_ids explicitly provided, must have at least one
|
|
# Use case: Prevent clearing all frontends (WAF rule would apply to nothing)
|
|
if not waf_rule.frontend_ids or len(waf_rule.frontend_ids) == 0:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(
|
|
status_code=400,
|
|
detail="At least one frontend must be selected for WAF rule. Cannot remove all frontend assignments. Please select target frontend(s)."
|
|
)
|
|
|
|
# Frontend IDs were explicitly provided in the request (could be [] or [1,2,3])
|
|
await conn.execute("DELETE FROM frontend_waf_rules WHERE waf_rule_id = $1", rule_id)
|
|
frontend_assignments, cluster_ids = await assign_frontends_and_get_clusters(conn, rule_id, waf_rule.frontend_ids or [], existing_rule.get("cluster_id"))
|
|
logger.info(f"WAF Update: Frontend associations updated for rule {rule_id}: {waf_rule.frontend_ids}")
|
|
else:
|
|
# Frontend IDs were not provided in the request → preserve existing associations
|
|
assigned_frontends = await conn.fetch("SELECT frontend_id FROM frontend_waf_rules WHERE waf_rule_id = $1", rule_id)
|
|
existing_frontend_ids = [af['frontend_id'] for af in assigned_frontends]
|
|
_ , cluster_ids = await assign_frontends_and_get_clusters(conn, rule_id, existing_frontend_ids, existing_rule.get("cluster_id"))
|
|
logger.info(f"WAF Update: Frontend associations preserved for rule {rule_id}: {existing_frontend_ids}")
|
|
|
|
# Pass entity snapshot to config version creation
|
|
sync_results = await create_pending_configs_for_clusters(conn, cluster_ids, "update", rule_id, entity_snapshot_metadata)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
if current_user and current_user.get('id'):
|
|
await log_activity(current_user, 'update', rule_id, waf_rule, frontend_assignments, cluster_ids, sync_results, request)
|
|
|
|
return {
|
|
"message": f"WAF rule '{waf_rule.name}' updated successfully. Go to Apply Changes to activate.",
|
|
"waf_rule": waf_rule.dict(),
|
|
"frontend_assignments": frontend_assignments,
|
|
"sync_results": sync_results
|
|
}
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Error updating WAF rule: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|
|
|
|
@router.post("/rules/{rule_id}/toggle")
|
|
async def toggle_waf_rule_status(
|
|
rule_id: int,
|
|
action: str = "toggle",
|
|
cluster_id: Optional[int] = None,
|
|
request: Request = None,
|
|
authorization: str = Header(None)
|
|
):
|
|
"""Toggle, delete or enable/disable a WAF rule and create a pending config change"""
|
|
try:
|
|
from auth_middleware import get_current_user_from_token, check_user_permission
|
|
current_user = await get_current_user_from_token(authorization)
|
|
|
|
# Check permission for WAF toggle (uses update permission for toggle/enable/disable, delete permission for delete action)
|
|
if action == "delete":
|
|
has_permission = await check_user_permission(current_user["id"], "waf", "delete")
|
|
if not has_permission:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Insufficient permissions: waf.delete required"
|
|
)
|
|
else:
|
|
has_permission = await check_user_permission(current_user["id"], "waf", "toggle")
|
|
if not has_permission:
|
|
raise HTTPException(
|
|
status_code=403,
|
|
detail="Insufficient permissions: waf.toggle required"
|
|
)
|
|
|
|
conn = await get_database_connection()
|
|
|
|
rule = await conn.fetchrow("SELECT * FROM waf_rules WHERE id = $1", rule_id)
|
|
if not rule:
|
|
await close_database_connection(conn)
|
|
raise HTTPException(status_code=404, detail="WAF rule not found")
|
|
|
|
# Bulgu #79 — validate cluster access for cluster-scoped
|
|
# WAF rules (global rules pass unconditionally).
|
|
rule_cluster_id = rule['cluster_id'] if 'cluster_id' in rule.keys() else None
|
|
if rule_cluster_id:
|
|
await validate_user_cluster_access(current_user['id'], rule_cluster_id, conn)
|
|
|
|
# Determine new status based on action
|
|
if action == "delete":
|
|
new_status = False
|
|
action_name = "delete"
|
|
elif action == "enable":
|
|
new_status = True
|
|
action_name = "enable"
|
|
elif action == "disable":
|
|
new_status = False
|
|
action_name = "disable"
|
|
else: # toggle
|
|
new_status = not rule['is_active']
|
|
action_name = "enable" if new_status else "disable"
|
|
|
|
async with conn.transaction():
|
|
await conn.execute(
|
|
"UPDATE waf_rules SET is_active = $1, updated_at = CURRENT_TIMESTAMP WHERE id = $2",
|
|
new_status, rule_id
|
|
)
|
|
|
|
# CRITICAL FIX: Mark PENDING *before* generating config
|
|
await conn.execute("UPDATE waf_rules SET last_config_status = 'PENDING' WHERE id = $1", rule_id)
|
|
|
|
# Get affected clusters to create pending changes
|
|
assigned_frontends = await conn.fetch("SELECT frontend_id FROM frontend_waf_rules WHERE waf_rule_id = $1", rule_id)
|
|
rule_cluster_id = await conn.fetchval("SELECT cluster_id FROM waf_rules WHERE id = $1", rule_id)
|
|
_, cluster_ids = await assign_frontends_and_get_clusters(conn, rule_id, [af['frontend_id'] for af in assigned_frontends], rule_cluster_id)
|
|
|
|
# If no clusters resolved via assignments or rule, but a cluster_id is provided by caller, use it
|
|
if not cluster_ids and cluster_id:
|
|
cluster_ids = {cluster_id}
|
|
|
|
# If no clusters resolved via assignments or rule, but a cluster_id is provided by caller, use it
|
|
if (not cluster_ids or len(cluster_ids) == 0) and cluster_id:
|
|
cluster_ids = {cluster_id}
|
|
|
|
sync_results = await create_pending_configs_for_clusters(conn, cluster_ids, action_name, rule_id)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
# Log the activity
|
|
if current_user and current_user.get('id'):
|
|
await log_user_activity(
|
|
user_id=current_user['id'],
|
|
action=action,
|
|
resource_type='waf_rule',
|
|
resource_id=str(rule_id),
|
|
details={
|
|
'waf_rule_name': rule['name'],
|
|
'new_status': 'active' if new_status else 'inactive',
|
|
'affected_clusters': list(cluster_ids)
|
|
},
|
|
ip_address=str(request.client.host) if request.client else None,
|
|
user_agent=request.headers.get('user-agent')
|
|
)
|
|
|
|
# Return appropriate message based on action
|
|
if action == "delete":
|
|
message = f"WAF rule '{rule['name']}' marked for deletion. Go to Apply Changes to remove it from agents."
|
|
elif action_name == "enable":
|
|
message = f"WAF rule '{rule['name']}' enabled. Go to Apply Changes to activate."
|
|
else:
|
|
message = f"WAF rule '{rule['name']}' disabled. Go to Apply Changes to deactivate."
|
|
|
|
return {
|
|
"message": message,
|
|
"new_status": new_status,
|
|
"sync_results": sync_results
|
|
}
|
|
except HTTPException:
|
|
raise
|
|
except Exception as e:
|
|
logger.error(f"Error toggling WAF rule: {e}", exc_info=True)
|
|
raise HTTPException(status_code=500, detail=str(e))
|
|
|