fix(monitor): distinguish notification recovery and backup outcomes

This commit is contained in:
martino
2026-10-01 08:38:32 +02:00
parent d137bb3d31
commit dcf00cc6b9
13 changed files with 616 additions and 77 deletions
+24 -3
View File
@@ -1037,6 +1037,21 @@ class EmailChannel(NotificationChannel):
# Determine group for section header
event_type = data.get('_event_type', '')
if event_type == 'error_resolved':
sev.update(self._SEV_DEFAULT)
sev['label'] = _runtime_text('email.severity.observation', data)
elif event_type == 'backup_complete':
outcome = data.get('backup_outcome')
if outcome == 'confirmed':
sev.update(self._SEV_STYLE['OK'])
status = 'completed'
elif outcome == 'failed':
sev.update(self._SEV_STYLE['CRITICAL'])
status = 'failed'
else:
sev.update(self._SEV_DEFAULT)
status = 'unconfirmed'
sev['label'] = _runtime_text(f'email.status.{status}', data)
group = data.get('_group', 'other')
# Keep unbroken recorded text inside the temperature email's table.
# Both properties are inline for mail clients; other events retain
@@ -1221,7 +1236,9 @@ class EmailChannel(NotificationChannel):
v = str(value).strip() if value else ''
if not v or v == '0' and original_label not in ('Failures',):
return
if fmt == 'severity':
if fmt == 'backup_error':
rows.append((esc(label), f'<span style="color:#dc2626;font-weight:600;">{esc(v)}</span>'))
elif fmt == 'severity':
sev_colors = {
'CRITICAL': '#dc2626', 'WARNING': '#d97706',
'INFO': '#2563eb', 'OK': '#16a34a',
@@ -1254,9 +1271,13 @@ class EmailChannel(NotificationChannel):
# tell which target the backup ran against. Reported gap: emails
# showed no way to distinguish which PBS failed with 2+ configured.
_add('Storage', data.get('storage') or data.get('storage_name'), 'code')
status_key = 'failed' if 'fail' in event_type else 'completed' if 'complete' in event_type else 'started'
if event_type == 'backup_complete' and data.get('backup_outcome') != 'confirmed':
status_key = ('failed' if data.get('backup_outcome') == 'failed'
else 'unconfirmed')
else:
status_key = 'failed' if 'fail' in event_type else 'completed' if 'complete' in event_type else 'started'
_add('Status', _runtime_text(f'email.status.{status_key}', language_data),
'severity' if 'fail' in event_type else '')
'backup_error' if status_key == 'failed' else '')
_add('Size', data.get('size'))
_add('Duration', data.get('duration'))
_add('Snapshot', data.get('snapshot_name'), 'code')
+56 -4
View File
@@ -3254,7 +3254,7 @@ class PollingCollector:
reason_lines = (reason or '').split('\n')
reason_summary = reason_lines[0] if reason_lines else ''
# Try to extract device info for a clean "Device: xxx (recovered)" line
# Keep the earlier device context without asserting recovery.
device_line = ''
for line in reason_lines:
if 'Device:' in line or 'Device not currently' in line or '/dev/' in line:
@@ -3267,11 +3267,11 @@ class PollingCollector:
break
if reason_summary and device_line:
clean_reason = f'{reason_summary}\n{device_line} (recovered)'
clean_reason = f'{reason_summary}\n{device_line} (no longer reported)'
elif reason_summary:
clean_reason = f'{reason_summary} (recovered)'
clean_reason = f'{reason_summary} (no longer reported)'
else:
clean_reason = 'Condition resolved'
clean_reason = 'Condition no longer reported'
# `original_severity` must match what the user actually saw
# in the most-recent notification for this error, not the
@@ -4289,6 +4289,51 @@ class ProxmoxHookWatcher:
def _hostname(self) -> str:
return _hostname()
@staticmethod
def _backup_outcome(severity: str, message: str) -> str:
"""Distinguish explicit failure, complete guest logs and unknown results."""
text = str(message or '')
if severity in ('error', 'err', 'critical') or re.search(
r'(?im)^\s*(?:ERROR:|TASK ERROR:|.*\bStatus\s+ERROR\b)', text):
return 'failed'
if severity not in ('info', 'ok', 'success') or re.search(
r'(?im)(?:^\s*WARNING:|\bWARNINGS\s*:\s*\d+)', text):
return 'unconfirmed'
starts = re.findall(r'(?im)\bStarting Backup of VM (\d+)\s*\(', text)
finished = re.findall(r'(?im)\bFinished Backup of VM (\d+)\s*\(', text)
lines = text.splitlines()
table_outcome = None
for index, header in enumerate(lines):
if not re.match(r'\s*VMID\s+Name\s+Status\b', header, re.IGNORECASE):
continue
status_start = header.find('Status')
status_end = header.find('Time', status_start)
if status_start < 0 or status_end < 0:
break
rows = []
for line in lines[index + 1:]:
if re.match(r'\s*Total\b', line, re.IGNORECASE):
table_outcome = ('confirmed' if rows and all(status == 'OK' for status in rows)
else 'unconfirmed')
break
if not line.strip():
break
if not re.match(r'\s*\d+\s+', line):
break
status = line[status_start:status_end].strip().upper()
if status == 'ERROR':
return 'failed'
rows.append(status)
break
if table_outcome == 'unconfirmed':
return 'unconfirmed'
if starts:
return 'confirmed' if sorted(starts) == sorted(finished) else 'unconfirmed'
if table_outcome == 'confirmed' or re.search(
r'(?im)^\s*(?:INFO:\s*)?TASK OK\s*$', text):
return 'confirmed'
return 'unconfirmed'
def process_webhook(self, payload: dict) -> dict:
"""Process an incoming Proxmox webhook payload.
@@ -4347,6 +4392,13 @@ class ProxmoxHookWatcher:
'title': title or event_type,
'job_id': pve_job_id,
}
if event_type in ('backup_complete', 'backup_fail'):
# This is presentation metadata, not a new event/toggle/delivery path.
data['backup_outcome'] = (
'failed' if event_type == 'backup_fail' else
self._backup_outcome(severity_raw, message) if pve_type == 'vzdump'
else 'unconfirmed'
)
if pve_type == 'replication':
replication = self._extract_replication_context(
+42 -23
View File
@@ -299,7 +299,7 @@ def _parse_vzdump_message(message: str) -> Optional[Dict[str, Any]]:
current_vm = {
'vmid': m_start.group(1),
'name': '',
'status': 'ok',
'status': 'unknown',
'time': '',
'size': '',
'filename': '',
@@ -338,15 +338,16 @@ def _parse_vzdump_message(message: str) -> Optional[Dict[str, Any]]:
# Finished -> duration
m_finish = re.match(
r'Finished Backup of VM (\d+)\s+\(([^)]+)\)', clean)
if m_finish:
if m_finish and m_finish.group(1) == current_vm['vmid']:
current_vm['time'] = m_finish.group(2)
current_vm['status'] = 'ok'
if current_vm['status'] != 'error':
current_vm['status'] = 'ok'
vms.append(current_vm)
current_vm = None
continue
# Error
if clean.startswith('ERROR:') or clean.startswith('TASK ERROR'):
if re.match(r'^\s*(?:ERROR:|TASK ERROR)', line, re.IGNORECASE):
if current_vm:
current_vm['status'] = 'error'
@@ -439,7 +440,7 @@ def _format_vzdump_body(parsed: Dict[str, Any], is_success: bool,
for vm in parsed.get('vms', []):
status = vm.get('status', '').lower()
icon = '\u2705' if status == 'ok' else '\u274C'
icon = '\u2705' if status == 'ok' else '\u274C' if status == 'error' else '\u2754'
# Determine VM/CT type prefix
vm_type = vm.get('type', '')
@@ -501,7 +502,8 @@ def _format_vzdump_body(parsed: Dict[str, Any], is_success: bool,
if vm_count > 0 or parsed.get('total_size'):
ok_count = sum(1 for v in parsed.get('vms', [])
if v.get('status', '').lower() == 'ok')
fail_count = vm_count - ok_count
fail_count = sum(1 for v in parsed.get('vms', [])
if v.get('status', '').lower() == 'error')
summary_parts = []
if vm_count:
@@ -779,11 +781,11 @@ TEMPLATES = {
# `{entity}` is populated by health_persistence.resolve_error()
# (via _entity_from_details) and by PollingCollector's spread of
# the original details blob. When absent, _SafeDict elides the
# placeholder and the title collapses back to "Resolved - <cat>"
# placeholder and the title collapses back to "No longer reported - <cat>"
# without a trailing dash.
'title': '{hostname}: Resolved - {category}{entity_suffix}',
'body': 'The {category} issue has been resolved.\n{reason}\n\U0001F6A6 Previous severity: {original_severity}\n\u23F1\uFE0F Duration: {duration}',
'label': 'Recovery notification',
'title': '{hostname}: No longer reported - {category}{entity_suffix}',
'body': 'The {category} issue is no longer in active health records.\n{reason}\n\U0001F6A6 Previous severity: {original_severity}\n\u23F1\uFE0F Time since first observation: {duration}',
'label': 'Health issue no longer reported',
'group': 'health',
'default_enabled': True,
},
@@ -999,9 +1001,9 @@ TEMPLATES = {
'default_enabled': False,
},
'backup_complete': {
'title': '{hostname} → {storage}: Backup complete — {vmname} ({vmid})',
'body': 'Backup of {vmname} (ID: {vmid}) completed successfully on {storage}.\nSize: {size}',
'label': 'Backup complete',
'title': '{hostname}: Backup outcome unconfirmed',
'body': 'The backup outcome could not be confirmed from this notice.',
'label': 'Backup report',
'group': 'backup',
'default_enabled': True,
},
@@ -1270,8 +1272,7 @@ TEMPLATES = {
'Stale node dirs removed: {stale_nodes}\n'
'Components reinstalled: {components}\n'
'Duration: {duration}\n'
'{warnings_block}\n'
'The node is now fully ready to use.'
'{warnings_block}'
),
'label': 'Host restore completed',
'group': 'services',
@@ -1855,6 +1856,16 @@ def render_template(event_type: str, data: Dict[str, Any],
)
if localized:
template[field] = localized
if event_type == 'backup_complete':
outcome = data.get('backup_outcome')
if outcome == 'confirmed':
template['title'] = runtime_message('backup.confirmedTitle', language,
hostname=data.get('hostname') or _get_hostname())
template['body'] = runtime_message('backup.confirmedBody', language)
elif outcome == 'failed':
template['title'] = runtime_message('backup.errorTitle', language,
hostname=data.get('hostname') or _get_hostname())
template['body'] = runtime_message('backup.errorBody', language)
# Ensure hostname is always available
variables = {
@@ -1997,7 +2008,6 @@ def render_template(event_type: str, data: Dict[str, Any],
# When the event came from PVE webhook with a full vzdump message,
# parse the table/logs and format a rich body instead of the sparse template.
pve_message = data.get('pve_message', '')
pve_title = data.get('pve_title', '')
# Check for custom formatter function
formatter_name = template.get('formatter')
@@ -2016,13 +2026,18 @@ def render_template(event_type: str, data: Dict[str, Any],
if parsed:
is_success = (event_type == 'backup_complete')
body_text = _format_vzdump_body(parsed, is_success, language=language)
# Preserve PVE's source title for English, but never leak it into a
# deterministic localized notification.
if pve_title and requested_language == 'en':
title = pve_title
if event_type == 'backup_complete' and data.get('backup_outcome') == 'failed':
error_lines = [line.strip() for line in pve_message.splitlines()
if re.match(r'^\s*(?:ERROR:|TASK ERROR)', line, re.IGNORECASE)]
if error_lines:
body_text += '\n' + '\n'.join(error_lines)
else:
# Couldn't parse -- use PVE raw message as body
body_text = pve_message.strip()
if event_type == 'backup_complete' and data.get('backup_outcome') != 'confirmed':
key = ('backup.errorBody' if data.get('backup_outcome') == 'failed'
else 'backup.unconfirmedBody')
body_text = runtime_message(key, language) + '\n' + body_text
elif event_type == 'system_mail' and pve_message:
# System mail -- use PVE message directly (mail bounce, cron, smartd)
body_text = pve_message.strip()[:1000]
@@ -2166,7 +2181,7 @@ EVENT_EMOJI = {
'host_backup_start': '\U0001F5C4️\U0001F680', # 🗄️🚀 cabinet + rocket
'host_backup_complete': '\U0001F5C4️✅', # 🗄️✅ cabinet + check
'host_backup_fail': '\U0001F5C4️❌', # 🗄️❌ cabinet + cross
'backup_complete': '\U0001F4BE\u2705', # 💾✅ floppy + check
'backup_complete': '\U0001F4BE', # 💾 neutral for digests without outcome metadata
'backup_warning': '\U0001F4BE\u26A0\uFE0F', # 💾⚠️ floppy + warning
'backup_fail': '\U0001F4BE\u274C', # 💾❌ floppy + cross
'snapshot_complete': '\U0001F4F8', # camera with flash
@@ -2204,14 +2219,14 @@ EVENT_EMOJI = {
'system_startup': '\U0001F680', # rocket (startup)
'system_shutdown': '\u23FB\uFE0F', # power symbol (Unicode)
'system_reboot': '\U0001F504',
'system_restore_completed': '✅', # check mark
'system_restore_completed': '\U0001F4CB', # post-restore task report (boot may have warnings)
'system_problem': '\u26A0\uFE0F',
'kernel_warning': '\u26A0\uFE0F',
'service_fail': '\u274C',
'oom_kill': '\U0001F4A3', # bomb
# Health
'new_error': '\U0001F198', # SOS
'error_resolved': '\u2705',
'error_resolved': '\U0001F4CB', # no longer active in health records, not proven recovery
'error_escalated': '\U0001F53A', # red triangle up
'health_degraded': '\u26A0\uFE0F',
'health_persistent': '\U0001F4CB', # clipboard
@@ -2363,6 +2378,10 @@ def enrich_with_emojis(event_type: str, title: str, body: str,
severity = data.get('severity', 'INFO')
icon = EVENT_EMOJI.get(event_type) or CATEGORY_EMOJI.get(group) or SEVERITY_ICONS.get(severity, '')
if event_type == 'backup_complete':
icon = {
'confirmed': '💾✅', 'failed': '💾❌',
}.get(str(data.get('backup_outcome') or ''), '💾❔')
# Build enriched title: replace severity circle with event-specific icon
# Current format: "hostname: Something" -> "ICON hostname: Something"
@@ -312,6 +312,7 @@ class RuntimeCatalogTests(unittest.TestCase):
"backup_complete",
{
"hostname": "pve01", "storage": "pbs-main", "vmname": "alpha", "vmid": "100",
"backup_outcome": "confirmed",
"pve_title": "Backup job finished",
"pve_message": (
"INFO: Starting Backup of VM 100 (qemu)\n"
@@ -322,7 +323,7 @@ class RuntimeCatalogTests(unittest.TestCase):
},
language="sk",
)
self.assertIn("Záloha dokončená", backup["title"])
self.assertIn("záloha dokončená", backup["title"])
self.assertNotIn("Backup job finished", backup["title"])
self.assertIn("Veľkosť: 1.5 GiB", backup["body"])
self.assertIn("Trvanie: 00:00:10", backup["body"])
@@ -627,6 +628,7 @@ class RuntimeCatalogTests(unittest.TestCase):
"_notification_language": "sk", "_event_type": event_type,
"_group": "backup", "hostname": "pve01", "vmid": "100",
"vmname": "alpha", "storage": "pbs-main",
"backup_outcome": "confirmed" if event_type == "backup_complete" else "unconfirmed",
},
)
self.assertIn(f">{localized_status}<", backup_html)
@@ -60,7 +60,7 @@ def _make_long_vzdump_report():
class VzdumpWebhookTruncationTests(unittest.TestCase):
def test_truncating_vzdump_report_at_4096_can_create_false_failed_backup(self):
def test_truncating_vzdump_report_at_4096_leaves_guest_unconfirmed(self):
full_message = _make_long_vzdump_report()
truncated_message = full_message[:4096]
@@ -82,8 +82,8 @@ class VzdumpWebhookTruncationTests(unittest.TestCase):
self.assertEqual(truncated_dockflare["name"], "dockflare")
self.assertEqual(truncated_dockflare["status"], "")
self.assertIn("❌ dockflare (129)", truncated_body)
self.assertIn("❌ 1 failed", truncated_body)
self.assertIn("❔ dockflare (129)", truncated_body)
self.assertNotIn("❌ 1 failed", truncated_body)
def test_webhook_handler_does_not_truncate_message_before_parsing(self):
source = (SCRIPTS_DIR / "flask_notification_routes.py").read_text()