Files
haproxy-openmanager/backend/haproxy_client.py
T
mustafa.ulukaya ef26860df9 feat(logging): unified request/response log with configurable retention (v1.11.0)
Until now the only record of what happened was `user_activity_logs`, which
stores non-GET 2xx operations with no bodies. When something failed you could
see that a counter went up, never what was sent or what came back.

This adds one queryable timeline covering both directions:

- inbound: every API call, including GETs and including 4xx/5xx, with the
  user, client IP, status, duration and — redacted, size-capped — the request
  and response bodies.
- outbound: every HTTP call the backend makes, tagged with who it went to
  (ACME/Let's Encrypt, Cloudflare, GoDaddy, HAProxy stats, agents, the ACME
  diagnostics probe).

Outbound rows inherit the inbound request's id, so one operator action and the
CA/DNS calls it triggered read as a single trace: opening a failed "Request
Certificate" shows the exact POST /acme/new-order and the CA's 429 underneath.

Implementation notes:

- Capture is a pure-ASGI middleware that TEES the request and response streams
  rather than draining them. `await request.body()` inside a BaseHTTPMiddleware
  would consume the receive channel and break the raw-body agent heartbeat
  handler. Registered last so it is outermost: it then sees the final
  client-visible response and seeds correlation_id_context before the error
  handler reads it.
- Rows are written by a batching background writer with a bounded queue, so the
  request path never awaits the database and a saturated logger drops rows
  visibly (surfaced on the page) instead of blocking. Redaction runs on the
  writer, off the request coroutine.
- Secrets never land: headers are an allowlist with Authorization/Cookie kept
  only as a presence marker; body keys and value shapes are redacted
  (passwords, tokens, api_token, API keys, private-key PEMs, JWTs); the ACME
  JWS request body is never stored, because a stored protected+signature pair
  is a replayable credential — a summary is logged instead; DNS-provider errors
  record only the exception type; the ACME HTTP-01 challenge endpoint is
  excluded so key_authorization is never captured.
- Retention is operator-configurable in Settings -> Request Log: separate day
  counts for successful and failed rows (7 / 30) plus a hard row cap (500k),
  whichever is reached first. Pruned in batches under a Postgres advisory lock,
  with the day counts bound as parameters, never interpolated.
- New permissions requestlog.read / requestlog.manage. super_admin and
  security_admin get both, operator gets read, viewer gets neither.

Schema: one new table (request_logs) plus its settings seed, SCHEMA_VERSION
10 -> 11, auto-migrated. No existing table altered, no agent or rendered-config
change. Kill switches: REQUEST_LOG_ENABLED=false (middleware never registered)
or the `enabled` toggle in Settings.

Tests: 245 new (7 backend files + 1 frontend), full suite 1655 backend +
17 frontend passing.
2026-08-11 02:36:03 +03:00

316 lines
13 KiB
Python

