feat: add multi-host monitoring and configuration
This commit is contained in:
+193
-12
@@ -1,16 +1,18 @@
|
||||
#!/usr/bin/env python3
|
||||
import configparser,ipaddress,json,re,secrets,shutil,subprocess,threading,time
|
||||
import configparser,ipaddress,json,re,secrets,shutil,subprocess,threading,time,socket
|
||||
from concurrent.futures import ThreadPoolExecutor,as_completed
|
||||
from http import HTTPStatus
|
||||
from http.server import BaseHTTPRequestHandler,ThreadingHTTPServer
|
||||
from pathlib import Path
|
||||
from urllib.parse import parse_qs,urlparse
|
||||
from urllib.parse import parse_qs,urlencode,urlparse
|
||||
from urllib.request import Request,urlopen
|
||||
|
||||
APP_VERSION="3.6.0"
|
||||
APP_VERSION="3.7.0"
|
||||
CFG=Path("/etc/pigway-pi-control.conf")
|
||||
STATUS=Path("/run/pigway-pi-control/status.json")
|
||||
TOKEN=Path("/etc/pigway-pi-control-api.token")
|
||||
WEB=Path("/usr/local/share/pigway-pi-control/web/index.html")
|
||||
ALLOWED_SECTIONS={"services","processes","alerts","fan","led","timing","dark_mode","display","api"}
|
||||
ALLOWED_SECTIONS={"device","hosts","services","processes","alerts","fan","led","timing","dark_mode","display","api"}
|
||||
NAME_RE=re.compile(r"^[A-Za-z0-9_.-]{1,64}$")
|
||||
UNIT_RE=re.compile(r"^[A-Za-z0-9_.@:-]{1,128}\.service$")
|
||||
RGB_KEY_RE=re.compile(r"^(cpu|power|memory|storage|network|service|normal)_[rgb]$")
|
||||
@@ -19,6 +21,67 @@ NORMAL_MODES={"off","solid","flash","flow","breathe","marquee","rainbow","colorf
|
||||
EFFECT_COLORS={"red","green","blue","yellow","purple","cyan","white"}
|
||||
PROTECTED_UNITS={"pigway-pi-control-api.service"}
|
||||
|
||||
def machine_identity():
|
||||
cfg=load_cfg(); identifier=cfg.get("device","identifier",fallback="auto").strip()
|
||||
if not identifier or identifier.lower()=="auto": identifier=socket.gethostname()
|
||||
name=cfg.get("device","name",fallback=identifier).strip() or identifier
|
||||
if name.lower()=="auto": name=identifier
|
||||
return {"id":identifier,"name":name,"hostname":socket.gethostname(),"local":True,"url":""}
|
||||
|
||||
def normalize_url(value):
|
||||
parsed=urlparse(str(value).strip())
|
||||
if parsed.scheme not in {"http","https"} or not parsed.hostname or parsed.username or parsed.password:
|
||||
raise ValueError("host URL must be an http(s) URL without credentials")
|
||||
return f"{parsed.scheme}://{parsed.netloc}".rstrip("/")
|
||||
|
||||
def host_registry():
|
||||
local=machine_identity(); result=[local]
|
||||
cfg=load_cfg()
|
||||
if cfg.has_section("hosts"):
|
||||
for host_id,url in cfg.items("hosts"):
|
||||
try: base=normalize_url(url)
|
||||
except ValueError: continue
|
||||
if host_id==local["id"]: continue
|
||||
result.append({"id":host_id,"name":host_id,"hostname":"","local":False,"url":base})
|
||||
return result
|
||||
|
||||
def find_host(host_id):
|
||||
for host in host_registry():
|
||||
if host["id"]==host_id:return host
|
||||
raise ValueError("unknown host")
|
||||
|
||||
def remote_json(base,path,method="GET",payload=None,token=""):
|
||||
data=None if payload is None else json.dumps(payload,separators=(",",":")).encode()
|
||||
headers={"Accept":"application/json"}
|
||||
if data is not None: headers["Content-Type"]="application/json"
|
||||
if token: headers["Authorization"]=f"Bearer {token}"
|
||||
request=Request(base+path,data=data,headers=headers,method=method)
|
||||
with urlopen(request,timeout=5) as response:
|
||||
return json.loads(response.read())
|
||||
|
||||
def fleet_status():
|
||||
hosts=host_registry(); results=[]
|
||||
try:
|
||||
status=json.loads(STATUS.read_text())
|
||||
status.setdefault("machine",machine_identity())
|
||||
results.append({"host":hosts[0],"online":True,"status":status})
|
||||
except Exception as exc:
|
||||
results.append({"host":hosts[0],"online":False,"error":str(exc)})
|
||||
def fetch(host):
|
||||
try:
|
||||
status=remote_json(host["url"],"/api/v1/status")
|
||||
machine=status.get("machine",{})
|
||||
merged={**host,"name":machine.get("name") or host["name"],"hostname":machine.get("hostname","")}
|
||||
return {"host":merged,"online":True,"status":status}
|
||||
except Exception as exc:return {"host":host,"online":False,"error":str(exc)}
|
||||
peers=hosts[1:]
|
||||
with ThreadPoolExecutor(max_workers=min(8,max(1,len(peers)))) as pool:
|
||||
futures=[pool.submit(fetch,h) for h in peers]
|
||||
for future in as_completed(futures): results.append(future.result())
|
||||
order={h["id"]:i for i,h in enumerate(hosts)}
|
||||
results.sort(key=lambda x:order.get(x["host"]["id"],999))
|
||||
return {"hosts":results,"timestamp":int(time.time())}
|
||||
|
||||
def load_cfg():
|
||||
cfg=configparser.ConfigParser(); cfg.read(CFG)
|
||||
return cfg
|
||||
@@ -32,7 +95,7 @@ def config_json():
|
||||
return {section:dict(cfg.items(section)) for section in cfg.sections() if section in ALLOWED_SECTIONS}
|
||||
|
||||
def validate_value(section,key,value,existing):
|
||||
if section not in {"services","processes"} and key not in existing.get(section,set()):
|
||||
if section not in {"services","processes","hosts"} and key not in existing.get(section,set()):
|
||||
raise ValueError(f"unsupported option: {section}.{key}")
|
||||
if RGB_KEY_RE.fullmatch(key):
|
||||
number=int(value)
|
||||
@@ -54,6 +117,42 @@ def validate_value(section,key,value,existing):
|
||||
if not 1<=number<=1000: raise ValueError("api.log_limit must be 1..1000")
|
||||
elif section=="api" and key=="bind":
|
||||
ipaddress.ip_address(value)
|
||||
elif section=="device" and key=="identifier" and value.lower()!="auto" and not NAME_RE.fullmatch(value):
|
||||
raise ValueError("device.identifier must be auto or 1..64 letters, numbers, dot, underscore or dash")
|
||||
elif section=="hosts":
|
||||
normalize_url(value)
|
||||
|
||||
def replace_ini_value(section,key,value):
|
||||
text=CFG.read_text()
|
||||
section_match=re.search(rf"(?mi)^\[{re.escape(section)}\]\s*$",text)
|
||||
if not section_match:
|
||||
if value is not None: text=text.rstrip()+f"\n\n[{section}]\n{key} = {value}\n"
|
||||
return text
|
||||
next_section=re.search(r"(?m)^\[.+\]\s*$",text[section_match.end():])
|
||||
end=section_match.end()+(next_section.start() if next_section else len(text)-section_match.end())
|
||||
chunk=text[section_match.end():end]
|
||||
key_match=re.search(rf"(?mi)^\s*{re.escape(key)}\s*=.*(?:\n|$)",chunk)
|
||||
if key_match:
|
||||
start=section_match.end()+key_match.start(); stop=section_match.end()+key_match.end()
|
||||
return text[:start]+((f"{key} = {value}\n") if value is not None else "")+text[stop:]
|
||||
if value is not None:return text[:end].rstrip()+f"\n{key} = {value}\n\n"+text[end:].lstrip("\n")
|
||||
return text
|
||||
|
||||
def save_ini_text(text,prefix="api"):
|
||||
parsed=configparser.ConfigParser(); parsed.read_string(text)
|
||||
backup=CFG.with_name(f"{CFG.name}.bak.{prefix}-{time.strftime('%Y%m%d-%H%M%S')}")
|
||||
shutil.copy2(CFG,backup)
|
||||
temporary=CFG.with_suffix(".tmp"); temporary.write_text(text); temporary.chmod(0o644); temporary.replace(CFG)
|
||||
return str(backup)
|
||||
|
||||
def edit_host(operation,host_id,url=""):
|
||||
if operation not in {"save","delete"}: raise ValueError("invalid host operation")
|
||||
if not NAME_RE.fullmatch(host_id): raise ValueError("invalid host identifier")
|
||||
local=machine_identity()["id"]
|
||||
if host_id==local: raise ValueError("local host cannot be changed in the remote host list")
|
||||
if operation=="save": url=normalize_url(url)
|
||||
text=replace_ini_value("hosts",host_id,url if operation=="save" else None)
|
||||
return save_ini_text(text,"hosts")
|
||||
|
||||
def update_ini(updates):
|
||||
if not isinstance(updates,dict): raise ValueError("updates must be an object")
|
||||
@@ -125,6 +224,50 @@ def journal_query(since=None,until=None,severity="",event="",search="",sort="tim
|
||||
return {"logs":rows[offset:offset+limit],"total":total,"offset":offset,"limit":limit,
|
||||
"event_types":event_types,"retention":"systemd-journal"}
|
||||
|
||||
def fleet_journal_query(since=None,until=None,severity="",event="",search="",machine="",
|
||||
sort="timestamp",order="desc",limit=100,offset=0):
|
||||
hosts=host_registry(); rows=[]; event_types=set(); machines=[]
|
||||
def decorate(payload,host):
|
||||
machine_info=payload.get("machine") or {"id":host["id"],"name":host["name"]}
|
||||
host_id=machine_info.get("id") or host["id"]
|
||||
host_name=machine_info.get("name") or host["name"]
|
||||
decorated=[]
|
||||
for row in payload.get("logs",[]):
|
||||
decorated.append({**row,"machine_id":host_id,"machine_name":host_name})
|
||||
return decorated,set(payload.get("event_types",[])),{"id":host_id,"name":host_name}
|
||||
local_payload=journal_query(since,until,severity,event,search,"timestamp","desc",50000,0)
|
||||
local_payload["machine"]=machine_identity()
|
||||
local_rows,local_events,local_machine=decorate(local_payload,hosts[0])
|
||||
rows.extend(local_rows); event_types.update(local_events); machines.append(local_machine)
|
||||
query={"limit":50000,"offset":0,"sort":"timestamp","order":"desc"}
|
||||
for key,value in (("since",since),("until",until),("severity",severity),("event",event),("search",search)):
|
||||
if value not in (None,""): query[key]=value
|
||||
def fetch(host):
|
||||
payload=remote_json(host["url"],"/api/v1/logs?"+urlencode(query))
|
||||
try:
|
||||
status=remote_json(host["url"],"/api/v1/status")
|
||||
payload["machine"]=status.get("machine",{})
|
||||
except Exception: pass
|
||||
return decorate(payload,host)
|
||||
peers=hosts[1:]
|
||||
with ThreadPoolExecutor(max_workers=min(8,max(1,len(peers)))) as pool:
|
||||
futures=[pool.submit(fetch,h) for h in peers]
|
||||
for future in as_completed(futures):
|
||||
try:
|
||||
remote_rows,remote_events,remote_machine=future.result()
|
||||
rows.extend(remote_rows); event_types.update(remote_events); machines.append(remote_machine)
|
||||
except Exception: pass
|
||||
if machine: rows=[row for row in rows if row["machine_id"]==machine]
|
||||
severity_rank={"EMERGENCY":0,"ALERT":1,"CRITICAL":2,"ERROR":3,"WARN":4,"NOTICE":5,"INFO":6,"DEBUG":7}
|
||||
keys={"timestamp":lambda x:x["timestamp"],"severity":lambda x:severity_rank.get(x["severity"],99),
|
||||
"event":lambda x:x["event"],"message":lambda x:x["message"],
|
||||
"machine":lambda x:(x["machine_name"],x["machine_id"])}
|
||||
rows.sort(key=keys[sort],reverse=order=="desc")
|
||||
total=len(rows)
|
||||
return {"logs":rows[offset:offset+limit],"total":total,"offset":offset,"limit":limit,
|
||||
"event_types":sorted(event_types),"machines":sorted(machines,key=lambda x:(x["name"],x["id"])),
|
||||
"retention":"systemd-journal"}
|
||||
|
||||
def service_inventory():
|
||||
cfg=load_cfg()
|
||||
monitored={target.lower():(name,target) for name,target in cfg.items("services")} if cfg.has_section("services") else {}
|
||||
@@ -205,7 +348,7 @@ def control_service(unit,action):
|
||||
if result.returncode: raise RuntimeError(result.stderr.strip() or f"systemctl {action} failed")
|
||||
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
server_version="PIGWayAPI/3.6"
|
||||
server_version="PIGWayAPI/3.7"
|
||||
|
||||
def log_message(self,fmt,*args):
|
||||
return
|
||||
@@ -240,37 +383,75 @@ class Handler(BaseHTTPRequestHandler):
|
||||
self.send_header("Content-Length",str(len(data))); self.security_headers(); self.end_headers(); self.wfile.write(data); return
|
||||
if parsed.path=="/api/v1/health":
|
||||
self.send_json(200,{"ok":True,"version":APP_VERSION,"status_available":STATUS.exists()}); return
|
||||
if parsed.path=="/api/v1/hosts":
|
||||
self.send_json(200,{"hosts":host_registry()}); return
|
||||
if parsed.path=="/api/v1/fleet/status":
|
||||
self.send_json(200,fleet_status()); return
|
||||
if parsed.path=="/api/v1/status":
|
||||
if not STATUS.exists(): self.send_json(503,{"error":"agent status unavailable"}); return
|
||||
self.send_json(200,json.loads(STATUS.read_text())); return
|
||||
if parsed.path=="/api/v1/logs":
|
||||
if parsed.path in {"/api/v1/logs","/api/v1/fleet/logs"}:
|
||||
query=parse_qs(parsed.query)
|
||||
value=lambda key,default="":query.get(key,[default])[0]
|
||||
limit=min(max(int(value("limit",api_setting("log_limit",200,int))),1),500)
|
||||
limit=min(max(int(value("limit",api_setting("log_limit",200,int))),1),50000)
|
||||
offset=max(int(value("offset",0)),0)
|
||||
since=int(value("since")) if value("since") else None
|
||||
until=int(value("until")) if value("until") else None
|
||||
if since is not None and until is not None and since>until: raise ValueError("since must not be after until")
|
||||
sort=value("sort","timestamp"); order=value("order","desc")
|
||||
if sort not in {"timestamp","severity","event","message"}: raise ValueError("invalid log sort")
|
||||
allowed_sort={"timestamp","severity","event","message"}
|
||||
if parsed.path=="/api/v1/fleet/logs": allowed_sort.add("machine")
|
||||
if sort not in allowed_sort: raise ValueError("invalid log sort")
|
||||
if order not in {"asc","desc"}: raise ValueError("invalid log order")
|
||||
severity=value("severity").upper(); event=value("event"); search=value("search")
|
||||
if len(event)>80 or len(search)>120: raise ValueError("log filter is too long")
|
||||
if parsed.path=="/api/v1/fleet/logs":
|
||||
self.send_json(200,fleet_journal_query(since,until,severity,event,search,value("machine"),sort,order,limit,offset)); return
|
||||
self.send_json(200,journal_query(since,until,severity,event,search,sort,order,limit,offset)); return
|
||||
if parsed.path=="/api/v1/config":
|
||||
self.send_json(200,{"config":config_json(),"write_requires_token":True}); return
|
||||
query=parse_qs(parsed.query); host_id=query.get("host",[""])[0]
|
||||
if host_id and host_id!=machine_identity()["id"]:
|
||||
host=find_host(host_id); self.send_json(200,remote_json(host["url"],"/api/v1/config")); return
|
||||
self.send_json(200,{"config":config_json(),"write_requires_token":True,"machine":machine_identity()}); return
|
||||
if parsed.path=="/api/v1/services":
|
||||
self.send_json(200,{"services":service_inventory()}); return
|
||||
query=parse_qs(parsed.query); host_id=query.get("host",[""])[0]
|
||||
if host_id and host_id!=machine_identity()["id"]:
|
||||
host=find_host(host_id); self.send_json(200,remote_json(host["url"],"/api/v1/services")); return
|
||||
self.send_json(200,{"services":service_inventory(),"machine":machine_identity()}); return
|
||||
self.send_json(404,{"error":"not found"})
|
||||
except ValueError as e: self.send_json(400,{"error":str(e)})
|
||||
except Exception as e: self.send_json(500,{"error":str(e)})
|
||||
|
||||
def do_PUT(self):
|
||||
path=urlparse(self.path).path
|
||||
if path not in {"/api/v1/config","/api/v1/services","/api/v1/services/control"}: self.send_json(404,{"error":"not found"}); return
|
||||
if path not in {"/api/v1/config","/api/v1/services","/api/v1/services/control","/api/v1/hosts","/api/v1/fleet/config"}: self.send_json(404,{"error":"not found"}); return
|
||||
if not self.authorized(): self.send_json(401,{"error":"bearer token required"}); return
|
||||
try:
|
||||
body=self.read_json()
|
||||
if path=="/api/v1/hosts":
|
||||
backup=edit_host(str(body.get("operation","")),str(body.get("id","")),str(body.get("url","")))
|
||||
self.send_json(200,{"ok":True,"backup":backup,"hosts":host_registry()}); return
|
||||
if path=="/api/v1/fleet/config":
|
||||
updates=body.get("updates"); targets=body.get("targets",[]); tokens=body.get("tokens",{})
|
||||
if not isinstance(targets,list) or not targets or len(targets)>32: raise ValueError("targets must contain 1..32 host identifiers")
|
||||
if not isinstance(tokens,dict): raise ValueError("tokens must be an object")
|
||||
results=[]; local_id=machine_identity()["id"]; restart_local_api=False
|
||||
for host_id in dict.fromkeys(str(x) for x in targets):
|
||||
try:
|
||||
if host_id==local_id:
|
||||
backup=update_ini(updates)
|
||||
subprocess.run(["systemctl","restart","pigway-pi-control.service"],check=True,timeout=10)
|
||||
restart_local_api="api" in updates
|
||||
results.append({"id":host_id,"ok":True,"backup":backup})
|
||||
else:
|
||||
host=find_host(host_id); token=str(tokens.get(host_id,""))
|
||||
if not token: raise ValueError("remote bearer token required")
|
||||
response=remote_json(host["url"],"/api/v1/config","PUT",{"updates":updates},token)
|
||||
results.append({"id":host_id,"ok":True,"backup":response.get("backup","")})
|
||||
except Exception as exc: results.append({"id":host_id,"ok":False,"error":str(exc)})
|
||||
self.send_json(200,{"ok":all(x["ok"] for x in results),"results":results})
|
||||
if restart_local_api: threading.Timer(0.2,self.server.shutdown).start()
|
||||
return
|
||||
if path=="/api/v1/services/control":
|
||||
control_service(str(body.get("unit","")),str(body.get("action","")))
|
||||
print(f"SERVICE_CONTROL action={body.get('action')} unit={body.get('unit')}",flush=True)
|
||||
|
||||
Reference in New Issue
Block a user