260 lines
9.5 KiB
Plaintext
260 lines
9.5 KiB
Plaintext
# ═══════════════════════════════════════════════════════════════════════════════
|
|||
|
|
# 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()
|