mirror of
https://github.com/taylanbakircioglu/haproxy-openmanager.git
synced 2026-09-16 23:55:13 +00:00
0ee227363e
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.
3527 lines
171 KiB
Python
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))
|