refactor: separate monitoring from hardware plugins
This commit is contained in:
@@ -0,0 +1,109 @@
|
||||
"""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/<plugin>/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()
|
||||
Reference in New Issue
Block a user