mirror of
https://github.com/taylanbakircioglu/haproxy-openmanager.git
synced 2026-09-16 15:45:11 +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.
273 lines
10 KiB
Python
273 lines
10 KiB
Python
import asyncio
|
|
import json
|
|
import logging
|
|
from datetime import datetime
|
|
from typing import Optional, Dict, Any
|
|
from database.connection import get_database_connection, close_database_connection
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
async def log_user_activity(
|
|
user_id: int,
|
|
action: str,
|
|
resource_type: str,
|
|
resource_id: Optional[str] = None,
|
|
details: Optional[Dict[str, Any]] = None,
|
|
ip_address: Optional[str] = None,
|
|
user_agent: Optional[str] = None
|
|
):
|
|
"""Log user activity to the database"""
|
|
try:
|
|
conn = await get_database_connection()
|
|
|
|
# Serialize details to JSON if it's a dict
|
|
details_json = json.dumps(details) if details and isinstance(details, dict) else details
|
|
|
|
# Try with created_at column first, then fallback
|
|
try:
|
|
await conn.execute("""
|
|
INSERT INTO user_activity_logs
|
|
(user_id, action, resource_type, resource_id, details, ip_address, user_agent, created_at)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
|
|
""", user_id, action, resource_type, resource_id,
|
|
details_json, ip_address, user_agent, datetime.utcnow())
|
|
except Exception as column_error:
|
|
logger.warning(f"created_at column not found, trying fallback: {column_error}")
|
|
# Fallback without created_at column
|
|
await conn.execute("""
|
|
INSERT INTO user_activity_logs
|
|
(user_id, action, resource_type, resource_id, details, ip_address, user_agent)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7)
|
|
""", user_id, action, resource_type, resource_id,
|
|
details_json, ip_address, user_agent)
|
|
|
|
await close_database_connection(conn)
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to log user activity: {e}")
|
|
# Don't raise exception - activity logging should not break the main flow
|
|
|
|
async def get_user_activity_logs(
|
|
user_id: Optional[int] = None,
|
|
limit: int = 100,
|
|
offset: int = 0
|
|
) -> list:
|
|
"""Get user activity logs"""
|
|
try:
|
|
conn = await get_database_connection()
|
|
|
|
# Schema-safe activity log queries
|
|
try:
|
|
if user_id:
|
|
logs = await conn.fetch("""
|
|
SELECT ual.*, u.username
|
|
FROM user_activity_logs ual
|
|
LEFT JOIN users u ON ual.user_id = u.id
|
|
WHERE ual.user_id = $1
|
|
ORDER BY ual.created_at DESC
|
|
LIMIT $2 OFFSET $3
|
|
""", user_id, limit, offset)
|
|
else:
|
|
logs = await conn.fetch("""
|
|
SELECT ual.*, u.username
|
|
FROM user_activity_logs ual
|
|
LEFT JOIN users u ON ual.user_id = u.id
|
|
ORDER BY ual.created_at DESC
|
|
LIMIT $1 OFFSET $2
|
|
""", limit, offset)
|
|
except Exception as schema_error:
|
|
logger.warning(f"Schema error in activity logs, using fallback: {schema_error}")
|
|
# Fallback without ORDER BY created_at
|
|
if user_id:
|
|
logs = await conn.fetch("""
|
|
SELECT ual.*, u.username
|
|
FROM user_activity_logs ual
|
|
LEFT JOIN users u ON ual.user_id = u.id
|
|
WHERE ual.user_id = $1
|
|
LIMIT $2 OFFSET $3
|
|
""", user_id, limit, offset)
|
|
else:
|
|
logs = await conn.fetch("""
|
|
SELECT ual.*, u.username
|
|
FROM user_activity_logs ual
|
|
LEFT JOIN users u ON ual.user_id = u.id
|
|
LIMIT $1 OFFSET $2
|
|
""", limit, offset)
|
|
|
|
await close_database_connection(conn)
|
|
return logs
|
|
|
|
except Exception as e:
|
|
logger.error(f"Failed to get user activity logs: {e}")
|
|
return []
|
|
|
|
|
|
# ----------------------------------------------------------------------------
|
|
# v1.5.0 Feature A (Issue #13): typed ACME order event log helper
|
|
# ----------------------------------------------------------------------------
|
|
|
|
|
|
async def record_event(
|
|
order_id: int,
|
|
event_type: str,
|
|
*,
|
|
severity: str = "INFO",
|
|
message: Optional[str] = None,
|
|
details: Optional[Dict[str, Any]] = None,
|
|
correlation_id: Optional[str] = None,
|
|
conn=None,
|
|
) -> Optional[int]:
|
|
"""Insert a typed event row into acme_order_events.
|
|
|
|
Wrapped in try/except so that an event-log DB failure NEVER breaks the
|
|
main ACME flow (Section 5.2 of the v1.5.0 plan).
|
|
|
|
M24: when `conn` is passed in (e.g. inside the
|
|
complete_pending_acme_orders pool-pressure-sensitive task), reuse the
|
|
existing connection instead of acquiring a new one from the pool.
|
|
|
|
Returns the inserted row id, or None on failure.
|
|
"""
|
|
own_conn = False
|
|
try:
|
|
if conn is None:
|
|
conn = await get_database_connection()
|
|
own_conn = True
|
|
|
|
details_json = json.dumps(details or {}) if not isinstance(details, str) else details
|
|
try:
|
|
row_id = await conn.fetchval(
|
|
"""
|
|
INSERT INTO acme_order_events (
|
|
order_id, event_type, severity, message, details, correlation_id
|
|
) VALUES ($1, $2, $3, $4, $5::jsonb, $6)
|
|
RETURNING id
|
|
""",
|
|
order_id,
|
|
event_type,
|
|
severity.upper() if severity else "INFO",
|
|
message,
|
|
details_json,
|
|
correlation_id,
|
|
)
|
|
return row_id
|
|
except Exception as e:
|
|
# Most likely cause: acme_order_events table missing in older
|
|
# deployments (migration not yet run). NEVER raise.
|
|
logger.debug(
|
|
f"record_event: insert failed (order_id={order_id}, "
|
|
f"event_type={event_type}): {e}"
|
|
)
|
|
return None
|
|
except Exception as e:
|
|
logger.debug(f"record_event: outer failure (order_id={order_id}): {e}")
|
|
return None
|
|
finally:
|
|
if own_conn and conn is not None:
|
|
try:
|
|
await close_database_connection(conn)
|
|
except Exception:
|
|
pass
|
|
|
|
|
|
async def prune_acme_events_and_drafts_if_due() -> Dict[str, int]:
|
|
"""Daily-watermarked TTL prune (M30 / Section 5.3 of v1.5.0 plan).
|
|
|
|
- acme_order_events: 90d retention.
|
|
- wizard_drafts: 30d retention (also pruned by expires_at < NOW() since
|
|
that column exists explicitly).
|
|
|
|
Watermarking via system_settings (dot-notation keys —
|
|
`acme.events_last_pruned_at` / `wizard.drafts_last_pruned_at`)
|
|
ensures multi-replica deployments only run the prune once per day.
|
|
|
|
Always returns a dict with the (possibly zero) prune counts. Never raises.
|
|
"""
|
|
counts = {"acme_events": 0, "wizard_drafts": 0}
|
|
conn = None
|
|
try:
|
|
conn = await get_database_connection()
|
|
|
|
async def _maybe_run(setting_key: str, ttl_query: str) -> int:
|
|
"""Returns # rows pruned, or 0 if not yet due."""
|
|
try:
|
|
row = await conn.fetchrow(
|
|
"SELECT value FROM system_settings WHERE key = $1",
|
|
setting_key,
|
|
)
|
|
last_at: Optional[datetime] = None
|
|
if row and row["value"] is not None:
|
|
raw = row["value"]
|
|
if isinstance(raw, str):
|
|
try:
|
|
raw = json.loads(raw)
|
|
except json.JSONDecodeError:
|
|
raw = None
|
|
if isinstance(raw, str):
|
|
try:
|
|
last_at = datetime.fromisoformat(raw.replace("Z", "+00:00"))
|
|
except ValueError:
|
|
last_at = None
|
|
|
|
if last_at is not None:
|
|
age_seconds = (datetime.utcnow() - last_at.replace(tzinfo=None)).total_seconds()
|
|
if age_seconds < 24 * 3600:
|
|
return 0
|
|
|
|
result = await conn.execute(ttl_query)
|
|
count = 0
|
|
if isinstance(result, str) and result.startswith("DELETE "):
|
|
try:
|
|
count = int(result.split()[-1])
|
|
except (ValueError, IndexError):
|
|
count = 0
|
|
|
|
ts_value = json.dumps(datetime.utcnow().isoformat() + "Z")
|
|
await conn.execute(
|
|
"""
|
|
INSERT INTO system_settings (key, value, category, description)
|
|
VALUES ($1, $2::jsonb, $3, $4)
|
|
ON CONFLICT (key) DO UPDATE
|
|
SET value = EXCLUDED.value, updated_at = CURRENT_TIMESTAMP
|
|
""",
|
|
setting_key,
|
|
ts_value,
|
|
"acme" if setting_key.startswith("acme.") else "wizard",
|
|
"Internal: last daily prune timestamp (v1.5.0)",
|
|
)
|
|
return count
|
|
except Exception as inner:
|
|
logger.debug(f"prune watermark step failed for {setting_key}: {inner}")
|
|
return 0
|
|
|
|
# acme_order_events 90d
|
|
counts["acme_events"] = await _maybe_run(
|
|
"acme.events_last_pruned_at",
|
|
"DELETE FROM acme_order_events WHERE created_at < NOW() - INTERVAL '90 days'",
|
|
)
|
|
# wizard_drafts 30d (also catches expires_at-passed rows)
|
|
counts["wizard_drafts"] = await _maybe_run(
|
|
"wizard.drafts_last_pruned_at",
|
|
"""
|
|
DELETE FROM wizard_drafts
|
|
WHERE created_at < NOW() - INTERVAL '30 days'
|
|
OR expires_at < NOW()
|
|
""",
|
|
)
|
|
|
|
if counts["acme_events"] or counts["wizard_drafts"]:
|
|
logger.info(
|
|
f"v1.5.0 daily prune: acme_events={counts['acme_events']} "
|
|
f"wizard_drafts={counts['wizard_drafts']}"
|
|
)
|
|
return counts
|
|
except Exception as e:
|
|
logger.debug(f"prune_acme_events_and_drafts_if_due: {e}")
|
|
return counts
|
|
finally:
|
|
if conn is not None:
|
|
try:
|
|
await close_database_connection(conn)
|
|
except Exception:
|
|
pass
|