fix(notifications): ship recovery helper and reject superseded proofs

This commit is contained in:
martino
2026-10-01 09:03:27 +02:00
parent 54eef11334
commit aac75a0753
6 changed files with 311 additions and 3 deletions
@@ -0,0 +1,83 @@
"""Execute the build's literal backend copies, then a clean shipped-only runtime."""
import os
from pathlib import Path
import re
import subprocess
import sys
import tempfile
import unittest
ROOT = Path(__file__).resolve().parents[3]
PROBE = r'''
import ast, pathlib, sys, types, typing
stage = pathlib.Path(sys.argv[1])
sys.path.insert(0, str(stage))
assert 'health_recovery' not in sys.modules
import notification_templates as templates
from notification_channels import EmailChannel
import health_recovery
for module in (templates, health_recovery, sys.modules['notification_channels']):
assert pathlib.Path(module.__file__).parent == stage
assert not any('projects/proxmox' in p for p in sys.path)
templates._get_hostname = lambda: 'node-a'
channel = object.__new__(EmailChannel)
channel.subject_prefix = '[ProxMenux]'
neutral = {'hostname': 'node-a', 'category': 'cpu', 'reason': 'CPU high', '_event_type': 'error_resolved'}
templates.render_template('error_resolved', neutral)
templates.enrich_with_emojis('error_resolved', 'Observation', 'Body', neutral)
channel._format_html('Observation', 'Body', 'OK', neutral)
templates.render_template('node_reconnect', {'hostname': 'node-a'})
now = 1000.0
calls = []
ns = dict(vars(typing), time=types.SimpleNamespace(time=lambda: now),
os=types.SimpleNamespace(cpu_count=lambda: 4),
psutil=types.SimpleNamespace(cpu_percent=lambda **kw: 20, cpu_count=lambda: 4),
health_persistence=types.SimpleNamespace(resolve_error=lambda *a, **kw: calls.append((a, kw))))
tree = ast.parse((stage/'health_monitor.py').read_text())
owner = next(n for n in tree.body if isinstance(n, ast.ClassDef) and n.name == 'HealthMonitor')
node = next(n for n in owner.body if isinstance(n, ast.FunctionDef) and n.name == '_check_cpu_with_hysteresis')
exec(compile(ast.Module(body=[node], type_ignores=[]), 'shipped_cpu', 'exec'), ns)
monitor = types.SimpleNamespace(state_history={'cpu_usage': [{'value':20,'time':now-i*10} for i in range(1,11)]},
CPU_WARNING=85, CPU_CRITICAL=95, CPU_RECOVERY=75, CPU_WARNING_DURATION=300,
CPU_CRITICAL_DURATION=300, CPU_RECOVERY_DURATION=120, _check_cpu_temperature=lambda:None)
assert ns['_check_cpu_with_hysteresis'](monitor)['status'] == 'OK'
assert len(calls) == 1 and calls[0][1]['check_evidence']['value'] == 20
proof = calls[0][1]['check_evidence']
native = dict(neutral, error_key='cpu_usage', check_evidence=proof, is_recovery=True, recovery_outcome='resolved')
# Actual body/icon/email consumers from the package, no source-module fixture rescue.
import unittest.mock
with unittest.mock.patch('health_recovery.time.time', return_value=now):
result = templates.render_template('error_resolved', native)
title, body = result['title'], result['body']
assert 'Resolved' in title and 'fresh health check' in body, (title, body, proof)
rich_title, rich_body = templates.enrich_with_emojis('error_resolved', title, body, native)
assert rich_title.startswith('✅')
assert 'background:#f0fdf4;' in channel._format_html(title, body, 'OK', native)
print('shipped-only neutral body/icon/email + CPU native measurement/proof consumers PASS')
'''
class PackagedRecoveryTests(unittest.TestCase):
def test_actual_copy_manifest_supports_isolated_recovery_runtime(self):
source = ROOT / 'AppImage/scripts'
build = (source / 'build_appimage.sh').read_text()
lines = [line for line in build.splitlines()
if re.match(r'^cp "\$SCRIPT_DIR/[^"/]+\.py" "\$APP_DIR/usr/bin/"', line)]
self.assertTrue(lines)
catalog_copy = re.search(r'^for locale in en de es fr it pt sk sv; do\n.*?^done$', build, re.MULTILINE | re.DOTALL)
if catalog_copy is None:
self.fail('The shipped locale-copy loop was not found in the actual build script')
lines.append(catalog_copy.group(0))
with tempfile.TemporaryDirectory(prefix='shipped-recovery-') as directory:
stage = Path(directory) / 'usr/bin'
stage.mkdir(parents=True)
copied = subprocess.run(['/bin/bash'], input='set -e\n'+'\n'.join(lines)+'\n', text=True,
capture_output=True, env={**os.environ, 'SCRIPT_DIR': str(source), 'APPIMAGE_ROOT': str(source.parent), 'APP_DIR': directory})
self.assertEqual(copied.returncode, 0, copied.stderr)
result = subprocess.run([sys.executable, '-I', '-B', '-c', PROBE, str(stage)],
cwd=directory, text=True, capture_output=True)
self.assertEqual(result.returncode, 0, result.stdout + result.stderr)
if __name__ == '__main__':
unittest.main()
@@ -0,0 +1,50 @@
"""Principal diagnostics must not be deduplicated against inventory substrings."""
import unittest
from notification_fixture import receive, LANGUAGES
from notification_final_fixture import deliver
def noisy_report(name='ordinary-web', storage='PBS', filename='null', cause='denied'):
# Frozen native Perl notifier/log-reader shape; fixed-column table is
# source-modeled from the official plaintext renderer, not a live PVE run.
return ('\nDetails\n=======\nVMID Name Status Time Size Filename\n'
f'100 {name} err 1m 1s 0 B {filename}\n\n'
'Total running time: 1m 1s\nTotal size: 0 B\n\nLogs\n====\n'
f'vzdump --all 1 --storage {storage} --mode snapshot\n\n' +
'\n'.join(f'100: 2026-09-29 17:00:00 ERROR: earlier diagnostic {i}' for i in range(12)) +
f'\n100: 2026-09-29 17:00:00 ERROR: {cause}\n\n')
class PrincipalInventoryTests(unittest.TestCase):
def assert_principal(self, raw, cause):
event = receive(raw, 'error', 'vzdump backup status (raw-host): backup failed: ' + cause)
for language in LANGUAGES:
for manual in (False, True):
with self.subTest(language=language, manual=manual):
result = deliver(event.event_type, {**event.data, 'hostname':'alias {rack.location}'},
event.severity, language, manual=manual)
exact = [line for line in result['body'].splitlines()
if line.strip() == cause or line.rstrip().endswith('ERROR: ' + cause)]
self.assertEqual(len(exact), 1)
diagnostics = [line for line in result['body'].splitlines() if 'ERROR:' in line]
self.assertLessEqual(len(diagnostics), 8)
self.assertLessEqual(len('\n'.join(diagnostics)), 1024)
self.assertTrue(all(len(line) <= 512 for line in diagnostics))
self.assertEqual(result['data']['pve_message'], raw)
self.assertNotIn('raw-host', result['text'])
self.assertIn('alias {rack.location}', result['text'])
self.assertEqual(result['text'].count('ERROR: ' + cause), 1)
def test_late_principal_survives_incidental_guest_storage_and_archive(self):
for name, storage, filename in (('ordinary-web','PBS','null'), ('denied','PBS','null'),
('ordinary-web','denied','null'), ('ordinary-web','PBS','denied.tar')):
with self.subTest(name=name, storage=storage, filename=filename):
self.assert_principal(noisy_report(name,storage,filename), 'denied')
def test_raw_principal_braces_and_markup_survive_once(self):
cause = 'denied {rack.location} <native>'
self.assert_principal(noisy_report('denied', cause=cause), cause)
if __name__ == '__main__':
unittest.main()
@@ -0,0 +1,170 @@
"""Durable native observation order, not wall-clock row identity, admits proof."""
import datetime
import json
import re
import sys
import types
import typing
import unittest
from unittest.mock import patch
from notification_recovery_fixture import case, Clock, BASE, cpu, service, sql, extract, scripts, TIME
from notification_fixture import LANGUAGES
from notification_final_fixture import deliver
def collector(store):
events = []
target = types.SimpleNamespace(_hostname='alias {rack.location}',
_ENTITY_MAP={'cpu':('node',''), 'pve_services':('node','')}, _first_poll_done=False,
_known_errors={}, _notified_severity={}, _last_notified={}, SAME_ERROR_COOLDOWN=86400,
_get_cooldown_from_db=lambda *a:BASE-1, _queue=types.SimpleNamespace(put=events.append),
_save_known_errors_meta=lambda:None)
ns = dict(vars(typing), time=TIME, json=json, re=re,
NotificationEvent=lambda *a, **kw:types.SimpleNamespace(event_type=a[0], severity=a[1], data=a[2], **kw),
startup_grace=types.SimpleNamespace(should_suppress_category=lambda *a:False))
target._guest_storage_error_is_now_foreign = extract(scripts/'notification_events.py',
'_guest_storage_error_is_now_foreign', 'PollingCollector', ns)
poll = extract(scripts/'notification_events.py', '_check_persistent_health', 'PollingCollector', ns)
def tick():
with patch.dict(sys.modules, {'health_persistence':types.SimpleNamespace(health_persistence=store),
'datetime':types.SimpleNamespace(**{**vars(datetime), 'datetime':Clock})}):
poll(target)
return target, events, tick
def abnormal(store, kind, value=99):
if kind == 'cpu':
return cpu(store, value, [{'value':value, 'time':Clock.epoch-i*5} for i in range(1,26)])
return service(store,3,'inactive\n')[0]
def normal(store, kind):
return cpu(store) if kind == 'cpu' else service(store)[0]
def initial_closure(store, kind):
key = 'cpu_usage' if kind == 'cpu' else 'pve_service_pvedaemon'
Clock.epoch = BASE-600
abnormal(store,kind,90)
target, events, tick = collector(store)
tick()
assert target._first_poll_done and key in target._known_errors and not events
snapshot = json.loads(json.dumps(target._known_errors))
Clock.epoch = BASE
assert normal(store,kind)['status'] == 'OK'
first = snapshot[key]['first_seen']
assert store.get_recovery_evidence(key,first)
return key, first, snapshot, target, events, tick
class RecoveryOrderTests(unittest.TestCase):
def assert_consumers(self, data, expected):
with patch('health_recovery.time.time', return_value=Clock.epoch):
for language in LANGUAGES:
for manual in (False,True):
with self.subTest(language=language, manual=manual):
result = deliver('error_resolved',data,'OK',language,manual=manual)
self.assertEqual('background:#f0fdf4;' in result['html'],expected)
self.assertIn('alias {rack.location}',result['text'])
quiet = deliver('error_resolved',data,'OK',language,quiet=True)
self.assertEqual(len(quiet['buffered']),1)
# The existing quiet digest stores the rendered title, not
# health body/proof metadata. Assert its exact outcome label.
from notification_fixture import templates
label = templates.render_template('error_resolved',data,language)['title'].split(': ',1)[-1]
self.assertIn(label, quiet['body'])
self.assertIn(label, quiet['buffered'][0][2])
def assert_superseded(self, kind):
for offset in (-100,0,1):
for value in ((90,99) if kind == 'cpu' else (99,)):
with self.subTest(kind=kind, offset=offset, value=value), case() as store:
key, first, snapshot, target, events, tick = initial_closure(store,kind)
oldrow = sql(store,'SELECT id,first_seen,resolved_at FROM errors')[0]
prior = sql(store,'SELECT id,event_type FROM events ORDER BY id')
Clock.epoch = BASE+offset
renewed = abnormal(store,kind,value)
self.assertEqual(renewed['status'],'WARNING' if kind == 'cpu' and value == 90 else 'CRITICAL')
self.assertEqual(sql(store,'SELECT id,first_seen,resolved_at FROM errors')[0],oldrow)
later = sql(store,'SELECT id,event_type FROM events ORDER BY id')[-1]
self.assertGreater(later[0],prior[-1][0])
self.assertEqual(later[1],'escalated' if kind == 'cpu' and value == 99 else 'updated')
if kind == 'cpu':
policy=json.loads(sql(store,'SELECT details FROM errors')[0][0])['cpu_policy']
self.assertEqual(policy,{'warning':85,'critical':95,'recovery':75})
Clock.epoch = BASE+2
tick()
self.assertEqual(len(events),1)
data = events[0].data
self.assertFalse(data['is_recovery'])
self.assertIsNone(store.get_recovery_evidence(key,first))
self.assert_consumers(data,False)
# Existing operations do not rearm the resolved row, so a
# later normal check cannot establish a NEW native closure.
before = sql(store,'SELECT id FROM events ORDER BY id')
Clock.epoch = BASE+3
self.assertEqual(normal(store,kind)['status'],'OK')
self.assertEqual(sql(store,'SELECT id FROM events ORDER BY id'),before)
self.assertIsNone(store.get_recovery_evidence(key,first))
def test_cpu_superseded_same_row_is_neutral_at_all_consumers(self):
self.assert_superseded('cpu')
def test_service_superseded_same_row_is_neutral_at_all_consumers(self):
self.assert_superseded('service')
def test_fresh_closure_and_repeated_noop_clear_keep_genuine_proof(self):
for kind in ('cpu','service'):
with self.subTest(kind=kind),case() as store:
key, first, snapshot, target, events, tick = initial_closure(store,kind)
before = sql(store,'SELECT id,event_type FROM events ORDER BY id')
Clock.epoch = BASE+1
for _ in range(3):
store.clear_error(key)
store.resolve_error(key,'generic repeat')
normal(store,kind)
self.assertEqual(sql(store,'SELECT id,event_type FROM events ORDER BY id'),before)
self.assertTrue(store.get_recovery_evidence(key,first))
tick()
self.assertEqual(len(events),1)
self.assertTrue(events[0].data['is_recovery'])
self.assert_consumers(events[0].data,True)
tick()
self.assertEqual(len(events),1)
def test_later_actual_generic_closure_blocks_older_proof_even_clock_rollback(self):
for kind in ('cpu','service'):
with self.subTest(kind=kind),case() as store:
key, first, *_ = initial_closure(store,kind)
Clock.epoch = BASE-100
with store._db_connection() as conn:
store._record_event(conn.cursor(),'cleared',key,{'reason':'generic actual closure','check_evidence':None})
conn.commit()
self.assertIsNone(store.get_recovery_evidence(key,first))
def test_explicit_acknowledged_row_suppresses_native_proof(self):
for kind in ('cpu','service'):
with self.subTest(kind=kind),case() as store:
key, first, snapshot, target, events, tick = initial_closure(store,kind)
store.acknowledge_error(key,suppression_hours=-1)
self.assertEqual(sql(store,'SELECT acknowledged FROM errors')[0][0],1)
self.assertIsNone(store.get_recovery_evidence(key,first))
tick()
self.assertEqual(events,[])
def test_new_incarnation_abnormal_order_blocks_previous_closure(self):
for kind in ('cpu','service'):
with self.subTest(kind=kind),case() as store:
key, first, *_ = initial_closure(store,kind)
old_id = sql(store,'SELECT id FROM errors')[0][0]
store.acknowledge_error(key,suppression_hours=-1)
store.clear_error(key)
Clock.epoch = BASE-100
abnormal(store,kind,99)
self.assertGreater(sql(store,'SELECT id FROM errors')[0][0],old_id)
self.assertEqual(sql(store,'SELECT event_type FROM events ORDER BY id DESC')[0][0],'new')
self.assertIsNone(store.get_recovery_evidence(key,first))
if __name__ == '__main__':
unittest.main()