MrMasterbay 77bfe71d98 security: enforce authz/validation gaps from the Aikido Testing-branch pentest (batch 2/2)
Second half of the adversarially-verified findings (batch 1 = 35d4078). Each fixed,
code-reviewed and covered by tests/test_aikido_batch2.py (17 new, full suite 401 green).

- power: only a global admin (effective_role) may overwrite the shared __default__
  power-rate row — a cluster.config holder edits only its own cluster
- metrics exporter: /api/metrics requires an admin-role token, not any valid token
  (it emits cluster-wide, cross-tenant infra gauges)
- portal: build_authz_user in _vm_power so a token's effective_role is honoured;
  invalidate the user's other sessions on portal password change
- cluster-groups: treat a global (tenant_id NULL) group as admin-only for the
  delete + balance-now writes, matching the earlier update fix
- vm-tags: reject a non-numeric vmid before the global DELETE+rewrite, and roll back
  save_vm_tags on error so a mid-loop failure can't persist a partial table wipe
- datacenter/multipath: allowlist the path_selector policy, and reject non-member
  nodes before SSH (no more `_get_node_ip(node) or node` fallback to a raw hostname)
- storage: pin http download-url fetches to the validated IP (DNS-rebind); https is
  left as the hostname since the node's TLS cert check already defeats a rebind
- multi-sdn: advertise an in-flight span's zone/controller so a concurrent purge
  can't tear down infra a create is still building (TOCTOU)
- ws-token validate: enforce node.shell for the standalone node SSH shell path
  (shell=node) — the VM termproxy path is unaffected
- SSE: scope vmware_vms / vmware_vm_detail to the server's linked_clusters instead
  of broadcasting guest_info/performance to every client; scope the portal audit
  task feed to the cluster it happened on (portal writers now set cluster=)
- LDAP: authoritative re-sync — rebuild LDAP-sourced perms/tenant_permissions from
  the current group mapping instead of only unioning them in, so group removal revokes
2026-08-07 17:31:57 +02:00

493 lines
21 KiB
Python

