Vier punten van de gebruiker, opgekomen tijdens het beproeven met de testclient. De getimede leerstand is eruit; de wachtlijst is de enige weg naar binnen. Zijn redenering: allebei de wegen vragen iemand die bij de app kan, dus het is dubbelop, en het venster is de zwakste omdat het iedereen toelaat die er toevallig in verbindt. Dat weegt zwaarder nu de relay op een publiek wss-adres kan staan. De oorspronkelijke reden voor de leerstand, dat je je eigen OwnerId nergens kon aflezen, verviel toen de weigerlijst dat id ging tonen. Daarmee verdwijnt ook de bug die hij dezelfde dag meldde: een geleerde eigenaar bleef in de weigerlijst staan terwijl hij al kon schrijven en lezen, want decideOwner haalde hem niet van die lijst af en de knop allow wel. "Refused owners" heet "Waiting list", met de badge Waiting en een teller waar de widget van het tijdvenster stond. Het veld op schijf blijft rejected: hernoemen zou een migratie zijn voor een woord dat niemand ziet. Het adres onderaan zei http:// en dat kan nergens werken, want de relay spreekt WebSocket en nooit HTTP. Nu ws://<host>:3852, met een regel over wss://<domein> zonder poort achter een reverse proxy. Dat is precies de fout die diezelfde dag een ronde kostte bij het koppelen van de testclient. De melding bij elke klik is weg. Die stond in de gewone stroom van de pagina, dus alles eronder schoof omlaag en weer omhoog. Nu gaan de knoppen in de lijsten even op slot tot de ronde de nieuwe stand heeft; foutmeldingen blijven wel staan, want die zeggen iets wat je nergens anders ziet. STATE_VERSION blijft 1 en een owners.json van 0.6.0 leest door: learning en learningUntil worden gelezen, genegeerd en niet teruggeschreven. Een verhoging zou store.js de allowlist van een werkende installatie opzij laten schuiven. Twee toetsen bewaken dat de leerstand niet terugsluipt: een onbekende eigenaar wordt geweigerd ook met learning: true in het bestand, en set-learning is een onbekende actie. Beide mutatie-getest. De compose staat op 0.7.0 zonder digest, zodat het hard faalt tot de image bestaat. Bouwen, duwen en pinnen ligt bij de gebruiker. Suite: 424 goed, 0 fout. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
439 lines
16 KiB
Python
439 lines
16 KiB
Python
# ═══════════════════════════════════════════════════════════════════════════════
|
|
# 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.
|
|
#
|
|
# Dit bestand zit in de image (tools/evolu-relay/Dockerfile) als /app/agent.py.
|
|
# Tot 0.5.6 was het een *.template in de app-map en haalde umbreld het bij elke
|
|
# start door envsubst, en daarom staat er geen dollarteken in en komen alle
|
|
# instellingen uit de omgeving. Dat laatste is zo gebleven: de compose zet ze in
|
|
# de omgeving neer, en dat is ook de nette weg. tests/test_relay_agent.py houdt
|
|
# het dollarteken buiten de deur, voor het geval dit ooit weer een template wordt.
|
|
# ═══════════════════════════════════════════════════════════════════════════════
|
|
|
|
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"
|
|
|
|
# De labels die de gebruiker aan een eigenaar-id hangt. Een eigen bestand, en dat
|
|
# is de kern van deze keuze: owners.json is van het relay-proces en labels.json is
|
|
# van de agent, dus er is per bestand precies één schrijver. Dat is dezelfde
|
|
# afspraak die de postbus hierboven oplevert, en de reden staat bovenaan dit
|
|
# bestand.
|
|
#
|
|
# Wat het bovendien oplevert: een label is meteen opgeslagen en niet pas als de
|
|
# relay de postbus leegmaakt, en labelen blijft werken als de relay omgevallen is.
|
|
# Dat mag, want de relay hoeft dit niet te weten: een label zegt niets over wie er
|
|
# binnen mag.
|
|
LABELS_FILE = STATE_DIR / "labels.json"
|
|
|
|
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"))
|
|
|
|
# De versie uit het manifest, voor de kop van de pagina. Tot 0.5.6 vulde umbreld
|
|
# die rechtstreeks in de pagina in; nu de pagina in de image zit, loopt het via de
|
|
# status. Leeg betekent dat de compose hem niet doorgeeft, en dan toont de pagina
|
|
# geen versie in plaats van een verzonnen.
|
|
APP_VERSION = os.environ.get("RELAY_APP_VERSION", "")
|
|
|
|
# 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 = ("block", "allow", "forget")
|
|
|
|
MAX_BODY_BYTES = 4096
|
|
MAX_OWNER_ID_LENGTH = 256
|
|
|
|
# Een label is een herkenpunt en geen aantekenveld: het staat op de pagina naast
|
|
# een id en moet daar op één regel passen.
|
|
MAX_LABEL_LENGTH = 48
|
|
|
|
# Een bovengrens op het aantal labels. De pagina zit achter de inlog van umbrelOS,
|
|
# dus dit is geen verdediging tegen een aanvaller maar tegen een lus die per
|
|
# ongeluk blijft schrijven. Ruim boven het aantal eigenaars dat iemand ooit heeft.
|
|
MAX_LABELS = 200
|
|
|
|
# Eén schrijver per bestand is de afspraak, maar de agent zelf is meerdradig:
|
|
# ThreadingHTTPServer geeft elk verzoek zijn eigen draad. Twee labels die op
|
|
# hetzelfde moment binnenkomen zouden elkaar dus kunnen overschrijven, want
|
|
# labelen is lezen-wijzigen-schrijven. Dit slot maakt daar één handeling van.
|
|
labels_lock = threading.Lock()
|
|
|
|
# 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 read_labels():
|
|
"""De labels, of een lege verzameling.
|
|
|
|
Bewust vergevingsgezind, en dat is het omgekeerde van hoe owners.json gelezen
|
|
wordt. Daar hangt aan een half begrepen bestand de vraag wie er binnen mag, en
|
|
dan is weigeren het antwoord. Hier gaat het om een naam naast een id: is het
|
|
onleesbaar, dan is het ergste gevolg dat je de rauwe ids ziet.
|
|
"""
|
|
try:
|
|
with LABELS_FILE.open("r", encoding="utf-8") as handle:
|
|
data = json.load(handle)
|
|
except (OSError, ValueError):
|
|
return {}
|
|
|
|
if not isinstance(data, dict):
|
|
return {}
|
|
|
|
schoon = {}
|
|
for owner_id, label in data.items():
|
|
if not isinstance(owner_id, str) or not isinstance(label, str):
|
|
continue
|
|
if not owner_id or len(owner_id) > MAX_OWNER_ID_LENGTH:
|
|
continue
|
|
label = label.strip()
|
|
if label:
|
|
schoon[owner_id] = label[:MAX_LABEL_LENGTH]
|
|
return schoon
|
|
|
|
|
|
def write_labels(labels):
|
|
"""Schrijft de labels. Eerst een tijdelijk bestand en dan hernoemen.
|
|
|
|
Hernoemen binnen dezelfde map is atomair, dus een onderbroken schrijfactie
|
|
laat geen half bestand achter. Dezelfde constructie als de postbus.
|
|
"""
|
|
STATE_DIR.mkdir(parents=True, exist_ok=True)
|
|
temporary = LABELS_FILE.with_suffix(".json.tmp")
|
|
with temporary.open("w", encoding="utf-8") as handle:
|
|
json.dump(labels, handle, indent=2, sort_keys=True)
|
|
handle.write("\n")
|
|
temporary.replace(LABELS_FILE)
|
|
|
|
|
|
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 met_label(entries, labels):
|
|
"""Hangt het label van de gebruiker aan elke regel.
|
|
|
|
Een kopie en niet ter plekke: wat hier binnenkomt komt uit het bestand van het
|
|
relay-proces, en daar horen wij niets aan toe te voegen.
|
|
"""
|
|
resultaat = []
|
|
for entry in entries:
|
|
regel = dict(entry)
|
|
regel["label"] = labels.get(entry.get("id"))
|
|
resultaat.append(regel)
|
|
return resultaat
|
|
|
|
|
|
def build_status():
|
|
owners = read_owners()
|
|
state = owners["state"] or {}
|
|
labels = read_labels()
|
|
|
|
return {
|
|
"version": APP_VERSION or None,
|
|
"relay": {
|
|
"reachable": probe["reachable"],
|
|
"checked": probe["checked"],
|
|
"publicPort": PUBLIC_PORT,
|
|
},
|
|
"owners": {
|
|
"problem": owners["problem"],
|
|
"allowed": met_label(
|
|
[
|
|
entry
|
|
for entry in state.get("owners", [])
|
|
if isinstance(entry, dict) and entry.get("allowed") is True
|
|
],
|
|
labels,
|
|
),
|
|
"blocked": met_label(
|
|
[
|
|
entry
|
|
for entry in state.get("owners", [])
|
|
if isinstance(entry, dict) and entry.get("allowed") is False
|
|
],
|
|
labels,
|
|
),
|
|
"rejected": met_label(
|
|
[entry for entry in state.get("rejected", []) if isinstance(entry, dict)],
|
|
labels,
|
|
),
|
|
},
|
|
"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"
|
|
|
|
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 valid_label(payload):
|
|
"""Geeft (ownerId, label) terug, of een foutmelding.
|
|
|
|
Een leeg label is geen fout maar de manier om er een weg te halen: dan hoeft er
|
|
geen tweede opdracht te bestaan voor iets dat de gebruiker als hetzelfde veld
|
|
ziet.
|
|
"""
|
|
if not isinstance(payload, dict):
|
|
return None, None, "geen object"
|
|
|
|
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, None, "ontbrekende of te lange ownerId"
|
|
|
|
label = payload.get("label")
|
|
if label is None:
|
|
label = ""
|
|
if not isinstance(label, str):
|
|
return None, None, "label moet tekst zijn"
|
|
|
|
# Regeleindes eruit: dit is één regel naast een id, en een label met een
|
|
# nieuwe regel erin zou de lijst uit elkaar trekken.
|
|
label = " ".join(label.split()).strip()
|
|
if len(label) > MAX_LABEL_LENGTH:
|
|
return None, None, "label is te lang"
|
|
|
|
return owner_id, label, None
|
|
|
|
|
|
def apply_label(owner_id, label):
|
|
"""Zet of haalt een label weg. Geeft een foutmelding terug, of None.
|
|
|
|
Onder het slot, want dit is lezen-wijzigen-schrijven en de agent bedient
|
|
meerdere verzoeken tegelijk.
|
|
"""
|
|
with labels_lock:
|
|
labels = read_labels()
|
|
if label:
|
|
if owner_id not in labels and len(labels) >= MAX_LABELS:
|
|
return "er zijn al te veel labels"
|
|
labels[owner_id] = label
|
|
else:
|
|
if owner_id not in labels:
|
|
# Niets te doen, en dat is geen fout: de pagina stuurt een leeg
|
|
# label als je het veld leegmaakt, ook als er nog niets stond.
|
|
return None
|
|
del labels[owner_id]
|
|
|
|
try:
|
|
write_labels(labels)
|
|
except OSError as error:
|
|
return "label kon niet worden weggeschreven: " + str(error)
|
|
return 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 _read_payload(self):
|
|
"""Het verzoek als object, of None als er al een fout verstuurd is."""
|
|
try:
|
|
length = int(self.headers.get("Content-Length", "0"))
|
|
except ValueError:
|
|
self._send(400, {"error": "lengte ontbreekt"})
|
|
return None
|
|
|
|
if length <= 0 or length > MAX_BODY_BYTES:
|
|
self._send(400, {"error": "lege of te grote opdracht"})
|
|
return None
|
|
|
|
try:
|
|
return json.loads(self.rfile.read(length).decode("utf-8"))
|
|
except (UnicodeDecodeError, ValueError):
|
|
self._send(400, {"error": "onleesbare opdracht"})
|
|
return None
|
|
|
|
def do_POST(self):
|
|
pad = self.path.rstrip("/")
|
|
|
|
# Een label gaat NIET via de postbus, en dat is de enige uitzondering op
|
|
# die regel. De reden dat opdrachten er wel door gaan, is dat het
|
|
# relay-proces beslist wie er binnen mag en dat twee schrijvers in die
|
|
# allowlist een wedloop zou zijn. Een label zegt niets over toegang, staat
|
|
# in een eigen bestand met de agent als enige schrijver, en is meteen
|
|
# opgeslagen in plaats van na de volgende ronde van de relay.
|
|
if pad in ("/api/label", "/label"):
|
|
payload = self._read_payload()
|
|
if payload is None:
|
|
return
|
|
|
|
owner_id, label, problem = valid_label(payload)
|
|
if problem is not None:
|
|
self._send(400, {"error": problem})
|
|
return
|
|
|
|
problem = apply_label(owner_id, label)
|
|
if problem is not None:
|
|
self._send(500, {"error": problem})
|
|
return
|
|
|
|
self._send(200, {"ownerId": owner_id, "label": label})
|
|
return
|
|
|
|
if pad not in ("/api/command", "/command"):
|
|
self._send(404, {"error": "onbekend pad"})
|
|
return
|
|
|
|
payload = self._read_payload()
|
|
if payload is None:
|
|
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()
|