refactor: separate fleet manager from local monitoring

This commit is contained in:
way
2026-09-27 22:28:03 +08:00
parent ced0140450
commit c8c12b7b70
8 changed files with 90 additions and 361 deletions
+15 -164
View File
@@ -14,8 +14,7 @@ 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={"device","hosts","services","processes","alerts","timing","dark_mode","api"}
ALLOWED_SECTIONS={"device","services","processes","alerts","timing","dark_mode","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]$")
@@ -65,60 +64,6 @@ def machine_identity():
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
@@ -132,8 +77,11 @@ def config_json():
return {section:{k:v for k,v in cfg.items(section) if not (section=="timing" and k in {"oled_refresh_interval","page_interval"})} for section in cfg.sections() if section in ALLOWED_SECTIONS}
def validate_value(section,key,value,existing):
if section=="api" and key=="enabled":
if value.lower() not in {"true","false"}:raise ValueError("enabled must be true or false")
return
if section=="timing" and key in {"oled_refresh_interval","page_interval"}:raise ValueError("display timing is managed by the hardware plugin")
if section not in {"services","processes","hosts"} and key not in existing.get(section,set()):
if section not in {"services","processes"} and key not in existing.get(section,set()):
raise ValueError(f"unsupported option: {section}.{key}")
if section=="hardware" and key=="cooling_hat_enabled":
if value.lower() not in ("true","false"): raise ValueError("cooling_hat_enabled must be true or false")
@@ -163,8 +111,7 @@ def validate_value(section,key,value,existing):
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()
@@ -172,7 +119,7 @@ def replace_ini_value(section,key,value):
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)^\[.+\][ \t]*$",text[section_match.end():])
next_section=re.search(r"(?m)^\[[^\]\r\n]+\]",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)^[ \t]*{re.escape(key)}\s*=.*(?:\n|$)",chunk)
@@ -189,15 +136,6 @@ def save_ini_text(text,prefix="api"):
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")
text=CFG.read_text()
@@ -217,7 +155,7 @@ def update_ini(updates):
if not section_match:
text=text.rstrip()+f"\n\n[{section}]\n{key} = {value}\n"
continue
next_section=re.search(r"(?m)^\[.+\][ \t]*$",text[section_match.end():])
next_section=re.search(r"(?m)^\[[^\]\r\n]+\]",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(str(key))}\s*=\s*).*$",chunk)
@@ -268,50 +206,6 @@ 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 {}
@@ -356,7 +250,7 @@ def edit_service_monitor(operation,name,target,note="",previous_name=""):
if not section_match:
if value is not None: text=text.rstrip()+f"\n\n[{section}]\n{key} = {value}\n"
return
next_section=re.search(r"(?m)^\[.+\][ \t]*$",text[section_match.end():])
next_section=re.search(r"(?m)^\[[^\]\r\n]+\]",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)^[ \t]*{re.escape(key)}\s*=.*(?:\n|$)",chunk)
@@ -422,19 +316,13 @@ class Handler(BaseHTTPRequestHandler):
def do_GET(self):
parsed=urlparse(self.path)
try:
if parsed.path=="/":
data=WEB.read_bytes(); self.send_response(200); self.send_header("Content-Type","text/html; charset=utf-8")
self.send_header("Content-Length",str(len(data))); self.security_headers(); self.end_headers(); self.wfile.write(data); return
if not self.authorized():self.send_json(401,{"error":"bearer token required"});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
self.send_json(200,{"ok":True,"version":APP_VERSION,"status_available":STATUS.exists()});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 in {"/api/v1/logs","/api/v1/fleet/logs"}:
if parsed.path=="/api/v1/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),50000)
@@ -444,33 +332,23 @@ class Handler(BaseHTTPRequestHandler):
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")
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 in {"/api/v1/plugins","/api/v1/plugin/config"}:
query=parse_qs(parsed.query);host_id=query.get("host",[""])[0]
plugin_id=query.get("id",[""])[0]
if host_id and host_id!=machine_identity()["id"]:
host=find_host(host_id)
self.send_json(200,remote_json(host["url"],parsed.path+"?"+urlencode({"id":plugin_id})));return
if parsed.path.endswith("/plugins"):
self.send_json(200,{"plugins":plugin_inventory()});return
item=find_plugin(plugin_id)
self.send_json(200,plugin_request(item["socket"],"GET","/v1/config"));return
if parsed.path=="/api/v1/config":
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":
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)})
@@ -478,44 +356,15 @@ class Handler(BaseHTTPRequestHandler):
def do_PUT(self):
path=urlparse(self.path).path
if path not in {"/api/v1/plugins","/api/v1/plugin/config","/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 path not in {"/api/v1/plugins","/api/v1/plugin/config","/api/v1/config","/api/v1/services","/api/v1/services/control"}: 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 in {"/api/v1/plugins","/api/v1/plugin/config"}:
host_id=str(body.get("host",""))
if host_id and host_id!=machine_identity()["id"]:
host=find_host(host_id);token=str(body.get("target_token",""))
if not token:raise ValueError("remote bearer token required")
payload={k:v for k,v in body.items() if k not in {"host","target_token"}}
self.send_json(200,remote_json(host["url"],path,"PUT",payload,token));return
if path.endswith("/plugins"):self.send_json(200,plugin_operation(body));return
item=find_plugin(str(body.get("id","")))
self.send_json(200,plugin_request(item["socket"],"PUT","/v1/config",{"updates":body.get("updates")}));return
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)
@@ -537,6 +386,8 @@ class Handler(BaseHTTPRequestHandler):
except Exception as e: self.send_json(500,{"error":str(e)})
def main():
if not load_cfg().getboolean("api","enabled",fallback=False):
print("API_DISABLED",flush=True);sys.exit(78)
host=api_setting("bind","0.0.0.0")
port=api_setting("port",6001,int)
server=ThreadingHTTPServer((host,port),Handler)