Files
UmbrelApps/whatsnext-evolu-relay/agent.py.template
T

260 lines
9.5 KiB
Plaintext
Raw Normal View History

# ═══════════════════════════════════════════════════════════════════════════════
# 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()