mirror of
https://github.com/PegaProx/project-pegaprox.git
synced 2026-08-12 15:27:47 +08:00
Builds on Phase 1's create+read layer with the lifecycle + drift half PDM lacks. Edit fan-out — PUT /api/multi-sdn/vnets/<id>: change alias + add/remove subnets, fanned out to every member (PVE vnet-alias PUT, subnet POST/DELETE with the <zone>-<cidr> id URL-encoded) + one apply per cluster. Structural fields (name/zone/vni/asn/controller) are immutable — a change there rebuilds the span, so it's rejected with a delete-and-recreate hint. Gated sdn.manage+admin.settings, per-member check_cluster_access, partial-failure status like create. Drift detect + reconcile: - Background scanner (6h, daemon thread started in api/__init__, mirrors drift.py) reads live per-member SDN state for each aggregate vnet and persists the drift status. DETECT-ONLY by default — it never writes to a cluster unless the new global setting multi_sdn_drift_reconcile is opted in. - Manual POST .../<id>/reconcile (re-assert desired + fix alias drift) and POST .../<id>/scan (read-only drift refresh, node.view) for on-demand use. - When auto-reconcile is opted in, it fixes ONLY the drifted members (never reloads an in-sync member), debounces 'missing' (needs two consecutive non-healthy passes so a transient partial read can't resurrect intentionally- removed SDN), and writes a log_audit trail for the unattended mutation. Frontend (MultiClusterEvpnView): edit modal (alias/subnets), per-cluster drift detail, Scan / Reconcile / Edit actions, and an admin-only auto-reconcile toggle. i18n all 7 langs. Adversarial review (4 dimensions x verify, 2 workflows) → 7 findings fixed: _apply_on_cluster now skips the cluster-wide apply when nothing changed (kills the over-broad blast radius: one drifted member no longer reloads the whole span, and reapply/reconcile of in-sync members is a no-op); reapply treats in_sync as done; auto-reconcile is per-member + missing-debounced + audited; the toggle is hidden from non-admins; edit subnet fold-back de-dupes; and the pre-existing DUPLICATE it: translation block (silently dropping 317 Italian keys) is merged into one — Italian is whole again. Verified in-process (no real same-ASN EVPN clusters — honest limit): 44 tests green across edit/reconcile/scan/scanner/gate/debounce/no-op-apply/subnet-dedup; app boots + 9 routes + scanner thread; build + translations parse; auth/injection review dimensions came back clean. Phase-2 auto-reconcile ships default-off.
571 lines
24 KiB
Python
571 lines
24 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""shared helpers for all api routes - split from monolith dec 2025, NS"""
|
|
|
|
import os
|
|
import json
|
|
import time
|
|
import logging
|
|
from datetime import datetime
|
|
|
|
from pegaprox.constants import (
|
|
SESSION_TIMEOUT, SERVER_SETTINGS_FILE,
|
|
LOGIN_MAX_ATTEMPTS, LOGIN_LOCKOUT_TIME, LOGIN_ATTEMPT_WINDOW,
|
|
TASK_USER_CACHE_TTL,
|
|
)
|
|
from pegaprox.globals import (
|
|
cluster_managers, active_sessions, users_db,
|
|
task_pegaprox_users_cache, task_pegaprox_users_lock,
|
|
)
|
|
from pegaprox.core.db import get_db
|
|
|
|
def effective_reverse_proxy(settings=None):
|
|
"""#614 — the frontend builds console (VNC/SSH) WebSocket URLs from
|
|
reverse_proxy_enabled. PEGAPROX_BEHIND_PROXY forces behind-proxy mode at boot
|
|
(app.py) but is never persisted, so report the OR of the persisted setting and
|
|
the env override — otherwise `PEGAPROX_BEHIND_PROXY=true` alone leaves consoles
|
|
dialing port+1/+2, which the reverse proxy can't route."""
|
|
if settings is None:
|
|
settings = load_server_settings()
|
|
return bool(settings.get('reverse_proxy_enabled', False)) or \
|
|
os.environ.get('PEGAPROX_BEHIND_PROXY', '').lower() in ('1', 'true', 'yes')
|
|
|
|
|
|
def load_server_settings():
|
|
"""Load server settings from SQLite database
|
|
|
|
SQLite migration
|
|
"""
|
|
defaults = {
|
|
'domain': '',
|
|
'port': 5000, # Web server port
|
|
'ssl_enabled': False,
|
|
# MK Jul 2026 — #612 Phase 2: auto-reconcile cross-cluster EVPN drift. OFF by
|
|
# default → the drift scanner is detect-only (never writes to a production SDN
|
|
# uninvited); opt in to let it re-push the desired definition to drifted members.
|
|
'multi_sdn_drift_reconcile': False,
|
|
# MK: Mar 2026 - ACME / Let's Encrypt auto-certs (#96)
|
|
'acme_enabled': False,
|
|
'acme_email': '',
|
|
'acme_staging': False, # use LE staging for testing
|
|
'acme_challenge_type': 'http-01',
|
|
'acme_dns_provider': 'manual',
|
|
'acme_dns_rfc2136_nameserver': '',
|
|
'acme_dns_rfc2136_port': 53,
|
|
'acme_dns_rfc2136_zone': '',
|
|
'acme_dns_rfc2136_key_name': '',
|
|
'acme_dns_rfc2136_secret': '',
|
|
'acme_dns_rfc2136_algorithm': 'hmac-sha512',
|
|
'acme_dns_rfc2136_ttl': 60,
|
|
'acme_dns_propagation_seconds': 30,
|
|
'logo_url': '',
|
|
'app_name': 'PegaProx',
|
|
# HTTP redirect port - NS Jan 2026
|
|
# Now that we have protocol detection on the main port, this is only needed
|
|
# if you want HTTP:80 → HTTPS:5000 redirect
|
|
# 0 = auto (80 if root, disabled otherwise), -1 = disabled, or specific port
|
|
'http_redirect_port': -1, # Disabled by default - protocol detection handles same-port redirect
|
|
# Brute force protection settings
|
|
'login_max_attempts': 5,
|
|
'login_lockout_time': 300, # 5 min
|
|
'login_attempt_window': 600, # 10 min
|
|
# Password policy settings
|
|
'password_min_length': 8,
|
|
'password_require_uppercase': True,
|
|
'password_require_lowercase': True,
|
|
'password_require_numbers': True,
|
|
'password_require_special': False, # too annoying for most users
|
|
# LW: Password expiry - Dec 2025
|
|
'password_expiry_enabled': False, # disabled by default
|
|
'password_expiry_days': 90, # days until password expires
|
|
'password_expiry_warning_days': 14, # warn this many days before
|
|
'password_expiry_email_enabled': True, # send email notifications
|
|
'password_expiry_include_admins': False, # MK: opt-in for admins, otherwise they could lock themselves out
|
|
# Session settings
|
|
'session_timeout': SESSION_TIMEOUT, # Use constant (8h HIPAA default)
|
|
# NS: SMTP Settings - Dec 2025
|
|
'smtp_enabled': False,
|
|
'smtp_host': '',
|
|
'smtp_port': 587,
|
|
'smtp_user': '',
|
|
'smtp_password': '', # stored encrypted ideally
|
|
'smtp_from_email': '',
|
|
'smtp_from_name': 'PegaProx Alerts',
|
|
'smtp_tls': True,
|
|
'smtp_ssl': False,
|
|
# Alert notification settings
|
|
'alert_email_recipients': [], # list of email addresses
|
|
'alert_cooldown': 300, # Don't send same alert within 5 min
|
|
# NS Apr 2026 (#331) — email notification when a new PegaProx release appears.
|
|
# Opt-in; re-uses alert_email_recipients. Dedupes via last-notified-version.
|
|
'alert_update_available': False,
|
|
'alert_last_notified_version': '',
|
|
# NS 2026-04-24 — when true, validate_session() invalidates a session if the
|
|
# source IP changes. Default off because mobile roaming / carrier NAT
|
|
# legitimately shifts IPs mid-session.
|
|
'strict_session_ip': False,
|
|
# MK Apr 2026 — when true, /api/metrics needs no auth. Useful for setups
|
|
# where a reverse proxy/mutual-TLS already gates scrapes. Default off.
|
|
'metrics_public': False,
|
|
# When enabled, the Syslog viewer only shows hostnames belonging to
|
|
# the currently selected cluster instead of all collected syslog rows.
|
|
'syslog_filter_by_selected_cluster': False,
|
|
# NS 2026-06-05 (audit N1): gate the syslog RECEIVER (UDP+TCP :1514).
|
|
# Default True preserves the always-on behaviour; set False to close the
|
|
# port on installs that don't ingest syslog. (DoS-safe either way now —
|
|
# ingestion is bounded-queue + batched off-hub.)
|
|
'syslog_enabled': True,
|
|
# NS 2026-06-05 (S1): retention for the syslog receiver DB (syslog.db). The
|
|
# receiver only INSERTs, so without a sweep it grows unbounded on the same
|
|
# volume as the main DB. Pruned ~hourly by the drain loop.
|
|
'syslog_retention_days': 30,
|
|
# Webhook alert channels (Slack, Discord, Teams, ntfy, generic)
|
|
# Each: {id, name, type, url, enabled, ...type-specific fields}
|
|
'alert_webhooks': [],
|
|
# IP Whitelisting - Jan 2026
|
|
'ip_whitelist_enabled': False,
|
|
'ip_whitelist': '', # Comma-separated IPs/CIDRs
|
|
'ip_blacklist': '', # Comma-separated IPs/CIDRs (always blocked)
|
|
# NS: Feb 2026 - LDAP defaults (must be here so get_ldap_settings always has values!)
|
|
# Without these, a partial save (e.g. only ldap_enabled=True) causes "LDAP not configured"
|
|
'ldap_enabled': False,
|
|
'ldap_server': '',
|
|
'ldap_port': 389,
|
|
'ldap_use_ssl': False,
|
|
'ldap_use_starttls': False,
|
|
'ldap_bind_dn': '',
|
|
'ldap_bind_password': '',
|
|
'ldap_base_dn': '',
|
|
'ldap_user_filter': '(&(objectClass=person)(sAMAccountName={username}))',
|
|
'ldap_username_attribute': 'sAMAccountName',
|
|
'ldap_email_attribute': 'mail',
|
|
'ldap_display_name_attribute': 'displayName',
|
|
'ldap_group_base_dn': '',
|
|
'ldap_group_filter': '(&(objectClass=group)(member={user_dn}))',
|
|
'ldap_admin_group': '',
|
|
'ldap_user_group': '',
|
|
'ldap_viewer_group': '',
|
|
'ldap_default_role': 'viewer',
|
|
'ldap_auto_create_users': True,
|
|
'ldap_group_mappings': [],
|
|
# NS: Mar 2026 - reverse proxy support (nginx/haproxy)
|
|
'reverse_proxy_enabled': False,
|
|
'trusted_proxies': '', # comma-separated IPs/CIDRs, empty = loopback only
|
|
'proxy_bind_address': '', # custom bind addr when behind proxy on different host
|
|
# OIDC defaults
|
|
'oidc_enabled': False,
|
|
'oidc_provider': 'entra',
|
|
'oidc_cloud_environment': 'commercial', # NS: commercial, gcc, gcc_high, dod
|
|
'oidc_client_id': '',
|
|
'oidc_client_secret': '',
|
|
'oidc_tenant_id': '',
|
|
'oidc_authority': '',
|
|
'oidc_scopes': 'openid profile email',
|
|
'oidc_redirect_uri': '',
|
|
'oidc_admin_group_id': '',
|
|
'oidc_user_group_id': '',
|
|
'oidc_viewer_group_id': '',
|
|
'oidc_default_role': 'viewer',
|
|
'oidc_auto_create_users': True,
|
|
'oidc_button_text': 'Sign in with Microsoft',
|
|
'oidc_group_mappings': [],
|
|
'oidc_skip_jwt_verification': False, # NS: disable JWT sig check for broken JWKS envs
|
|
'oidc_skip_ssl_verify': False, # NS Apr 2026 (#188): self-signed-cert escape hatch
|
|
# MK May 2026 (#412 SeeJayEmm): SSRF guard's default behaviour rejects
|
|
# any discovery URL that resolves to a private/loopback IP. Internal IdPs
|
|
# (Keycloak/Authentik/Authentik-on-LAN at 10.x or 192.168.x) are the
|
|
# exact use case that breaks. Opt-in knob to relax the guard for the
|
|
# OIDC discovery path SPECIFICALLY — metadata IPs (169.254.169.254
|
|
# etc.) are still rejected, and the guard remains on for all other
|
|
# outbound paths (webhook, SAML metadata fetch, plugin upstream).
|
|
'oidc_allow_private_ip': False,
|
|
# NS May 2026 (PVE 9.2 parity) — extra audiences (comma-separated)
|
|
# accepted on the JWT verify alongside the client_id.
|
|
'oidc_audiences': '',
|
|
}
|
|
|
|
try:
|
|
db = get_db()
|
|
saved = db.get_server_settings()
|
|
if saved:
|
|
# Merge with defaults (so new fields are always present)
|
|
return {**defaults, **saved}
|
|
except Exception as e:
|
|
logging.error(f"Error loading server settings from database: {e}")
|
|
# NS May 2026 - plain-JSON SERVER_SETTINGS_FILE fallback removed (encrypted DB only).
|
|
|
|
return defaults
|
|
|
|
|
|
def decrypt_secret_setting(value, *, label='secret'):
|
|
"""Decrypt an encrypted server setting, preserving legacy plaintext values."""
|
|
if not value or value == '********':
|
|
return ''
|
|
try:
|
|
return get_db()._decrypt(str(value))
|
|
except RuntimeError as e:
|
|
logging.error(f"Failed to decrypt {label}: {e}")
|
|
return ''
|
|
except Exception as e:
|
|
if str(value).startswith(('aes256:', 'gAAAA')):
|
|
logging.error(f"Could not decrypt encrypted {label}: {e}")
|
|
return ''
|
|
logging.warning(f"Could not decrypt {label}; treating as legacy plaintext: {e}")
|
|
return str(value)
|
|
|
|
|
|
def acme_dns_config_from_settings(settings):
|
|
"""Build an RFC 2136 DNS config, decrypting the TSIG secret only for use."""
|
|
settings = settings or {}
|
|
return {
|
|
'nameserver': settings.get('acme_dns_rfc2136_nameserver', ''),
|
|
'port': settings.get('acme_dns_rfc2136_port', 53),
|
|
'zone': settings.get('acme_dns_rfc2136_zone', ''),
|
|
'key_name': settings.get('acme_dns_rfc2136_key_name', ''),
|
|
'secret': decrypt_secret_setting(
|
|
settings.get('acme_dns_rfc2136_secret', ''),
|
|
label='ACME RFC 2136 secret'
|
|
),
|
|
'algorithm': settings.get('acme_dns_rfc2136_algorithm', 'hmac-sha512'),
|
|
'ttl': settings.get('acme_dns_rfc2136_ttl', 60),
|
|
'propagation_seconds': settings.get('acme_dns_propagation_seconds', 30),
|
|
}
|
|
|
|
|
|
def save_server_settings(settings):
|
|
"""Save server settings to SQLite database
|
|
|
|
SQLite migration
|
|
"""
|
|
try:
|
|
db = get_db()
|
|
db.save_server_settings(settings)
|
|
return True
|
|
except Exception as e:
|
|
logging.error(f"Error saving server settings: {e}")
|
|
return False
|
|
|
|
|
|
def get_session_timeout():
|
|
# get timeout from settings
|
|
try:
|
|
settings = load_server_settings()
|
|
return settings.get('session_timeout', SESSION_TIMEOUT)
|
|
except:
|
|
return SESSION_TIMEOUT # fallback
|
|
|
|
def _fmt_size(size_bytes):
|
|
# NS: simple bytes formatter, nothing fancy
|
|
if size_bytes < 1024:
|
|
return f"{size_bytes} B"
|
|
elif size_bytes < 1024**2:
|
|
return f"{size_bytes/1024:.1f} KB"
|
|
elif size_bytes < 1024**3:
|
|
return f"{size_bytes/1024**2:.1f} MB"
|
|
else:
|
|
return f"{size_bytes/1024**3:.1f} GB"
|
|
# TODO: add TB support? probably overkill
|
|
|
|
def get_login_settings():
|
|
# MK: pulled these out to be configurable via settings
|
|
try:
|
|
settings = load_server_settings()
|
|
except:
|
|
settings = {} # w/e just use defaults
|
|
return {
|
|
'max_attempts': settings.get('login_max_attempts', LOGIN_MAX_ATTEMPTS),
|
|
'lockout_time': settings.get('login_lockout_time', LOGIN_LOCKOUT_TIME),
|
|
'attempt_window': settings.get('login_attempt_window', LOGIN_ATTEMPT_WINDOW)
|
|
}
|
|
|
|
def register_task_user(upid: str, username: str, cluster_id: str = None):
|
|
"""Register which PegaProx user initiated a task - persists to database"""
|
|
if not upid or not username:
|
|
return
|
|
|
|
# Update in-memory cache
|
|
with task_pegaprox_users_lock:
|
|
task_pegaprox_users_cache[upid] = {'user': username, 'timestamp': time.time()}
|
|
# S2 (regression scan): this dict was never evicted (TASK_USER_CACHE_TTL was
|
|
# dead) → slow unbounded RSS creep over weeks. Bound it: when over the cap,
|
|
# drop the oldest ~10% by timestamp. The DB row remains the source of truth
|
|
# (get_task_user falls back to the DB on a cache miss).
|
|
if len(task_pegaprox_users_cache) > 50000:
|
|
try:
|
|
_old = sorted(task_pegaprox_users_cache.items(),
|
|
key=lambda kv: kv[1].get('timestamp', 0))[:5000]
|
|
for _k, _ in _old:
|
|
task_pegaprox_users_cache.pop(_k, None)
|
|
except Exception:
|
|
task_pegaprox_users_cache.clear()
|
|
|
|
# Persist to database
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('''
|
|
INSERT OR REPLACE INTO task_users (upid, username, cluster_id, created_at)
|
|
VALUES (?, ?, ?, ?)
|
|
''', (upid, username, cluster_id, datetime.now().isoformat()))
|
|
db.conn.commit()
|
|
except Exception as e:
|
|
logging.debug(f"Failed to persist task user to DB: {e}")
|
|
|
|
def get_task_user(upid: str) -> str:
|
|
"""Get PegaProx user who initiated a task - checks cache first, then database"""
|
|
if not upid:
|
|
return None
|
|
|
|
# Check in-memory cache first (fast path)
|
|
with task_pegaprox_users_lock:
|
|
data = task_pegaprox_users_cache.get(upid)
|
|
if data:
|
|
return data.get('user')
|
|
|
|
# Check database (slow path, but persists across restarts)
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('SELECT username FROM task_users WHERE upid = ?', (upid,))
|
|
row = cursor.fetchone()
|
|
if row:
|
|
username = row[0]
|
|
# Update cache for future lookups
|
|
with task_pegaprox_users_lock:
|
|
task_pegaprox_users_cache[upid] = {'user': username, 'timestamp': time.time()}
|
|
return username
|
|
except Exception as e:
|
|
logging.debug(f"Failed to get task user from DB: {e}")
|
|
|
|
return None
|
|
|
|
|
|
|
|
def get_connected_manager(cluster_id):
|
|
"""Get a cluster manager, return (manager, None) if connected, (None, error_response) if not"""
|
|
from flask import jsonify
|
|
if cluster_id not in cluster_managers:
|
|
return None, (jsonify({'error': 'Cluster not found'}), 404)
|
|
manager = cluster_managers[cluster_id]
|
|
if not manager.is_connected:
|
|
return None, (jsonify({
|
|
'error': 'Cluster not connected',
|
|
'offline': True,
|
|
'connection_error': manager.connection_error
|
|
}), 503)
|
|
return manager, None
|
|
|
|
def check_cluster_access(cluster_id):
|
|
"""Check if current user can access a cluster based on tenant or VM ACLs.
|
|
Returns (True, None) if allowed, (False, error_response) if not.
|
|
"""
|
|
from flask import request, jsonify, g
|
|
from pegaprox.utils.rbac import get_user_clusters
|
|
# H2 (scale audit): reuse the acting user require_auth already fetched (g.current_user),
|
|
# else fetch just that one user — don't re-scan the whole users table per cluster route.
|
|
user = getattr(g, 'current_user', None)
|
|
if user is None:
|
|
try:
|
|
user = get_db().get_user(request.session['user']) or {}
|
|
except Exception:
|
|
from pegaprox.utils.auth import load_users
|
|
user = load_users().get(request.session['user'], {})
|
|
allowed = get_user_clusters(user)
|
|
if allowed is not None and cluster_id not in allowed:
|
|
# #248: check VM ACLs as fallback — users with VM-level access can reach the cluster
|
|
username = request.session.get('user', '')
|
|
from pegaprox.utils.rbac import load_vm_acls
|
|
cluster_acls = load_vm_acls().get(cluster_id, {})
|
|
for vmid, acl in cluster_acls.items():
|
|
if username in acl.get('users', []) or '*' in acl.get('users', []):
|
|
return True, None
|
|
# #555: pool fallback — any pool grant in THIS cluster lets the user reach it
|
|
# (per-VM gating still runs downstream via user_can_access_vm)
|
|
try:
|
|
groups = user.get('groups', []) if isinstance(user, dict) else []
|
|
# sec-review: a {pool: []} row is truthy as a dict but grants nothing — match
|
|
# the rest of the pool model (user_has_any_pool_access) and require a real perm.
|
|
_pp = get_db().get_user_pool_permissions(cluster_id, username, groups)
|
|
if any(p for p in _pp.values()):
|
|
return True, None
|
|
except Exception:
|
|
pass
|
|
return False, (jsonify({'error': 'Access denied to this cluster'}), 403)
|
|
return True, None
|
|
|
|
|
|
def check_pbs_access(pbs_id):
|
|
"""Check if current user can access a PBS server based on its linked clusters.
|
|
Returns (True, None) if allowed, (False, error_response) if not.
|
|
|
|
A PBS server is accessible if:
|
|
- User is admin (full access), OR
|
|
- PBS has no linked_clusters (backward compatibility - accessible to all), OR
|
|
- User has access to at least one of the PBS's linked clusters
|
|
"""
|
|
from flask import request, jsonify
|
|
from pegaprox.utils.auth import load_users
|
|
from pegaprox.utils.rbac import get_user_clusters
|
|
from pegaprox.globals import pbs_managers
|
|
from pegaprox.models.permissions import ROLE_ADMIN
|
|
|
|
# Check if PBS exists
|
|
if pbs_id not in pbs_managers:
|
|
return False, (jsonify({'error': 'PBS server not found'}), 404)
|
|
|
|
pbs_mgr = pbs_managers[pbs_id]
|
|
users = load_users()
|
|
user = users.get(request.session['user'], {})
|
|
|
|
# Admins have full access
|
|
if user.get('role') == ROLE_ADMIN:
|
|
return True, None
|
|
|
|
# Get PBS linked clusters
|
|
pbs_linked = pbs_mgr.linked_clusters or []
|
|
|
|
# If PBS has no linked clusters, allow access (backward compatibility)
|
|
if not pbs_linked:
|
|
return True, None
|
|
|
|
# Get user's allowed clusters
|
|
user_clusters = get_user_clusters(user)
|
|
|
|
# If user has access to all clusters (None), allow
|
|
if user_clusters is None:
|
|
return True, None
|
|
|
|
# Check if user has access to at least one linked cluster
|
|
for cluster_id in pbs_linked:
|
|
if cluster_id in user_clusters:
|
|
return True, None
|
|
|
|
return False, (jsonify({'error': 'Access denied to this PBS server'}), 403)
|
|
|
|
|
|
def check_vmware_access(vmware_id):
|
|
"""NS Jul 2026 (CodeAnt re-scan IDOR) — tenant gate for a VMware/ESXi server, mirroring
|
|
check_pbs_access. Most vmware.py routes only had a role perm and never scoped to tenant, so
|
|
any vmware.* holder could read/act on ANOTHER tenant's ESXi. A server is accessible if the
|
|
caller is a global admin, the server has no linked_clusters (backward-compat), the caller is
|
|
all-cluster (get_user_clusters None), or the caller reaches one of the server's linked clusters.
|
|
Returns (True, None) or (False, error_response)."""
|
|
from flask import request, jsonify
|
|
from pegaprox.utils.auth import load_users
|
|
from pegaprox.utils.rbac import get_user_clusters
|
|
from pegaprox.globals import vmware_managers
|
|
from pegaprox.models.permissions import ROLE_ADMIN
|
|
|
|
if vmware_id not in vmware_managers:
|
|
return False, (jsonify({'error': 'VMware server not found'}), 404)
|
|
user = load_users().get(request.session.get('user', ''), {})
|
|
if user.get('role') == ROLE_ADMIN:
|
|
return True, None
|
|
linked = getattr(vmware_managers[vmware_id], 'linked_clusters', None) or []
|
|
if not linked:
|
|
return True, None # backward-compat: unlinked server is accessible to all
|
|
uc = get_user_clusters(user)
|
|
if uc is None:
|
|
return True, None
|
|
if any(c in uc for c in linked):
|
|
return True, None
|
|
return False, (jsonify({'error': 'Access denied to this VMware server'}), 403)
|
|
|
|
|
|
def safe_error(e, default_msg='An internal error occurred'):
|
|
"""Return a safe error message for API responses.
|
|
MK Feb 2026 - logs full exception but returns generic message to client.
|
|
Prevents leaking internal paths, stack traces, and DB details.
|
|
"""
|
|
logging.error(f"[API] {default_msg}: {e}", exc_info=True)
|
|
return default_msg
|
|
|
|
|
|
def parse_pve_error(response_text, fallback='Proxmox API error'):
|
|
"""Extract user-friendly error from Proxmox API response.
|
|
PVE returns JSON like {"data":null,"message":"some error\\n"} or plain text.
|
|
|
|
MK May 2026 — defense-in-depth: HTML-escape the extracted message before
|
|
returning. Reflecting raw upstream response text into our JSON error
|
|
field gets flagged by Snyk Code as reflected-XSS-via-JSON even though
|
|
Flask's jsonify sets Content-Type: application/json (which prevents
|
|
browser execution). Escaping makes the trace clean and gives us a
|
|
safety net if a future code path returns this string as text/html.
|
|
"""
|
|
import html
|
|
if not response_text:
|
|
return fallback
|
|
try:
|
|
import json
|
|
# PVE often has literal newlines in JSON strings — strip them
|
|
cleaned = response_text.replace('\n', ' ').replace('\r', '')
|
|
data = json.loads(cleaned)
|
|
msg = data.get('message') or data.get('errors') or data.get('error')
|
|
if isinstance(msg, dict):
|
|
msg = '; '.join(f"{k}: {v}" for k, v in msg.items())
|
|
if msg:
|
|
return html.escape(str(msg).strip()[:500])
|
|
except (json.JSONDecodeError, ValueError, AttributeError):
|
|
pass
|
|
# plain text — truncate and clean
|
|
text = response_text.strip()[:200]
|
|
if '<html' in text.lower():
|
|
return fallback
|
|
return html.escape(text) if text else fallback
|
|
|
|
|
|
# NS 2026-06-04 — shared metrics_history loader for insights/cost/power.
|
|
# The expensive part of these three endpoints isn't the SQL fetch, it's
|
|
# json.loads()'ing every snapshot blob (8.6k rows over 30d). Tier-1 moved the
|
|
# fetch off the gevent hub; this moves the PARSE off too AND caches the parsed
|
|
# result. The parse runs inside run_heavy_read's transform (worker thread), so
|
|
# even the cold-cache caller doesn't block the hub, and concurrent callers for
|
|
# the same window coalesce onto one fetch+parse (single-flight). Returns a list
|
|
# of (ts_unix, clusters_dict) for ALL clusters; callers filter for their own id.
|
|
def _history_stride(days):
|
|
# Snapshots land ~every 5 min. Parsing thousands of them is the GIL-bound
|
|
# ceiling (json.loads holds the GIL even in a worker thread), so for long
|
|
# windows we decimate to a coarser cadence. All three consumers are
|
|
# ratio/average/percentile based — sample COUNT doesn't change the result,
|
|
# only the resolution — so this is lossless for cost/power numbers and only
|
|
# smooths insights trends. Recent (<=2d) views keep full 5-min detail.
|
|
if days <= 2:
|
|
return 1 # 5-min, full resolution
|
|
if days <= 14:
|
|
return 3 # ~15-min
|
|
return 12 # ~hourly for month+ windows
|
|
|
|
|
|
def load_metrics_window(days):
|
|
from datetime import timedelta
|
|
from pegaprox.core.dbcrypto import run_heavy_read
|
|
cutoff = (datetime.now() - timedelta(days=days)).isoformat()
|
|
stride = _history_stride(days)
|
|
|
|
def _parse(rows):
|
|
out = []
|
|
for row in rows:
|
|
try:
|
|
d = json.loads(row['data'])
|
|
ts_unix = int(datetime.fromisoformat(row['timestamp']).timestamp())
|
|
out.append((ts_unix, d.get('clusters') or {}))
|
|
except Exception:
|
|
continue
|
|
return out
|
|
|
|
# `id % stride = 0` picks ~every Nth snapshot. rowid lives in the timestamp
|
|
# index, so SQLite evaluates the modulo without a table lookup and only
|
|
# decrypts the `data` blob for rows it keeps — cuts BOTH decrypt and parse.
|
|
if stride > 1:
|
|
sql = ("SELECT timestamp, data FROM metrics_history "
|
|
"WHERE timestamp >= ? AND id % ? = 0 ORDER BY timestamp ASC")
|
|
sql_params = (cutoff, stride)
|
|
else:
|
|
sql = ("SELECT timestamp, data FROM metrics_history "
|
|
"WHERE timestamp >= ? ORDER BY timestamp ASC")
|
|
sql_params = (cutoff,)
|
|
|
|
# cache_key shared across every cluster + across insights/cost/power.
|
|
# NOTE: the returned dicts are the cached parsed structure — callers must
|
|
# treat them as read-only (the aggregation paths only read, never mutate).
|
|
return run_heavy_read(sql, sql_params, cache_key=f"mh_parsed:{days}", transform=_parse)
|