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

318 lines
12 KiB
Python

# -*- coding: utf-8 -*-
"""
PegaProx Realtime Updates - Layer 4
WebSocket and SSE broadcasting utilities.
"""
import time
import json
import logging
import threading
import base64
import os
from datetime import datetime
from pegaprox.constants import SSE_TOKEN_TTL
from pegaprox.globals import (
cluster_managers, ws_clients, ws_clients_lock,
sse_tokens, sse_tokens_lock,
sse_clients, sse_clients_lock,
ws_tokens, ws_tokens_lock,
)
# NS 2026-06-05 (#528 scaling): max SSE/WS broadcast message size. The old hard
# 500KB cap silently dropped any broadcast above it — a cluster with thousands
# of VMs has a `resources` payload well over 500KB, so its live UI just stopped
# updating with only a log warning. Raised to 5MB, env-overridable. (The real
# long-term fix is per-cluster subscription so a client only gets its own data.)
_MAX_BROADCAST_BYTES = int(os.environ.get('PEGAPROX_MAX_BROADCAST_BYTES', str(5_000_000)))
def watched_clusters():
"""Cluster IDs at least one live SSE/WS client is subscribed to, or None if
any client has all-access (clusters=None → poll everything). Shared by the
broadcast loop AND the per-cluster background refreshers so they skip work
for clusters nobody is viewing. NS 2026-06-05 (scale audit H4 / #528)."""
watched = set()
with sse_clients_lock:
for c in list(sse_clients.values()):
sub = c.get('clusters')
if sub is None:
return None
watched.update(sub)
with ws_clients_lock:
for c in list(ws_clients.values()):
sub = c.get('clusters')
if sub is None:
return None
watched.update(sub)
return watched
def is_cluster_watched(cluster_id):
"""True if any live client is viewing this cluster (or has all-access)."""
w = watched_clusters()
return w is None or cluster_id in w
def push_immediate_update(cluster_id: str, delay: float = 0.3):
"""NS: push immediate SSE update after VM actions for faster UI feedback"""
def _push():
time.sleep(delay)
try:
if cluster_id not in cluster_managers:
return
manager = cluster_managers[cluster_id]
if not manager.is_connected:
return
# Push resources
# NS: Fixed - was calling get_all_resources() which doesn't exist
resources = manager.get_vm_resources()
if resources:
broadcast_sse('resources', resources, cluster_id)
# Push tasks — force=True bypasses the 3s result cache so the action's
# just-started task shows up immediately (N-2), not on the next tick.
tasks = manager.get_tasks(limit=50, force=True)
if tasks:
broadcast_sse('tasks', tasks, cluster_id)
except Exception as e:
logging.debug(f"[SSE] Immediate push failed for {cluster_id}: {e}")
threading.Thread(target=_push, daemon=True).start()
def broadcast_update(update_type: str, data: dict, cluster_id: str = None):
"""Broadcast update to all connected WebSocket clients"""
try:
message = json.dumps({
'type': update_type,
'data': data,
'cluster_id': cluster_id,
'timestamp': datetime.now().isoformat()
})
# Limit message size
if len(message) > _MAX_BROADCAST_BYTES:
logging.warning(f"Broadcast message too large ({len(message)} bytes), skipping")
return
disconnected = []
# Get clients list under lock, then send outside lock
clients_to_send = []
with ws_clients_lock:
for client_id, client_info in list(ws_clients.items()):
ws = client_info.get('ws')
client_lock = client_info.get('lock')
if ws is None or client_lock is None:
disconnected.append(client_id)
continue
# Only send if client is subscribed to this cluster or all clusters
subscribed = client_info.get('clusters')
if cluster_id is None or subscribed is None or cluster_id in subscribed:
clients_to_send.append((client_id, ws, client_lock))
# Send to clients outside the main lock
for client_id, ws, client_lock in clients_to_send:
try:
with client_lock:
ws.send(message)
except Exception as e:
logging.debug(f"Failed to send to client {client_id}: {e}")
disconnected.append(client_id)
# Remove disconnected clients
if disconnected:
with ws_clients_lock:
for client_id in set(disconnected): # Use set to avoid duplicates
if client_id in ws_clients:
del ws_clients[client_id]
logging.info(f"Removed disconnected client: {client_id}")
except Exception as e:
logging.error(f"Broadcast error: {e}")
def broadcast_action(action: str, resource_type: str, resource_id: str, details: dict = None, cluster_id: str = None, user: str = None):
"""Broadcast an action event to all clients for real-time UI updates"""
broadcast_update('action', {
'action': action,
'resource_type': resource_type,
'resource_id': resource_id,
'details': details or {},
'user': user
}, cluster_id)
def create_sse_token(username: str, allowed_clusters: list) -> str:
"""Create SSE token - avoids session ID in URL"""
token = base64.urlsafe_b64encode(os.urandom(24)).decode('utf-8')
expires = time.time() + SSE_TOKEN_TTL
with sse_tokens_lock:
# cleanup expired
now = time.time()
expired = [t for t, data in sse_tokens.items() if data['expires'] < now]
for t in expired:
del sse_tokens[t]
sse_tokens[token] = {
'user': username,
'expires': expires,
'allowed_clusters': allowed_clusters
}
return token
def validate_sse_token(token: str) -> dict:
"""Validate an SSE token and return user info or None"""
if not token:
return None
with sse_tokens_lock:
token_data = sse_tokens.get(token)
if not token_data:
return None
if token_data['expires'] < time.time():
del sse_tokens[token]
return None
return token_data
# MK: Mar 2026 - WS tokens for VNC/SSH, avoids putting session_id in WebSocket URLs
# These are single-use and expire after 60s
WS_TOKEN_TTL = 60
def create_ws_token(username: str, role: str) -> str:
"""Create a short-lived single-use WebSocket auth token"""
token = base64.urlsafe_b64encode(os.urandom(24)).decode('utf-8')
expires = time.time() + WS_TOKEN_TTL
with ws_tokens_lock:
# cleanup old ones
now = time.time()
expired = [t for t, d in ws_tokens.items() if d['expires'] < now]
for t in expired:
del ws_tokens[t]
ws_tokens[token] = {
'user': username,
'role': role,
'expires': expires,
}
return token
def validate_ws_token(token: str) -> dict:
"""Validate and consume a WS token (single-use). Returns user info or None."""
if not token:
return None
with ws_tokens_lock:
token_data = ws_tokens.pop(token, None)
if not token_data:
return None
if token_data['expires'] < time.time():
return None
return token_data
def broadcast_sse(update_type: str, data: dict, cluster_id: str = None, target_clusters=None):
"""Broadcast update to SSE clients
For cluster-specific events (node_status, vm_update, etc.), only sends to clients
subscribed to that cluster. Global events (update_type starting with 'global_')
are sent to all clients.
NS Aug 2026 (Aikido pentest) — target_clusters scopes an event that maps to a SET of
clusters (e.g. a VMware/ESXi server's linked_clusters) rather than a single cluster_id.
When provided (not None) it takes precedence: deliver to all-access clients (subscribed
is None) and to any client whose subscription intersects target_clusters. An empty list
means "not linked to any cluster" → global, mirroring check_vmware_access's backward-compat
rule. Without it (default None) the classic cluster_id / global logic below is unchanged.
"""
try:
# MK 2026-05-31 — `default=str` so a datetime / set / bytes / custom
# object slipping into `data` doesn't TypeError and silently lose the
# broadcast. Caller's intent was "best-effort dispatch", not "verify
# data shape" — that's a stability/observability win for broadcasts
# like #413 layer 1 where a wrong arg shape killed the publisher.
try:
message = json.dumps({
'type': update_type,
'data': data,
'cluster_id': cluster_id,
'timestamp': datetime.now().isoformat()
}, default=str)
except (TypeError, ValueError) as _ser_err:
# If even default=str can't coerce, log enough context to find
# the bad caller, then drop. Don't take the broadcaster down.
logging.warning(
f"[SSE] broadcast '{update_type}' (cluster={cluster_id}) "
f"unserialisable, skipped: {_ser_err}"
)
return
# Limit message size
if len(message) > _MAX_BROADCAST_BYTES:
logging.warning(f"SSE message too large ({len(message)} bytes), skipping")
return
# Determine if this is a cluster-specific event
# NS: Added 'tasks' and 'resources' - broadcast loop sends these types
cluster_specific_events = ['node_status', 'vm_update', 'task_update', 'tasks',
'metrics', 'resources', 'migration', 'maintenance',
'ha_event', 'alert', 'ha_status']
is_cluster_specific = update_type in cluster_specific_events or cluster_id is not None
with sse_clients_lock:
for client_id, client_info in list(sse_clients.items()):
try:
q = client_info.get('queue')
subscribed = client_info.get('clusters')
should_send = False
if target_clusters is not None:
# NS Aug 2026 (Aikido pentest) — multi-cluster-scoped event (VMware
# linked_clusters). Empty → unlinked server → global (matches REST).
if not target_clusters:
should_send = True
elif subscribed is None:
should_send = True # admin / all-access
elif subscribed and any(c in subscribed for c in target_clusters):
should_send = True
elif not is_cluster_specific:
# Global event - send to everyone
should_send = True
elif cluster_id and subscribed is None:
# NS: subscribed=None means admin/all-access -> send everything
# Was previously blocking ALL SSE events for admin users!
should_send = True
elif cluster_id and subscribed and cluster_id in subscribed:
# Cluster-specific event and client is subscribed
should_send = True
if q and should_send:
try:
q.put_nowait(message)
except Exception:
# R3 (regression scan): a slow client's queue is full, so
# this frame is dropped — make it OBSERVABLE instead of
# silent (its VM grid goes stale otherwise with no signal).
n = client_info['dropped'] = client_info.get('dropped', 0) + 1
if n == 1 or n % 100 == 0:
logging.warning(f"[SSE] client {client_id} queue full — dropped {n} frames (slow consumer)")
except:
pass
except Exception as e:
logging.error(f"SSE broadcast error: {e}")