import asyncio
import aiohttp
import socket
import csv
import io
from typing import Dict, List, Optional, Any
import logging
logger = logging.getLogger(__name__)
class HAProxyClient:
"""Client for connecting to HAProxy instances and fetching stats"""
def __init__(self, cluster_or_host, port: Optional[int] = None, connection_type: str = 'stats_socket',
stats_socket_path: str = '/run/haproxy/admin.sock',
stats_username: Optional[str] = None,
stats_password: Optional[str] = None):
# Support both dict (cluster) and string (host) parameters
if isinstance(cluster_or_host, dict):
# Initialize from cluster dict
cluster = cluster_or_host
self.host = cluster.get('host', 'localhost')
self.port = cluster.get('port', 8404)
self.connection_type = cluster.get('connection_type', 'local')
self.stats_socket_path = cluster.get('stats_socket_path', '/run/haproxy/admin.sock')
self.stats_username = cluster.get('stats_username')
self.stats_password = cluster.get('stats_password')
else:
# Initialize from individual parameters (backward compatibility)
self.host = cluster_or_host
self.port = port or 8404
self.connection_type = connection_type
self.stats_socket_path = stats_socket_path
self.stats_username = stats_username
self.stats_password = stats_password
async def get_stats(self) -> Dict[str, Any]:
"""Get HAProxy statistics"""
try:
if self.connection_type == 'stats_socket':
# Try to get real stats, but for demo purposes always use fallback
# return await self._get_stats_via_socket()
return self._get_fallback_stats()
elif self.connection_type == 'local':
# Try to get real stats, but for demo purposes always use fallback
# return await self._get_stats_via_http()
return self._get_fallback_stats()
else:
# For other connection types, use HTTP stats page
# return await self._get_stats_via_http()
return self._get_fallback_stats()
except Exception as e:
logger.error(f"Failed to get HAProxy stats: {e}")
return self._get_fallback_stats()
async def _get_stats_via_http(self) -> Dict[str, Any]:
"""Get stats via HTTP stats page"""
try:
url = f"http://{self.host}:{self.port}/stats?stats;csv"
auth = None
if self.stats_username and self.stats_password:
auth = aiohttp.BasicAuth(self.stats_username, self.stats_password)
# v1.11.0: instrumented for completeness. NOTE the CSV body is
# deliberately NOT handed to the span — a full stats dump is large,
# changes every poll, and has no diagnostic value in an audit row;
# status + duration is what matters. `auth` is likewise never logged:
# aiohttp.BasicAuth is a NamedTuple whose repr contains the cleartext
# password.
from utils.http_instrumentation import outbound_span, TARGET_HAPROXY_STATS
async with aiohttp.ClientSession() as session:
async with outbound_span(
target=TARGET_HAPROXY_STATS, method="GET", url=url,
capture_body=False, capture_response_body=False,
) as span:
async with session.get(url, auth=auth, timeout=aiohttp.ClientTimeout(total=10)) as response:
span.set_response(response.status, getattr(response, "headers", None))
if response.status == 200:
csv_data = await response.text()
return self._parse_csv_stats(csv_data)
else:
logger.warning(f"HTTP stats request failed with status {response.status}")
return self._get_fallback_stats()
except Exception as e:
logger.error(f"HTTP stats request failed: {e}")
return self._get_fallback_stats()
async def _get_stats_via_socket(self) -> Dict[str, Any]:
"""Get stats via Unix socket (for Docker containers)"""
try:
# For Docker containers, we'll use HTTP as socket access is more complex
# In a real deployment, you'd use actual socket connections
return await self._get_stats_via_http()
except Exception as e:
logger.error(f"Socket stats request failed: {e}")
return self._get_fallback_stats()
def _parse_csv_stats(self, csv_data: str) -> Dict[str, Any]:
"""Parse HAProxy CSV stats format"""
try:
lines = csv_data.strip().split('\n')
if not lines:
return self._get_fallback_stats()
# Skip comment lines and get header
data_lines = [line for line in lines if not line.startswith('#')]
if not data_lines:
return self._get_fallback_stats()
reader = csv.DictReader(io.StringIO('\n'.join(data_lines)))
rows = list(reader)
return self._process_haproxy_stats(rows)
except Exception as e:
logger.error(f"Failed to parse CSV stats: {e}")
return self._get_fallback_stats()
def _process_haproxy_stats(self, rows: List[Dict]) -> Dict[str, Any]:
"""Process HAProxy stats into our format"""
try:
backends = []
servers = []
frontend_stats = {"requests_total": 0, "requests_rate": 0, "status": "OPEN"}
for row in rows:
pxname = row.get('# pxname', '')
svname = row.get('svname', '')
status = row.get('status', 'UP')
# Skip empty rows
if not pxname or not svname:
continue
# Process backend stats
if svname == 'BACKEND':
total_servers = int(row.get('act', 0)) + int(row.get('bck', 0))
active_servers = int(row.get('act', 0))
backends.append({
"name": pxname,
"status": "UP" if status in ['UP', 'OPEN'] else "DOWN",
"active_servers": active_servers,
"total_servers": max(total_servers, 1), # Avoid division by zero
"requests": int(row.get('stot', 0)),
"response_time": int(float(row.get('rtime', 0)) or 0)
})
# Process individual server stats
elif svname != 'FRONTEND' and svname != 'BACKEND':
servers.append({
"name": svname,
"status": "UP" if status == 'UP' else "DOWN",
"weight": int(row.get('weight', 100)),
"requests": int(row.get('stot', 0)),
"backend": pxname,
"address": f"{row.get('addr', 'unknown')}:{row.get('port', '0')}"
})
# Process frontend stats
elif svname == 'FRONTEND':
frontend_stats["requests_total"] += int(row.get('stot', 0))
frontend_stats["requests_rate"] += float(row.get('req_rate', 0) or 0)
if status not in ['UP', 'OPEN']:
frontend_stats["status"] = "MAINT"
return {
"frontend": frontend_stats,
"backends": backends,
"servers": servers
}
except Exception as e:
logger.error(f"Failed to process HAProxy stats: {e}")
return self._get_fallback_stats()
def _get_fallback_stats(self) -> Dict[str, Any]:
"""Return fallback stats when connection fails"""
import random
from datetime import datetime
# Generate realistic demo data for development
current_time = datetime.now()
base_requests = 1000 + random.randint(0, 500)
backends = []
servers = []
# Generate sample backends with different health statuses
backend_names = ["web_servers", "api_servers", "db_pool", "cache_cluster"]
for i, backend_name in enumerate(backend_names):
total_servers = random.randint(2, 4)
active_servers = random.randint(max(1, total_servers-1), total_servers)
status = "UP" if active_servers > 0 else "DOWN"
backend_requests = max(0, base_requests + random.randint(-200, 200))
response_time = random.randint(50, 300) if status == "UP" else 0
backends.append({
"name": backend_name,
"status": status,
"active_servers": active_servers,
"total_servers": total_servers,
"requests": backend_requests,
"response_time": response_time,
"health_percentage": int((active_servers / total_servers) * 100)
})
# Generate servers for this backend
for j in range(total_servers):
server_name = f"{backend_name}-{j+1}"
server_status = "UP" if j < active_servers else "DOWN"
server_weight = random.randint(50, 150)
server_requests = int(backend_requests / total_servers) + random.randint(-50, 50)
server_port = 80 if "web" in backend_name else (3000 if "api" in backend_name else (5432 if "db" in backend_name else 6379))
servers.append({
"name": server_name,
"status": server_status,
"weight": server_weight,
"requests": server_requests if server_status == "UP" else 0,
"backend": backend_name,
"address": f"192.168.1.{10 + len(servers)}:{server_port}",
"check_status": "L4OK" if server_status == "UP" else "L4TOUT",
"last_check": "0ms" if server_status == "UP" else "timeout"
})
return {
"frontend": {
"requests_total": base_requests,
"requests_rate": round(random.uniform(10.0, 50.0), 2),
"status": "OPEN"
},
"backends": backends,
"servers": servers,
"traffic_distribution": [
{"name": backend["name"], "value": backend["requests"], "percentage": int((backend["requests"] / base_requests) * 100)}
for backend in backends
]
}
async def test_connection(self) -> Dict[str, Any]:
"""Test connection to HAProxy instance"""
try:
stats = await self.get_stats()
if stats["frontend"]["status"] != "MAINT":
return {
"status": "connected",
"error": None,
"version": "2.4.0" # Could be parsed from stats if available
}
else:
return {
"status": "error",
"error": "HAProxy appears to be in maintenance mode",
"version": None
}
except Exception as e:
return {
"status": "error",
"error": str(e),
"version": None
}
async def get_config(self, config_path: str = "/etc/haproxy/haproxy.cfg") -> str:
"""Get HAProxy configuration file content"""
try:
if self.connection_type == 'local':
# For local connections, read from local file system
import os
if os.path.exists(config_path):
with open(config_path, 'r') as f:
return f.read()
else:
# Try common paths
for path in ["/usr/local/etc/haproxy/haproxy.cfg", "/etc/haproxy.cfg", "/opt/haproxy/haproxy.cfg"]:
if os.path.exists(path):
with open(path, 'r') as f:
return f.read()
raise Exception(f"HAProxy config file not found at {config_path}")
else:
# For remote connections, this would need SSH access
# For now, return a demo config
return """# HAProxy Configuration
# This is a demo configuration - implement SSH access for remote configs
global
daemon
log stdout local0
defaults
mode http
log global
option httplog
timeout connect 5s
timeout client 50s
timeout server 50s
frontend stats
bind *:8404
stats enable
stats uri /stats
stats admin if TRUE
frontend main
bind *:80
default_backend servers
backend servers
balance roundrobin
server web1 127.0.0.1:8001 check
server web2 127.0.0.1:8002 check
"""
except Exception as e:
logger.error(f"Failed to get HAProxy config: {e}")
raise Exception(f"Could not retrieve HAProxy configuration: {str(e)}")