"""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 now1 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 wanted18: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);last=state;last_time=self.clock() 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)