Files
taylanbakircioglu 0ee227363e fix(agent): adoption cannot hide a node; Config Import on fresh installs (v1.11.1)
Six fixes to the agent's discovery path and one long-standing parity gap, found
by auditing it in loops against a real fleet.

DISCOVERY REPORTS THAT NEVER REACHED THE SERVER

The agent posts an unmanaged keepalived.conf for adoption and caches the hash of
what it sent, so the file - which carries the VRRP password - is re-posted only
when it changes. Delivery was judged by curl's exit code, and curl without -f
exits 0 on 5xx too, so a report the server REJECTED was recorded as delivered.
Since a hand-maintained config does not change on its own, that node dropped out
of "Unmanaged keepalived detected" permanently; the only cure was deleting a
cache file on the node by hand.

  - the report is cached only on a 2xx;
  - GET /agents/{name}/keepalived-config now reports whether the server actually
    holds a discovery for that agent, and the cache may only suppress while it
    says yes - which is what lets nodes stuck from earlier releases recover on
    their own, with nobody touching them;
  - the flag is parsed with has() + tostring, not `// empty`: jq's alternative
    operator returns the alternative for **false** as well as null, so the naive
    form could not tell "no record" from "older backend" and the recovery would
    have been completely inert;
  - a 400/413/422 records the refusal so identical bytes are not re-posted
    forever - 4xx and 5xx agent calls are never sampled out of the request log,
    so an unattended loop would write a row carrying the whole config every
    cycle - while 401 and 404 keep retrying, because here they mean a token
    rotation or an agent row briefly absent, not a bad payload;
  - the CLEAR path had the same exit-code defect, where it left a stale row
    offering a managed node for adoption with nothing to ever retry it.

CONFIG IMPORT WAS A NO-OP ON FRESHLY INSTALLED AGENTS

check_config_requests uploads a node's live haproxy.cfg on request. It was
defined in the installer body and in the self-upgrade daemon, but not in the
heredoc a fresh install writes, and its call site is guarded by `type` - so on
such a node the operator asked for a config and nothing arrived, with no error
anywhere. Any agent that had self-upgraded at least once already had it, which
is why it went unnoticed. The self-upgrade definition is copied verbatim
(verified line-for-line). A freshly installed agent now polls that endpoint once
per cycle exactly as every upgraded agent already does; no node running today
changes behaviour.

DETERMINISTIC CONFIG PATH

A pool may hold several clusters and the join that resolves keepalived_config_path
was unordered, so the path handed to an agent could differ between polls whenever
two clusters disagreed - the agent would inspect a file that is not there and the
node would never appear, intermittently. A customised path now wins over the
shipped default, then the lowest cluster id. Verified against a real PostgreSQL
over seven arrangements: with one cluster per pool, or when every cluster carries
the default, the value is byte-identical to before.

Verified end to end on a production fleet and, for each decision, against the
real _kp_discover block rather than a paraphrase.

Backend suite: 1674 passed, 152 skipped. bash -n passes on the whole file and on
the fresh-install body in isolation. The keepalived path is logic-identical
across both daemon copies, now pinned by a test.
2026-08-15 19:37:53 +03:00

3527 lines
171 KiB
Python

