mirror of
https://github.com/PegaProx/project-pegaprox.git
synced 2026-08-12 15:27:47 +08:00
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.
210 lines
7.6 KiB
Python
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
|
|
|