# -*- coding: utf-8 -*-
"""
PegaProx Realtime API Routes - Layer 6
WebSocket, SSE, and email test endpoints.
"""
import json
import logging
import threading
import uuid
import queue as queue_module
from datetime import datetime
from flask import Blueprint, jsonify, request, Response
from flask_sock import Sock
# MK 2026-06-04 (CWE-117 log-injection scanner findings): strip CR/LF/U+2028/9
# from anything user-controlled before f-stringing into a logger.
from pegaprox.utils.sanitization import sanitize_log_message as _sl
from pegaprox.constants import *
from pegaprox.globals import (
cluster_managers, vmware_managers,
ws_clients, ws_clients_lock,
sse_clients, sse_clients_lock,
)
from pegaprox.utils.auth import require_auth, validate_session, load_users
from pegaprox.utils.rbac import get_user_clusters
from pegaprox.utils.realtime import (
broadcast_update, broadcast_sse, broadcast_action,
create_sse_token, validate_sse_token,
create_ws_token, validate_ws_token,
push_immediate_update,
)
from pegaprox.utils.email import send_email
from pegaprox.api.helpers import load_server_settings, get_connected_manager
from pegaprox.models.permissions import ROLE_ADMIN
bp = Blueprint('realtime', __name__)
sock = Sock()
@sock.route('/api/ws/updates')
def ws_live_updates(ws):
"""WebSocket endpoint for live updates"""
client_id = str(uuid.uuid4())
client_lock = threading.Lock()
# Authenticate via first message
try:
auth_msg = ws.receive(timeout=3)
auth_data = json.loads(auth_msg)
session_id = auth_data.get('session_id')
session = validate_session(session_id)
if not session:
ws.send(json.dumps({'type': 'error', 'message': 'Authentication required'}))
return
username = session['user']
subscribed_clusters = auth_data.get('clusters', None)
with ws_clients_lock:
ws_clients[client_id] = {
'ws': ws,
'lock': client_lock,
'user': username,
'clusters': subscribed_clusters,
'connected_at': datetime.now().isoformat()
}
logging.info(f"WebSocket client connected: {_sl(username)} ({client_id})")
ws.send(json.dumps({'type': 'connected', 'client_id': client_id}))
# Keep connection alive
while True:
try:
# Wait for incoming messages with timeout
msg = ws.receive(timeout=30)
if msg is None:
break
data = json.loads(msg)
msg_type = data.get('type')
if msg_type == 'ping':
with client_lock:
ws.send(json.dumps({'type': 'pong'}))
elif msg_type == 'pong':
pass
elif msg_type == 'subscribe':
with ws_clients_lock:
if client_id in ws_clients:
ws_clients[client_id]['clusters'] = data.get('clusters')
except Exception as e:
err_str = str(e).lower()
if 'timed out' in err_str:
# Send ping on timeout
try:
with client_lock:
ws.send(json.dumps({'type': 'ping'}))
except:
break
else:
logging.debug(f"WebSocket error for {client_id}: {e}")
break
except Exception as e:
logging.error(f"WebSocket connection error: {e}")
finally:
with ws_clients_lock:
if client_id in ws_clients:
del ws_clients[client_id]
logging.info(f"WebSocket client disconnected: {client_id}")
@bp.route('/api/sse/token', methods=['POST'])
@require_auth()
def get_sse_token():
"""Get SSE token for URL param auth"""
user = request.session.get('user', 'unknown')
users = load_users()
user_data = users.get(user, {})
allowed_clusters = get_user_clusters(user_data)
token = create_sse_token(user, allowed_clusters)
return jsonify({
'token': token,
'expires_in': SSE_TOKEN_TTL,
'hint': 'Use this token in /api/sse/updates?token=...'
})
# NS: Mar 2026 - WebSocket auth tokens (single-use, 60s TTL)
# VNC/SSH WebSocket servers call /api/ws/token/validate instead of trusting session in URL
@bp.route('/api/ws/token', methods=['POST'])
@require_auth()
def get_ws_token():
"""Get a single-use WebSocket auth token - avoids session_id in URLs"""
user = request.session.get('user', 'unknown')
role = request.session.get('role', 'viewer')
token = create_ws_token(user, role)
return jsonify({'token': token, 'expires_in': 60})
@bp.route('/api/ws/token/validate')
def validate_ws_token_api():
"""Validate a WS token - called by standalone VNC/SSH servers
MK: internal endpoint, consumes the token (single-use)
MK May 2026 (CodeAnt CWE-285) - if the caller passes ?cluster_id=<id>, also
verify the token's user has access to that cluster. Closes the cross-cluster
BOLA case where someone with cluster-A access could open a WS for cluster-B
and trust the token alone to gate it. cluster_id is OPTIONAL for back-compat
(the VNC paths in vms.py / VM-level shells call without it today).
"""
token = request.args.get('token')
if not token:
return jsonify({'error': 'Token required'}), 401
data = validate_ws_token(token)
if not data:
return jsonify({'error': 'Invalid or expired token'}), 401
requested_cluster = (request.args.get('cluster_id') or '').strip()
cluster_context = None
if requested_cluster:
try:
from pegaprox.utils.auth import load_users
from pegaprox.utils.rbac import get_user_clusters, load_vm_acls
from pegaprox.core.db import get_db
# MK Aug 2026 — resolve the token's user by its indexed row, not a whole-table
# load_users() reload. That read decrypts every user's TOTP; on a transient
# failure (SQLite/WAL contention under gevent, a bad TOTP row) it degrades to {},
# which get_user_clusters() then reads as a default-tenant viewer — silently
# dropping admin/all-access and 403-ing a valid node console ("No access to
# cluster", intermittent). An unresolvable identity is a retryable auth failure
# (401), not a cluster denial; a genuinely unauthorized user still resolves + 403s.
try:
user = get_db().get_user(data['user'])
except Exception:
user = load_users().get(data['user'])
if not user:
return jsonify({'error': 'Invalid or expired token'}), 401
allowed = get_user_clusters(user)
access_ok = allowed is None or requested_cluster in allowed
if not access_ok:
# VM-ACL fallback (mirrors api/helpers.py check_cluster_access)
cluster_acls = load_vm_acls().get(requested_cluster, {}) or {}
for _vmid, acl in cluster_acls.items():
if data['user'] in (acl.get('users') or []) or '*' in (acl.get('users') or []):
access_ok = True
break
if not access_ok:
# #555: pool fallback — a pool grant in this cluster lets the WS token open
# (cluster-level reach only; the proxy + user_can_access_vm gate the VM)
try:
from pegaprox.core.db import get_db
# sec-review: ignore empty {pool: []} grants (truthy dict, no perm)
_pp = get_db().get_user_pool_permissions(requested_cluster, data['user'], user.get('groups', []))
if any(p for p in _pp.values()):
access_ok = True
except Exception:
pass
if not access_ok:
logging.warning(f"[WS-TOKEN] user '{_sl(data['user'])}' has no access to cluster '{_sl(requested_cluster)}'")
return jsonify({'error': 'Access denied to this cluster'}), 403
# NS Aug 2026 (Aikido pentest) — the standalone SSH shell server (mainPort+2) calls
# this with &shell=node. Unlike the in-process node_shell_websocket_proxy (vms.py),
# this validate path never checked node.shell, so a user with mere cluster access
# could open a root node shell. Enforce node.shell here. The VM termproxy path does
# NOT send shell=node (it's a guest console gated separately), so it's unaffected.
if (request.args.get('shell') or '') == 'node':
from pegaprox.utils.rbac import has_permission
if not has_permission(user, 'node.shell'):
logging.warning(f"[WS-TOKEN] user '{_sl(data['user'])}' lacks node.shell for a node shell on '{_sl(requested_cluster)}'")
return jsonify({'error': 'node.shell permission required'}), 403
# MK May 2026 - lightweight cluster context for the SSH/VNC proxy.
# We intentionally do NOT call mgr._get_node_ip() here: that has a
# network-probing first-call path (up to ~15s) which would hang the
# validate endpoint when this same Flask process is busy. The proxy
# gets the cluster's primary host + any fallback_hosts already known
# from the DB. Multi-node clusters where the frontend prefetched a
# per-node IP that isn't in this set will need to re-resolve through
# the normal cluster-creds endpoint with a session cookie.
try:
from pegaprox.globals import cluster_managers
mgr = cluster_managers.get(requested_cluster)
if mgr is not None:
cluster_host = getattr(mgr, 'host', None)
cfg = getattr(mgr, 'config', None)
node_ips = {}
# Include fallback_hosts from the DB — those were discovered by
# the manager at connect time and persist across restarts. Cheap.
fb = getattr(cfg, 'fallback_hosts', None) or []
for i, host in enumerate(fb):
if host:
node_ips[f'_fallback_{i}'] = host
cluster_context = {
'host': cluster_host,
'node_ips': node_ips,
'ssh_port': getattr(cfg, 'ssh_port', 22) or 22,
}
# NS 2026-06-05 (C-1): hand the PVE session cookie to the WS
# subprocess server-side (it used to come from the browser).
# Mint fresh; None for token-only clusters. termproxy is the
# only consumer — VNC/SSH ignore it.
try:
_tk = mgr.mint_console_auth_ticket()
if _tk:
cluster_context['pve_auth_ticket'] = _tk
except Exception:
pass
except Exception as e:
logging.debug(f"[WS-TOKEN] cluster-context build soft-fail: {e}")
except Exception as e:
logging.error(f"[WS-TOKEN] cluster-access check failed: {e}")
# fail closed
return jsonify({'error': 'Authorization check failed'}), 500
resp = {'valid': True, 'user': data['user'], 'role': data['role']}
if cluster_context is not None:
resp['cluster_context'] = cluster_context
return jsonify(resp)
@bp.route('/api/sse/updates')
def sse_updates():
"""SSE endpoint for live updates
NS: accepts ?token= (preferred) or ?session= (legacy)
MK: token is better because session IDs in URLs can leak to logs
"""
# token auth first (preferred)
sse_token = request.args.get('token')
session_id = request.args.get('session')
user = None
allowed_clusters = None
auth_method = None
if sse_token:
# Validate SSE token
token_data = validate_sse_token(sse_token)
if token_data:
user = token_data['user']
allowed_clusters = token_data['allowed_clusters']
auth_method = 'token'
# NS Mar 2026 - removed session_id fallback, token-only auth for SSE
if not user:
return jsonify({'error': 'Authentication required. Provide a valid SSE token.'}), 401
client_id = str(uuid.uuid4())
message_queue = queue_module.Queue(maxsize=100)
# Get cluster subscription from query params
clusters_param = request.args.get('clusters')
requested_clusters = clusters_param.split(',') if clusters_param else None
# MK: only let users subscribe to clusters they have access to
if requested_clusters:
if allowed_clusters is None:
# admin - all clusters allowed
subscribed_clusters = requested_clusters
else:
# filter to allowed only
subscribed_clusters = [c for c in requested_clusters if c in allowed_clusters]
if not subscribed_clusters:
logging.warning(f"[SSE] User {user} tried to subscribe to unauthorized clusters")
subscribed_clusters = allowed_clusters
else:
subscribed_clusters = allowed_clusters
with sse_clients_lock:
sse_clients[client_id] = {
'queue': message_queue,
'user': user,
'clusters': subscribed_clusters,
'connected_at': datetime.now().isoformat(),
'auth_method': auth_method
}
logging.info(f"[SSE] Client connected: {client_id} (user: {user}, auth: {auth_method}) - Total: {len(sse_clients)}")
def generate():
try:
# Send initial connected message
yield f"data: {json.dumps({'type': 'connected', 'client_id': client_id})}\n\n"
while True:
try:
# Wait for message with timeout
message = message_queue.get(timeout=30)
yield f"data: {message}\n\n"
except queue_module.Empty:
# Send keepalive
yield f": keepalive\n\n"
except GeneratorExit:
pass
finally:
with sse_clients_lock:
if client_id in sse_clients:
del sse_clients[client_id]
logging.info(f"[SSE] Client disconnected: {client_id} - Remaining clients: {len(sse_clients)}")
response = Response(generate(), mimetype='text/event-stream')
response.headers['Cache-Control'] = 'no-cache'
response.headers['X-Accel-Buffering'] = 'no'
response.headers['Connection'] = 'keep-alive'
return response
@bp.route('/api/sse/subscribe', methods=['POST'])
@require_auth()
def update_sse_subscription():
"""Update cluster subscription for an active SSE client without reconnecting.
NS: Mar 2026 - avoids 200-500ms data gap on sidebar toggle
"""
data = request.json or {}
client_id = data.get('client_id')
requested = data.get('clusters') # list of cluster IDs or None for all
if not client_id:
return jsonify({'error': 'client_id required'}), 400
username = request.session.get('user', 'unknown')
# RBAC: what clusters is this user allowed to see?
users = load_users()
user_data = users.get(username, {})
allowed = get_user_clusters(user_data) # None = admin
# filter requested against allowed
if requested and len(requested) > 0:
if allowed is not None:
filtered = [c for c in requested if c in allowed]
new_sub = filtered if filtered else allowed
else:
new_sub = requested # admin sees all
else:
new_sub = allowed # None = everything user can see
with sse_clients_lock:
client = sse_clients.get(client_id)
if not client:
# NS: return 200 not 404 — client may have reconnected with a new ID
# (token refresh cycle), frontend treats subscribe as best-effort anyway
return jsonify({'ok': False, 'reason': 'client_not_found'})
if client.get('user') != username:
return jsonify({'error': 'Unauthorized'}), 403
old_sub = client.get('clusters')
client['clusters'] = new_sub
# R2 (regression fix): the IP/disk refresh loop is gated on watched-clusters
# and only re-runs ~every 40s, so a freshly-selected cluster would show stale
# guest IPs/disk for that window. Kick a one-shot refresh for clusters that
# just became watched (explicit→explicit transition only; all-access already
# refreshes everything).
if old_sub is not None and new_sub is not None:
newly = set(new_sub) - set(old_sub)
for cid in newly:
mgr = cluster_managers.get(cid)
if mgr and getattr(mgr, 'is_connected', False) and hasattr(mgr, 'refresh_ip_cache'):
try:
import gevent
gevent.spawn(mgr.refresh_ip_cache)
except Exception:
pass
logging.debug(f"[SSE] Subscription updated for {client_id}: {new_sub}")
return jsonify({'ok': True, 'clusters': new_sub})
@bp.route('/api/settings/smtp/test', methods=['POST'])
@require_auth(perms=['admin.settings'])
def test_smtp():
"""Send a test email to verify SMTP settings
NS: Uses the same send_email function for consistency
"""
data = request.json or {}
test_email = data.get('email', '')
logging.info(f"[SMTP Test] Received data: {list(data.keys())}")
if not test_email:
return jsonify({'error': 'Email address required'}), 400
# Load saved settings first (we might need the real password)
saved_settings = load_server_settings()
# Build SMTP settings from request or use saved
smtp_host = data.get('smtp_host', '')
if smtp_host:
# Use provided settings for testing (before save)
# But if password is masked (********), use the saved password
provided_password = data.get('smtp_password', '')
if provided_password == '********' or not provided_password:
# NS Aug 2026 (Aikido pentest) — only reuse the saved SMTP password when the test host
# matches the saved host; otherwise a settings admin could point the stored credential
# at an attacker-controlled MX = credential exfil. A different host needs its own password.
if smtp_host != (saved_settings.get('smtp_host') or ''):
return jsonify({'error': 'A password is required to test a host other than the saved one'}), 400
# Use saved password - NS: Feb 2026: now encrypted in DB, must decrypt
raw_password = saved_settings.get('smtp_password', '')
try:
from pegaprox.core.db import get_db
real_password = get_db()._decrypt(raw_password) if raw_password else ''
except Exception:
real_password = raw_password # Fallback for unencrypted legacy values
logging.info("[SMTP Test] Using saved password (frontend sent masked value)")
else:
real_password = provided_password
smtp_settings = {
'smtp_host': smtp_host,
'smtp_port': data.get('smtp_port', 587),
'smtp_user': data.get('smtp_user', ''),
'smtp_password': real_password,
'smtp_from_email': data.get('smtp_from_email', ''),
'smtp_from_name': data.get('smtp_from_name', 'PegaProx'),
'smtp_tls': data.get('smtp_tls', True),
'smtp_ssl': data.get('smtp_ssl', False),
}
if not smtp_settings['smtp_from_email']:
return jsonify({'error': 'From email address is required'}), 400
logging.info(f"[SMTP Test] Using settings: host={smtp_host}, user={smtp_settings['smtp_user']}, has_password={bool(real_password)}")
else:
# Use saved settings
smtp_settings = None # send_email will load from database
if not saved_settings.get('smtp_enabled'):
return jsonify({'error': 'SMTP not enabled'}), 400
# Send test email using the same function as alerts
success, error = send_email(
to_addresses=[test_email],
subject='PegaProx Test Email',
body='This is a test email from PegaProx to verify your SMTP settings are working correctly.',
html_body='<h2>PegaProx Test Email</h2><p>This is a test email to verify your SMTP settings.</p><p style="color: green;">Your SMTP configuration is working!</p>',
smtp_settings=smtp_settings
)
if success:
return jsonify({'success': True, 'message': f'Test email sent to {test_email}'})
else:
return jsonify({'error': error or 'Failed to send test email'}), 400