MrMasterbay 35d40781bd security: enforce authz gaps from the Aikido Testing-branch pentest (batch 1/2)
Adversarially verified all 45 AI-pentest findings against the current code; this batch
fixes the 10 HIGH + 2 MED confirmed REAL (17 were by-design/known admin-wide debt, 2
already fixed, 1 false positive). Full in-process test suite is green — the one case that
failed was a test that encoded the very BMC credential-exfil this now guards; corrected it
and added the negative case.

- rbac.py: user_can_access_vmware_vm honours effective_role (token-scoped), like its Proxmox
  twin, so an admin-owned viewer-scoped API token no longer gets the full-admin VMware bypass.
- nodes.py / vmware.py: reject a masked password paired with a CHANGED host - the stored BMC /
  VMware secret can no longer be shipped to a caller-chosen host (credential exfiltration).
- schedules.py: create/update enforce the same per-VM ACL as the live action (build_authz_user
  + user_can_access_vm), not just cluster reachability.
- pbs.py: restore (overwrite/test) enforces per-VM authorization on the destination VMID;
  backup-diff now requires pbs.datastore.view instead of plain pbs.view.
- users.py: create/update cap delegated permissions and custom roles to what the caller holds
  (a tenant admin.users delegate can no longer mint accounts with admin.* perms via the
  permissions[] list or the custom-role tier fallback); update_user rejects a foreign-tenant
  custom role; update_tenant is tenant-scoped like get_tenant_quota.
- groups.py: a global (tenant_id NULL) cluster group is admin-only for writes.
2026-08-07 15:31:13 +02:00

