Initial commit: core-api service extraction from portainer-core
Build and Push / build (release) Successful in 43s

This commit is contained in:
2025-12-11 15:52:59 +01:00
commit 488a4e8a91
49 changed files with 9478 additions and 0 deletions
+5
View File
@@ -0,0 +1,5 @@
"""
API Clients package for Core-API
Provides HTTP/WebSocket clients for external infrastructure services.
"""
+197
View File
@@ -0,0 +1,197 @@
"""
Core-AI HTTP Client
Provides interface to Core-AI service for AI performance metrics.
"""
import httpx
from typing import Optional, Dict, List, Any
from src.logging_config import get_logger
from src.config import get_settings
logger = get_logger(__name__)
settings = get_settings()
class CoreAIClient:
"""
HTTP client for Core-AI service
Provides access to AI performance metrics, tool execution stats,
and memory system monitoring.
"""
def __init__(
self,
base_url: Optional[str] = None,
timeout: int = 10
):
"""
Initialize Core-AI client
Args:
base_url: Core-AI base URL (default from settings)
timeout: Request timeout in seconds
"""
self.base_url = (base_url or getattr(settings, 'core_ai_base_url', 'http://core-ai:8086')).rstrip("/")
self.timeout = timeout
self.client = httpx.AsyncClient(timeout=self.timeout)
async def close(self):
"""Close the HTTP client"""
await self.client.aclose()
async def health_check(self) -> bool:
"""
Check if Core-AI service is accessible
Returns:
True if accessible, False otherwise
"""
try:
response = await self.client.get(f"{self.base_url}/health")
return response.status_code == 200
except Exception as e:
logger.error(f"Core-AI health check failed: {e}")
return False
async def get_metrics(self) -> Dict[str, Any]:
"""
Get comprehensive AI performance metrics
Returns:
Dict with agent performance, tool execution, memory stats
Example:
{
"uptime_seconds": 3600,
"timestamp": "2025-12-03T20:00:00Z",
"agent": {
"total_requests": 100,
"avg_response_time_ms": 1250.5,
"p95_response_time_ms": 3200.0,
...
},
"tools": {
"total_calls": 250,
"success_rate": 0.98,
"top_tools": {...}
},
"memory": {
"tier1_hit_rate": 0.85,
...
},
...
}
"""
try:
response = await self.client.get(f"{self.base_url}/metrics")
response.raise_for_status()
return response.json()
except httpx.HTTPStatusError as e:
logger.error(f"Failed to get metrics: HTTP {e.response.status_code}")
raise
except Exception as e:
logger.error(f"Failed to get metrics: {e}")
raise
async def get_recent_errors(self, limit: int = 20) -> List[Dict[str, Any]]:
"""
Get recent request errors
Args:
limit: Maximum number of errors to return
Returns:
List of error records with timestamps
Example:
[
{
"timestamp": "2025-12-03T19:45:12Z",
"agent_type": "pydantic",
"error": "Connection timeout",
"duration_ms": 5000
},
...
]
"""
try:
response = await self.client.get(
f"{self.base_url}/metrics/errors",
params={"limit": limit}
)
response.raise_for_status()
data = response.json()
return data.get("errors", [])
except Exception as e:
logger.error(f"Failed to get recent errors: {e}")
raise
async def get_tool_failures(self, limit: int = 20) -> List[Dict[str, Any]]:
"""
Get recent tool execution failures
Args:
limit: Maximum number of failures to return
Returns:
List of tool failure records
Example:
[
{
"timestamp": "2025-12-03T19:50:30Z",
"tool_name": "list_containers",
"error": "Connection refused",
"duration_ms": 150
},
...
]
"""
try:
response = await self.client.get(
f"{self.base_url}/metrics/tool-failures",
params={"limit": limit}
)
response.raise_for_status()
data = response.json()
return data.get("failures", [])
except Exception as e:
logger.error(f"Failed to get tool failures: {e}")
raise
async def reset_metrics(self) -> bool:
"""
Reset all metrics (admin operation)
Returns:
True if successful
"""
try:
response = await self.client.post(f"{self.base_url}/metrics/reset")
response.raise_for_status()
logger.info("Successfully reset Core-AI metrics")
return True
except Exception as e:
logger.error(f"Failed to reset metrics: {e}")
raise
async def __aenter__(self):
"""Async context manager entry"""
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
"""Async context manager exit"""
await self.close()
# Singleton instance
_ai_client: Optional[CoreAIClient] = None
def get_ai_client() -> CoreAIClient:
"""Get singleton Core-AI client instance"""
global _ai_client
if _ai_client is None:
_ai_client = CoreAIClient()
return _ai_client
+302
View File
@@ -0,0 +1,302 @@
"""
Authentik API Client
Provides methods for interacting with Authentik Identity Provider API.
Used for managing applications, providers, and authentication flows.
"""
import httpx
from typing import Dict, List, Any, Optional
from functools import lru_cache
from src.logging_config import get_logger
logger = get_logger(__name__)
class AuthentikClient:
"""Client for Authentik API operations"""
def __init__(self, base_url: str, api_token: str):
"""
Initialize Authentik client
Args:
base_url: Authentik base URL (e.g., http://authentik-server:9000)
api_token: API token for authentication
"""
self.base_url = base_url.rstrip('/')
self.api_token = api_token
self.client = httpx.AsyncClient(timeout=30.0)
async def _request(self, method: str, endpoint: str, **kwargs) -> Dict:
"""Make authenticated API request using token auth"""
headers = kwargs.pop("headers", {})
headers["Authorization"] = f"Bearer {self.api_token}"
response = await self.client.request(
method,
f"{self.base_url}/api/v3/{endpoint.lstrip('/')}",
headers=headers,
**kwargs
)
if not response.is_success:
logger.error(f"API request failed: {response.status_code}")
logger.error(f"Response body: {response.text}")
response.raise_for_status()
return response.json()
async def health_check(self) -> bool:
"""Check if Authentik is accessible"""
try:
response = await self.client.get(f"{self.base_url}/-/health/live/")
return response.status_code == 200
except Exception as e:
logger.error(f"Authentik health check failed: {e}")
return False
async def create_oauth2_provider(
self,
name: str,
client_id: str,
redirect_uris: List[str],
authorization_flow_slug: str = "default-provider-authorization-implicit-consent",
signing_key: Optional[str] = None
) -> Dict:
"""
Create an OAuth2/OIDC provider
Args:
name: Provider name
client_id: OAuth2 client ID
redirect_uris: List of allowed redirect URIs
authorization_flow_slug: Authorization flow slug (will be resolved to UUID)
signing_key: Signing key UUID (defaults to auto-selected)
Returns:
Created provider data including client_secret
"""
# Get authorization flow UUID from slug
flows = await self.list_flows()
auth_flow_uuid = None
invalidation_flow_uuid = None
for flow in flows:
if flow.get("slug") == authorization_flow_slug:
auth_flow_uuid = flow.get("pk")
if flow.get("slug") == "default-provider-invalidation-flow":
invalidation_flow_uuid = flow.get("pk")
if not auth_flow_uuid:
raise ValueError(f"Authorization flow '{authorization_flow_slug}' not found")
if not invalidation_flow_uuid:
raise ValueError("Invalidation flow not found")
# Get signing key if not provided
if not signing_key:
keys = await self._request("GET", "crypto/certificatekeypairs/")
# Find the self-signed cert
for key in keys.get("results", []):
if "authentik" in key.get("name", "").lower():
signing_key = key.get("pk")
break
if not signing_key and keys.get("results"):
signing_key = keys["results"][0]["pk"]
# Format redirect URIs as objects with matching_mode
formatted_redirect_uris = [
{"url": uri, "matching_mode": "strict"}
for uri in redirect_uris
]
provider_data = {
"name": name,
"authorization_flow": auth_flow_uuid,
"invalidation_flow": invalidation_flow_uuid,
"client_type": "confidential",
"client_id": client_id,
"redirect_uris": formatted_redirect_uris,
"signing_key": signing_key,
"sub_mode": "hashed_user_id",
"include_claims_in_id_token": True,
"issuer_mode": "per_provider",
"access_token_validity": "minutes=60",
"refresh_token_validity": "days=30",
"property_mappings": [] # Will use default mappings
}
result = await self._request("POST", "providers/oauth2/", json=provider_data)
logger.info(f"Created OAuth2 provider: {name} (ID: {result.get('pk')})")
return result
async def create_application(
self,
name: str,
slug: str,
provider_pk: int,
launch_url: Optional[str] = None,
icon_url: Optional[str] = None
) -> Dict:
"""
Create an application
Args:
name: Application display name
slug: Application slug (URL-safe identifier)
provider_pk: Primary key of the provider to use
launch_url: Optional launch URL
icon_url: Optional icon URL
Returns:
Created application data
"""
app_data = {
"name": name,
"slug": slug,
"provider": provider_pk,
"meta_launch_url": launch_url or "",
"meta_icon": icon_url or "",
"policy_engine_mode": "any",
"open_in_new_tab": False
}
result = await self._request("POST", "core/applications/", json=app_data)
logger.info(f"Created application: {name} (slug: {slug})")
return result
async def get_provider_by_name(self, name: str) -> Optional[Dict]:
"""Get OAuth2 provider by name"""
providers = await self._request("GET", "providers/oauth2/", params={"name": name})
results = providers.get("results", [])
return results[0] if results else None
async def get_application_by_slug(self, slug: str) -> Optional[Dict]:
"""Get application by slug"""
apps = await self._request("GET", "core/applications/", params={"slug": slug})
results = apps.get("results", [])
return results[0] if results else None
async def list_flows(self) -> List[Dict]:
"""List all authentication flows"""
result = await self._request("GET", "flows/instances/")
return result.get("results", [])
async def create_proxy_provider(
self,
name: str,
external_host: str,
authorization_flow_slug: str = "default-provider-authorization-implicit-consent",
mode: str = "forward_single",
token_validity: int = 480 # 8 hours in minutes
) -> Dict:
"""
Create a Proxy Provider for forward authentication
Args:
name: Provider name
external_host: External URL (e.g., https://auth.schweitz.net)
authorization_flow_slug: Authorization flow slug
mode: Proxy mode (forward_single for forward auth)
token_validity: Token validity in minutes (default: 480 = 8 hours)
Returns:
Created provider data
"""
# Get authorization flow UUID from slug
flows = await self.list_flows()
auth_flow_uuid = None
invalidation_flow_uuid = None
for flow in flows:
if flow.get("slug") == authorization_flow_slug:
auth_flow_uuid = flow.get("pk")
if flow.get("slug") == "default-provider-invalidation-flow":
invalidation_flow_uuid = flow.get("pk")
if not auth_flow_uuid:
raise ValueError(f"Authorization flow '{authorization_flow_slug}' not found")
if not invalidation_flow_uuid:
raise ValueError("Invalidation flow not found")
provider_data = {
"name": name,
"authorization_flow": auth_flow_uuid,
"invalidation_flow": invalidation_flow_uuid,
"mode": mode,
"external_host": external_host,
"access_token_validity": f"minutes={token_validity}",
"refresh_token_validity": f"minutes={token_validity}",
"session_duration": f"seconds={token_validity * 60}",
"cookie_domain": "", # Will use the domain of each proxied site
"property_mappings": []
}
result = await self._request("POST", "providers/proxy/", json=provider_data)
logger.info(f"Created Proxy provider: {name} (ID: {result.get('pk')})")
return result
async def get_provider_by_name_proxy(self, name: str) -> Optional[Dict]:
"""Get Proxy provider by name"""
providers = await self._request("GET", "providers/proxy/", params={"name": name})
results = providers.get("results", [])
return results[0] if results else None
async def create_outpost(
self,
name: str,
type: str,
providers: List[int],
config: Optional[Dict] = None
) -> Dict:
"""
Create an Authentik Outpost
Args:
name: Outpost name
type: Outpost type (e.g., "proxy")
providers: List of provider PKs
config: Optional configuration overrides
Returns:
Created outpost data
"""
outpost_data = {
"name": name,
"type": type,
"providers": providers,
"config": config or {},
"service_connection": None # Will use local Docker
}
result = await self._request("POST", "outposts/instances/", json=outpost_data)
logger.info(f"Created outpost: {name} (ID: {result.get('pk')})")
return result
async def get_outpost_by_name(self, name: str) -> Optional[Dict]:
"""Get outpost by name"""
outposts = await self._request("GET", "outposts/instances/", params={"name": name})
results = outposts.get("results", [])
return results[0] if results else None
async def close(self):
"""Close HTTP client"""
await self.client.aclose()
@lru_cache()
def get_authentik_client() -> AuthentikClient:
"""Get cached Authentik client instance"""
# Import credentials from gitignored module
try:
from src.credentials import AUTHENTIK_URL, AUTHENTIK_CORE_API_TOKEN
except ImportError:
# Fallback to environment variables if credentials.py doesn't exist
import os
AUTHENTIK_URL = os.getenv("AUTHENTIK_URL", "http://authentik-server:9000")
AUTHENTIK_CORE_API_TOKEN = os.getenv("AUTHENTIK_API_TOKEN", "")
return AuthentikClient(
base_url=AUTHENTIK_URL,
api_token=AUTHENTIK_CORE_API_TOKEN
)
+561
View File
@@ -0,0 +1,561 @@
"""
Uptime Kuma Socket.IO Client
Provides interface to Uptime Kuma via Socket.IO for monitor management.
Also provides metrics API access for real-time status data.
"""
import socketio
import asyncio
import httpx
import re
from typing import Optional, Dict, List, Any
from src.logging_config import get_logger
from src.config import get_settings
logger = get_logger(__name__)
settings = get_settings()
class KumaClient:
"""
Socket.IO client for Uptime Kuma
Uses Socket.IO for real-time communication with Uptime Kuma.
"""
def __init__(
self,
base_url: Optional[str] = None,
username: Optional[str] = None,
password: Optional[str] = None,
timeout: int = 30
):
"""
Initialize Kuma client
Args:
base_url: Kuma base URL (default from settings)
username: Kuma username (default from settings)
password: Kuma password (default from settings)
timeout: Request timeout in seconds
"""
self.base_url = (base_url or settings.kuma_url).rstrip("/")
self.username = username or settings.kuma_username
self.password = password or settings.kuma_password
self.timeout = timeout
self.sio = socketio.AsyncClient(
reconnection=True,
reconnection_attempts=3,
reconnection_delay=1,
)
self._connected = False
self._authenticated = False
self._monitors_cache: Dict[int, Dict[str, Any]] = {}
if not self.username or not self.password:
logger.warning("Uptime Kuma credentials not configured")
async def _ensure_connected(self):
"""Ensure we have an active connection and authentication"""
if not self._connected:
await self.connect()
if not self._authenticated:
await self.login()
async def connect(self):
"""Connect to Uptime Kuma Socket.IO server"""
if self._connected:
return
try:
await self.sio.connect(self.base_url, transports=['websocket'])
self._connected = True
logger.info(f"Connected to Uptime Kuma at {self.base_url}")
except Exception as e:
logger.error(f"Failed to connect to Uptime Kuma: {e}")
raise
async def disconnect(self):
"""Disconnect from Uptime Kuma"""
if self._connected:
await self.sio.disconnect()
self._connected = False
self._authenticated = False
logger.info("Disconnected from Uptime Kuma")
async def login(self):
"""Authenticate with Uptime Kuma"""
if not self._connected:
await self.connect()
try:
# Uptime Kuma login event
login_response = await self.sio.call(
'login',
{
'username': self.username,
'password': self.password,
'token': None
},
timeout=self.timeout
)
if login_response and login_response.get('ok'):
self._authenticated = True
logger.info("Successfully authenticated with Uptime Kuma")
else:
error_msg = login_response.get('msg', 'Unknown error') if login_response else 'No response'
raise Exception(f"Login failed: {error_msg}")
except Exception as e:
logger.error(f"Failed to authenticate with Uptime Kuma: {e}")
raise
async def health_check(self) -> bool:
"""
Check if Uptime Kuma is accessible
Returns:
True if accessible, False otherwise
"""
try:
await self._ensure_connected()
return self._authenticated
except Exception as e:
logger.error(f"Uptime Kuma health check failed: {e}")
return False
async def get_monitors(self) -> List[Dict[str, Any]]:
"""
List all monitors with uptime data
Returns:
List of monitor configurations with uptime_24h field
"""
await self._ensure_connected()
try:
# Storage for monitor list and uptime data received via events
monitor_list_data = {}
uptime_list_data = {}
monitor_event_received = asyncio.Event()
uptime_event_received = asyncio.Event()
# Register event handler for monitorList
@self.sio.event
async def monitorList(data):
nonlocal monitor_list_data
monitor_list_data = data
monitor_event_received.set()
# Register event handler for uptimeList (24h uptime percentages)
@self.sio.event
async def uptimeList(monitor_id, uptime_data):
nonlocal uptime_list_data
# uptime_data is typically a dict with time periods: {"24": 99.5, "720": 98.2, ...}
uptime_list_data[str(monitor_id)] = uptime_data
# Don't set event here as we'll get multiple calls
# Request monitor list - this triggers the server to send monitorList event
response = await self.sio.call('getMonitorList', timeout=self.timeout)
logger.info(f"getMonitorList call response: {response}")
# Wait for the monitorList event (with timeout)
try:
await asyncio.wait_for(monitor_event_received.wait(), timeout=5.0)
logger.info(f"Received monitorList event with {len(monitor_list_data)} items")
# Give time for uptimeList events to arrive
await asyncio.sleep(0.5)
logger.info(f"Received uptime data for {len(uptime_list_data)} monitors")
except asyncio.TimeoutError:
logger.warning("Timeout waiting for monitorList event")
# Process the monitor list data
if monitor_list_data and isinstance(monitor_list_data, dict):
monitors = []
for monitor_id, monitor_data in monitor_list_data.items():
if isinstance(monitor_data, dict):
monitor_data['id'] = int(monitor_id)
# Add uptime data if available
uptime_info = uptime_list_data.get(str(monitor_id), {})
if isinstance(uptime_info, dict):
# Uptime Kuma provides 24h uptime as key "24"
monitor_data['uptime_24h'] = float(uptime_info.get('24', 0))
else:
monitor_data['uptime_24h'] = 0.0
monitors.append(monitor_data)
self._monitors_cache[int(monitor_id)] = monitor_data
logger.info(f"Found {len(monitors)} monitors total")
return monitors
logger.warning(f"No valid monitor data received")
return []
except Exception as e:
logger.error(f"Failed to get monitors: {e}", exc_info=True)
raise
async def get_monitor(self, monitor_id: int) -> Dict[str, Any]:
"""
Get details of a specific monitor
Args:
monitor_id: Monitor identifier
Returns:
Monitor configuration details
"""
await self._ensure_connected()
try:
response = await self.sio.call('getMonitor', monitor_id, timeout=self.timeout)
if response:
self._monitors_cache[monitor_id] = response
return response
raise Exception(f"Monitor {monitor_id} not found")
except Exception as e:
logger.error(f"Failed to get monitor {monitor_id}: {e}")
raise
async def find_monitor_by_name(self, name: str) -> Optional[Dict[str, Any]]:
"""
Find a monitor by its name (case-insensitive)
Args:
name: Monitor name to search for
Returns:
Monitor object if found, None otherwise
"""
monitors = await self.get_monitors()
name_lower = name.lower()
for monitor in monitors:
if monitor.get("name", "").lower() == name_lower:
return monitor
return None
async def find_monitors_by_tag(self, tag: str) -> List[Dict[str, Any]]:
"""
Find all monitors with a specific tag
Args:
tag: Tag name to search for
Returns:
List of monitors with the tag
"""
monitors = await self.get_monitors()
tagged_monitors = []
for monitor in monitors:
monitor_tags = monitor.get("tags", [])
if any(t.get("name", "").lower() == tag.lower() for t in monitor_tags):
tagged_monitors.append(monitor)
return tagged_monitors
async def pause_monitor(self, monitor_id: int) -> bool:
"""
Pause a monitor (disable monitoring)
Args:
monitor_id: Monitor identifier
Returns:
True if successful
"""
await self._ensure_connected()
try:
# Uptime Kuma pause event
response = await self.sio.call('pauseMonitor', monitor_id, timeout=self.timeout)
if response and response.get('ok'):
logger.info(f"Paused monitor {monitor_id}")
return True
error_msg = response.get('msg', 'Unknown error') if response else 'No response'
raise Exception(f"Failed to pause monitor: {error_msg}")
except Exception as e:
logger.error(f"Failed to pause monitor {monitor_id}: {e}")
raise
async def resume_monitor(self, monitor_id: int) -> bool:
"""
Resume a monitor (enable monitoring)
Args:
monitor_id: Monitor identifier
Returns:
True if successful
"""
await self._ensure_connected()
try:
# Uptime Kuma resume event
response = await self.sio.call('resumeMonitor', monitor_id, timeout=self.timeout)
if response and response.get('ok'):
logger.info(f"Resumed monitor {monitor_id}")
return True
error_msg = response.get('msg', 'Unknown error') if response else 'No response'
raise Exception(f"Failed to resume monitor: {error_msg}")
except Exception as e:
logger.error(f"Failed to resume monitor {monitor_id}: {e}")
raise
async def pause_monitor_by_name(self, name: str) -> bool:
"""
Pause a monitor by its name
Args:
name: Monitor name
Returns:
True if successful, False if monitor not found
"""
monitor = await self.find_monitor_by_name(name)
if not monitor:
logger.warning(f"Monitor '{name}' not found")
return False
await self.pause_monitor(monitor["id"])
return True
async def resume_monitor_by_name(self, name: str) -> bool:
"""
Resume a monitor by its name
Args:
name: Monitor name
Returns:
True if successful, False if monitor not found
"""
monitor = await self.find_monitor_by_name(name)
if not monitor:
logger.warning(f"Monitor '{name}' not found")
return False
await self.resume_monitor(monitor["id"])
return True
async def add_monitor(self, monitor_config: Dict[str, Any]) -> Dict[str, Any]:
"""
Create a new monitor
Args:
monitor_config: Monitor configuration dict
Returns:
Created monitor details including ID
"""
await self._ensure_connected()
try:
# Uptime Kuma add monitor event
response = await self.sio.call('add', monitor_config, timeout=self.timeout)
if response and response.get('ok'):
monitor_id = response.get('monitorID')
logger.info(f"Created monitor '{monitor_config.get('name')}' with ID {monitor_id}")
# Get full monitor details
monitor = await self.get_monitor(monitor_id)
return monitor
error_msg = response.get('msg', 'Unknown error') if response else 'No response'
raise Exception(f"Failed to create monitor: {error_msg}")
except Exception as e:
logger.error(f"Failed to create monitor '{monitor_config.get('name')}': {e}")
raise
async def update_monitor(self, monitor_id: int, monitor_config: Dict[str, Any]) -> Dict[str, Any]:
"""
Update an existing monitor
Args:
monitor_id: Monitor identifier
monitor_config: Updated monitor configuration
Returns:
Updated monitor details
"""
await self._ensure_connected()
try:
# Ensure ID is in the config
monitor_config['id'] = monitor_id
# Uptime Kuma edit monitor event
response = await self.sio.call('editMonitor', monitor_config, timeout=self.timeout)
if response and response.get('ok'):
logger.info(f"Updated monitor {monitor_id}")
# Get updated monitor details
monitor = await self.get_monitor(monitor_id)
return monitor
error_msg = response.get('msg', 'Unknown error') if response else 'No response'
raise Exception(f"Failed to update monitor: {error_msg}")
except Exception as e:
logger.error(f"Failed to update monitor {monitor_id}: {e}")
raise
async def delete_monitor(self, monitor_id: int) -> bool:
"""
Delete a monitor
Args:
monitor_id: Monitor identifier
Returns:
True if successful
"""
await self._ensure_connected()
try:
# Uptime Kuma delete monitor event
response = await self.sio.call('deleteMonitor', monitor_id, timeout=self.timeout)
if response and response.get('ok'):
logger.info(f"Deleted monitor {monitor_id}")
# Remove from cache
self._monitors_cache.pop(monitor_id, None)
return True
error_msg = response.get('msg', 'Unknown error') if response else 'No response'
raise Exception(f"Failed to delete monitor: {error_msg}")
except Exception as e:
logger.error(f"Failed to delete monitor {monitor_id}: {e}")
raise
async def delete_monitor_by_name(self, name: str) -> bool:
"""
Delete a monitor by its name
Args:
name: Monitor name
Returns:
True if successful, False if monitor not found
"""
monitor = await self.find_monitor_by_name(name)
if not monitor:
logger.warning(f"Monitor '{name}' not found")
return False
await self.delete_monitor(monitor["id"])
return True
async def get_metrics_status(self) -> Dict[str, Dict[str, Any]]:
"""
Get monitor status from Prometheus metrics endpoint
This is simpler and more reliable than Socket.IO for getting current status.
Returns real-time UP/DOWN status but not historical uptime percentages.
Returns:
Dict mapping monitor names to status info:
{
"Portainer": {
"status": 1, # 1=UP, 0=DOWN, 2=PENDING, 3=MAINTENANCE
"response_time": 5, # ms
"monitor_type": "http",
"url": "http://192.168.86.149:8001"
},
...
}
"""
try:
# Use API key authentication
api_key = settings.kuma_api_key
if not api_key:
logger.warning("Kuma API key not configured")
return {}
# Fetch metrics with HTTP Basic Auth (empty username, API key as password)
async with httpx.AsyncClient(timeout=10.0) as client:
response = await client.get(
f"{self.base_url}/metrics",
auth=("", api_key)
)
response.raise_for_status()
metrics_text = response.text
# Parse Prometheus format metrics
# Format: metric_name{label1="value1",label2="value2"} value
monitor_data = {}
# Parse monitor_status lines
status_pattern = r'monitor_status\{monitor_name="([^"]+)",.*?\} (\d+)'
for match in re.finditer(status_pattern, metrics_text):
monitor_name = match.group(1)
status = int(match.group(2))
if monitor_name not in monitor_data:
monitor_data[monitor_name] = {}
monitor_data[monitor_name]['status'] = status
# Parse monitor_response_time lines
response_pattern = r'monitor_response_time\{monitor_name="([^"]+)",monitor_type="([^"]+)",monitor_url="([^"]+)",.*?\} ([\d.]+)'
for match in re.finditer(response_pattern, metrics_text):
monitor_name = match.group(1)
monitor_type = match.group(2)
monitor_url = match.group(3)
response_time = float(match.group(4))
if monitor_name not in monitor_data:
monitor_data[monitor_name] = {}
monitor_data[monitor_name].update({
'response_time': response_time,
'monitor_type': monitor_type,
'url': monitor_url
})
logger.info(f"Fetched metrics for {len(monitor_data)} monitors")
return monitor_data
except Exception as e:
logger.error(f"Failed to fetch metrics: {e}")
return {}
async def __aenter__(self):
"""Async context manager entry"""
await self._ensure_connected()
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
"""Async context manager exit"""
await self.disconnect()
# Singleton instance
_kuma_client: Optional[KumaClient] = None
def get_kuma_client() -> KumaClient:
"""Get singleton Kuma client instance"""
global _kuma_client
if _kuma_client is None:
_kuma_client = KumaClient()
return _kuma_client
+383
View File
@@ -0,0 +1,383 @@
"""
Nginx Proxy Manager API Client
Provides interface to NPM REST API for proxy host and SSL certificate management.
"""
import httpx
from typing import Optional, Dict, List, Any
from datetime import datetime, timedelta
from src.logging_config import get_logger
from src.config import get_settings
logger = get_logger(__name__)
settings = get_settings()
class NPMClient:
"""
HTTP client for Nginx Proxy Manager API
Uses JWT Bearer token authentication with automatic token refresh.
Tokens expire after ~24 hours.
"""
def __init__(
self,
base_url: Optional[str] = None,
email: Optional[str] = None,
password: Optional[str] = None,
timeout: int = 30
):
"""
Initialize NPM client
Args:
base_url: NPM base URL (default from settings)
email: NPM admin email (default from settings)
password: NPM admin password (default from settings)
timeout: Request timeout in seconds
"""
self.base_url = (base_url or settings.npm_url).rstrip("/")
self.email = email or settings.npm_email
self.password = password or settings.npm_password
self.timeout = timeout
self._token: Optional[str] = None
self._token_expires: Optional[datetime] = None
if not self.email or not self.password:
logger.warning("NPM credentials not configured")
async def _ensure_token(self):
"""Ensure we have a valid token, refresh if needed"""
if self._token and self._token_expires:
# If token expires in less than 1 hour, refresh it
if datetime.now() + timedelta(hours=1) < self._token_expires:
return
# Get new token
await self._refresh_token()
async def _refresh_token(self):
"""Get a new authentication token"""
try:
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/tokens",
json={
"identity": self.email,
"secret": self.password
}
)
response.raise_for_status()
data = response.json()
self._token = data.get("token")
# Assume 23-hour expiration to be safe
self._token_expires = datetime.now() + timedelta(hours=23)
logger.info("NPM token refreshed successfully")
except Exception as e:
logger.error(f"Failed to refresh NPM token: {e}")
raise
def _get_headers(self) -> Dict[str, str]:
"""Get request headers with authentication"""
if not self._token:
raise RuntimeError("No NPM token available. Call _ensure_token() first.")
return {
"Authorization": f"Bearer {self._token}",
"Content-Type": "application/json"
}
async def health_check(self) -> bool:
"""
Check if NPM API is accessible
Returns:
True if accessible, False otherwise
"""
try:
async with httpx.AsyncClient(timeout=self.timeout, follow_redirects=True) as client:
response = await client.get(f"{self.base_url}/api")
# Accept any successful response (2xx) or redirect (3xx) as healthy
# A redirect indicates the service is up and responding
return 200 <= response.status_code < 400
except Exception as e:
logger.error(f"NPM health check failed: {e}")
return False
async def get_proxy_hosts(self) -> List[Dict[str, Any]]:
"""
List all proxy hosts
Returns:
List of proxy host configurations
"""
await self._ensure_token()
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/nginx/proxy-hosts",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def get_proxy_host(self, host_id: int) -> Dict[str, Any]:
"""
Get details of a specific proxy host
Args:
host_id: Proxy host identifier
Returns:
Proxy host configuration
"""
await self._ensure_token()
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/nginx/proxy-hosts/{host_id}",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def create_proxy_host(
self,
domain_names: List[str],
forward_host: str,
forward_port: int,
forward_scheme: str = "http",
certificate_id: int = 0,
ssl_forced: bool = False,
block_exploits: bool = True,
caching_enabled: bool = True,
websocket_upgrade: bool = True,
http2_support: bool = True,
hsts_enabled: bool = True,
advanced_config: str = ""
) -> Dict[str, Any]:
"""
Create a new proxy host
Args:
domain_names: List of domain names for this proxy
forward_host: Target host to proxy to
forward_port: Target port to proxy to
forward_scheme: http or https
certificate_id: SSL certificate ID (0 for none)
ssl_forced: Force HTTPS redirect
block_exploits: Enable exploit blocking
caching_enabled: Enable response caching
websocket_upgrade: Allow WebSocket upgrades
http2_support: Enable HTTP/2
hsts_enabled: Enable HSTS headers
advanced_config: Custom nginx configuration
Returns:
Created proxy host details
"""
await self._ensure_token()
payload = {
"domain_names": domain_names,
"forward_scheme": forward_scheme,
"forward_host": forward_host,
"forward_port": forward_port,
"certificate_id": certificate_id,
"ssl_forced": ssl_forced,
"block_exploits": block_exploits,
"caching_enabled": caching_enabled,
"allow_websocket_upgrade": websocket_upgrade,
"http2_support": http2_support,
"hsts_enabled": hsts_enabled,
"hsts_subdomains": False,
"advanced_config": advanced_config,
"access_list_id": 0,
"meta": {}
}
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/nginx/proxy-hosts",
headers=self._get_headers(),
json=payload
)
response.raise_for_status()
return response.json()
async def update_proxy_host(
self,
proxy_id: int,
config: Dict[str, Any]
) -> Dict[str, Any]:
"""
Update an existing proxy host configuration
Args:
proxy_id: Proxy host ID to update
config: Full proxy host configuration (get from get_proxy_host, modify, then update)
Returns:
Updated proxy host details
"""
await self._ensure_token()
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.put(
f"{self.base_url}/api/nginx/proxy-hosts/{proxy_id}",
headers=self._get_headers(),
json=config
)
if not response.is_success:
logger.error(f"Update failed: {response.status_code}")
logger.error(f"Response: {response.text}")
response.raise_for_status()
return response.json()
async def enable_authentik_forward_auth(
self,
proxy_id: int,
authentik_url: str = "http://authentik-server:9000"
) -> Dict[str, Any]:
"""
Enable Authentik forward authentication on a proxy host
Args:
proxy_id: Proxy host ID to update
authentik_url: Authentik server URL (default: http://authentik-server:9000)
Returns:
Updated proxy host details
"""
# Get current config
proxy_host = await self.get_proxy_host(proxy_id)
# Authentik forward auth configuration
auth_config = f"""# Authentik Forward Authentication
# Send authentication requests to Authentik
auth_request /outpost.goauthentik.io/auth/nginx;
# Preserve authentication cookies
auth_request_set $auth_cookie $upstream_http_set_cookie;
add_header Set-Cookie $auth_cookie;
# Get user information from Authentik
auth_request_set $authentik_username $upstream_http_x_authentik_username;
auth_request_set $authentik_groups $upstream_http_x_authentik_groups;
auth_request_set $authentik_email $upstream_http_x_authentik_email;
auth_request_set $authentik_name $upstream_http_x_authentik_name;
auth_request_set $authentik_uid $upstream_http_x_authentik_uid;
# Pass user info to backend
proxy_set_header X-authentik-username $authentik_username;
proxy_set_header X-authentik-groups $authentik_groups;
proxy_set_header X-authentik-email $authentik_email;
proxy_set_header X-authentik-name $authentik_name;
proxy_set_header X-authentik-uid $authentik_uid;
# On authentication failure, redirect to Authentik login
error_page 401 = @authentik_proxy_signin;
location @authentik_proxy_signin {{
internal;
add_header Set-Cookie $auth_cookie;
return 302 /outpost.goauthentik.io/start?rd=$scheme://$http_host$request_uri;
}}
# Authentik authentication endpoint
location /outpost.goauthentik.io {{
proxy_pass {authentik_url}/outpost.goauthentik.io;
proxy_set_header X-Original-URL $scheme://$http_host$request_uri;
proxy_pass_request_body off;
proxy_set_header Content-Length "";
proxy_set_header Host $host;
}}
"""
# Update the advanced config
proxy_host["advanced_config"] = auth_config
# Remove read-only fields that NPM doesn't accept in updates
readonly_fields = [
"id", "created_on", "modified_on", "owner", "owner_user_id",
"certificate", "use_default_location", "ipv6", "meta", "nginx_online",
"nginx_err", "access_list", "certificate_id"
]
clean_config = {k: v for k, v in proxy_host.items() if k not in readonly_fields}
# Ensure locations is an array (required field)
if "locations" not in clean_config or clean_config["locations"] is None:
clean_config["locations"] = []
# Update the proxy host
return await self.update_proxy_host(proxy_id, clean_config)
async def get_certificates(self) -> List[Dict[str, Any]]:
"""
List all SSL certificates
Returns:
List of certificate details
"""
await self._ensure_token()
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/nginx/certificates",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def create_certificate(
self,
domain_names: List[str],
provider: str = "letsencrypt"
) -> Dict[str, Any]:
"""
Request a new SSL certificate from Let's Encrypt
Args:
domain_names: List of domains for the certificate
provider: Certificate provider (default: letsencrypt)
Returns:
Certificate details
"""
await self._ensure_token()
payload = {
"provider": provider,
"domain_names": domain_names,
"meta": {
"dns_challenge": False
}
}
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/nginx/certificates",
headers=self._get_headers(),
json=payload
)
response.raise_for_status()
return response.json()
# Singleton instance
_npm_client: Optional[NPMClient] = None
def get_npm_client() -> NPMClient:
"""Get singleton NPM client instance"""
global _npm_client
if _npm_client is None:
_npm_client = NPMClient()
return _npm_client
+450
View File
@@ -0,0 +1,450 @@
"""
Portainer API Client
Provides interface to Portainer REST API for stack and container management.
Includes fallback to Docker socket for containers not managed by Portainer.
"""
import httpx
import json
from typing import Optional, Dict, List, Any
from src.logging_config import get_logger
from src.config import get_settings
logger = get_logger(__name__)
settings = get_settings()
class PortainerClient:
"""
HTTP client for Portainer API
Uses access token authentication (X-API-Key header)
for long-lived API access without session management.
"""
def __init__(
self,
base_url: Optional[str] = None,
api_key: Optional[str] = None,
timeout: int = 30
):
"""
Initialize Portainer client
Args:
base_url: Portainer base URL (default from settings)
api_key: Portainer API access token (default from settings)
timeout: Request timeout in seconds
"""
self.base_url = (base_url or settings.portainer_url).rstrip("/")
self.api_key = api_key or settings.portainer_api_key
self.timeout = timeout
if not self.api_key:
logger.warning("Portainer API key not configured")
def _get_headers(self) -> Dict[str, str]:
"""Get request headers with authentication"""
return {
"X-API-Key": self.api_key,
"Content-Type": "application/json"
}
async def health_check(self) -> bool:
"""
Check if Portainer API is accessible
Returns:
True if accessible, False otherwise
"""
try:
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(f"{self.base_url}/api/status")
return response.status_code == 200
except Exception as e:
logger.error(f"Portainer health check failed: {e}")
return False
async def get_endpoints(self) -> List[Dict[str, Any]]:
"""
List all Portainer endpoints (Docker environments)
Returns:
List of endpoint configurations
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/endpoints",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def get_stacks(self, endpoint_id: Optional[int] = None) -> List[Dict[str, Any]]:
"""
List all stacks
Args:
endpoint_id: Filter by specific endpoint (optional)
Returns:
List of stack configurations
"""
params = {}
if endpoint_id:
params["endpointId"] = endpoint_id
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/stacks",
headers=self._get_headers(),
params=params
)
response.raise_for_status()
return response.json()
async def get_stack(self, stack_id: int) -> Dict[str, Any]:
"""
Get details of a specific stack
Args:
stack_id: Stack identifier
Returns:
Stack configuration details
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/stacks/{stack_id}",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def create_stack(
self,
name: str,
stack_file_content: str,
endpoint_id: int
) -> Dict[str, Any]:
"""
Create a new stack from compose file content
Args:
name: Stack name
stack_file_content: Docker Compose YAML content
endpoint_id: Portainer endpoint to deploy to
Returns:
Created stack details
"""
payload = {
"name": name,
"stackFileContent": stack_file_content
}
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/stacks/create/standalone/string",
headers=self._get_headers(),
params={"endpointId": endpoint_id},
json=payload
)
response.raise_for_status()
return response.json()
async def update_stack(
self,
stack_id: int,
stack_file_content: str,
endpoint_id: int,
prune: bool = False,
pull_image: bool = False
) -> Dict[str, Any]:
"""
Update an existing stack
Args:
stack_id: Stack identifier
stack_file_content: New Docker Compose YAML content
endpoint_id: Portainer endpoint
prune: Remove services no longer defined
pull_image: Pull latest images before deployment
Returns:
Updated stack details
"""
payload = {
"stackFileContent": stack_file_content,
"prune": prune,
"pullImage": pull_image
}
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.put(
f"{self.base_url}/api/stacks/{stack_id}",
headers=self._get_headers(),
params={"endpointId": endpoint_id},
json=payload
)
response.raise_for_status()
return response.json()
async def delete_stack(self, stack_id: int, endpoint_id: int) -> bool:
"""
Delete a stack
Args:
stack_id: Stack identifier
endpoint_id: Portainer endpoint
Returns:
True if successful
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.delete(
f"{self.base_url}/api/stacks/{stack_id}",
headers=self._get_headers(),
params={"endpointId": endpoint_id}
)
response.raise_for_status()
return True
async def get_containers(self, endpoint_id: int, all_containers: bool = True) -> List[Dict[str, Any]]:
"""
List containers on a specific endpoint
Args:
endpoint_id: Portainer endpoint identifier
all_containers: Include stopped containers (default: True)
Returns:
List of container details
"""
params = {"all": 1 if all_containers else 0}
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/endpoints/{endpoint_id}/docker/containers/json",
headers=self._get_headers(),
params=params
)
response.raise_for_status()
return response.json()
async def get_container(self, endpoint_id: int, container_id: str) -> Dict[str, Any]:
"""
Get detailed information about a specific container
Args:
endpoint_id: Portainer endpoint identifier
container_id: Container ID or name
Returns:
Container details including network and port information
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.get(
f"{self.base_url}/api/endpoints/{endpoint_id}/docker/containers/{container_id}/json",
headers=self._get_headers()
)
response.raise_for_status()
return response.json()
async def stop_container(self, endpoint_id: int, container_id: str) -> bool:
"""
Stop a container
Args:
endpoint_id: Portainer endpoint identifier
container_id: Container ID or name
Returns:
True if successful
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/endpoints/{endpoint_id}/docker/containers/{container_id}/stop",
headers=self._get_headers()
)
response.raise_for_status()
logger.info(f"Stopped container {container_id}")
return True
async def start_container(self, endpoint_id: int, container_id: str) -> bool:
"""
Start a container
Args:
endpoint_id: Portainer endpoint identifier
container_id: Container ID or name
Returns:
True if successful
"""
async with httpx.AsyncClient(timeout=self.timeout) as client:
response = await client.post(
f"{self.base_url}/api/endpoints/{endpoint_id}/docker/containers/{container_id}/start",
headers=self._get_headers()
)
response.raise_for_status()
logger.info(f"Started container {container_id}")
return True
# ========================================================================
# Docker Socket Fallback (for containers not managed by Portainer)
# ========================================================================
async def _list_containers_via_socket(self, all_containers: bool = True) -> List[Dict[str, Any]]:
"""
Fallback: List containers directly via Docker socket
Used when Portainer API doesn't return complete data (e.g., containers
started outside Portainer, AMP game servers, etc.)
Args:
all_containers: Include stopped containers
Returns:
List of container details in Docker API format
"""
try:
# Docker socket is mounted at /var/run/docker.sock
# Use httpx with unix socket transport
transport = httpx.AsyncHTTPTransport(uds="/var/run/docker.sock")
async with httpx.AsyncClient(transport=transport, timeout=10) as client:
params = {"all": 1 if all_containers else 0}
response = await client.get(
"http://localhost/v1.41/containers/json",
params=params
)
response.raise_for_status()
return response.json()
except Exception as e:
logger.warning(f"Docker socket fallback failed: {e}")
return []
async def _inspect_container_via_socket(self, container_id_or_name: str) -> Optional[Dict[str, Any]]:
"""
Fallback: Inspect container directly via Docker socket
Args:
container_id_or_name: Container ID or name
Returns:
Container details or None
"""
try:
transport = httpx.AsyncHTTPTransport(uds="/var/run/docker.sock")
async with httpx.AsyncClient(transport=transport, timeout=10) as client:
response = await client.get(
f"http://localhost/v1.41/containers/{container_id_or_name}/json"
)
response.raise_for_status()
return response.json()
except Exception as e:
logger.warning(f"Docker socket inspect fallback failed for '{container_id_or_name}': {e}")
return None
# ========================================================================
# Helper methods for agent tools (auto-detect endpoint + fallback)
# ========================================================================
async def list_containers(self, all_containers: bool = True) -> List[Dict[str, Any]]:
"""
List containers using auto-detected endpoint with Docker socket fallback
This is a convenience wrapper that automatically uses the first/default endpoint.
If Portainer doesn't have complete data, falls back to Docker socket.
Args:
all_containers: Include stopped containers (default: True)
Returns:
List of container details
"""
try:
# Try Portainer first
endpoints = await self.get_endpoints()
if endpoints:
endpoint_id = endpoints[0]["Id"]
containers = await self.get_containers(endpoint_id, all_containers)
if containers:
return containers
# Fallback to Docker socket
logger.info("Portainer returned no containers, trying Docker socket fallback...")
return await self._list_containers_via_socket(all_containers)
except Exception as e:
logger.error(f"Error listing containers: {e}")
# Try fallback even on exception
try:
return await self._list_containers_via_socket(all_containers)
except Exception as fallback_error:
logger.error(f"Fallback also failed: {fallback_error}")
return []
async def inspect_container(self, container_name: str) -> Optional[Dict[str, Any]]:
"""
Inspect a container by name using auto-detected endpoint with Docker socket fallback
This is a convenience wrapper that automatically uses the first/default endpoint.
If Portainer doesn't find the container, falls back to Docker socket.
Args:
container_name: Container name (e.g., "jellyfin", "ollama")
Returns:
Container details or None if not found
"""
try:
# Try Portainer first
endpoints = await self.get_endpoints()
if endpoints:
endpoint_id = endpoints[0]["Id"]
# First list all containers to find the one matching the name
all_containers = await self.get_containers(endpoint_id, all_containers=True)
matching_container = None
for container in all_containers:
# Container names come as array like ['/jellyfin']
names = container.get('Names', [])
for name in names:
clean_name = name.lstrip('/')
if clean_name == container_name or clean_name.lower() == container_name.lower():
matching_container = container
break
if matching_container:
break
if matching_container:
# Get detailed info using container ID
container_id = matching_container['Id']
return await self.get_container(endpoint_id, container_id)
# Not found in Portainer, try Docker socket fallback
logger.info(f"Container '{container_name}' not found in Portainer, trying Docker socket fallback...")
return await self._inspect_container_via_socket(container_name)
except Exception as e:
logger.error(f"Error inspecting container '{container_name}': {e}")
# Try fallback even on exception
try:
return await self._inspect_container_via_socket(container_name)
except Exception as fallback_error:
logger.error(f"Fallback also failed: {fallback_error}")
return None
# Singleton instance
_portainer_client: Optional[PortainerClient] = None
def get_portainer_client() -> PortainerClient:
"""Get singleton Portainer client instance"""
global _portainer_client
if _portainer_client is None:
_portainer_client = PortainerClient()
return _portainer_client