Fix principal cause caps and bind measured recovery to incident policy

This commit is contained in:
martino
2026-10-01 08:38:32 +02:00
parent 8e396f957f
commit 3be15211b1
9 changed files with 537 additions and 140 deletions
@@ -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)
@@ -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()
@@ -1,100 +1,68 @@
"""Fresh existing-check provenance; extracted consumers, real disposable SQLite.""" """Fresh existing-check provenance; native initializer and disposable SQLite."""
import contextlib
import datetime
import json import json
import os
import sqlite3
import tempfile
import time import time
import types
import unittest import unittest
from unittest.mock import patch 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_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): class RecoveryEvidenceTests(unittest.TestCase):
def test_cpu_success_provenance_is_persisted_only_after_normal_samples(self): def test_cpu_success_provenance_is_persisted_only_after_normal_samples(self):
events=[] with case() as store:
with tempfile.TemporaryDirectory() as scratch: first = original_cpu(store)
db=scratch+'/health.sqlite' self.assertEqual(cpu(store)['status'], 'OK')
conn=sqlite3.connect(db) proof = store.get_recovery_evidence('cpu_usage', first)
conn.execute('CREATE TABLE errors(id INTEGER PRIMARY KEY,error_key TEXT,details TEXT,resolved_at TEXT,resolution_type TEXT,resolution_reason TEXT)') self.assertTrue(proof)
conn.execute("INSERT INTO errors(error_key,details) VALUES ('cpu_usage','{}')") self.assertEqual(proof['check'], 'cpu_usage')
conn.commit();conn.close() self.assertEqual(proof['checked_at'], BASE)
@contextlib.contextmanager # A generic closure never gains proof; use another actual native row.
def connection(): store.record_error('pve_service_test','pve_services','CRITICAL','inactive')
c=sqlite3.connect(db) store.resolve_error('pve_service_test','No longer present')
try:yield c self.assertFalse(json.loads(sql(store,"SELECT data FROM events ORDER BY id DESC LIMIT 1")[0][0]).get('check_evidence'))
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'))
def test_recovery_query_requires_fresh_same_incident_proof(self): def test_recovery_query_requires_fresh_same_incident_proof(self):
with tempfile.TemporaryDirectory() as scratch: with case() as store:
db=scratch+'/health.sqlite' first = original_cpu(store); cpu(store)
conn=sqlite3.connect(db) proof = store.get_recovery_evidence('cpu_usage',first)
conn.execute('CREATE TABLE errors(id INTEGER PRIMARY KEY,error_key TEXT,first_seen TEXT,last_seen TEXT,resolved_at TEXT,acknowledged INTEGER)') self.assertTrue(proof)
conn.execute('CREATE TABLE events(id INTEGER PRIMARY KEY,event_type TEXT,error_key TEXT,timestamp TEXT,data TEXT)') self.assertIsNone(store.get_recovery_evidence('cpu_usage','different incident'))
now=datetime.datetime.now(); first=(now-datetime.timedelta(minutes=10)).isoformat(); last=(now-datetime.timedelta(minutes=1)).isoformat(); resolved=now.isoformat() saved = sql(store,'SELECT last_seen,resolved_at FROM errors')[0]
proof={'check':'cpu_usage','checked_at':now.timestamp()} for field,value in [('acknowledged',1),('resolved_at',None),('last_seen',Clock.fromtimestamp(BASE+1).isoformat())]:
conn.execute('INSERT INTO errors VALUES(1,?,?,?,?,0)',('cpu_usage',first,last,resolved)) sql(store,f'UPDATE errors SET {field}=?',(value,))
conn.execute('INSERT INTO events VALUES(1,?,?,?,?)',('resolved','cpu_usage',resolved,json.dumps({'check_evidence':proof}))) self.assertIsNone(store.get_recovery_evidence('cpu_usage',first))
conn.commit();conn.close() sql(store,'UPDATE errors SET acknowledged=0,last_seen=?,resolved_at=?',saved)
@contextlib.contextmanager data = json.loads(sql(store,"SELECT data FROM events WHERE event_type='resolved'")[0][0])
def connection(**kwargs): 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}):
c=sqlite3.connect(db) sql(store,"UPDATE events SET data=? WHERE event_type='resolved'",(json.dumps({**data,'check_evidence':bad}),))
try:yield c self.assertIsNone(store.get_recovery_evidence('cpu_usage',first))
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))
def test_poller_and_all_consumers_distinguish_proven_recovery_from_disappearance(self): def test_poller_and_all_consumers_distinguish_proven_recovery_from_disappearance(self):
import sys for proved in (False,True):
ns={'time':time,'json':json,'Dict':dict,'NotificationEvent':lambda *a,**kw:types.SimpleNamespace(event_type=a[0],severity=a[1],data=a[2])} with case() as store:
poll=extract(SCRIPTS/'notification_events.py','_check_persistent_health','PollingCollector',ns) first = original_cpu(store)
for proof in (None,{'check':'cpu_usage','checked_at':time.time()}): if proved: cpu(store)
events=[] else: store.resolve_error('cpu_usage','No longer present')
store=types.SimpleNamespace(get_active_errors=lambda:[],is_error_acknowledged=lambda key:False,get_recovery_evidence=lambda *a:proof) data = poll(store,first,reason='CPU high')[0][0]
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) self.assertEqual(data['is_recovery'],proved)
with patch.dict(sys.modules,{'health_persistence':types.SimpleNamespace(health_persistence=store)}):poll(collector) self.assertEqual(data['recovery_outcome'],'resolved' if proved else 'no_longer_reported')
self.assertEqual(len(events),1) with patch('health_recovery.time.time',return_value=BASE):
event=events[0] for lang in LANGUAGES:
self.assertEqual(event.data.get('recovery_outcome'),'resolved' if proof else 'no_longer_reported') for manual in (False,True):
self.assertEqual(event.data['is_recovery'],bool(proof)) result = deliver('error_resolved',data,'OK',lang,manual=manual)
for lang in LANGUAGES: self.assertEqual('background:#f0fdf4;' in result['html'],proved)
for manual in (False,True): if proved:
result=deliver(event.event_type,event.data,event.severity,lang,manual=manual) self.assertIn(templates.runtime_message('healthRecovery.title',lang,hostname='node-a',category='cpu',entity_suffix=''),result['title'])
if proof: self.assertEqual(result['text'].count(data['reason']),1)
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)
def test_manual_recovery_flag_alone_is_not_authoritative_evidence(self): 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'} data={'hostname':'node-a','category':'cpu','reason':'Observation disappeared','duration':'1h','original_severity':'WARNING','recovery_outcome':'resolved'}
@@ -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()
+25 -3
View File
@@ -1427,6 +1427,8 @@ class HealthMonitor:
'details': f'Sustained for {actual_duration}s above {self.CPU_CRITICAL}%.', 'details': f'Sustained for {actual_duration}s above {self.CPU_CRITICAL}%.',
'cpu_percent': cpu_percent, 'cpu_percent': cpu_percent,
'duration': actual_duration, '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: 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}%.', 'details': f'Sustained for {actual_duration}s above {self.CPU_WARNING}%.',
'cpu_percent': cpu_percent, 'cpu_percent': cpu_percent,
'duration': actual_duration, 'duration': actual_duration,
'cpu_policy': {'warning': self.CPU_WARNING, 'critical': self.CPU_CRITICAL,
'recovery': self.CPU_RECOVERY},
}, },
) )
else: else:
@@ -1453,8 +1457,21 @@ class HealthMonitor:
reason = None reason = None
# CPU is normal - auto-resolve any existing CPU errors # CPU is normal - auto-resolve any existing CPU errors
evidence = None evidence = None
if cpu_percent < self.CPU_RECOVERY and len(recovery_samples) >= RECOVERY_MIN_SAMPLES: # Presentation proof is stricter than operational hysteresis:
evidence = {'check': 'cpu_usage', 'checked_at': current_time} # 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', health_persistence.resolve_error('cpu_usage', 'CPU usage returned to normal',
check_evidence=evidence) check_evidence=evidence)
@@ -3904,6 +3921,7 @@ class HealthMonitor:
failed_services = [] failed_services = []
service_details = {} service_details = {}
active_evidence = {}
for service in services_to_check: for service in services_to_check:
try: try:
@@ -3918,6 +3936,10 @@ class HealthMonitor:
if result.returncode != 0 or status != 'active': if result.returncode != 0 or status != 'active':
failed_services.append(service) failed_services.append(service)
service_details[service] = status or 'inactive' 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: except Exception:
failed_services.append(service) failed_services.append(service)
service_details[service] = 'error' service_details[service] = 'error'
@@ -3932,7 +3954,7 @@ class HealthMonitor:
if svc not in failed_services: if svc not in failed_services:
error_key = f'pve_service_{svc}' error_key = f'pve_service_{svc}'
if health_persistence.is_error_active(error_key): 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 # Build checks dict with status per service
checks = {} checks = {}
+49 -28
View File
@@ -600,7 +600,7 @@ class HealthPersistence:
cursor.execute(''' cursor.execute('''
SELECT id, acknowledged, resolved_at, category, severity, first_seen, 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 = ? FROM errors WHERE error_key = ?
''', (error_key,)) ''', (error_key,))
existing = cursor.fetchone() existing = cursor.fetchone()
@@ -609,7 +609,7 @@ class HealthPersistence:
if existing: if existing:
(err_id, ack, resolved_at, old_cat, old_severity, first_seen, (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: if ack == 1:
# SAFETY OVERRIDE: Critical CPU temperature ALWAYS re-triggers # SAFETY OVERRIDE: Critical CPU temperature ALWAYS re-triggers
@@ -680,6 +680,18 @@ class HealthPersistence:
conn.commit() conn.commit()
return event_info 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 # Not acknowledged - update existing active error
cursor.execute(''' cursor.execute('''
UPDATE errors UPDATE errors
@@ -776,7 +788,7 @@ class HealthPersistence:
# was created — otherwise "Storage 'Tuxis' unavailable" # was created — otherwise "Storage 'Tuxis' unavailable"
# comes back as "Resolved - Storage" with no identity. # comes back as "Resolved - Storage" with no identity.
cursor.execute( 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,), (error_key,),
) )
row = cursor.fetchone() row = cursor.fetchone()
@@ -786,11 +798,18 @@ class HealthPersistence:
stored_details = json.loads(row[0]) stored_details = json.loads(row[0])
except Exception: except Exception:
stored_details = None 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, { self._record_event(cursor, 'resolved', error_key, {
'reason': reason, 'reason': reason,
# Only explicit current-check callers attach this proof. # Only explicit current-check callers attach this proof.
# Generic resolve/cleanup remains neutral. # Generic resolve/cleanup remains neutral.
'check_evidence': check_evidence, 'check_evidence': check_evidence,
'incident': {'id': row[1], 'first_seen': row[2], 'resolved_at': row[3]},
'entity': self._entity_from_details(stored_details), 'entity': self._entity_from_details(stored_details),
'details': stored_details or {}, 'details': stored_details or {},
}) })
@@ -800,41 +819,40 @@ class HealthPersistence:
def get_recovery_evidence(self, error_key: str, first_seen: str): def get_recovery_evidence(self, error_key: str, first_seen: str):
"""Return fresh same-incident native check proof, never absence of errors. """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 checks, generic clears, excluded/deleted records and legacy events stay
neutral until they have equivalent per-condition provenance. 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 return None
try: 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(''' row = conn.execute('''
SELECT first_seen, last_seen, resolved_at, acknowledged SELECT e.first_seen, e.last_seen, e.resolved_at, e.acknowledged,
FROM errors WHERE error_key = ? ORDER BY id DESC LIMIT 1 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() ''', (error_key,)).fetchone()
if not row or row[0] != first_seen or not row[2] or row[3]: 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:
return None return None
proof = json.loads(event[1]).get('check_evidence') event_data = json.loads(row[6])
if not isinstance(proof, dict) or proof.get('check') != error_key: if event_data.get('incident') != {'id': row[4], 'first_seen': row[0], 'resolved_at': row[2]}:
return None return None
checked = proof.get('checked_at') proof = event_data.get('check_evidence')
if not isinstance(checked, (int, float)) or isinstance(checked, bool): from health_recovery import valid_check_evidence
if not valid_check_evidence(error_key, proof, now=datetime.now().timestamp()):
return None return None
checked = float(checked) if error_key == 'cpu_usage' and proof.get('policy') != event_data.get('details', {}).get('cpu_policy'):
# Reuse the collector's existing two-hour freshness boundary.
now = datetime.now().timestamp()
if not 0 <= now - checked <= 7200:
return None return None
checked = float(proof['checked_at'])
last_seen = datetime.fromisoformat(row[1]).timestamp() last_seen = datetime.fromisoformat(row[1]).timestamp()
resolved = datetime.fromisoformat(row[2]).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: if not last_seen <= checked <= resolved <= recorded:
return None return None
return proof return proof
@@ -906,7 +924,7 @@ class HealthPersistence:
return False 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. Remove/resolve a specific error immediately.
Used when the condition that caused the error no longer exists Used when the condition that caused the error no longer exists
@@ -928,7 +946,7 @@ class HealthPersistence:
# Check if this error was acknowledged (dismissed) # Check if this error was acknowledged (dismissed)
cursor.execute(''' cursor.execute('''
SELECT acknowledged FROM errors WHERE error_key = ? SELECT acknowledged, id, first_seen FROM errors WHERE error_key = ?
''', (error_key,)) ''', (error_key,))
row = cursor.fetchone() row = cursor.fetchone()
@@ -947,7 +965,10 @@ class HealthPersistence:
''', (now, error_key)) ''', (now, error_key))
if cursor.rowcount > 0: 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() conn.commit()
+51
View File
@@ -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')))
+7 -6
View File
@@ -1038,10 +1038,8 @@ class EmailChannel(NotificationChannel):
# Determine group for section header # Determine group for section header
event_type = data.get('_event_type', '') event_type = data.get('_event_type', '')
if event_type == 'error_resolved': if event_type == 'error_resolved':
if (data.get('recovery_outcome') == 'resolved' from health_recovery import presents_recovery
and data.get('is_recovery') is True if presents_recovery(data):
and isinstance(data.get('check_evidence'), dict)
and data['check_evidence'].get('check') == 'cpu_usage'):
sev.update(self._SEV_STYLE['OK']) sev.update(self._SEV_STYLE['OK'])
sev['label'] = _runtime_notification_text('healthRecovery.status', data) sev['label'] = _runtime_notification_text('healthRecovery.status', data)
else: else:
@@ -1070,8 +1068,11 @@ class EmailChannel(NotificationChannel):
or backup_email or data.get('_restore_summary') or data.get('_backup_summary')) 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_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 '' temp_table_layout = 'table-layout:fixed;' if wrap_body else ''
backup_title_wrap = temp_cell_wrap if backup_email else '' # Recovery exposes the same literal host context as backup notices.
backup_metadata_layout = 'table-layout:fixed;' if backup_email else '' # 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) section_label = _runtime_text(f'email.groups.{group}', data)
report_label = _runtime_text('email.report', data, group=section_label) report_label = _runtime_text('email.report', data, group=section_label)
host_label = _runtime_text('email.host', data) host_label = _runtime_text('email.host', data)
+17 -18
View File
@@ -2070,10 +2070,8 @@ def render_template(event_type: str, data: Dict[str, Any],
return '' return ''
safe_vars = _SafeDict(variables) safe_vars = _SafeDict(variables)
if (event_type == 'error_resolved' and data.get('recovery_outcome') == 'resolved' from health_recovery import presents_recovery
and data.get('is_recovery') is True if event_type == 'error_resolved' and presents_recovery(data):
and isinstance(data.get('check_evidence'), dict)
and data['check_evidence'].get('check') == 'cpu_usage'):
safe_vars['_health_title'] = runtime_message('healthRecovery.title', language, **variables) safe_vars['_health_title'] = runtime_message('healthRecovery.title', language, **variables)
safe_vars['_health_body'] = runtime_message('healthRecovery.body', language, **variables) safe_vars['_health_body'] = runtime_message('healthRecovery.body', language, **variables)
template['title'] = '{_health_title}' 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. # parse the table/logs and format a rich body instead of the sparse template.
pve_message = data.get('pve_message', '') pve_message = data.get('pve_message', '')
backup_diagnostics = [] 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; # 1024 chars matches the repository's small-channel message convention;
# 8 lines keeps repeated producer warnings readable. Inventory/title # 8 lines keeps repeated producer warnings readable. Inventory/title
# size is separate: this is not a one-Telegram-message guarantee. # size is separate: this is not a one-Telegram-message guarantee.
unique = list(dict.fromkeys(line for line in lines if line.strip())) 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: if principal:
unique.remove(principal) unique.remove(principal)
unique.insert(0, 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 '' source_subject = cause.group(1).strip() if cause else ''
if source_subject.lower() == 'multiple problems': if source_subject.lower() == 'multiple problems':
source_subject = '' source_subject = ''
cause_in_diagnostics = any( if source_subject and source_subject not in body_text:
line.strip() == source_subject or # Reserve the subject-equivalent diagnostic BEFORE the cap. Finding
re.split(r'\b(?:TASK ERROR:|ERROR:)\s*', line, maxsplit=1, flags=re.IGNORECASE)[-1].strip() == source_subject # it in uncapped logs is not enough: that late line could be omitted.
for line in backup_diagnostics) principal_cause = next((line for line in backup_diagnostics
if source_subject and source_subject not in body_text and not cause_in_diagnostics: if line.strip() == source_subject or
# A unique job/setup cause must survive a warning-heavy report. re.split(r'\b(?:TASK ERROR:|ERROR:)\s*', line, maxsplit=1, flags=re.IGNORECASE)[-1].strip() == source_subject), None)
backup_diagnostics.insert(0, source_subject) if not principal_cause:
principal_cause = source_subject
backup_diagnostics.insert(0, source_subject)
if backup_diagnostics: 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 # Clean up: collapse runs of 3+ blank lines into 1, remove trailing whitespace
import re as _re import re as _re
@@ -2531,10 +2532,8 @@ def enrich_with_emojis(event_type: str, title: str, body: str,
severity = data.get('severity', 'INFO') severity = data.get('severity', 'INFO')
icon = EVENT_EMOJI.get(event_type) or CATEGORY_EMOJI.get(group) or SEVERITY_ICONS.get(severity, '') 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' from health_recovery import presents_recovery
and data.get('is_recovery') is True if event_type == 'error_resolved' and presents_recovery(data):
and isinstance(data.get('check_evidence'), dict)
and data['check_evidence'].get('check') == 'cpu_usage'):
icon = '✅' icon = '✅'
if event_type == 'backup_complete': if event_type == 'backup_complete':
icon = { icon = {