"""Hardware-independent client for the versioned local plugin API.""" import copy from contextlib import closing import http.client import json import socket import threading import uuid from pathlib import Path PLUGIN_ROOT=Path('/run/pigway-plugins') def socket_path(value): path=Path(value) if not path.is_absolute() or '..' in path.parts or path.name!='api.sock' or PLUGIN_ROOT not in path.parents: raise ValueError('plugin socket must be /run/pigway-plugins//api.sock') return str(path) class Connection(http.client.HTTPConnection): def __init__(self,path): super().__init__('localhost',timeout=.8);self.path=socket_path(path) def connect(self): self.sock=socket.socket(socket.AF_UNIX,socket.SOCK_STREAM) self.sock.settimeout(self.timeout) try:self.sock.connect(self.path) except Exception:self.sock.close();raise def request(path,method,endpoint,payload=None): with closing(Connection(path)) as conn: body=None if payload is None else json.dumps(payload,allow_nan=False).encode() conn.request(method,endpoint,body,{'Content-Type':'application/json'}) response=conn.getresponse();raw=response.read(262145) if len(raw)>262144:raise ValueError('plugin response too large') value=json.loads(raw) if response.status>=400:raise ValueError(value.get('error','plugin request failed')) return value def registered(cfg): if not cfg.has_section('plugins'):return {} return {k:socket_path(v) for k,v in cfg.items('plugins') if v.strip()} def discover(cfg): linked=registered(cfg);paths={v for v in linked.values()} paths.update(str(p) for p in PLUGIN_ROOT.glob('*/api.sock')) result=[] for path in sorted(paths)[:32]: key=next((k for k,v in linked.items() if v==path),None) try: d=request(path,'GET','/v1/descriptor') if d.get('api_version')!=1:raise ValueError('incompatible plugin API version') status=request(path,'GET','/v1/status') result.append({'id':key or d['id'],'socket':path,'linked':key is not None,'online':True,'descriptor':d,'status':status}) except Exception as exc: result.append({'id':key or Path(path).parent.name,'socket':path,'linked':key is not None,'online':False,'error':str(exc)}) return result class Publisher: """One replaceable snapshot, not an unbounded queue of animation commands.""" def __init__(self,plugins,source,log): self.plugins=plugins;self.source=source;self.log=log self.session=str(uuid.uuid4());self.revision=0;self.snapshot=None self.lock=threading.Lock();self.stop_event=threading.Event();self.status={} self.threads=[threading.Thread(target=self.run,args=(key,path),name='plugin-'+key,daemon=True) for key,path in plugins.items()] for thread in self.threads:thread.start() def publish(self,state): with self.lock:self.snapshot=copy.deepcopy(state) def current_status(self): with self.lock:return copy.deepcopy(self.status) def run(self,key,path): previous_error=None;revision=0;descriptor=None while not self.stop_event.is_set(): with self.lock:state=self.snapshot if state is not None: revision+=1 payload={'api_version':1,'source':self.source,'session':self.session,'revision':revision, 'ttl_seconds':10,'system':state['system'],'network':state['network'],'alerts':state['display']['alerts']} try: if descriptor is None: descriptor=request(path,'GET','/v1/descriptor') if descriptor.get('api_version')!=1:raise ValueError('incompatible plugin API version') request(path,'PUT','/v1/state',payload) status=request(path,'GET','/v1/status') with self.lock:self.status[key]={'online':True,'status':status,'name':descriptor.get('name',key),'capabilities':descriptor.get('capabilities',{})} if previous_error is not None:self.log('INFO','PLUGIN_RECOVERED',plugin=key) previous_error=None except Exception as exc: descriptor=None error=str(exc) with self.lock:self.status[key]={'online':False,'error':error} if previous_error!=error:self.log('WARN','PLUGIN_UNAVAILABLE',plugin=key,error=error) previous_error=error self.stop_event.wait(1) try:request(path,'DELETE','/v1/state',{'source':self.source,'session':self.session}) except Exception:pass # Server-side leases expire if transport is unavailable. def close(self): self.stop_event.set() for thread in self.threads:thread.join()