# ═══════════════════════════════════════════════════════════════════════════════ # De agent van Evolu Relay. # # Hij bedient de statuspagina en doet verder niets: lezen wat het relay-proces # heeft opgeschreven, en opdrachten van de pagina in de postbus leggen. # # Waarom de agent niet zelf beslist wie er binnen mag: dat beleid hoort bij het # proces dat de verbindingen aanneemt, en dat is de relay. Twee processen die in # dezelfde allowlist schrijven is een wedloop die je een keer per jaar treft en # dan niet kunt reproduceren. De agent schrijft daarom uitsluitend command.json, # en het relay-proces past hem toe en ruimt hem op. # # LET OP: umbreld haalt dit bestand bij elke start door envsubst. Er mag dus geen # dollarteken in staan, ook niet in een regex of een tekst. Een accolade-variabele # die niet bestaat wordt leeg, en dat sloopt Python-code zonder foutmelding. # tests/test_appstore_vorm.py controleert dat. # ═══════════════════════════════════════════════════════════════════════════════ import json import os import socket import threading import time from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path STATE_DIR = Path(os.environ.get("RELAY_STATE_DIR", "/var/lib/relay")) OWNERS_FILE = STATE_DIR / "owners.json" COMMAND_FILE = STATE_DIR / "command.json" DATABASE_FILE = STATE_DIR / "evolu-relay.db" RELAY_HOST = os.environ.get("RELAY_HOST", "") RELAY_PORT = int(os.environ.get("RELAY_PORT", "4000")) PUBLIC_PORT = int(os.environ.get("RELAY_PUBLIC_PORT", "3852")) API_PORT = int(os.environ.get("RELAY_API_PORT", "8000")) PROBE_INTERVAL = int(os.environ.get("RELAY_PROBE_INTERVAL", "15")) # Wat de pagina mag vragen. Expliciet en niet doorgeven wat er binnenkomt: dit # bestand wordt door een ander proces uitgevoerd, en een onbekende actie hoort # hier te stranden en niet daar. ALLOWED_ACTIONS = ("set-learning", "block", "allow", "forget") MAX_BODY_BYTES = 4096 MAX_OWNER_ID_LENGTH = 256 # Door de achtergrondlus bijgewerkt, door de webserver gelezen. Een dict wordt in # zijn geheel vervangen en nooit ter plekke aangepast, zodat een lezer altijd een # samenhangend beeld heeft zonder slot. probe = {"reachable": None, "checked": None} def relay_reachable(): """Kan de relay een TCP-verbinding aannemen. Bewust niet meer dan dat. De relay is een WebSocket-server en antwoordt niet op een gewoon verzoek; een handdruk nabouwen om de pagina groen te krijgen is meer code dan het waard is. Wat dit wel uitsluit is de meest voorkomende storing: het proces is omgevallen. """ if not RELAY_HOST: return None try: with socket.create_connection((RELAY_HOST, RELAY_PORT), timeout=3): return True except OSError: return False def probe_loop(): global probe while True: # In zijn geheel vervangen en niet ter plekke aanpassen: een lezer ziet # dan altijd een samenhangend beeld, zonder dat er een slot nodig is. probe = { "reachable": relay_reachable(), "checked": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()), } time.sleep(PROBE_INTERVAL) def read_owners(): """De allowlist zoals het relay-proces hem heeft achtergelaten.""" try: with OWNERS_FILE.open("r", encoding="utf-8") as handle: data = json.load(handle) except FileNotFoundError: # Nog nooit geschreven. Dat is de normale toestand vlak na een # installatie: het relay-proces schrijft pas bij de eerste wijziging. return {"state": None, "problem": "nog-niet-aangemaakt"} except (OSError, ValueError): return {"state": None, "problem": "onleesbaar"} if not isinstance(data, dict): return {"state": None, "problem": "onleesbaar"} return {"state": data, "problem": None} def database_facts(): try: stat = DATABASE_FILE.stat() except OSError: return {"bytes": None, "modified": None} return { "bytes": stat.st_size, "modified": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime(stat.st_mtime)), } def build_status(): owners = read_owners() state = owners["state"] or {} return { "relay": { "reachable": probe["reachable"], "checked": probe["checked"], "publicPort": PUBLIC_PORT, }, "owners": { "problem": owners["problem"], # Ontbreekt de staat, dan is 'learning' onbekend en niet 'false'. De # pagina hoort dat verschil te tonen: onbekend is een reden om te # kijken, uit is een keuze. "learning": state.get("learning"), "allowed": [ entry for entry in state.get("owners", []) if isinstance(entry, dict) and entry.get("allowed") is True ], "blocked": [ entry for entry in state.get("owners", []) if isinstance(entry, dict) and entry.get("allowed") is False ], "rejected": [ entry for entry in state.get("rejected", []) if isinstance(entry, dict) ], }, "database": database_facts(), # Ligt er nog een opdracht, dan heeft het relay-proces hem nog niet # opgepakt. De pagina kan dat tonen in plaats van te doen alsof er niets # gebeurd is. "pendingCommand": COMMAND_FILE.exists(), } def valid_command(payload): """Geeft de opdracht terug, of een foutmelding. Streng aan deze kant, want dit is de enige plek waar iets van buiten in de postbus belandt. """ if not isinstance(payload, dict): return None, "geen object" action = payload.get("action") if action not in ALLOWED_ACTIONS: return None, "onbekende actie" if action == "set-learning": value = payload.get("value") if not isinstance(value, bool): return None, "waarde moet true of false zijn" return {"action": action, "value": value}, None owner_id = payload.get("ownerId") if not isinstance(owner_id, str) or not owner_id or len(owner_id) > MAX_OWNER_ID_LENGTH: return None, "ontbrekende of te lange ownerId" return {"action": action, "ownerId": owner_id}, None def write_command(command): """Legt de opdracht in de postbus. Eerst een tijdelijk bestand en dan hernoemen: het relay-proces kijkt op zijn eigen moment en mag geen half bestand aantreffen. """ STATE_DIR.mkdir(parents=True, exist_ok=True) temporary = COMMAND_FILE.with_suffix(".json.tmp") with temporary.open("w", encoding="utf-8") as handle: json.dump(command, handle) handle.write("\n") temporary.replace(COMMAND_FILE) class Handler(BaseHTTPRequestHandler): # De standaardregel van BaseHTTPRequestHandler noemt de naam van de server en # de Python-versie. Dat hoeft niemand te weten. server_version = "evolu-relay-agent" sys_version = "" def log_message(self, format, *args): # Geen toegangslog. Elke regel zou het adres van de bezoeker bevatten en # de pagina zit achter de inlog van umbrelOS; er valt niets te zien wat # het bewaren waard is. return def _send(self, status, payload): body = json.dumps(payload).encode("utf-8") self.send_response(status) self.send_header("Content-Type", "application/json") self.send_header("Content-Length", str(len(body))) self.send_header("Cache-Control", "no-store") self.end_headers() self.wfile.write(body) def do_GET(self): if self.path.rstrip("/") in ("/api/status", "/status"): self._send(200, build_status()) return self._send(404, {"error": "onbekend pad"}) def do_POST(self): if self.path.rstrip("/") not in ("/api/command", "/command"): self._send(404, {"error": "onbekend pad"}) return try: length = int(self.headers.get("Content-Length", "0")) except ValueError: self._send(400, {"error": "lengte ontbreekt"}) return if length <= 0 or length > MAX_BODY_BYTES: self._send(400, {"error": "lege of te grote opdracht"}) return try: payload = json.loads(self.rfile.read(length).decode("utf-8")) except (UnicodeDecodeError, ValueError): self._send(400, {"error": "onleesbare opdracht"}) return command, problem = valid_command(payload) if command is None: self._send(400, {"error": problem}) return # Eén opdracht tegelijk. Ligt er nog een, dan zou schrijven hem stil # overschrijven en verdwijnt de vorige zonder dat iemand het merkt. if COMMAND_FILE.exists(): self._send(409, {"error": "vorige opdracht is nog niet verwerkt"}) return try: write_command(command) except OSError as error: self._send(500, {"error": "opdracht kon niet worden weggeschreven: " + str(error)}) return self._send(202, {"accepted": command}) def main(): threading.Thread(target=probe_loop, daemon=True).start() ThreadingHTTPServer(("0.0.0.0", API_PORT), Handler).serve_forever() if __name__ == "__main__": main()