mirror of
https://github.com/PegaProx/project-pegaprox.git
synced 2026-08-12 15:27:47 +08:00
Cluster-wide in-band BMC health, surfaced + alertable — mirrors the proven #601 temperature pipeline (5-min collector populates a per-manager cache; the 60s alert loop only READS it, never SSHes). Backend: * manager.py: _node_hw_cache/_lock/_backoff + get_cached_node_hardware() and get_cluster_hw_rollup() -> {health, available, checked, counts, degraded[]}. * metrics.py: _node_hw_summary() SSH probe (~1h backoff for nodes without ipmitool/BMC), populated in the 5-min collector via run_per_node(cap 8, 90s), GATED on the compliance consent + proxmox-only. Off the hot-path. * alerts.py: 'hardware_health' alert metric (cluster worst + per-node), ok/warning/ critical mapped to 0/1/2, auto-severity, cache-only. * nodes.py: GET /clusters/<cid>/hardware/health rollup (consent-gated, cache-only). * vms.py: compact hardware rollup injected into datacenter/status (cache-only, consent-gated) so the overview badge is free. Frontend: 'hardware_health' alert metric with a Warning-or-worse / Critical-only threshold selector; degraded-hardware badge in the corporate sidebar, the warning banner, and the cloud overview. i18n across all 7 languages. Adversarial review (12 agents): 3 rejected (bounded/cosmetic), 5 confirmed & fixed: * MED — hwHealth prop was missing at the 2nd (ungrouped) ClusterSidebarItem call site, so the sidebar badge never showed for ungrouped clusters — now passed. * LOW — the rollup endpoint 500'd on non-proxmox clusters — proxmox/callable guard now degrades to an empty 'unknown' rollup. * LOW — a '<' operator on hardware_health builds a silent no-alert rule — the UI now only offers '>' and the create/update API pins hardware_health to '>'. * NIT — alerts now show OK/WARNING/CRITICAL instead of the raw 0/1/2 code. * NIT — corrected a stale '>=' code comment. Tests: +11 (rollup endpoint gates, non-proxmox graceful, operator coercion, manager rollup/cache unit tests). Full suite 322 passed; frontend build clean.
779 lines
30 KiB
Python
779 lines
30 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""alerts & cluster affinity rules routes - split from monolith dec 2025, NS/MK"""
|
|
|
|
import os
|
|
import json
|
|
import logging
|
|
import uuid
|
|
from datetime import datetime
|
|
from flask import Blueprint, jsonify, request
|
|
|
|
from pegaprox.constants import *
|
|
from pegaprox.globals import *
|
|
from pegaprox.models.permissions import *
|
|
from pegaprox.core.db import get_db
|
|
|
|
from pegaprox.utils.auth import require_auth
|
|
from pegaprox.utils.audit import log_audit
|
|
from pegaprox.api.helpers import check_cluster_access, safe_error
|
|
from pegaprox.background.alerts import load_alerts_config, save_alerts_config
|
|
|
|
bp = Blueprint('alerts', __name__)
|
|
|
|
# NOTE: get_cluster_report_summary is in reports.py (no duplicate here)
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/reports/top-vms', methods=['GET'])
|
|
@require_auth()
|
|
def get_cluster_top_vms(cluster_id):
|
|
"""Get top VMs by resource usage for a specific 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]
|
|
metric = request.args.get('metric', 'cpu')
|
|
limit = int(request.args.get('limit', 10))
|
|
|
|
if not mgr.is_connected:
|
|
return jsonify([])
|
|
|
|
vms = []
|
|
try:
|
|
resources = mgr.get_vm_resources()
|
|
for r in resources:
|
|
if r.get('status') != 'running':
|
|
continue
|
|
|
|
vm_data = {
|
|
'vmid': r.get('vmid'),
|
|
'name': r.get('name'),
|
|
'node': r.get('node'),
|
|
'type': r.get('type'),
|
|
'cpu': r.get('cpu', 0),
|
|
'mem': r.get('mem', 0),
|
|
'maxmem': r.get('maxmem', 0),
|
|
'mem_percent': round(r.get('mem', 0) / max(r.get('maxmem', 1), 1) * 100, 1)
|
|
}
|
|
vms.append(vm_data)
|
|
except:
|
|
pass
|
|
|
|
# Sort by metric
|
|
if metric == 'memory':
|
|
vms.sort(key=lambda x: x.get('mem_percent', 0), reverse=True)
|
|
else:
|
|
vms.sort(key=lambda x: x.get('cpu', 0), reverse=True)
|
|
|
|
return jsonify(vms[:limit])
|
|
|
|
|
|
# ============================================
|
|
# Cluster-Based Alerts Endpoints
|
|
# moved to per-cluster
|
|
# ============================================
|
|
|
|
def load_cluster_alerts():
|
|
"""Load alerts config from SQLite database
|
|
|
|
NS: Migrated from JSON to SQLite Jan 2026
|
|
MK: keeps falling back to json if db fails which is kinda nice for debugging
|
|
|
|
Returns: {cluster_id: [list of alert objects]}
|
|
"""
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('SELECT * FROM cluster_alerts WHERE enabled = 1')
|
|
|
|
alerts = {}
|
|
for row in cursor.fetchall():
|
|
cluster_id = row['cluster_id']
|
|
if cluster_id not in alerts:
|
|
alerts[cluster_id] = []
|
|
|
|
try:
|
|
# config contains the full alert object as JSON
|
|
alert_data = json.loads(row['config'] or '{}')
|
|
# ensure id is present
|
|
if 'id' not in alert_data:
|
|
alert_data['id'] = row['alert_type']
|
|
alerts[cluster_id].append(alert_data)
|
|
except:
|
|
# fallback for old format where config was just settings
|
|
alerts[cluster_id].append({
|
|
'id': row['alert_type'],
|
|
'name': row['alert_type'],
|
|
'config': row['config']
|
|
})
|
|
|
|
return alerts
|
|
except Exception as e:
|
|
logging.error(f"Error loading cluster alerts from DB: {e}")
|
|
# Fallback to JSON for backwards compat
|
|
try:
|
|
alerts_file = os.path.join(CONFIG_DIR, 'cluster_alerts.json')
|
|
if os.path.exists(alerts_file):
|
|
with open(alerts_file, 'r') as f:
|
|
return json.load(f)
|
|
except:
|
|
pass
|
|
return {}
|
|
|
|
def save_cluster_alerts(alerts):
|
|
"""Save alerts config to SQLite database
|
|
|
|
NS: stores each alert as a row with config containing full alert object
|
|
|
|
Expects: {cluster_id: [list of alert objects]}
|
|
"""
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
now = datetime.now().isoformat()
|
|
|
|
for cluster_id, alert_list in alerts.items():
|
|
# Handle list format (from API)
|
|
if isinstance(alert_list, list):
|
|
for alert in alert_list:
|
|
alert_id = alert.get('id', str(uuid.uuid4())[:8])
|
|
cursor.execute('''
|
|
INSERT OR REPLACE INTO cluster_alerts
|
|
(cluster_id, alert_type, config, enabled, updated_at)
|
|
VALUES (?, ?, ?, ?, ?)
|
|
''', (
|
|
cluster_id,
|
|
alert_id,
|
|
json.dumps(alert),
|
|
1 if alert.get('enabled', True) else 0,
|
|
now
|
|
))
|
|
# Handle dict format (legacy)
|
|
elif isinstance(alert_list, dict):
|
|
for alert_type, config in alert_list.items():
|
|
cursor.execute('''
|
|
INSERT OR REPLACE INTO cluster_alerts
|
|
(cluster_id, alert_type, config, enabled, updated_at)
|
|
VALUES (?, ?, ?, 1, ?)
|
|
''', (
|
|
cluster_id,
|
|
alert_type,
|
|
json.dumps(config) if isinstance(config, dict) else str(config),
|
|
now
|
|
))
|
|
|
|
db.conn.commit()
|
|
except Exception as e:
|
|
logging.error(f"Error saving cluster alerts to DB: {e}")
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/alerts', methods=['GET'])
|
|
@require_auth()
|
|
def get_cluster_alerts(cluster_id):
|
|
"""Get alerts for a specific cluster"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
try:
|
|
alerts = load_cluster_alerts()
|
|
cluster_alerts = alerts.get(cluster_id, [])
|
|
return jsonify({'alerts': cluster_alerts})
|
|
except Exception as e:
|
|
logging.error(f"Error getting cluster alerts: {e}")
|
|
return jsonify({'alerts': [], 'error': safe_error(e, 'Alert operation failed')})
|
|
|
|
# NS #501 hardening (F1): bound the escalation chain + channel fan-out so a runaway
|
|
# rule can't flood notification channels or stall the single-threaded alert loop.
|
|
MAX_ESCALATION_STEPS = 10
|
|
MAX_STEP_CHANNELS = 20
|
|
MAX_AFTER_MINUTES = 10080 # cap a step delay at one week
|
|
|
|
|
|
def _sanitize_channels(raw):
|
|
if not isinstance(raw, list):
|
|
return []
|
|
return [str(c) for c in raw if isinstance(c, (str, int))][:MAX_STEP_CHANNELS]
|
|
|
|
|
|
def _sanitize_escalation(raw):
|
|
"""Normalise an escalation chain to a bounded list of {after_minutes, channels}."""
|
|
if not isinstance(raw, list):
|
|
return []
|
|
out = []
|
|
for s in raw:
|
|
if not isinstance(s, dict):
|
|
continue
|
|
try:
|
|
after = int(s.get('after_minutes', 0) or 0)
|
|
except (TypeError, ValueError):
|
|
after = 0
|
|
out.append({
|
|
'after_minutes': max(0, min(after, MAX_AFTER_MINUTES)),
|
|
'channels': _sanitize_channels(s.get('channels')),
|
|
})
|
|
if len(out) >= MAX_ESCALATION_STEPS:
|
|
break
|
|
return out
|
|
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/alerts', methods=['POST'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def create_cluster_alert(cluster_id):
|
|
"""Create a new alert for a cluster"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
data = request.get_json()
|
|
if not data:
|
|
return jsonify({'error': 'No data provided'}), 400
|
|
|
|
alerts = load_cluster_alerts()
|
|
if cluster_id not in alerts:
|
|
alerts[cluster_id] = []
|
|
|
|
# NS May 2026 — used to silently drop `channels` and `cluster_id`, so the
|
|
# background loop never knew where to dispatch and which cluster the alert
|
|
# belonged to. Persist both in the JSON config.
|
|
channels = _sanitize_channels(data.get('channels'))
|
|
alert = {
|
|
'id': str(uuid.uuid4())[:8],
|
|
'name': data.get('name', 'Unnamed Alert'),
|
|
'cluster_id': cluster_id,
|
|
'metric': data.get('metric', 'cpu'),
|
|
# #609: the categorical hardware_health code (0/1/2) is only meaningful with '>';
|
|
# a '<' rule would silently never fire on degraded hardware — pin it to '>'.
|
|
'operator': '>' if data.get('metric') == 'hardware_health' else data.get('operator', '>'),
|
|
'threshold': data.get('threshold', 80),
|
|
'target_type': data.get('target_type', 'cluster'),
|
|
'target_id': data.get('target_id'),
|
|
'channels': channels,
|
|
'severity': data.get('severity', 'auto'), # NS #501: 'auto' | critical | warning | info
|
|
'escalation': _sanitize_escalation(data.get('escalation')), # NS #501 (F1: bounded chain)
|
|
'action': data.get('action', 'log'), # legacy fallback
|
|
'enabled': data.get('enabled', True),
|
|
'created_at': datetime.now().isoformat()
|
|
}
|
|
|
|
alerts[cluster_id].append(alert)
|
|
save_cluster_alerts(alerts)
|
|
|
|
return jsonify({'success': True, 'alert': alert})
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/alerts/<alert_id>', methods=['PUT'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def update_cluster_alert(cluster_id, alert_id):
|
|
"""Update an alert"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
data = request.get_json()
|
|
alerts = load_cluster_alerts()
|
|
cluster_alerts = alerts.get(cluster_id, [])
|
|
|
|
for alert in cluster_alerts:
|
|
if alert['id'] == alert_id:
|
|
for k in ('enabled', 'name', 'threshold', 'metric', 'operator',
|
|
'target_type', 'target_id', 'action', 'severity'):
|
|
if k in data:
|
|
alert[k] = data[k]
|
|
if 'channels' in data:
|
|
alert['channels'] = _sanitize_channels(data['channels'])
|
|
if 'escalation' in data:
|
|
alert['escalation'] = _sanitize_escalation(data['escalation']) # NS #501 (F1: bounded)
|
|
# #609: keep hardware_health rules on '>' (the categorical 0/1/2 ladder)
|
|
if alert.get('metric') == 'hardware_health':
|
|
alert['operator'] = '>'
|
|
# ensure cluster_id is always present for older rows
|
|
alert.setdefault('cluster_id', cluster_id)
|
|
save_cluster_alerts(alerts)
|
|
return jsonify({'success': True, 'alert': alert})
|
|
|
|
return jsonify({'error': 'Alert not found'}), 404
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/alerts/<alert_id>', methods=['DELETE'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def delete_cluster_alert(cluster_id, alert_id):
|
|
"""Delete an alert"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
# NS: delete directly from DB for efficiency
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('DELETE FROM cluster_alerts WHERE cluster_id = ? AND alert_type = ?',
|
|
(cluster_id, alert_id))
|
|
db.conn.commit()
|
|
deleted = cursor.rowcount > 0
|
|
return jsonify({'success': True, 'deleted': deleted})
|
|
except Exception as e:
|
|
logging.error(f"Error deleting cluster alert: {e}")
|
|
return jsonify({'error': safe_error(e, 'Alert operation failed')}), 500
|
|
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/active-alerts', methods=['GET'])
|
|
@require_auth()
|
|
def get_active_alerts(cluster_id):
|
|
"""List currently-firing (unresolved) alert incidents for a cluster (#501)."""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
try:
|
|
db = get_db()
|
|
cur = db.conn.cursor()
|
|
cols = ['id', 'alert_id', 'severity', 'message', 'target_type', 'target_name', 'metric',
|
|
'current_value', 'threshold', 'operator', 'triggered_at', 'last_fired_at',
|
|
'acked_at', 'acked_by', 'escalation_step']
|
|
rows = cur.execute(
|
|
f"SELECT {', '.join(cols)} FROM active_alerts "
|
|
"WHERE cluster_id = ? AND resolved_at IS NULL ORDER BY triggered_at DESC",
|
|
(cluster_id,)).fetchall()
|
|
return jsonify({'active_alerts': [dict(zip(cols, r)) for r in rows]})
|
|
except Exception as e:
|
|
logging.error(f"Error listing active alerts: {e}")
|
|
return jsonify({'active_alerts': [], 'error': safe_error(e, 'Alert operation failed')})
|
|
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/active-alerts/<fired_id>/ack', methods=['POST'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def ack_active_alert(cluster_id, fired_id):
|
|
"""Acknowledge a firing alert — stops escalation for this incident (#501)."""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
user = getattr(request, 'username', None) or (request.session.get('user', 'unknown') if hasattr(request, 'session') else 'unknown')
|
|
try:
|
|
db = get_db()
|
|
cur = db.conn.cursor()
|
|
cur.execute(
|
|
"UPDATE active_alerts SET acked_at=?, acked_by=? WHERE id=? AND cluster_id=? AND resolved_at IS NULL",
|
|
(datetime.now().isoformat(), user, fired_id, cluster_id))
|
|
db.conn.commit()
|
|
if cur.rowcount > 0:
|
|
return jsonify({'success': True, 'acked_by': user})
|
|
return jsonify({'error': 'Active alert not found or already resolved'}), 404
|
|
except Exception as e:
|
|
logging.error(f"Error acknowledging alert: {e}")
|
|
return jsonify({'error': safe_error(e, 'Alert operation failed')}), 500
|
|
|
|
|
|
# ============================================
|
|
# Cluster-Based Affinity Rules Endpoints
|
|
# moved to per-cluster
|
|
# ============================================
|
|
|
|
def load_cluster_affinity_rules():
|
|
"""Load affinity rules from SQLite database
|
|
|
|
MK: affinity = keep VMs together, anti-affinity = keep them apart
|
|
useful for HA where you want replicas on different hosts
|
|
NS: reuses the affinity_rules table we already had
|
|
"""
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('SELECT * FROM affinity_rules WHERE enabled = 1')
|
|
|
|
rules = {}
|
|
for row in cursor.fetchall():
|
|
cluster_id = row['cluster_id']
|
|
if cluster_id not in rules:
|
|
rules[cluster_id] = []
|
|
|
|
vms_list = json.loads(row['vms'] or '[]')
|
|
rules[cluster_id].append({
|
|
'id': row['id'],
|
|
'name': row['name'],
|
|
'type': row['type'],
|
|
'vms': vms_list,
|
|
'vm_ids': vms_list, # NS: alias for frontend compatibility
|
|
'enabled': bool(row['enabled']),
|
|
'enforce': bool(row['enforce']) if 'enforce' in row.keys() else False,
|
|
'created_at': row['created_at']
|
|
})
|
|
|
|
return rules
|
|
except Exception as e:
|
|
logging.error(f"Error loading affinity rules from DB: {e}")
|
|
# Fallback to JSON for backwards compat
|
|
try:
|
|
rules_file = os.path.join(CONFIG_DIR, 'cluster_affinity_rules.json')
|
|
if os.path.exists(rules_file):
|
|
with open(rules_file, 'r') as f:
|
|
return json.load(f)
|
|
except:
|
|
pass
|
|
return {}
|
|
|
|
def save_cluster_affinity_rules(rules):
|
|
"""Save affinity rules to SQLite database
|
|
|
|
NS: uses upsert pattern, handles both 'vms' and 'vm_ids' field names
|
|
"""
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
now = datetime.now().isoformat()
|
|
|
|
for cluster_id, cluster_rules in rules.items():
|
|
for rule in cluster_rules:
|
|
# generate id if missing (old rules might not have one)
|
|
rule_id = rule.get('id', str(uuid.uuid4()))
|
|
# NS: handle both 'vms' and 'vm_ids' field names
|
|
vms_data = rule.get('vms') or rule.get('vm_ids') or []
|
|
cursor.execute('''
|
|
INSERT OR REPLACE INTO affinity_rules
|
|
(id, cluster_id, name, type, vms, enabled, enforce, created_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
|
''', (
|
|
rule_id,
|
|
cluster_id,
|
|
rule.get('name', ''),
|
|
rule.get('type', 'affinity'),
|
|
json.dumps(vms_data),
|
|
1 if rule.get('enabled', True) else 0,
|
|
1 if rule.get('enforce', False) else 0,
|
|
rule.get('created_at', now)
|
|
))
|
|
|
|
db.conn.commit()
|
|
except Exception as e:
|
|
logging.error(f"Error saving affinity rules to DB: {e}")
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/affinity-rules', methods=['GET'])
|
|
@require_auth()
|
|
def get_cluster_affinity_rules(cluster_id):
|
|
"""Get affinity rules for a specific cluster"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
try:
|
|
rules = load_cluster_affinity_rules()
|
|
cluster_rules = rules.get(cluster_id, [])
|
|
return jsonify({'rules': cluster_rules})
|
|
except Exception as e:
|
|
logging.error(f"Error getting affinity rules: {e}")
|
|
return jsonify({'rules': [], 'error': safe_error(e, 'Alert operation failed')})
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/affinity-rules', methods=['POST'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def create_cluster_affinity_rule(cluster_id):
|
|
"""Create a new affinity rule for a cluster"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
data = request.get_json()
|
|
if not data:
|
|
return jsonify({'error': 'No data provided'}), 400
|
|
|
|
rules = load_cluster_affinity_rules()
|
|
if cluster_id not in rules:
|
|
rules[cluster_id] = []
|
|
|
|
# NS: get vms from either 'vm_ids' or 'vms' field
|
|
vms_data = data.get('vm_ids') or data.get('vms') or []
|
|
|
|
rule = {
|
|
'id': str(uuid.uuid4())[:8],
|
|
'name': data.get('name', f"Rule {len(rules[cluster_id]) + 1}"),
|
|
'type': data.get('type', 'together'), # 'together' or 'separate'
|
|
'vms': vms_data,
|
|
'vm_ids': vms_data, # alias for frontend
|
|
'enforce': data.get('enforce', False),
|
|
'enabled': True,
|
|
'created_at': datetime.now().isoformat()
|
|
}
|
|
|
|
rules[cluster_id].append(rule)
|
|
save_cluster_affinity_rules(rules)
|
|
|
|
return jsonify({'success': True, 'rule': rule})
|
|
|
|
@bp.route('/api/clusters/<cluster_id>/affinity-rules/<rule_id>', methods=['DELETE'])
|
|
@require_auth(perms=['cluster.config'])
|
|
def delete_cluster_affinity_rule(cluster_id, rule_id):
|
|
"""Delete an affinity rule"""
|
|
ok, err = check_cluster_access(cluster_id)
|
|
if not ok:
|
|
return err
|
|
|
|
# NS: Delete directly from DB instead of load/filter/save
|
|
try:
|
|
db = get_db()
|
|
cursor = db.conn.cursor()
|
|
cursor.execute('DELETE FROM affinity_rules WHERE id = ? AND cluster_id = ?',
|
|
(rule_id, cluster_id))
|
|
db.conn.commit()
|
|
deleted = cursor.rowcount > 0
|
|
return jsonify({'success': True, 'deleted': deleted})
|
|
except Exception as e:
|
|
logging.error(f"Error deleting affinity rule: {e}")
|
|
return jsonify({'error': safe_error(e, 'Alert operation failed')}), 500
|
|
|
|
|
|
# ============================================
|
|
|
|
@bp.route('/api/alerts', methods=['GET'])
|
|
@require_auth(perms=['cluster.view'])
|
|
def get_alerts():
|
|
"""Get all alert configurations"""
|
|
cfg = load_alerts_config()
|
|
# NS Jul 2026 (CodeAnt IDOR) — scope alert configs to the caller's reachable clusters
|
|
# (was exposing every tenant's alert rules — names, cluster/VM targets — to any viewer).
|
|
from pegaprox.utils.rbac import get_user_clusters
|
|
from flask import g as _g
|
|
_allowed = get_user_clusters(getattr(_g, 'current_user', None) or {})
|
|
if _allowed is not None:
|
|
cfg = dict(cfg)
|
|
cfg['alerts'] = [a for a in cfg.get('alerts', []) if a.get('cluster_id') in _allowed]
|
|
return jsonify(cfg)
|
|
|
|
@bp.route('/api/alerts', methods=['POST'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def create_alert():
|
|
"""Create a new alert"""
|
|
data = request.json or {}
|
|
config = load_alerts_config()
|
|
|
|
import uuid
|
|
new_alert = {
|
|
'id': str(uuid.uuid4())[:8],
|
|
'name': data.get('name', 'New Alert'),
|
|
'cluster_id': data.get('cluster_id', ''),
|
|
'target_type': data.get('target_type', 'cluster'),
|
|
'target_id': data.get('target_id', ''),
|
|
'metric': data.get('metric', 'cpu'),
|
|
'operator': data.get('operator', '>'),
|
|
'threshold': data.get('threshold', 80),
|
|
'enabled': data.get('enabled', True),
|
|
'created': datetime.now().isoformat()
|
|
}
|
|
|
|
config['alerts'].append(new_alert)
|
|
save_alerts_config(config)
|
|
|
|
user = request.session.get('user', 'unknown')
|
|
log_audit(user, 'alert.created', f"Created alert: {new_alert['name']}")
|
|
|
|
return jsonify(new_alert), 201
|
|
|
|
@bp.route('/api/alerts/<alert_id>', methods=['PUT'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def update_alert(alert_id):
|
|
"""Update an alert"""
|
|
data = request.json or {}
|
|
config = load_alerts_config()
|
|
|
|
for alert in config['alerts']:
|
|
if alert['id'] == alert_id:
|
|
alert.update({
|
|
'name': data.get('name', alert['name']),
|
|
'cluster_id': data.get('cluster_id', alert['cluster_id']),
|
|
'target_type': data.get('target_type', alert['target_type']),
|
|
'target_id': data.get('target_id', alert['target_id']),
|
|
'metric': data.get('metric', alert['metric']),
|
|
'operator': data.get('operator', alert['operator']),
|
|
'threshold': data.get('threshold', alert['threshold']),
|
|
'enabled': data.get('enabled', alert['enabled']),
|
|
})
|
|
save_alerts_config(config)
|
|
return jsonify(alert)
|
|
|
|
return jsonify({'error': 'Alert not found'}), 404
|
|
|
|
@bp.route('/api/alerts/<alert_id>', methods=['DELETE'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def delete_alert(alert_id):
|
|
"""Delete an alert"""
|
|
config = load_alerts_config()
|
|
config['alerts'] = [a for a in config['alerts'] if a['id'] != alert_id]
|
|
save_alerts_config(config)
|
|
|
|
user = request.session.get('user', 'unknown')
|
|
log_audit(user, 'alert.deleted', f"Deleted alert: {alert_id}")
|
|
|
|
return jsonify({'success': True})
|
|
|
|
|
|
# ─────────────────────────────────────────────────────────
|
|
# MK Apr 2026 — webhook alert channels (Slack, Discord, Teams, ntfy, generic)
|
|
# ─────────────────────────────────────────────────────────
|
|
|
|
@bp.route('/api/alert-channels', methods=['GET'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def list_alert_channels():
|
|
from pegaprox.api.helpers import load_server_settings
|
|
channels = (load_server_settings() or {}).get('alert_webhooks') or []
|
|
# MK: scrub url secrets on read — ?full=1 bypasses for edit flows
|
|
if request.args.get('full', '').lower() in ('1', 'true', 'yes'):
|
|
return jsonify(channels)
|
|
masked = []
|
|
for ch in channels:
|
|
c = dict(ch)
|
|
u = c.get('url') or ''
|
|
if len(u) > 32:
|
|
c['url'] = u[:24] + '…' + u[-6:]
|
|
if c.get('token'):
|
|
c['token'] = '********'
|
|
masked.append(c)
|
|
return jsonify(masked)
|
|
|
|
|
|
@bp.route('/api/alert-channels', methods=['POST'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def create_alert_channel():
|
|
from pegaprox.api.helpers import load_server_settings, save_server_settings
|
|
from pegaprox.utils.webhooks import new_channel
|
|
settings = load_server_settings()
|
|
channels = list(settings.get('alert_webhooks') or [])
|
|
ch = new_channel(request.get_json() or {})
|
|
if not ch.get('url'):
|
|
return jsonify({'error': 'url required'}), 400
|
|
channels.append(ch)
|
|
settings['alert_webhooks'] = channels
|
|
save_server_settings(settings)
|
|
log_audit(request.session.get('user', 'admin'), 'alerts.channel_create', f"added webhook '{ch.get('name')}' ({ch.get('type')})")
|
|
return jsonify({'success': True, 'channel': ch})
|
|
|
|
|
|
@bp.route('/api/alert-channels/<cid>', methods=['PUT'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def update_alert_channel(cid):
|
|
from pegaprox.api.helpers import load_server_settings, save_server_settings
|
|
settings = load_server_settings()
|
|
channels = list(settings.get('alert_webhooks') or [])
|
|
data = request.get_json() or {}
|
|
for i, ch in enumerate(channels):
|
|
if ch.get('id') != cid:
|
|
continue
|
|
# Merge the allowed fields. Skip url/token if caller sent the masked placeholder
|
|
# (admin UI shows dots; don't wipe secret because of a round-trip).
|
|
updated = dict(ch)
|
|
for k in ('name', 'type', 'enabled', 'topic', 'url', 'token'):
|
|
if k in data:
|
|
v = data[k]
|
|
if k in ('url', 'token') and isinstance(v, str) and ('…' in v or v == '********'):
|
|
continue # untouched
|
|
updated[k] = v
|
|
channels[i] = updated
|
|
settings['alert_webhooks'] = channels
|
|
save_server_settings(settings)
|
|
log_audit(request.session.get('user', 'admin'), 'alerts.channel_update', f"updated webhook '{updated.get('name')}'")
|
|
return jsonify({'success': True, 'channel': updated})
|
|
return jsonify({'error': 'channel not found'}), 404
|
|
|
|
|
|
@bp.route('/api/alert-channels/<cid>', methods=['DELETE'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def delete_alert_channel(cid):
|
|
from pegaprox.api.helpers import load_server_settings, save_server_settings
|
|
settings = load_server_settings()
|
|
before = settings.get('alert_webhooks') or []
|
|
after = [c for c in before if c.get('id') != cid]
|
|
if len(after) == len(before):
|
|
return jsonify({'error': 'channel not found'}), 404
|
|
settings['alert_webhooks'] = after
|
|
save_server_settings(settings)
|
|
log_audit(request.session.get('user', 'admin'), 'alerts.channel_delete', f"removed webhook {cid}")
|
|
return jsonify({'success': True})
|
|
|
|
|
|
@bp.route('/api/alert-channels/<cid>/test', methods=['POST'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def test_alert_channel(cid):
|
|
from pegaprox.api.helpers import load_server_settings
|
|
from pegaprox.utils.webhooks import send_to_channel
|
|
channels = (load_server_settings() or {}).get('alert_webhooks') or []
|
|
ch = next((c for c in channels if c.get('id') == cid), None)
|
|
if not ch:
|
|
return jsonify({'error': 'channel not found'}), 404
|
|
alert = {
|
|
'alert_name': 'PegaProx test alert',
|
|
'metric': 'test',
|
|
'current_value': 42,
|
|
'threshold': 0,
|
|
'operator': '>',
|
|
'target_type': 'server',
|
|
'target_name': 'pegaprox',
|
|
'cluster_id': 'test-cluster',
|
|
'severity': 'info',
|
|
'message': 'This is a test alert triggered from PegaProx settings.',
|
|
'timestamp': datetime.now().isoformat(),
|
|
}
|
|
ok, _detail = send_to_channel(ch, alert)
|
|
# M-7: don't echo the precise upstream HTTP status / connect-vs-refuse — that
|
|
# turned this admin endpoint into an SSRF port-scan oracle. Coarse pass/fail
|
|
# still lets an admin validate a real Slack/Discord/ntfy URL.
|
|
detail = 'Delivered' if ok else 'Delivery failed — check the channel URL and that the endpoint is reachable'
|
|
return jsonify({'success': ok, 'detail': detail})
|
|
|
|
|
|
# NS May 2026 — when customers report "alerts don't fire", the previous
|
|
# diagnostic story was: nothing. The check loop ran silently every 60s and
|
|
# you couldn't tell *why* a rule didn't trigger. These two endpoints give
|
|
# admins a way to see what the loop saw and force a re-check on demand.
|
|
|
|
@bp.route('/api/alerts/diagnostics', methods=['GET'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def alerts_diagnostics():
|
|
from pegaprox.api.helpers import load_server_settings
|
|
from pegaprox.background import alerts as A
|
|
settings = load_server_settings() or {}
|
|
cfg = A.load_alerts_config()
|
|
return jsonify({
|
|
'last_tick_at': A._last_tick_at,
|
|
'tick_interval_seconds': 60,
|
|
'alerts_in_config': len(cfg.get('alerts', [])),
|
|
'enabled': cfg.get('enabled', True),
|
|
'cooldown_seconds': settings.get('alert_cooldown', 300),
|
|
'email_recipients': len(settings.get('alert_email_recipients') or []),
|
|
'webhook_channels': [
|
|
{'id': c.get('id'), 'name': c.get('name'), 'type': c.get('type'),
|
|
'enabled': c.get('enabled', True)}
|
|
for c in (settings.get('alert_webhooks') or [])
|
|
],
|
|
'clusters_loaded': sorted([
|
|
{'id': cid, 'connected': bool(getattr(m, 'is_connected', False))}
|
|
for cid, m in cluster_managers.items()
|
|
], key=lambda r: r['id']),
|
|
'alerts': [
|
|
{'id': a.get('id'), 'name': a.get('name'),
|
|
'cluster_id': a.get('cluster_id'),
|
|
'metric': a.get('metric'),
|
|
'operator': a.get('operator'),
|
|
'threshold': a.get('threshold'),
|
|
'target_type': a.get('target_type'),
|
|
'target_id': a.get('target_id'),
|
|
'channels': a.get('channels'),
|
|
'enabled': a.get('enabled', True),
|
|
'last_evaluation': A._last_eval.get(a.get('id'))}
|
|
for a in cfg.get('alerts', [])
|
|
],
|
|
})
|
|
|
|
|
|
@bp.route('/api/alerts/force-check', methods=['POST'])
|
|
@require_auth(perms=['alert.manage'])
|
|
def alerts_force_check():
|
|
"""Run check_and_send_alerts() once, optionally clearing the cooldown
|
|
map so an alert that already fired in this process can re-fire."""
|
|
from pegaprox.background import alerts as A
|
|
if (request.args.get('reset_cooldown', '').lower() in ('1', 'true', 'yes')):
|
|
A._alert_last_sent.clear()
|
|
try:
|
|
A.check_and_send_alerts()
|
|
return jsonify({'ok': True, 'evaluations': A._last_eval})
|
|
except Exception as e:
|
|
return jsonify({'ok': False, 'error': safe_error(e)}), 500
|
|
|
|
|
|
# =====================================================
|
|
|