256 lines
14 KiB
Python
256 lines
14 KiB
Python
"""Local display arbitration and autonomous cooling, independent of the monitor."""
|
|
import copy
|
|
import math
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
from settings import COLORS,EFFECTS,GROUPS
|
|
from local_state import LocalState
|
|
|
|
|
|
def log(event,**fields):
|
|
print('event='+event+' '+ ' '.join(f'{k}={str(v).replace(chr(32),"_")}' for k,v in fields.items()),flush=True)
|
|
|
|
|
|
def temperature():
|
|
return float(Path('/sys/class/thermal/thermal_zone0/temp').read_text())/1000
|
|
|
|
|
|
def finite(value,low,high):
|
|
if isinstance(value,bool) or not isinstance(value,(int,float)) or not math.isfinite(value) or not low<=value<=high:
|
|
raise ValueError('invalid numeric state')
|
|
return value
|
|
|
|
|
|
def validate_state(payload):
|
|
if payload.get('api_version')!=1:raise ValueError('unsupported API version')
|
|
for key in ('source','session'):
|
|
if not isinstance(payload.get(key),str) or not 1<=len(payload[key])<=128:raise ValueError('invalid '+key)
|
|
revision=payload.get('revision')
|
|
if type(revision) is not int or revision<1:raise ValueError('invalid revision')
|
|
ttl=finite(payload.get('ttl_seconds',10),3,30)
|
|
alerts=payload.get('alerts',[])
|
|
if not isinstance(alerts,list) or len(alerts)>128:raise ValueError('too many alerts')
|
|
clean=[];ids=set()
|
|
for item in alerts:
|
|
if not isinstance(item,dict):raise ValueError('invalid alert')
|
|
a={}
|
|
for k in ('id','title','l2','l3','l4'):
|
|
v=item.get(k,'')
|
|
if not isinstance(v,str) or len(v)>128:raise ValueError('invalid alert text')
|
|
a[k]=v
|
|
if not a['id'] or a['id'] in ids:raise ValueError('duplicate or empty alert id')
|
|
ids.add(a['id'])
|
|
a['severity']=finite(item.get('severity'),1,3)
|
|
if int(a['severity'])!=a['severity']:raise ValueError('invalid severity')
|
|
a['priority']=finite(item.get('priority'),1,1000)
|
|
a['category']=item.get('category')
|
|
if a['category'] not in GROUPS:raise ValueError('invalid alert category')
|
|
clean.append(a)
|
|
system=payload.get('system',{});network=payload.get('network',{})
|
|
if not isinstance(system,dict) or not isinstance(network,dict):raise ValueError('invalid home state')
|
|
system={k:finite(system.get(k,0),-50 if k=='temperature_c' else 0,150 if k=='temperature_c' else 100)
|
|
for k in ('cpu_percent','temperature_c','memory_percent','disk_percent')}
|
|
network={k:str(network.get(k,''))[:128] for k in ('kind','metric','ip_label','ip')}
|
|
return {"source":payload['source'],"session":payload['session'],"revision":revision,"ttl_seconds":ttl,
|
|
"alerts":clean,"system":system,"network":network}
|
|
|
|
|
|
class Controller:
|
|
def __init__(self,config,driver=None,clock=time.monotonic,temp_reader=temperature):
|
|
self.config=config;self.driver=driver;self.clock=clock;self.temp_reader=temp_reader
|
|
self.lock=threading.RLock();self.stop_event=threading.Event();self.threads=[]
|
|
self.local_reader=LocalState();self.local_state={'system':{},'network':{}};self.local_sampled_at=None
|
|
self.state=None;self.expires=0;self.ever_connected=False;self.first_seen={}
|
|
self.owner='NORMAL_HOME';self.page='NORMAL_HOME';self.page_since=clock()
|
|
self.fan_level=0;self.temp=None;self.errors={};self.i2c_errors=0;self.rgb_mode='OFF'
|
|
|
|
def accept(self,payload):
|
|
value=validate_state(payload);now=self.clock()
|
|
with self.lock:
|
|
old=self.state
|
|
if old and now<self.expires:
|
|
if (old['source'],old['session'])!=(value['source'],value['session']):raise ValueError('another source holds the display lease')
|
|
if value['revision']<=old['revision']:raise ValueError('stale revision')
|
|
ids={a['id'] for a in value['alerts']}
|
|
self.first_seen={i:self.first_seen.get(i,now) for i in ids}
|
|
self.state=value;self.expires=now+value['ttl_seconds'];self.ever_connected=True
|
|
self._view_locked()
|
|
|
|
def release(self,source=None,session=None):
|
|
with self.lock:
|
|
if source is not None and self.state and (source,session)!=(self.state['source'],self.state['session']):
|
|
raise ValueError('display lease belongs to another source')
|
|
self.state=None;self.expires=0;self.first_seen={};self.ever_connected=False
|
|
self._view_locked()
|
|
|
|
def _view_locked(self):
|
|
now=self.clock();connected=self.state is not None and now<self.expires
|
|
data=self.state if connected else {'alerts':[],'system':{},'network':{}}
|
|
alerts=sorted(data['alerts'],key=lambda a:(-a['priority'],self.first_seen.get(a['id'],now),a['id']))
|
|
byid={a['id']:a for a in alerts};owner=alerts[0]['id'] if alerts else 'NORMAL_HOME'
|
|
if self.owner in byid and byid[self.owner]['priority']==byid[owner]['priority']:owner=self.owner
|
|
# External lease expiry clears only external alerts, never the local home.
|
|
if owner!=self.owner:
|
|
log('DISPLAY_OWNER',old=self.owner,new=owner)
|
|
self.owner=owner;self.page=owner;self.page_since=now
|
|
pages=['NORMAL_HOME']+[a['id'] for a in alerts] if connected else [owner]
|
|
if self.page not in pages:self.page=owner;self.page_since=now
|
|
if len(pages)>1 and now-self.page_since>=self.config['oled']['page_seconds']:
|
|
self.page=pages[(pages.index(self.page)+1)%len(pages)];self.page_since=now
|
|
rgb=byid.get(self.page) or byid.get(owner)
|
|
return {'owner':owner,'page':self.page,'connected':connected,'alerts':alerts,'rgb':rgb,
|
|
'system':self.local_state['system'],'network':self.local_state['network'],'config':self.config}
|
|
|
|
def view(self):
|
|
with self.lock:return copy.deepcopy(self._view_locked())
|
|
|
|
def status(self):
|
|
v=self.view()
|
|
with self.lock:
|
|
return {'available':self.driver is not None,'connected':v['connected'],'owner':v['owner'],
|
|
'oled_page':v['page'],'rgb_mode':self.rgb_mode,'fan_level':self.fan_level,
|
|
'fan_name':('OFF','L1','L2','L3','L4','MAX')[self.fan_level] if self.driver else 'UNAVAILABLE',
|
|
'temperature_c':self.temp,'local_state':copy.deepcopy(self.local_state),
|
|
'local_sample_age_seconds':None if self.local_sampled_at is None else max(0,self.clock()-self.local_sampled_at),'errors':dict(self.errors),'i2c_errors':self.i2c_errors}
|
|
|
|
def error(self,worker,exc):
|
|
with self.lock:
|
|
value=str(exc)
|
|
if self.errors.get(worker)!=value:log('HARDWARE_ERROR',worker=worker,error=value)
|
|
self.errors[worker]=value
|
|
|
|
def recovered(self,worker):
|
|
with self.lock:
|
|
if self.errors.pop(worker,None) is not None:log('HARDWARE_RECOVERED',worker=worker)
|
|
|
|
def start(self):
|
|
if self.driver is None:return
|
|
self.driver.off();self.driver.oled_init()
|
|
for name,work in (('telemetry',self.telemetry_loop),('fan',self.fan_loop),('oled',self.oled_loop),('rgb',self.rgb_loop)):
|
|
thread=threading.Thread(target=work,name=name,daemon=True);self.threads.append(thread);thread.start()
|
|
|
|
def telemetry_loop(self):
|
|
while not self.stop_event.is_set():
|
|
try:
|
|
state,errors=self.local_reader.sample()
|
|
with self.lock:
|
|
self.local_state=state;self.local_sampled_at=self.clock()
|
|
for key in ('cpu_percent','memory_percent','disk_percent','temperature_c','network'):
|
|
if key in errors:self.error('telemetry_'+key,errors[key])
|
|
else:self.recovered('telemetry_'+key)
|
|
self.recovered('telemetry')
|
|
except Exception as exc:self.error('telemetry',exc)
|
|
self.stop_event.wait(1)
|
|
|
|
def fan_loop(self):
|
|
initialized=False
|
|
while not self.stop_event.is_set():
|
|
try:
|
|
temp=self.temp_reader()
|
|
if not math.isfinite(temp):raise ValueError('temperature unavailable')
|
|
self.recovered('temperature')
|
|
except Exception as exc:
|
|
self.error('temperature',exc);temp=None
|
|
with self.lock:
|
|
c=dict(self.config['fan']);old=self.fan_level
|
|
curve=[c[k] for k in ('start','level2','level3','level4','max')]
|
|
wanted=5 if temp is None else sum(temp>=x for x in curve)
|
|
if initialized and wanted<old:
|
|
wanted=old-1 if temp<curve[max(0,old-1)]-c['hysteresis'] else old
|
|
try:
|
|
self.driver.fan(wanted)
|
|
with self.lock:self.fan_level=wanted;self.temp=temp
|
|
if not initialized or wanted!=old:log('FAN_LEVEL',old=old,new=wanted)
|
|
initialized=True;self.recovered('fan')
|
|
except Exception as exc:self.error('fan',exc)
|
|
self.stop_event.wait(1)
|
|
|
|
def oled_loop(self):
|
|
last=None;last_time=0;last_alert_ids=None
|
|
while not self.stop_event.is_set():
|
|
try:
|
|
v=self.view();page=v['page'];system=v['system'];network=v['network'];home=False
|
|
alert_ids=tuple(a['id'] for a in v['alerts'])
|
|
if page=='NORMAL_HOME':
|
|
home=True
|
|
ip=network.get('ip','NO IP');label=network.get('ip_label','IP4')
|
|
if label=='IP6' and len(ip)>18:ip=ip[:7]+'..'+ip[-7:]
|
|
metric=network.get('metric','--');kind=network.get('kind','NET')
|
|
try:signal=max(0,min(100,round((float(metric)+100)*2)))
|
|
except (ValueError,TypeError):signal=None
|
|
def value(key):
|
|
n=system.get(key)
|
|
return f'{n:.1f}' if isinstance(n,(float,int)) and math.isfinite(n) else '--'
|
|
lines=[('CPU '+value('cpu_percent')+'%','MEM '+value('memory_percent')+'%'),
|
|
('TMP '+value('temperature_c')+'C','FAN '+('OFF','L1','L2','L3','L4','MAX')[self.fan_level]),
|
|
('DSK '+value('disk_percent')+'%',f'WIF {signal}%' if kind=='WIF' and signal is not None else f'{kind} {metric}'),(label+' '+ip,'')]
|
|
else:
|
|
a=next(a for a in v['alerts'] if a['id']==page);index=v['alerts'].index(a)+1
|
|
lines=[(a['title'],f"{index}/{len(v['alerts'])}"),(a['l2'],''),(a['l3'],''),(a['l4'],'')]
|
|
state=(lines,home)
|
|
if state!=last or self.clock()-last_time>=v['config']['oled']['refresh_seconds']:
|
|
self.driver.oled(lines,home)
|
|
# Alert recovery changes only a few header pixels (for
|
|
# example 1/2 -> 1/1). Confirm that rare topology change
|
|
# with a second complete frame: during undervoltage an
|
|
# OLED transfer can be only partly applied without an
|
|
# I2C exception being reported by the controller.
|
|
if last_alert_ids is not None and alert_ids!=last_alert_ids:
|
|
self.stop_event.wait(.05)
|
|
if not self.stop_event.is_set():self.driver.oled(lines,home)
|
|
last=state
|
|
last_time=self.clock()
|
|
last_alert_ids=alert_ids
|
|
self.recovered('oled')
|
|
except Exception as exc:self.error('oled',exc)
|
|
self.stop_event.wait(.1)
|
|
|
|
@staticmethod
|
|
def profile(view):
|
|
c=view['config']['rgb'];a=view['rgb']
|
|
if a:
|
|
group=a['category'];level=('warning','critical','emergency')[int(a['severity'])-1]
|
|
label=f"{group}_{level}_{c['alert_mode']}".upper()
|
|
if c['alert_mode']=='breathe':return ('effect',1,c[level+'_speed'],COLORS[c.get(group+'_breathe_color',GROUPS[group])],c['write_delay_ms']),label
|
|
mode='flash';prefix=level
|
|
else:
|
|
group='normal';prefix='normal';mode=c['normal_mode'];label='NORMAL_'+mode.upper()
|
|
if mode=='off':return ('off',),'OFF'
|
|
if mode in EFFECTS:return ('effect',EFFECTS[mode],c['normal_effect_speed'],COLORS[c['normal_effect_color']],c['write_delay_ms']),label
|
|
channels=tuple(c[group+'_'+k] for k in 'rgb')
|
|
if mode=='solid':return ('solid',channels),label
|
|
return ('flash',channels,max(c['custom_min_hold_ms'],c[prefix+'_flash_on_ms'])/1000,
|
|
max(c['custom_min_hold_ms'],c[prefix+'_flash_off_ms'])/1000),label
|
|
|
|
def rgb_loop(self):
|
|
applied=None;next_edge=0;on=True
|
|
while not self.stop_event.is_set():
|
|
try:
|
|
wanted,label=self.profile(self.view())
|
|
if wanted!=applied:
|
|
if wanted[0]=='off':self.driver.off()
|
|
elif wanted[0]=='effect':self.driver.effect(*wanted[1:])
|
|
else:self.driver.color(wanted[1])
|
|
applied=wanted;on=True
|
|
if wanted[0]=='flash':next_edge=self.clock()+wanted[2]
|
|
log('RGB_MODE',old=self.rgb_mode,new=label)
|
|
self.rgb_mode=label
|
|
if applied[0]=='flash' and self.clock()>=next_edge:
|
|
self.driver.color((0,0,0) if on else applied[1]);on=not on
|
|
next_edge=self.clock()+applied[2 if on else 3]
|
|
self.recovered('rgb')
|
|
except Exception as exc:
|
|
self.error('rgb',exc)
|
|
# Do not spin writes at 20Hz on an I2C error.
|
|
self.stop_event.wait(1)
|
|
self.stop_event.wait(.05)
|
|
|
|
def stop(self):
|
|
self.stop_event.set()
|
|
for thread in self.threads:thread.join()
|
|
if self.driver:
|
|
try:self.driver.fan(5)
|
|
except Exception as exc:self.error('shutdown_fan',exc)
|
|
for name,error in self.driver.close():self.error(name,error)
|