Nico Schmidt d5215eecd2 perf(scale): raise concurrency limits for large fleets (30+ clusters / 100+ nodes)
Hardening for big deployments (#526/#528 class — single gevent hub, everything
funnels through it):
- WSGIServer request pool default max(8,cpu*4) -> max(32,cpu*16). Each live
  SSE/WebSocket stream holds a pool slot for its whole lifetime, so the old
  default (16 on a 4c box) could be drained by ~16 open dashboard tabs and
  starve all other API traffic. Still PEGAPROX_WORKERS-overridable.
- GEVENT_POOL (node-status / IP-refresh fan-out) 50 -> 100, env
  PEGAPROX_NODE_POOL_SIZE. Safe to raise now that managers reuse keep-alive
  sessions: connections are pooled, not one fd per call — that fd pressure is
  what forced 100->50 before.
- manager session pool_maxsize 16 -> 64 to match the node fan-out, so a cluster
  polling many nodes at once reuses connections instead of churning throwaway
  ones (pool_block=False).
- at startup: raise RLIMIT_NOFILE to the hard cap, and bump gevent's threadpool
  from its default 10 (shared by the DNS resolver AND off-hub DB reads — the
  #528 contention point). env PEGAPROX_NOFILE / PEGAPROX_THREADPOOL_SIZE.

verified: boot prints new limits (fd 1048576, threadpool 50, 192 greenlets on
12c); pools applied (GEVENT_POOL 100, session maxsize 64); 20 concurrent
resources calls -> all 200, health 51ms, ZERO "pool is full" churn warnings.
2026-06-05 09:01:45 +02:00

210 lines
7.6 KiB
Python

# -*- coding: utf-8 -*-
"""
PegaProx Concurrency Helpers - Layer 2
"""
import os
import logging
from typing import Dict
GEVENT_AVAILABLE = False
GEVENT_PATCHED = False
GEVENT_POOL = None
# NS 2026-06-05 — env-tunable (was hard 50). Default raised to 100 now that
# managers reuse keep-alive sessions (#528): the fd pressure that forced 100→50
# on the old fresh-session-per-call model is gone since connections are pooled.
_NODE_POOL_SIZE = int(os.environ.get('PEGAPROX_NODE_POOL_SIZE', '100'))
try:
from gevent.pool import Pool as GeventPool
GEVENT_POOL = GeventPool(size=_NODE_POOL_SIZE)
GEVENT_AVAILABLE = True
# Check if gevent has actually monkey-patched the socket module
import gevent.monkey
GEVENT_PATCHED = gevent.monkey.is_module_patched('socket')
except ImportError:
pass
def get_paramiko():
"""lazy import for paramiko, its optional"""
# MK: paramiko takes forever to import so we only do it when needed
try:
import paramiko
return paramiko
except ImportError:
return None
# ============================================
# Concurrent API Helpers - added late 2025
# Use gevent pool for parallel requests when available
# MK: This made the dashboard like 5x faster, totally worth it
# ============================================
def run_concurrent(tasks: list, timeout: float = 30.0) -> list:
"""Run tasks concurrently with gevent pool"""
# NS: chatgpt helped with this one, i was mass confused about greenlets
# TODO: maybe add retry logic? - MK
#
# MK 2026-05-31 — CRITICAL FIX. The original check `if GEVENT_POOL and
# GEVENT_AVAILABLE` was always-False on entry: gevent.pool.Pool overrides
# __bool__ to len() == 0. So every call silently fell through to the
# sequential branch from day one. The "5x faster" comment above was
# aspiration, not reality. Switching to `is not None` actually wires up
# the parallel path the helper was designed for.
if not tasks:
return []
if GEVENT_POOL is not None and GEVENT_AVAILABLE:
# Use gevent pool for concurrent execution
try:
greenlets = [GEVENT_POOL.spawn(task) for task in tasks]
# Wait for all with timeout
from gevent import joinall
joinall(greenlets, timeout=timeout)
results = []
for g in greenlets:
try:
results.append(g.value if g.successful() else None)
except Exception as e:
logging.error(f"Concurrent task failed: {e}")
results.append(None)
return results
except Exception as e:
logging.error(f"Concurrent execution failed: {e}")
# Fall through to sequential execution
# Fallback: sequential execution (when gevent not available)
results = []
for task in tasks:
try:
results.append(task())
except Exception as e:
logging.error(f"Task failed: {e}")
results.append(None)
return results
def run_concurrent_dict(tasks: dict, timeout: float = 30.0) -> dict:
"""same as run_concurrent but takes/returns a dict of {key: callable} -> {key: result}"""
if not tasks:
return {}
keys = list(tasks.keys())
callables = [tasks[k] for k in keys]
results = run_concurrent(callables, timeout)
return dict(zip(keys, results))
# MK: exponential backoff helper for retryable SSH/API ops
# used by predictive analysis engine and cross-cluster sync
def retry_with_backoff(fn, max_retries=3, base_delay=0.5, jitter=True):
"""Retry a callable with exponential backoff. Returns (success, result)."""
import time, random
last_err = None
for attempt in range(max_retries):
try:
result = fn()
return True, result
except Exception as e:
last_err = e
delay = base_delay * (2 ** attempt)
if jitter:
delay += random.uniform(0, delay * 0.3)
# NS: don't log first attempt failure, its noisy
if attempt > 0:
logging.debug(f"retry_with_backoff attempt {attempt+1}/{max_retries}: {e}")
time.sleep(delay)
return False, last_err
# NS Apr 2026 — SSH-aware multi-node fanout for big clusters (15+ nodes).
# Bounded concurrency so we don't open 30 simultaneous SSH connections (which
# triggers AccountLockFailures on hardened nodes — we hit this on ESXi already).
#
# CRITICAL: This helper is for NEW multi-node fanouts only (custom-scripts on
# many nodes, hardening-multi, compliance-dashboard backend aggregation).
# HA SSH paths (HA monitor, fence operations, evacuation) MUST NOT go through
# this — they have their own latency requirements and bypass any throttle.
# That's why it lives next to run_concurrent and not in ssh.py.
#
# Uses gevent pool (size-bounded) when gevent is available, otherwise falls
# back to a thread pool with a Semaphore.
def run_per_node(node_callables, max_concurrent=8, timeout=120):
"""Fan out per-node callables with bounded concurrency.
Args:
node_callables: dict {node_name: callable(node_name) -> any}
max_concurrent: hard ceiling on parallel SSH workers (default 8).
Tuned conservatively — going higher than 8 risks per-host SSH
rate-limits on busier nodes. Per-cluster, NOT global.
timeout: per-task wall-clock timeout in seconds.
Returns:
dict {node_name: result_or_None}. Failed/timed-out tasks return None,
the exception is logged at debug level.
"""
if not node_callables:
return {}
# Cap concurrency at the lesser of node count and max_concurrent
n = len(node_callables)
workers = max(1, min(int(max_concurrent), n))
# Path 1: gevent pool — preferred since pegaprox is gevent-monkey-patched
if GEVENT_AVAILABLE:
try:
from gevent.pool import Pool as GP
pool = GP(size=workers)
jobs = {}
for node, fn in node_callables.items():
# bind node name into the closure so the callable receives it
jobs[node] = pool.spawn(_run_node_safe, node, fn)
from gevent import joinall
joinall(list(jobs.values()), timeout=timeout)
results = {}
for node, g in jobs.items():
try:
results[node] = g.value if g.successful() else None
except Exception as e:
logging.debug(f"run_per_node[{node}] failed: {e}")
results[node] = None
return results
except Exception as e:
logging.warning(f"run_per_node gevent path failed, falling back: {e}")
# Path 2: stdlib threading + Semaphore — fallback when gevent isn't available
import threading
sem = threading.BoundedSemaphore(workers)
results = {}
threads = []
lock = threading.Lock()
def _worker(node, fn):
with sem:
r = _run_node_safe(node, fn)
with lock:
results[node] = r
for node, fn in node_callables.items():
t = threading.Thread(target=_worker, args=(node, fn), daemon=True)
t.start()
threads.append(t)
for t in threads:
t.join(timeout=timeout)
# Any thread still alive after timeout → that node is None
for node in node_callables:
results.setdefault(node, None)
return results
def _run_node_safe(node, fn):
"""Internal wrapper: invoke fn(node), swallow exceptions, return result or None."""
try:
return fn(node)
except Exception as e:
logging.debug(f"_run_node_safe[{node}] exception: {e}")
return None