mirror of
https://github.com/PegaProx/project-pegaprox.git
synced 2026-08-12 15:27:47 +08:00
test_ssl_bootstrap was skipped in its entirety under a root euid, but only the cases that chmod a cert/dir to 0 actually need a non-root user (root ignores the permission bits and they would false-fail). Moved the skipif onto just those; the content/ENOENT/pyOpenSSL/config-guard/migration fail-closed cases now execute in a root CI container or under `sudo pytest` too, so the #633 contract keeps real coverage there instead of silently skipping. 19 pass locally (non-root, all run). Also corrected the author tags on the #632/#633 comments to MK (app.py, constants.py, core/manager.py, the test header) - they were carrying the wrong initials.
1732 lines
77 KiB
Python
1732 lines
77 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""
|
|
PegaProx Flask App Factory - Layer 8
|
|
Creates and configures the Flask application.
|
|
"""
|
|
|
|
import os
|
|
import sys
|
|
import time
|
|
import errno
|
|
import stat
|
|
import logging
|
|
import threading
|
|
import signal
|
|
import gc
|
|
import multiprocessing
|
|
import ssl
|
|
import socket
|
|
|
|
from flask import Flask, jsonify, request
|
|
from flask_cors import CORS
|
|
from flask_sock import Sock
|
|
from flask_compress import Compress
|
|
from pathlib import Path
|
|
|
|
from pegaprox.constants import (
|
|
PEGAPROX_VERSION, PEGAPROX_BUILD,
|
|
SESSION_TIMEOUT, SSL_CERT_FILE, SSL_KEY_FILE,
|
|
API_RATE_LIMIT, API_RATE_WINDOW, SSH_MAX_CONCURRENT,
|
|
)
|
|
from pegaprox import globals as g
|
|
from pegaprox.api import register_blueprints
|
|
|
|
|
|
def get_allowed_origins():
|
|
"""Get list of allowed CORS origins (dynamic for Open Source)"""
|
|
origins = set()
|
|
|
|
# 1. Environment variable origins (highest priority)
|
|
if g._cors_origins_env:
|
|
for origin in g._cors_origins_env.split(','):
|
|
origin = origin.strip()
|
|
if origin and origin != '*':
|
|
origins.add(origin)
|
|
|
|
# 2. Auto-detected origins from successful logins
|
|
origins.update(g._auto_allowed_origins)
|
|
|
|
# 3. If nothing configured, allow requests without Origin header (same-origin)
|
|
# This is safe because browsers always send Origin header for cross-origin requests
|
|
if not origins:
|
|
return None # None = no CORS headers = same-origin only
|
|
|
|
return list(origins)
|
|
|
|
|
|
def add_allowed_origin(origin: str):
|
|
"""Add an origin to the auto-allowed list (called on successful login)"""
|
|
if origin and origin.startswith(('http://', 'https://')) and origin != '*':
|
|
g._auto_allowed_origins.add(origin)
|
|
logging.info(f"Auto-allowed CORS origin: {origin}")
|
|
|
|
|
|
def create_app():
|
|
"""Flask application factory."""
|
|
# root_path must point to the project root (parent of pegaprox/)
|
|
# so that send_from_directory('web', ...) and other relative paths work
|
|
project_root = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
|
app = Flask(__name__, root_path=project_root)
|
|
|
|
# CORS Configuration - NS: Feb 2026 - only enable if origins are explicitly set
|
|
if g._cors_origins_env:
|
|
allowed_origins = [o.strip() for o in g._cors_origins_env.split(',') if o.strip() and o.strip() != '*']
|
|
if allowed_origins:
|
|
CORS(app, supports_credentials=True, resources={
|
|
r"/api/*": {
|
|
"origins": allowed_origins,
|
|
"methods": ["GET", "POST", "PUT", "DELETE", "OPTIONS"],
|
|
"allow_headers": ["Content-Type", "Authorization", "X-Username", "X-Session-Id"],
|
|
"expose_headers": ["Content-Type"],
|
|
"supports_credentials": True
|
|
}
|
|
})
|
|
# else: no CORS init = browser same-origin policy applies (safest default)
|
|
|
|
# Gzip compression
|
|
app.config['COMPRESS_MIMETYPES'] = [
|
|
'text/html', 'text/css', 'text/xml', 'text/plain',
|
|
'application/json', 'application/javascript', 'application/xml'
|
|
]
|
|
app.config['COMPRESS_LEVEL'] = 6
|
|
app.config['COMPRESS_MIN_SIZE'] = 500
|
|
Compress(app)
|
|
|
|
# Max request size - NS: Feb 2026 - separate limit for file uploads (#82)
|
|
_default_max = int(os.environ.get('PEGAPROX_MAX_REQUEST_SIZE', 10 * 1024 * 1024)) # 10 MB default for API
|
|
_upload_max = int(os.environ.get('PEGAPROX_MAX_UPLOAD_SIZE', 100 * 1024 * 1024 * 1024)) # MK: 100 GB for uploads (#116)
|
|
app.config['MAX_CONTENT_LENGTH'] = _upload_max # set high, we check per-route below
|
|
|
|
# Request validation & rate limiting
|
|
# LW: Mar 2026 - ACME HTTP-01 challenge route, must be unauthenticated (#96)
|
|
@app.route('/.well-known/acme-challenge/<token>')
|
|
def acme_challenge(token):
|
|
from pegaprox.core.acme import get_challenge_response
|
|
response = get_challenge_response(token)
|
|
if response:
|
|
return response, 200, {'Content-Type': 'text/plain'}
|
|
return '', 404
|
|
|
|
@app.before_request
|
|
def validate_request():
|
|
if request.path.startswith('/static/') or request.path.startswith('/images/'):
|
|
return None
|
|
if request.path.startswith('/ws'):
|
|
return None
|
|
# MK: Mar 2026 - ACME challenges must bypass all security checks (#96)
|
|
if request.path.startswith('/.well-known/'):
|
|
return None
|
|
|
|
# NS: Feb 2026 - per-route size limits: uploads get the big limit, everything else 10MB
|
|
# MK: Mar 2026 - removed global config mutation, was causing 413s on subsequent uploads (#119)
|
|
is_upload = request.path.endswith('/upload')
|
|
max_size = _upload_max if is_upload else _default_max
|
|
if request.content_length and request.content_length > max_size:
|
|
return jsonify({'error': f'Request too large. Max {max_size // (1024*1024)} MB'}), 413
|
|
# H-6 (security audit): a chunked Transfer-Encoding request carries NO
|
|
# Content-Length, so the check above is skipped and an unauth client could
|
|
# stream an unbounded body → OOM DoS. Pin werkzeug's per-request cap so the
|
|
# limit is enforced when the body is actually read (counts real bytes,
|
|
# works for chunked too). Per-request, not a global config mutation — so
|
|
# uploads keep their big ceiling without the #119 cross-request race.
|
|
try:
|
|
request.max_content_length = max_size
|
|
except Exception:
|
|
pass
|
|
|
|
if request.path.startswith('/api/'):
|
|
skip_paths = ['/api/auth/login', '/api/auth/check', '/api/events', '/api/health', '/api/sse',
|
|
'/api/vmware/migrations']
|
|
if not any(request.path.startswith(p) for p in skip_paths):
|
|
# NS: Mar 2026 - use centralized get_client_ip, respects trusted_proxies
|
|
from pegaprox.utils.audit import get_client_ip
|
|
client_ip = get_client_ip()
|
|
|
|
if not _check_api_rate_limit(client_ip):
|
|
logging.warning(f"Rate limit exceeded for {client_ip}")
|
|
return jsonify({
|
|
'error': 'Rate limit exceeded. Please slow down.',
|
|
'retry_after': API_RATE_WINDOW
|
|
}), 429
|
|
|
|
if request.method in ['POST', 'PUT', 'PATCH'] and request.content_length:
|
|
content_type = request.content_type or ''
|
|
allowed_types = ['application/json', 'multipart/form-data', 'application/x-www-form-urlencoded']
|
|
if not any(t in content_type for t in allowed_types):
|
|
if request.content_length > 0:
|
|
return jsonify({'error': 'Invalid Content-Type'}), 415
|
|
|
|
# NS: Mar 2026 — CSRF check for multipart uploads.
|
|
# MK May 2026 (audit fix H-1) — also enforced for application/json
|
|
# POST/PUT/PATCH/DELETE. Earlier the assumption was "JSON triggers
|
|
# CORS preflight which blocks cross-origin", which is true for the
|
|
# browser path but doesn't help against subdomain takeover, mis-
|
|
# configured trusted-proxy reflecting Origin, or non-browser tools
|
|
# that already have a session cookie. So now: every state-changing
|
|
# /api/* request must come with X-Requested-With or a matching Origin.
|
|
# Exempt: unauth flows (login, OIDC redirects) where we have no
|
|
# session yet to protect.
|
|
_CSRF_EXEMPT = (
|
|
'/api/auth/login',
|
|
'/api/auth/setup', # MK May 2026 — first-run wizard, no session yet
|
|
'/api/auth/oidc/authorize',
|
|
'/api/auth/oidc/callback',
|
|
'/api/auth/oidc/config',
|
|
'/api/auth/check',
|
|
'/api/auth/validate',
|
|
'/api/auth/logout', # logout is idempotent + harmless
|
|
'/api/health',
|
|
'/api/webauthn/auth/begin',
|
|
'/api/webauthn/auth/finish',
|
|
)
|
|
if (request.method in ('POST', 'PUT', 'PATCH', 'DELETE')
|
|
and request.path.startswith('/api/')
|
|
and request.path not in _CSRF_EXEMPT):
|
|
# NS Jul 2026 (CodeAnt CSRF) — the CSRF check must run for EVERY state-changing
|
|
# non-exempt /api/* request, not only JSON/form bodies: a cross-site form with
|
|
# enctype=text/plain is a browser "simple request" that previously skipped this gate.
|
|
if True:
|
|
has_xhr = request.headers.get('X-Requested-With') == 'XMLHttpRequest'
|
|
origin = request.headers.get('Origin', '')
|
|
referer = request.headers.get('Referer', '')
|
|
allowed_origins = get_allowed_origins() or []
|
|
fwd_host = request.headers.get('X-Forwarded-Host', '')
|
|
|
|
# NS May 2026 (#382 follow-up) — safer Origin matcher.
|
|
# The previous version used `value.startswith(f"{scheme}://{host}")`
|
|
# which (a) had a suffix-confusion bug — `https://pegaprox.com`
|
|
# would match an Origin of `https://pegaprox.com.attacker.com`
|
|
# because that string really does start with the substring —
|
|
# and (b) was strict about scheme, which broke users behind
|
|
# Apache/nginx reverse proxies that don't forward
|
|
# X-Forwarded-Proto (cklabautermann's report).
|
|
# New approach: parse the URL, compare *hostname* (and port
|
|
# if both sides specify one). Scheme is irrelevant for CSRF;
|
|
# the browser controls Origin and won't lie about hostname.
|
|
# HTTPS enforcement happens elsewhere (HSTS, secure cookie flag).
|
|
from urllib.parse import urlparse
|
|
|
|
def _collapse_folded_host(hp):
|
|
# NS Jul 2026 (#626) — a duplicated `Host` header (e.g. a reverse
|
|
# proxy injecting its own Host on top of one the client already
|
|
# sent) is folded by the WSGI layer into a comma-joined value such
|
|
# as "example.com, example.com". Collapse it ONLY when every part
|
|
# is identical; a value carrying genuinely different hosts is
|
|
# ambiguous/hostile and is left intact so it fails the match below
|
|
# (fail-closed — we never pick one host out of a conflicting set).
|
|
if ',' not in hp:
|
|
return hp
|
|
parts = [p.strip() for p in hp.split(',') if p.strip()]
|
|
# Host names are case-insensitive, so "Example.com, example.com"
|
|
# is still one host — collapse it. Genuinely different hosts
|
|
# (even ignoring case) are left intact and fail the match.
|
|
if parts and all(p.lower() == parts[0].lower() for p in parts):
|
|
return parts[0]
|
|
return hp
|
|
|
|
def _host_port(hp):
|
|
# split request.host or fwd_host into (host, port|None)
|
|
if not hp: return ('', None)
|
|
hp = _collapse_folded_host(hp)
|
|
if ':' in hp:
|
|
h, _, p = hp.rpartition(':')
|
|
try: return (h.lower(), int(p))
|
|
except ValueError: return (hp.lower(), None)
|
|
return (hp.lower(), None)
|
|
|
|
req_host, req_port = _host_port(request.host)
|
|
fwd_h, fwd_p = _host_port(fwd_host)
|
|
|
|
def _origin_ok(value):
|
|
if not value: return False
|
|
if value in allowed_origins:
|
|
return True
|
|
# NS May 2026 (pentest finding) — Python's urlparse silently
|
|
# normalises tabs/whitespace inside the scheme: 'ht\ttp://x'
|
|
# parses as scheme='http'. Browsers never produce that, but
|
|
# an attacker with raw HTTP control could craft it. Lock the
|
|
# scheme prefix down with a strict, byte-exact check before
|
|
# parsing — only the two browser-realistic prefixes pass.
|
|
if not (value.startswith('http://') or value.startswith('https://')):
|
|
return False
|
|
try:
|
|
u = urlparse(value)
|
|
# MK May 2026: u.port can ValueError for malformed authority
|
|
# like "localhost:5000.attacker.com" — guard explicitly.
|
|
try:
|
|
cand_port = u.port
|
|
except (ValueError, TypeError):
|
|
return False
|
|
except Exception:
|
|
return False
|
|
# Defensive: reject userinfo. RFC 6454 origins have no userinfo;
|
|
# `http://evil.com:80@localhost` parses with hostname=localhost,
|
|
# which would otherwise slip through.
|
|
if u.username or u.password:
|
|
return False
|
|
if u.scheme not in ('http', 'https'): # belt + braces
|
|
return False
|
|
if not u.hostname:
|
|
return False
|
|
cand_host = u.hostname.lower()
|
|
# A portless Origin implies its scheme's default port
|
|
# (https -> 443, http -> 80). Comparing that *effective* port
|
|
# (NS Jul 2026, #626 hardening) keeps an https Origin from ever
|
|
# matching a :80 target, while the common reverse-proxy cases
|
|
# (portless Origin vs the site's default-port Host, or vs a
|
|
# proxy-dropped unknown port) still pass. https://host:9999 is
|
|
# never accepted against an unknown-port target.
|
|
eff_cand = cand_port if cand_port is not None else (443 if u.scheme == 'https' else 80)
|
|
# accept against request host or proxy-forwarded host
|
|
targets = [(req_host, req_port)]
|
|
if fwd_h:
|
|
targets.append((fwd_h, fwd_p))
|
|
for t_host, t_port in targets:
|
|
if cand_host != t_host:
|
|
continue
|
|
if t_port is None:
|
|
# target port unknown (proxy dropped it) — accept only a
|
|
# standard :80/:443 origin, never e.g. https://host:9999.
|
|
if eff_cand in (80, 443):
|
|
return True
|
|
elif eff_cand == t_port:
|
|
return True
|
|
# any other combination is a real port mismatch → reject
|
|
return False
|
|
|
|
# accept either a same-origin Origin/Referer OR XHR + same-origin
|
|
# (XHR alone is not enough — fetch() lets attacker set X-R-W on
|
|
# same-origin, but cross-origin requests can also set it freely
|
|
# in non-browser contexts).
|
|
# MK May 2026: Referer parses as a full URL — pass it directly to
|
|
# _origin_ok which now uses urlparse, no manual splitting needed.
|
|
ok_origin = _origin_ok(origin) or _origin_ok(referer)
|
|
if not ok_origin:
|
|
# If neither Origin nor Referer matches, only allow when XHR
|
|
# marker is set AND there's no foreign Origin/Referer.
|
|
if not has_xhr:
|
|
return jsonify({'error': 'CSRF validation failed'}), 403
|
|
if origin and not _origin_ok(origin):
|
|
return jsonify({'error': 'CSRF validation failed'}), 403
|
|
if referer and not _origin_ok(referer):
|
|
return jsonify({'error': 'CSRF validation failed'}), 403
|
|
|
|
return None
|
|
|
|
# Security headers
|
|
@app.after_request
|
|
def add_security_headers(response):
|
|
response.headers['X-Content-Type-Options'] = 'nosniff'
|
|
# NS May 2026 (#381) — relaxed from DENY to SAMEORIGIN so plugins
|
|
# can ship a frontend UI that the dashboard embeds in an iframe tab.
|
|
# Cross-origin clickjacking remains prevented; same-origin embedding
|
|
# is the documented plugin-frontend contract.
|
|
response.headers['X-Frame-Options'] = 'SAMEORIGIN'
|
|
response.headers['X-XSS-Protection'] = '1; mode=block'
|
|
response.headers['Referrer-Policy'] = 'strict-origin-when-cross-origin'
|
|
response.headers['Permissions-Policy'] = 'geolocation=(), microphone=(), camera=()'
|
|
|
|
# MK: Mar 2026 - tightened CSP, removed dead tailwindcss CDN ref (#118)
|
|
# NS May 2026 (#381) — frame-ancestors 'self' instead of 'none' to
|
|
# match the X-Frame-Options switch above. This is the modern equivalent.
|
|
csp = (
|
|
"default-src 'self'; "
|
|
# NS Jul 2026 (CodeAnt config) — 'unsafe-eval' dropped: Babel is pre-compiled at build
|
|
# time (web/Dev/build.sh) and never runs in the browser, so nothing needs eval().
|
|
"script-src 'self' 'unsafe-inline' "
|
|
"https://cdn.jsdelivr.net; "
|
|
"style-src 'self' 'unsafe-inline' "
|
|
"https://fonts.googleapis.com https://cdn.jsdelivr.net; "
|
|
"font-src 'self' data: https://fonts.gstatic.com https://fonts.googleapis.com; "
|
|
"img-src 'self' data: blob:; "
|
|
"connect-src 'self' wss: ws: https://cdn.jsdelivr.net; "
|
|
"frame-ancestors 'self'; "
|
|
"base-uri 'self'; "
|
|
"form-action 'self'"
|
|
)
|
|
response.headers['Content-Security-Policy'] = csp
|
|
|
|
# LW: Mar 2026 - only trust X-Forwarded-Proto from trusted proxies
|
|
from pegaprox.utils.audit import _is_trusted_proxy
|
|
is_https = request.is_secure or (_is_trusted_proxy(request.remote_addr) and request.headers.get('X-Forwarded-Proto') == 'https')
|
|
if is_https:
|
|
response.headers['Strict-Transport-Security'] = 'max-age=31536000; includeSubDomains'
|
|
|
|
# NS: kein Cache fuer API/auth-stuff, sonst leakt session-state via shared
|
|
# caches (browser-cache nach logout, reverse-proxy mit zu generouser
|
|
# cache config, browser-back-button mit cred response). semgrep findung
|
|
# vom 2026-05-06. /static/* darf weiter gecacht werden, das sind die
|
|
# JS-libs.
|
|
path = request.path or ''
|
|
if path.startswith('/api/') or path in ('/', '/portal', '/oidc/callback'):
|
|
# don't override if a route explicitly set its own Cache-Control
|
|
if 'Cache-Control' not in response.headers:
|
|
response.headers['Cache-Control'] = 'no-store, private'
|
|
response.headers['Pragma'] = 'no-cache' # http/1.0 fallback, harmless
|
|
|
|
return response
|
|
|
|
# Register all API blueprints
|
|
register_blueprints(app)
|
|
|
|
# Load enabled plugins
|
|
from pegaprox.api.plugins import load_enabled_plugins
|
|
load_enabled_plugins(app)
|
|
|
|
return app
|
|
|
|
|
|
def _check_api_rate_limit(client_ip: str) -> bool:
|
|
"""Simple sliding window rate limiter."""
|
|
if API_RATE_LIMIT <= 0:
|
|
return True
|
|
|
|
current_time = time.time()
|
|
|
|
with g.api_rate_limit_lock:
|
|
if client_ip not in g.api_request_counts:
|
|
g.api_request_counts[client_ip] = {'count': 1, 'window_start': current_time}
|
|
return True
|
|
|
|
info = g.api_request_counts[client_ip]
|
|
|
|
if current_time - info['window_start'] > API_RATE_WINDOW:
|
|
info['count'] = 1
|
|
info['window_start'] = current_time
|
|
return True
|
|
|
|
if info['count'] >= API_RATE_LIMIT:
|
|
return False
|
|
|
|
info['count'] += 1
|
|
return True
|
|
|
|
|
|
def download_static_files():
|
|
"""Download all required static files for offline operation."""
|
|
import urllib.request
|
|
import re as _re
|
|
|
|
print("=" * 60)
|
|
print("PegaProx Static Files Downloader")
|
|
print("=" * 60)
|
|
print()
|
|
|
|
static_files = {
|
|
'js': [
|
|
('react.production.min.js', 'https://cdn.jsdelivr.net/npm/react@18/umd/react.production.min.js'),
|
|
('react-dom.production.min.js', 'https://cdn.jsdelivr.net/npm/react-dom@18/umd/react-dom.production.min.js'),
|
|
('babel.min.js', 'https://cdn.jsdelivr.net/npm/@babel/standalone@7/babel.min.js'),
|
|
('chart.umd.min.js', 'https://cdn.jsdelivr.net/npm/chart.js@4/dist/chart.umd.min.js'),
|
|
('xterm.min.js', 'https://cdn.jsdelivr.net/npm/xterm@5.3.0/lib/xterm.min.js'),
|
|
('xterm-addon-fit.min.js', 'https://cdn.jsdelivr.net/npm/xterm-addon-fit@0.8.0/lib/xterm-addon-fit.min.js'),
|
|
],
|
|
'css': [
|
|
('xterm.min.css', 'https://cdn.jsdelivr.net/npm/xterm@5.3.0/css/xterm.min.css'),
|
|
]
|
|
}
|
|
|
|
os.makedirs('static/js', exist_ok=True)
|
|
os.makedirs('static/css', exist_ok=True)
|
|
|
|
ctx = ssl.create_default_context() # NS: Feb 2026 - use default SSL verification for downloads
|
|
|
|
success = 0
|
|
failed = 0
|
|
|
|
for subdir, files in static_files.items():
|
|
print(f"Downloading {subdir} files...")
|
|
for filename, url in files:
|
|
dest = f'static/{subdir}/{filename}'
|
|
print(f" {filename}...", end=' ')
|
|
try:
|
|
req = urllib.request.Request(url, headers={
|
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
|
|
})
|
|
with urllib.request.urlopen(req, timeout=30, context=ctx) as response:
|
|
data = response.read()
|
|
with open(dest, 'wb') as f:
|
|
f.write(data)
|
|
print(f"OK ({len(data):,} bytes)")
|
|
success += 1
|
|
except Exception as e:
|
|
print(f"FAILED: {e}")
|
|
failed += 1
|
|
|
|
# MK: Mar 2026 - tailwind.min.css is now a full CLI build, don't overwrite it (#118)
|
|
if os.path.exists('static/css/tailwind.min.css'):
|
|
sz = os.path.getsize('static/css/tailwind.min.css')
|
|
print(f"\n tailwind.min.css already exists ({sz:,} bytes), skipping")
|
|
print(" (rebuild with: npx tailwindcss -i input.css -o static/css/tailwind.min.css --minify)")
|
|
else:
|
|
print("\n WARNING: static/css/tailwind.min.css missing!")
|
|
print(" Run: npx tailwindcss -i input.css -o static/css/tailwind.min.css --minify")
|
|
failed += 1
|
|
|
|
# LW: Mar 2026 - download Google Fonts for offline (#118)
|
|
print("\nDownloading Google Fonts for offline use...")
|
|
os.makedirs('static/fonts', exist_ok=True)
|
|
|
|
_gfonts = {
|
|
'plus-jakarta-sans': {
|
|
'family': 'Plus Jakarta Sans',
|
|
'weights': {
|
|
'400': 'https://fonts.gstatic.com/s/plusjakartasans/v8/LDIbaomQNQcsA88c7O9yZ4KMCoOg4IA6-91aHEjcWuA_KU7NShXUEKi4Rw.woff2',
|
|
'500': 'https://fonts.gstatic.com/s/plusjakartasans/v8/LDIbaomQNQcsA88c7O9yZ4KMCoOg4IA6-91aHEjcWuA_AU7NShXUEKi4Rw.woff2',
|
|
'600': 'https://fonts.gstatic.com/s/plusjakartasans/v8/LDIbaomQNQcsA88c7O9yZ4KMCoOg4IA6-91aHEjcWuA_zUnNShXUEKi4Rw.woff2',
|
|
'700': 'https://fonts.gstatic.com/s/plusjakartasans/v8/LDIbaomQNQcsA88c7O9yZ4KMCoOg4IA6-91aHEjcWuA_9EnNShXUEKi4Rw.woff2',
|
|
'800': 'https://fonts.gstatic.com/s/plusjakartasans/v8/LDIbaomQNQcsA88c7O9yZ4KMCoOg4IA6-91aHEjcWuA_KUnNShXUEKi4Rw.woff2',
|
|
}
|
|
},
|
|
'jetbrains-mono': {
|
|
'family': 'JetBrains Mono',
|
|
'weights': {
|
|
'400': 'https://fonts.gstatic.com/s/jetbrainsmono/v18/tDbY2o-flEEny0FZhsfKu5WU4zr3E_BX0PnT8RD8yKxjPVmUsaaDhw.woff2',
|
|
'500': 'https://fonts.gstatic.com/s/jetbrainsmono/v18/tDbY2o-flEEny0FZhsfKu5WU4zr3E_BX0PnT8RD8-axjPVmUsaaDhw.woff2',
|
|
'600': 'https://fonts.gstatic.com/s/jetbrainsmono/v18/tDbY2o-flEEny0FZhsfKu5WU4zr3E_BX0PnT8RD8FapjPVmUsaaDhw.woff2',
|
|
'700': 'https://fonts.gstatic.com/s/jetbrainsmono/v18/tDbY2o-flEEny0FZhsfKu5WU4zr3E_BX0PnT8RD8LapjPVmUsaaDhw.woff2',
|
|
}
|
|
}
|
|
}
|
|
|
|
font_css = "/* LW: Mar 2026 - local Google Fonts for offline mode (#118) */\n"
|
|
for font_id, font_info in _gfonts.items():
|
|
for weight, url in font_info['weights'].items():
|
|
fname = f"{font_id}-{weight}.woff2"
|
|
dest = f"static/fonts/{fname}"
|
|
print(f" {fname}...", end=' ')
|
|
try:
|
|
req = urllib.request.Request(url, headers={
|
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
|
|
})
|
|
with urllib.request.urlopen(req, timeout=30, context=ctx) as response:
|
|
data = response.read()
|
|
with open(dest, 'wb') as f:
|
|
f.write(data)
|
|
print(f"OK ({len(data):,} bytes)")
|
|
success += 1
|
|
except Exception as e:
|
|
print(f"FAILED: {e}")
|
|
failed += 1
|
|
|
|
font_css += f"""@font-face {{
|
|
font-family: '{font_info['family']}';
|
|
font-style: normal;
|
|
font-weight: {weight};
|
|
font-display: swap;
|
|
src: url('/static/fonts/{fname}') format('woff2');
|
|
}}
|
|
"""
|
|
|
|
try:
|
|
with open('static/css/fonts.css', 'w') as f:
|
|
f.write(font_css)
|
|
print(" fonts.css... OK")
|
|
success += 1
|
|
except Exception as e:
|
|
print(f" fonts.css... FAILED: {e}")
|
|
failed += 1
|
|
|
|
# Download noVNC for offline VNC console
|
|
print("\nDownloading noVNC for offline VNC console...")
|
|
novnc_base = 'https://cdn.jsdelivr.net/npm/@novnc/novnc@1.4.0'
|
|
novnc_files = [
|
|
'core/rfb.js', 'core/display.js', 'core/inflator.js', 'core/deflator.js',
|
|
'core/websock.js', 'core/encodings.js', 'core/des.js', 'core/ra2.js', 'core/base64.js',
|
|
'core/decoders/copyrect.js', 'core/decoders/hextile.js', 'core/decoders/raw.js',
|
|
'core/decoders/rre.js', 'core/decoders/tight.js', 'core/decoders/tightpng.js',
|
|
'core/decoders/zrle.js', 'core/decoders/jpeg.js',
|
|
'core/input/keyboard.js', 'core/input/keysym.js', 'core/input/keysymdef.js',
|
|
'core/input/gesturehandler.js', 'core/input/domkeytable.js', 'core/input/util.js',
|
|
'core/input/vkeys.js', 'core/input/xtscancodes.js', 'core/input/fixedkeys.js',
|
|
'core/util/browser.js', 'core/util/cursor.js', 'core/util/element.js',
|
|
'core/util/events.js', 'core/util/eventtarget.js', 'core/util/int.js',
|
|
'core/util/logging.js', 'core/util/strings.js', 'core/util/md5.js',
|
|
'vendor/pako/lib/zlib/inflate.js', 'vendor/pako/lib/zlib/zstream.js',
|
|
'vendor/pako/lib/zlib/deflate.js', 'vendor/pako/lib/zlib/messages.js',
|
|
'vendor/pako/lib/zlib/trees.js', 'vendor/pako/lib/zlib/adler32.js',
|
|
'vendor/pako/lib/zlib/crc32.js', 'vendor/pako/lib/zlib/inffast.js',
|
|
'vendor/pako/lib/zlib/inftrees.js', 'vendor/pako/lib/utils/common.js',
|
|
]
|
|
|
|
for subdir in ['core', 'core/decoders', 'core/input', 'core/util',
|
|
'vendor/pako/lib/zlib', 'vendor/pako/lib/utils']:
|
|
os.makedirs(f'static/js/novnc/{subdir}', exist_ok=True)
|
|
|
|
novnc_success = 0
|
|
novnc_failed = 0
|
|
|
|
for filepath in novnc_files:
|
|
url = f"{novnc_base}/{filepath}"
|
|
dest = f"static/js/novnc/{filepath}"
|
|
filename = filepath.split('/')[-1]
|
|
print(f" {filename}...", end=' ')
|
|
try:
|
|
req = urllib.request.Request(url, headers={
|
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36'
|
|
})
|
|
with urllib.request.urlopen(req, timeout=30, context=ctx) as response:
|
|
content = response.read().decode('utf-8')
|
|
|
|
file_dir = '/'.join(filepath.split('/')[:-1])
|
|
pattern = r'''from\s+(['"])(\.{1,2}/[^'"]+)\1'''
|
|
|
|
def rewrite_import(match):
|
|
quote = match.group(1)
|
|
rel_path = match.group(2)
|
|
if rel_path.startswith('./'):
|
|
resolved = f"/static/js/novnc/{file_dir}/{rel_path[2:]}"
|
|
elif rel_path.startswith('../'):
|
|
parts = file_dir.split('/') if file_dir else []
|
|
rest = rel_path
|
|
while rest.startswith('../'):
|
|
if parts:
|
|
parts.pop()
|
|
rest = rest[3:]
|
|
parent = '/'.join(parts)
|
|
resolved = f"/static/js/novnc/{parent}/{rest}" if parent else f"/static/js/novnc/{rest}"
|
|
else:
|
|
resolved = rel_path
|
|
while '//' in resolved:
|
|
resolved = resolved.replace('//', '/')
|
|
return f"from {quote}{resolved}{quote}"
|
|
|
|
content = _re.sub(pattern, rewrite_import, content)
|
|
|
|
with open(dest, 'w') as f:
|
|
f.write(content)
|
|
print("OK")
|
|
novnc_success += 1
|
|
success += 1
|
|
except Exception as e:
|
|
print(f"FAILED: {e}")
|
|
novnc_failed += 1
|
|
failed += 1
|
|
|
|
rfb_entry = '''// noVNC entry point for PegaProx offline mode
|
|
// Auto-generated by --download-static
|
|
export { default } from '/static/js/novnc/core/rfb.js';
|
|
export * from '/static/js/novnc/core/rfb.js';
|
|
'''
|
|
try:
|
|
with open('static/js/novnc/rfb.min.js', 'w') as f:
|
|
f.write(rfb_entry)
|
|
print(" rfb.min.js (entry point)... OK")
|
|
success += 1
|
|
except Exception as e:
|
|
print(f" rfb.min.js... FAILED: {e}")
|
|
failed += 1
|
|
|
|
print(f"\n noVNC: {novnc_success}/{len(novnc_files)} files downloaded")
|
|
print()
|
|
print("=" * 60)
|
|
print(f"Done: {success} succeeded, {failed} failed")
|
|
print("=" * 60)
|
|
|
|
if failed == 0:
|
|
print("\nAll static files downloaded!")
|
|
print(" PegaProx can run fully offline now (including VNC console)")
|
|
else:
|
|
print("\nSome downloads failed, will use CDN fallback")
|
|
|
|
return failed == 0
|
|
|
|
|
|
def _path_diagnostics(path):
|
|
"""owner/mode of a path and of its parent, plus who we are.
|
|
|
|
#633: the operator needs to compare the two. The whole bug was a cert the
|
|
service user could not reach, and the log said nothing about who owned it.
|
|
"""
|
|
try:
|
|
import pwd
|
|
import grp
|
|
except ImportError: # non-POSIX, ids only
|
|
pwd = grp = None
|
|
|
|
def _name(getter, attr, num):
|
|
if getter is None:
|
|
return str(num)
|
|
try:
|
|
return getattr(getter(num), attr)
|
|
except (KeyError, OSError):
|
|
return str(num)
|
|
|
|
lines = []
|
|
target = os.path.abspath(path)
|
|
for p in (target, os.path.dirname(target)):
|
|
try:
|
|
st = os.stat(p)
|
|
lines.append(" %s owner=%s:%s mode=0o%03o" % (
|
|
p,
|
|
_name(pwd and pwd.getpwuid, 'pw_name', st.st_uid),
|
|
_name(grp and grp.getgrgid, 'gr_name', st.st_gid),
|
|
stat.S_IMODE(st.st_mode)))
|
|
except OSError as e:
|
|
lines.append(" %s cannot stat: %s" % (p, e.strerror))
|
|
lines.append(" this process uid=%s(%s) gid=%s(%s)" % (
|
|
os.geteuid(), _name(pwd and pwd.getpwuid, 'pw_name', os.geteuid()),
|
|
os.getegid(), _name(grp and grp.getgrgid, 'gr_name', os.getegid())))
|
|
return "\n".join(lines)
|
|
|
|
|
|
def _tls_setup_failed(reason, path):
|
|
"""Fail closed, or return None if plaintext was asked for explicitly.
|
|
|
|
#633: this used to print a WARNING and fall through with ssl_context=None,
|
|
which bound cleartext HTTP on the port that was meant to be TLS. TLS clients
|
|
then got "Invalid http version: '\\x16\\x03\\x01...'" and the dashboard was
|
|
down while the service looked healthy. A downgrade has to be a decision, not
|
|
an accident.
|
|
"""
|
|
detail = "TLS is enabled but there is no usable certificate: %s\n%s" % (
|
|
reason, _path_diagnostics(path))
|
|
if os.environ.get('PEGAPROX_ALLOW_PLAINTEXT', '').strip().lower() in ('1', 'true', 'yes', 'on'):
|
|
logging.getLogger(__name__).error(
|
|
"%s\n PEGAPROX_ALLOW_PLAINTEXT is set, so serving PLAINTEXT HTTP on the "
|
|
"TLS port anyway - every TLS client will fail against it.", detail)
|
|
return None
|
|
raise SystemExit(
|
|
"%s\n Refusing to serve plaintext on the TLS port. Fix the certificate, or put "
|
|
"PegaProx behind a reverse proxy (the reverse_proxy setting), or set "
|
|
"PEGAPROX_ALLOW_PLAINTEXT=1 to serve cleartext there on purpose." % detail)
|
|
|
|
|
|
def _unreadable(path):
|
|
"""The OSError from opening path, or None if it opens fine.
|
|
|
|
os.path.exists() is not enough: it is False for a cert in a directory we
|
|
cannot search, which is how EACCES ended up in the "no certificates" branch
|
|
and got a valid cert overwritten (#633).
|
|
"""
|
|
try:
|
|
with open(path, 'rb'):
|
|
return None
|
|
except OSError as e:
|
|
return e
|
|
|
|
|
|
def _unloadable(cert_file, key_file):
|
|
"""The exception from actually loading the cert+key into an SSL context, or
|
|
None if the pair parses.
|
|
|
|
MK #633 follow-up (adversarial review of #637): readable != loadable. A
|
|
truncated, empty, corrupt or mismatched cert/key opens fine (so _unreadable
|
|
says it's ok) but blows up later at ssl_ctx.load_cert_chain() as an UNCAUGHT
|
|
ssl.SSLError - a raw traceback instead of the actionable message, for what is
|
|
the most common non-permission cert failure. Load it here so a bad-but-present
|
|
pair routes through the same fail-closed path as a missing one.
|
|
"""
|
|
try:
|
|
ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER).load_cert_chain(cert_file, key_file)
|
|
return None
|
|
except (ssl.SSLError, OSError, ValueError) as e:
|
|
return e
|
|
|
|
|
|
def _generate_self_signed(cert_file, key_file, domain, app_name):
|
|
from OpenSSL import crypto
|
|
key = crypto.PKey()
|
|
key.generate_key(crypto.TYPE_RSA, 2048)
|
|
cert = crypto.X509()
|
|
cert.get_subject().C = "DE"
|
|
cert.get_subject().ST = "State"
|
|
cert.get_subject().L = "City"
|
|
cert.get_subject().O = app_name or "PegaProx"
|
|
cert.get_subject().OU = app_name or "PegaProx"
|
|
cert.get_subject().CN = domain or app_name or "PegaProx"
|
|
cert.set_serial_number(1000)
|
|
cert.gmtime_adj_notBefore(0)
|
|
cert.gmtime_adj_notAfter(365 * 24 * 60 * 60)
|
|
cert.set_issuer(cert.get_subject())
|
|
cert.set_pubkey(key)
|
|
cert.sign(key, 'sha256')
|
|
with open(cert_file, "wb") as f:
|
|
f.write(crypto.dump_certificate(crypto.FILETYPE_PEM, cert))
|
|
with open(key_file, "wb") as f:
|
|
f.write(crypto.dump_privatekey(crypto.FILETYPE_PEM, key))
|
|
os.chmod(key_file, 0o600)
|
|
|
|
|
|
def _resolve_ssl_context(reverse_proxy, domain='', app_name='PegaProx',
|
|
cert_file=SSL_CERT_FILE, key_file=SSL_KEY_FILE):
|
|
"""(cert, key) for the TLS listener, or None for plaintext.
|
|
|
|
Posture (#633): TLS unless a reverse proxy terminates it for us. If TLS is
|
|
the posture and we cannot load or generate a usable pair, we do not come up.
|
|
"""
|
|
if reverse_proxy:
|
|
return None # nginx/haproxy/traefik owns TLS, plain HTTP on the bind
|
|
|
|
cert_err, key_err = _unreadable(cert_file), _unreadable(key_file)
|
|
if cert_err is None and key_err is None:
|
|
# MK #633 follow-up: readable isn't enough - a corrupt/mismatched pair has
|
|
# to fail here with a clear reason, not crash later at load_cert_chain.
|
|
bad = _unloadable(cert_file, key_file)
|
|
if bad is not None:
|
|
return _tls_setup_failed(
|
|
"%s and %s are present but do not load as a cert/key pair: %s"
|
|
% (cert_file, key_file, bad), cert_file)
|
|
print("SSL certificates found - starting with HTTPS")
|
|
return (cert_file, key_file)
|
|
|
|
# Present but not readable is NOT "missing". Never generate over it - that
|
|
# would report the wrong problem, and destroy a working cert if the write
|
|
# happened to succeed.
|
|
for path, err in ((cert_file, cert_err), (key_file, key_err)):
|
|
if err is not None and err.errno != errno.ENOENT:
|
|
return _tls_setup_failed("cannot read %s: %s" % (path, err.strerror), path)
|
|
|
|
# MK #633 follow-up: only a genuinely empty pair (BOTH missing) may be
|
|
# generated. If one half is present and only the other is ENOENT, generation
|
|
# would clobber the surviving half - so fail closed and name the missing one.
|
|
if cert_err is None or key_err is None:
|
|
present, missing = (cert_file, key_file) if key_err is not None else (key_file, cert_file)
|
|
return _tls_setup_failed(
|
|
"%s is present but its counterpart %s is missing - refusing to regenerate "
|
|
"over the existing half" % (present, missing), missing)
|
|
|
|
# MK 2026-06-08 (#531): generate into the persisted config/ssl dir, and
|
|
# create the dir first - on a fresh container the old target did not exist,
|
|
# generation failed with ENOENT and we fell back to plain HTTP.
|
|
print("No SSL certificates found. Generating self-signed certificate...")
|
|
try:
|
|
os.makedirs(os.path.dirname(os.path.abspath(cert_file)), exist_ok=True)
|
|
_generate_self_signed(cert_file, key_file, domain, app_name)
|
|
except ImportError:
|
|
return _tls_setup_failed("pyOpenSSL is not installed (pip install pyOpenSSL)", cert_file)
|
|
except Exception as e:
|
|
return _tls_setup_failed("could not generate one: %s" % e, cert_file)
|
|
bad = _unloadable(cert_file, key_file)
|
|
if bad is not None: # a generator that emits a pair we can't load is a bug, not plaintext
|
|
return _tls_setup_failed("generated a certificate that will not load: %s" % bad, cert_file)
|
|
print("Self-signed certificate generated: %s" % cert_file)
|
|
return (cert_file, key_file)
|
|
|
|
|
|
def main(debug_mode=False):
|
|
"""Main entry point - starts PegaProx server."""
|
|
from pegaprox.utils.auth import load_users, load_sessions, backfill_initialized_marker, is_initialized
|
|
from pegaprox.utils.audit import load_audit_log
|
|
from pegaprox.core.config import load_config
|
|
from pegaprox.core.pbs import load_pbs_servers
|
|
from pegaprox.core.vmware import load_vmware_servers
|
|
from pegaprox.models.tasks import PegaProxConfig
|
|
from pegaprox.core.manager import PegaProxManager
|
|
from pegaprox.background.broadcast import start_broadcast_thread
|
|
from pegaprox.background.alerts import start_alert_thread
|
|
from pegaprox.background.scheduler import start_scheduler_thread
|
|
from pegaprox.background.password_expiry import start_password_expiry_thread
|
|
from pegaprox.background.cross_cluster_lb import start_cross_cluster_lb_thread
|
|
from pegaprox.background.cross_cluster_replication import start_cross_cluster_replication_thread
|
|
from pegaprox.background.syslog_server import start_syslog_server
|
|
from pegaprox.api.schedules import start_scheduler as start_actions_scheduler
|
|
from pegaprox.api.helpers import load_server_settings, acme_dns_config_from_settings
|
|
from pegaprox.utils.rbac import get_pool_membership_cache
|
|
from pegaprox.constants import AUDIT_RETENTION_DAYS
|
|
|
|
# Initialize SSH semaphore
|
|
g.init_ssh_semaphore(SSH_MAX_CONCURRENT)
|
|
|
|
# Configure logging
|
|
# MK May 2026 (#357): env-var override (PEGAPROX_LOG_LEVEL) wins over default
|
|
# but --debug still forces DEBUG so the troubleshooting path doesn't need an
|
|
# extra knob. Unset env + no --debug → previous WARNING default.
|
|
from pegaprox.constants import LOG_LEVEL as _ENV_LOG_LEVEL
|
|
if debug_mode:
|
|
log_level = logging.DEBUG
|
|
elif _ENV_LOG_LEVEL is not None:
|
|
log_level = _ENV_LOG_LEVEL
|
|
else:
|
|
log_level = logging.WARNING
|
|
logging.basicConfig(
|
|
level=log_level,
|
|
format='%(asctime)s [%(name)s] %(levelname)s: %(message)s' if debug_mode else '%(message)s',
|
|
datefmt='%Y-%m-%d %H:%M:%S'
|
|
)
|
|
|
|
if not debug_mode:
|
|
logging.getLogger('werkzeug').setLevel(logging.ERROR)
|
|
logging.getLogger('gevent').setLevel(logging.ERROR)
|
|
logging.getLogger('urllib3').setLevel(logging.ERROR)
|
|
|
|
if debug_mode:
|
|
print("=" * 50)
|
|
print("DEBUG MODE ENABLED")
|
|
print("=" * 50)
|
|
|
|
# Check optional libraries
|
|
print("\nChecking optional libraries...")
|
|
missing_libs = []
|
|
try:
|
|
import websockets
|
|
print(" ✓ websockets (VNC/SSH console)")
|
|
except ImportError:
|
|
missing_libs.append('websockets')
|
|
print(" ✗ websockets - VNC/SSH console will NOT work!")
|
|
|
|
try:
|
|
import paramiko
|
|
print(" ✓ paramiko (SSH features)")
|
|
except ImportError:
|
|
missing_libs.append('paramiko')
|
|
print(" ✗ paramiko - SSH features disabled")
|
|
|
|
GEVENT_AVAILABLE = False
|
|
try:
|
|
from gevent.pywsgi import WSGIServer
|
|
GEVENT_AVAILABLE = True
|
|
print(" ✓ gevent (high performance)")
|
|
except ImportError:
|
|
print(" ✗ gevent - using Flask dev server (slower)")
|
|
|
|
ARGON2_AVAILABLE = False
|
|
try:
|
|
import argon2
|
|
ARGON2_AVAILABLE = True
|
|
print(" ✓ argon2-cffi (secure password hashing)")
|
|
except ImportError:
|
|
print(" ⚠ argon2-cffi - using PBKDF2 fallback")
|
|
|
|
try:
|
|
import XenAPI
|
|
print(" ✓ XenAPI (XCP-ng integration)")
|
|
except ImportError:
|
|
print(" ✗ XenAPI - XCP-ng clusters disabled (pip install XenAPI)")
|
|
|
|
if missing_libs:
|
|
print(f"\n Install missing: pip install {' '.join(missing_libs)}")
|
|
print()
|
|
|
|
# NS May 2026 — auto-encrypt every DB under CONFIG_DIR on first boot if
|
|
# sqlcipher3 is available. MUST run BEFORE create_app(): plugin loader +
|
|
# push-inbox initialiser open DB connections inside create_app(), so by
|
|
# the time we'd hit the post-create_app point it's already too late and
|
|
# SQLCipher fails the PRAGMA key handshake against a plain file.
|
|
# Idempotent on subsequent boots (state == 'encrypted' short-circuits).
|
|
#
|
|
# Covers pegaprox.db (main) + syslog.db (Apr 2026 syslog server) — both
|
|
# get opened through dbcrypto.connect() so both must be encrypted in lock-
|
|
# step. ensure_db_encrypted() handles 'missing' state cleanly (noop) so
|
|
# only-pegaprox-no-syslog deployments are fine.
|
|
try:
|
|
from pegaprox.core import dbcrypto as _dbcrypto
|
|
from pegaprox.constants import CONFIG_DIR as _CFGDIR
|
|
for _db_name in ('pegaprox.db', 'syslog.db'):
|
|
_db_path = os.path.join(_CFGDIR, _db_name)
|
|
_r = _dbcrypto.ensure_db_encrypted(_db_path)
|
|
if _r.get('action') == 'migrated':
|
|
print(f" ✓ {_db_name} auto-encrypted ({_r['rows_copied']} rows, {_r['duration_s']}s)")
|
|
print(f" backup: {_r['backup_path']}")
|
|
elif _r.get('action') == 'no-backend':
|
|
print(f" ⚠ sqlcipher3 not installed — {_db_name} stays plain (field-level Fernet still active)")
|
|
except RuntimeError as _e:
|
|
# corrupt / unknown-key — refuse to start
|
|
print(f"\n[FATAL] {_e}\n")
|
|
sys.exit(1)
|
|
except Exception as _e:
|
|
# don't take down boot for a non-fatal hiccup — log and continue
|
|
logging.error(f"[DBCRYPTO] auto-encrypt check failed: {_e}", exc_info=True)
|
|
|
|
# Create Flask app (plugins + push inbox will hit the DB here)
|
|
app = create_app()
|
|
|
|
# Init user system
|
|
print("Initializing user system...")
|
|
g.users_db = load_users()
|
|
print(f"Loaded {len(g.users_db)} users")
|
|
|
|
# Init audit log
|
|
print("Initializing audit log...")
|
|
load_audit_log()
|
|
print(f"Loaded {len(g.audit_log)} audit entries (retention: {AUDIT_RETENTION_DAYS} days)")
|
|
|
|
# Load sessions
|
|
print("Loading sessions...")
|
|
load_sessions()
|
|
print(f"Loaded {len(g.active_sessions)} active sessions")
|
|
|
|
# MK May 2026 — backfill initialized marker for upgrades from pre-setup-wizard
|
|
# builds (pre-init installs already had users; we just stamp the marker so the
|
|
# /login path's is_initialized() doesn't fall through to NOT_INITIALIZED).
|
|
backfill_initialized_marker()
|
|
|
|
if not is_initialized():
|
|
print("\n" + "=" * 50)
|
|
print("FIRST-RUN SETUP REQUIRED")
|
|
print(" No admin account exists yet — open the PegaProx URL")
|
|
print(" in a browser to create the first administrator via the")
|
|
print(" setup wizard. /api/auth/login is disabled until that")
|
|
print(" is done.")
|
|
print("=" * 50 + "\n")
|
|
|
|
# Load existing configuration
|
|
config = load_config()
|
|
|
|
# Start managers for existing clusters
|
|
for cluster_id, cluster_data in config.items():
|
|
config_obj = PegaProxConfig(cluster_data)
|
|
ctype = cluster_data.get('cluster_type', 'proxmox')
|
|
if ctype == 'xcpng':
|
|
from pegaprox.core.xcpng import XcpngManager
|
|
manager = XcpngManager(cluster_id, config_obj)
|
|
manager.start()
|
|
g.cluster_managers[cluster_id] = manager
|
|
print(f"Started XCP-ng manager for pool: {cluster_data['name']}")
|
|
else:
|
|
manager = PegaProxManager(cluster_id, config_obj)
|
|
manager.start()
|
|
g.cluster_managers[cluster_id] = manager
|
|
print(f"Started PegaProx manager for cluster: {cluster_data['name']}")
|
|
|
|
# Start background threads
|
|
start_broadcast_thread()
|
|
print("Started WebSocket live updates broadcast thread")
|
|
|
|
try:
|
|
load_pbs_servers()
|
|
except Exception as e:
|
|
logging.warning(f"Failed to load PBS servers at startup: {e}")
|
|
|
|
try:
|
|
load_vmware_servers()
|
|
# NS: register ESXi hosts as XHM-capable clusters
|
|
from pegaprox.core.esxi_cluster import ESXiClusterManager
|
|
for vmw_id, vmw_mgr in g.vmware_managers.items():
|
|
if getattr(vmw_mgr, 'server_type', '') == 'esxi':
|
|
g.cluster_managers[vmw_id] = ESXiClusterManager(vmw_id, vmw_mgr)
|
|
logging.info(f"Registered ESXi host '{vmw_mgr.name}' as XHM cluster {vmw_id}")
|
|
except Exception as e:
|
|
logging.warning(f"Failed to load VMware servers at startup: {e}")
|
|
|
|
start_alert_thread()
|
|
print("Started alert monitoring thread")
|
|
|
|
start_scheduler_thread()
|
|
print("Started task scheduler thread")
|
|
|
|
# NS: Mar 2026 - the scheduled_actions scheduler (UI-created schedules, #134)
|
|
# background/scheduler.py only handles the old scheduled_tasks table
|
|
start_actions_scheduler()
|
|
print("Started scheduled actions thread")
|
|
|
|
start_password_expiry_thread()
|
|
print("Started password expiry check thread")
|
|
|
|
start_cross_cluster_lb_thread()
|
|
print("Started cross-cluster load balancer thread")
|
|
|
|
start_cross_cluster_replication_thread()
|
|
print("Started cross-cluster replication scheduler thread")
|
|
|
|
try:
|
|
start_syslog_server()
|
|
print("Started integrated syslog server")
|
|
except Exception as e:
|
|
logging.warning(f"Syslog server failed to start: {e}")
|
|
|
|
# #238: reset stuck DR plans from a previous crash/restart
|
|
try:
|
|
from datetime import datetime as _dt
|
|
from pegaprox.core.db import get_db
|
|
_db = get_db()
|
|
stuck = _db.query("SELECT id, name FROM site_recovery_plans WHERE status IN ('running', 'testing')")
|
|
for p in (stuck or []):
|
|
_db.execute("UPDATE site_recovery_plans SET status = 'failed', updated_at = ? WHERE id = ?",
|
|
(_dt.now().isoformat(), p['id']))
|
|
print(f" Reset stuck DR plan '{p['name']}' → failed")
|
|
except Exception as e:
|
|
print(f" DR plan reset check failed: {e}")
|
|
|
|
from pegaprox.background.site_recovery import start_heartbeat
|
|
start_heartbeat()
|
|
print("Started site recovery heartbeat monitor")
|
|
|
|
# Start plugin background tasks
|
|
from pegaprox.api.plugins import start_plugin_backgrounds
|
|
start_plugin_backgrounds()
|
|
|
|
# Warm up pool cache
|
|
def warmup_pool_cache():
|
|
time.sleep(5)
|
|
for cluster_id in g.cluster_managers:
|
|
try:
|
|
get_pool_membership_cache(cluster_id)
|
|
print(f" Pool cache warmed for cluster: {cluster_id}")
|
|
except Exception as e:
|
|
print(f" Warning: Could not warm pool cache for {cluster_id}: {e}")
|
|
|
|
threading.Thread(target=warmup_pool_cache, daemon=True).start()
|
|
print("Started pool cache warmup thread")
|
|
|
|
# MK: Mar 2026 - ACME auto-renewal thread (#96)
|
|
def acme_renewal_loop():
|
|
time.sleep(30) # wait for server to fully start
|
|
while True:
|
|
try:
|
|
_settings = load_server_settings()
|
|
if _settings.get('acme_enabled') and _settings.get('domain'):
|
|
from pegaprox.core.acme import check_and_renew
|
|
if Path("/usr/lib/pegaprox").exists():
|
|
_ssl = str(Path("/var/lib/pegaprox/ssl"))
|
|
else:
|
|
_ssl = str(Path(__file__).resolve().parent.parent / 'ssl')
|
|
_challenge_type = _settings.get('acme_challenge_type') or 'http-01'
|
|
_dns_provider = _settings.get('acme_dns_provider') or 'manual'
|
|
renewed = check_and_renew(
|
|
_settings['domain'], _settings.get('acme_email', ''),
|
|
_ssl, staging=_settings.get('acme_staging', False),
|
|
directory_url=_settings.get('acme_directory_url', ''),
|
|
challenge_type=_challenge_type,
|
|
dns_provider=_dns_provider,
|
|
dns_config=acme_dns_config_from_settings(_settings)
|
|
)
|
|
if renewed:
|
|
logging.info("[ACME] Certificate renewed, restart required for new cert")
|
|
except Exception as e:
|
|
logging.debug(f"[ACME] Renewal check error: {e}")
|
|
time.sleep(86400) # check once per day
|
|
|
|
threading.Thread(target=acme_renewal_loop, daemon=True).start()
|
|
print("Started ACME auto-renewal thread")
|
|
|
|
# Load server settings
|
|
server_settings = load_server_settings()
|
|
port = server_settings.get('port', 5000)
|
|
bind_host = os.environ.get('PEGAPROX_HOST')
|
|
|
|
# NS Mar 2026 - reverse proxy mode: skip SSL, bind localhost, trust proxy headers
|
|
reverse_proxy = server_settings.get('reverse_proxy_enabled', False)
|
|
if os.environ.get('PEGAPROX_BEHIND_PROXY', '').lower() in ('1', 'true', 'yes'):
|
|
reverse_proxy = True
|
|
|
|
# load trusted proxy IPs for X-Forwarded-For (loopback always trusted)
|
|
from pegaprox.utils.audit import load_trusted_proxies
|
|
trusted = os.environ.get('PEGAPROX_TRUSTED_PROXIES', '') or server_settings.get('trusted_proxies', '')
|
|
load_trusted_proxies(trusted)
|
|
if trusted:
|
|
print(f"Trusted proxies: {trusted}")
|
|
|
|
if not bind_host:
|
|
if reverse_proxy:
|
|
custom_bind = server_settings.get('proxy_bind_address', '').strip()
|
|
if custom_bind:
|
|
bind_host = custom_bind
|
|
print(f"Reverse proxy mode — custom bind: {bind_host}")
|
|
else:
|
|
bind_host = '127.0.0.1'
|
|
print("Reverse proxy mode — binding to 127.0.0.1 only")
|
|
elif _test_ipv6_available():
|
|
bind_host = '::'
|
|
print("IPv6 available — binding dual-stack (::)")
|
|
else:
|
|
bind_host = '0.0.0.0'
|
|
print("IPv6 not available — binding IPv4 only (0.0.0.0)")
|
|
else:
|
|
if ':' in bind_host and not _test_ipv6_available():
|
|
print(f"WARNING: IPv6 bind address '{bind_host}' requested but IPv6 not available")
|
|
print("Falling back to 0.0.0.0")
|
|
bind_host = '0.0.0.0'
|
|
|
|
# MK: when behind proxy, SSL is handled by nginx/haproxy - we run plain HTTP
|
|
#
|
|
# MK Aug 2026: we deliberately do NOT gate on the `ssl_enabled` setting here.
|
|
# The pre-#633 code only ever used it as a fast path - its else-branch generated
|
|
# a cert and served HTTPS regardless - so with the toggle off (which is the
|
|
# api/helpers.py default) every existing install is in fact running TLS. Honouring
|
|
# it now would silently downgrade all of them to cleartext on upgrade. Posture
|
|
# stays "TLS unless a reverse proxy terminates it"; wiring the toggle up properly
|
|
# is its own change (#638).
|
|
domain = server_settings.get('domain', '')
|
|
app_name = server_settings.get('app_name', 'PegaProx')
|
|
if reverse_proxy:
|
|
print("SSL disabled (handled by reverse proxy)")
|
|
|
|
# Check for SSL certificates (skipped entirely behind a reverse proxy).
|
|
# MK Aug 2026 (#633): this fails closed now - see _resolve_ssl_context(). It used
|
|
# to warn and fall through to plain HTTP on the port that was supposed to be TLS.
|
|
ssl_context = _resolve_ssl_context(reverse_proxy, domain, app_name)
|
|
|
|
# Start HTTP redirect server if SSL is enabled (not needed behind reverse proxy)
|
|
http_redirect_port = server_settings.get('http_redirect_port', 0)
|
|
if http_redirect_port == 0:
|
|
http_redirect_port = 80 if os.geteuid() == 0 else -1
|
|
http_redirect_port = int(os.environ.get('PEGAPROX_HTTP_PORT', http_redirect_port))
|
|
|
|
if ssl_context and http_redirect_port > 0 and not reverse_proxy:
|
|
redirect_thread = threading.Thread(
|
|
target=_start_http_redirect,
|
|
args=(bind_host, http_redirect_port, port, domain),
|
|
daemon=True
|
|
)
|
|
redirect_thread.start()
|
|
print(f"Started additional HTTP -> HTTPS redirect on port {http_redirect_port}")
|
|
|
|
# Determine workers
|
|
# MK 2026-05-31 (v2) — auto-scale with CPU, no hardcoded cap.
|
|
# Two-part fix:
|
|
# (1) `workers` was previously just a log-label — _start_gevent_server
|
|
# prints "(N greenlets)" but never passed `spawn=Pool(N)` to
|
|
# WSGIServer, so the server actually spawned UNLIMITED greenlets.
|
|
# Now plumbed through (see _start_gevent_server below).
|
|
# (2) Formula changed: `min(cpu_count*2, 16)` capped huge customer
|
|
# boxes at 16. `max(8, cpu_count * 4)` gives:
|
|
# 1c VM: 8 workers
|
|
# 4c: 16 workers (same as old default)
|
|
# 8c: 32 workers
|
|
# 32c: 128 workers
|
|
# gevent greenlets are extremely cheap (a few KB stack) so 100s
|
|
# per request handler are fine; the I/O-bound workload benefits
|
|
# from a generous pool when /health + /vms-backup-status + a
|
|
# dashboard refresh all fire at the same time.
|
|
# NS 2026-06-05 — raised floor + multiplier: EACH live SSE/WebSocket stream
|
|
# holds a pool slot for its whole lifetime, so max(8, cpu*4) (=16 on a 4c
|
|
# box) could be consumed by ~16 open dashboard tabs and starve all other API
|
|
# traffic (root of #526's "health spammed, absurdly large time"). Greenlets
|
|
# are cheap so a big pool is fine. Still PEGAPROX_WORKERS-overridable.
|
|
# 1c VM: 32 4c: 64 8c: 128 32c: 512
|
|
cpu_count = multiprocessing.cpu_count()
|
|
workers = int(os.environ.get('PEGAPROX_WORKERS', max(32, cpu_count * 16)))
|
|
|
|
print(f"System: {cpu_count} CPU cores detected")
|
|
print(f"Memory optimization: Garbage collection tuned for {workers} workers")
|
|
gc.set_threshold(700, 10, 10)
|
|
|
|
# Start with Gevent if available
|
|
use_gevent = os.environ.get('PEGAPROX_SERVER', 'auto').lower()
|
|
|
|
if use_gevent == 'gevent' or (use_gevent == 'auto' and GEVENT_AVAILABLE):
|
|
if GEVENT_AVAILABLE:
|
|
_start_gevent_server(app, bind_host, port, ssl_context, domain, workers, http_redirect_port)
|
|
return
|
|
|
|
# Fallback to Flask development server
|
|
print("Starting PegaProx with Flask development server")
|
|
print("WARNING: Not recommended for production!")
|
|
print("Install gevent for better performance: pip install gevent")
|
|
|
|
vnc_ws_port = port + 1
|
|
ssh_ws_port = port + 2
|
|
|
|
# Start VNC/SSH WebSocket servers
|
|
_start_console_servers(bind_host, port, ssl_context)
|
|
|
|
if ssl_context:
|
|
print(f"HTTPS on https://{bind_host}:{port}")
|
|
app.run(host=bind_host, port=port, debug=False, ssl_context=ssl_context, threaded=True)
|
|
else:
|
|
print(f"HTTP on http://{bind_host}:{port}")
|
|
app.run(host=bind_host, port=port, debug=False, threaded=True)
|
|
|
|
|
|
def _start_console_servers(bind_host, port, ssl_context):
|
|
"""Start VNC and SSH WebSocket servers on port+1 and port+2."""
|
|
vnc_ws_port = port + 1
|
|
ssh_ws_port = port + 2
|
|
|
|
try:
|
|
from pegaprox.api.vms import start_vnc_websocket_server, start_ssh_websocket_server
|
|
except ImportError as e:
|
|
print(f"WARNING: Console WebSocket servers not available: {e}")
|
|
return
|
|
|
|
# NS Feb 2026 - asyncio/websockets creates IPv6-only socket for '::' (#95)
|
|
# Use '' so asyncio binds to ALL interfaces (creates both IPv4 + IPv6 listeners)
|
|
console_host = '' if bind_host == '::' else bind_host
|
|
|
|
# MK Feb 2026 - start each server independently so one failure doesn't block the other
|
|
for name, start_fn, ws_port in [
|
|
("VNC", start_vnc_websocket_server, vnc_ws_port),
|
|
("SSH", start_ssh_websocket_server, ssh_ws_port),
|
|
]:
|
|
try:
|
|
if ssl_context:
|
|
start_fn(ws_port, ssl_cert=ssl_context[0], ssl_key=ssl_context[1], host=console_host)
|
|
else:
|
|
start_fn(ws_port, host=console_host)
|
|
except Exception as e:
|
|
print(f"ERROR: {name} WebSocket server (port {ws_port}) failed to start: {e}")
|
|
logging.error(f"{name} WebSocket server startup failed: {e}", exc_info=True)
|
|
|
|
|
|
def _test_ipv6_available():
|
|
"""Test if the system supports IPv6 sockets - Issue #71"""
|
|
try:
|
|
s = socket.socket(socket.AF_INET6, socket.SOCK_STREAM)
|
|
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
s.bind(('::', 0))
|
|
s.close()
|
|
return True
|
|
except (OSError, socket.error):
|
|
return False
|
|
|
|
|
|
def _start_http_redirect(bind_host, http_redirect_port, https_port, domain):
|
|
"""Start a simple HTTP server that redirects to HTTPS using raw sockets"""
|
|
try:
|
|
use_ipv6 = ':' in bind_host
|
|
af = socket.AF_INET6 if use_ipv6 else socket.AF_INET
|
|
sock = socket.socket(af, socket.SOCK_STREAM)
|
|
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
if use_ipv6:
|
|
sock.setsockopt(socket.IPPROTO_IPV6, socket.IPV6_V6ONLY, 0)
|
|
sock.bind((bind_host, http_redirect_port))
|
|
sock.listen(100)
|
|
sock.settimeout(1.0)
|
|
|
|
print(f"HTTP redirect server listening on port {http_redirect_port}")
|
|
|
|
while True:
|
|
try:
|
|
client, addr = sock.accept()
|
|
client.settimeout(5.0)
|
|
try:
|
|
request = client.recv(4096).decode('utf-8', errors='ignore')
|
|
path = '/'
|
|
if request:
|
|
first_line = request.split('\r\n')[0]
|
|
parts = first_line.split(' ')
|
|
if len(parts) >= 2:
|
|
path = parts[1].replace('\r', '').replace('\n', '')
|
|
|
|
# MK: Mar 2026 - serve ACME challenges on port 80 instead of redirecting (#96)
|
|
if path.startswith('/.well-known/acme-challenge/'):
|
|
acme_token = path.split('/')[-1]
|
|
from pegaprox.core.acme import get_challenge_response
|
|
challenge_resp = get_challenge_response(acme_token)
|
|
if challenge_resp:
|
|
http_resp = (
|
|
f"HTTP/1.1 200 OK\r\n"
|
|
f"Content-Type: text/plain\r\n"
|
|
f"Content-Length: {len(challenge_resp)}\r\n"
|
|
f"Connection: close\r\n"
|
|
f"\r\n"
|
|
f"{challenge_resp}"
|
|
)
|
|
client.sendall(http_resp.encode())
|
|
client.close()
|
|
continue
|
|
|
|
host_header = ''
|
|
fwd_proto = ''
|
|
fwd_port = ''
|
|
for line in request.split('\r\n'):
|
|
low = line.lower()
|
|
if low.startswith('host:'):
|
|
host_value = line.split(':', 1)[1].strip()
|
|
if ':' in host_value:
|
|
host_header = host_value.rsplit(':', 1)[0]
|
|
else:
|
|
host_header = host_value
|
|
elif low.startswith('x-forwarded-proto:'):
|
|
fwd_proto = line.split(':', 1)[1].strip().lower()
|
|
elif low.startswith('x-forwarded-port:'):
|
|
fwd_port = line.split(':', 1)[1].strip()
|
|
|
|
# behind reverse proxy with SSL termination? skip redirect (#125)
|
|
from pegaprox.utils.audit import _is_trusted_proxy
|
|
if fwd_proto == 'https' and _is_trusted_proxy(addr[0]):
|
|
continue
|
|
|
|
# NS Jul 2026 (CodeAnt http-response-splitting) — host_header is untrusted;
|
|
# strip CR/LF + reject non-hostname chars before it can reach the Location header.
|
|
import re as _re
|
|
redirect_host = (host_header or 'localhost').split('/')[0].strip()
|
|
if not _re.match(r'^[A-Za-z0-9._\-\[\]:]+$', redirect_host):
|
|
redirect_host = 'localhost'
|
|
if domain:
|
|
if ':' in domain and not domain.startswith('['):
|
|
redirect_host = domain.rsplit(':', 1)[0]
|
|
else:
|
|
redirect_host = domain
|
|
|
|
port = int(fwd_port) if fwd_port else https_port
|
|
if port == 443:
|
|
redirect_url = f'https://{redirect_host}{path}'
|
|
else:
|
|
redirect_url = f'https://{redirect_host}:{port}{path}'
|
|
|
|
response = (
|
|
f"HTTP/1.1 301 Moved Permanently\r\n"
|
|
f"Location: {redirect_url}\r\n"
|
|
f"Content-Length: 0\r\n"
|
|
f"Connection: close\r\n"
|
|
f"\r\n"
|
|
)
|
|
client.sendall(response.encode())
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
try:
|
|
client.close()
|
|
except Exception:
|
|
pass
|
|
except socket.timeout:
|
|
continue
|
|
except Exception as e:
|
|
if 'Bad file descriptor' not in str(e):
|
|
logging.debug(f"HTTP redirect accept error: {e}")
|
|
continue
|
|
except PermissionError:
|
|
print(f"WARNING: Cannot bind to port {http_redirect_port} (requires root). HTTP redirect not available.")
|
|
except OSError as e:
|
|
if 'Address already in use' in str(e):
|
|
print(f"WARNING: Port {http_redirect_port} already in use. HTTP redirect not available.")
|
|
else:
|
|
print(f"WARNING: HTTP redirect server failed: {e}")
|
|
except Exception as e:
|
|
print(f"WARNING: HTTP redirect server failed: {e}")
|
|
|
|
|
|
def _create_listener(bind_host, port_num):
|
|
"""Create a listener socket, IPv6 dual-stack if needed - Issue #71"""
|
|
is_ipv6 = ':' in bind_host
|
|
if is_ipv6:
|
|
try:
|
|
listener = socket.socket(socket.AF_INET6, socket.SOCK_STREAM)
|
|
listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
|
listener.setsockopt(socket.IPPROTO_IPV6, socket.IPV6_V6ONLY, 0)
|
|
listener.bind((bind_host, port_num))
|
|
listener.listen(128)
|
|
listener.setblocking(False)
|
|
return listener
|
|
except OSError as e:
|
|
print(f"WARNING: IPv6 listener on port {port_num} failed ({e}), using IPv4")
|
|
return ('0.0.0.0', port_num)
|
|
else:
|
|
return (bind_host, port_num)
|
|
|
|
|
|
def _start_gevent_server(app, bind_host, port, ssl_context, domain, workers, http_redirect_port=-1):
|
|
"""Start production server with Gevent."""
|
|
from gevent.pywsgi import WSGIServer
|
|
|
|
print(f"Starting PegaProx with Gevent WSGIServer ({workers} greenlets)", flush=True)
|
|
print("Mode: Production (async I/O optimized)", flush=True)
|
|
|
|
# NS: Suppress noisy errors from bots/scanners/disconnects
|
|
import logging as log_module
|
|
log_module.getLogger('gevent').setLevel(log_module.CRITICAL)
|
|
log_module.getLogger('gevent.pywsgi').setLevel(log_module.CRITICAL)
|
|
log_module.getLogger('websockets').setLevel(log_module.CRITICAL)
|
|
log_module.getLogger('websockets.server').setLevel(log_module.CRITICAL)
|
|
log_module.getLogger('websockets.asyncio').setLevel(log_module.CRITICAL)
|
|
|
|
# LW: Monkey-patch traceback to suppress SSL errors
|
|
# gevent uses traceback.print_exception directly, bypassing logging
|
|
import traceback as tb_module
|
|
_original_print_exception = tb_module.print_exception
|
|
_original_print_exc = tb_module.print_exc
|
|
_original_format_exception = tb_module.format_exception
|
|
|
|
def quiet_print_exception(exc, value=None, tb=None, limit=None, file=None, chain=True):
|
|
exc_type = exc if isinstance(exc, type) else type(exc)
|
|
if exc_type and 'ssl' in exc_type.__name__.lower():
|
|
return
|
|
if value and 'ssl' in str(value).lower():
|
|
return
|
|
_original_print_exception(exc, value, tb, limit, file, chain)
|
|
|
|
def quiet_print_exc(limit=None, file=None, chain=True):
|
|
exc_type, exc_value, exc_tb = sys.exc_info()
|
|
if exc_type and 'ssl' in exc_type.__name__.lower():
|
|
return
|
|
_original_print_exc(limit, file, chain)
|
|
|
|
def quiet_format_exception(exc, value=None, tb=None, limit=None, chain=True):
|
|
exc_type = exc if isinstance(exc, type) else type(exc)
|
|
if exc_type and 'ssl' in exc_type.__name__.lower():
|
|
return []
|
|
return _original_format_exception(exc, value, tb, limit, chain)
|
|
|
|
tb_module.print_exception = quiet_print_exception
|
|
tb_module.print_exc = quiet_print_exc
|
|
tb_module.format_exception = quiet_format_exception
|
|
|
|
# NS: Also filter stderr directly as last resort
|
|
import io
|
|
class SSLFilteredStderr:
|
|
def __init__(self, original):
|
|
self._original = original
|
|
self._buffer = []
|
|
self._in_ssl_traceback = False
|
|
|
|
def write(self, text):
|
|
if 'Traceback (most recent call last):' in text:
|
|
self._in_ssl_traceback = False
|
|
self._buffer = [text]
|
|
return len(text)
|
|
if self._buffer:
|
|
self._buffer.append(text)
|
|
full_text = ''.join(self._buffer)
|
|
if 'SSLEOFError' in full_text or 'ssl.SSL' in full_text:
|
|
self._in_ssl_traceback = True
|
|
if text.strip() and not text.startswith(' ') and not text.startswith('Traceback'):
|
|
if self._in_ssl_traceback:
|
|
self._buffer = []
|
|
self._in_ssl_traceback = False
|
|
return len(text)
|
|
else:
|
|
for line in self._buffer:
|
|
self._original.write(line)
|
|
self._buffer = []
|
|
return len(text)
|
|
return self._original.write(text)
|
|
|
|
def flush(self):
|
|
if self._buffer and not self._in_ssl_traceback:
|
|
for line in self._buffer:
|
|
self._original.write(line)
|
|
self._buffer = []
|
|
self._original.flush()
|
|
|
|
def __getattr__(self, name):
|
|
return getattr(self._original, name)
|
|
|
|
sys.stderr = SSLFilteredStderr(sys.stderr)
|
|
|
|
os.environ['GEVENT_DEBUG'] = 'off'
|
|
|
|
# WebSocket handler
|
|
use_websocket_handler = False
|
|
try:
|
|
from geventwebsocket.handler import WebSocketHandler
|
|
use_websocket_handler = True
|
|
print("WebSocket support: geventwebsocket enabled")
|
|
except ImportError:
|
|
use_websocket_handler = False
|
|
print("WebSocket support: geventwebsocket NOT installed")
|
|
print(" Install with: pip install gevent-websocket")
|
|
|
|
# NS: Custom handler to suppress SSL error tracebacks completely
|
|
# These happen when users close browser tabs - totally normal
|
|
if use_websocket_handler:
|
|
class QuietWebSocketHandler(WebSocketHandler):
|
|
def handle_one_response(self):
|
|
try:
|
|
return super().handle_one_response()
|
|
except Exception as e:
|
|
if 'ssl' in type(e).__name__.lower() or 'ssl' in str(e).lower():
|
|
return
|
|
raise
|
|
|
|
def log_error(self, msg, *args):
|
|
if 'ssl' in str(msg).lower() or 'eof' in str(msg).lower():
|
|
return
|
|
super().log_error(msg, *args)
|
|
else:
|
|
QuietWebSocketHandler = None
|
|
|
|
# Custom error handler to suppress SSL errors (from bots/scanners/disconnects)
|
|
class QuietWSGIServer(WSGIServer):
|
|
def wrap_socket_and_handle(self, client_socket, address):
|
|
"""Override to catch SSL errors during handshake"""
|
|
try:
|
|
return super().wrap_socket_and_handle(client_socket, address)
|
|
except Exception as e:
|
|
if 'ssl' in str(type(e).__name__).lower() or 'ssl' in str(e).lower():
|
|
pass
|
|
else:
|
|
raise
|
|
|
|
def handle_error(self, *args):
|
|
"""Suppress SSL errors - they're normal with self-signed certs"""
|
|
exc_info = sys.exc_info()
|
|
exc_type = exc_info[0]
|
|
if exc_type is not None:
|
|
if 'ssl' in exc_type.__name__.lower():
|
|
return
|
|
pass
|
|
|
|
def log_error(self, msg, *args):
|
|
"""Suppress SSL error logging"""
|
|
msg_lower = str(msg).lower()
|
|
if 'ssl' in msg_lower or 'eof' in msg_lower or 'broken pipe' in msg_lower:
|
|
return
|
|
print(f"[Server Error] {msg % args if args else msg}")
|
|
|
|
# DualProtocolWSGIServer - HTTP and HTTPS on same port
|
|
# If someone visits http://server:5000, they get redirected to https://server:5000
|
|
# MK: Claude helped with the TLS detection logic - checking for 0x16/0x80 bytes
|
|
class DualProtocolWSGIServer(QuietWSGIServer):
|
|
"""WSGI Server that detects HTTP vs HTTPS and redirects HTTP to HTTPS"""
|
|
|
|
def __init__(self, *args, redirect_domain=None, **kwargs):
|
|
self._redirect_domain = redirect_domain
|
|
super().__init__(*args, **kwargs)
|
|
|
|
def wrap_socket_and_handle(self, client_socket, address):
|
|
"""Peek at first bytes to detect protocol"""
|
|
if not self.ssl_args:
|
|
return super().wrap_socket_and_handle(client_socket, address)
|
|
try:
|
|
first_byte = client_socket.recv(1, socket.MSG_PEEK)
|
|
if not first_byte:
|
|
client_socket.close()
|
|
return
|
|
if first_byte[0] == 0x16 or first_byte[0] == 0x80:
|
|
return super().wrap_socket_and_handle(client_socket, address)
|
|
else:
|
|
# NS: #125 - reverse proxy with SSL termination? serve as plain HTTP
|
|
# only trust forwarded headers from loopback / configured trusted proxies
|
|
from pegaprox.utils.audit import _is_trusted_proxy
|
|
if _is_trusted_proxy(address[0]):
|
|
try:
|
|
peek = client_socket.recv(8192, socket.MSG_PEEK)
|
|
if b'x-forwarded-proto' in peek.lower():
|
|
for hdr in peek.decode('utf-8', errors='ignore').split('\r\n'):
|
|
if hdr.lower().startswith('x-forwarded-proto:'):
|
|
if hdr.split(':', 1)[1].strip().lower() == 'https':
|
|
return self.handle(client_socket, address)
|
|
break
|
|
except Exception:
|
|
pass
|
|
self._handle_http_redirect(client_socket, address)
|
|
return
|
|
except Exception as e:
|
|
if 'ssl' in str(type(e).__name__).lower():
|
|
return
|
|
try:
|
|
return super().wrap_socket_and_handle(client_socket, address)
|
|
except Exception:
|
|
pass
|
|
|
|
def _handle_http_redirect(self, client_socket, address):
|
|
"""Send HTTP 301 redirect to HTTPS version"""
|
|
try:
|
|
client_socket.settimeout(5.0)
|
|
request_data = b''
|
|
while b'\r\n\r\n' not in request_data and len(request_data) < 8192:
|
|
chunk = client_socket.recv(1024)
|
|
if not chunk:
|
|
break
|
|
request_data += chunk
|
|
|
|
request = request_data.decode('utf-8', errors='ignore')
|
|
path = '/'
|
|
if request:
|
|
first_line = request.split('\r\n')[0]
|
|
parts = first_line.split(' ')
|
|
if len(parts) >= 2:
|
|
path = parts[1].replace('\r', '').replace('\n', '')
|
|
|
|
host = self._redirect_domain or 'localhost'
|
|
for line in request.split('\r\n'):
|
|
if line.lower().startswith('host:'):
|
|
host_value = line.split(':', 1)[1].strip()
|
|
if host_value.startswith('['):
|
|
if ']:' in host_value:
|
|
host = host_value.rsplit(':', 1)[0]
|
|
else:
|
|
host = host_value
|
|
elif ':' in host_value:
|
|
host = host_value.rsplit(':', 1)[0]
|
|
else:
|
|
host = host_value
|
|
break
|
|
|
|
# NS Jul 2026 (CodeAnt http-response-splitting) — the Host header is untrusted;
|
|
# reject non-hostname chars before it can reach the Location header (open-redirect
|
|
# / header injection). A configured _redirect_domain (below) always wins.
|
|
import re as _re
|
|
if not _re.match(r'^[A-Za-z0-9._\-\[\]:]+$', host or ''):
|
|
host = 'localhost'
|
|
|
|
if self._redirect_domain:
|
|
d = self._redirect_domain
|
|
if ':' in d and not d.startswith('['):
|
|
host = d.rsplit(':', 1)[0]
|
|
else:
|
|
host = d
|
|
|
|
# NS: #125 - respect proxy headers so we don't redirect to internal port
|
|
fwd_proto = ''
|
|
fwd_port = ''
|
|
for line in request.split('\r\n'):
|
|
lower = line.lower()
|
|
if lower.startswith('x-forwarded-proto:'):
|
|
fwd_proto = line.split(':', 1)[1].strip().lower()
|
|
elif lower.startswith('x-forwarded-port:'):
|
|
fwd_port = line.split(':', 1)[1].strip()
|
|
|
|
if fwd_proto == 'https':
|
|
# already behind SSL-terminating proxy, don't redirect
|
|
return
|
|
|
|
port = int(fwd_port) if fwd_port else self.server_port
|
|
if port == 443:
|
|
redirect_url = f'https://{host}{path}'
|
|
else:
|
|
redirect_url = f'https://{host}:{port}{path}'
|
|
|
|
response = (
|
|
f"HTTP/1.1 301 Moved Permanently\r\n"
|
|
f"Location: {redirect_url}\r\n"
|
|
f"Content-Type: text/html\r\n"
|
|
f"Content-Length: 0\r\n"
|
|
f"Connection: close\r\n"
|
|
f"\r\n"
|
|
)
|
|
client_socket.sendall(response.encode())
|
|
except Exception:
|
|
pass
|
|
finally:
|
|
try:
|
|
client_socket.close()
|
|
except Exception:
|
|
pass
|
|
|
|
# Server args - add WebSocket handler if available
|
|
# MK 2026-05-31 — actually wire `workers` into the request-handler pool.
|
|
# gevent.pywsgi.WSGIServer defaults to `spawn=None` which spawns an
|
|
# unlimited greenlet per request. PEGAPROX_WORKERS was a startup-log
|
|
# label only — never enforced. Now caps the request-handling pool at
|
|
# `workers`; per-request fanouts (storage scan, PBS scan, SSH calls)
|
|
# still spawn inside their own request handler.
|
|
from gevent.pool import Pool as _RequestPool
|
|
server_kwargs = {'log': None, 'spawn': _RequestPool(workers)}
|
|
if use_websocket_handler and QuietWebSocketHandler:
|
|
server_kwargs['handler_class'] = QuietWebSocketHandler
|
|
|
|
is_ipv6_bind = ':' in bind_host
|
|
|
|
if ssl_context:
|
|
print(f"HTTPS on https://{bind_host}:{port}", flush=True)
|
|
ssl_ctx = ssl.SSLContext(ssl.PROTOCOL_TLS_SERVER)
|
|
ssl_ctx.minimum_version = ssl.TLSVersion.TLSv1_2
|
|
ssl_ctx.load_cert_chain(ssl_context[0], ssl_context[1])
|
|
# NS: http_redirect_port == -1 disables ALL http→https redirect (#125)
|
|
# including the dual-protocol detection on the main port
|
|
if http_redirect_port < 0:
|
|
http_server = QuietWSGIServer(
|
|
_create_listener(bind_host, port), app,
|
|
ssl_context=ssl_ctx,
|
|
**server_kwargs
|
|
)
|
|
else:
|
|
http_server = DualProtocolWSGIServer(
|
|
_create_listener(bind_host, port), app,
|
|
ssl_context=ssl_ctx,
|
|
redirect_domain=domain,
|
|
**server_kwargs
|
|
)
|
|
else:
|
|
print(f"HTTP on http://{bind_host}:{port}", flush=True)
|
|
print("WARNING: Running without HTTPS - noVNC console may not work!", flush=True)
|
|
http_server = QuietWSGIServer(_create_listener(bind_host, port), app, **server_kwargs)
|
|
|
|
# Start VNC/SSH WebSocket servers
|
|
_start_console_servers(bind_host, port, ssl_context)
|
|
|
|
# Handle graceful shutdown
|
|
def signal_handler(signum, frame):
|
|
print("\nShutting down gracefully...")
|
|
http_server.stop()
|
|
sys.exit(0)
|
|
|
|
signal.signal(signal.SIGINT, signal_handler)
|
|
signal.signal(signal.SIGTERM, signal_handler)
|
|
|
|
print("SSL/WebSocket errors (bots, scanners, disconnects) are suppressed")
|
|
http_server.serve_forever()
|