2833 lines
116 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- coding: utf-8 -*-
"""node management, updates, SMBIOS & scripts routes - split from monolith dec 2025, NS/MK"""
import os
import json
import time
import logging
import uuid
import re
import shlex
from datetime import datetime, timedelta
from flask import Blueprint, jsonify, request
from pegaprox.constants import *
from pegaprox.globals import *
from pegaprox.models.permissions import *
from pegaprox.utils.sanitization import sanitize_log_message as _sl # CWE-117
from pegaprox.core.db import get_db
from pegaprox.utils.auth import require_auth, load_users, verify_password
from pegaprox.utils.audit import log_audit
from pegaprox.api.helpers import check_cluster_access, safe_error
bp = Blueprint('nodes', __name__)
# ==================== NODE MANAGEMENT API ENDPOINTS ====================
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/summary', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_summary_api(cluster_id, node):
"""Get node summary"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_summary(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/ip', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_ip_api(cluster_id, node):
"""Get node IP address for SSH connections
MK: Tries multiple methods to get the node's IP:
1. Cluster status API (for clustered nodes)
2. Network configuration API
3. Fallback to cluster host
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
cluster_host, cluster_port = mgr.host, mgr.api_port
node_ip = None
source = None
# NS Mar 2026: XCP-ng uses XAPI host.get_address instead of Proxmox REST
if getattr(mgr, 'cluster_type', 'proxmox') == 'xcpng':
try:
node_ip = mgr._get_host_ip(node)
source = 'xapi_host_address'
except Exception as e:
logging.error(f"XCP-ng get_node_ip: {e}")
node_ip = cluster_host
source = 'xcpng_fallback'
else:
# NS Apr 2026 (PR #324): let manager._get_node_ip do the heavy lifting
# (scores interfaces, filters corosync IPs out of the mgmt net, probes
# the SSH port). This endpoint is informational -- fall back to
# cluster_host only when we truly couldn't resolve anything.
try:
node_ip = mgr._get_node_ip(node)
if node_ip:
source = 'manager_get_node_ip'
except Exception as e:
logging.error(f"Error getting node IP: {e}")
if not node_ip:
node_ip = cluster_host
source = 'cluster_host_fallback'
return jsonify({
'ip': node_ip,
'node': node,
'source': source
})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/rrddata', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_rrddata_api(cluster_id, node):
"""Get node performance metrics (RRD data) for charts
Query params:
- timeframe: hour, day, week, month, year (default: hour)
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
timeframe = request.args.get('timeframe', 'hour')
return jsonify(manager.get_node_rrddata(node, timeframe))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_network_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_network_config(node))
# MK May 2026 — node-name regex (mirrors vms.py:2718 hardening). RFC-1035-ish
# DNS rules: letter-led, then letters/digits/hyphen/dot, max ~63. Everything
# else gets a 400 before the value flows into _get_node_ip() / PVE API URLs.
_NODE_NAME_RE = re.compile(r'^[a-zA-Z][a-zA-Z0-9.\-]{0,62}$')
def _reject_bad_node(node):
"""Return (None, None) if node looks valid, (response, 400) otherwise.
Used by the SSH-fronted node endpoints below as a defense-in-depth gate
even though the SSH commands themselves don't interpolate `node`.
Non-strings (e.g. a crafted JSON list `{"nodes":[123]}`) are rejected as
invalid rather than blowing up the regex with a TypeError -> clean 400."""
if not isinstance(node, str) or not _NODE_NAME_RE.match(node):
return jsonify({'error': 'Invalid node name'}), 400
return None, None
# MK May 2026 — per-NIC error/drop counters (SSH /proc/net/dev).
# PVE's own /network endpoint only carries bridge+bond config; for the
# actual rx_err/rx_drop/tx_err counters operators need we have to read
# /proc/net/dev directly. Same auth path the syslog viewer uses.
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/netstats', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_netstats_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
result = cluster_managers[cluster_id].get_node_netstats(node)
if isinstance(result, dict) and result.get('error'):
return jsonify(result), 502
return jsonify(result)
# MK May 2026 — cluster health: corosync ring + pvecm quorum + service state.
# /cluster/status gives "are we quorate" but not ring latency, per-service
# uptime, or per-ring health — chasing a flapping cluster needs all three.
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/cluster-health', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_cluster_health_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
result = cluster_managers[cluster_id].get_node_cluster_health(node)
if isinstance(result, dict) and result.get('error'):
return jsonify(result), 502
return jsonify(result)
# MK May 2026 — lm-sensors readings (CPU temp, fan rpm, voltages).
# Graceful on hosts without lm-sensors installed (returns error string the
# UI can show as "not available").
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/sensors', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_sensors_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
result = cluster_managers[cluster_id].get_node_sensors(node)
if isinstance(result, dict) and result.get('error'):
return jsonify(result), 502
return jsonify(result)
# #601 — per-node hottest-temperature history for the chart-over-time ask. Reads the
# background collector's persisted 5-min metrics_history snapshots (off-hub, cached);
# no live SSH here, so it stays cheap even at fleet scale.
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/temperature-history', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_temperature_history_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
from pegaprox.background.metrics import load_metrics_history
series = []
try:
hist = load_metrics_history() or {}
for snap in (hist.get('snapshots') or []):
try:
nodes = ((snap.get('clusters') or {}).get(cluster_id) or {}).get('nodes') or {}
temp = (nodes.get(node) or {}).get('temp')
ts = snap.get('timestamp')
if temp is not None and ts:
series.append({'ts': ts, 'temp': temp})
except Exception:
continue
except Exception as e:
return jsonify({'error': f'history read failed: {e}'}), 500
series.sort(key=lambda x: x['ts'])
return jsonify({'series': series, 'unit': '°C', 'count': len(series)})
# ── In-band hardware monitoring / BMC (#609) ──────────────────────────────────
# Credential-free reads via local ipmitool on the node (see pegaprox/core/bmc.py).
# Gated behind a one-time, versioned compliance acknowledgement that is written to
# the audit log — so "I never enabled that" cannot hold ("dann heißt es nicht 'Ich
# hab es nie bekommen'"). The read route refuses to serve until consent is on.
def _hw_int(v):
"""Coerce a stored/posted ack_version to int, degrading to 0 on any bad value
(e.g. a corrupt/crafted settings blob with ack_version='x' or [1]) so the gate
fails CLOSED — returns disabled — rather than raising a 500 on every read."""
try:
return int(v)
except (TypeError, ValueError):
return 0
def _hw_consent_state():
"""(enabled, ack_record) for in-band hardware monitoring from server settings.
enabled == acknowledged AND the stored ack is at/above the current warning ver."""
from pegaprox.api.helpers import load_server_settings
from pegaprox.core import bmc
ack = (load_server_settings() or {}).get('hardware_monitoring') or {}
enabled = bool(ack.get('enabled')) and _hw_int(ack.get('ack_version')) >= bmc.HW_CONSENT_VERSION
return enabled, ack
def _redfish_consent_state():
"""(enabled, ack_record) for OUT-OF-BAND Redfish monitoring — a separate, sharper
opt-in from the in-band one (its own versioned warning + 5s delay)."""
from pegaprox.api.helpers import load_server_settings
from pegaprox.core import bmc
ack = (load_server_settings() or {}).get('hardware_monitoring_redfish') or {}
enabled = bool(ack.get('enabled')) and _hw_int(ack.get('ack_version')) >= bmc.REDFISH_CONSENT_VERSION
return enabled, ack
@bp.route('/api/hardware-monitoring/consent', methods=['GET'])
@require_auth(perms=['node.view'])
def get_hw_monitoring_consent():
"""Current consent status + the versioned warning the UI must present, in the
requested language (?lang=de|en|fr|es|pt|ko|it; falls back to English)."""
from pegaprox.core import bmc
enabled, ack = _hw_consent_state()
return jsonify({
'enabled': enabled,
'acknowledged_by': ack.get('acknowledged_by'),
'acknowledged_at': ack.get('acknowledged_at'),
'ack_version': ack.get('ack_version'),
'ack_lang': ack.get('ack_lang'),
'current_version': bmc.HW_CONSENT_VERSION,
'warning': bmc.hw_consent_warning(request.args.get('lang')),
})
@bp.route('/api/hardware-monitoring/consent', methods=['POST'])
@require_auth(perms=['admin.settings'])
def set_hw_monitoring_consent():
"""Enable (acknowledge current warning) or disable in-band hardware monitoring.
Enabling requires acknowledge=true AND ack_version == the current warning
version, so a stale UI cannot silently opt in under an old warning. The
acknowledgement is persisted to the audit log with who/when/which-version.
"""
from pegaprox.api.helpers import load_server_settings, save_server_settings
from pegaprox.core import bmc
data = request.get_json(silent=True) or {}
usr = request.session.get('user', 'admin')
settings = load_server_settings() or {}
if data.get('enabled') is False:
settings['hardware_monitoring'] = {
'enabled': False,
'acknowledged_by': usr,
'acknowledged_at': datetime.now().isoformat(),
'ack_version': 0,
}
save_server_settings(settings)
log_audit(usr, 'hardware_monitoring.disabled',
'In-band hardware monitoring disabled')
return jsonify({'enabled': False})
if data.get('acknowledge') is True and _hw_int(data.get('ack_version')) == bmc.HW_CONSENT_VERSION:
# which localized text the user actually saw (part of the non-repudiation record)
ack_lang = (str(data.get('lang') or 'en').split('-')[0].lower())[:8]
if ack_lang not in bmc._HW_CONSENT_TEXT:
ack_lang = 'en'
rec = {
'enabled': True,
'acknowledged_by': usr,
'acknowledged_at': datetime.now().isoformat(),
'ack_version': bmc.HW_CONSENT_VERSION,
'ack_lang': ack_lang,
}
settings['hardware_monitoring'] = rec
save_server_settings(settings)
# Durable, attributed non-repudiation record of the acknowledgement.
log_audit(usr, 'hardware_monitoring.enabled',
'In-band hardware monitoring enabled; acknowledged compliance '
'warning v%d [%s]' % (bmc.HW_CONSENT_VERSION, ack_lang))
return jsonify(rec)
return jsonify({
'error': 'Must acknowledge the current compliance warning to enable',
'code': 'ACK_REQUIRED',
'current_version': bmc.HW_CONSENT_VERSION,
}), 400
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/hardware', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_hardware_api(cluster_id, node):
"""Hardware health for one node (sensors/power/FRU/SEL rollup). Prefers in-band
ipmitool; falls back to a configured out-of-band Redfish endpoint (#609 phase 3).
Refuses with CONSENT_REQUIRED until at least one source is enabled + acknowledged.
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
from pegaprox.core import bmc, redfish
inband_ok, _ = _hw_consent_state()
redfish_ok, _ = _redfish_consent_state()
if not (inband_ok or redfish_ok):
return jsonify({
'error': 'Hardware monitoring is not enabled',
'code': 'CONSENT_REQUIRED',
'current_version': bmc.HW_CONSENT_VERSION,
}), 403
result = redfish.read_node_hardware(cluster_managers[cluster_id], cluster_id, node,
inband_ok=inband_ok, redfish_ok=redfish_ok)
return jsonify(result)
# ── Out-of-band Redfish consent (#609 phase 3) ────────────────────────────────
# A SEPARATE, sharper opt-in from the in-band one (its own versioned warning with
# an enforced 5s delay). Same non-repudiation contract: the ack (with the language
# the admin saw) is written to the audit log.
@bp.route('/api/hardware-monitoring/redfish-consent', methods=['GET'])
@require_auth(perms=['node.view'])
def get_redfish_consent():
"""Redfish consent status + the versioned out-of-band warning (?lang=...)."""
from pegaprox.core import bmc
enabled, ack = _redfish_consent_state()
return jsonify({
'enabled': enabled,
'acknowledged_by': ack.get('acknowledged_by'),
'acknowledged_at': ack.get('acknowledged_at'),
'ack_version': ack.get('ack_version'),
'ack_lang': ack.get('ack_lang'),
'current_version': bmc.REDFISH_CONSENT_VERSION,
'warning': bmc.redfish_consent_warning(request.args.get('lang')),
})
@bp.route('/api/hardware-monitoring/redfish-consent', methods=['POST'])
@require_auth(perms=['admin.settings'])
def set_redfish_consent():
"""Enable/disable out-of-band Redfish monitoring. Enabling requires the current
warning version to be acknowledged; the acknowledgement + language are audited."""
from pegaprox.api.helpers import load_server_settings, save_server_settings
from pegaprox.core import bmc
data = request.get_json(silent=True) or {}
usr = request.session.get('user', 'admin')
settings = load_server_settings() or {}
if data.get('enabled') is False:
settings['hardware_monitoring_redfish'] = {
'enabled': False, 'acknowledged_by': usr,
'acknowledged_at': datetime.now().isoformat(), 'ack_version': 0,
}
save_server_settings(settings)
log_audit(usr, 'hardware_monitoring_redfish.disabled',
'Out-of-band Redfish hardware monitoring disabled')
return jsonify({'enabled': False})
if data.get('acknowledge') is True and _hw_int(data.get('ack_version')) == bmc.REDFISH_CONSENT_VERSION:
ack_lang = (str(data.get('lang') or 'en').split('-')[0].lower())[:8]
if ack_lang not in bmc._REDFISH_CONSENT_TEXT:
ack_lang = 'en'
rec = {
'enabled': True, 'acknowledged_by': usr,
'acknowledged_at': datetime.now().isoformat(),
'ack_version': bmc.REDFISH_CONSENT_VERSION, 'ack_lang': ack_lang,
}
settings['hardware_monitoring_redfish'] = rec
save_server_settings(settings)
log_audit(usr, 'hardware_monitoring_redfish.enabled',
'Out-of-band Redfish hardware monitoring enabled; acknowledged '
'compliance warning v%d [%s]' % (bmc.REDFISH_CONSENT_VERSION, ack_lang))
return jsonify(rec)
return jsonify({'error': 'Must acknowledge the current compliance warning to enable',
'code': 'ACK_REQUIRED',
'current_version': bmc.REDFISH_CONSENT_VERSION}), 400
# ── Per-node BMC (Redfish) endpoint config (#609 phase 3) ─────────────────────
# admin.settings only; the password is stored encrypted and NEVER returned to the
# browser (masked as ********). Configuring an endpoint requires Redfish consent.
def _require_redfish_consent():
from pegaprox.core import bmc
enabled, _ = _redfish_consent_state()
if not enabled:
return jsonify({'error': 'Out-of-band (Redfish) monitoring is not enabled',
'code': 'CONSENT_REQUIRED',
'current_version': bmc.REDFISH_CONSENT_VERSION}), 403
return None
@bp.route('/api/clusters/<cluster_id>/bmc-endpoints', methods=['GET'])
@require_auth(perms=['admin.settings'])
def list_bmc_endpoints_api(cluster_id):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
return jsonify({'endpoints': get_db().list_bmc_endpoints(cluster_id)})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/bmc-endpoint', methods=['GET'])
@require_auth(perms=['admin.settings'])
def get_bmc_endpoint_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
ep = get_db().get_bmc_endpoint(cluster_id, node, decrypt=False)
if not ep:
return jsonify({'configured': False})
ep['configured'] = True
ep['password'] = '********' if ep.pop('has_password', False) else ''
return jsonify(ep)
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/bmc-endpoint', methods=['POST'])
@require_auth(perms=['admin.settings'])
def set_bmc_endpoint_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
consent_err = _require_redfish_consent()
if consent_err is not None:
return consent_err
from pegaprox.core import redfish
data = request.get_json(silent=True) or {}
host = (data.get('host') or '').strip()
base, why = redfish._validate_host(host)
if base is None:
return jsonify({'error': why or 'invalid BMC host'}), 400
user = (data.get('user') or '').strip()[:128]
verify_ssl = bool(data.get('verify_ssl'))
enabled = data.get('enabled', True)
pw = data.get('password')
# keep the existing password when the UI submits the mask (never overwrite with ****).
# NS Aug 2026 (Aikido pentest) — only when the host is unchanged: the 5-min collector
# auto-connects (Basic auth) to whatever host is stored, so a masked password paired with a
# changed host would exfil the saved secret to a caller-chosen box. Same guard the /test route has.
if pw in (None, '', '********'):
existing = get_db().get_bmc_endpoint(cluster_id, node)
if host == ((existing or {}).get('host') or '').strip():
pw = (existing or {}).get('password', '')
else:
return jsonify({'error': 'A full password is required when changing the BMC host.'}), 400
if not pw:
return jsonify({'error': 'BMC password required'}), 400
if len(str(pw)) > 512:
return jsonify({'error': 'BMC password too long'}), 400
get_db().save_bmc_endpoint(cluster_id, node, host, user, str(pw),
verify_ssl=verify_ssl, enabled=bool(enabled))
usr = request.session.get('user', 'admin')
log_audit(usr, 'hardware_monitoring_redfish.endpoint_set',
f'BMC/Redfish endpoint set for node {node} (host={host}, verify_ssl={verify_ssl})',
cluster=getattr(getattr(cluster_managers.get(cluster_id), 'config', None), 'name', cluster_id))
return jsonify({'success': True})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/bmc-endpoint', methods=['DELETE'])
@require_auth(perms=['admin.settings'])
def delete_bmc_endpoint_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
removed = get_db().delete_bmc_endpoint(cluster_id, node)
if removed:
usr = request.session.get('user', 'admin')
log_audit(usr, 'hardware_monitoring_redfish.endpoint_removed',
f'BMC/Redfish endpoint removed for node {node}')
return jsonify({'success': True, 'removed': removed})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/bmc-endpoint/test', methods=['POST'])
@require_auth(perms=['admin.settings'])
def test_bmc_endpoint_api(cluster_id, node):
"""Live read the stored (or posted) BMC endpoint over Redfish to validate it.
Returns {available, health?, reason?} — read-only, no config change."""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
bad, code = _reject_bad_node(node)
if bad is not None: return bad, code
consent_err = _require_redfish_consent()
if consent_err is not None:
return consent_err
from pegaprox.core import redfish
data = request.get_json(silent=True) or {}
stored = get_db().get_bmc_endpoint(cluster_id, node) or {}
host = (data.get('host') or stored.get('host') or '').strip()
if not host:
return jsonify({'available': False, 'reason': 'no BMC host configured'}), 400
user = (data.get('user') or stored.get('user') or '').strip()
pw = data.get('password')
# SECURITY (pentest 2026-07-25): never pair the stored BMC secret with a caller-CHOSEN host.
# A masked/empty password falls back to the stored plaintext, but if the caller also overrode
# `host`, that credential would be sent (preemptive Basic-auth) to an arbitrary off-box host of
# their choosing = credential exfiltration. Only reuse the stored password when the host being
# tested IS the stored host; testing any other host must carry its own full password.
if pw in (None, '', '********'):
if host == (stored.get('host') or '').strip():
pw = stored.get('password', '')
else:
return jsonify({'available': False,
'reason': 'a full password is required to test a host other than the stored one'}), 400
verify_ssl = bool(data.get('verify_ssl', stored.get('verify_ssl', False)))
res = redfish.read_node_bmc_redfish(host, user, str(pw or ''), verify_ssl=verify_ssl, timeout=8)
return jsonify({'available': res.get('available', False),
'health': res.get('health'), 'reason': res.get('reason'),
'source': 'redfish'})
@bp.route('/api/clusters/<cluster_id>/hardware/health', methods=['GET'])
@require_auth(perms=['node.view'])
def get_cluster_hardware_health_api(cluster_id):
"""Cluster-wide degraded-hardware rollup from the cached per-node BMC health
(#609 phase 2). Cache-only — no SSH on the request; the 5-min metrics collector
populates the cache. Consent-gated like the per-node read."""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
from pegaprox.core import bmc
inband_ok, _ = _hw_consent_state()
redfish_ok, _ = _redfish_consent_state()
if not (inband_ok or redfish_ok): # #609 phase 3 — a Redfish-only deployment must not 403
return jsonify({'error': 'Hardware monitoring is not enabled',
'code': 'CONSENT_REQUIRED',
'current_version': bmc.HW_CONSENT_VERSION}), 403
mgr = cluster_managers[cluster_id]
# in-band BMC is proxmox-only; non-proxmox managers don't have the rollup helper.
rollup = getattr(mgr, 'get_cluster_hw_rollup', None)
if getattr(mgr, 'cluster_type', 'proxmox') != 'proxmox' or not callable(rollup):
return jsonify({'health': 'unknown', 'available': False, 'checked': 0,
'counts': {'ok': 0, 'warning': 0, 'critical': 0}, 'degraded': []})
return jsonify(rollup())
# ipmitool one-click install (#609 step 3). Fully static script — no user-controlled
# input reaches the shell (unlike the StarWind installer there is no repo/key URL), so
# no SSRF/url-allowlist is needed. Read-only in-band IPMI (local KCS, no BMC creds).
# Gated behind the SAME compliance acknowledgement as the read path: the install
# MUTATES the node (package + IPMI kernel modules) so it must not run before the
# warning is acknowledged. admin.settings + audited + idempotent + bounded-parallel.
IPMITOOL_INSTALL_SCRIPT = """#!/usr/bin/env bash
# PegaProx (#609) — install ipmitool for in-band hardware monitoring on this node.
set -uo pipefail
if command -v ipmitool >/dev/null 2>&1; then
ver="$(ipmitool -V 2>/dev/null | head -1 || echo ipmitool)"
echo "PP_OK already_installed $ver"; exit 0
fi
if ! command -v apt-get >/dev/null 2>&1; then echo 'PP_ERR no-apt-get'; exit 3; fi
apt-get update -o Acquire::Retries=2 >/tmp/pp_ipmitool_apt.log 2>&1 || true
if ! DEBIAN_FRONTEND=noninteractive apt-get install -y ipmitool >>/tmp/pp_ipmitool_apt.log 2>&1; then
echo 'PP_ERR apt-install-failed'; exit 5
fi
command -v ipmitool >/dev/null 2>&1 || { echo 'PP_ERR not-installed-after-apt'; exit 6; }
ver="$(ipmitool -V 2>/dev/null | head -1 || echo ipmitool)"
echo "PP_OK installed $ver"
"""
@bp.route('/api/clusters/<cluster_id>/hardware/ipmitool/install', methods=['POST'])
@require_auth(perms=['admin.settings']) # root apt-install on nodes -> admin-only, like starlvm
def install_ipmitool_api(cluster_id):
"""Install ipmitool on cluster node(s) over SSH for in-band hardware monitoring.
Body (optional): {nodes: [names]} — defaults to all cluster nodes. Idempotent
per node (skips if ipmitool present). Requires the hardware-monitoring feature
to be enabled (compliance acknowledged) first — the install changes the node.
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if getattr(mgr, 'cluster_type', 'proxmox') != 'proxmox':
return jsonify({'error': 'ipmitool install is Proxmox-only'}), 400
# consent gate — do not mutate a node before the warning is acknowledged
from pegaprox.core import bmc
enabled, _ack = _hw_consent_state()
if not enabled:
return jsonify({'error': 'Hardware monitoring is not enabled',
'code': 'CONSENT_REQUIRED',
'current_version': bmc.HW_CONSENT_VERSION}), 403
raw_nodes = (request.get_json(silent=True) or {}).get('nodes') or []
if not isinstance(raw_nodes, list):
return jsonify({'error': 'nodes must be a list of node names'}), 400
# validate each entry BEFORE set() — a dict/list entry would raise in set()
for n in raw_nodes:
bad, code = _reject_bad_node(n)
if bad is not None:
return bad, code
only = set(raw_nodes)
nodes = _cluster_node_names(mgr)
if only:
nodes = [n for n in nodes if n in only]
if not nodes:
return jsonify({'error': 'No nodes found in cluster'}), 404
def _install_one(node):
node_ip = mgr._get_node_ip(node)
if not node_ip:
return {'node': node, 'success': False, 'error': 'no SSH-reachable IP'}
ssh = None
try:
for attempt in range(3):
ssh = mgr._ssh_connect(node_ip)
if ssh:
break
if attempt < 2:
time.sleep(1.5)
if not ssh:
return {'node': node, 'success': False, 'error': 'SSH connect failed after 3 tries'}
_ssh_write_file(ssh, '/tmp/pegaprox-ipmitool-install.sh', IPMITOOL_INSTALL_SCRIPT, 0o755)
out, _e = _ssh_run_checked(ssh, 'bash /tmp/pegaprox-ipmitool-install.sh', timeout=180)
try:
ssh.exec_command('rm -f /tmp/pegaprox-ipmitool-install.sh')
except Exception:
pass
last = (out or '').strip().splitlines()[-1] if (out or '').strip() else ''
if 'PP_OK already_installed' in out:
return {'node': node, 'success': True, 'already_installed': True, 'detail': last}
if 'PP_OK installed' in out:
return {'node': node, 'success': True, 'already_installed': False, 'detail': last}
return {'node': node, 'success': False, 'error': last or 'unexpected installer output'}
except Exception as e:
return {'node': node, 'success': False, 'error': safe_error(e, 'install failed')}
finally:
if ssh:
try: ssh.close()
except Exception: pass
from pegaprox.utils.concurrent import run_per_node
res_map = run_per_node({n: _install_one for n in nodes}, max_concurrent=8, timeout=200)
results = [res_map.get(n) or {'node': n, 'success': False, 'error': 'timed out'} for n in nodes]
installed = sum(1 for r in results if r.get('success'))
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'hardware_monitoring.ipmitool_install',
f"ipmitool install: {installed}/{len(nodes)} node(s)",
cluster=mgr.config.name)
return jsonify({
'success': installed == len(nodes),
'installed': installed,
'total': len(nodes),
'results': results,
})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network/<iface>', methods=['PUT'])
@require_auth(perms=['node.network'])
def update_node_network_api(cluster_id, node, iface):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.update_node_network(node, iface, request.json or {})
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network', methods=['POST'])
@require_auth(perms=['node.network'])
def create_node_network_api(cluster_id, node):
"""Create a new network interface"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
data = request.json or {}
iface = data.get('iface', '')
iface_type = data.get('type', 'bridge')
if not iface:
return jsonify({'error': 'Interface name required'}), 400
config = {k: v for k, v in data.items() if k not in ['iface', 'type']}
result = mgr.create_node_network(node, iface, iface_type, config)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network/<iface>', methods=['DELETE'])
@require_auth(perms=['node.network'])
def delete_node_network_api(cluster_id, node, iface):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.delete_node_network(node, iface)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network', methods=['PUT'])
@require_auth(perms=['node.network'])
def apply_node_network_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.apply_node_network(node)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/network', methods=['DELETE'])
@require_auth(perms=['node.network'])
def revert_node_network_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.revert_node_network(node)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/networks', methods=['GET'])
@require_auth(perms=['node.view'])
def get_cluster_networks_api(cluster_id):
"""NS: Mar 2026 - Cluster-wide network overview with VM assignments"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
return jsonify(mgr.get_cluster_networks())
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/dns', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_dns_api(cluster_id, node):
# NS Jul 2026 (CodeAnt re-scan auth-bypass/IDOR) — cluster-scoped route was missing the tenant gate
ok, err = check_cluster_access(cluster_id)
if not ok:
return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_dns(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/dns', methods=['PUT'])
@require_auth(perms=['node.network'])
def update_node_dns_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.update_node_dns(node, request.json or {})
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/hosts', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_hosts_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
return jsonify({'data': cluster_managers[cluster_id].get_node_hosts(node)})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/hosts', methods=['POST'])
@require_auth(perms=['node.network'])
def update_node_hosts_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
data = request.json or {}
result = mgr.update_node_hosts(node, data.get('data', ''))
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/time', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_time_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
return jsonify(cluster_managers[cluster_id].get_node_time(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/time', methods=['PUT'])
@require_auth(perms=['node.network'])
def update_node_time_api(cluster_id, node):
"""Update node timezone"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
data = request.json or {}
result = mgr.update_node_time(node, data.get('timezone', 'UTC'))
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/syslog', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_syslog_api(cluster_id, node):
"""Get node system log"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
start = request.args.get('start', 0, type=int)
limit = request.args.get('limit', 500, type=int)
return jsonify(manager.get_node_syslog(node, start, limit))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/certificates', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_certificates_api(cluster_id, node):
"""Get node certificates"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_certificates(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/certificates/renew', methods=['POST'])
@require_auth(perms=['node.network'])
def renew_node_certificate_api(cluster_id, node):
"""Renew node certificate"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
force = request.json.get('force', False) if request.json else False
result = manager.renew_node_certificate(node, force)
if result['success']:
return jsonify(result)
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/certificates/custom', methods=['POST'])
@require_auth(perms=['node.network'])
def upload_node_certificate_api(cluster_id, node):
"""Upload custom certificate to node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
data = request.json or {}
certificates = data.get('certificates', '')
key = data.get('key', '')
restart = data.get('restart', True)
force = data.get('force', False)
if not certificates or not key:
return jsonify({'error': 'Certificate and key are required'}), 400
result = manager.upload_node_certificate(node, certificates, key, restart, force)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/certificates/custom', methods=['DELETE'])
@require_auth(perms=['node.network'])
def delete_node_certificate_api(cluster_id, node):
"""Delete custom certificate from node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
restart = request.args.get('restart', 'true').lower() == 'true'
result = mgr.delete_node_certificate(node, restart)
if result['success']:
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/replication', methods=['GET'])
@require_auth(perms=['cluster.view'])
def get_node_replication_api(cluster_id, node):
"""Get replication jobs for node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_replication(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/tasks', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_tasks_api(cluster_id, node):
"""Get task history for node, optionally filtered by vmid"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
start = request.args.get('start', 0, type=int)
limit = request.args.get('limit', 50, type=int)
errors = request.args.get('errors', 'false').lower() == 'true'
vmid = request.args.get('vmid', None, type=int)
tasks = manager.get_node_tasks(node, start, limit * 3 if vmid else limit, errors) # Get more if filtering
# Filter by vmid if specified
if vmid and tasks:
filtered = [t for t in tasks if t.get('id') == str(vmid) or str(vmid) in str(t.get('upid', ''))]
return jsonify(filtered[:limit])
return jsonify(tasks)
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/tasks/<path:upid>/log', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_task_log_api(cluster_id, node, upid):
"""Get log for a specific task
NS: Fixed Dec 2025 - frontend expects { log: "..." } format
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
start = request.args.get('start', 0, type=int)
limit = request.args.get('limit', 500, type=int)
log_lines = manager.get_node_task_log(node, upid, start, limit)
# Join lines into a single string for display
log_text = '\n'.join(log_lines) if log_lines else ''
return jsonify({'log': log_text})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/subscription', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_subscription_api(cluster_id, node):
"""Get node subscription status"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
manager = cluster_managers[cluster_id]
return jsonify(manager.get_node_subscription(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/subscription', methods=['PUT'])
@require_auth(perms=['admin.settings'])
def update_node_subscription_api(cluster_id, node):
"""Update subscription key - admin only"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
data = request.json or {}
result = mgr.update_node_subscription(node, data.get('key', ''))
if result['success']:
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'subscription.updated', f"Subscription key updated for {node}", cluster=mgr.config.name)
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/subscription', methods=['POST'])
@require_auth(perms=['admin.settings'])
def check_node_subscription_api(cluster_id, node):
"""Refresh subscription status - admin only"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
data = request.json or {}
result = mgr.check_node_subscription(node, bool(data.get('force', False)))
if result['success']:
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'subscription.checked', f"Subscription checked for {node}", cluster=mgr.config.name)
return jsonify({'message': 'Subscription check started', 'data': result.get('data')})
return jsonify({'error': result['error']}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/subscription', methods=['DELETE'])
@require_auth(perms=['admin.settings'])
def delete_node_subscription_api(cluster_id, node):
"""Delete subscription key - admin only"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
result = mgr.delete_node_subscription(node)
if result['success']:
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'subscription.deleted', f"Subscription key deleted for {node}", cluster=mgr.config.name)
return jsonify({'message': result['message']})
return jsonify({'error': result['error']}), 500
# =============================================================================
# SMBIOS Auto-Configurator Feature
# MK: Automatically sets SMBIOS data on new VMs for Windows licensing etc.
# MK: this was surprisingly tricky to get right, proxmox smbios format is picky
# =============================================================================
# NS 2026-04-24 — sudo detection is shared across SMBIOS helpers. We query once per
# SSH connection and cache on the client object so every helper sees the same answer.
def _ssh_sudo_prefix(ssh):
"""Returns '' when the SSH login user is already root, 'sudo -n ' otherwise.
Cached on the ssh object so we don't re-query for every command."""
cached = getattr(ssh, '_pegaprox_sudo_prefix', None)
if cached is not None:
return cached
try:
stdin, stdout, _ = ssh.exec_command('id -u', timeout=10)
uid = stdout.read().decode().strip()
prefix = '' if uid == '0' else 'sudo -n '
except Exception:
prefix = ''
try:
ssh._pegaprox_sudo_prefix = prefix
except Exception:
pass
return prefix
def _ssh_run_checked(ssh, cmd, timeout=30):
"""Run a command via SSH, raise on non-zero exit. NS Apr 2026 — used to be
scattered all over with `stdout.read()` and no check, which silently swallowed
systemctl failures during smbios deploy (#317 aftermath).
NS 2026-04-24 — auto-prefix with sudo when login isn't root (pegaprox@pam)."""
prefix = _ssh_sudo_prefix(ssh)
# Wrap multi-line / redirect commands through `sudo bash -c` via base64 so
# `>` and heredocs work. Single-token commands just get a plain prefix.
if prefix and (cmd.lstrip().startswith(('cat ', 'sh -c', 'bash -c')) or ' > ' in cmd or '\n' in cmd):
import base64 as _b64
enc = _b64.b64encode(cmd.encode('utf-8')).decode('ascii')
full = f"echo {enc} | base64 -d | sudo -n bash"
else:
full = f"{prefix}{cmd}"
stdin, stdout, stderr = ssh.exec_command(full, timeout=timeout)
out = stdout.read().decode('utf-8', errors='replace')
rc = stdout.channel.recv_exit_status()
err = stderr.read().decode('utf-8', errors='replace').strip()
if rc != 0:
raise RuntimeError(f"`{cmd}` failed (rc={rc}): {err or out[:200] or 'no output'}")
return out, err
def _ssh_write_file(ssh, path, content, mode=None):
"""Write file via SSH. Works for both root and non-root (pegaprox@pam + sudo) logins.
Strategy:
- root login: SFTP direct (fastest), fall back to `cat >`
- non-root login: upload to /tmp first (no sudo needed), then `sudo mv` into place.
SFTP-writing to /etc via sudo is a mess because paramiko's SFTP runs as the
login user — we bypass by staging in /tmp.
NS 2026-04-24 — pegaprox@pam + NOPASSWD sudo deployments
"""
import os
prefix = _ssh_sudo_prefix(ssh)
parent = os.path.dirname(path)
q_parent = shlex.quote(parent)
q_path = shlex.quote(path)
# mkdir first — mkdir -p is a no-op if the dir exists already
stdin, stdout, stderr = ssh.exec_command(f"{prefix}mkdir -p {q_parent}")
rc = stdout.channel.recv_exit_status()
if rc != 0:
err = stderr.read().decode('utf-8', errors='replace').strip()
raise RuntimeError(f"mkdir -p {parent} failed (rc={rc}): {err or 'permission denied?'}")
if not prefix:
# root login — original simple SFTP path
try:
sftp = ssh.open_sftp()
with sftp.file(path, 'w') as f:
f.write(content)
if mode is not None:
sftp.chmod(path, mode)
sftp.close()
except (IOError, OSError) as e:
logging.warning(f"SFTP write to {path} failed ({e}), falling back to exec_command")
stdin, stdout, stderr = ssh.exec_command(f"cat > {q_path}")
stdin.write(content)
stdin.channel.shutdown_write()
rc = stdout.channel.recv_exit_status()
if rc != 0:
err = stderr.read().decode('utf-8', errors='replace').strip()
raise RuntimeError(f"write {path} failed (rc={rc}): {err or 'unknown'}")
if mode is not None:
_ssh_run_checked(ssh, f"chmod {oct(mode)[2:]} {q_path}")
else:
# non-root — stage in /tmp, then sudo-mv into place
import uuid as _uuid
tmp_path = f"/tmp/pegaprox-stage-{_uuid.uuid4().hex[:12]}"
q_tmp = shlex.quote(tmp_path)
try:
sftp = ssh.open_sftp()
with sftp.file(tmp_path, 'w') as f:
f.write(content)
sftp.close()
except (IOError, OSError) as e:
logging.warning(f"SFTP stage to {tmp_path} failed ({e}), using stdin pipe")
stdin, stdout, stderr = ssh.exec_command(f"cat > {q_tmp}")
stdin.write(content); stdin.channel.shutdown_write()
rc = stdout.channel.recv_exit_status()
if rc != 0:
raise RuntimeError(f"staging {tmp_path} failed (rc={rc})")
# Move to final location as root, set mode if requested
mv_cmd = f"sudo -n mv {q_tmp} {q_path}"
if mode is not None:
mv_cmd += f" && sudo -n chmod {oct(mode)[2:]} {q_path}"
stdin, stdout, stderr = ssh.exec_command(mv_cmd)
rc = stdout.channel.recv_exit_status()
if rc != 0:
err = stderr.read().decode('utf-8', errors='replace').strip()
# best-effort cleanup so /tmp doesn't stay littered on failure
try: ssh.exec_command(f"rm -f {q_tmp}")
except Exception: pass
raise RuntimeError(f"sudo mv → {path} failed (rc={rc}): {err or 'permission denied?'}")
# Verify — file actually on disk with non-zero size
stdin, stdout, stderr = ssh.exec_command(f"{prefix}test -s {q_path}")
if stdout.channel.recv_exit_status() != 0:
raise RuntimeError(f"Write verification failed: {path} missing or empty")
SMBIOS_SCRIPT_TEMPLATE = '''#!/usr/bin/env python3
"""
SMBIOS Auto-Configurator for Proxmox VE
Deployed by PegaProx - automatically configures SMBIOS for new VMs
Runs as a systemd service, monitors for new VMs and sets SMBIOS data.
"""
import subprocess
import time
import os
import random
from datetime import datetime
# Configuration - set by PegaProx when deployed
MANUFACTURER = "{manufacturer}"
PRODUCT = "{product}"
VERSION = "{version}"
FAMILY = "{family}"
# Paths
LOG_FILE = "/var/log/pegaprox-smbios.log"
PROCESSED_VMS_FILE = "/var/lib/pegaprox-smbios-processed.txt" # keeps track of what we already did
def log_message(message):
"""write to log file, nothing fancy"""
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
log_entry = f"[{{timestamp}}] {{message}}"
print(log_entry)
try:
with open(LOG_FILE, 'a') as f:
f.write(log_entry + "\\n")
except:
pass # if we cant log, oh well
def get_all_vms():
"""read vmids from /etc/pve/.vmlist instead of spawning qm list (way cheaper)"""
try:
import json as _json
with open('/etc/pve/.vmlist', 'r') as f:
data = _json.load(f)
return [str(vmid) for vmid, info in data.get('ids', {{}}).items() if info.get('type') == 'qemu']
except Exception as e:
log_message(f"Error reading vmlist: {{e}}")
return []
def load_processed_vms():
"""Load already processed VMs"""
if os.path.exists(PROCESSED_VMS_FILE):
with open(PROCESSED_VMS_FILE, 'r') as f:
return set(f.read().splitlines())
return set()
def save_processed_vm(vmid):
"""Mark VM as processed"""
with open(PROCESSED_VMS_FILE, 'a') as f:
f.write(f"{{vmid}}\\n")
def get_current_smbios(vmid):
"""read smbios from conf file directly — no perl overhead"""
try:
conf_path = f"/etc/pve/qemu-server/{{vmid}}.conf"
if not os.path.exists(conf_path):
return None
with open(conf_path, 'r') as f:
for line in f:
if line.startswith('smbios1:'):
return line.split(':', 1)[1].strip()
return None
except:
return None
def parse_smbios_string(smbios_str):
"""Parse SMBIOS string into dictionary"""
params = {{}}
if not smbios_str:
return params
for part in smbios_str.split(','):
if '=' in part:
key, value = part.split('=', 1)
params[key.strip()] = value.strip()
return params
def needs_smbios_update(vmid):
"""Check if VM needs SMBIOS configuration"""
smbios_str = get_current_smbios(vmid)
if not smbios_str:
return True
params = parse_smbios_string(smbios_str)
relevant_params = {{k: v for k, v in params.items() if k != 'uuid'}}
if not relevant_params:
return True
if ('manufacturer' in params or 'product' in params or
'version' in params or 'serial' in params or 'family' in params):
log_message(f"VM {{vmid}} already has SMBIOS configuration")
return False
return True
def generate_unique_serial():
"""Generate unique serial number"""
timestamp = datetime.now().strftime("%y%m%d%H%M%S")
random_part = random.randint(1000, 9999)
return f"PVE{{timestamp}}{{random_part}}"
def set_smbios(vmid):
"""Set SMBIOS configuration"""
current_smbios_str = get_current_smbios(vmid)
current_params = parse_smbios_string(current_smbios_str) if current_smbios_str else {{}}
uuid_value = current_params.get('uuid', '')
serial = generate_unique_serial()
smbios_parts = [
f"manufacturer={{MANUFACTURER}}",
f"product={{PRODUCT}}",
f"version={{VERSION}}",
f"serial={{serial}}",
]
if uuid_value:
smbios_parts.append(f"uuid={{uuid_value}}")
smbios_parts.append(f"family={{FAMILY}}")
smbios_string = ",".join(smbios_parts)
cmd = ['qm', 'set', vmid, '-smbios1', smbios_string]
try:
subprocess.run(cmd, capture_output=True, text=True, check=True)
log_message(f"Set SMBIOS for VM {{vmid}} | Serial: {{serial}}")
return True
except subprocess.CalledProcessError as e:
log_message(f"Error setting SMBIOS for VM {{vmid}}: {{e.stderr if e.stderr else e}}")
return False
def check_vm_exists(vmid):
"""Check if VM config exists on this node"""
return os.path.exists(f"/etc/pve/qemu-server/{{vmid}}.conf")
def cleanup_processed_list(processed):
"""Remove VMs that no longer exist"""
current_vms = set(get_all_vms())
removed_vms = processed - current_vms
if removed_vms:
for vmid in removed_vms:
log_message(f"VM {{vmid}} no longer exists, removing from tracking")
processed.remove(vmid)
with open(PROCESSED_VMS_FILE, 'w') as f:
for vmid in processed:
f.write(f"{{vmid}}\\n")
return processed
def main():
log_message("=== PegaProx SMBIOS Auto-Configurator started ===")
log_message(f"Config: {{MANUFACTURER}} | {{PRODUCT}} | {{VERSION}} | {{FAMILY}}")
processed = load_processed_vms()
cleanup_counter = 0
while True:
try:
cleanup_counter += 1
if cleanup_counter >= 30:
processed = cleanup_processed_list(processed)
cleanup_counter = 0
current_vms = get_all_vms()
for vmid in current_vms:
if vmid not in processed:
if check_vm_exists(vmid):
if needs_smbios_update(vmid):
log_message(f"Configuring SMBIOS for new VM {{vmid}}")
if set_smbios(vmid):
save_processed_vm(vmid)
processed.add(vmid)
else:
save_processed_vm(vmid)
processed.add(vmid)
elif vmid in processed and check_vm_exists(vmid) and needs_smbios_update(vmid):
# VM was deleted + recreated with same ID
log_message(f"VM {{vmid}} re-created, reconfiguring SMBIOS")
set_smbios(vmid)
time.sleep(30)
except KeyboardInterrupt:
log_message("=== SMBIOS Auto-Configurator stopped ===")
break
except Exception as e:
log_message(f"Error: {{e}}")
time.sleep(10)
if __name__ == "__main__":
main()
'''
SMBIOS_SERVICE_TEMPLATE = '''[Unit]
Description=PegaProx SMBIOS Auto-Configurator
After=pve-cluster.service
Wants=pve-cluster.service
[Service]
Type=simple
ExecStart=/usr/bin/python3 /opt/pegaprox-smbios-autoconfig.py
Restart=always
RestartSec=10
StandardOutput=journal
StandardError=journal
[Install]
WantedBy=multi-user.target
'''
@bp.route('/api/clusters/<cluster_id>/smbios-autoconfig', methods=['GET'])
@require_auth(perms=['node.view'])
def get_smbios_autoconfig(cluster_id):
"""get smbios settings for the cluster, returns defaults if not configured yet"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
try:
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
# defaults if nothing configured - NS: proxmox doesnt allow underscores so no spaces either
mgr = cluster_managers[cluster_id]
settings = getattr(mgr.config, 'smbios_autoconfig', None) or {
'enabled': False,
'manufacturer': 'Proxmox',
'product': 'PegaProxManagment',
'version': 'v1',
'family': 'ProxmoxVE'
}
return jsonify(settings)
except Exception as e:
logging.error(f"Error getting SMBIOS config: {e}")
return jsonify({'error': safe_error(e, 'Failed to get SMBIOS config')}), 500
@bp.route('/api/clusters/<cluster_id>/smbios-autoconfig', methods=['PUT'])
@require_auth(perms=['admin.settings'])
def update_smbios_autoconfig(cluster_id):
"""save smbios settings - gets deployed to nodes when they click deploy"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
data = request.json or {}
mgr = cluster_managers[cluster_id]
# NS Feb 2026: Validate SMBIOS values to prevent template injection
for key in ['manufacturer', 'product', 'version', 'family']:
val = data.get(key, '')
if val and not re.match(r'^[a-zA-Z0-9 ._-]{1,64}$', val):
return jsonify({'error': f'Invalid {key}: only alphanumeric, spaces, dots, hyphens allowed (max 64 chars)'}), 400
# Update settings
mgr.config.smbios_autoconfig = {
'enabled': data.get('enabled', False),
'manufacturer': data.get('manufacturer', 'Proxmox'),
'product': data.get('product', 'Virtual Machine'),
'version': data.get('version', 'PVE8'),
'family': data.get('family', 'ProxmoxVE')
}
# Save to database
db = get_db()
db.update_cluster(cluster_id, {'smbios_autoconfig': json.dumps(mgr.config.smbios_autoconfig)})
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'smbios_autoconfig.updated', f"SMBIOS auto-config updated", cluster=mgr.config.name)
return jsonify({'success': True, 'message': 'Settings saved'})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/smbios-autoconfig/status', methods=['GET'])
@require_auth(perms=['node.view'])
def get_smbios_autoconfig_status(cluster_id, node):
"""Check if SMBIOS auto-config service is running on node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
try:
# NS Apr 2026 (PR #324): one real resolver instead of 3 copies of the
# cluster/status hack. cluster/status can point at the corosync IP.
node_ip = mgr._get_node_ip(node)
if not node_ip:
return jsonify({'installed': False, 'running': False, 'error': f'Could not determine SSH-reachable IP for node {node}'})
ssh = mgr._ssh_connect(node_ip)
if not ssh:
return jsonify({'installed': False, 'running': False, 'error': 'SSH not available - check SSH key in cluster settings'})
# Check if script exists
stdin, stdout, stderr = ssh.exec_command('test -f /opt/pegaprox-smbios-autoconfig.py && echo exists')
installed = 'exists' in stdout.read().decode()
# Check if service is running
stdin, stdout, stderr = ssh.exec_command('systemctl is-active pegaprox-smbios-autoconfig 2>/dev/null || echo inactive')
status = stdout.read().decode().strip()
running = status == 'active'
# Get last log entries
stdin, stdout, stderr = ssh.exec_command('tail -5 /var/log/pegaprox-smbios.log 2>/dev/null || echo "No logs yet"')
logs = stdout.read().decode().strip()
ssh.close()
return jsonify({
'installed': installed,
'running': running,
'status': status,
'logs': logs
})
except Exception as e:
logging.error(f"Error checking SMBIOS autoconfig status: {e}")
return jsonify({'installed': False, 'running': False, 'error': safe_error(e, 'Failed to get node status')})
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/smbios-autoconfig/deploy', methods=['POST'])
@require_auth(perms=['admin.settings'])
def deploy_smbios_autoconfig(cluster_id, node):
"""Deploy SMBIOS auto-config script to node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
settings = getattr(mgr.config, 'smbios_autoconfig', None) or {}
ssh = None
try:
# PR #324: resolve SSH-reachable mgmt IP via manager (not cluster/status)
node_ip = mgr._get_node_ip(node)
if not node_ip:
return jsonify({'error': f'Could not determine SSH-reachable IP for node {node}'}), 502
# NS Apr 2026 - retry SSH 3x, transient failures shouldn't kill the whole deploy
for attempt in range(3):
ssh = mgr._ssh_connect(node_ip)
if ssh:
break
if attempt < 2:
time.sleep(1.5)
if not ssh:
return jsonify({'error': 'SSH connection failed after 3 attempts - check SSH key in cluster settings'}), 500
# Generate script with settings (defense-in-depth: strip quotes/backslashes - NS Feb 2026)
def _sanitize_smbios(val):
"""Strip characters dangerous in Python string literals as defense-in-depth."""
return re.sub(r'[^a-zA-Z0-9 ._-]', '', str(val))[:64]
script = SMBIOS_SCRIPT_TEMPLATE.format(
manufacturer=_sanitize_smbios(settings.get('manufacturer', 'Proxmox')),
product=_sanitize_smbios(settings.get('product', 'PegaProxManagment')),
version=_sanitize_smbios(settings.get('version', 'v1')),
family=_sanitize_smbios(settings.get('family', 'ProxmoxVE'))
)
# Write script and service to node (SFTP with exec_command fallback) — raises on failure
_ssh_write_file(ssh, '/opt/pegaprox-smbios-autoconfig.py', script, 0o755)
_ssh_write_file(ssh, '/etc/systemd/system/pegaprox-smbios-autoconfig.service', SMBIOS_SERVICE_TEMPLATE)
# clear processed list + enable/start. Each command is checked; a failed
# systemctl used to be swallowed and we'd return success anyway
_ssh_run_checked(ssh, 'rm -f /var/lib/pegaprox-smbios-processed.txt')
_ssh_run_checked(ssh, 'systemctl daemon-reload')
_ssh_run_checked(ssh, 'systemctl enable pegaprox-smbios-autoconfig')
_ssh_run_checked(ssh, 'systemctl restart pegaprox-smbios-autoconfig')
# confirm it actually became active (systemctl start/restart can succeed even when unit fails)
stdin, stdout, stderr = ssh.exec_command('systemctl is-active pegaprox-smbios-autoconfig')
active = stdout.read().decode('utf-8', errors='replace').strip()
stdout.channel.recv_exit_status()
if active != 'active':
# grab last log lines for context
stdin, stdout, stderr = ssh.exec_command('journalctl -u pegaprox-smbios-autoconfig -n 10 --no-pager 2>/dev/null | tail -10')
log_tail = stdout.read().decode('utf-8', errors='replace').strip()
raise RuntimeError(f"service not active (state={active}). Last log: {log_tail[:400]}")
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'smbios_autoconfig.deployed', f"SMBIOS auto-config deployed to {node}", cluster=mgr.config.name)
return jsonify({'success': True, 'message': f'SMBIOS Auto-Config deployed to {node}'})
except Exception as e:
logging.error(f"Error deploying SMBIOS autoconfig to {_sl(node)}: {e}")
return jsonify({'error': safe_error(e, 'SMBIOS deploy failed')}), 500
finally:
if ssh:
try: ssh.close()
except: pass
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/smbios-autoconfig', methods=['DELETE'])
@require_auth(perms=['admin.settings'])
def remove_smbios_autoconfig(cluster_id, node):
"""Remove SMBIOS auto-config from node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
try:
node_ip = mgr._get_node_ip(node)
if not node_ip:
return jsonify({'error': f'Could not determine SSH-reachable IP for node {node}'}), 502
ssh = mgr._ssh_connect(node_ip)
if not ssh:
return jsonify({'error': 'SSH connection failed - check SSH key in cluster settings'}), 500
# Stop and disable service, remove files. _ssh_run_checked handles sudo
# for pegaprox@pam deployments (NS 2026-04-24).
commands = [
'systemctl stop pegaprox-smbios-autoconfig 2>/dev/null || true',
'systemctl disable pegaprox-smbios-autoconfig 2>/dev/null || true',
'rm -f /etc/systemd/system/pegaprox-smbios-autoconfig.service',
'rm -f /opt/pegaprox-smbios-autoconfig.py',
'rm -f /var/lib/pegaprox-smbios-processed.txt',
'systemctl daemon-reload'
]
for cmd in commands:
try:
_ssh_run_checked(ssh, cmd)
except RuntimeError as rerr:
# `|| true` lines stay silent; only bubble hard failures on the unconditional ones
if cmd.endswith('|| true'):
continue
logging.warning(f"SMBIOS remove: `{cmd}` → {rerr}")
ssh.close()
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'smbios_autoconfig.removed', f"SMBIOS auto-config removed from {node}", cluster=mgr.config.name)
return jsonify({'success': True, 'message': f'SMBIOS Auto-Config removed from {node}'})
except Exception as e:
logging.error(f"Error removing SMBIOS autoconfig: {e}")
return jsonify({'error': safe_error(e, 'SMBIOS removal failed')}), 500
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/smbios-autoconfig/control', methods=['POST'])
@require_auth(perms=['admin.settings'])
def control_smbios_autoconfig(cluster_id, node):
"""Start/Stop/Rescan SMBIOS auto-config service on node"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
data = request.get_json() or {}
action = data.get('action') # 'start', 'stop', 'restart', 'rescan'
if action not in ['start', 'stop', 'restart', 'rescan']:
return jsonify({'error': 'Invalid action. Use start, stop, restart, or rescan'}), 400
mgr = cluster_managers[cluster_id]
try:
node_ip = mgr._get_node_ip(node)
if not node_ip:
return jsonify({'error': f'Could not determine SSH-reachable IP for node {node}'}), 502
ssh = mgr._ssh_connect(node_ip)
if not ssh:
return jsonify({'error': 'SSH connection failed'}), 500
# NS: rescan = nuke the processed list and restart, forces re-check of all VMs
if action == 'rescan':
cmd = 'rm -f /var/lib/pegaprox-smbios-processed.txt && systemctl restart pegaprox-smbios-autoconfig'
else:
cmd = f'systemctl {action} pegaprox-smbios-autoconfig'
err_output = ''
try:
_ssh_run_checked(ssh, cmd)
except RuntimeError as rerr:
err_output = str(rerr)
ssh.close()
if err_output and 'not found' in err_output.lower():
return jsonify({'error': 'Service not installed on this node'}), 404
if err_output:
return jsonify({'error': err_output}), 500
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, f'smbios_autoconfig.{action}', f"SMBIOS auto-config {action} on {node}", cluster=mgr.config.name)
return jsonify({'success': True, 'message': f'Service {action}ed on {node}'})
except Exception as e:
logging.error(f"Error controlling SMBIOS autoconfig: {e}")
return jsonify({'error': safe_error(e, 'SMBIOS service control failed')}), 500
@bp.route('/api/clusters/<cluster_id>/smbios-autoconfig/status-all', methods=['GET'])
@require_auth(perms=['node.view'])
def get_smbios_autoconfig_status_all(cluster_id):
"""Get SMBIOS auto-config status for ALL nodes in cluster"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
try:
# #198: get nodes with online/offline status
node_status = {}
try:
node_status = mgr.get_node_status() or {}
except:
pass
node_names = list(node_status.keys()) if node_status else []
if not node_names:
try:
nodes = mgr.get_nodes()
node_names = [n.get('node', n.get('name', '')) for n in nodes if n]
except:
node_names = []
if not node_names:
return jsonify({'error': 'No nodes available'}), 400
# PR #324: _get_node_ip() per node -- cluster/status returns the
# corosync ring IP on separate-VLAN setups which we can't reach via SSH
results = {}
for node_name in node_names:
if not node_name:
continue
# #198: skip offline nodes early
ns = node_status.get(node_name, {})
if ns.get('offline') or ns.get('status') == 'offline':
results[node_name] = {'installed': False, 'running': False, 'error': 'Node offline'}
continue
try:
node_ip = mgr._get_node_ip(node_name)
if not node_ip:
results[node_name] = {'installed': False, 'running': False, 'error': 'Could not determine node IP'}
continue
ssh = mgr._ssh_connect(node_ip, retries=1)
if not ssh:
results[node_name] = {'installed': False, 'running': False, 'error': 'SSH not available'}
continue
try:
# Check if script exists
stdin, stdout, stderr = ssh.exec_command('test -f /opt/pegaprox-smbios-autoconfig.py && echo exists')
installed = 'exists' in stdout.read().decode()
# Check if service is running
stdin, stdout, stderr = ssh.exec_command('systemctl is-active pegaprox-smbios-autoconfig 2>/dev/null || echo inactive')
status = stdout.read().decode().strip()
running = status == 'active'
results[node_name] = {
'installed': installed,
'running': running,
'status': status
}
finally:
ssh.close()
except Exception as e:
results[node_name] = {'installed': False, 'running': False, 'error': safe_error(e, 'Failed to get node status')}
return jsonify(results)
except Exception as e:
logging.error(f"Error getting SMBIOS autoconfig status: {e}")
return jsonify({'error': safe_error(e, 'Failed to get SMBIOS status')}), 500
@bp.route('/api/clusters/<cluster_id>/smbios-autoconfig/deploy-all', methods=['POST'])
@require_auth(perms=['admin.settings'])
def deploy_smbios_autoconfig_all(cluster_id):
"""Deploy SMBIOS auto-config script to ALL nodes in cluster"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
settings = getattr(mgr.config, 'smbios_autoconfig', None) or {}
# Get all nodes in cluster. PR #324: we no longer pre-pin IPs from
# cluster/status (that's the corosync ring IP on VLAN setups) --
# each node gets resolved via _get_node_ip inside the loop.
nodes = []
try:
cluster_host, cluster_port = mgr.host, mgr.api_port
status_url = f"https://{cluster_host}:{cluster_port}/api2/json/cluster/status"
r = mgr._create_session().get(status_url, timeout=10)
if r.status_code == 200:
for item in r.json().get('data', []):
if item.get('type') == 'node':
nodes.append(item.get('name'))
else:
# Single node cluster - just use cluster host
nodes = [mgr.config.host.split('.')[0]]
except Exception as e:
logging.error(f"Error getting cluster nodes: {e}")
return jsonify({'error': f'Could not get cluster nodes: {e}'}), 500
if not nodes:
return jsonify({'error': 'No nodes found in cluster'}), 404
results = []
def _sanitize_smbios_val(val):
"""Strip characters dangerous in Python string literals as defense-in-depth."""
return re.sub(r'[^a-zA-Z0-9 ._-]', '', str(val))[:64]
script = SMBIOS_SCRIPT_TEMPLATE.format(
manufacturer=_sanitize_smbios_val(settings.get('manufacturer', 'Proxmox')),
product=_sanitize_smbios_val(settings.get('product', 'PegaProxManagment')),
version=_sanitize_smbios_val(settings.get('version', 'v1')),
family=_sanitize_smbios_val(settings.get('family', 'ProxmoxVE'))
)
for node in nodes:
node_ip = mgr._get_node_ip(node)
if not node_ip:
results.append({'node': node, 'success': False, 'error': 'Could not determine SSH-reachable IP'})
continue
ssh = None
try:
# NS: Staggered connections to prevent SSH server overload
if results: # Not the first node
time.sleep(1.0)
# retry transient SSH failures before giving up
for attempt in range(3):
ssh = mgr._ssh_connect(node_ip)
if ssh:
break
if attempt < 2:
time.sleep(1.5)
if not ssh:
results.append({'node': node, 'success': False, 'error': 'SSH connection failed after 3 attempts'})
continue
# Write script and service - _ssh_write_file now raises on failure
_ssh_write_file(ssh, '/opt/pegaprox-smbios-autoconfig.py', script, 0o755)
_ssh_write_file(ssh, '/etc/systemd/system/pegaprox-smbios-autoconfig.service', SMBIOS_SERVICE_TEMPLATE)
# NS: clear processed list + checked systemctl commands
_ssh_run_checked(ssh, 'rm -f /var/lib/pegaprox-smbios-processed.txt')
_ssh_run_checked(ssh, 'systemctl daemon-reload')
_ssh_run_checked(ssh, 'systemctl enable pegaprox-smbios-autoconfig')
_ssh_run_checked(ssh, 'systemctl restart pegaprox-smbios-autoconfig')
# verify active
stdin, stdout, stderr = ssh.exec_command('systemctl is-active pegaprox-smbios-autoconfig')
active = stdout.read().decode('utf-8', errors='replace').strip()
stdout.channel.recv_exit_status()
if active != 'active':
raise RuntimeError(f"service not active (state={active})")
results.append({'node': node, 'success': True})
except Exception as e:
results.append({'node': node, 'success': False, 'error': safe_error(e, 'SMBIOS deploy failed')})
finally:
if ssh:
try: ssh.close()
except: pass
success_count = sum(1 for r in results if r['success'])
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'smbios_autoconfig.deployed_all', f"SMBIOS auto-config deployed to {success_count}/{len(nodes)} nodes", cluster=mgr.config.name)
return jsonify({
'success': success_count == len(nodes),
'message': f'Deployed to {success_count}/{len(nodes)} nodes',
'results': results
})
# =============================================================================
# StarWind "starlvm" SAN plugin install
# MK: one-click install of the StarWind x Proxmox SAN Integration plugin (the
# custom `starlvm` storage type — thin snapshots on a shared SAN LUN) across the
# cluster. We deliberately do NOT run StarWind's own installer script: it falls
# back to an UNSIGNED (trusted=yes) apt repo if its signed-repo test hiccups. We
# build a SIGNED deb822 source ourselves (key over HTTPS, Signed-By) so apt
# enforces the signature on install — no unsigned fallback. Repo/key are
# admin-overridable for air-gapped mirrors. Per-node, idempotent, bounded-parallel.
# =============================================================================
STARWIND_REPO_DEFAULT = 'http://repo.starwind.com/proxmox/'
STARWIND_KEY_DEFAULT = 'https://repo.starwind.com/keys/repo_public.key'
# these strings land inside a root-run bash script, so keep the charset tight
def _safe_repo_url(u, default):
u = (u or '').strip() or default
if len(u) > 300 or not re.match(r'^https?://[A-Za-z0-9._~:/\-]+$', u):
raise ValueError(f'unsafe repo/key url: {u[:60]}')
# NS Jul 2026 (CodeAnt SSRF) — this URL is curl'd / apt-fetched inside a ROOT bash script on
# the node; the charset check alone still allows internal/metadata IPs. Gate it (allow_private
# keeps internal package mirrors working; cloud-metadata endpoints stay blocked).
from pegaprox.utils.url_security import is_safe_outbound_url
_ok, _why = is_safe_outbound_url(u, allowed_schemes=('http', 'https'), allow_private=True)
if not _ok:
raise ValueError(f'repo/key url rejected by SSRF guard: {_why}')
return u
STARLVM_INSTALL_SCRIPT = """#!/usr/bin/env bash
# PegaProx — install the StarWind x Proxmox SAN plugin (starlvm) on this node.
set -uo pipefail
REPO_URL='__REPO_URL__'
KEY_URL='__KEY_URL__'
FORCE='__FORCE__'
KEYRING='/usr/share/keyrings/starwind-proxmox.gpg'
SRC='/etc/apt/sources.list.d/starwind-proxmox.sources'
pv="$(pveversion 2>/dev/null | grep -oE 'pve-manager/[0-9]+' | grep -oE '[0-9]+' | head -1 || true)"
if [ -z "${pv:-}" ]; then echo 'PP_ERR not-a-proxmox-node'; exit 3; fi
if [ "$pv" -ge 9 ]; then PKG='starwind-proxmox-plugin-pve9'; SUITE='trixie'; else PKG='starwind-proxmox-plugin'; SUITE='master'; fi
# idempotent: skip if already there (unless forced)
if [ "$FORCE" != '1' ] && dpkg-query -W -f='${Status}' "$PKG" 2>/dev/null | grep -q 'install ok installed'; then
ver="$(dpkg-query -W -f='${Version}' "$PKG" 2>/dev/null || echo '?')"
echo "PP_OK already_installed $PKG $ver pve$pv"; exit 0
fi
command -v curl >/dev/null 2>&1 || { echo 'PP_ERR curl-missing'; exit 3; }
command -v gpg >/dev/null 2>&1 || { echo 'PP_ERR gpg-missing'; exit 3; }
# key over HTTPS, dearmored into its own keyring
if ! curl -fsSL "$KEY_URL" | gpg --dearmor --yes --output "$KEYRING" 2>/dev/null; then
echo 'PP_ERR key-fetch-failed'; exit 4
fi
chmod 0644 "$KEYRING"
# SIGNED deb822 source — we never write trusted=yes
cat > "$SRC" <<EOF
Types: deb
URIs: ${REPO_URL}
Suites: ${SUITE}
Components: main
Signed-By: ${KEYRING}
EOF
# scrub any legacy unsigned config a previous StarWind install may have left
rm -f /etc/apt/sources.list.d/starwind-proxmox.list /etc/apt/trusted.gpg.d/starwind-proxmox.gpg 2>/dev/null || true
apt-get update -o Acquire::Retries=2 >/tmp/pp_starlvm_apt.log 2>&1 || true
# on PVE9 drop the old bookworm package (StarWind's documented 8->9 upgrade step)
if [ "$pv" -ge 9 ]; then DEBIAN_FRONTEND=noninteractive apt-get remove -y starwind-proxmox-plugin >/dev/null 2>&1 || true; fi
# apt refuses an unverifiable Signed-By source, so this is the real security gate
if ! DEBIAN_FRONTEND=noninteractive apt-get install -y "$PKG" >>/tmp/pp_starlvm_apt.log 2>&1; then
echo 'PP_ERR apt-install-failed'; exit 5
fi
# register the custom storage plugin with PVE
systemctl restart pvedaemon pveproxy >/dev/null 2>&1 || true
dpkg-query -W -f='${Status}' "$PKG" 2>/dev/null | grep -q 'install ok installed' || { echo 'PP_ERR not-installed-after-apt'; exit 6; }
[ -f /usr/share/perl5/PVE/Storage/Custom/StarLvmPlugin.pm ] || { echo 'PP_ERR plugin-pm-missing'; exit 6; }
ver="$(dpkg-query -W -f='${Version}' "$PKG" 2>/dev/null || echo '?')"
echo "PP_OK installed $PKG $ver pve$pv"
"""
STARLVM_STATUS_SCRIPT = """set -uo pipefail
pv="$(pveversion 2>/dev/null | grep -oE 'pve-manager/[0-9]+' | grep -oE '[0-9]+' | head -1 || true)"
pm='no'; [ -f /usr/share/perl5/PVE/Storage/Custom/StarLvmPlugin.pm ] && pm='yes'
inst='no'; pkg='none'; ver='none'
for p in starwind-proxmox-plugin-pve9 starwind-proxmox-plugin; do
if dpkg-query -W -f='${Status}' "$p" 2>/dev/null | grep -q 'install ok installed'; then
inst='yes'; pkg="$p"; ver="$(dpkg-query -W -f='${Version}' "$p" 2>/dev/null)"; break
fi
done
echo "PP_STATUS installed=$inst pkg=$pkg ver=$ver pm=$pm pve=${pv:-?}"
"""
def _cluster_node_names(mgr):
"""Node short-names from cluster/status; falls back to the config host for single-node."""
try:
r = mgr._create_session().get(f"https://{mgr.host}:{mgr.api_port}/api2/json/cluster/status", timeout=10)
if r.status_code == 200:
names = [it.get('name') for it in r.json().get('data', []) if it.get('type') == 'node']
if names:
return names
except Exception as e:
logging.warning(f"[starlvm] node enumerate failed: {e}")
return [mgr.config.host.split('.')[0]]
@bp.route('/api/clusters/<cluster_id>/storage/starlvm/install', methods=['POST'])
@require_auth(perms=['admin.settings']) # MK: root apt-install on every node + caller-overridable repo/key → admin-only (not node.maintenance, which tenant_operator holds)
def install_starlvm_plugin(cluster_id):
"""Install the StarWind SAN plugin (starlvm storage type) on cluster nodes over SSH.
Body (all optional): {repo_url, key_url, force: bool, nodes: [names]}.
Signed deb822 source only — no unsigned fallback. Idempotent per node."""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if getattr(mgr, 'cluster_type', 'proxmox') != 'proxmox':
return jsonify({'error': 'StarWind plugin install is Proxmox-only'}), 400
body = request.get_json(silent=True) or {}
try:
repo_url = _safe_repo_url(body.get('repo_url'), STARWIND_REPO_DEFAULT)
key_url = _safe_repo_url(body.get('key_url'), STARWIND_KEY_DEFAULT)
except ValueError as e:
return jsonify({'error': str(e)}), 400
force = '1' if body.get('force') else '0'
only = set(body.get('nodes') or [])
script = (STARLVM_INSTALL_SCRIPT
.replace('__REPO_URL__', repo_url)
.replace('__KEY_URL__', key_url)
.replace('__FORCE__', force))
nodes = _cluster_node_names(mgr)
if only:
nodes = [n for n in nodes if n in only]
if not nodes:
return jsonify({'error': 'No nodes found in cluster'}), 404
def _install_one(node):
node_ip = mgr._get_node_ip(node)
if not node_ip:
return {'node': node, 'success': False, 'error': 'no SSH-reachable IP'}
ssh = None
try:
for attempt in range(3):
ssh = mgr._ssh_connect(node_ip)
if ssh:
break
if attempt < 2:
time.sleep(1.5)
if not ssh:
return {'node': node, 'success': False, 'error': 'SSH connect failed after 3 tries'}
_ssh_write_file(ssh, '/tmp/pegaprox-starlvm-install.sh', script, 0o755)
out, _e = _ssh_run_checked(ssh, 'bash /tmp/pegaprox-starlvm-install.sh', timeout=200)
try:
ssh.exec_command('rm -f /tmp/pegaprox-starlvm-install.sh')
except Exception:
pass
last = (out or '').strip().splitlines()[-1] if (out or '').strip() else ''
if 'PP_OK already_installed' in out:
return {'node': node, 'success': True, 'already_installed': True, 'detail': last}
if 'PP_OK installed' in out:
return {'node': node, 'success': True, 'already_installed': False, 'detail': last}
return {'node': node, 'success': False, 'error': last or 'unexpected installer output'}
except Exception as e:
return {'node': node, 'success': False, 'error': safe_error(e, 'install failed')}
finally:
if ssh:
try: ssh.close()
except Exception: pass
from pegaprox.utils.concurrent import run_per_node
res_map = run_per_node({n: _install_one for n in nodes}, max_concurrent=8, timeout=240)
results = [res_map.get(n) or {'node': n, 'success': False, 'error': 'timed out'} for n in nodes]
installed = sum(1 for r in results if r.get('success'))
usr = getattr(request, 'session', {}).get('user', 'system')
log_audit(usr, 'storage.starlvm_install',
f"StarWind SAN plugin install: {installed}/{len(nodes)} nodes (repo={repo_url})",
cluster=mgr.config.name)
return jsonify({
'success': installed == len(nodes),
'installed': installed,
'total': len(nodes),
'results': results,
})
@bp.route('/api/clusters/<cluster_id>/storage/starlvm/status', methods=['GET'])
@require_auth(perms=['storage.view'])
def starlvm_plugin_status(cluster_id):
"""Per-node install state of the StarWind starlvm plugin (read-only probe)."""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if getattr(mgr, 'cluster_type', 'proxmox') != 'proxmox':
return jsonify({'error': 'Proxmox-only'}), 400
import base64 as _b64
enc = _b64.b64encode(STARLVM_STATUS_SCRIPT.encode('utf-8')).decode('ascii')
def _probe(node):
node_ip = mgr._get_node_ip(node)
if not node_ip:
return {'node': node, 'reachable': False}
ssh = None
try:
ssh = mgr._ssh_connect(node_ip)
if not ssh:
return {'node': node, 'reachable': False}
# read-only, no sudo needed (dpkg-query / file test are world-readable)
_in, _out, _err = ssh.exec_command(f"echo {enc} | base64 -d | bash", timeout=20)
out = _out.read().decode('utf-8', errors='replace')
_out.channel.recv_exit_status()
d = {}
for tok in out.split():
if '=' in tok:
k, v = tok.split('=', 1)
d[k] = v
return {
'node': node,
'reachable': True,
'installed': d.get('installed') == 'yes',
'package': None if d.get('pkg', 'none') == 'none' else d.get('pkg'),
'version': None if d.get('ver', 'none') == 'none' else d.get('ver'),
'plugin_file': d.get('pm') == 'yes',
'pve_major': d.get('pve'),
}
except Exception as e:
return {'node': node, 'reachable': False, 'error': safe_error(e, 'probe failed')}
finally:
if ssh:
try: ssh.close()
except Exception: pass
nodes = _cluster_node_names(mgr)
from pegaprox.utils.concurrent import run_per_node
res_map = run_per_node({n: _probe for n in nodes}, max_concurrent=8, timeout=60)
results = [res_map.get(n) or {'node': n, 'reachable': False} for n in nodes]
return jsonify({
'nodes': results,
'installed_all': bool(results) and all(r.get('installed') for r in results),
})
# =============================================================================
# Custom Scripts Feature
# MK: Run custom .sh/.py scripts on cluster nodes with permission control
# =============================================================================
@bp.route('/api/clusters/<cluster_id>/scripts', methods=['GET'])
@require_auth(perms=['admin.scripts'])
def get_custom_scripts(cluster_id):
"""Get all custom scripts for a cluster (excludes soft-deleted)"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
try:
db = get_db()
# Ensure table exists with soft delete support
db.execute('''
CREATE TABLE IF NOT EXISTS custom_scripts (
id TEXT PRIMARY KEY,
cluster_id TEXT NOT NULL,
name TEXT NOT NULL,
description TEXT DEFAULT '',
type TEXT DEFAULT 'bash',
content TEXT NOT NULL,
target_nodes TEXT DEFAULT 'all',
enabled INTEGER DEFAULT 1,
last_run TEXT,
last_status TEXT,
last_output TEXT,
created_at TEXT,
updated_at TEXT,
created_by TEXT,
deleted_at TEXT,
deleted_by TEXT
)
''')
# Add columns if they don't exist (migration for existing tables)
try:
db.execute('ALTER TABLE custom_scripts ADD COLUMN deleted_at TEXT')
except: pass
try:
db.execute('ALTER TABLE custom_scripts ADD COLUMN deleted_by TEXT')
except: pass
try:
db.execute('ALTER TABLE custom_scripts ADD COLUMN created_by TEXT')
except: pass
try:
db.execute('ALTER TABLE custom_scripts ADD COLUMN last_output TEXT')
except: pass
# Only return non-deleted scripts
scripts = db.query(
'SELECT * FROM custom_scripts WHERE cluster_id = ? AND deleted_at IS NULL ORDER BY name',
(cluster_id,)
)
return jsonify([dict(s) for s in scripts] if scripts else [])
except Exception as e:
logging.error(f"Error loading scripts: {e}")
return jsonify({'error': safe_error(e, 'Failed to load scripts')}), 500
# Cleanup job for permanently deleting scripts after 20 days
def cleanup_deleted_scripts():
"""Permanently delete scripts that have been soft-deleted for 20+ days"""
try:
db = get_db()
cutoff = (datetime.now() - timedelta(days=20)).isoformat()
deleted = db.query(
'SELECT id, name, cluster_id, deleted_by FROM custom_scripts WHERE deleted_at IS NOT NULL AND deleted_at < ?',
(cutoff,)
)
for script in deleted:
db.execute('DELETE FROM custom_scripts WHERE id = ?', (script['id'],))
log_audit('system', 'script.purged', f"Permanently deleted script '{script['name']}' after 20-day retention", cluster=script['cluster_id'])
if deleted:
logging.info(f"Purged {len(deleted)} scripts after 20-day retention period")
except Exception as e:
logging.error(f"Error cleaning up deleted scripts: {e}")
def cleanup_orphaned_excluded_vms():
"""Remove excluded VM entries for VMs that no longer exist
MK: This runs daily to clean up stale entries from balancing_excluded_vms
when VMs are deleted through other means (e.g. directly in Proxmox UI)
"""
try:
db = get_db()
cursor = db.conn.cursor()
# Get all excluded VM entries
cursor.execute('SELECT cluster_id, vmid FROM balancing_excluded_vms')
excluded_entries = cursor.fetchall()
if not excluded_entries:
return
removed_count = 0
for entry in excluded_entries:
cluster_id = entry['cluster_id']
vmid = entry['vmid']
# Check if cluster still exists and is connected
if cluster_id not in cluster_managers:
# Cluster no longer exists, remove entry
cursor.execute(
'DELETE FROM balancing_excluded_vms WHERE cluster_id = ? AND vmid = ?',
(cluster_id, vmid)
)
removed_count += 1
continue
mgr = cluster_managers[cluster_id]
if not mgr.is_connected:
continue # Skip if we can't verify
# Check if VM still exists
try:
vms = mgr.get_vm_resources()
vm_exists = any(vm.get('vmid') == vmid for vm in vms)
if not vm_exists:
cursor.execute(
'DELETE FROM balancing_excluded_vms WHERE cluster_id = ? AND vmid = ?',
(cluster_id, vmid)
)
removed_count += 1
logging.info(f"Removed orphaned excluded VM entry: cluster={cluster_id}, vmid={vmid}")
except Exception as e:
logging.debug(f"Could not verify VM {vmid} in cluster {cluster_id}: {e}")
if removed_count > 0:
db.conn.commit()
logging.info(f"[CLEANUP] Removed {removed_count} orphaned excluded VM entries")
except Exception as e:
logging.error(f"Error cleaning up orphaned excluded VMs: {e}")
@bp.route('/api/clusters/<cluster_id>/scripts', methods=['POST'])
@require_auth(perms=['admin.scripts'])
def create_custom_script(cluster_id):
"""Create a new custom script - requires admin.scripts permission"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
data = request.json or {}
if not data.get('name') or not data.get('content'):
return jsonify({'error': 'Name and content required'}), 400
script_type = data.get('type', 'bash')
if script_type not in ['bash', 'python']:
return jsonify({'error': 'Type must be bash or python'}), 400
usr = getattr(request, 'session', {}).get('user', 'system')
db = get_db()
script_id = str(uuid.uuid4())[:8]
db.execute('''
INSERT INTO custom_scripts (id, cluster_id, name, description, type, content, target_nodes, enabled, created_at, created_by)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
''', (
script_id,
cluster_id,
data.get('name'),
data.get('description', ''),
script_type,
data.get('content'),
data.get('target_nodes', 'all'),
1 if data.get('enabled', True) else 0,
datetime.now().isoformat(),
usr
))
# Get cluster name for audit log
cluster_name = cluster_managers.get(cluster_id, {})
if hasattr(cluster_name, 'config'):
cluster_name = cluster_name.config.name
else:
cluster_name = cluster_id
log_audit(usr, 'script.created', f"Created script '{data.get('name')}' (ID: {script_id}, Type: {script_type})", cluster=cluster_name)
return jsonify({'success': True, 'id': script_id})
@bp.route('/api/clusters/<cluster_id>/scripts/<script_id>', methods=['PUT'])
@require_auth(perms=['admin.scripts'])
def update_custom_script(cluster_id, script_id):
"""Update a custom script - requires admin.scripts permission"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
data = request.json or {}
db = get_db()
# Check script exists and not deleted
script = db.query_one('SELECT * FROM custom_scripts WHERE id = ? AND cluster_id = ? AND deleted_at IS NULL', (script_id, cluster_id))
if not script:
return jsonify({'error': 'Script not found'}), 404
usr = getattr(request, 'session', {}).get('user', 'system')
db.execute('''
UPDATE custom_scripts SET
name = ?,
description = ?,
type = ?,
content = ?,
target_nodes = ?,
enabled = ?,
updated_at = ?
WHERE id = ? AND cluster_id = ?
''', (
data.get('name', script['name']),
data.get('description', script['description']),
data.get('type', script['type']),
data.get('content', script['content']),
data.get('target_nodes', script['target_nodes']),
1 if data.get('enabled', script['enabled']) else 0,
datetime.now().isoformat(),
script_id,
cluster_id
))
# Get cluster name for audit log
cluster_name = cluster_managers.get(cluster_id, {})
if hasattr(cluster_name, 'config'):
cluster_name = cluster_name.config.name
else:
cluster_name = cluster_id
log_audit(usr, 'script.updated', f"Updated script '{data.get('name', script['name'])}' (ID: {script_id})", cluster=cluster_name)
return jsonify({'success': True})
@bp.route('/api/clusters/<cluster_id>/scripts/<script_id>', methods=['DELETE'])
@require_auth(perms=['admin.scripts'])
def delete_custom_script(cluster_id, script_id):
"""Soft-delete a custom script - will be permanently deleted after 20 days"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
db = get_db()
# Check script exists
script = db.query_one('SELECT * FROM custom_scripts WHERE id = ? AND cluster_id = ? AND deleted_at IS NULL', (script_id, cluster_id))
if not script:
return jsonify({'error': 'Script not found'}), 404
usr = getattr(request, 'session', {}).get('user', 'system')
# Soft delete - mark as deleted but keep for 20 days
db.execute('''
UPDATE custom_scripts SET deleted_at = ?, deleted_by = ? WHERE id = ? AND cluster_id = ?
''', (datetime.now().isoformat(), usr, script_id, cluster_id))
# Get cluster name for audit log
cluster_name = cluster_managers.get(cluster_id, {})
if hasattr(cluster_name, 'config'):
cluster_name = cluster_name.config.name
else:
cluster_name = cluster_id
log_audit(usr, 'script.deleted', f"Soft-deleted script '{script['name']}' (ID: {script_id}) - will be purged in 20 days", cluster=cluster_name)
return jsonify({'success': True, 'message': 'Script marked for deletion. Will be permanently removed in 20 days.'})
@bp.route('/api/clusters/<cluster_id>/scripts/<script_id>/run', methods=['POST'])
@require_auth(perms=['admin.scripts'])
def run_custom_script(cluster_id, script_id):
"""Run a custom script on target nodes - REQUIRES PASSWORD CONFIRMATION
This is a sensitive operation that executes arbitrary code on nodes.
Password confirmation is required to prevent accidental or unauthorized execution.
"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
# SECURITY: Require password confirmation before running any script
data = request.json or {}
password = data.get('password')
if not password:
return jsonify({'error': 'Password confirmation required to run scripts'}), 401
# Verify password against current user
usr = getattr(request, 'session', {}).get('user', 'system')
users = load_users()
user_data = users.get(usr)
if not user_data:
return jsonify({'error': 'User not found'}), 401
# Check password
stored_salt = user_data.get('password_salt', '')
stored_hash = user_data.get('password_hash', '')
if not stored_salt or not stored_hash or not verify_password(password, stored_salt, stored_hash):
cluster_name = cluster_managers[cluster_id].config.name if cluster_id in cluster_managers else cluster_id
log_audit(usr, 'script.run_denied', f"Failed password verification for script execution (ID: {script_id})", cluster=cluster_name)
return jsonify({'error': 'Invalid password'}), 401
db = get_db()
script = db.query_one('SELECT * FROM custom_scripts WHERE id = ? AND cluster_id = ? AND deleted_at IS NULL', (script_id, cluster_id))
if not script:
return jsonify({'error': 'Script not found'}), 404
if not script['enabled']:
return jsonify({'error': 'Script is disabled'}), 400
mgr = cluster_managers[cluster_id]
# Get cluster name for audit log
cluster_name = mgr.config.name if hasattr(mgr, 'config') else cluster_id
# Get target nodes (PR #324: IP resolution moved into the per-node loop
# via _get_node_ip -- cluster/status IPs can be corosync ring addresses).
target_nodes = script['target_nodes']
nodes_to_run = []
try:
cluster_host, cluster_port = mgr.host, mgr.api_port
status_url = f"https://{cluster_host}:{cluster_port}/api2/json/cluster/status"
r = mgr._create_session().get(status_url, timeout=10)
if r.status_code == 200:
for item in r.json().get('data', []):
if item.get('type') == 'node':
node_name = item.get('name')
if target_nodes == 'all' or node_name in target_nodes.split(','):
nodes_to_run.append(node_name)
except Exception as e:
logging.error(f"Error getting cluster nodes: {e}")
return jsonify({'error': f'Could not get cluster nodes: {e}'}), 500
if not nodes_to_run:
return jsonify({'error': 'No target nodes found'}), 404
# Log the execution attempt BEFORE running
log_audit(usr, 'script.execution_started', f"Starting execution of script '{script['name']}' (ID: {script_id}) on {len(nodes_to_run)} nodes: {', '.join(nodes_to_run)}", cluster=cluster_name)
script_ext = '.py' if script['type'] == 'python' else '.sh'
interpreter = 'python3' if script['type'] == 'python' else 'bash'
# NS Apr 2026 — was a serial `for node in nodes_to_run` loop. With 15-node
# clusters that meant up to 15 × 300s timeout = 75 min worst-case for one
# script run. Now bounded-parallel via run_per_node (default 8 workers).
# IMPORTANT: HA paths are completely separate (manager.py _ssh_run_command_*)
# and DO NOT go through this helper — HA latency is unaffected.
def _exec_on_node(node):
node_ip = mgr._get_node_ip(node)
if not node_ip:
return {'node': node, 'success': False, 'error': 'Could not determine SSH-reachable IP', 'output': ''}
try:
ssh = mgr._ssh_connect(node_ip)
if not ssh:
return {'node': node, 'success': False, 'error': 'SSH connection failed', 'output': ''}
script_path = f'/tmp/pegaprox_script_{script_id}{script_ext}'
sftp = ssh.open_sftp()
with sftp.file(script_path, 'w') as f:
f.write(script['content'])
sftp.chmod(script_path, 0o755)
sftp.close()
stdin, stdout, stderr = ssh.exec_command(f'{interpreter} {script_path} 2>&1', timeout=300)
output = stdout.read().decode('utf-8', errors='replace')
exit_code = stdout.channel.recv_exit_status()
ssh.exec_command(f'rm -f {script_path}')
ssh.close()
return {
'node': node,
'success': exit_code == 0,
'exit_code': exit_code,
'output': output[:10000] if output else ''
}
except Exception as e:
return {'node': node, 'success': False, 'error': str(e), 'output': ''}
from pegaprox.utils.concurrent import run_per_node
# Per-cluster cap of 8 workers — keeps total connections bounded even if
# multiple clusters run scripts at once, AND leaves headroom for HA SSH
# which goes through the separate untracked path.
raw = run_per_node(
{n: _exec_on_node for n in nodes_to_run},
max_concurrent=8,
timeout=320, # script timeout is 300, give the wrapper a small slack
)
# Preserve original input order in results list
results = []
all_output = []
for node in nodes_to_run:
r = raw.get(node) or {'node': node, 'success': False, 'error': 'Timed out or no result', 'output': ''}
results.append(r)
if r.get('output'):
all_output.append(f"=== {node} (exit: {r.get('exit_code', '?')}) ===\n{r['output']}\n")
elif r.get('error'):
all_output.append(f"=== {node} ===\n{r['error']}\n")
# Update last run info with output
success_count = sum(1 for r in results if r['success'])
status = 'success' if success_count == len(nodes_to_run) else ('partial' if success_count > 0 else 'failed')
combined_output = '\n'.join(all_output)[:50000] # Limit stored output
db.execute('''
UPDATE custom_scripts SET last_run = ?, last_status = ?, last_output = ? WHERE id = ?
''', (datetime.now().isoformat(), status, combined_output, script_id))
# Detailed audit log of execution result
log_audit(usr, 'script.executed', f"Script '{script['name']}' completed: {success_count}/{len(nodes_to_run)} nodes succeeded ({status})", cluster=cluster_name)
return jsonify({
'success': success_count == len(nodes_to_run),
'message': f'Ran on {success_count}/{len(nodes_to_run)} nodes',
'status': status,
'results': results
})
@bp.route('/api/clusters/<cluster_id>/scripts/<script_id>/output', methods=['GET'])
@require_auth(perms=['admin.scripts'])
def get_script_output(cluster_id, script_id):
"""Get the last execution output of a script"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
db = get_db()
script = db.query_one('SELECT name, last_run, last_status, last_output FROM custom_scripts WHERE id = ? AND cluster_id = ? AND deleted_at IS NULL', (script_id, cluster_id))
if not script:
return jsonify({'error': 'Script not found'}), 404
return jsonify({
'name': script['name'],
'last_run': script['last_run'],
'last_status': script['last_status'],
'output': script['last_output'] or 'No output available'
})
@bp.route('/api/clusters/<cluster_id>/scripts/deleted', methods=['GET'])
@require_auth(perms=['admin.scripts'])
def get_deleted_scripts(cluster_id):
"""Get list of soft-deleted scripts (pending permanent deletion)"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
try:
db = get_db()
scripts = db.query(
'''SELECT id, name, description, type, deleted_at, deleted_by,
datetime(deleted_at, '+20 days') as purge_date
FROM custom_scripts
WHERE cluster_id = ? AND deleted_at IS NOT NULL
ORDER BY deleted_at DESC''',
(cluster_id,)
)
return jsonify([dict(s) for s in scripts] if scripts else [])
except Exception as e:
return jsonify({'error': safe_error(e, 'Failed to load deleted scripts')}), 500
@bp.route('/api/clusters/<cluster_id>/scripts/<script_id>/restore', methods=['POST'])
@require_auth(perms=['admin.scripts'])
def restore_deleted_script(cluster_id, script_id):
"""Restore a soft-deleted script before it's permanently purged"""
ok, err = check_cluster_access(cluster_id)
if not ok: return err
db = get_db()
script = db.query_one('SELECT * FROM custom_scripts WHERE id = ? AND cluster_id = ? AND deleted_at IS NOT NULL', (script_id, cluster_id))
if not script:
return jsonify({'error': 'Deleted script not found'}), 404
usr = getattr(request, 'session', {}).get('user', 'system')
db.execute('''
UPDATE custom_scripts SET deleted_at = NULL, deleted_by = NULL WHERE id = ? AND cluster_id = ?
''', (script_id, cluster_id))
# Get cluster name for audit log
cluster_name = cluster_managers.get(cluster_id, {})
if hasattr(cluster_name, 'config'):
cluster_name = cluster_name.config.name
else:
cluster_name = cluster_id
log_audit(usr, 'script.restored', f"Restored deleted script '{script['name']}' (ID: {script_id})", cluster=cluster_name)
return jsonify({'success': True, 'message': f"Script '{script['name']}' restored"})
# ──────────────────────────────────────────
# XCP-ng specific: PIFs, bonds, guest metrics, pool HA
# ──────────────────────────────────────────
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/pifs', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_pifs_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if not hasattr(mgr, 'get_host_pifs'):
return jsonify([])
return jsonify(mgr.get_host_pifs(node))
@bp.route('/api/clusters/<cluster_id>/nodes/<node>/bonds', methods=['GET'])
@require_auth(perms=['node.view'])
def get_node_bonds_api(cluster_id, node):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if not hasattr(mgr, 'get_bonds'):
return jsonify([])
return jsonify(mgr.get_bonds(node))
@bp.route('/api/clusters/<cluster_id>/ha', methods=['GET'])
@require_auth(perms=['cluster.view'])
def get_pool_ha_api(cluster_id):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if not hasattr(mgr, 'get_ha_status'):
return jsonify({'enabled': False})
return jsonify(mgr.get_ha_status())
@bp.route('/api/clusters/<cluster_id>/ha/enable', methods=['POST'])
@require_auth(perms=['ha.config'])
def enable_pool_ha_api(cluster_id):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if getattr(mgr, 'cluster_type', 'proxmox') != 'xcpng':
return jsonify({'error': 'Only supported for XCP-ng pools'}), 400
data = request.get_json(silent=True) or {}
result = mgr.enable_pool_ha(
heartbeat_srs=data.get('heartbeat_srs'),
host_failures_to_tolerate=int(data.get('host_failures_to_tolerate', 1))
)
if result.get('success'):
return jsonify(result)
return jsonify(result), 500
@bp.route('/api/clusters/<cluster_id>/ha/disable', methods=['POST'])
@require_auth(perms=['ha.config'])
def disable_pool_ha_api(cluster_id):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if getattr(mgr, 'cluster_type', 'proxmox') != 'xcpng':
return jsonify({'error': 'Only supported for XCP-ng pools'}), 400
result = mgr.disable_pool_ha()
if result.get('success'):
return jsonify(result)
return jsonify(result), 500
@bp.route('/api/clusters/<cluster_id>/vms/<int:vmid>/ha-priority', methods=['PUT'])
@require_auth(perms=['vm.config'])
def set_vm_ha_priority_api(cluster_id, vmid):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if not hasattr(mgr, 'set_vm_ha_restart_priority'):
return jsonify({'error': 'Not supported for this cluster type'}), 400
data = request.get_json(silent=True) or {}
priority = data.get('priority', '')
if priority not in ('restart', 'best-effort', ''):
return jsonify({'error': 'Invalid priority. Use: restart, best-effort, or empty string'}), 400
result = mgr.set_vm_ha_restart_priority(vmid, priority)
if result.get('success'):
return jsonify(result)
return jsonify(result), 500
# guest metrics (works for XCP-ng, Proxmox uses qemu-guest-agent differently)
@bp.route('/api/clusters/<cluster_id>/vms/<int:vmid>/guest-metrics', methods=['GET'])
@require_auth(perms=['vm.view'])
def get_vm_guest_metrics_api(cluster_id, vmid):
ok, err = check_cluster_access(cluster_id)
if not ok: return err
if cluster_id not in cluster_managers:
return jsonify({'error': 'Cluster not found'}), 404
mgr = cluster_managers[cluster_id]
if not hasattr(mgr, 'get_guest_metrics'):
return jsonify({})
return jsonify(mgr.get_guest_metrics(None, vmid))