from fastapi import APIRouter, HTTPException, Header, Request, Depends
from typing import Optional
import logging
import time
import secrets
import base64
import string
import os
import json
import ipaddress
import hashlib
import re
# Pipeline trigger - force backend redeploy v2
# v1.10.4 — a discovered keepalived.conf is stored and served to the UI, so the VRRP password is
# masked out of the stored copy (the real value lives Fernet-encrypted in its own column). Mask
# the WHOLE remainder of the line, mirroring vip.py's version-diff masking, so a password
# containing whitespace cannot partially leak.
_AUTH_PASS_MASK_RE = re.compile(r"(auth_pass\s+).*")
from models import AgentCreate
from models.agent import AgentToggle, AgentHeartbeat, AgentScriptRequest, AgentUpgradeRequest
from database.connection import get_database_connection, close_database_connection
from utils.activity_log import log_user_activity
from auth_middleware import get_current_user_from_token, validate_agent_api_key, require_permission, check_user_permission
from utils.auth import authenticate_user
from config import MANAGEMENT_BASE_URL
router = APIRouter(prefix="/api/agents", tags=["agents"])
logger = logging.getLogger(__name__)
# ==== AGENT VERSION MANAGEMENT - HELPER FUNCTIONS ====
# Global version storage (acts as in-memory database)
AGENT_VERSIONS = {
"macos": "2.1.0", # Updated via endpoint
"linux": "2.1.0"
}
def _sanitize_agent_json(body_str: str):
"""Repair the common malformed-JSON patterns a hand-built agent heartbeat can emit.
Agents assemble their heartbeat JSON as text in bash, so an empty interpolated value can leave
a structurally-invalid comma (issue #31). Returns (possibly_repaired_str, was_changed). The
repairs are conservative and target only structural artifacts an agent produces; they never
alter this endpoint's legitimate string values (the agent emits no string containing ',,'
haproxy_stats_csv is base64/comma-free and the rest are constrained os/kernel/ip/version text).
"""
import re
sanitized = False
# Fix 1: empty value before a comma ("server_statuses": ,)
if re.search(r':\s*,', body_str):
body_str = re.sub(r':\s*,', ': null,', body_str); sanitized = True
# Fix 2: empty value before a closing brace ("field":})
if re.search(r':\s*}', body_str):
body_str = re.sub(r':\s*}', ': null}', body_str); sanitized = True
# Fix 3: trailing comma before } or ]
if re.search(r',(\s*[}\]])', body_str):
body_str = re.sub(r',(\s*[}\]])', r'\1', body_str); sanitized = True
# Fix 4: leading comma run right after an opening brace/bracket (issue #31): an empty
# $system_info as the first member collapses to '{ , "name": ...'. The ': ,' fix above cannot
# catch this because there is no key/colon before the comma.
if re.search(r'([{\[])(\s*,)+', body_str):
body_str = re.sub(r'([{\[])(\s*,)+', r'\1', body_str); sanitized = True
# Fix 5: a run of commas between members (issue #31): an empty $system_info between two fields
# produces '"version": "x",\n ,\n "haproxy_status": ...'. Runs after Fix 1/3 so only
# structural commas remain; collapse any comma run to a single comma.
if re.search(r',(\s*,)+', body_str):
body_str = re.sub(r',(\s*,)+', ',', body_str); sanitized = True
return body_str, sanitized
def get_platform_key(agent_platform: str) -> str:
"""Convert agent platform to standardized platform key - fixed empty platform fallback"""
platform = agent_platform.lower() if agent_platform else 'unknown'
if platform in ['darwin', 'macos', 'osx', 'mac']:
return 'macos'
elif platform in ['linux', 'ubuntu', 'debian', 'centos', 'rhel', 'fedora', 'alpine']:
return 'linux'
else:
# Default fallback for empty/unknown platforms - use linux as most common
return 'linux'
async def get_version_for_platform(platform: str) -> str:
"""Get latest version for specified platform from database"""
try:
conn = await get_database_connection()
platform_key = get_platform_key(platform)
# CRITICAL FIX: Order by updated_at DESC instead of created_at DESC
# This ensures "Reset to Default" updates are properly reflected
# because ON CONFLICT updates updated_at but not created_at
version_info = await conn.fetchrow("""
SELECT version, created_at, updated_at
FROM agent_versions
WHERE platform = $1 AND is_active = true
ORDER BY updated_at DESC, created_at DESC
LIMIT 1
""", platform_key)
await close_database_connection(conn)
if version_info:
logger.info(f"VERSION LOOKUP: platform={platform_key}, version={version_info['version']}, "
f"created_at={version_info['created_at']}, updated_at={version_info['updated_at']}")
return version_info['version']
else:
# Fallback to global storage if no database entry
logger.warning(f"VERSION LOOKUP: No active version in DB for platform={platform_key}, using fallback")
return AGENT_VERSIONS.get(platform_key, 'unknown')
except Exception as e:
logger.error(f"Error getting version for platform {platform}: {e}")
# Fallback to global storage
platform_key = get_platform_key(platform)
return AGENT_VERSIONS.get(platform_key, 'unknown')
# ==== FRONTEND SYNC HELPER ====
def _should_sync_frontend(frontend_name: str) -> bool:
"""
Determine if a frontend should be synced from agent to management system.
Returns False for system/manual frontends that should not be managed.
"""
# List of frontend names/patterns that should NOT be synced
skip_patterns = [
'stats', # HAProxy stats frontend
'haproxy-stats', # Alternative stats name
'admin', # Admin interfaces
'monitoring', # Monitoring frontends
'health', # Health check frontends
'prometheus', # Prometheus exporters
'socket', # Socket frontends
]
frontend_lower = frontend_name.lower()
# Check exact matches and patterns
for pattern in skip_patterns:
if pattern in frontend_lower:
return False
# Skip frontends that start with underscore (convention for system frontends)
if frontend_name.startswith('_'):
return False
return True
# ==== BACKEND SYNC HELPER ====
def _should_sync_backend(backend_name: str) -> bool:
"""
Determine if a backend should be synced from agent to management system.
Returns False for system/auto-managed backends that must not be persisted in DB.
Issue #11: '_acme_challenge_backend' is auto-injected by haproxy_config.py
on every Apply when ACME is enabled. Persisting it in DB caused duplicate
sections on next Apply (HAProxy "Duplicate Name" validation failure).
"""
# System/auto-managed backend names that must not be synced
skip_patterns = [
'_acme_challenge_backend', # Issue #11: auto-managed ACME challenge backend
]
if backend_name in skip_patterns:
return False
# Skip backends that start with underscore (system convention)
if backend_name.startswith('_'):
return False
return True
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
def calculate_agent_health(status, last_seen):
"""Calculate agent health status based on status and last_seen timestamp"""
if not last_seen:
return "unknown"
time_diff = time.time() - last_seen.timestamp()
if status == "online" and time_diff < 30:
return "healthy"
elif status == "online" and time_diff < 120:
return "warning"
else:
return "offline"
@router.get("", summary="Get All Agents", response_description="List of all agents")
async def get_agents(pool_id: Optional[int] = None, authorization: str = Header(None), x_api_key: Optional[str] = Header(None)):
"""
# Get All Agents
Retrieve list of all registered agents. Optionally filter by pool_id (cluster scope).
## Query Parameters
- **pool_id** (optional): Filter agents by pool ID
## Example Request - All Agents
```bash
curl -X GET "{BASE_URL}/api/agents" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..."
```
## Example Request - Filtered by Pool
```bash
curl -X GET "{BASE_URL}/api/agents?pool_id=1" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..."
```
## Example Response
```json
[
{
"id": 1,
"name": "production-agent-01",
"pool_id": 1,
"pool_name": "production-pool",
"pool_environment": "production",
"platform": "linux",
"architecture": "x86_64",
"version": "2.0.0",
"hostname": "haproxy-prod-01",
"ip_address": "10.0.1.10",
"operating_system": "Ubuntu 22.04",
"status": "online",
"enabled": true,
"haproxy_status": "running",
"haproxy_version": "2.8.3",
"last_seen": "2024-01-15T10:30:00Z",
"created_at": "2024-01-10T08:00:00Z"
}
]
```
## Response Fields
- **status**: Agent connection status (online/offline)
- **enabled**: Whether agent is enabled to receive tasks
- **haproxy_status**: Status of HAProxy service on agent's server
- **last_seen**: Last heartbeat timestamp
"""
# SECURITY (GHSA-3p5c-m5m4-mjpx): the agent inventory (names, hostnames, IPs,
# pools, OS) is operator data and was previously served unauthenticated — it is
# also the read-back channel used in the RCE exfil PoC. Require EITHER a valid
# operator JWT OR a valid agent X-API-Key: deployed agents poll this endpoint
# (with their key, not a JWT) to read their own applied_config_version and avoid
# re-applying config on restart, so a JWT-only gate would break them. Checked
# before the try so the 401 is not swallowed by the generic handler.
if authorization:
current_user = await get_current_user_from_token(authorization) # raises 401 on invalid JWT
else:
from auth_middleware import validate_agent_api_key
if not await validate_agent_api_key(x_api_key):
raise HTTPException(status_code=401, detail="Authentication required")
try:
conn = await get_database_connection()
try:
if pool_id:
agents = await conn.fetch("""
SELECT a.id, a.name, a.pool_id, a.platform, a.architecture, a.version,
a.hostname, a.ip_address, a.operating_system, a.kernel_version,
a.uptime, a.cpu_count, a.memory_total, a.disk_space,
a.network_interfaces, a.capabilities, a.status, a.last_seen, a.last_action_time,
a.haproxy_status, a.haproxy_version, a.config_version, a.applied_config_version,
a.keepalive_state, a.keepalive_ip,
COALESCE(a.enabled, TRUE) as enabled,
a.created_at, a.updated_at,
p.name as pool_name, p.environment as pool_environment
FROM agents a
LEFT JOIN haproxy_cluster_pools p ON a.pool_id = p.id
WHERE a.pool_id = $1 AND a.name NOT LIKE 'token_%'
ORDER BY a.name
""", pool_id)
else:
agents = await conn.fetch("""
SELECT a.id, a.name, a.pool_id, a.platform, a.architecture, a.version,
a.hostname, a.ip_address, a.operating_system, a.kernel_version,
a.uptime, a.cpu_count, a.memory_total, a.disk_space,
a.network_interfaces, a.capabilities, a.status, a.last_seen, a.last_action_time,
a.haproxy_status, a.haproxy_version, a.config_version, a.applied_config_version,
a.keepalive_state, a.keepalive_ip,
COALESCE(a.enabled, TRUE) as enabled,
a.created_at, a.updated_at,
p.name as pool_name, p.environment as pool_environment
FROM agents a
LEFT JOIN haproxy_cluster_pools p ON a.pool_id = p.id
WHERE a.name NOT LIKE 'token_%'
ORDER BY a.name
""")
except Exception as schema_error:
logger.warning(f"Schema error in agents query, using fallback: {schema_error}")
agents = await conn.fetch("""
SELECT a.id, a.name, a.status, a.last_seen,
COALESCE(a.enabled, TRUE) as enabled,
a.created_at, a.updated_at
FROM agents a
WHERE a.name NOT LIKE 'token_%'
ORDER BY a.name
""")
await close_database_connection(conn)
# Check if user has permission to view version info
user_can_view_versions = False
if authorization:
try:
current_user = await get_current_user_from_token(authorization)
user_can_view_versions = await check_user_permission(current_user["id"], "agents", "version")
except:
pass # User not authenticated or no permission, continue without version info
# Build agents list with platform-aware versions
agents_list = []
for agent in agents:
agent_dict = {
"id": agent["id"],
"name": agent["name"],
"pool_id": agent.get("pool_id"),
"pool_name": agent.get("pool_name"),
"platform": agent.get("platform", "unknown"),
"architecture": agent.get("architecture", "unknown"),
"version": agent.get("version", "unknown"),
"hostname": agent.get("hostname"),
"ip_address": str(agent["ip_address"]) if agent.get("ip_address") else None,
"operating_system": agent.get("operating_system"),
"kernel_version": agent.get("kernel_version"),
"uptime": agent.get("uptime"),
"cpu_count": agent.get("cpu_count"),
"memory_total": agent.get("memory_total"),
"disk_space": agent.get("disk_space"),
"network_interfaces": agent.get("network_interfaces", []),
"capabilities": agent.get("capabilities", []),
"status": agent.get("status", "unknown"),
"haproxy_status": agent.get("haproxy_status", "unknown"),
"haproxy_version": agent.get("haproxy_version", "unknown"),
"keepalive_state": agent.get("keepalive_state"),
"keepalive_ip": agent.get("keepalive_ip"),
"config_version": agent.get("config_version"),
"applied_config_version": agent.get("applied_config_version"),
"pool_environment": agent.get("pool_environment"),
"enabled": agent.get("enabled", True),
"last_seen": agent["last_seen"].isoformat().replace('+00:00', 'Z') if agent.get("last_seen") else None,
"last_action_time": agent["last_action_time"].isoformat().replace('+00:00', 'Z') if agent.get("last_action_time") else None,
"health": calculate_agent_health(agent.get("status"), agent.get("last_seen")),
"created_at": agent["created_at"].isoformat().replace('+00:00', 'Z') if agent.get("created_at") else None,
"updated_at": agent["updated_at"].isoformat().replace('+00:00', 'Z') if agent.get("updated_at") else None,
}
# Add version info only if user has permission
if user_can_view_versions:
# Get platform-specific latest version
latest_version = await get_version_for_platform(agent.get('platform', 'unknown'))
# Upgrade available if versions differ (regardless of direction)
current_ver = agent.get("version", "unknown")
upgrade_available = (current_ver != latest_version and
current_ver != "unknown" and
latest_version != "unknown")
agent_dict.update({
"current_version": current_ver, # Agent's current script version
"available_version": latest_version, # Platform-specific latest version
"upgrade_available": upgrade_available,
})
agents_list.append(agent_dict)
return {
"agents": agents_list
}
except Exception as e:
logger.error(f"Error fetching agents: {e}")
return {"agents": []}
@router.get("/platforms")
async def get_agent_platforms():
"""Get available agent platforms and architectures"""
return {
"platforms": {
"linux": {
"name": "Linux",
"architectures": {
"amd64": {"display": "x86-64 (amd64)"},
"arm64": {"display": "ARM64 (aarch64)"}
}
},
"darwin": {
"name": "macOS",
"architectures": {
"amd64": {"display": "Intel (amd64)"},
"arm64": {"display": "Apple Silicon (arm64)"}
}
}
}
}
@router.post("/generate-install-script", summary="Generate Agent Install Script", response_description="Installation script")
async def generate_install_script(req_data: AgentScriptRequest, request: Request, authorization: str = Header(None), x_api_key: Optional[str] = Header(None)):
"""
# Generate Agent Installation Script
Generate platform-specific installation script for deploying agent on HAProxy servers.
This is the primary way to create and install agents.
## Request Body
- **pool_id**: Agent pool ID (required) - Links agent to cluster
- **platform**: Platform type: "macos" or "linux" (required)
- **agent_name** (optional): Custom agent name (auto-generated if not provided)
## Example Request
```bash
curl -X POST "{BASE_URL}/api/agents/generate-install-script" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..." \\
-H "Content-Type: application/json" \\
-d '{
"pool_id": 1,
"platform": "linux",
"agent_name": "production-agent-01"
}'
```
## Example Response
```json
{
"script": "#!/bin/bash\\n\\n# HAProxy Agent Installation Script\\n# Generated: 2024-01-15...\\n\\nAGENT_NAME=\\"production-agent-01\\"\\nAPI_KEY=\\"agt_abc123xyz789...\\n",
"agent_name": "production-agent-01",
"api_key": "agt_abc123xyz789...",
"platform": "linux"
}
```
## Installation Steps
1. Save the returned script to a file (e.g., `install-agent.sh`)
2. Make it executable: `chmod +x install-agent.sh`
3. Run on target HAProxy server: `sudo ./install-agent.sh`
4. Agent will start and appear online in UI
## Agent Features
- Polls backend every 10-30 seconds for tasks
- Applies configuration changes to local haproxy.cfg
- Reloads HAProxy service automatically
- Reports metrics and status
- Self-updates when new versions available
"""
try:
# Support both user authentication and agent API key
current_user = None
if authorization:
from auth_middleware import get_current_user_from_token, check_user_permission
current_user = await get_current_user_from_token(authorization)
# SECURITY: Check permission for agent script generation
has_permission = await check_user_permission(current_user["id"], "agents", "create")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.create required"
)
elif x_api_key:
# Agent API key authentication for upgrade process
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
raise HTTPException(status_code=401, detail="Invalid agent API key")
else:
raise HTTPException(status_code=401, detail="Authorization header or X-API-Key required")
# Debug logging for agent upgrade requests
logger.info(f"SCRIPT GENERATION: Request from agent - platform: {req_data.platform}, cluster_id: {req_data.cluster_id}, pool_id: {req_data.pool_id}")
logger.info(f"SCRIPT GENERATION: agent_name: {req_data.agent_name}, hostname_prefix: {req_data.hostname_prefix}")
logger.info(f"SCRIPT GENERATION: haproxy_bin_path: {req_data.haproxy_bin_path}, haproxy_config_path: {req_data.haproxy_config_path}")
# Use the specific cluster_id sent from frontend instead of searching by pool_id
cluster_id = req_data.cluster_id
# Validate that the cluster exists and belongs to the specified pool
conn = await get_database_connection()
# Bulgu #82 (round-22 audit) — pre-fix any operator with
# `agents.create` permission could generate an install
# script for ANY cluster, regardless of pool/cluster
# scope. Skip the access check during agent
# auto-upgrade (no `current_user`; agent uses its own
# API key and already proved it owns this cluster's
# config-apply pipeline).
if current_user and cluster_id:
await validate_user_cluster_access(current_user['id'], cluster_id, conn)
# CRITICAL: Check if this is an agent upgrade FIRST
# Skip strict validation for upgrades - agent may have old/fallback config
is_agent_upgrade = (x_api_key is not None and authorization is None)
# For NEW agent creation, verify pool has at least one cluster
if not is_agent_upgrade:
pool_cluster_count = await conn.fetchval(
"SELECT COUNT(*) FROM haproxy_clusters WHERE pool_id = $1",
req_data.pool_id
)
if pool_cluster_count == 0:
await close_database_connection(conn)
raise HTTPException(
status_code=400,
detail=f"Pool {req_data.pool_id} has no clusters assigned. Please create a cluster for this pool or select a different pool."
)
cluster_result = await conn.fetchrow(
"SELECT id, pool_id FROM haproxy_clusters WHERE id = $1",
req_data.cluster_id
)
if not cluster_result:
await close_database_connection(conn)
raise HTTPException(
status_code=404,
detail=f"Cluster {req_data.cluster_id} not found"
)
if not is_agent_upgrade:
# NEW AGENT CREATION: Strict validation
# Cluster must have a pool_id assigned
if not cluster_result['pool_id']:
await close_database_connection(conn)
raise HTTPException(
status_code=400,
detail=f"Cluster {req_data.cluster_id} does not have a pool assigned. Please edit the cluster and assign a pool before creating agents."
)
# Validate pool_id matches
if cluster_result['pool_id'] != req_data.pool_id:
await close_database_connection(conn)
raise HTTPException(
status_code=400,
detail=f"Cluster {req_data.cluster_id} belongs to pool {cluster_result['pool_id']}, not pool {req_data.pool_id}. Please select the correct cluster."
)
else:
# AGENT UPGRADE: Allow upgrade even if pool_id mismatch (will be fixed by heartbeat)
logger.info(f"AGENT UPGRADE: Skipping pool validation for agent upgrade request (cluster_id={req_data.cluster_id})")
# If cluster has pool but agent sent different pool, log warning
if cluster_result['pool_id'] and cluster_result['pool_id'] != req_data.pool_id:
logger.warning(f"AGENT UPGRADE: Pool mismatch - cluster pool={cluster_result['pool_id']}, agent sent={req_data.pool_id}. Will be auto-corrected on next heartbeat.")
await close_database_connection(conn)
# Dynamic platform detection - more flexible approach
platform = req_data.platform.lower()
# macOS variants (Darwin-based systems)
macos_variants = ['macos', 'darwin', 'osx', 'mac']
# Linux variants
linux_variants = ['linux', 'ubuntu', 'debian', 'centos', 'rhel', 'fedora', 'alpine', 'amazon']
if any(variant in platform for variant in macos_variants):
template_filename = "macos_install.sh"
logger.info(f"PLATFORM: Detected macOS variant: {platform} -> using macos_install.sh")
elif any(variant in platform for variant in linux_variants):
template_filename = "linux_install.sh"
logger.info(f"PLATFORM: Detected Linux variant: {platform} -> using linux_install.sh")
else:
# Try to detect based on architecture or other hints
if req_data.architecture and 'arm' in req_data.architecture.lower():
# ARM architecture often indicates Apple Silicon (macOS)
template_filename = "macos_install.sh"
logger.info(f"PLATFORM: Unknown platform '{platform}' but ARM architecture detected -> using macos_install.sh")
else:
# Default to Linux for unknown platforms
template_filename = "linux_install.sh"
logger.info(f"PLATFORM: Unknown platform '{platform}' -> defaulting to linux_install.sh")
# Try to get script template from database first
conn = await get_database_connection()
# Get platform key for database lookup
platform_key = get_platform_key(req_data.platform)
# Get latest version for this platform (function is defined later in this file)
latest_version = await get_version_for_platform(req_data.platform)
# CRITICAL FIX: Prefer database template (UI edits), fallback to file
logger.info(f"TEMPLATE SEARCH: Looking for platform={platform_key}, version={latest_version}")
template = await conn.fetchrow("""
SELECT script_content, version, created_at, updated_at
FROM agent_script_templates
WHERE platform = $1 AND version = $2 AND is_active = true
LIMIT 1
""", platform_key, latest_version)
# Debug: Show what's actually in the database
all_templates = await conn.fetch("""
SELECT platform, version, is_active, created_at, updated_at
FROM agent_script_templates
WHERE platform = $1
ORDER BY updated_at DESC, created_at DESC
""", platform_key)
logger.info(f"ALL TEMPLATES for {platform_key}: {[(t['version'], t['is_active'], str(t['updated_at'])[:19]) for t in all_templates]}")
if template:
script_template = template['script_content']
script_length = len(script_template)
logger.info(f"SCRIPT SOURCE: Using database template for {platform_key} version {latest_version} "
f"(size: {script_length} bytes, updated_at: {template['updated_at']})")
else:
# Fallback to file-based template (initial setup or database empty)
script_template_path = os.path.join(os.path.dirname(__file__), '..', 'utils', 'agent_scripts', template_filename)
if not os.path.exists(script_template_path):
await close_database_connection(conn)
raise HTTPException(status_code=500, detail=f"Script template not found: {template_filename}")
with open(script_template_path, 'r') as f:
script_template = f.read()
logger.info(f"SCRIPT SOURCE: Using file template for {platform_key} (database empty, initial setup)")
# Sync file template to database for future use
try:
file_hash = hashlib.sha256(script_template.encode()).hexdigest()
await conn.execute("""
INSERT INTO agent_script_templates (platform, version, script_content, source_file_hash, is_active)
VALUES ($1, $2, $3, $4, true)
ON CONFLICT (platform, version) DO NOTHING
""", platform_key, latest_version, script_template, file_hash)
logger.info(f"SCRIPT SYNC: Synced file template to database for {platform_key} version {latest_version}")
except Exception as sync_error:
logger.warning(f"Could not sync template to database: {sync_error}")
await close_database_connection(conn)
# Debug each replacement value
logger.info(f"REPLACEMENTS: cluster_id={cluster_id}, agent_name={req_data.agent_name}")
logger.info(f"REPLACEMENTS: hostname_prefix={req_data.hostname_prefix}, haproxy_config_path={req_data.haproxy_config_path}")
logger.info(f"REPLACEMENTS: haproxy_bin_path={req_data.haproxy_bin_path}, stats_socket_path={req_data.stats_socket_path}")
logger.info(f"REPLACEMENTS: latest_version={latest_version}")
replacements = {
"{{MANAGEMENT_URL}}": MANAGEMENT_BASE_URL,
"{{CLUSTER_ID}}": str(cluster_id),
"{{AGENT_NAME}}": req_data.agent_name,
"{{HOSTNAME_PREFIX}}": req_data.hostname_prefix,
"{{HAPROXY_CONFIG_PATH}}": req_data.haproxy_config_path,
"{{HAPROXY_BIN_PATH}}": req_data.haproxy_bin_path,
"{{STATS_SOCKET_PATH}}": req_data.stats_socket_path,
"{{AGENT_VERSION}}": latest_version,
}
script_content = script_template
for key, value in replacements.items():
script_content = script_content.replace(key, value)
# Log activity (only if user authentication, not for agent API key)
if current_user:
await log_user_activity(
user_id=current_user["id"],
action='generate_script',
resource_type='agent',
resource_id=str(cluster_id),
details={'platform': req_data.platform, 'agent_name': req_data.agent_name, 'hostname_prefix': req_data.hostname_prefix}
)
return {
"script": script_content,
"platform": req_data.platform,
"cluster_id": cluster_id,
"filename": f"install-agent-{req_data.platform}.sh"
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Failed to generate install script: {e}")
raise HTTPException(status_code=500, detail=f"Failed to generate install script: {str(e)}")
@router.get("/generate-uninstall-script/{platform}", summary="Generate Agent Uninstall Script", response_description="Uninstall script")
async def generate_uninstall_script(platform: str, authorization: str = Header(None), x_api_key: Optional[str] = Header(None)):
"""
# Generate Agent Uninstall Script
Generate platform-specific uninstall script for completely removing agent from HAProxy servers.
This script removes ONLY agent-related files and services - HAProxy itself remains untouched.
## Path Parameters
- **platform**: Platform type: "macos" or "linux" (required)
## What Gets Removed
- Agent service (systemd/launchd)
- Agent binary (/usr/local/bin/haproxy-agent)
- Agent configuration (/etc/haproxy-agent)
- Agent logs (/var/log/haproxy-agent)
- Agent temp files and upgrade markers
- Agent backup files
## What Is NOT Touched
- HAProxy service
- HAProxy configuration (/etc/haproxy/haproxy.cfg)
- HAProxy logs
- Any other HAProxy-related files
## Example Request
```bash
curl -X GET "{BASE_URL}/api/agents/generate-uninstall-script/linux" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..."
```
## Usage on Remote Server
```bash
# Save the script
curl -o uninstall-agent.sh "{BASE_URL}/api/agents/generate-uninstall-script/linux"
chmod +x uninstall-agent.sh
sudo ./uninstall-agent.sh
```
"""
# SECURITY (GHSA-3p5c-m5m4-mjpx): require authentication (operator JWT or agent
# key), consistent with generate-install-script. The uninstall script itself is
# generic (no secrets/topology), but an agent-management endpoint should not be
# anonymously reachable. Checked before the try so the 401 is not swallowed.
if authorization:
await get_current_user_from_token(authorization)
elif not await validate_agent_api_key(x_api_key):
raise HTTPException(status_code=401, detail="Authentication required")
try:
# Normalize platform to a canonical key (always 'linux' or 'macos').
# macOS agents register with platform 'darwin' (from `uname -s`), so the
# strict ['linux','macos'] check used to 400 on the UI's uninstall flow
# (GET /generate-uninstall-script/darwin). Reuse the same get_platform_key()
# helper the install-script generator uses, so darwin/osx/mac and the
# linux distro variants all resolve correctly. Backward-compatible:
# 'linux'/'macos' still map to themselves.
platform_lower = get_platform_key(platform)
# Read uninstall script from agent_scripts directory (same as install scripts)
import os
script_filename = f"uninstall-agent-{platform_lower}.sh"
# Try multiple possible paths - prioritize agent_scripts directory
possible_paths = [
# Primary: Same directory as install scripts
os.path.join(os.path.dirname(__file__), "..", "utils", "agent_scripts", script_filename),
# Container path
f"/app/backend/utils/agent_scripts/{script_filename}",
f"/app/utils/agent_scripts/{script_filename}",
# Fallback: Legacy utils directory at project root
os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(__file__))), "utils", script_filename),
]
script_content = None
script_path_used = None
for path in possible_paths:
normalized_path = os.path.normpath(path)
if os.path.exists(normalized_path):
with open(normalized_path, 'r') as f:
script_content = f.read()
script_path_used = normalized_path
logger.info(f"Uninstall script found at: {normalized_path}")
break
if not script_content:
logger.error(f"Uninstall script not found for platform: {platform_lower}. Searched paths: {possible_paths}")
raise HTTPException(
status_code=404,
detail=f"Uninstall script not found for platform: {platform_lower}"
)
return {
"platform": platform_lower,
"script": script_content,
"filename": script_filename,
"usage": f"chmod +x {script_filename} && sudo ./{script_filename}"
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Failed to generate uninstall script: {e}")
raise HTTPException(status_code=500, detail=f"Failed to generate uninstall script: {str(e)}")
@router.delete("/{agent_id}", summary="Delete Agent", response_description="Agent deleted successfully")
async def delete_agent(agent_id: int, authorization: str = Header(None)):
"""
# Delete Agent
Delete an agent from the system. The agent service on the remote server will stop receiving tasks.
## Path Parameters
- **agent_id**: Agent ID to delete
## Important Notes
- Agent service on remote server should be uninstalled manually
- Configuration will no longer be sent to this agent
- Historical logs and activity will be preserved
## Example Request
```bash
curl -X DELETE "{BASE_URL}/api/agents/1" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..."
```
## Example Response
```json
{
"message": "Agent deleted successfully"
}
```
## Manual Uninstall on Remote Server
Linux:
```bash
sudo systemctl stop haproxy-agent
sudo systemctl disable haproxy-agent
sudo rm /etc/systemd/system/haproxy-agent.service
sudo rm -rf /opt/haproxy-agent
```
macOS:
```bash
sudo launchctl unload /Library/LaunchDaemons/com.haproxy.agent.plist
sudo rm /Library/LaunchDaemons/com.haproxy.agent.plist
sudo rm -rf /opt/haproxy-agent
```
## Error Responses
- **403**: Insufficient permissions
- **404**: Agent not found
- **500**: Server error
"""
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 agent delete
has_permission = await check_user_permission(current_user["id"], "agents", "delete")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.delete required"
)
conn = await get_database_connection()
# Get agent info including cluster relationship
agent = await conn.fetchrow("""
SELECT a.name, a.pool_id, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON a.pool_id = hc.pool_id
WHERE a.id = $1
""", agent_id)
if not agent:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail="Agent not found")
# Validate cluster access if agent belongs to a cluster
if agent['cluster_id']:
await validate_user_cluster_access(current_user['id'], agent['cluster_id'], conn)
# HA/VIP (Issue #27): block deleting a node that's still a member of an active VIP.
# Otherwise the CASCADE would silently drop it from the VIP (breaking the one-MASTER
# topology with no signal) and the still-running node would keep advertising the VIP
# with no way to be told to tear down (review MED-3). Make the operator remove it
# from the VIP first — that stages a clean PENDING change they can apply.
try:
vip_member = await conn.fetchrow(
"SELECT v.name FROM vip_members vm JOIN vip_instances v ON v.id = vm.vip_id "
"WHERE vm.agent_id = $1 AND v.is_active = TRUE LIMIT 1", agent_id)
except Exception: # noqa: BLE001 — vip_* may not exist on older schemas; don't block delete
vip_member = None
if vip_member:
await close_database_connection(conn)
raise HTTPException(
status_code=409,
detail=(f"This node is a member of VIP '{vip_member['name']}'. Remove it from the "
f"VIP on the HA / VIP page (and apply) before deleting the agent."))
await conn.execute("DELETE FROM agents WHERE id = $1", agent_id)
await close_database_connection(conn)
if current_user and current_user.get('id'):
await log_user_activity(
user_id=current_user['id'],
action='delete',
resource_type='agent',
resource_id=str(agent_id),
details={'agent_name': agent['name']}
)
return {"message": f"Agent '{agent['name']}' deleted successfully"}
except HTTPException:
raise
except Exception as e:
logger.error(f"Failed to delete agent: {e}")
raise HTTPException(status_code=500, detail=f"Failed to delete agent: {str(e)}")
def _safe_ip_for_inet(ip_str: Optional[str]) -> Optional[str]:
"""Validate IP string for PostgreSQL INET column. Returns None for invalid/empty."""
if not ip_str or not ip_str.strip():
return None
try:
ipaddress.ip_address(ip_str.strip())
return ip_str.strip()
except ValueError:
return None
def _extract_agent_ip(heartbeat_data: AgentHeartbeat) -> Optional[str]:
"""Extract validated IP from heartbeat - top-level first, then system_info fallback."""
ip = _safe_ip_for_inet(heartbeat_data.ip_address)
if ip:
return ip
if heartbeat_data.system_info and isinstance(heartbeat_data.system_info, dict):
ip = _safe_ip_for_inet(heartbeat_data.system_info.get("ip_address"))
if ip:
return ip
return None
@router.post("/{agent_id}/heartbeat")
async def agent_heartbeat(agent_id: int, heartbeat_data: AgentHeartbeat, x_api_key: Optional[str] = Header(None)):
"""Receive agent heartbeat and update status."""
# Agent authentication is MANDATORY (GHSA-3p5c-m5m4-mjpx). This legacy by-ID
# heartbeat previously had NO auth, allowing unauthenticated state spoofing of
# any agent row. Deployed agents use the by-name heartbeat; a valid global
# agent token is now required here too. NOTE: raised BEFORE the try below so
# the 401 is not swallowed by the generic `except Exception` handler.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(f"Missing/invalid API key on by-id heartbeat for agent ID {agent_id}")
raise HTTPException(status_code=401, detail="Authentication required")
try:
conn = await get_database_connection()
await conn.execute("""
UPDATE agents
SET status = 'online',
last_seen = CURRENT_TIMESTAMP,
hostname = COALESCE($2, hostname),
haproxy_status = COALESCE($3, haproxy_status),
haproxy_version = COALESCE($4, haproxy_version),
keepalive_state = CASE WHEN $5::text IS NOT NULL THEN NULLIF($5::text, '') ELSE keepalive_state END,
keepalive_ip = CASE WHEN $6::text IS NOT NULL THEN NULLIF($6::text, '') ELSE keepalive_ip END,
ip_address = COALESCE($7::inet, ip_address),
updated_at = CURRENT_TIMESTAMP
WHERE id = $1
""", agent_id, heartbeat_data.hostname, heartbeat_data.haproxy_status, heartbeat_data.haproxy_version,
heartbeat_data.keepalive_state, heartbeat_data.keepalive_ip, _extract_agent_ip(heartbeat_data))
await close_database_connection(conn)
return {"status": "ok", "message": "Heartbeat received"}
except Exception as e:
logger.error(f"Heartbeat processing failed for agent ID {agent_id}: {e}", exc_info=True)
raise HTTPException(status_code=500, detail=f"Internal server error during heartbeat processing for agent ID {agent_id}")
@router.post("/{agent_name}/config-applied")
async def agent_config_applied_notification(agent_name: str, notification_data: dict, x_api_key: Optional[str] = Header(None)):
"""Receive instant notification when agent applies configuration - for real-time UI sync"""
try:
# Bulgu #75 (round-22 audit) — pre-fix the guard read:
#
# if x_api_key and not agent_auth:
# raise HTTPException(401, "Invalid API key")
#
# which accepted requests with NO `x_api_key` header at
# all (the `and` short-circuits). An unauthenticated
# attacker could therefore POST fabricated
# `config-applied`, `config-validation-failed`,
# `config-sync` and `upgrade-complete` notifications,
# poisoning the control-plane's view of agent state and —
# for `config-sync` — overwriting entity rows in the DB
# to match attacker-supplied HAProxy fragments. The new
# guard requires a present-AND-valid agent API key.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(
f"Rejected config-applied call for agent {agent_name!r}: "
f"missing or invalid x-api-key"
)
raise HTTPException(status_code=401, detail="Invalid or missing API key")
# GLOBAL TOKEN: Token can be used by multiple agents across different pools/clusters
conn = await get_database_connection()
# Update agent's applied config version for real-time sync
if notification_data.get("status") == "applied" and notification_data.get("version"):
version = notification_data["version"]
# CRITICAL: Update agent's applied_config_version for Agent Sync calculation
# Config versions are already APPLIED immediately after apply-changes
# Agent Sync is tracked separately via applied_config_version vs latest_config_version
logger.info(f"DEBUG: Updating applied_config_version for agent '{agent_name}' to '{version}'")
try:
# Use explicit transaction for asyncpg
async with conn.transaction():
# First check if agent exists (moved inside transaction for atomic operation)
agent_exists = await conn.fetchrow("SELECT id, name FROM agents WHERE name = $1", agent_name)
logger.info(f"DEBUG: Agent exists check: {agent_exists}")
if not agent_exists:
logger.error(f"DEBUG: Agent '{agent_name}' not found in database")
await close_database_connection(conn)
return {"status": "error", "message": f"Agent {agent_name} not found"}
# SIMPLIFIED: Direct UPDATE without column check (column should exist from migration)
# ALSO: Clear last_validation_error when config is successfully applied
result = await conn.execute("""
UPDATE agents SET
applied_config_version = $1,
last_validation_error = NULL,
last_validation_error_at = NULL,
updated_at = CURRENT_TIMESTAMP
WHERE name = $2
""", version, agent_name)
logger.info(f"DEBUG: Database UPDATE result: {result}")
# CRITICAL UX FIX: Clear validation_error on the config_version when successfully applied
# This ensures UI doesn't show stale validation errors
try:
agent_cluster = await conn.fetchrow("""
SELECT hc.id as cluster_id FROM agents a
JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
""", agent_name)
if agent_cluster:
await conn.execute("""
UPDATE config_versions
SET validation_error = NULL, validation_error_reported_at = NULL
WHERE cluster_id = $1 AND version_name = $2
""", agent_cluster['cluster_id'], version)
logger.info(f"CONFIG-APPLIED: Cleared validation error on version '{version}' after successful apply")
except Exception as clear_err:
logger.warning(f"CONFIG-APPLIED: Could not clear validation error: {clear_err}")
# Verify the update was successful within the same transaction
updated_agent = await conn.fetchrow("SELECT applied_config_version FROM agents WHERE name = $1", agent_name)
logger.info(f"DEBUG: Agent '{agent_name}' applied_config_version after update: {updated_agent['applied_config_version'] if updated_agent else 'NOT_FOUND'}")
logger.info(f"DEBUG: Transaction committed successfully for agent '{agent_name}'")
logger.info(f"INSTANT SYNC: Agent '{agent_name}' applied config version '{version}' - UI sync updated")
except Exception as transaction_error:
logger.error(f"TRANSACTION ERROR: Failed to update applied_config_version for agent '{agent_name}': {transaction_error}", exc_info=True)
await close_database_connection(conn)
return {"status": "error", "message": f"Transaction failed: {str(transaction_error)}"}
# Log activity
await log_agent_activity(agent_name, "config_applied", {"version": version})
else:
logger.warning(f"CONFIG-APPLIED: Invalid notification format - status: {notification_data.get('status')}, version: {notification_data.get('version')}")
await close_database_connection(conn)
return {"status": "ok", "message": "Config applied notification received"}
except HTTPException:
raise # let auth 401/403 propagate (do not turn it into a 200 error body)
except Exception as e:
logger.error(f"Failed to process config applied notification from agent '{agent_name}': {e}")
return {"status": "error", "message": str(e)}
@router.post("/{agent_name}/config-validation-failed")
async def agent_config_validation_failed(agent_name: str, notification_data: dict, x_api_key: Optional[str] = Header(None)):
"""Receive notification when agent's HAProxy config validation fails - for UI error display"""
try:
# Bulgu #75 (round-22 audit) — see config-applied above.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(
f"Rejected validation-failed call for agent {agent_name!r}: "
f"missing or invalid x-api-key"
)
raise HTTPException(status_code=401, detail="Invalid or missing API key")
# GLOBAL TOKEN: Token can be used by multiple agents across different pools/clusters
conn = await get_database_connection()
# Get agent and cluster info
agent_info = await conn.fetchrow("""
SELECT a.id, a.name, a.pool_id, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
""", agent_name)
if not agent_info or not agent_info['cluster_id']:
await close_database_connection(conn)
logger.error(f"Agent '{agent_name}' or cluster not found for validation error notification")
return {"status": "error", "message": "Agent or cluster not found"}
cluster_id = agent_info['cluster_id']
version = notification_data.get("version", "unknown")
validation_error = notification_data.get("validation_error", "Unknown validation error")
logger.error(f"VALIDATION FAILED: Agent '{agent_name}' (cluster {cluster_id}) - HAProxy config validation failed")
logger.error(f"VALIDATION ERROR: {validation_error}")
# Store validation error in config_versions table for UI display
async with conn.transaction():
# CRITICAL DEBUG: Log exact version being searched
logger.info(f"VALIDATION-FAILED: Searching for version_name='{version}' in cluster_id={cluster_id}")
# First, check what versions exist for this cluster (for debugging)
existing_versions = await conn.fetch("""
SELECT version_name, status, is_active FROM config_versions
WHERE cluster_id = $1 AND is_active = TRUE
ORDER BY created_at DESC LIMIT 3
""", cluster_id)
logger.info(f"VALIDATION-FAILED: Active versions in cluster: {[v['version_name'] for v in existing_versions]}")
# Update the config_version with validation error
result = await conn.execute("""
UPDATE config_versions
SET validation_error = $1,
validation_error_reported_at = CURRENT_TIMESTAMP,
updated_at = CURRENT_TIMESTAMP
WHERE cluster_id = $2 AND version_name = $3
""", validation_error, cluster_id, version)
# CRITICAL: Check if UPDATE actually affected any rows
rows_affected = int(result.split()[-1]) if result else 0
if rows_affected == 0:
logger.error(f"VALIDATION-FAILED: UPDATE affected 0 rows! version_name '{version}' NOT FOUND in cluster {cluster_id}")
# Try to update the active version instead as fallback
fallback_result = await conn.execute("""
UPDATE config_versions
SET validation_error = $1,
validation_error_reported_at = CURRENT_TIMESTAMP,
updated_at = CURRENT_TIMESTAMP
WHERE cluster_id = $2 AND is_active = TRUE
""", validation_error, cluster_id)
fallback_rows = int(fallback_result.split()[-1]) if fallback_result else 0
logger.info(f"VALIDATION-FAILED: Fallback UPDATE to active version affected {fallback_rows} rows")
else:
logger.info(f"VALIDATION-FAILED: UPDATE affected {rows_affected} rows for version '{version}'")
# Also update agent status to indicate validation failure
await conn.execute("""
UPDATE agents
SET last_validation_error = $1,
last_validation_error_at = CURRENT_TIMESTAMP,
updated_at = CURRENT_TIMESTAMP
WHERE id = $2
""", validation_error, agent_info['id'])
logger.info(f"VALIDATION ERROR STORED: Agent '{agent_name}', version '{version}', cluster {cluster_id}")
await close_database_connection(conn)
return {"status": "ok", "message": "Validation error notification received"}
except HTTPException:
raise # let auth 401/403 propagate (do not turn it into a 200 error body)
except Exception as e:
logger.error(f"Failed to process validation error notification from agent '{agent_name}': {e}")
return {"status": "error", "message": str(e)}
@router.post("/{agent_name}/config-sync")
async def agent_config_sync(agent_name: str, sync_data: dict, x_api_key: Optional[str] = Header(None)):
"""Receive agent's current config content and sync database entities accordingly"""
try:
# Bulgu #75 (round-22 audit) — see config-applied above.
# config-sync is the highest-impact agent webhook: it
# mutates the entity rows in the DB based on the agent's
# reported HAProxy fragments. Accepting no-API-key
# requests here let an attacker rewrite arbitrary rows.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(
f"Rejected config-sync call for agent {agent_name!r}: "
f"missing or invalid x-api-key"
)
raise HTTPException(status_code=401, detail="Invalid or missing API key")
# GLOBAL TOKEN: Token can be used by multiple agents across different pools/clusters
conn = await get_database_connection()
# Get agent and cluster info
agent_info = await conn.fetchrow("""
SELECT a.id, a.name, a.pool_id, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
""", agent_name)
if not agent_info or not agent_info['cluster_id']:
await close_database_connection(conn)
return {"status": "error", "message": "Agent or cluster not found"}
cluster_id = agent_info['cluster_id']
agent_id = agent_info['id']
config_content = sync_data.get('config_content', '')
if not config_content:
await close_database_connection(conn)
return {"status": "error", "message": "No config content provided"}
# ========================================================================
# COLLISION PREVENTION: Extract and save listen block names
# Agent's local listen blocks can conflict with frontend/backend names
# By storing them, we can detect collisions BEFORE entity creation
# ========================================================================
preserved_listen_blocks = []
for line in config_content.split('\n'):
line_stripped = line.strip()
# Parse listen block definitions (e.g., "listen stats")
if line_stripped.startswith('listen '):
listen_name = line_stripped.split()[1] if len(line_stripped.split()) > 1 else None
if listen_name:
preserved_listen_blocks.append(listen_name)
logger.debug(f"CONFIG-SYNC: Found listen block '{listen_name}' in agent {agent_name}")
# Save preserved listen blocks to agents table for collision detection
if preserved_listen_blocks:
try:
import json
await conn.execute("""
UPDATE agents
SET preserved_listen_blocks = $1::jsonb
WHERE id = $2
""", json.dumps(preserved_listen_blocks), agent_id)
logger.info(f"CONFIG-SYNC: Agent '{agent_name}' has {len(preserved_listen_blocks)} listen blocks: {preserved_listen_blocks}")
except Exception as e:
logger.warning(f"CONFIG-SYNC: Failed to save listen blocks for agent {agent_name}: {e}")
else:
# Clear listen blocks if agent has none
try:
await conn.execute("""
UPDATE agents
SET preserved_listen_blocks = '[]'::jsonb
WHERE id = $1
""", agent_id)
except:
pass
# ========================================================================
# END COLLISION PREVENTION
# ========================================================================
# Parse config content to extract all configuration entities
active_servers = []
active_backends = []
active_frontends = []
current_backend = None
current_frontend = None
for line in config_content.split('\n'):
line = line.strip()
# Parse backend definitions
if line.startswith('backend '):
backend_name = line.split(' ', 1)[1]
current_frontend = None
# Issue #11: Skip system/auto-managed backends (e.g. _acme_challenge_backend).
# Setting current_backend = None ensures subsequent `server` lines under
# a skipped backend are NOT attached to a previous user backend (orphan rows).
if not _should_sync_backend(backend_name):
logger.info(f"🚫 AGENT SYNC: Skipping system/auto-managed backend '{backend_name}' from sync")
current_backend = None
continue
current_backend = backend_name
active_backends.append({
'name': backend_name,
'is_commented': False
})
continue
elif line.startswith('# DISABLED: backend '):
backend_name = line[20:].strip() # Remove "# DISABLED: backend "
# Skip system/auto-managed backends even if disabled
if not _should_sync_backend(backend_name):
continue
active_backends.append({
'name': backend_name,
'is_commented': True
})
continue
# Parse frontend definitions - Filter out stats and manual frontends
if line.startswith('frontend '):
frontend_name = line.split(' ', 1)[1]
current_frontend = frontend_name
current_backend = None
# Skip stats frontend and other system frontends
if not _should_sync_frontend(frontend_name):
logger.info(f"🚫 AGENT SYNC: Skipping system/manual frontend '{frontend_name}' from sync")
continue
active_frontends.append({
'name': frontend_name,
'is_commented': False
})
continue
elif line.startswith('# DISABLED: frontend '):
frontend_name = line[21:].strip() # Remove "# DISABLED: frontend "
# Skip stats frontend and other system frontends even if disabled
if not _should_sync_frontend(frontend_name):
continue
active_frontends.append({
'name': frontend_name,
'is_commented': True
})
continue
# Parse server lines (existing logic)
if current_backend and line.startswith('server '):
parts = line.split()
if len(parts) >= 3:
server_name = parts[1]
server_address = parts[2]
active_servers.append({
'backend_name': current_backend,
'server_name': server_name,
'server_address': server_address,
'is_commented': False
})
elif current_backend and line.startswith('# DISABLED: server '):
# Parse commented out servers
parts = line[12:].split() # Remove "# DISABLED: "
if len(parts) >= 3:
server_name = parts[1]
server_address = parts[2]
active_servers.append({
'backend_name': current_backend,
'server_name': server_name,
'server_address': server_address,
'is_commented': True
})
# Sync database with agent's actual config
async with conn.transaction():
# 1. SYNC BACKENDS
# Mark all backends in this cluster as potentially inactive
await conn.execute("""
UPDATE backends
SET is_active = FALSE, updated_at = CURRENT_TIMESTAMP
WHERE cluster_id = $1
""", cluster_id)
# Update backends that exist in agent's config
for backend_info in active_backends:
is_active = not backend_info['is_commented']
# Try to update existing backend
result = await conn.execute("""
UPDATE backends
SET is_active = $1, last_config_status = 'APPLIED', updated_at = CURRENT_TIMESTAMP
WHERE name = $2 AND cluster_id = $3
""", is_active, backend_info['name'], cluster_id)
# If backend doesn't exist, create it (this handles restore scenarios)
if result == "UPDATE 0":
await conn.execute("""
INSERT INTO backends
(name, balance_method, mode, cluster_id, is_active)
VALUES ($1, 'roundrobin', 'http', $2, $3)
ON CONFLICT (name, cluster_id) DO UPDATE SET
is_active = EXCLUDED.is_active,
updated_at = CURRENT_TIMESTAMP
""", backend_info['name'], cluster_id, is_active)
# 2. SYNC FRONTENDS
# Mark all frontends in this cluster as potentially inactive
await conn.execute("""
UPDATE frontends
SET is_active = FALSE, updated_at = CURRENT_TIMESTAMP
WHERE cluster_id = $1
""", cluster_id)
# Update frontends that exist in agent's config
for frontend_info in active_frontends:
is_active = not frontend_info['is_commented']
# CRITICAL FIX: Only update status fields, preserve all other configuration
# including SSL settings, ports, backends, etc. that are managed via UI
result = await conn.execute("""
UPDATE frontends
SET is_active = $1, last_config_status = 'APPLIED', updated_at = CURRENT_TIMESTAMP
WHERE name = $2 AND cluster_id = $3
""", is_active, frontend_info['name'], cluster_id)
# If frontend doesn't exist, create it (this handles restore scenarios)
if result == "UPDATE 0":
await conn.execute("""
INSERT INTO frontends
(name, bind_address, bind_port, mode, cluster_id, is_active)
VALUES ($1, '*', 80, 'http', $2, $3)
ON CONFLICT (name, cluster_id) DO UPDATE SET
is_active = EXCLUDED.is_active,
updated_at = CURRENT_TIMESTAMP
""", frontend_info['name'], cluster_id, is_active)
# 3. SYNC SERVERS (existing logic)
# Mark all servers in this cluster as potentially inactive
await conn.execute("""
UPDATE backend_servers
SET is_active = FALSE, updated_at = CURRENT_TIMESTAMP
WHERE cluster_id = $1
""", cluster_id)
# Update servers that exist in agent's config
for server_info in active_servers:
is_active = not server_info['is_commented']
# Try to update existing server
result = await conn.execute("""
UPDATE backend_servers
SET is_active = $1, last_config_status = 'APPLIED', updated_at = CURRENT_TIMESTAMP
WHERE backend_name = $2 AND server_name = $3 AND cluster_id = $4
""", is_active, server_info['backend_name'], server_info['server_name'], cluster_id)
# If server doesn't exist, create it ONLY if it's not commented (deleted)
# Commented servers that don't exist in DB were intentionally hard deleted
if result == "UPDATE 0" and not server_info['is_commented']:
# Extract IP and port from server_address
address_parts = server_info['server_address'].split(':')
server_ip = address_parts[0]
server_port = int(address_parts[1]) if len(address_parts) > 1 else 80
# Insert new server without ON CONFLICT since constraint doesn't exist
await conn.execute("""
INSERT INTO backend_servers
(backend_name, server_name, server_address, server_port, weight, is_active, cluster_id)
VALUES ($1, $2, $3, $4, 100, $5, $6)
""", server_info['backend_name'], server_info['server_name'],
server_ip, server_port, is_active, cluster_id)
elif result == "UPDATE 0" and server_info['is_commented']:
logger.info(f"🗑️ AGENT SYNC: Skipping creation of commented (deleted) server {server_info['server_name']} in backend {server_info['backend_name']}")
await close_database_connection(conn)
logger.info(f"CONFIG SYNC: Agent '{agent_name}' synced {len(active_backends)} backends, {len(active_frontends)} frontends, {len(active_servers)} servers with database")
return {"status": "ok", "message": f"Config synced - {len(active_backends)} backends, {len(active_frontends)} frontends, {len(active_servers)} servers processed"}
except HTTPException:
raise # let auth 401/403 propagate (do not turn it into a 200 error body)
except Exception as e:
logger.error(f"Failed to process config sync from agent '{agent_name}': {e}")
return {"status": "error", "message": str(e)}
@router.post("/heartbeat")
async def agent_heartbeat_by_name(
request: Request,
x_api_key: Optional[str] = Header(None)
):
"""Receive agent heartbeat by agent name, with auto-registration and data validation.
Enhanced with robust JSON sanitization to handle malformed agent payloads.
"""
import re
import json
from pydantic import ValidationError
# Read raw body. Parse VALID JSON as-is (the normal case for every agent version) and only
# fall back to the malformed-JSON repair when the body does not parse. This guarantees a healthy
# heartbeat from any agent version is byte-for-byte untouched — the repair regexes can never run
# against a well-formed payload (issue #31; strictly safer than repairing unconditionally).
try:
raw_body = await request.body()
body_str = raw_body.decode('utf-8')
try:
heartbeat_dict = json.loads(body_str)
except json.JSONDecodeError:
# Malformed body (would otherwise be a hard 400). Attempt a conservative repair of the
# comma artifacts a hand-built agent heartbeat can emit, then re-parse.
repaired, changed = _sanitize_agent_json(body_str)
if changed:
agent_name = "unknown"
try:
name_match = re.search(r'"name"\s*:\s*"([^"]+)"', repaired)
if name_match:
agent_name = name_match.group(1)
except Exception:
pass
logger.info(f"Repaired malformed JSON from agent '{agent_name}' before parsing")
logger.debug(f"Original JSON (preview): {body_str[:300]}")
logger.debug(f"Repaired JSON (preview): {repaired[:300]}")
heartbeat_dict = json.loads(repaired) # may still raise -> handled as 400 below
# DEBUG: Log cluster_id for auto-register troubleshooting
if heartbeat_dict.get('name'):
logger.info(f"HEARTBEAT DEBUG: agent={heartbeat_dict.get('name')}, cluster_id={heartbeat_dict.get('cluster_id')}, has_cluster_id={bool(heartbeat_dict.get('cluster_id'))}")
heartbeat_data = AgentHeartbeat(**heartbeat_dict)
except json.JSONDecodeError as e:
logger.error(f"JSON decode error even after sanitization: {e}")
raise HTTPException(status_code=400, detail=f"Invalid JSON: {str(e)}")
except ValidationError as e:
logger.error(f"Pydantic validation error after sanitization: {e}")
raise HTTPException(status_code=422, detail=f"Validation error: {str(e)}")
except Exception as e:
logger.error(f"Unexpected error processing heartbeat: {e}")
raise HTTPException(status_code=500, detail="Internal server error")
# Agent authentication is MANDATORY (GHSA-3p5c-m5m4-mjpx). A valid global agent
# token is required to heartbeat OR auto-register. Deployed agents always send
# X-API-Key; an absent/invalid key is an unauthenticated caller. This is done
# OUTSIDE the processing try below (whose generic `except Exception` would
# otherwise convert the 401 into a 500), and before opening a DB connection
# (validate_agent_api_key(None) needs no DB). Closes keyless heartbeat spoofing
# and keyless rogue-agent auto-registration (the `elif not agent` keyless path
# below is now unreachable, since agent_auth is guaranteed truthy past here).
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(f"Missing/invalid API key on heartbeat for agent '{heartbeat_data.name}'")
raise HTTPException(status_code=401, detail="Authentication required")
# Continue with normal heartbeat processing
try:
conn = await get_database_connection()
agent_name = heartbeat_data.name
agent = await conn.fetchrow("SELECT id, pool_id, api_key FROM agents WHERE name = $1", agent_name)
# If agent exists and is using a different API key, update the token association
# This handles: 1) agent switching to different token, 2) agent getting token for first time
current_agent_api_key = agent['api_key'] if agent else None
if agent and x_api_key and agent_auth and current_agent_api_key != x_api_key:
# Get token info from the new API key (from placeholder or another agent)
token_info = await conn.fetchrow("""
SELECT api_key_name, api_key_created_at, api_key_expires_at, api_key_created_by
FROM agents WHERE api_key = $1 LIMIT 1
""", x_api_key)
if token_info:
# Update agent's token association
await conn.execute("""
UPDATE agents SET
api_key = $1,
api_key_name = $2,
api_key_created_at = $3,
api_key_expires_at = $4,
api_key_created_by = $5,
updated_at = CURRENT_TIMESTAMP
WHERE id = $6
""", x_api_key, token_info['api_key_name'], token_info['api_key_created_at'],
token_info['api_key_expires_at'], token_info['api_key_created_by'], agent['id'])
logger.info(f"TOKEN UPDATE: Agent '{agent_name}' token updated to '{token_info['api_key_name']}' (previously used different token)")
else:
# Token not found in system, just update the api_key
await conn.execute("""
UPDATE agents SET
api_key = $1,
updated_at = CURRENT_TIMESTAMP
WHERE id = $2
""", x_api_key, agent['id'])
logger.info(f"TOKEN UPDATE: Agent '{agent_name}' api_key updated (token info not found)")
if not agent and x_api_key and agent_auth:
# Check if there's a placeholder agent with this API key
# Check if API key already used by another agent
existing_agent = await conn.fetchrow("SELECT id, name, pool_id FROM agents WHERE api_key = $1 AND name = $2", x_api_key, agent_name)
if existing_agent:
# Agent already exists with this API key
agent = {'id': existing_agent['id'], 'pool_id': existing_agent['pool_id']}
else:
# Check if there's a placeholder agent with this API key
placeholder_agent = await conn.fetchrow("SELECT id, name, pool_id FROM agents WHERE api_key = $1 AND name LIKE 'token_%'", x_api_key)
if placeholder_agent:
# Create new agent instead of updating placeholder (allows multiple agents per token)
logger.info(f"Creating new agent '{agent_name}' using token '{placeholder_agent['name']}'")
# Get pool_id from cluster_id if provided in heartbeat
pool_id_from_cluster = None
if heartbeat_dict.get('cluster_id'):
logger.info(f"AUTO-REGISTER DEBUG: Querying cluster_id {heartbeat_dict['cluster_id']} for pool_id")
cluster_pool = await conn.fetchrow("""
SELECT pool_id FROM haproxy_clusters WHERE id = $1
""", heartbeat_dict['cluster_id'])
logger.info(f"AUTO-REGISTER DEBUG: Cluster query result: {cluster_pool}")
if cluster_pool:
pool_id_from_cluster = cluster_pool['pool_id']
logger.info(f"Auto-register: Got pool_id {pool_id_from_cluster} from cluster_id {heartbeat_dict['cluster_id']}")
else:
logger.warning(f"AUTO-REGISTER DEBUG: No cluster_id in heartbeat for agent '{agent_name}'")
agent_id = await conn.fetchval("""
INSERT INTO agents (
name, api_key, api_key_name, api_key_created_at,
api_key_expires_at, api_key_created_by, status,
enabled, created_at, updated_at, pool_id
)
SELECT $1, api_key, api_key_name, api_key_created_at,
api_key_expires_at, api_key_created_by, 'online',
TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, $3
FROM agents WHERE id = $2
RETURNING id
""", agent_name, placeholder_agent['id'], pool_id_from_cluster)
logger.info(f"Auto-registered new agent '{agent_name}' with ID {agent_id} and pool_id {pool_id_from_cluster}")
agent = {'id': agent_id, 'pool_id': pool_id_from_cluster}
else:
# Auto-register new agent
logger.info(f"Agent '{agent_name}' not found. Auto-registering...")
logger.info(f"AUTO-REGISTER DEBUG: heartbeat_dict has cluster_id={heartbeat_dict.get('cluster_id')}")
try:
# Get pool_id from cluster_id if provided in heartbeat
pool_id_from_cluster = None
if heartbeat_dict.get('cluster_id'):
logger.info(f"AUTO-REGISTER: Querying cluster_id {heartbeat_dict['cluster_id']} for pool_id")
cluster_pool = await conn.fetchrow("""
SELECT pool_id FROM haproxy_clusters WHERE id = $1
""", heartbeat_dict['cluster_id'])
logger.info(f"AUTO-REGISTER: Cluster query result: {cluster_pool}")
if cluster_pool:
pool_id_from_cluster = cluster_pool['pool_id']
logger.info(f"Auto-register: Got pool_id {pool_id_from_cluster} from cluster_id {heartbeat_dict['cluster_id']}")
else:
logger.warning(f"AUTO-REGISTER: No cluster_id in heartbeat for agent '{agent_name}'")
agent_id = await conn.fetchval("""
INSERT INTO agents (name, status, last_seen, api_key, pool_id)
VALUES ($1, 'online', CURRENT_TIMESTAMP, $2, $3) RETURNING id
""", agent_name, x_api_key, pool_id_from_cluster)
logger.info(f"Auto-registered new agent '{agent_name}' with ID {agent_id} and pool_id {pool_id_from_cluster}")
agent = {'id': agent_id, 'pool_id': pool_id_from_cluster}
except Exception as e:
await close_database_connection(conn)
logger.error(f"Failed to auto-register agent '{agent_name}': {e}", exc_info=True)
raise HTTPException(status_code=500, detail=f"Could not create agent: {agent_name}")
elif not agent:
# Auto-register new agent without API key
logger.info(f"Agent '{agent_name}' not found. Auto-registering (no API key path)...")
logger.info(f"AUTO-REGISTER DEBUG (no key): heartbeat_dict has cluster_id={heartbeat_dict.get('cluster_id')}")
try:
# Get pool_id from cluster_id if provided in heartbeat
pool_id_from_cluster = None
if heartbeat_dict.get('cluster_id'):
logger.info(f"AUTO-REGISTER (no key): Querying cluster_id {heartbeat_dict['cluster_id']} for pool_id")
cluster_pool = await conn.fetchrow("""
SELECT pool_id FROM haproxy_clusters WHERE id = $1
""", heartbeat_dict['cluster_id'])
logger.info(f"AUTO-REGISTER (no key): Cluster query result: {cluster_pool}")
if cluster_pool:
pool_id_from_cluster = cluster_pool['pool_id']
logger.info(f"Auto-register: Got pool_id {pool_id_from_cluster} from cluster_id {heartbeat_dict['cluster_id']}")
else:
logger.warning(f"AUTO-REGISTER (no key): No cluster_id in heartbeat for agent '{agent_name}'")
agent_id = await conn.fetchval("""
INSERT INTO agents (name, status, last_seen, pool_id)
VALUES ($1, 'online', CURRENT_TIMESTAMP, $2) RETURNING id
""", agent_name, pool_id_from_cluster)
logger.info(f"Auto-registered new agent '{agent_name}' with ID {agent_id} and pool_id {pool_id_from_cluster}")
agent = {'id': agent_id, 'pool_id': pool_id_from_cluster}
except Exception as e:
await close_database_connection(conn)
logger.error(f"Failed to auto-register agent '{agent_name}': {e}", exc_info=True)
raise HTTPException(status_code=500, detail=f"Could not create agent: {agent_name}")
agent_id = agent['id']
# Update API key last used timestamp for all agents using this token
if x_api_key:
await conn.execute("""
UPDATE agents
SET api_key_last_used = CURRENT_TIMESTAMP
WHERE api_key = $1
""", x_api_key)
# Check if agent was in upgrading status and version has changed
# (v1.8.6: one round-trip instead of three; row is None exactly when the
# per-column fetchvals would each have returned None)
current_agent_row = await conn.fetchrow(
"SELECT status, version, upgrade_status FROM agents WHERE id = $1", agent_id)
current_agent_status = current_agent_row['status'] if current_agent_row else None
current_agent_version = current_agent_row['version'] if current_agent_row else None
current_upgrade_status = current_agent_row['upgrade_status'] if current_agent_row else None
# Determine new status - preserve upgrading status unless version actually changed
new_status = current_agent_status or 'online'
if current_agent_status == 'upgrading':
# CRITICAL FIX: If upgrade_status is NULL but status is 'upgrading', this is an error state
# This happens after cleanup job or manual intervention - agent should be online
if not current_upgrade_status:
new_status = 'online'
logger.info(f"AGENT STATUS FIX: Agent '{heartbeat_data.name}' status was 'upgrading' but upgrade_status is NULL - marking as 'online'")
# Only mark as online if version actually changed (upgrade completed)
elif heartbeat_data.version and heartbeat_data.version != current_agent_version:
new_status = 'online'
logger.info(f"AGENT UPGRADE: Agent '{heartbeat_data.name}' completed upgrade to version '{heartbeat_data.version}' (was: {current_agent_version})")
logger.info(f"AGENT UPGRADE: Status changed from 'upgrading' to 'online' - upgrade cycle complete")
else:
# Keep upgrading status if version hasn't changed yet
new_status = 'upgrading'
logger.debug(f"AGENT UPGRADE: Agent '{heartbeat_data.name}' still upgrading (version unchanged: {heartbeat_data.version})")
else:
new_status = 'online'
# Auto-detect platform if not provided or empty
detected_platform = heartbeat_data.platform
if not detected_platform or detected_platform.strip() == "":
# Try to detect from hostname or operating system
if heartbeat_data.hostname and 'mac' in heartbeat_data.hostname.lower():
detected_platform = 'darwin'
logger.info(f"PLATFORM AUTO-DETECT: Agent '{agent_name}' platform detected as 'darwin' from hostname")
elif heartbeat_data.operating_system and any(os_hint in heartbeat_data.operating_system.lower() for os_hint in ['macos', 'darwin', 'mac']):
detected_platform = 'darwin'
logger.info(f"PLATFORM AUTO-DETECT: Agent '{agent_name}' platform detected as 'darwin' from OS")
elif heartbeat_data.operating_system and any(os_hint in heartbeat_data.operating_system.lower() for os_hint in ['linux', 'ubuntu', 'centos', 'rhel', 'debian']):
detected_platform = 'linux'
logger.info(f"PLATFORM AUTO-DETECT: Agent '{agent_name}' platform detected as 'linux' from OS")
else:
# Default to linux for unknown platforms
detected_platform = 'linux'
logger.info(f"PLATFORM AUTO-DETECT: Agent '{agent_name}' platform defaulted to 'linux' (unknown)")
# Extract validated agent IP (returns None for invalid/empty)
agent_ip = _extract_agent_ip(heartbeat_data)
# IP/VIP change detection (non-critical logging)
try:
if agent_ip or heartbeat_data.keepalive_ip:
current_agent = await conn.fetchrow(
"SELECT ip_address, keepalive_ip FROM agents WHERE id = $1", agent_id)
if current_agent:
old_ip = str(current_agent['ip_address']) if current_agent['ip_address'] else None
old_vip = current_agent['keepalive_ip']
if agent_ip and old_ip and agent_ip != old_ip:
logger.info(f"IP CHANGE: Agent '{agent_name}' (id={agent_id}) IP changed: {old_ip} -> {agent_ip}")
if heartbeat_data.keepalive_ip and old_vip and heartbeat_data.keepalive_ip != old_vip:
logger.info(f"VIP CHANGE: Agent '{agent_name}' (id={agent_id}) VIP changed: {old_vip} -> {heartbeat_data.keepalive_ip}")
except Exception:
pass # Non-critical, never break heartbeat processing
# CRITICAL FIX: Only update applied_config_version if agent sends a VALID value
# Agent sends "none" on startup before any config is applied - don't override DB value with this
update_applied_version = heartbeat_data.applied_config_version and heartbeat_data.applied_config_version not in ["none", ""]
await conn.execute("""
UPDATE agents
SET status = $18,
last_seen = CURRENT_TIMESTAMP,
hostname = COALESCE($2, hostname),
platform = $3,
architecture = COALESCE($4, architecture),
version = COALESCE($5, version),
operating_system = COALESCE($6, operating_system),
kernel_version = COALESCE($7, kernel_version),
uptime = COALESCE($8, uptime),
cpu_count = COALESCE($9, cpu_count),
memory_total = COALESCE($10, memory_total),
disk_space = COALESCE($11, disk_space),
network_interfaces = COALESCE($12, network_interfaces),
capabilities = COALESCE($13, capabilities),
ip_address = COALESCE($14, ip_address),
haproxy_status = COALESCE($15, haproxy_status),
haproxy_version = COALESCE($16, haproxy_version),
applied_config_version = CASE WHEN $19 THEN $17 ELSE applied_config_version END,
keepalive_state = CASE WHEN $20::text IS NOT NULL THEN NULLIF($20::text, '') ELSE keepalive_state END,
keepalive_ip = CASE WHEN $21::text IS NOT NULL THEN NULLIF($21::text, '') ELSE keepalive_ip END,
updated_at = CURRENT_TIMESTAMP
WHERE id = $1
""", agent_id, heartbeat_data.hostname, detected_platform,
heartbeat_data.architecture, heartbeat_data.version,
heartbeat_data.operating_system, heartbeat_data.kernel_version,
heartbeat_data.uptime, heartbeat_data.cpu_count, heartbeat_data.memory_total,
heartbeat_data.disk_space,
# Convert lists to JSON for JSONB columns. Don't WIPE network_interfaces/capabilities
# when a heartbeat omits them: send NULL so COALESCE keeps the existing value (a bare
# `or []` would store "[]" and erase the reported NICs / keepalived_management on every
# daemon heartbeat that doesn't include them — issue #27 corporate test).
(heartbeat_data.network_interfaces if isinstance(heartbeat_data.network_interfaces, str)
else (json.dumps(heartbeat_data.network_interfaces) if heartbeat_data.network_interfaces else None)),
(heartbeat_data.capabilities if isinstance(heartbeat_data.capabilities, str)
else (json.dumps(heartbeat_data.capabilities) if heartbeat_data.capabilities else None)),
agent_ip, heartbeat_data.haproxy_status, heartbeat_data.haproxy_version,
heartbeat_data.applied_config_version, new_status, update_applied_version,
heartbeat_data.keepalive_state, heartbeat_data.keepalive_ip)
# Cache keepalive state in Redis for fast dashboard access
_keepalive_state = heartbeat_data.keepalive_state
if _keepalive_state is not None:
try:
from database.connection import redis_client
_redis_key = f"haproxy:agent:{agent_id}:keepalive"
if _keepalive_state not in ("", "NONE"):
redis_client.setex(
_redis_key,
90,
json.dumps({"state": _keepalive_state, "ip": heartbeat_data.keepalive_ip or ""})
)
else:
redis_client.delete(_redis_key)
except Exception:
pass # Redis failure should never break heartbeat
# API key validation is already done at the beginning - no need for final check
# Multi-agent token system allows multiple agents to use the same token
# CRITICAL FIX: Auto-heal agent pool_id BEFORE checking server_statuses
# This ensures pool_id is set even on first heartbeat (when haproxy config is empty)
if not agent['pool_id'] and heartbeat_data.cluster_id:
# Agent has no pool_id but sends cluster_id - auto-heal it
cluster_pool = await conn.fetchval(
"SELECT pool_id FROM haproxy_clusters WHERE id = $1",
heartbeat_data.cluster_id
)
if cluster_pool:
await conn.execute(
"UPDATE agents SET pool_id = $1 WHERE id = $2",
cluster_pool, agent_id
)
logger.info(f"AGENT AUTO-HEAL: Fixed agent '{agent_name}' pool_id to {cluster_pool} from heartbeat cluster_id {heartbeat_data.cluster_id}")
# Update in-memory agent dict for subsequent logic
agent['pool_id'] = cluster_pool
else:
logger.warning(f"AGENT AUTO-HEAL: Cluster {heartbeat_data.cluster_id} not found or has no pool assigned")
if heartbeat_data.server_statuses:
logger.info(f"AGENT HEARTBEAT: Agent '{agent_name}' reported server statuses: {heartbeat_data.server_statuses}")
# CRITICAL FIX: Find cluster_id even if agent pool_id is NULL
# This handles the case where cluster was created before pool, then pool was assigned later
cluster_id = None
# Method 1: Try to find cluster via agent's pool_id (if set)
if agent['pool_id']:
cluster_result = await conn.fetchrow(
"SELECT id FROM haproxy_clusters WHERE pool_id = $1 LIMIT 1",
agent['pool_id']
)
if cluster_result:
cluster_id = cluster_result['id']
logger.debug(f"CLUSTER LOOKUP: Found cluster {cluster_id} via agent pool_id {agent['pool_id']}")
# Method 2: If cluster_id not found and agent sends cluster_id in heartbeat, use that
# This provides fallback when agent pool_id is NULL (e.g. cluster created before pool)
if not cluster_id and heartbeat_data.cluster_id:
# Verify this cluster actually exists and get its pool
cluster_pool = await conn.fetchval(
"SELECT pool_id FROM haproxy_clusters WHERE id = $1",
heartbeat_data.cluster_id
)
if cluster_pool:
# Cluster exists, use it for stats processing
cluster_id = heartbeat_data.cluster_id
logger.info(f"CLUSTER LOOKUP: Using cluster_id {cluster_id} from agent heartbeat data")
else:
logger.warning(f"CLUSTER LOOKUP: Cluster {heartbeat_data.cluster_id} not found or has no pool assigned")
if cluster_id:
# Process stats for dashboard cache
try:
from services.dashboard_stats_service import dashboard_stats_service
import base64
# Log what agent sent
logger.info(f"HEARTBEAT: Agent '{agent_name}' heartbeat received")
logger.info(f"HEARTBEAT: Has haproxy_stats_csv: {heartbeat_data.haproxy_stats_csv is not None}")
logger.info(f"HEARTBEAT: Has server_statuses: {heartbeat_data.server_statuses is not None}")
# Decode CSV stats if provided (supports both base64 and plain text)
raw_csv = None
if heartbeat_data.haproxy_stats_csv:
try:
# Try base64 decode first (for new agents)
raw_csv = base64.b64decode(heartbeat_data.haproxy_stats_csv).decode('utf-8')
logger.info(f"HEARTBEAT: Decoded base64 CSV stats from agent '{agent_name}' ({len(raw_csv)} bytes)")
except Exception as decode_error:
# If base64 decode fails, treat as plain text (for old agents)
logger.info(f"HEARTBEAT: Base64 decode failed, treating as plain text CSV from agent '{agent_name}'")
raw_csv = heartbeat_data.haproxy_stats_csv
if raw_csv:
logger.debug(f"HEARTBEAT: CSV preview: {raw_csv[:200]}...")
else:
logger.warning(f" HEARTBEAT: Agent '{agent_name}' sent NO haproxy_stats_csv!")
await dashboard_stats_service.process_agent_stats(
agent_name=agent_name,
cluster_id=cluster_id,
raw_stats_csv=raw_csv,
server_statuses=heartbeat_data.server_statuses
)
logger.debug(f"STATS CACHE: Processed stats for agent '{agent_name}' cluster {cluster_id}")
except Exception as stats_error:
logger.error(f"Failed to process stats for dashboard: {stats_error}")
# Update database server statuses
for backend_name, servers in heartbeat_data.server_statuses.items():
for server_name, status in servers.items():
try:
# Try to update by exact server_name first
result = await conn.execute("""
UPDATE backend_servers
SET haproxy_status = $1, haproxy_status_updated_at = CURRENT_TIMESTAMP
WHERE backend_name = $2 AND server_name = $3 AND cluster_id = $4
""", status, backend_name, server_name, cluster_id)
logger.info(f"SERVER STATUS UPDATE: {backend_name}/{server_name} = {status}, cluster_id={cluster_id}, rows_updated={result}")
# If no exact match, try to update any server in the backend (agent parsing might be wrong)
if result == "UPDATE 0":
result2 = await conn.execute("""
UPDATE backend_servers
SET haproxy_status = $1, haproxy_status_updated_at = CURRENT_TIMESTAMP
WHERE backend_name = $2 AND cluster_id = $3
AND id = (SELECT MIN(id) FROM backend_servers WHERE backend_name = $2 AND cluster_id = $3 LIMIT 1)
""", status, backend_name, cluster_id)
logger.info(f"SERVER STATUS FALLBACK: {backend_name}/ANY_SERVER = {status}, cluster_id={cluster_id}, rows_updated={result2}")
else:
logger.debug(f"Updated server status: {backend_name}/{server_name} = {status}")
except Exception as e:
logger.warning(f"Failed to update server status for {backend_name}/{server_name}: {e}")
else:
logger.warning(f"Could not determine cluster_id for agent '{agent_name}' to update server statuses")
# NOTE: Pool assignment is now handled above in the cluster_id lookup logic
await close_database_connection(conn)
return {"status": "ok", "message": "Heartbeat received", "agent_id": agent_id}
except Exception as e:
logger.error(f"Heartbeat processing failed for agent '{heartbeat_data.name}': {e}", exc_info=True)
raise HTTPException(status_code=500, detail=f"Internal server error during heartbeat processing for agent '{heartbeat_data.name}'")
@router.get("/{agent_name}/config")
async def get_agent_config(agent_name: str, x_api_key: Optional[str] = Header(None)):
"""Get HAProxy configuration for specific agent"""
# Validate agent API key — MANDATORY (GHSA-3p5c-m5m4-mjpx). Checked BEFORE any
# DB work and before the existence check, so an unauthenticated caller learns
# neither the full haproxy.cfg nor whether the agent exists. Raised before the
# try so it is not swallowed by the generic handler; validate_agent_api_key(None)
# returns None without touching the DB.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(f"Missing/invalid API key for agent '{agent_name}' config fetch")
raise HTTPException(status_code=401, detail="Authentication required")
try:
conn = await get_database_connection()
# Get agent info first to check pool
# CRITICAL: Include cluster's haproxy_bin_path, haproxy_config_path, stats_socket_path
# These are needed for dynamic validation - cluster admin can change paths without reinstalling agent
agent_info = await conn.fetchrow("""
SELECT a.id, a.name, a.pool_id, hc.id as cluster_id, hc.name as cluster_name,
COALESCE(a.enabled, TRUE) as enabled,
hc.haproxy_bin_path, hc.haproxy_config_path, hc.stats_socket_path
FROM agents a
LEFT JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
""", agent_name)
if not agent_info:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail=f"Agent '{agent_name}' not found")
# Log which agent's API key was used (for audit trail)
if agent_auth['name'] == agent_name:
logger.info(f"Agent '{agent_name}' fetching config using its own API key")
else:
logger.info(f"Agent '{agent_name}' fetching config using API key from agent '{agent_auth['name']}'")
logger.debug(f"Config fetch authorized for agent '{agent_name}'")
if not agent_info['enabled']:
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": agent_info['cluster_id'],
"cluster_name": agent_info['cluster_name'],
"config_content": "# Agent is disabled - no configuration available\n",
"version": "disabled",
"status": "disabled"
}
cluster_id = agent_info['cluster_id']
if not cluster_id:
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": None,
"cluster_name": None,
"config_content": "# No cluster assigned to this agent\n",
"version": "no-cluster",
"status": "no_cluster"
}
# CRITICAL SECURITY: ONLY return APPLIED versions to agents
# NEVER return PENDING, REJECTED, or VALIDATION_FAILED versions
# Agents should only receive configurations that have been explicitly approved
config_version = await conn.fetchrow("""
SELECT version_name, config_content, checksum, is_active
FROM config_versions
WHERE cluster_id = $1 AND status = 'APPLIED' AND is_active = TRUE
ORDER BY created_at DESC
LIMIT 1
""", cluster_id)
# CRITICAL: Fallback also MUST require status = 'APPLIED'
# Previous code had a security gap where fallback could return PENDING versions
if not config_version:
config_version = await conn.fetchrow("""
SELECT version_name, config_content, checksum, is_active
FROM config_versions
WHERE cluster_id = $1 AND status = 'APPLIED'
ORDER BY created_at DESC
LIMIT 1
""", cluster_id)
if config_version:
logger.warning(f"GET CONFIG: Using inactive APPLIED version for cluster {cluster_id} (no active version found)")
if not config_version:
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": cluster_id,
"cluster_name": agent_info['cluster_name'],
"config_content": "# No configuration available for this cluster\n",
"version": "no-config",
"status": "no_config"
}
await conn.execute("""
UPDATE agents
SET config_version = $1, updated_at = CURRENT_TIMESTAMP
WHERE id = $2
""", config_version['version_name'], agent_info['id'])
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": cluster_id,
"cluster_name": agent_info['cluster_name'],
"config_content": config_version['config_content'],
"version": config_version['version_name'],
"checksum": config_version['checksum'],
"status": "available",
# Dynamic paths from cluster configuration - agent should use these for validation
"haproxy_bin_path": agent_info.get('haproxy_bin_path', '/usr/sbin/haproxy'),
"haproxy_config_path": agent_info.get('haproxy_config_path', '/etc/haproxy/haproxy.cfg'),
"stats_socket_path": agent_info.get('stats_socket_path', '/var/run/haproxy/admin.sock')
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Config retrieval failed: {e}")
raise HTTPException(status_code=500, detail=str(e))
@router.get("/{agent_name}/ssl-certificates")
async def get_agent_ssl_certificates(agent_name: str, since: Optional[str] = None, x_api_key: Optional[str] = Header(None)):
"""Get SSL certificates for specific agent's cluster"""
# Validate agent API key — MANDATORY (GHSA-3p5c-m5m4-mjpx). This response
# returns SSL private_key_content, so authentication is checked BEFORE any DB
# work and before the existence check. Raised before the try so the 401 is not
# swallowed; validate_agent_api_key(None) returns None without a DB hit.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(f"Missing/invalid API key for agent '{agent_name}' SSL certificates")
raise HTTPException(status_code=401, detail="Authentication required")
try:
conn = await get_database_connection()
# Get agent and cluster info first
agent_info = await conn.fetchrow("""
SELECT a.id, a.name, a.pool_id, hc.id as cluster_id, hc.name as cluster_name,
COALESCE(a.enabled, TRUE) as enabled
FROM agents a
LEFT JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
""", agent_name)
if not agent_info:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail=f"Agent '{agent_name}' not found")
# Log which agent's API key was used (for audit trail)
if agent_auth['name'] == agent_name:
logger.info(f"Agent '{agent_name}' fetching SSL certificates using its own API key")
else:
logger.info(f"Agent '{agent_name}' fetching SSL certificates using API key from agent '{agent_auth['name']}'")
logger.debug(f"SSL fetch authorized for agent '{agent_name}'")
if not agent_info['enabled']:
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": agent_info['cluster_id'],
"ssl_certificates": [],
"status": "disabled"
}
cluster_id = agent_info['cluster_id']
if not cluster_id:
await close_database_connection(conn)
return {
"agent_name": agent_name,
"cluster_id": None,
"ssl_certificates": [],
"status": "no_cluster"
}
# Build SSL certificates query with optional timestamp filtering
ssl_query = """
SELECT DISTINCT s.id, s.name, s.primary_domain as domain, s.certificate_content,
s.private_key_content, s.chain_content, s.expiry_date, s.status, s.fingerprint,
s.usage_type, s.created_at, s.updated_at, s.last_config_status
FROM ssl_certificates s
LEFT JOIN ssl_certificate_clusters scc ON s.id = scc.ssl_certificate_id
WHERE s.is_active = TRUE
AND (
-- Global SSLs: No cluster associations (not in junction table)
NOT EXISTS (SELECT 1 FROM ssl_certificate_clusters WHERE ssl_certificate_id = s.id)
-- Cluster-specific SSLs: Only for this specific cluster
OR scc.cluster_id = $1
)
-- Only return SSL certificates that have been APPLIED (not PENDING)
AND (s.last_config_status = 'APPLIED' OR s.last_config_status IS NULL)
"""
params = [cluster_id]
# Add timestamp filter if provided (for incremental updates)
if since:
try:
from datetime import datetime
import pytz
# Parse the ISO timestamp and convert to UTC naive datetime for database comparison
since_datetime = datetime.fromisoformat(since.replace('Z', '+00:00'))
since_naive = since_datetime.astimezone(pytz.UTC).replace(tzinfo=None)
ssl_query += " AND (s.updated_at > $2 OR s.created_at > $2)"
params.append(since_naive)
logger.info(f"SSL INCREMENTAL: Filtering SSL certificates since {since_naive} UTC")
except (ValueError, ImportError) as e:
logger.warning(f"SSL INCREMENTAL: Invalid timestamp format '{since}' or timezone error: {e}, ignoring filter")
ssl_query += " ORDER BY s.created_at DESC"
# Get SSL certificates for this cluster (both global and cluster-specific)
# Only return SSL certificates that have been APPLIED (not PENDING)
logger.info(f"SSL QUERY DEBUG: cluster_id={cluster_id}, params={params}")
logger.info(f"SSL QUERY: {ssl_query}")
ssl_certificates = await conn.fetch(ssl_query, *params)
logger.info(f"SSL RESULTS: Found {len(ssl_certificates)} certificates")
await close_database_connection(conn)
# Format SSL certificates for agent deployment
ssl_certs_data = []
latest_timestamp = None
for cert in ssl_certificates:
# Track the latest timestamp for incremental updates
cert_timestamp = max(cert['created_at'], cert['updated_at']) if cert['updated_at'] else cert['created_at']
if not latest_timestamp or cert_timestamp > latest_timestamp:
latest_timestamp = cert_timestamp
cert_data = {
"id": cert['id'],
"name": cert['name'],
"domain": cert['domain'],
"certificate_content": cert['certificate_content'],
"private_key_content": cert['private_key_content'],
"chain_content": cert['chain_content'],
"usage_type": cert.get('usage_type', 'frontend'), # Default to frontend for backward compatibility
"file_path": f"/etc/ssl/haproxy/{cert['name']}.pem",
"expiry_date": cert['expiry_date'].isoformat() if cert['expiry_date'] else None,
"status": cert['status'],
"fingerprint": cert['fingerprint'],
"created_at": cert['created_at'].isoformat() if cert['created_at'] else None,
"updated_at": cert['updated_at'].isoformat() if cert['updated_at'] else None
}
ssl_certs_data.append(cert_data)
response_data = {
"agent_name": agent_name,
"cluster_id": cluster_id,
"cluster_name": agent_info['cluster_name'],
"ssl_certificates": ssl_certs_data,
"status": "available",
"total_certificates": len(ssl_certs_data),
"latest_update": latest_timestamp.isoformat() if latest_timestamp else None,
"incremental_since": since
}
if since and ssl_certs_data:
logger.info(f"SSL INCREMENTAL: Returned {len(ssl_certs_data)} updated certificates since {since}")
elif since:
logger.info(f"SSL INCREMENTAL: No certificates updated since {since}")
return response_data
except HTTPException:
raise
except Exception as e:
logger.error(f"SSL certificates retrieval failed: {e}")
raise HTTPException(status_code=500, detail=str(e))
@router.get("/{agent_name}/keepalived-config")
async def get_agent_keepalived_config(agent_name: str, x_api_key: Optional[str] = Header(None)):
"""Issue #27 (v1.7.0) — deliver this agent's APPLIED keepalived snapshot, or a
teardown / no-op signal.
Snapshot-based (T-1): keys off vip_instances.is_active + the member's
applied_config_content, NOT the live last_config_status — so a PENDING edit never
flips a running member to not_configured (no mid-edit teardown). Auth is MANDATORY (a
valid agent key is required), but — like the /config and /ssl-certificates endpoints —
the agent API key is a SHARED/global install token, so the config is resolved by the
requested agent_name and a token/name mismatch is an advisory audit log, NOT a 403
(a hard 403 would break every VIP member whose name isn't the one row the shared token
resolves to — review HIGH-1). Any unexpected error degrades to not_configured (B-7) so
the agent stays inert; a node with no membership row always gets not_configured.
"""
conn = None
try:
# Auth FIRST and MANDATORY: a valid agent key is REQUIRED (the response carries the
# VRRP secret). The token is shared/global, so resolve by agent_name and only LOG a
# name mismatch — do not 403 (HIGH-1). Done before the agent lookup so an
# unauthenticated caller can't probe which agent names exist.
from auth_middleware import validate_agent_api_key
if not x_api_key:
raise HTTPException(status_code=401, detail="Agent API key required")
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
raise HTTPException(status_code=401, detail="Invalid API key")
# The agent API key is a SHARED/global install token (many agent rows per token),
# so validate_agent_api_key resolves it to one arbitrary agent for that token. Resolve
# the keepalived config strictly by the requested agent_name and treat a token/name
# mismatch as an advisory audit log — exactly like the /config and /ssl-certificates
# endpoints. (A hard 403 here would reject every VIP member whose name isn't the one
# row the shared token happens to return, so the VIP could never converge — review HIGH-1.)
if agent_auth['name'] == agent_name:
logger.info(f"Agent '{agent_name}' fetching keepalived config using its own API key")
else:
logger.info(f"Agent '{agent_name}' fetching keepalived config using API key from agent '{agent_auth['name']}'")
conn = await get_database_connection()
# Resolve the agent + its cluster's keepalived.conf path (cluster-driven, like the
# HAProxy paths). config_path is returned in EVERY response so the agent knows where
# to write/own-marker-check even on not_configured/teardown.
agent = await conn.fetchrow("""
SELECT a.id, a.name, COALESCE(a.enabled, TRUE) AS enabled,
hc.keepalived_config_path,
-- v1.11.1: does the server already hold a discovery for this node? The agent
-- caches the hash of its last discovery report next to the config and skips
-- re-posting while it matches. That cache used to be written even when the
-- POST was REJECTED, so a node could be hidden from the adoption panel for
-- good: the file never changes, so the agent never speaks again. Telling it
-- what we actually hold lets it recover on its own, with no extra request and
-- no one having to touch the node.
EXISTS (SELECT 1 FROM vip_discoveries vd WHERE vd.agent_id = a.id)
AS discovery_known
FROM agents a
LEFT JOIN haproxy_clusters hc ON hc.pool_id = a.pool_id
WHERE a.name = $1
-- A pool may hold more than one cluster, and the join then multiplies this row. With
-- no ordering the fetch took an arbitrary one, so the keepalived.conf PATH handed to
-- the agent was non-deterministic whenever two clusters in a pool disagreed on it:
-- the agent would look at the wrong file, find nothing there, and the node would
-- never appear for adoption — intermittently, which is the worst way to fail.
--
-- A CUSTOMISED path wins over the shipped default, then the lowest cluster id. The
-- column defaults to '/etc/keepalived/keepalived.conf' rather than NULL, so ordering
-- by id alone could have picked a default-valued row over one the operator had
-- deliberately set — turning "undefined" into "reliably wrong" for that install.
-- When every cluster in the pool carries the default the string is identical, so the
-- ordering cannot change what any working deployment already receives.
ORDER BY (hc.keepalived_config_path IS NULL
OR hc.keepalived_config_path = '/etc/keepalived/keepalived.conf'),
hc.id
LIMIT 1
""", agent_name)
if not agent:
raise HTTPException(status_code=404, detail=f"Agent '{agent_name}' not found")
config_path = agent['keepalived_config_path'] or '/etc/keepalived/keepalived.conf'
if not agent['enabled']:
return {"agent_name": agent_name, "status": "not_configured", "config_path": config_path, "discovery_known": bool(agent["discovery_known"]), "keepalived": None}
row = await conn.fetchrow("""
SELECT v.id AS vip_id, v.name AS vip_name, v.is_active, v.track_haproxy,
v.purge_on_teardown,
m.applied_config_content, m.applied_config_hash, m.takeover_expected_hash
FROM vip_members m JOIN vip_instances v ON v.id = m.vip_id
WHERE m.agent_id = $1
-- Active VIP first (an agent has at most one). With NO active VIP, pick the most
-- RECENTLY updated inactive membership so a teardown reflects the latest delete
-- (incl. its purge flag) — not a stale older VIP the node was once part of.
ORDER BY v.is_active DESC, v.updated_at DESC, v.id DESC
LIMIT 1
""", agent['id'])
if not row:
return {"agent_name": agent_name, "status": "not_configured", "config_path": config_path, "discovery_known": bool(agent["discovery_known"]), "keepalived": None}
if not row['is_active']:
# Soft-deleted VIP → teardown. purge carries the operator's opt-in package removal;
# the agent still only purges on nodes where IT installed keepalived (install marker).
return {"agent_name": agent_name, "status": "teardown", "vip_id": row['vip_id'],
"config_path": config_path, "discovery_known": bool(agent["discovery_known"]), "keepalived": None,
"purge": bool(row['purge_on_teardown'])}
if not row['applied_config_content']:
return {"agent_name": agent_name, "status": "not_configured", "config_path": config_path, "discovery_known": bool(agent["discovery_known"]), "keepalived": None}
from services.keepalived_config import build_haproxy_check_script
check_script = build_haproxy_check_script() if row['track_haproxy'] else ""
return {
"agent_name": agent_name,
"status": "available",
"config_path": config_path, "discovery_known": bool(agent["discovery_known"]),
"keepalived": {
"desired_state": "enabled",
"install_if_missing": True,
"vip_id": row['vip_id'],
"vip_name": row['vip_name'],
"config_content": row['applied_config_content'],
"config_hash": row['applied_config_hash'],
"check_script": check_script,
# v1.10.4 adoption handoff. The agent refuses to overwrite a keepalived.conf
# without our ownership marker — the guard that protects a hand-maintained
# setup. Adoption does not weaken it: it authorises exactly ONE takeover, of
# exactly the file we analysed, by pinning its hash. If the file changed since
# adoption the hashes differ and the agent keeps refusing, so an edit made
# between adoption and Apply can never be silently overwritten.
"allow_takeover": bool(row['takeover_expected_hash']),
"takeover_expected_hash": row['takeover_expected_hash'],
},
}
except HTTPException:
raise
except Exception as e:
logger.error(f"keepalived-config delivery failed for '{agent_name}': {e}")
# Degrade to no-op rather than 500 (B-7) — keeps the fleet inert on any error.
return {"agent_name": agent_name, "status": "not_configured",
"config_path": "/etc/keepalived/keepalived.conf", "keepalived": None}
finally:
if conn:
await close_database_connection(conn)
@router.post("/{agent_name}/keepalived-status")
async def agent_keepalived_status(agent_name: str, status_data: dict, x_api_key: Optional[str] = Header(None)):
"""Issue #27 (v1.7.0) — agent reports the outcome of a keepalived deploy/teardown.
Auth mirrors config-applied's post-Bulgu-#75 guard (reject a MISSING key — never the
`and` short-circuit that accepted no-key requests). The token is a shared/global install
token, so the status is recorded strictly for the requested agent_name and a token/name
mismatch is an advisory audit log (like /config-applied), not a 403 — review HIGH-1.
"""
conn = None
try:
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not x_api_key or not agent_auth:
raise HTTPException(status_code=401, detail="Invalid API key")
if agent_auth['name'] != agent_name:
logger.info(f"Agent '{agent_name}' reporting keepalived status using API key from agent '{agent_auth['name']}'")
conn = await get_database_connection()
agent = await conn.fetchrow("SELECT id FROM agents WHERE name = $1", agent_name)
if not agent:
raise HTTPException(status_code=404, detail=f"Agent '{agent_name}' not found")
vip_id = status_data.get("vip_id")
state = (status_data.get("state") or "").strip()[:24]
config_hash = (status_data.get("config_hash") or "")[:64]
message = status_data.get("message")
# v1.10.9 — retire the adoption takeover authorisation once the node CONFIRMS it is
# running our rendered config. `takeover_expected_hash` is the permission to overwrite a
# keepalived.conf that lacks our ownership marker; it was written at adoption and never
# cleared, so it stayed valid indefinitely and "one-shot" was only true in the sense of
# "for exactly that file content". Clearing it the moment the member acks OUR hash makes
# the claim real: if the file is replaced by hand afterwards the agent refuses and reports
# "externally managed", which is the visible behaviour an operator should get.
#
# Gated on the acked hash MATCHING applied_config_hash, so a partial or failed deploy
# never drops the authorisation and leaves the VIP unable to converge.
# The hash is passed TWICE on purpose. Reusing one placeholder for both the assignment
# (`last_deploy_hash=$n`, a VARCHAR column) and the comparison inside the CASE made
# PostgreSQL deduce two different types for it and asyncpg refused the whole statement
# with AmbiguousParameterError ("text versus character varying"). Because the failure is
# in the UPDATE itself, not in one column, EVERY status ack was lost and every VIP sat at
# SYNCING forever — including teardown acks. A separate placeholder is only ever compared
# against the column, so its type is unambiguous.
_retire_takeover = ("takeover_expected_hash = CASE WHEN applied_config_hash IS NOT NULL "
"AND applied_config_hash = {p} THEN NULL ELSE takeover_expected_hash END")
if vip_id is None:
# No specific VIP (e.g. a teardown ack) — update all this agent's memberships.
await conn.execute(f"""
UPDATE vip_members SET last_deploy_state=$2, last_deploy_message=$3,
last_deploy_hash=$4, last_deploy_at=CURRENT_TIMESTAMP, updated_at=CURRENT_TIMESTAMP,
{_retire_takeover.format(p="$5")}
WHERE agent_id=$1
""", agent['id'], state, message, config_hash, config_hash)
else:
await conn.execute(f"""
UPDATE vip_members SET last_deploy_state=$3, last_deploy_message=$4,
last_deploy_hash=$5, last_deploy_at=CURRENT_TIMESTAMP, updated_at=CURRENT_TIMESTAMP,
{_retire_takeover.format(p="$6")}
WHERE agent_id=$1 AND vip_id=$2
""", agent['id'], int(vip_id), state, message, config_hash, config_hash)
return {"status": "ok"}
except HTTPException:
raise
except Exception as e:
logger.error(f"keepalived-status update failed for '{agent_name}': {e}")
raise HTTPException(status_code=500, detail="keepalived-status update failed")
finally:
if conn:
await close_database_connection(conn)
@router.post("/{agent_name}/keepalived-discovery")
async def agent_keepalived_discovery(agent_name: str, payload: dict, x_api_key: Optional[str] = Header(None)):
"""v1.10.4 — the agent reports a keepalived.conf it found on the node but does NOT own.
This is what makes adopting a hand-maintained VIP possible: the heartbeat only carries the
VIP address and a best-effort MASTER/BACKUP, while rendering a node's config needs eleven
fields, so the file itself has to be read. Read-only on the agent side — reporting never
changes anything on the node.
Auth mirrors /keepalived-status: a MISSING key is rejected outright, and because the token
is a shared install token a name mismatch is an advisory audit log rather than a 403.
SECRETS: the reported content may contain the VRRP `auth_pass`. It is split immediately —
the password is Fernet-encrypted into its own column and the stored copy of the file has it
masked, so nothing readable through the API or a DB dump carries it in cleartext. The
parse result is never logged.
"""
conn = None
try:
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not x_api_key or not agent_auth:
raise HTTPException(status_code=401, detail="Invalid API key")
if agent_auth['name'] != agent_name:
logger.info(f"Agent '{agent_name}' reporting keepalived discovery using API key "
f"from agent '{agent_auth['name']}'")
conn = await get_database_connection()
agent = await conn.fetchrow("SELECT id FROM agents WHERE name = $1", agent_name)
if not agent:
raise HTTPException(status_code=404, detail=f"Agent '{agent_name}' not found")
config_path = (payload.get("config_path") or "/etc/keepalived/keepalived.conf")[:500]
exists = bool(payload.get("exists"))
if not exists:
# The file is gone (keepalived removed, or we adopted and now own it) — drop the row
# so the UI stops offering a stale candidate.
await conn.execute("DELETE FROM vip_discoveries WHERE agent_id = $1", agent['id'])
return {"status": "cleared"}
content = payload.get("config_content") or ""
if len(content) > 256_000:
raise HTTPException(status_code=413, detail="keepalived.conf too large to analyse")
is_managed = bool(payload.get("is_managed"))
config_hash = hashlib.md5(content.encode("utf-8", "replace")).hexdigest()
from services.keepalived_parser import analyse_keepalived_conf, KeepalivedParseError
from services.keepalived_config import encrypt_vrrp_secret
parse_error = None
analysis = None
auth_enc = None
try:
analysis = analyse_keepalived_conf(content)
# Split the secret out of everything we persist or serve.
for cand in analysis.get("candidates", []):
secret = (cand.get("vip") or {}).pop("auth_pass", None)
cand["vip"]["has_auth_pass"] = bool(secret)
if secret and auth_enc is None:
auth_enc = encrypt_vrrp_secret(secret)
except KeepalivedParseError as exc:
parse_error = str(exc)[:500]
except Exception as exc: # noqa: BLE001 — a malformed file must not 500 the agent loop
parse_error = f"could not analyse the config ({type(exc).__name__})"
masked = _AUTH_PASS_MASK_RE.sub(r"\1********", content)
await conn.execute("""
INSERT INTO vip_discoveries
(agent_id, config_path, config_hash, is_managed, raw_config_masked,
auth_pass_encrypted, analysis, parse_error, reported_at)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,CURRENT_TIMESTAMP)
ON CONFLICT (agent_id) DO UPDATE SET
config_path = EXCLUDED.config_path,
config_hash = EXCLUDED.config_hash,
is_managed = EXCLUDED.is_managed,
raw_config_masked = EXCLUDED.raw_config_masked,
auth_pass_encrypted = EXCLUDED.auth_pass_encrypted,
analysis = EXCLUDED.analysis,
parse_error = EXCLUDED.parse_error,
reported_at = CURRENT_TIMESTAMP
""", agent['id'], config_path, config_hash, is_managed, masked, auth_enc,
json.dumps(analysis) if analysis is not None else None, parse_error)
return {"status": "recorded", "config_hash": config_hash}
except HTTPException:
raise
except Exception as e:
logger.error(f"keepalived-discovery failed for '{agent_name}': {e}")
raise HTTPException(status_code=500, detail="keepalived-discovery failed")
finally:
if conn:
await close_database_connection(conn)
@router.get("/script-version")
async def get_latest_script_version(platform: str = "macos"):
"""Get the latest available agent script version for specified platform"""
try:
conn = await get_database_connection()
# Get latest version for platform
version_info = await conn.fetchrow("""
SELECT version, release_date, changelog
FROM agent_versions
WHERE platform = $1 AND is_active = true
ORDER BY created_at DESC
LIMIT 1
""", platform.lower())
await close_database_connection(conn)
if version_info:
return {
"latest_version": version_info['version'],
"release_date": version_info['release_date'].isoformat() if version_info['release_date'] else None,
"changelog": version_info['changelog'] or []
}
else:
# Fallback to default if no database entry
# Use global version storage
latest_version = AGENT_VERSIONS.get(platform.lower(), 'unknown')
return {
"latest_version": latest_version,
"release_date": "2024-01-23",
"changelog": ["Version from global storage"]
}
except Exception as e:
logger.error(f"Error getting script version: {e}")
# Fallback to hardcoded version
# Use global version storage as fallback
latest_version = AGENT_VERSIONS.get(platform.lower(), 'unknown')
return {
"latest_version": latest_version,
"release_date": "2024-01-23",
"changelog": ["Fallback version from global storage"]
}
@router.get("/{agent_name}/upgrade-status")
async def get_agent_upgrade_status(agent_name: str, x_api_key: Optional[str] = Header(None)):
"""Get agent upgrade status - used by agents to check if they should upgrade"""
# Validate agent API key — MANDATORY (GHSA-3p5c-m5m4-mjpx). Checked before any
# DB work; deployed agents always send X-API-Key. Raised before the try so the
# 401 is not swallowed; validate_agent_api_key(None) returns None without a DB hit.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(f"Missing/invalid API key for agent '{agent_name}' upgrade status")
raise HTTPException(status_code=401, detail="Authentication required")
try:
conn = await get_database_connection()
# Check if agent has upgrade pending (include platform and pool for validation)
agent = await conn.fetchrow("""
SELECT status, version as current_version, platform, pool_id
FROM agents
WHERE name = $1
""", agent_name)
if not agent:
await close_database_connection(conn)
return {
"should_upgrade": False,
"target_version": "",
"message": "Agent not found"
}
# Log which agent's API key was used (for audit trail)
if agent_auth['name'] == agent_name:
logger.debug(f"Agent '{agent_name}' checking upgrade status using its own API key")
else:
logger.info(f"Agent '{agent_name}' checking upgrade status using API key from agent '{agent_auth['name']}'")
logger.debug(f"Upgrade status check authorized for agent '{agent_name}'")
await close_database_connection(conn)
# Agent should upgrade if status is 'upgrading'
should_upgrade = agent['status'] == 'upgrading'
# Debug logging for upgrade status
logger.info(f"UPGRADE STATUS: agent='{agent_name}', status='{agent['status']}', current_version='{agent['current_version']}', should_upgrade={should_upgrade}")
# Get target version from global storage based on agent platform
try:
if should_upgrade:
agent_platform = agent.get('platform', 'unknown')
target_version = await get_version_for_platform(agent_platform)
platform_key = get_platform_key(agent_platform)
logger.info(f"UPGRADE TARGET: agent='{agent_name}', platform='{platform_key}', target_version='{target_version}', current_version='{agent['current_version']}'")
else:
target_version = ""
logger.info(f"NO UPGRADE: agent='{agent_name}' status is '{agent['status']}' (not 'upgrading')")
except Exception as version_error:
logger.warning(f"Could not get target version for agent {agent_name}: {version_error}")
target_version = "unknown" if should_upgrade else ""
return {
"should_upgrade": should_upgrade,
"target_version": target_version,
"current_version": agent['current_version'],
"message": "Upgrade required" if should_upgrade else "No upgrade needed"
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error checking upgrade status for agent {agent_name}: {e}")
raise HTTPException(status_code=500, detail="Failed to check upgrade status")
@router.post("/{agent_name}/upgrade-complete")
async def agent_upgrade_complete(agent_name: str, completion_data: dict, x_api_key: Optional[str] = Header(None)):
"""Receive notification when agent completes or fails upgrade"""
try:
# Bulgu #75 (round-22 audit) — see config-applied above.
from auth_middleware import validate_agent_api_key
agent_auth = await validate_agent_api_key(x_api_key)
if not agent_auth:
logger.warning(
f"Rejected upgrade-complete call for agent {agent_name!r}: "
f"missing or invalid x-api-key"
)
raise HTTPException(status_code=401, detail="Invalid or missing API key")
# GLOBAL TOKEN: Token can be used by multiple agents across different pools/clusters
conn = await get_database_connection()
upgrade_status = completion_data.get("status", "unknown")
new_version = completion_data.get("version", "")
if upgrade_status == "completed":
# Update agent version and reset status to online
await conn.execute("""
UPDATE agents
SET version = $1, status = 'online', updated_at = CURRENT_TIMESTAMP
WHERE name = $2
""", new_version, agent_name)
logger.info(f"AGENT UPGRADE: Agent '{agent_name}' completed upgrade to version '{new_version}'")
# Log activity (always "upgrade")
await log_agent_activity(agent_name, 'upgrade', {"version": new_version, "status": "completed"})
elif upgrade_status == "failed":
# Reset status to online but don't update version
await conn.execute("""
UPDATE agents
SET status = 'online', updated_at = CURRENT_TIMESTAMP
WHERE name = $1
""", agent_name)
logger.error(f"AGENT UPGRADE: Agent '{agent_name}' failed to upgrade to version '{new_version}'")
await close_database_connection(conn)
return {"status": "ok", "message": f"Upgrade completion status received: {upgrade_status}"}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error processing upgrade completion for agent {agent_name}: {e}")
return {"status": "error", "message": str(e)}
@router.post("/{agent_id}/upgrade")
async def upgrade_agent(agent_id: int, authorization: str = Header(None)):
"""Initiate agent upgrade process"""
try:
# Get current user and check permission
current_user = await get_current_user_from_token(authorization)
has_permission = await check_user_permission(current_user["id"], "agents", "upgrade")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.upgrade required"
)
conn = await get_database_connection()
# Get agent info including cluster relationship and platform
agent = await conn.fetchrow("""
SELECT a.name, a.pool_id, a.version as current_version, a.platform, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON a.pool_id = hc.pool_id
WHERE a.id = $1
""", agent_id)
if not agent:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail="Agent not found")
# Validate cluster access if agent belongs to a cluster
if agent['cluster_id']:
await validate_user_cluster_access(current_user['id'], agent['cluster_id'], conn)
# Get latest version from global storage based on agent platform
try:
agent_platform = agent.get('platform', 'unknown')
logger.info(f"AGENT PLATFORM DEBUG: Raw platform='{agent_platform}', agent_name='{agent['name']}'")
latest_version = await get_version_for_platform(agent_platform)
platform_key = get_platform_key(agent_platform)
logger.info(f"VERSION: Using global version {latest_version} for {platform_key} upgrade (platform: {agent_platform})")
except Exception as version_error:
logger.warning(f"Could not get latest version: {version_error}")
latest_version = "unknown" # No fallback
# SIMPLIFIED: No upgrade/downgrade distinction - just sync to latest version
# If versions differ, agent needs to sync to the latest script from database
if agent['current_version'] == latest_version:
await close_database_connection(conn)
raise HTTPException(status_code=400, detail="Agent is already on this version")
# Update agent status to 'upgrading' with full upgrade tracking
await conn.execute("""
UPDATE agents
SET status = 'upgrading',
upgrade_status = 'upgrading',
upgrade_target_version = $2,
upgraded_at = CURRENT_TIMESTAMP,
updated_at = CURRENT_TIMESTAMP
WHERE id = $1
""", agent_id, latest_version)
await close_database_connection(conn)
# Log the upgrade activity (always "upgrade", regardless of version direction)
await log_user_activity(
user_id=current_user["id"],
action='upgrade',
resource_type='agent',
resource_id=str(agent_id),
details={'agent_name': agent['name'], 'from_version': agent['current_version'], 'to_version': latest_version}
)
return {
"message": f"Agent '{agent['name']}' upgrade initiated",
"agent_id": agent_id,
"from_version": agent['current_version'],
"to_version": latest_version,
"status": "upgrading"
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error upgrading agent {agent_id}: {e}")
raise HTTPException(status_code=500, detail="Failed to initiate agent upgrade")
@router.put("/{agent_id}/toggle", summary="Toggle Agent Status", response_description="Agent status toggled")
async def toggle_agent(agent_id: int, toggle_data: AgentToggle, authorization: str = Header(None)):
"""
# Toggle Agent Enabled/Disabled Status
Enable or disable an agent. Disabled agents will not receive configuration tasks but remain registered.
## Path Parameters
- **agent_id**: Agent ID to toggle
## Request Body
- **enabled**: Boolean (true to enable, false to disable)
## Example Request - Disable Agent
```bash
curl -X PUT "{BASE_URL}/api/agents/1/toggle" \\
-H "Authorization: Bearer eyJhbGciOiJIUz..." \\
-H "Content-Type: application/json" \\
-d '{"enabled": false}'
```
## Example Response
```json
{
"message": "Agent 'production-agent-01' disabled successfully",
"enabled": false
}
```
## Use Cases
- **Disable**: Temporarily stop sending config to agent during maintenance
- **Enable**: Resume normal operation
Note: Agent service keeps running but won't receive new tasks when disabled.
## Error Responses
- **403**: Insufficient permissions
- **404**: Agent not found
- **500**: Server error
"""
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()
# Get agent info including cluster relationship
agent = await conn.fetchrow("""
SELECT a.name, a.pool_id, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON a.pool_id = hc.pool_id
WHERE a.id = $1
""", agent_id)
if not agent:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail="Agent not found")
# Validate cluster access if agent belongs to a cluster
if agent['cluster_id']:
await validate_user_cluster_access(current_user['id'], agent['cluster_id'], conn)
agent = await conn.fetchrow("SELECT * FROM agents WHERE id = $1", agent_id)
if not agent:
raise HTTPException(status_code=404, detail="Agent not found")
await conn.execute(
"UPDATE agents SET enabled = $1, updated_at = CURRENT_TIMESTAMP WHERE id = $2",
toggle_data.enabled, agent_id
)
await close_database_connection(conn)
return {
"success": True,
"message": f"Agent {'enabled' if toggle_data.enabled else 'disabled'} successfully",
"agent_id": agent_id,
"enabled": toggle_data.enabled
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error toggling agent {agent_id}: {e}")
raise HTTPException(status_code=500, detail="Failed to toggle agent status")
# ==== AGENT VERSION MANAGEMENT ENDPOINTS ====
@router.get("/versions")
async def get_agent_versions(authorization: str = Header(None)):
"""Get all agent versions by platform"""
try:
current_user = await get_current_user_from_token(authorization)
# Get versions from database
conn = await get_database_connection()
versions = await conn.fetch("""
SELECT platform, version, release_date, changelog, is_active, created_at
FROM agent_versions
WHERE is_active = true
ORDER BY platform, created_at DESC
""")
# Group by platform
result = {}
for version_row in versions:
platform = version_row['platform']
if platform not in result:
result[platform] = []
result[platform].append({
"version": version_row['version'],
"release_date": version_row['release_date'].isoformat() if version_row['release_date'] else "2024-01-23",
"changelog": version_row['changelog'] or [
"Database-driven version management",
"Self-updating agent scripts",
"Production-ready upgrade system"
],
"is_active": version_row['is_active'],
"created_at": version_row['created_at'].isoformat() if version_row['created_at'] else "2024-01-23T00:00:00Z"
})
# Add fallback for platforms not in database
for platform, version in AGENT_VERSIONS.items():
if platform not in result:
result[platform] = [{
"version": version,
"release_date": "2024-01-23",
"changelog": ["Fallback from global storage"],
"is_active": True,
"created_at": "2024-01-23T00:00:00Z"
}]
# Detect if on-disk agent scripts differ from database (new version shipped)
script_update_available = False
try:
script_files = {'linux': 'linux_install.sh', 'macos': 'macos_install.sh'}
for platform_key, filename in script_files.items():
file_path = os.path.join(os.path.dirname(__file__), '..', 'utils', 'agent_scripts', filename)
if not os.path.exists(file_path):
continue
with open(file_path, 'r') as f:
current_file_hash = hashlib.sha256(f.read().encode()).hexdigest()
db_row = await conn.fetchrow("""
SELECT source_file_hash FROM agent_script_templates
WHERE platform = $1 AND is_active = true
ORDER BY updated_at DESC LIMIT 1
""", platform_key)
if db_row and db_row['source_file_hash']:
if current_file_hash != db_row['source_file_hash']:
script_update_available = True
break
else:
db_content_row = await conn.fetchrow("""
SELECT script_content FROM agent_script_templates
WHERE platform = $1 AND is_active = true
ORDER BY updated_at DESC LIMIT 1
""", platform_key)
if db_content_row:
db_hash = hashlib.sha256(db_content_row['script_content'].encode()).hexdigest()
if current_file_hash != db_hash:
script_update_available = True
break
except Exception as e:
logger.warning(f"Script update detection failed: {e}")
await close_database_connection(conn)
return {"platforms": result, "script_update_available": script_update_available}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error getting agent versions: {e}")
raise HTTPException(status_code=500, detail=str(e))
@router.post("/versions")
async def create_agent_version(version_data: dict, authorization: str = Header(None)):
"""Create new agent version for specified platform"""
try:
current_user = await get_current_user_from_token(authorization)
# Check permission for agent version management
has_permission = await check_user_permission(current_user["id"], "agents", "version")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.version required"
)
platform = version_data.get('platform', '').lower()
version = version_data.get('version', '')
changelog = version_data.get('changelog', [])
if not platform or not version:
raise HTTPException(status_code=400, detail="Platform and version are required")
# Save to database
conn = await get_database_connection()
# Insert new version (or update if exists)
await conn.execute("""
INSERT INTO agent_versions (platform, version, changelog, is_active)
VALUES ($1, $2, $3, true)
ON CONFLICT (platform, version)
DO UPDATE SET
changelog = EXCLUDED.changelog,
is_active = true,
updated_at = CURRENT_TIMESTAMP
""", platform, version, changelog or [])
# Deactivate old versions for this platform
await conn.execute("""
UPDATE agent_versions
SET is_active = false
WHERE platform = $1 AND version != $2
""", platform, version)
await close_database_connection(conn)
# Also update global storage for backward compatibility
global AGENT_VERSIONS
AGENT_VERSIONS[platform] = version
# Log activity
await log_user_activity(
user_id=current_user["id"],
action='create',
resource_type='agent_version',
resource_id=f"{platform}-{version}",
details={
'platform': platform,
'version': version,
'changelog_items': len(changelog)
}
)
logger.info(f"VERSION UPDATE: {platform} version updated to {version}")
return {
"message": f"Agent version {version} created for {platform}",
"platform": platform,
"version": version,
"current_versions": AGENT_VERSIONS
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error creating agent version: {e}")
raise HTTPException(status_code=500, detail=str(e))
# ==== AGENT SCRIPT TEMPLATE MANAGEMENT ENDPOINTS ====
@router.get("/script-templates/{platform}")
async def get_agent_script_template(platform: str, authorization: str = Header(None)):
"""Get the latest script template for specified platform from database"""
try:
current_user = await get_current_user_from_token(authorization)
# SECURITY (GHSA-7rhv-c5pc-69r8): the raw install/upgrade script is a
# version-management surface. Gate reads with agents.version too, matching
# the write path above (operator/security_admin/super_admin retain access).
has_permission = await check_user_permission(current_user["id"], "agents", "version")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.version required"
)
conn = await get_database_connection()
# Get latest script template for platform
template = await conn.fetchrow("""
SELECT ast.script_content, ast.version, av.changelog
FROM agent_script_templates ast
JOIN agent_versions av ON ast.platform = av.platform AND ast.version = av.version
WHERE ast.platform = $1 AND ast.is_active = true AND av.is_active = true
ORDER BY ast.created_at DESC
LIMIT 1
""", platform.lower())
await close_database_connection(conn)
if template:
return {
"platform": platform,
"version": template['version'],
"script_content": template['script_content'],
"changelog": template['changelog'] or []
}
else:
# Fallback to file-based template if not in database
import os
script_template_path = os.path.join(os.path.dirname(__file__), '..', 'utils', 'agent_scripts', f"{platform.lower()}_install.sh")
if os.path.exists(script_template_path):
with open(script_template_path, 'r') as f:
script_content = f.read()
return {
"platform": platform,
"version": "1.0.0",
"script_content": script_content,
"changelog": ["Loaded from file template"]
}
else:
raise HTTPException(status_code=404, detail=f"Script template not found for platform: {platform}")
except HTTPException:
raise
except Exception as e:
logger.error(f"Error getting script template for {platform}: {e}")
raise HTTPException(status_code=500, detail=str(e))
@router.post("/script-templates/{platform}")
async def save_agent_script_template(platform: str, template_data: dict, authorization: str = Header(None)):
"""Save updated script template to database using shared helper function"""
try:
current_user = await get_current_user_from_token(authorization)
# SECURITY (GHSA-7rhv-c5pc-69r8): agent script templates become the
# install/self-upgrade script executed as root on HAProxy nodes. A poisoned
# template is RCE. Authentication alone is NOT enough — require the same
# agents.version permission as POST /versions; otherwise any JWT holder
# (including viewer) could overwrite the active script.
has_permission = await check_user_permission(current_user["id"], "agents", "version")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.version required"
)
script_content = template_data.get('script_content', '')
version = template_data.get('version', '')
if not script_content or not version:
raise HTTPException(status_code=400, detail="Script content and version are required")
conn = await get_database_connection()
# Use the same update logic as Reset to Default operation
await update_agent_script_in_database(
conn=conn,
platform=platform.lower(),
version=version,
script_content=script_content,
changelog=['Script template updated via UI editor']
)
await close_database_connection(conn)
# Log activity
await log_user_activity(
user_id=current_user["id"],
action='update',
resource_type='agent_script_template',
resource_id=f"{platform}-{version}",
details={
'platform': platform,
'version': version,
'script_size': len(script_content),
'source': 'ui_editor'
}
)
logger.info(f"SCRIPT TEMPLATE: {platform} template updated to version {version}")
return {
"message": f"Script template saved for {platform}",
"platform": platform,
"version": version
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error saving script template for {platform}: {e}")
raise HTTPException(status_code=500, detail=str(e))
# ==== AGENT ACTIVITY LOGGING ====
async def log_agent_activity(agent_name: str, action_type: str, action_details: dict = None):
"""Log meaningful agent actions to database"""
try:
conn = await get_database_connection()
# Get agent_id
agent = await conn.fetchrow("""
SELECT id FROM agents WHERE name = $1
""", agent_name)
if not agent:
logger.warning(f"Cannot log activity for unknown agent: {agent_name}")
await close_database_connection(conn)
return
# Insert activity log
await conn.execute("""
INSERT INTO agent_activity_logs (agent_id, agent_name, action_type, action_details, timestamp)
VALUES ($1, $2, $3, $4, CURRENT_TIMESTAMP)
""", agent['id'], agent_name, action_type, json.dumps(action_details) if action_details else None)
# Update last_action_time in agents table
await conn.execute("""
UPDATE agents SET last_action_time = CURRENT_TIMESTAMP WHERE name = $1
""", agent_name)
await close_database_connection(conn)
logger.info(f"ACTIVITY LOG: {agent_name} - {action_type}")
except Exception as e:
logger.error(f"Failed to log agent activity: {e}")
# Don't raise - logging failure shouldn't break the main flow
@router.get("/{agent_name}/activity-logs")
async def get_agent_activity_logs(agent_name: str, limit: int = 50, authorization: str = Header(None)):
"""Get activity logs for an agent"""
try:
from auth_middleware import check_user_permission
current_user = await get_current_user_from_token(authorization)
# Check permission for agent logs
has_permission = await check_user_permission(current_user["id"], "agents", "logs")
if not has_permission:
raise HTTPException(
status_code=403,
detail="Insufficient permissions: agents.logs required"
)
conn = await get_database_connection()
# Get agent to validate it exists and get cluster for access control
agent = await conn.fetchrow("""
SELECT a.id, a.pool_id, hc.id as cluster_id
FROM agents a
LEFT JOIN haproxy_clusters hc ON a.pool_id = hc.pool_id
WHERE a.name = $1
""", agent_name)
if not agent:
await close_database_connection(conn)
raise HTTPException(status_code=404, detail=f"Agent {agent_name} not found")
# Validate cluster access
if agent['cluster_id']:
await validate_user_cluster_access(current_user['id'], agent['cluster_id'], conn)
# Get activity logs
logs = await conn.fetch("""
SELECT
id,
action_type,
action_details,
timestamp,
created_at
FROM agent_activity_logs
WHERE agent_name = $1
ORDER BY timestamp DESC
LIMIT $2
""", agent_name, limit)
await close_database_connection(conn)
# Format logs
activity_logs = []
for log in logs:
activity_logs.append({
"id": log['id'],
"action_type": log['action_type'],
"action_details": json.loads(log['action_details']) if log['action_details'] else {},
"timestamp": log['timestamp'].isoformat().replace('+00:00', 'Z') if log['timestamp'] else None,
"created_at": log['created_at'].isoformat().replace('+00:00', 'Z') if log['created_at'] else None
})
return {
"agent_name": agent_name,
"logs": activity_logs,
"total": len(activity_logs)
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error getting activity logs for {agent_name}: {e}")
raise HTTPException(status_code=500, detail=str(e))
# ==== HELPER FUNCTION FOR SCRIPT UPDATES ====
async def update_agent_script_in_database(conn, platform: str, version: str, script_content: str, changelog: list = None, source_file_hash: str = None):
"""
Helper function to update agent script in database
Used by both Edit and Reset to Default operations
"""
# Deactivate all existing versions for this platform in agent_script_templates
await conn.execute("""
UPDATE agent_script_templates
SET is_active = false
WHERE platform = $1
""", platform)
# Insert/update version with script content in agent_script_templates
await conn.execute("""
INSERT INTO agent_script_templates (platform, version, script_content, source_file_hash, is_active)
VALUES ($1, $2, $3, $4, true)
ON CONFLICT (platform, version) DO UPDATE SET
script_content = EXCLUDED.script_content,
source_file_hash = COALESCE(EXCLUDED.source_file_hash, agent_script_templates.source_file_hash),
is_active = true,
updated_at = CURRENT_TIMESTAMP
""", platform, version, script_content, source_file_hash)
# CRITICAL: Also update agent_versions table for UI to display correct Available version
# Deactivate all existing versions in agent_versions
await conn.execute("""
UPDATE agent_versions
SET is_active = false
WHERE platform = $1
""", platform)
# Insert/update version in agent_versions table
# CRITICAL: ON CONFLICT must update both is_active and updated_at
# updated_at is used for version ordering in get_version_for_platform
if changelog is None:
changelog = ['Script template updated']
await conn.execute("""
INSERT INTO agent_versions (platform, version, release_date, changelog, is_active, created_at, updated_at)
VALUES ($1, $2, CURRENT_DATE, $3, true, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
ON CONFLICT (platform, version) DO UPDATE SET
is_active = true,
updated_at = CURRENT_TIMESTAMP,
changelog = EXCLUDED.changelog,
release_date = CURRENT_DATE
""", platform, version, changelog)
# Also update global storage for backward compatibility
global AGENT_VERSIONS
AGENT_VERSIONS[platform] = version
@router.post("/sync-scripts-from-files")
async def sync_scripts_from_files(
target_version: Optional[str] = None,
authorization: str = Header(None),
current_user: dict = Depends(authenticate_user)
):
"""
Force sync agent scripts from files to database
This will UPDATE database with latest file versions using the same Edit logic
Use cases:
- After updating linux_install.sh or macos_install.sh in code
- To reset scripts to file-based versions (1.0.0)
- To apply bug fixes from code to database
"""
try:
# Admin check - include super_admin and admin roles from user_roles table
conn = await get_database_connection()
# Check if user has admin role in database (consistent with config.py)
user_roles = await conn.fetch("""
SELECT r.name
FROM user_roles ur
JOIN roles r ON ur.role_id = r.id
WHERE ur.user_id = $1 AND r.is_active = true
""", current_user["id"])
role_names = [role['name'] for role in user_roles]
is_admin = (
current_user.get("role") in ["super_admin", "admin"] or
"super_admin" in role_names or
"admin" in role_names
)
if not is_admin:
await close_database_connection(conn)
logger.warning(f"User {current_user.get('username')} (role: {current_user.get('role')}, roles: {role_names}) denied access to sync scripts - admin access required")
raise HTTPException(status_code=403, detail="Admin access required")
logger.info(f"User {current_user.get('username')} (role: {current_user.get('role')}, roles: {role_names}) initiating agent script sync from files")
script_files = {
'linux': 'linux_install.sh',
'macos': 'macos_install.sh'
}
sync_results = []
for platform_key, filename in script_files.items():
# Read file from project
script_path = os.path.join(os.path.dirname(__file__), '..', 'utils', 'agent_scripts', filename)
if not os.path.exists(script_path):
sync_results.append({
"platform": platform_key,
"status": "error",
"message": f"File not found: {filename}"
})
continue
with open(script_path, 'r') as f:
file_content = f.read()
# CRITICAL DEBUG: Log file content size and check for "local listen_lines"
has_bug = "local listen_lines" in file_content
logger.info(f"FILE READ: {filename} -> {len(file_content)} bytes, has_bug={has_bug}, path={script_path}")
# Determine version to use
if target_version:
# User explicitly specified version (from UI)
new_version = target_version
logger.info(f"USER VERSION: Using specified version {new_version} for {platform_key}")
else:
# Auto-increment current version to trigger upgrades
current_version_row = await conn.fetchrow("""
SELECT version FROM agent_versions
WHERE platform = $1 AND is_active = true
ORDER BY updated_at DESC LIMIT 1
""", platform_key)
if current_version_row:
current_version = current_version_row['version']
# Auto-increment patch version (e.g., 1.0.0 → 1.0.1)
try:
parts = current_version.split('.')
major, minor, patch = int(parts[0]), int(parts[1]), int(parts[2])
new_version = f"{major}.{minor}.{patch + 1}"
logger.info(f"AUTO VERSION BUMP: {current_version}{new_version} for {platform_key}")
except:
# Fallback if version format is unexpected
new_version = "1.0.1"
else:
# First time sync
new_version = "1.0.0"
# Use the same update logic as Edit operation
file_hash = hashlib.sha256(file_content.encode()).hexdigest()
await update_agent_script_in_database(
conn=conn,
platform=platform_key,
version=new_version,
script_content=file_content,
source_file_hash=file_hash,
changelog=["Reset to file-based defaults", "Original stable version"]
)
sync_results.append({
"platform": platform_key,
"status": "success",
"version": new_version,
"message": f"Synced {filename} to database as version {new_version}"
})
logger.info(f"SCRIPT SYNC: Synced {platform_key} from file to database (version {new_version})")
# Log activity
await log_user_activity(
user_id=current_user["id"],
action='reset_to_default',
resource_type='agent_script_template',
resource_id=f"{platform_key}-{new_version}",
details={
'platform': platform_key,
'version': new_version,
'script_size': len(file_content),
'source': 'project_file'
}
)
await close_database_connection(conn)
return {
"status": "success",
"message": "Scripts synced from files to database",
"results": sync_results
}
except HTTPException:
raise
except Exception as e:
logger.error(f"Error syncing scripts from files: {e}")
raise HTTPException(status_code=500, detail=str(e))