177 lines
9.6 KiB
Python
177 lines
9.6 KiB
Python
#!/usr/bin/env python3
|
|
"""Standalone fleet manager. No local sensors, systemd or hardware dependencies."""
|
|
import concurrent.futures
|
|
import json
|
|
import os
|
|
import re
|
|
import secrets
|
|
import threading
|
|
import time
|
|
from http.server import BaseHTTPRequestHandler,ThreadingHTTPServer
|
|
from pathlib import Path
|
|
from urllib.parse import urlparse,parse_qs,urlencode
|
|
from urllib.request import Request,build_opener,HTTPRedirectHandler
|
|
from urllib.error import HTTPError
|
|
|
|
DATA=Path(os.environ.get('PIGWAY_DATA','/data'))
|
|
WEB=Path(__file__).resolve().parent.parent/'web/index.html'
|
|
LOCK=threading.RLock()
|
|
MAX_HOSTS=32
|
|
|
|
|
|
def initialize():
|
|
DATA.mkdir(parents=True,exist_ok=True,mode=0o700)
|
|
token=DATA/'manager.token'
|
|
if not token.exists():
|
|
fd=os.open(token,os.O_WRONLY|os.O_CREAT|os.O_EXCL,0o600)
|
|
with os.fdopen(fd,'w') as f:f.write(secrets.token_urlsafe(32)+'\n')
|
|
return token.read_text().strip()
|
|
|
|
|
|
def registry():
|
|
with LOCK:
|
|
path=DATA/'hosts.json'
|
|
return json.loads(path.read_text()) if path.exists() else []
|
|
|
|
|
|
def public(host):return {k:v for k,v in host.items() if k!='token'}
|
|
|
|
|
|
def host_by_id(identifier):
|
|
return next((h for h in registry() if h['id']==identifier),None) or fail('unknown host')
|
|
|
|
|
|
def fail(message):raise ValueError(message)
|
|
|
|
|
|
class NoRedirect(HTTPRedirectHandler):
|
|
def redirect_request(self,*args,**kwargs):return None
|
|
|
|
|
|
def remote(host,path,method='GET',payload=None):
|
|
body=None if payload is None else json.dumps(payload).encode()
|
|
req=Request(host['url']+path,data=body,method=method,headers={'Authorization':'Bearer '+host['token'],'Content-Type':'application/json'})
|
|
try:
|
|
with build_opener(NoRedirect).open(req,timeout=5) as response:
|
|
raw=response.read(16*1024*1024+1)
|
|
if len(raw)>16*1024*1024:raise ValueError('device response exceeds size limit')
|
|
return json.loads(raw)
|
|
except HTTPError as exc:raise ValueError('device API returned HTTP '+str(exc.code)) from None
|
|
|
|
|
|
def edit_host(body):
|
|
identifier=body.get('id','')
|
|
if not isinstance(identifier,str) or not re.fullmatch(r'[A-Za-z0-9_.-]{1,64}',identifier):fail('invalid host identifier')
|
|
with LOCK:
|
|
hosts=registry();old=next((h for h in hosts if h['id']==identifier),None)
|
|
operation=body.get('operation')
|
|
if operation=='save':
|
|
url=str(body.get('url','')).rstrip('/');parsed=urlparse(url)
|
|
if parsed.scheme not in ('http','https') or not parsed.hostname or parsed.username or parsed.password or parsed.path or parsed.query or parsed.fragment:fail('use an HTTP(S) origin without credentials or path')
|
|
token=body.get('token') or (old or {}).get('token')
|
|
if not isinstance(token,str) or not token or len(token)>1024 or '\n' in token or '\r' in token:fail('device token required')
|
|
if not old and len(hosts)>=MAX_HOSTS:fail('at most 32 devices')
|
|
entry={'id':identifier,'name':identifier,'url':url,'token':token,'local':False}
|
|
hosts=[entry if h['id']==identifier else h for h in hosts] if old else hosts+[entry]
|
|
elif operation=='delete':hosts=[h for h in hosts if h['id']!=identifier]
|
|
else:fail('invalid operation')
|
|
path=DATA/'hosts.tmp'
|
|
with path.open('w') as f:json.dump(hosts,f)
|
|
path.chmod(0o600);path.replace(DATA/'hosts.json')
|
|
return {'hosts':[public(h) for h in hosts]}
|
|
|
|
|
|
def parallel(hosts,fn):
|
|
with concurrent.futures.ThreadPoolExecutor(max_workers=min(8,max(1,len(hosts)))) as pool:
|
|
return list(pool.map(fn,hosts))
|
|
|
|
|
|
def fleet_status():
|
|
def fetch(host):
|
|
try:return {'host':public(host),'online':True,'status':remote(host,'/api/v1/status')}
|
|
except Exception:return {'host':public(host),'online':False,'error':'Device unavailable; check address, token and API status'}
|
|
return {'hosts':parallel(registry(),fetch),'timestamp':int(time.time())}
|
|
|
|
|
|
def fleet_logs(query):
|
|
hosts=registry();selected=query.get('machine','')
|
|
if selected:hosts=[h for h in hosts if h['id']==selected]
|
|
params={k:v for k,v in query.items() if k in ('since','until','severity','event','search')}
|
|
params.update(limit=1000,sort='timestamp',order='desc')
|
|
def fetch(host):
|
|
try:
|
|
rows=[];offset=0;events=set()
|
|
while True:
|
|
value=remote(host,'/api/v1/logs?'+urlencode({**params,'offset':offset}))
|
|
batch=value['logs'];rows.extend({**r,'machine_id':host['id'],'machine_name':host['name']} for r in batch)
|
|
events.update(value.get('event_types',[]));offset+=len(batch)
|
|
if offset>=value['total'] or not batch:break
|
|
if offset>=50000:return rows,events,{'id':host['id'],'error':'Device log query exceeds 50000 rows; narrow date range'}
|
|
return rows,events,None
|
|
except Exception:return [],set(),{'id':host['id'],'error':'Device log query failed'}
|
|
results=parallel(hosts,fetch);rows=[row for result in results for row in result[0]]
|
|
ranks={'EMERGENCY':0,'ALERT':1,'CRITICAL':2,'ERROR':3,'WARN':4,'NOTICE':5,'INFO':6,'DEBUG':7}
|
|
sort=query.get('sort','timestamp');order=query.get('order','desc')
|
|
if sort not in ('timestamp','machine','severity','event','message') or order not in ('asc','desc'):fail('invalid sorting')
|
|
rows.sort(key=lambda r: ranks.get(r['severity'],99) if sort=='severity' else r['machine_name'] if sort=='machine' else r[sort],reverse=order=='desc')
|
|
offset=max(0,int(query.get('offset',0)));limit=min(1000,max(1,int(query.get('limit',100))))
|
|
return {'logs':rows[offset:offset+limit],'total':len(rows),'offset':offset,'limit':limit,'event_types':sorted(set().union(*(r[1] for r in results))),
|
|
'machines':[{'id':h['id'],'name':h['name']} for h in registry()],'errors':[r[2] for r in results if r[2]],'retention':'device-systemd-journal'}
|
|
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
def setup(self):super().setup();self.connection.settimeout(15)
|
|
def log_message(self,*args):pass
|
|
def send(self,code,data,html=False):
|
|
raw=data if html else json.dumps(data,ensure_ascii=False).encode()
|
|
self.send_response(code);self.send_header('Content-Type','text/html; charset=utf-8' if html else 'application/json')
|
|
self.send_header('Content-Length',str(len(raw)));self.send_header('Cache-Control','no-store');self.send_header('X-Content-Type-Options','nosniff');self.send_header('X-Frame-Options','DENY')
|
|
self.send_header('Content-Security-Policy',"default-src 'self'; style-src 'self' 'unsafe-inline'; script-src 'self' 'unsafe-inline'")
|
|
self.end_headers();self.wfile.write(raw)
|
|
def authorized(self):
|
|
return secrets.compare_digest(self.headers.get('Authorization',''),'Bearer '+self.server.token)
|
|
def handle_request(self,method):
|
|
parsed=urlparse(self.path);path=parsed.path;q={k:v[0] for k,v in parse_qs(parsed.query).items()}
|
|
if method=='GET' and path=='/':self.send(200,WEB.read_bytes(),True);return
|
|
if not self.authorized():self.send(401,{'error':'Manager token required'});return
|
|
try:
|
|
if method=='GET':
|
|
if path=='/api/v1/hosts':result={'hosts':[public(h) for h in registry()]}
|
|
elif path=='/api/v1/fleet/status':result=fleet_status()
|
|
elif path=='/api/v1/fleet/logs':result=fleet_logs(q)
|
|
elif path in ('/api/v1/config','/api/v1/services','/api/v1/plugins','/api/v1/plugin/config'):
|
|
host=host_by_id(q.get('host',''));result=remote(host,path+('?' +urlencode({'id':q['id']}) if 'id' in q else ''))
|
|
else:self.send(404,{'error':'not found'});return
|
|
else:
|
|
length=int(self.headers.get('Content-Length',0))
|
|
if not 0<length<=65536:fail('invalid body length')
|
|
body=json.loads(self.rfile.read(length))
|
|
if not isinstance(body,dict):fail('object required')
|
|
if path=='/api/v1/hosts':result=edit_host(body)
|
|
elif path=='/api/v1/hosts/check':result=remote(host_by_id(body.get('id')),'/api/v1/health')
|
|
elif path=='/api/v1/fleet/config':
|
|
targets=body.get('targets',[])
|
|
if not isinstance(targets,list) or not 1<=len(targets)<=32:fail('select 1..32 devices')
|
|
results=[]
|
|
for identifier in dict.fromkeys(targets):
|
|
try:remote(host_by_id(identifier),'/api/v1/config','PUT',{'updates':body.get('updates')});results.append({'id':identifier,'ok':True})
|
|
except Exception:results.append({'id':identifier,'ok':False,'error':'Device update failed; check connection and configuration'})
|
|
result={'results':results}
|
|
elif path in ('/api/v1/services','/api/v1/services/control','/api/v1/plugins','/api/v1/plugin/config'):
|
|
result=remote(host_by_id(body.get('host','')),path,'PUT',{k:v for k,v in body.items() if k not in ('host','target_token')})
|
|
else:self.send(404,{'error':'not found'});return
|
|
self.send(200,result)
|
|
except (ValueError,TypeError,KeyError) as exc:self.send(400,{'error':str(exc)})
|
|
except Exception:self.send(502,{'error':'Device unavailable or request failed'})
|
|
def do_GET(self):self.handle_request('GET')
|
|
def do_PUT(self):self.handle_request('PUT')
|
|
|
|
|
|
def main():
|
|
token=initialize();server=ThreadingHTTPServer((os.environ.get('PIGWAY_BIND','0.0.0.0'),int(os.environ.get('PIGWAY_PORT','6001'))),Handler);server.token=token
|
|
print('Manager ready. Read access token from '+str(DATA/'manager.token'),flush=True)
|
|
try:server.serve_forever()
|
|
finally:server.server_close()
|
|
|
|
if __name__=='__main__':main()
|