From 3be15211b1180ccb7d566a68d4adc2a87c0e26a4 Mon Sep 17 00:00:00 2001 From: martino <32328813+f3rs3n@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:23:36 +0200 Subject: [PATCH] Fix principal cause caps and bind measured recovery to incident policy --- .../tests/notification_recovery_fixture.py | 86 +++++++ .../test_notification_recovery_corrections.py | 225 ++++++++++++++++++ .../test_notification_recovery_evidence.py | 138 +++++------ .../test_notification_review_corrections.py | 24 ++ AppImage/scripts/health_monitor.py | 28 ++- AppImage/scripts/health_persistence.py | 77 +++--- AppImage/scripts/health_recovery.py | 51 ++++ AppImage/scripts/notification_channels.py | 13 +- AppImage/scripts/notification_templates.py | 35 ++- 9 files changed, 537 insertions(+), 140 deletions(-) create mode 100644 .github/scripts/tests/notification_recovery_fixture.py create mode 100644 .github/scripts/tests/test_notification_recovery_corrections.py create mode 100644 .github/scripts/tests/test_notification_review_corrections.py create mode 100644 AppImage/scripts/health_recovery.py diff --git a/.github/scripts/tests/notification_recovery_fixture.py b/.github/scripts/tests/notification_recovery_fixture.py new file mode 100644 index 00000000..8e4f91d6 --- /dev/null +++ b/.github/scripts/tests/notification_recovery_fixture.py @@ -0,0 +1,86 @@ +"""Pinned, inert recovery review. No operational module imports or threads. +Run with python3 -S under bwrap --unshare-net; DBs and export use TMPDIR. +""" +import ast, contextlib, datetime, hashlib, io, json, os, pathlib, sqlite3, subprocess, sys, tarfile, tempfile, threading, types, typing +from unittest.mock import patch +REPO=pathlib.Path('/home/martino/projects/proxmox/notification-maintainer-followup') +BASE=datetime.datetime(2026,9,30,20,0,0).timestamp() +class Clock(datetime.datetime): + epoch=BASE + tick=0.0 + @classmethod + def now(cls,tz=None): + value=cls.fromtimestamp(cls.epoch,tz) + cls.epoch+=cls.tick + return value +TIME=types.SimpleNamespace(time=lambda:Clock.epoch) + +def extract(path,name,owner,ns): + tree=ast.parse(path.read_text()) + nodes=tree.body if owner is None else next(n.body for n in tree.body if isinstance(n,ast.ClassDef) and n.name==owner) + node=next(n for n in nodes if isinstance(n,(ast.FunctionDef,ast.AsyncFunctionDef)) and n.name==name) + node.decorator_list=[] + exec(compile(ast.Module(body=[node],type_ignores=[]),str(path),'exec'),ns) + return ns[name] + +from notification_fixture import SCRIPTS as scripts +pns=dict(vars(typing),datetime=Clock,timedelta=datetime.timedelta,json=json,sqlite3=sqlite3,contextmanager=contextlib.contextmanager,re=__import__('re'),_re_disk_base=__import__('re')) +pns['disk_base_name']=extract(scripts/'health_persistence.py','disk_base_name',None,pns) +methods={n:extract(scripts/'health_persistence.py',n,'HealthPersistence',pns) for n in ('_get_conn','_db_connection','_init_database','record_error','_record_error_impl','resolve_error','_resolve_error_impl','get_recovery_evidence','_record_event','_entity_from_details','clear_error','get_active_errors','is_error_active','is_error_acknowledged','_get_setting_impl','get_setting','set_setting','acknowledge_error','_acknowledge_error_impl','get_excluded_interface_names')} +def make_store(directory): + s=types.SimpleNamespace(db_path=pathlib.Path(directory)/'health.sqlite',_db_lock=threading.RLock(),DEFAULT_SUPPRESSION_HOURS=24,CATEGORY_SETTING_MAP={}) + for n,f in methods.items(): + if n=='_entity_from_details':setattr(s,n,f) + elif n=='_db_connection':setattr(s,n,types.MethodType(contextlib.contextmanager(f),s)) + else:setattr(s,n,types.MethodType(f,s)) + s._init_database() + return s +def sql(s,query,args=()): + with s._db_connection() as c: + data=c.execute(query,args).fetchall();c.commit();return data +def record(s,key='cpu_usage',category='cpu',reason='CPU high',details=None): + Clock.epoch=BASE-600 + with patch.dict(sys.modules,{'os':types.SimpleNamespace(path=types.SimpleNamespace(exists=lambda p:True))}):s.record_error(key,category,'WARNING',reason,details) + Clock.epoch=BASE-10 + with patch.dict(sys.modules,{'os':types.SimpleNamespace(path=types.SimpleNamespace(exists=lambda p:True))}):s.record_error(key,category,'WARNING',reason,details) + Clock.epoch=BASE + assert s.get_active_errors() + return sql(s,'SELECT first_seen FROM errors WHERE error_key=?',(key,))[0][0] +def cpu(s,current=20,history=None,warning=85,critical=95): + ns=dict(vars(typing),time=TIME,os=types.SimpleNamespace(cpu_count=lambda:4),health_persistence=s,psutil=types.SimpleNamespace(cpu_percent=lambda **kw:current,cpu_count=lambda:4)) + fn=extract(scripts/'health_monitor.py','_check_cpu_with_hysteresis','HealthMonitor',ns) + if history is None:history=[{'value':20,'time':BASE-i*10} for i in range(1,11)] + target=types.SimpleNamespace(state_history={'cpu_usage':list(history)},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) + refresh=extract(scripts/'health_monitor.py','_refresh_thresholds','HealthMonitor',{}) + with patch.dict(sys.modules,{'health_thresholds':types.SimpleNamespace(get=lambda section,key:({'warning':warning,'critical':critical}.get(key) if section=='cpu' else None))}):refresh(target) + return fn(target) +def poll(s,first,key='cpu_usage',category='cpu',reason='CPU high',details=None,first_done=True,foreign=False,restored=False): + events=[] + meta={'category':category,'reason':reason,'severity':'WARNING','first_seen':first,'details':details} + c=types.SimpleNamespace(_hostname='node-a',_ENTITY_MAP={'cpu':('node',''),'pve_services':('node',''),'network':('node','')},_first_poll_done=first_done,_known_errors={key:meta},_notified_severity={key:'WARNING'},_last_notified={key:BASE-1},SAME_ERROR_COOLDOWN=86400,_get_cooldown_from_db=lambda *a:BASE-1,_queue=types.SimpleNamespace(put=events.append),_guest_storage_error_is_now_foreign=lambda *a:foreign,_save_known_errors_meta=lambda:None) + ns=dict(vars(typing),time=TIME,json=json,re=__import__('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)) + c._guest_storage_error_is_now_foreign=extract(scripts/'notification_events.py','_guest_storage_error_is_now_foreign','PollingCollector',ns) + fn=extract(scripts/'notification_events.py','_check_persistent_health','PollingCollector',ns) + with patch.dict(sys.modules,{'health_persistence':types.SimpleNamespace(health_persistence=s),'flask_server':types.SimpleNamespace(get_proxmox_node_name=lambda:'node-a',get_cached_pvesh_cluster_resources_vm=lambda:[{'vmid':100,'type':'lxc','node':'node-b' if foreign else 'node-a'}]),'datetime':types.SimpleNamespace(**{**vars(datetime),'datetime':Clock})}): + if restored: + c._KNOWN_ERRORS_SETTING_KEY='pollingcollector_known_errors_v1' + s.set_setting(c._KNOWN_ERRORS_SETTING_KEY,json.dumps(c._known_errors));c._known_errors={} + extract(scripts/'notification_events.py','_load_known_errors_meta','PollingCollector',ns)(c) + c._first_poll_done=bool(c._known_errors) + fn(c) + return [e.data for e in events],c +def service(store, rc=0, stdout='active\n', raised=False, services=('pvedaemon',), clustered=False): + calls=[] + def run(argv, **kw): + calls.append((argv, kw)) + assert argv[:2] == ['systemctl', 'is-active'] + if raised: raise TimeoutError('inert timeout') + return types.SimpleNamespace(returncode=rc, stdout=stdout) + ns=dict(vars(typing),time=TIME,os=types.SimpleNamespace(path=types.SimpleNamespace(exists=lambda p:clustered)),subprocess=types.SimpleNamespace(run=run),health_persistence=store) + result=extract(scripts/'health_monitor.py','_check_pve_services','HealthMonitor',ns)(types.SimpleNamespace(PVE_SERVICES=list(services))) + return result,calls + +@contextlib.contextmanager +def case(): + Clock.epoch=BASE + with tempfile.TemporaryDirectory(prefix='recovery-review-db-',dir=os.environ['TMPDIR']) as d:yield make_store(d) diff --git a/.github/scripts/tests/test_notification_recovery_corrections.py b/.github/scripts/tests/test_notification_recovery_corrections.py new file mode 100644 index 00000000..6908cbfb --- /dev/null +++ b/.github/scripts/tests/test_notification_recovery_corrections.py @@ -0,0 +1,225 @@ +"""Native initializer, measurement methods, SQL writers/readers and collector.""" +import unittest +from notification_recovery_fixture import case, Clock, BASE, cpu, sql, poll + +class RecoveryCorrectionTests(unittest.TestCase): + def test_supported_low_warning_current_violation_is_neutral(self): + with case() as store: + Clock.epoch = BASE - 600 + initial = cpu(store, 80, [{'value':80, 'time':Clock.epoch-i*5} for i in range(1,26)], warning=50) + self.assertEqual(initial['status'], 'WARNING') + first = sql(store, 'SELECT first_seen FROM errors')[0][0] + Clock.epoch = BASE + result = cpu(store, 60, [{'value':60, 'time':BASE-i*5} for i in range(1,26)], warning=50) + events, _ = poll(store, first, reason=initial['reason']) + # Preserve operational clear/hysteresis behavior, not its factual claim. + self.assertEqual(result['status'], 'OK') + self.assertFalse(events[0]['is_recovery'], events) + self.assertIsNone(store.get_recovery_evidence('cpu_usage', first)) + + def test_original_policy_survives_repeated_native_updates(self): + import json + for later_warning in (85, 40): + with self.subTest(later_warning=later_warning), case() as store: + Clock.epoch = BASE - 600 + cpu(store, 80, [{'value':80, 'time':Clock.epoch-i*5} for i in range(1,26)], warning=50) + first = sql(store, 'SELECT first_seen FROM errors')[0][0] + for step in (400, 200): + Clock.epoch = BASE-step + cpu(store, 90, [{'value':90, 'time':Clock.epoch-i*5} for i in range(1,26)], warning=later_warning) + original = json.loads(sql(store, 'SELECT details FROM errors')[0][0])['cpu_policy'] + self.assertEqual(original['warning'], 50) + Clock.epoch = BASE + cpu(store, 20, warning=later_warning) + self.assertIsNone(store.get_recovery_evidence('cpu_usage', first)) + self.assertFalse(poll(store, first)[0][0]['is_recovery']) + + def test_clock_rollback_new_row_generic_clear_cannot_inherit_old_proof(self): + for rollback, reuse_first in ((True,False),(False,False),(True,True)): + with self.subTest(rollback=rollback,reuse_first=reuse_first),case() as store: + Clock.epoch = BASE-600 + cpu(store, 90, [{'value':90, 'time':Clock.epoch-i*5} for i in range(1,26)]) + first = sql(store, 'SELECT first_seen FROM errors')[0][0] + Clock.epoch = BASE; Clock.tick = .001 + try: cpu(store) + finally: Clock.tick = 0 + self.assertTrue(store.get_recovery_evidence('cpu_usage', first)) + store.acknowledge_error('cpu_usage', suppression_hours=-1) + store.clear_error('cpu_usage') + Clock.epoch = BASE-100 if rollback else BASE+1 + store.record_error('cpu_usage','cpu','WARNING','new incident after clock step',{}) + second = sql(store, 'SELECT first_seen FROM errors')[0][0] + self.assertNotEqual(first, second) + if reuse_first: + # Restored malformed snapshot with reused wall-clock identity; + # native row id and latest closure still prevent replay. + sql(store,'UPDATE errors SET first_seen=?',(first,));second=first + Clock.epoch = BASE+.0005 if rollback else BASE+2 + store.clear_error('cpu_usage'); Clock.epoch = BASE+3 + self.assertIsNone(store.get_recovery_evidence('cpu_usage', second)) + self.assertFalse(poll(store, second)[0][0]['is_recovery']) + + def test_malformed_native_and_manual_proof_is_neutral_at_all_consumers(self): + import copy, json, time + from unittest.mock import patch + from notification_fixture import LANGUAGES + from notification_final_fixture import deliver + with case() as store: + Clock.epoch = BASE-600 + cpu(store, 90, [{'value':90,'time':Clock.epoch-i*5} for i in range(1,26)]) + first = sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch = BASE; cpu(store) + event_data = json.loads(sql(store,"SELECT data FROM events WHERE event_type='resolved'")[0][0]) + data = poll(store,first)[0][0] + for field, bad in [('checked_at','bad'),('checked_at',True),('checked_at',float('inf')), + ('checked_at',10**400),('value',60),('max_sample',float('nan')), + ('normal_samples',True),('normal_samples',9),('checked_at',BASE+1),('checked_at',BASE-7201), + ('policy',{'warning':False,'critical':95,'recovery':75}), + ('policy',{'warning':96,'critical':95,'recovery':75}), + ('policy',{'warning':85,'critical':95}), + ('policy',{'warning':85,'critical':95,'recovery':float('nan')}), + ('policy',{'warning':85,'critical':95,'recovery':75,'extra':0})]: + broken = copy.deepcopy(event_data) + broken['check_evidence'][field] = bad + if field == 'value': broken['check_evidence']['policy']['warning'] = 50 + sql(store,"UPDATE events SET data=? WHERE event_type='resolved'",(json.dumps(broken),)) + with self.subTest(field=field,bad=str(bad)),patch('health_recovery.time.time',return_value=BASE): + self.assertIsNone(store.get_recovery_evidence('cpu_usage',first)) + for language in LANGUAGES: + for manual in (False,True): + result = deliver('error_resolved',{**data,'check_evidence':broken['check_evidence']},'OK',language,manual=manual) + self.assertNotIn('background:#f0fdf4;', result['html']) + # Caller content is trusted, not authenticated native proof; even + # well-shaped assertions require a valid time/type/numeric contract. + for proof in ({'check':'cpu_usage'}, {'check':'cpu_usage','checked_at':time.time()+1}): + self.assertNotIn('background:#f0fdf4;', deliver('error_resolved',{**data,'check_evidence':proof},'OK',manual=True)['html']) + + def test_exact_service_active_native_clear_reaches_recovery_consumers(self): + from unittest.mock import patch + from notification_recovery_fixture import service + from notification_fixture import LANGUAGES + from notification_final_fixture import deliver + with case() as store: + Clock.epoch = BASE-600 + service(store,3,'inactive\n') + first = sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch = BASE + result, calls = service(store) + self.assertEqual(result['status'],'OK') + self.assertEqual(calls,[(['systemctl','is-active','pvedaemon'], {'capture_output':True,'text':True,'timeout':2})]) + proof = store.get_recovery_evidence('pve_service_pvedaemon',first) + self.assertTrue(proof) + data = poll(store,first,'pve_service_pvedaemon','pve_services','PVE service pvedaemon is inactive',details={'service':'pvedaemon'})[0][0] + self.assertTrue(data['is_recovery']) + with patch('health_recovery.time.time',return_value=BASE): + for language in LANGUAGES: + for manual in (False,True): + rendered = deliver('error_resolved',data,'OK',language,manual=manual) + self.assertIn('background:#f0fdf4;',rendered['html']) + self.assertIn('pvedaemon', rendered['text']) + + def test_cpu_positive_default_low_policy_and_neutral_history_controls(self): + from notification_recovery_fixture import record + for warning,current,history,legacy,expected in ( + (85,20,None,False,True), (50,20,None,False,True), + (85,20,None,True,False), (85,99,[],False,False), + (85,20,[],False,False), + (85,20,[{'value':20,'time':BASE+i*5} for i in range(1,10)],False,False), + (85,20,[{'value':20,'time':BASE-121-i} for i in range(12)],False,False), + (85,99,None,False,False), + (85,float('nan'),None,False,False), (85,float('inf'),None,False,False), + (85,True,None,False,False), (85,10**400,None,False,False)): + with self.subTest(warning=warning,current=str(current),legacy=legacy), case() as store: + if legacy: first = record(store) + else: + Clock.epoch=BASE-600 + cpu(store,90,[{'value':90,'time':Clock.epoch-i*5} for i in range(1,26)],warning=warning) + first=sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch=BASE; cpu(store,current,history,warning=warning) + self.assertEqual(bool(store.get_recovery_evidence('cpu_usage',first)),expected) + + def test_service_unavailable_removed_overall_ok_and_ack_controls(self): + from notification_recovery_fixture import service, record + for rc,stdout,raised in ((3,'inactive\n',False),(4,'unknown\n',False),(0,'active extra\n',False),(1,'active\n',False),(0,'',True)): + with self.subTest(rc=rc,stdout=stdout,raised=raised),case() as store: + Clock.epoch=BASE-600; service(store,3,'inactive\n') + first=sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch=BASE; service(store,rc,stdout,raised) + self.assertIsNone(store.get_recovery_evidence('pve_service_pvedaemon',first)) + self.assertFalse(poll(store,first,'pve_service_pvedaemon','pve_services')[0]) + for removed in ('pvedaemon','corosync'): + with self.subTest(removed=removed),case() as store: + first=record(store,'pve_service_'+removed,'pve_services','service inactive',{'service':removed}) + result,calls=service(store,services=(),clustered=False) + self.assertEqual(result['status'],'OK'); self.assertEqual(calls,[]) + store.clear_error('pve_service_'+removed) + self.assertFalse(poll(store,first,'pve_service_'+removed,'pve_services')[0][0]['is_recovery']) + with case() as store: + first=record(store,'pve_service_corosync','pve_services','corosync inactive',{'service':'corosync'}) + result,calls=service(store,services=('pvedaemon',),clustered=False) + self.assertEqual(result['status'],'OK') + self.assertEqual([c[0][-1] for c in calls],['pvedaemon']) + self.assertTrue(store.is_error_active('pve_service_corosync')) + self.assertIsNone(store.get_recovery_evidence('pve_service_corosync',first)) + with case() as store: + first=record(store,'pve_service_pvedaemon','pve_services','service inactive') + store.acknowledge_error('pve_service_pvedaemon',suppression_hours=-1) + store.clear_error('pve_service_pvedaemon',check_evidence={'check':'pve_service_pvedaemon','checked_at':BASE,'service':'pvedaemon','state':'active','returncode':0}) + self.assertEqual(sql(store,'SELECT id FROM errors'),[]) + self.assertIsNone(store.get_recovery_evidence('pve_service_pvedaemon',first)) + + def test_native_binding_latest_closure_rollbacks_and_consistent_ack_read(self): + import contextlib, json, sqlite3 + for mutation in ('row_id','first_seen','closure','last_seen','latest_clear','latest_resolve','ack','event_insert_failure','ack_before_join'): + with self.subTest(mutation=mutation),case() as store: + Clock.epoch=BASE-600 + cpu(store,90,[{'value':90,'time':Clock.epoch-i*5} for i in range(1,26)]) + first=sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch=BASE + if mutation=='event_insert_failure': + sql(store,"CREATE TRIGGER fail_resolve BEFORE INSERT ON events WHEN NEW.event_type='resolved' BEGIN SELECT RAISE(ABORT,'fixture'); END") + self.assertEqual(cpu(store)['status'],'UNKNOWN') + self.assertIsNone(sql(store,'SELECT resolved_at FROM errors')[0][0]) + continue + cpu(store); self.assertTrue(store.get_recovery_evidence('cpu_usage',first)) + if mutation in ('row_id','first_seen','closure'): + data=json.loads(sql(store,"SELECT data FROM events WHERE event_type='resolved'")[0][0]) + field={'row_id':'id','first_seen':'first_seen','closure':'resolved_at'}[mutation] + data['incident'][field]='wrong' + sql(store,"UPDATE events SET data=? WHERE event_type='resolved'",(json.dumps(data),)) + elif mutation=='last_seen': sql(store,'UPDATE errors SET last_seen=?',(Clock.fromtimestamp(BASE+1).isoformat(),)) + elif mutation.startswith('latest_'): + sql(store,"INSERT INTO events(event_type,error_key,timestamp,data) VALUES(?,'cpu_usage',?,'{}')",('cleared' if mutation=='latest_clear' else 'resolved',Clock.now().isoformat())) + elif mutation=='ack': store.acknowledge_error('cpu_usage',suppression_hours=-1) + else: + original=store._db_connection; triggered=[] + class Proxy: + def __init__(self,connection):self.connection=connection + def __getattr__(self,name):return getattr(self.connection,name) + def execute(self,query,args=()): + if 'FROM errors e JOIN events' in query and not triggered: + triggered.append(True) + store.acknowledge_error('cpu_usage',suppression_hours=-1) + return self.connection.execute(query,args) + @contextlib.contextmanager + def interleaved(**kw): + with original(**kw) as conn: yield Proxy(conn) + store._db_connection=interleaved + self.assertIsNone(store.get_recovery_evidence('cpu_usage',first)) + if mutation=='ack_before_join': self.assertTrue(triggered) + with case() as store: + sql(store,"INSERT INTO errors(error_key,category,severity,reason,first_seen,last_seen) VALUES('cpu_usage','cpu','WARNING','fixture','x','x')") + with self.assertRaises(sqlite3.IntegrityError): + sql(store,"INSERT INTO errors(error_key,category,severity,reason,first_seen,last_seen) VALUES('cpu_usage','cpu','WARNING','duplicate','x','x')") + + def test_malformed_history_declines_proof_without_changing_operational_clear(self): + with case() as store: + Clock.epoch=BASE-600 + cpu(store,90,[{'value':90,'time':Clock.epoch-i*5} for i in range(1,26)]) + first=sql(store,'SELECT first_seen FROM errors')[0][0] + Clock.epoch=BASE + result=cpu(store,20,[{'value':10**400,'time':BASE-i*5} for i in range(1,10)]) + self.assertEqual(result['status'],'OK') + self.assertIsNone(store.get_recovery_evidence('cpu_usage',first)) + +if __name__ == '__main__': unittest.main() diff --git a/.github/scripts/tests/test_notification_recovery_evidence.py b/.github/scripts/tests/test_notification_recovery_evidence.py index 45f316f8..096feed2 100644 --- a/.github/scripts/tests/test_notification_recovery_evidence.py +++ b/.github/scripts/tests/test_notification_recovery_evidence.py @@ -1,100 +1,68 @@ -"""Fresh existing-check provenance; extracted consumers, real disposable SQLite.""" -import contextlib -import datetime +"""Fresh existing-check provenance; native initializer and disposable SQLite.""" import json -import os -import sqlite3 -import tempfile import time -import types import unittest from unittest.mock import patch -from notification_fixture import extract, SCRIPTS, templates, LANGUAGES +from notification_fixture import templates, LANGUAGES from notification_final_fixture import deliver +from notification_recovery_fixture import case, Clock, BASE, cpu, sql, poll + + +def original_cpu(store): + Clock.epoch = BASE-600 + result = cpu(store, 90, [{'value':90,'time':Clock.epoch-i*5} for i in range(1,26)]) + assert result['status'] == 'WARNING' + Clock.epoch = BASE + return sql(store, 'SELECT first_seen FROM errors')[0][0] + class RecoveryEvidenceTests(unittest.TestCase): def test_cpu_success_provenance_is_persisted_only_after_normal_samples(self): - events=[] - with tempfile.TemporaryDirectory() as scratch: - db=scratch+'/health.sqlite' - conn=sqlite3.connect(db) - conn.execute('CREATE TABLE errors(id INTEGER PRIMARY KEY,error_key TEXT,details TEXT,resolved_at TEXT,resolution_type TEXT,resolution_reason TEXT)') - conn.execute("INSERT INTO errors(error_key,details) VALUES ('cpu_usage','{}')") - conn.commit();conn.close() - @contextlib.contextmanager - def connection(): - c=sqlite3.connect(db) - try:yield c - finally:c.close() - ns={'datetime':datetime.datetime,'json':json} - resolve=extract(SCRIPTS/'health_persistence.py','_resolve_error_impl','HealthPersistence',ns) - store=types.SimpleNamespace(_db_connection=connection,_entity_from_details=lambda details:'',_record_event=lambda cursor,kind,key,data:events.append(data)) - store.resolve_error=lambda key,reason,**kw:resolve(store,key,reason,**kw) - ns={'Dict':dict,'Any':object,'os':os,'time':time,'health_persistence':store,'psutil':types.SimpleNamespace(cpu_percent=lambda **kw:20,cpu_count=lambda:4)} - check=extract(SCRIPTS/'health_monitor.py','_check_cpu_with_hysteresis','HealthMonitor',ns) - target=types.SimpleNamespace(state_history={'cpu_usage':[{'value':20,'time':time.time()-i*10} for i in range(10)]},CPU_CRITICAL=95,CPU_WARNING=85,CPU_RECOVERY=75,CPU_CRITICAL_DURATION=300,CPU_WARNING_DURATION=300,CPU_RECOVERY_DURATION=120,_check_cpu_temperature=lambda:None) - result=check(target) - self.assertEqual(result['status'],'OK') - self.assertTrue(events[-1].get('check_evidence'),events) - proof=events[-1]['check_evidence'] - self.assertEqual(proof['check'],'cpu_usage') - self.assertGreaterEqual(proof['checked_at'],time.time()-5) - # Existing generic resolve callers (cleanup/exclusion) get no proof. - conn=sqlite3.connect(db);conn.execute('UPDATE errors SET resolved_at=NULL');conn.commit();conn.close() - resolve(store,'cpu_usage','No longer present') - self.assertFalse(events[-1].get('check_evidence')) + with case() as store: + first = original_cpu(store) + self.assertEqual(cpu(store)['status'], 'OK') + proof = store.get_recovery_evidence('cpu_usage', first) + self.assertTrue(proof) + self.assertEqual(proof['check'], 'cpu_usage') + self.assertEqual(proof['checked_at'], BASE) + # A generic closure never gains proof; use another actual native row. + store.record_error('pve_service_test','pve_services','CRITICAL','inactive') + store.resolve_error('pve_service_test','No longer present') + self.assertFalse(json.loads(sql(store,"SELECT data FROM events ORDER BY id DESC LIMIT 1")[0][0]).get('check_evidence')) def test_recovery_query_requires_fresh_same_incident_proof(self): - with tempfile.TemporaryDirectory() as scratch: - db=scratch+'/health.sqlite' - conn=sqlite3.connect(db) - conn.execute('CREATE TABLE errors(id INTEGER PRIMARY KEY,error_key TEXT,first_seen TEXT,last_seen TEXT,resolved_at TEXT,acknowledged INTEGER)') - conn.execute('CREATE TABLE events(id INTEGER PRIMARY KEY,event_type TEXT,error_key TEXT,timestamp TEXT,data TEXT)') - now=datetime.datetime.now(); first=(now-datetime.timedelta(minutes=10)).isoformat(); last=(now-datetime.timedelta(minutes=1)).isoformat(); resolved=now.isoformat() - proof={'check':'cpu_usage','checked_at':now.timestamp()} - conn.execute('INSERT INTO errors VALUES(1,?,?,?,?,0)',('cpu_usage',first,last,resolved)) - conn.execute('INSERT INTO events VALUES(1,?,?,?,?)',('resolved','cpu_usage',resolved,json.dumps({'check_evidence':proof}))) - conn.commit();conn.close() - @contextlib.contextmanager - def connection(**kwargs): - c=sqlite3.connect(db) - try:yield c - finally:c.close() - ns={'datetime':datetime.datetime,'json':json,'time':time} - tree=(SCRIPTS/'health_persistence.py').read_text() - query=extract(SCRIPTS/'health_persistence.py','get_recovery_evidence','HealthPersistence',ns) if 'def get_recovery_evidence(' in tree else lambda *args:None - store=types.SimpleNamespace(_db_connection=connection) - self.assertEqual(query(store,'cpu_usage',first),proof) - self.assertIsNone(query(store,'cpu_usage','different incident')) - for field,value in [('acknowledged',1),('resolved_at',None),('last_seen',(now+datetime.timedelta(seconds=1)).isoformat())]: - conn=sqlite3.connect(db);conn.execute(f'UPDATE errors SET {field}=?',(value,));conn.commit();conn.close() - self.assertIsNone(query(store,'cpu_usage',first)) - conn=sqlite3.connect(db);conn.execute('UPDATE errors SET acknowledged=0,resolved_at=?,last_seen=?',(resolved,last));conn.commit();conn.close() - for bad in (None,{'check':'cpu_usage','checked_at':now.timestamp()-7201},{'check':'storage_removed','checked_at':now.timestamp()},{'check':'cpu_usage','checked_at':float('inf')},{'check':'cpu_usage','checked_at':10**400}): - conn=sqlite3.connect(db);conn.execute('UPDATE events SET data=?',(json.dumps({'check_evidence':bad}),));conn.commit();conn.close() - self.assertIsNone(query(store,'cpu_usage',first)) + with case() as store: + first = original_cpu(store); cpu(store) + proof = store.get_recovery_evidence('cpu_usage',first) + self.assertTrue(proof) + self.assertIsNone(store.get_recovery_evidence('cpu_usage','different incident')) + saved = sql(store,'SELECT last_seen,resolved_at FROM errors')[0] + for field,value in [('acknowledged',1),('resolved_at',None),('last_seen',Clock.fromtimestamp(BASE+1).isoformat())]: + sql(store,f'UPDATE errors SET {field}=?',(value,)) + self.assertIsNone(store.get_recovery_evidence('cpu_usage',first)) + sql(store,'UPDATE errors SET acknowledged=0,last_seen=?,resolved_at=?',saved) + data = json.loads(sql(store,"SELECT data FROM events WHERE event_type='resolved'")[0][0]) + for bad in (None,{'check':'cpu_usage','checked_at':BASE-7201}, {'check':'storage_removed','checked_at':BASE}, {'check':'cpu_usage','checked_at':float('inf')}, {'check':'cpu_usage','checked_at':10**400}): + sql(store,"UPDATE events SET data=? WHERE event_type='resolved'",(json.dumps({**data,'check_evidence':bad}),)) + self.assertIsNone(store.get_recovery_evidence('cpu_usage',first)) def test_poller_and_all_consumers_distinguish_proven_recovery_from_disappearance(self): - import sys - ns={'time':time,'json':json,'Dict':dict,'NotificationEvent':lambda *a,**kw:types.SimpleNamespace(event_type=a[0],severity=a[1],data=a[2])} - poll=extract(SCRIPTS/'notification_events.py','_check_persistent_health','PollingCollector',ns) - for proof in (None,{'check':'cpu_usage','checked_at':time.time()}): - events=[] - store=types.SimpleNamespace(get_active_errors=lambda:[],is_error_acknowledged=lambda key:False,get_recovery_evidence=lambda *a:proof) - collector=types.SimpleNamespace(_hostname='node-a',_ENTITY_MAP={'cpu':('node','')},_first_poll_done=True,_known_errors={'cpu_usage':{'category':'cpu','reason':'CPU high','severity':'WARNING','first_seen':'2026-09-30T00:00:00'}},_notified_severity={'cpu_usage':'WARNING'},_last_notified={'cpu_usage':1},_queue=types.SimpleNamespace(put=events.append),_guest_storage_error_is_now_foreign=lambda *a:False,_save_known_errors_meta=lambda:None) - with patch.dict(sys.modules,{'health_persistence':types.SimpleNamespace(health_persistence=store)}):poll(collector) - self.assertEqual(len(events),1) - event=events[0] - self.assertEqual(event.data.get('recovery_outcome'),'resolved' if proof else 'no_longer_reported') - self.assertEqual(event.data['is_recovery'],bool(proof)) - for lang in LANGUAGES: - for manual in (False,True): - result=deliver(event.event_type,event.data,event.severity,lang,manual=manual) - if proof: - self.assertIn(templates.runtime_message('healthRecovery.title',lang,hostname='node-a',category='cpu',entity_suffix=''),result['title']) - self.assertIn('background:#f0fdf4;',result['html']) - else:self.assertNotIn('background:#f0fdf4;',result['html']) - self.assertEqual(result['text'].count(event.data['reason']),1) + for proved in (False,True): + with case() as store: + first = original_cpu(store) + if proved: cpu(store) + else: store.resolve_error('cpu_usage','No longer present') + data = poll(store,first,reason='CPU high')[0][0] + self.assertEqual(data['is_recovery'],proved) + self.assertEqual(data['recovery_outcome'],'resolved' if proved else 'no_longer_reported') + with patch('health_recovery.time.time',return_value=BASE): + for lang in LANGUAGES: + for manual in (False,True): + result = deliver('error_resolved',data,'OK',lang,manual=manual) + self.assertEqual('background:#f0fdf4;' in result['html'],proved) + if proved: + self.assertIn(templates.runtime_message('healthRecovery.title',lang,hostname='node-a',category='cpu',entity_suffix=''),result['title']) + self.assertEqual(result['text'].count(data['reason']),1) def test_manual_recovery_flag_alone_is_not_authoritative_evidence(self): data={'hostname':'node-a','category':'cpu','reason':'Observation disappeared','duration':'1h','original_severity':'WARNING','recovery_outcome':'resolved'} diff --git a/.github/scripts/tests/test_notification_review_corrections.py b/.github/scripts/tests/test_notification_review_corrections.py new file mode 100644 index 00000000..3928a203 --- /dev/null +++ b/.github/scripts/tests/test_notification_review_corrections.py @@ -0,0 +1,24 @@ +"""Review regressions at actual consumer seams; no operational imports.""" +import unittest +from notification_fixture import receive, LANGUAGES +from notification_final_fixture import deliver + +class ReviewCorrectionTests(unittest.TestCase): + def test_subject_equivalent_late_error_is_reserved_before_cap(self): + cause = 'job-end hook denied' + # Frozen native Perl notifier/log-reader output; Rust table source-modeled. + raw = '\nDetails\n=======\nVMID Name Status Time Size Filename \n100 web err 1m 1s 0 B null \n\nTotal running time: 1m 1s\nTotal size: 0 B\n\nLogs\n====\nvzdump --all 1 --storage PBS --mode snapshot\n\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 0\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 1\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 2\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 3\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 4\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 5\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 6\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 7\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 8\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 9\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 10\n100: 2026-09-29 17:00:00 ERROR: earlier diagnostic 11\n100: 2026-09-29 17:00:00 ERROR: job-end hook denied\n\n\n' + 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) + self.assertEqual(result['text'].count(cause), 1) + self.assertNotIn('raw-host', result['text']) + 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) + +if __name__ == '__main__': unittest.main() diff --git a/AppImage/scripts/health_monitor.py b/AppImage/scripts/health_monitor.py index 565f68e3..fdfe931f 100644 --- a/AppImage/scripts/health_monitor.py +++ b/AppImage/scripts/health_monitor.py @@ -1427,6 +1427,8 @@ class HealthMonitor: 'details': f'Sustained for {actual_duration}s above {self.CPU_CRITICAL}%.', 'cpu_percent': cpu_percent, 'duration': actual_duration, + 'cpu_policy': {'warning': self.CPU_WARNING, 'critical': self.CPU_CRITICAL, + 'recovery': self.CPU_RECOVERY}, }, ) elif len(warning_samples) >= WARNING_MIN_SAMPLES and len(recovery_samples) < RECOVERY_MIN_SAMPLES: @@ -1446,6 +1448,8 @@ class HealthMonitor: 'details': f'Sustained for {actual_duration}s above {self.CPU_WARNING}%.', 'cpu_percent': cpu_percent, 'duration': actual_duration, + 'cpu_policy': {'warning': self.CPU_WARNING, 'critical': self.CPU_CRITICAL, + 'recovery': self.CPU_RECOVERY}, }, ) else: @@ -1453,8 +1457,21 @@ class HealthMonitor: reason = None # CPU is normal - auto-resolve any existing CPU errors evidence = None - if cpu_percent < self.CPU_RECOVERY and len(recovery_samples) >= RECOVERY_MIN_SAMPLES: - evidence = {'check': 'cpu_usage', 'checked_at': current_time} + # Presentation proof is stricter than operational hysteresis: + # a supported warning can be below the fixed recovery cutoff. + from health_recovery import _finite_number + criterion = min(self.CPU_WARNING, self.CPU_RECOVERY) + normal_samples = [entry for entry in self.state_history[state_key] + if _finite_number(entry['value']) and 0 <= entry['value'] < criterion + and _finite_number(entry['time']) + and 0 <= current_time - entry['time'] <= self.CPU_RECOVERY_DURATION] + if (_finite_number(cpu_percent) and 0 <= cpu_percent < criterion + and len(normal_samples) >= RECOVERY_MIN_SAMPLES): + evidence = {'check': 'cpu_usage', 'checked_at': current_time, + 'value': cpu_percent, 'normal_samples': len(normal_samples), + 'max_sample': max(entry['value'] for entry in normal_samples), + 'policy': {'warning': self.CPU_WARNING, 'critical': self.CPU_CRITICAL, + 'recovery': self.CPU_RECOVERY}} health_persistence.resolve_error('cpu_usage', 'CPU usage returned to normal', check_evidence=evidence) @@ -3904,6 +3921,7 @@ class HealthMonitor: failed_services = [] service_details = {} + active_evidence = {} for service in services_to_check: try: @@ -3918,6 +3936,10 @@ class HealthMonitor: if result.returncode != 0 or status != 'active': failed_services.append(service) service_details[service] = status or 'inactive' + else: + active_evidence[service] = {'check': f'pve_service_{service}', + 'checked_at': time.time(), 'service': service, 'state': status, + 'returncode': result.returncode} except Exception: failed_services.append(service) service_details[service] = 'error' @@ -3932,7 +3954,7 @@ class HealthMonitor: if svc not in failed_services: error_key = f'pve_service_{svc}' if health_persistence.is_error_active(error_key): - health_persistence.clear_error(error_key) + health_persistence.clear_error(error_key, check_evidence=active_evidence.get(svc)) # Build checks dict with status per service checks = {} diff --git a/AppImage/scripts/health_persistence.py b/AppImage/scripts/health_persistence.py index 35458418..8c108c13 100644 --- a/AppImage/scripts/health_persistence.py +++ b/AppImage/scripts/health_persistence.py @@ -600,7 +600,7 @@ class HealthPersistence: cursor.execute(''' SELECT id, acknowledged, resolved_at, category, severity, first_seen, - notification_sent, suppression_hours, acknowledged_at + notification_sent, suppression_hours, acknowledged_at, details FROM errors WHERE error_key = ? ''', (error_key,)) existing = cursor.fetchone() @@ -609,7 +609,7 @@ class HealthPersistence: if existing: (err_id, ack, resolved_at, old_cat, old_severity, first_seen, - notif_sent, stored_suppression, acknowledged_at) = existing + notif_sent, stored_suppression, acknowledged_at, old_details_json) = existing if ack == 1: # SAFETY OVERRIDE: Critical CPU temperature ALWAYS re-triggers @@ -680,6 +680,18 @@ class HealthPersistence: conn.commit() return event_info + # Original CPU policy is immutable for this row's incident. + # Never upgrade a legacy row from later/current settings. + if error_key == 'cpu_usage': + try: + old_details = json.loads(old_details_json or '{}') + except (ValueError, TypeError): + old_details = {} + details = dict(details) if isinstance(details, dict) else {} + details.pop('cpu_policy', None) + if isinstance(old_details, dict) and 'cpu_policy' in old_details: + details['cpu_policy'] = old_details['cpu_policy'] + details_json = json.dumps(details) # Not acknowledged - update existing active error cursor.execute(''' UPDATE errors @@ -776,7 +788,7 @@ class HealthPersistence: # was created — otherwise "Storage 'Tuxis' unavailable" # comes back as "Resolved - Storage" with no identity. cursor.execute( - 'SELECT details FROM errors WHERE error_key = ? ORDER BY id DESC LIMIT 1', + 'SELECT details, id, first_seen, resolved_at FROM errors WHERE error_key = ? ORDER BY id DESC LIMIT 1', (error_key,), ) row = cursor.fetchone() @@ -786,11 +798,18 @@ class HealthPersistence: stored_details = json.loads(row[0]) except Exception: stored_details = None + # Legacy/mismatched original policy cannot certify normality. + if error_key == 'cpu_usage' and (not isinstance(stored_details, dict) + or not isinstance(check_evidence, dict) + or check_evidence.get('policy') != stored_details.get('cpu_policy') + or not stored_details.get('cpu_policy')): + check_evidence = None self._record_event(cursor, 'resolved', error_key, { 'reason': reason, # Only explicit current-check callers attach this proof. # Generic resolve/cleanup remains neutral. 'check_evidence': check_evidence, + 'incident': {'id': row[1], 'first_seen': row[2], 'resolved_at': row[3]}, 'entity': self._entity_from_details(stored_details), 'details': stored_details or {}, }) @@ -800,41 +819,40 @@ class HealthPersistence: def get_recovery_evidence(self, error_key: str, first_seen: str): """Return fresh same-incident native check proof, never absence of errors. - Initially only the host CPU check has a stable condition identity. Other + Host CPU and exact per-service active checks carry provenance. Other checks, generic clears, excluded/deleted records and legacy events stay neutral until they have equivalent per-condition provenance. """ - if error_key != 'cpu_usage' or not first_seen: + if not isinstance(error_key, str) or not first_seen or not ( + error_key == 'cpu_usage' or error_key.startswith('pve_service_')): return None try: - with self._db_connection() as conn: + # One SQLite statement is one consistent row/ack/closure snapshot. + # Latest closure by event id, never search past a generic clear. + with self._db_lock, self._db_connection() as conn: row = conn.execute(''' - SELECT first_seen, last_seen, resolved_at, acknowledged - FROM errors WHERE error_key = ? ORDER BY id DESC LIMIT 1 + SELECT e.first_seen, e.last_seen, e.resolved_at, e.acknowledged, + e.id, v.timestamp, v.data + FROM errors e JOIN events v ON v.id = ( + SELECT id FROM events WHERE error_key = e.error_key + AND event_type IN ('resolved', 'cleared') ORDER BY id DESC LIMIT 1 + ) WHERE e.error_key = ? ''', (error_key,)).fetchone() - if not row or row[0] != first_seen or not row[2] or row[3]: - return None - event = conn.execute(''' - SELECT timestamp, data FROM events - WHERE error_key = ? AND event_type = 'resolved' - ORDER BY id DESC LIMIT 1 - ''', (error_key,)).fetchone() - if not event: + if not row or row[0] != first_seen or not row[2] or row[3]: return None - proof = json.loads(event[1]).get('check_evidence') - if not isinstance(proof, dict) or proof.get('check') != error_key: + event_data = json.loads(row[6]) + if event_data.get('incident') != {'id': row[4], 'first_seen': row[0], 'resolved_at': row[2]}: return None - checked = proof.get('checked_at') - if not isinstance(checked, (int, float)) or isinstance(checked, bool): + proof = event_data.get('check_evidence') + from health_recovery import valid_check_evidence + if not valid_check_evidence(error_key, proof, now=datetime.now().timestamp()): return None - checked = float(checked) - # Reuse the collector's existing two-hour freshness boundary. - now = datetime.now().timestamp() - if not 0 <= now - checked <= 7200: + if error_key == 'cpu_usage' and proof.get('policy') != event_data.get('details', {}).get('cpu_policy'): return None + checked = float(proof['checked_at']) last_seen = datetime.fromisoformat(row[1]).timestamp() resolved = datetime.fromisoformat(row[2]).timestamp() - recorded = datetime.fromisoformat(event[0]).timestamp() + recorded = datetime.fromisoformat(row[5]).timestamp() if not last_seen <= checked <= resolved <= recorded: return None return proof @@ -906,7 +924,7 @@ class HealthPersistence: return False - def clear_error(self, error_key: str): + def clear_error(self, error_key: str, *, check_evidence=None): """ Remove/resolve a specific error immediately. Used when the condition that caused the error no longer exists @@ -928,7 +946,7 @@ class HealthPersistence: # Check if this error was acknowledged (dismissed) cursor.execute(''' - SELECT acknowledged FROM errors WHERE error_key = ? + SELECT acknowledged, id, first_seen FROM errors WHERE error_key = ? ''', (error_key,)) row = cursor.fetchone() @@ -947,7 +965,10 @@ class HealthPersistence: ''', (now, error_key)) if cursor.rowcount > 0: - self._record_event(cursor, 'cleared', error_key, {'reason': 'condition_resolved'}) + self._record_event(cursor, 'cleared', error_key, { + 'reason': 'condition_resolved', 'check_evidence': check_evidence, + 'incident': {'id': row[1], 'first_seen': row[2], 'resolved_at': now}, + }) conn.commit() diff --git a/AppImage/scripts/health_recovery.py b/AppImage/scripts/health_recovery.py new file mode 100644 index 00000000..9a8a00df --- /dev/null +++ b/AppImage/scripts/health_recovery.py @@ -0,0 +1,51 @@ +"""Bounded recovery metadata admission, without importing monitor singletons. + +Native persistence additionally binds the exact incident and closure. Manual +notifications remain authenticated caller assertions: shape validation cannot +establish that an asserted measurement actually happened. +""" +import math +import time +from typing import TypeGuard + + +def _finite_number(value) -> TypeGuard[int | float]: + try: + return isinstance(value, (int, float)) and not isinstance(value, bool) and math.isfinite(value) + except (OverflowError, ValueError, TypeError): + return False + + +def valid_check_evidence(error_key, proof, *, now=None): + """Validate a supported measurement contract, not its external authenticity.""" + if not isinstance(proof, dict) or proof.get('check') != error_key: + return False + checked = proof.get('checked_at') + now = time.time() if now is None else now + if not _finite_number(checked) or not _finite_number(now) or not 0 <= now-checked <= 7200: + return False + if error_key == 'cpu_usage': + policy = proof.get('policy') + if not isinstance(policy, dict) or set(policy) != {'warning', 'critical', 'recovery'}: + return False + if not all(_finite_number(v) and 1 <= v <= 100 for v in policy.values()): + return False + if policy['warning'] > policy['critical']: + return False + value, maximum, count = proof.get('value'), proof.get('max_sample'), proof.get('normal_samples') + return (_finite_number(value) and _finite_number(maximum) + and 0 <= value <= maximum < min(policy['warning'], policy['recovery']) + and isinstance(count, int) and not isinstance(count, bool) and count >= 10) + if isinstance(error_key, str) and error_key.startswith('pve_service_'): + service = error_key[len('pve_service_'):] + return (bool(service) and proof.get('service') == service + and proof.get('state') == 'active' + and type(proof.get('returncode')) is int and proof['returncode'] == 0) + return False + + +def presents_recovery(data): + """One presentation predicate shared by template, icon and email badge.""" + return (isinstance(data, dict) and data.get('recovery_outcome') == 'resolved' + and data.get('is_recovery') is True + and valid_check_evidence(data.get('error_key'), data.get('check_evidence'))) diff --git a/AppImage/scripts/notification_channels.py b/AppImage/scripts/notification_channels.py index b0b106ab..584ad7a1 100644 --- a/AppImage/scripts/notification_channels.py +++ b/AppImage/scripts/notification_channels.py @@ -1038,10 +1038,8 @@ class EmailChannel(NotificationChannel): # Determine group for section header event_type = data.get('_event_type', '') if event_type == 'error_resolved': - if (data.get('recovery_outcome') == 'resolved' - and data.get('is_recovery') is True - and isinstance(data.get('check_evidence'), dict) - and data['check_evidence'].get('check') == 'cpu_usage'): + from health_recovery import presents_recovery + if presents_recovery(data): sev.update(self._SEV_STYLE['OK']) sev['label'] = _runtime_notification_text('healthRecovery.status', data) else: @@ -1070,8 +1068,11 @@ class EmailChannel(NotificationChannel): or backup_email or data.get('_restore_summary') or data.get('_backup_summary')) temp_cell_wrap = 'word-wrap:break-word;overflow-wrap:break-word;word-break:break-word;' if wrap_body else '' temp_table_layout = 'table-layout:fixed;' if wrap_body else '' - backup_title_wrap = temp_cell_wrap if backup_email else '' - backup_metadata_layout = 'table-layout:fixed;' if backup_email else '' + # Recovery exposes the same literal host context as backup notices. + # Keep wrapping event-scoped; unrelated mail remains byte-identical. + context_email = backup_email or event_type == 'error_resolved' + backup_title_wrap = temp_cell_wrap if context_email else '' + backup_metadata_layout = 'table-layout:fixed;' if context_email else '' section_label = _runtime_text(f'email.groups.{group}', data) report_label = _runtime_text('email.report', data, group=section_label) host_label = _runtime_text('email.host', data) diff --git a/AppImage/scripts/notification_templates.py b/AppImage/scripts/notification_templates.py index 99ae05ff..5cc9f2b6 100644 --- a/AppImage/scripts/notification_templates.py +++ b/AppImage/scripts/notification_templates.py @@ -2070,10 +2070,8 @@ def render_template(event_type: str, data: Dict[str, Any], return '' safe_vars = _SafeDict(variables) - if (event_type == 'error_resolved' and data.get('recovery_outcome') == 'resolved' - and data.get('is_recovery') is True - and isinstance(data.get('check_evidence'), dict) - and data['check_evidence'].get('check') == 'cpu_usage'): + from health_recovery import presents_recovery + if event_type == 'error_resolved' and presents_recovery(data): safe_vars['_health_title'] = runtime_message('healthRecovery.title', language, **variables) safe_vars['_health_body'] = runtime_message('healthRecovery.body', language, **variables) template['title'] = '{_health_title}' @@ -2089,13 +2087,14 @@ def render_template(event_type: str, data: Dict[str, Any], # parse the table/logs and format a rich body instead of the sparse template. pve_message = data.get('pve_message', '') backup_diagnostics = [] + principal_cause = None - def bounded_backup_diagnostics(lines): + def bounded_backup_diagnostics(lines, principal_cause=None): # 1024 chars matches the repository's small-channel message convention; # 8 lines keeps repeated producer warnings readable. Inventory/title # size is separate: this is not a one-Telegram-message guarantee. unique = list(dict.fromkeys(line for line in lines if line.strip())) - principal = next((line for line in unique if re.search(r'\b(?:ERROR:|TASK ERROR:)', line, re.IGNORECASE)), None) + principal = principal_cause or next((line for line in unique if re.search(r'\b(?:ERROR:|TASK ERROR:)', line, re.IGNORECASE)), None) if principal: unique.remove(principal) unique.insert(0, principal) @@ -2189,16 +2188,18 @@ def render_template(event_type: str, data: Dict[str, Any], source_subject = cause.group(1).strip() if cause else '' if source_subject.lower() == 'multiple problems': source_subject = '' - cause_in_diagnostics = any( - line.strip() == source_subject or - re.split(r'\b(?:TASK ERROR:|ERROR:)\s*', line, maxsplit=1, flags=re.IGNORECASE)[-1].strip() == source_subject - for line in backup_diagnostics) - if source_subject and source_subject not in body_text and not cause_in_diagnostics: - # A unique job/setup cause must survive a warning-heavy report. - backup_diagnostics.insert(0, source_subject) + if source_subject and source_subject not in body_text: + # Reserve the subject-equivalent diagnostic BEFORE the cap. Finding + # it in uncapped logs is not enough: that late line could be omitted. + principal_cause = next((line for line in backup_diagnostics + if line.strip() == source_subject or + re.split(r'\b(?:TASK ERROR:|ERROR:)\s*', line, maxsplit=1, flags=re.IGNORECASE)[-1].strip() == source_subject), None) + if not principal_cause: + principal_cause = source_subject + backup_diagnostics.insert(0, source_subject) if backup_diagnostics: - body_text += '\n' + bounded_backup_diagnostics(backup_diagnostics) + body_text += '\n' + bounded_backup_diagnostics(backup_diagnostics, principal_cause) # Clean up: collapse runs of 3+ blank lines into 1, remove trailing whitespace import re as _re @@ -2531,10 +2532,8 @@ 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 == 'error_resolved' and data.get('recovery_outcome') == 'resolved' - and data.get('is_recovery') is True - and isinstance(data.get('check_evidence'), dict) - and data['check_evidence'].get('check') == 'cpu_usage'): + from health_recovery import presents_recovery + if event_type == 'error_resolved' and presents_recovery(data): icon = '✅' if event_type == 'backup_complete': icon = {