#!/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