262 lines
15 KiB
Python
262 lines
15 KiB
Python
#!/usr/bin/env python3
|
|
import configparser,ipaddress,json,re,secrets,shutil,subprocess,threading,time
|
|
from http import HTTPStatus
|
|
from http.server import BaseHTTPRequestHandler,ThreadingHTTPServer
|
|
from pathlib import Path
|
|
from urllib.parse import parse_qs,urlparse
|
|
|
|
APP_VERSION="3.6.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"}
|
|
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]$")
|
|
FLASH_KEY_RE=re.compile(r"^(warning|critical|emergency|normal)_flash_(on|off)_ms$")
|
|
NORMAL_MODES={"off","solid","flash","flow","breathe","marquee","rainbow","colorful"}
|
|
EFFECT_COLORS={"red","green","blue","yellow","purple","cyan","white"}
|
|
PROTECTED_UNITS={"ssh.service","sshd.service","runawesun.service","tailscaled.service",
|
|
"pigway-pi-control.service","pigway-pi-control-api.service"}
|
|
|
|
def load_cfg():
|
|
cfg=configparser.ConfigParser(); cfg.read(CFG)
|
|
return cfg
|
|
|
|
def api_setting(key,default,cast=str):
|
|
try:return cast(load_cfg().get("api",key))
|
|
except Exception:return default
|
|
|
|
def config_json():
|
|
cfg=load_cfg()
|
|
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()):
|
|
raise ValueError(f"unsupported option: {section}.{key}")
|
|
if RGB_KEY_RE.fullmatch(key):
|
|
number=int(value)
|
|
if not 0<=number<=255: raise ValueError(f"{key} must be 0..255")
|
|
elif FLASH_KEY_RE.fullmatch(key):
|
|
number=int(value)
|
|
if not 50<=number<=60000: raise ValueError(f"{key} must be 50..60000")
|
|
elif section=="led" and key=="normal_mode":
|
|
if value.lower() not in NORMAL_MODES: raise ValueError("normal_mode must be off, solid, flash, flow, breathe, marquee, rainbow or colorful")
|
|
elif section=="led" and key=="normal_effect_speed":
|
|
if int(value) not in (1,2,3): raise ValueError("normal_effect_speed must be 1, 2 or 3")
|
|
elif section=="led" and key=="normal_effect_color":
|
|
if value.lower() not in EFFECT_COLORS: raise ValueError("normal_effect_color must be red, green, blue, yellow, purple, cyan or white")
|
|
elif section=="api" and key=="port":
|
|
number=int(value)
|
|
if not 1<=number<=65535: raise ValueError("api.port must be 1..65535")
|
|
elif section=="api" and key=="log_limit":
|
|
number=int(value)
|
|
if not 1<=number<=1000: raise ValueError("api.log_limit must be 1..1000")
|
|
elif section=="api" and key=="bind":
|
|
ipaddress.ip_address(value)
|
|
|
|
def update_ini(updates):
|
|
if not isinstance(updates,dict): raise ValueError("updates must be an object")
|
|
text=CFG.read_text()
|
|
current=load_cfg(); existing={section:set(current.options(section)) for section in current.sections()}
|
|
for section,values in updates.items():
|
|
if section not in ALLOWED_SECTIONS or not isinstance(values,dict):
|
|
raise ValueError(f"unsupported section: {section}")
|
|
for key,value in values.items():
|
|
if not NAME_RE.fullmatch(str(key)): raise ValueError(f"invalid key: {key}")
|
|
if value is None:
|
|
raise ValueError(f"null is not allowed: {section}.{key}")
|
|
value=str(value).strip()
|
|
if not value or len(value)>256 or any(c in value for c in "\r\n\x00"):
|
|
raise ValueError(f"invalid value: {section}.{key}")
|
|
validate_value(section,str(key),value,existing)
|
|
section_match=re.search(rf"(?mi)^\[{re.escape(section)}\]\s*$",text)
|
|
if not section_match:
|
|
text=text.rstrip()+f"\n\n[{section}]\n{key} = {value}\n"
|
|
continue
|
|
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(str(key))}\s*=\s*).*$",chunk)
|
|
if key_match:
|
|
start=section_match.end()+key_match.start(); stop=section_match.end()+key_match.end()
|
|
replacement=key_match.group(1)+value
|
|
text=text[:start]+replacement+text[stop:]
|
|
else:
|
|
text=text[:end].rstrip()+f"\n{key} = {value}\n\n"+text[end:].lstrip("\n")
|
|
parsed=configparser.ConfigParser(); parsed.read_string(text)
|
|
backup=CFG.with_name(f"{CFG.name}.bak.api-{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 journal(limit):
|
|
result=subprocess.run(["journalctl","-u","pigway-pi-control.service","-n",str(limit),
|
|
"--no-pager","-o","json"],capture_output=True,text=True,timeout=5)
|
|
if result.returncode: raise RuntimeError(result.stderr.strip() or "journalctl failed")
|
|
rows=[]
|
|
for line in result.stdout.splitlines():
|
|
try:
|
|
item=json.loads(line)
|
|
rows.append({"timestamp":item.get("__REALTIME_TIMESTAMP"),"priority":item.get("PRIORITY"),
|
|
"message":item.get("MESSAGE","")})
|
|
except json.JSONDecodeError: pass
|
|
return rows
|
|
|
|
def service_inventory():
|
|
cfg=load_cfg()
|
|
monitored={target.lower():(name,target) for name,target in cfg.items("services")} if cfg.has_section("services") else {}
|
|
notes=dict(cfg.items("service_notes")) if cfg.has_section("service_notes") else {}
|
|
files=subprocess.run(["systemctl","list-unit-files","--type=service","--no-legend","--no-pager"],
|
|
capture_output=True,text=True,timeout=10,check=True)
|
|
units={}
|
|
for line in files.stdout.splitlines():
|
|
parts=line.split()
|
|
if len(parts)>=2 and UNIT_RE.fullmatch(parts[0]):
|
|
units[parts[0]]={"unit":parts[0],"enabled":parts[1],"active":"inactive","description":""}
|
|
states=subprocess.run(["systemctl","list-units","--all","--type=service","--no-legend","--no-pager","--plain"],
|
|
capture_output=True,text=True,timeout=10,check=True)
|
|
for line in states.stdout.splitlines():
|
|
parts=line.split(None,4)
|
|
if len(parts)>=4 and UNIT_RE.fullmatch(parts[0]):
|
|
item=units.setdefault(parts[0],{"unit":parts[0],"enabled":"unknown"})
|
|
item.update({"active":parts[2],"description":parts[4] if len(parts)>4 else ""})
|
|
for name,target in monitored.values():
|
|
item=units.setdefault(target,{"unit":target,"enabled":"not-found","active":"inactive","description":""})
|
|
item.update({"monitored":True,"monitor_name":name.upper(),"note":notes.get(name,"")})
|
|
for item in units.values():
|
|
item.setdefault("monitored",False); item.setdefault("monitor_name",""); item.setdefault("note","")
|
|
item["protected"]=item["unit"] in PROTECTED_UNITS
|
|
return sorted(units.values(),key=lambda x:(not x["monitored"],x["unit"]))
|
|
|
|
def edit_service_monitor(operation,name,target,note="",previous_name=""):
|
|
if operation not in {"save","delete"}: raise ValueError("operation must be save or delete")
|
|
if not NAME_RE.fullmatch(name): raise ValueError("invalid monitor name")
|
|
if not UNIT_RE.fullmatch(target): raise ValueError("invalid systemd service name")
|
|
if len(note)>120 or any(c in note for c in "\r\n\x00"): raise ValueError("invalid note")
|
|
if previous_name and not NAME_RE.fullmatch(previous_name): raise ValueError("invalid previous monitor name")
|
|
cfg=load_cfg(); old_target=cfg.get("services",name,fallback=None)
|
|
text=CFG.read_text()
|
|
def change(section,key,value):
|
|
nonlocal 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
|
|
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()
|
|
text=text[:start]+((f"{key} = {value}\n") if value is not None else "")+text[stop:]
|
|
elif value is not None:
|
|
text=text[:end].rstrip()+f"\n{key} = {value}\n\n"+text[end:].lstrip("\n")
|
|
if operation=="delete":
|
|
if old_target is None: raise ValueError("monitor not found")
|
|
change("services",name,None); change("service_notes",name,None)
|
|
else:
|
|
if previous_name and previous_name.lower()!=name.lower():
|
|
change("services",previous_name,None); change("service_notes",previous_name,None)
|
|
change("services",name,target); change("service_notes",name,note or None)
|
|
parsed=configparser.ConfigParser(); parsed.read_string(text)
|
|
backup=CFG.with_name(f"{CFG.name}.bak.api-{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)
|
|
subprocess.run(["systemctl","restart","pigway-pi-control.service"],check=True,timeout=10)
|
|
return str(backup)
|
|
|
|
def control_service(unit,action):
|
|
if not UNIT_RE.fullmatch(unit): raise ValueError("invalid systemd service name")
|
|
if action not in {"start","stop"}: raise ValueError("action must be start or stop")
|
|
if unit in PROTECTED_UNITS and action=="stop": raise ValueError("protected remote-management service cannot be stopped here")
|
|
result=subprocess.run(["systemctl",action,unit],capture_output=True,text=True,timeout=20)
|
|
if result.returncode: raise RuntimeError(result.stderr.strip() or f"systemctl {action} failed")
|
|
|
|
class Handler(BaseHTTPRequestHandler):
|
|
server_version="PIGWayAPI/3.6"
|
|
|
|
def log_message(self,fmt,*args):
|
|
return
|
|
|
|
def send_json(self,status,payload):
|
|
data=json.dumps(payload,ensure_ascii=False,separators=(",",":")).encode()
|
|
self.send_response(status); self.send_header("Content-Type","application/json; charset=utf-8")
|
|
self.send_header("Content-Length",str(len(data))); self.security_headers(); self.end_headers(); self.wfile.write(data)
|
|
|
|
def security_headers(self):
|
|
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'")
|
|
|
|
def authorized(self):
|
|
try: expected=TOKEN.read_text().strip()
|
|
except OSError: return False
|
|
supplied=self.headers.get("Authorization","")
|
|
return supplied.startswith("Bearer ") and secrets.compare_digest(supplied[7:],expected)
|
|
|
|
def read_json(self):
|
|
length=int(self.headers.get("Content-Length","0"))
|
|
if length<=0 or length>65536: raise ValueError("invalid request size")
|
|
return json.loads(self.rfile.read(length))
|
|
|
|
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 parsed.path=="/api/v1/health":
|
|
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=="/api/v1/logs":
|
|
limit=min(max(int(parse_qs(parsed.query).get("limit",[api_setting("log_limit",200,int)])[0]),1),1000)
|
|
self.send_json(200,{"logs":journal(limit)}); return
|
|
if parsed.path=="/api/v1/config":
|
|
self.send_json(200,{"config":config_json(),"write_requires_token":True}); return
|
|
if parsed.path=="/api/v1/services":
|
|
self.send_json(200,{"services":service_inventory()}); return
|
|
self.send_json(404,{"error":"not found"})
|
|
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 not self.authorized(): self.send_json(401,{"error":"bearer token required"}); return
|
|
try:
|
|
body=self.read_json()
|
|
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)
|
|
self.send_json(200,{"ok":True}); return
|
|
if path=="/api/v1/services":
|
|
backup=edit_service_monitor(str(body.get("operation","")),str(body.get("name","")),
|
|
str(body.get("target","")),str(body.get("note","")),
|
|
str(body.get("previous_name","")))
|
|
print(f"SERVICE_MONITOR_UPDATED operation={body.get('operation')} target={body.get('target')} backup={backup}",flush=True)
|
|
self.send_json(200,{"ok":True,"backup":backup}); return
|
|
updates=body.get("updates")
|
|
backup=update_ini(updates)
|
|
subprocess.run(["systemctl","restart","pigway-pi-control.service"],check=True,timeout=10)
|
|
restart_api="api" in updates
|
|
print(f"CONFIG_UPDATED sections={','.join(sorted(updates))} backup={backup}",flush=True)
|
|
self.send_json(200,{"ok":True,"backup":backup,"api_restart":restart_api})
|
|
if restart_api: threading.Timer(0.2,self.server.shutdown).start()
|
|
except (ValueError,json.JSONDecodeError) as e: self.send_json(400,{"error":str(e)})
|
|
except Exception as e: self.send_json(500,{"error":str(e)})
|
|
|
|
def main():
|
|
host=api_setting("bind","0.0.0.0")
|
|
port=api_setting("port",6001,int)
|
|
server=ThreadingHTTPServer((host,port),Handler)
|
|
print(f"API_START version={APP_VERSION} bind={host} port={port}",flush=True)
|
|
try: server.serve_forever()
|
|
finally: server.server_close()
|
|
|
|
if __name__=="__main__": main()
|