Files
ARIA-AGENT/satellite/satellite.py
T
duffyduckandClaude Opus 4.8 f9cc10433a feat(satellite): Credential-Store + snmp.ports/info + FritzBox-Reader
Baut SNMP zur generellen Geraete-Auskunft aus und macht Zugangsdaten pro Geraet
hinterlegbar, damit der Satellit tiefer lesen kann.

Satellit:
- Verschluesselter Credential-Store (Fernet) pro IP im Bind-Volume ./data:
  SNMP (community/v2c oder v3 user/auth/priv), HTTP-Basic, FritzBox-Login.
  Neue Messages sat_creds_set/list/delete + *_result; Secrets werden NIE
  zurueckgeliefert (list gibt nur Typen). snmp.*/http/fritzbox nutzen die
  Creds automatisch (params haben Vorrang).
- snmp.ports: ifTable -> aktive/freie Ports, Linkspeed ('sind noch Ports frei').
- snmp.info: sysName/Descr + Entity-MIB Modell/Serial/Firmware (installierte
  Version; Update-Check ist Hersteller-Sache, kein SNMP).
- fritzbox.info/hosts: TR-064 (SOAP+Digest) — Verbindung/Datenrate/externe IP
  bzw. verbundene Geraete. Bewusst als Reader, weil TR-064 fuer on-the-fly
  http.post zu fummelig ist.
- CONTROL_ALLOWLIST-Default + .env.example erweitert; data/ ge-gitignore-t.

RVS: sat_creds_* in ALLOWED_TYPES (sonst verworfen).

Diagnostic: pro entdecktem Geraet ein Schluessel-Button -> Modal (SNMP v2c/v3,
HTTP-Basic, FritzBox); Speichern/Loeschen via RVS an den Satelliten; gesetzte
Typen werden pro Geraet angezeigt (🔑 ✓). Server relayed die Cred-Messages.

Brain: satellite_command-Tool um snmp.ports/info + fritzbox.* erweitert; ARIA
muss keine Passwoerter mitgeben (Satellit nutzt den Store).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-09-19 17:16:06 +02:00

1291 lines
53 KiB
Python

"""
ARIA Satellit — Info-/Gateway-Aussenposten in einem fremden Netz.
Laeuft eigenstaendig (z.B. im Buero) und verbindet sich als RVS-Client in
Stefans Raum (gleicher Token). Gibt ARIA damit Augen + Haende in DIESEM Netz:
Augen: entdeckt Geraete (mDNS/Zeroconf, SSDP/UPnP + DIAL, ARP-Tabelle) und
meldet ein Inventar → sat_devices.
Haende: steuert Geraete (DIAL-App-Launch z.B. YouTube auf Fire TV, Wake-on-
LAN, generisches HTTP) → sat_command / sat_result. Nur wenn
CONTROL_ENABLED=true, Aktion in der Allowlist, alles geloggt.
Adressierung: mehrere Satelliten haengen im selben RVS-Raum. Jeder hat eine
SATELLITE_ID (technisch, eindeutig) + SATELLITE_LOCATION (menschlich, "Buero").
ARIA spricht einen Satelliten ueber seine ID/Location an.
Message-Typen (RVS, Base64/JSON-Relay wie der Rest):
raus: sat_hello {id, location, caps, ts}
sat_devices {requestId, satellite, devices:[...]}
sat_result {requestId, satellite, ok, result|error}
rein: sat_discover {satellite?, requestId}
sat_command {satellite?, requestId, device, action, params}
Sicherheit: reagiert nur auf den eigenen RVS-Raum (Token). Commands brauchen
CONTROL_ENABLED + Allowlist. Discovery ist read-only. Keine offenen Ports.
"""
from __future__ import annotations
import asyncio
import json
import logging
import os
import re
import socket
import struct
import time
from typing import Optional
import websockets
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [satellite] %(levelname)s %(message)s",
)
logger = logging.getLogger("satellite")
def _load_dotenv() -> None:
"""Laedt eine .env neben dem Script (oder im CWD) in os.environ — fuer den
NATIVEN Start (`python satellite.py`). In Docker sind die Variablen via
env_file schon gesetzt; bereits gesetzte Werte gewinnen (werden NICHT
ueberschrieben). Kein python-dotenv noetig."""
here = os.path.dirname(os.path.abspath(__file__))
for path in (os.path.join(here, ".env"), os.path.join(os.getcwd(), ".env")):
if not os.path.isfile(path):
continue
try:
with open(path, "r", encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, _, val = line.partition("=")
key = key.strip()
if key.startswith("export "):
key = key[len("export "):].strip()
val = val.strip()
if val[:1] in ("'", '"'):
# Gequotet: Inhalt bis zum schliessenden Quote, Rest (Kommentar) egal.
q = val[0]
end = val.find(q, 1)
val = val[1:end] if end != -1 else val[1:]
else:
# Ungequotet: Inline-Kommentar (Whitespace + #) abschneiden.
m = re.search(r"\s+#", val)
if m:
val = val[:m.start()]
val = val.strip()
if key and key not in os.environ:
os.environ[key] = val
except Exception as exc:
logging.getLogger("satellite").warning(".env laden fehlgeschlagen (%s): %s", path, exc)
break # erste gefundene .env gewinnt
_load_dotenv()
# ─── Konfiguration ──────────────────────────────────────────────────
def _env_bool(name: str, default: bool) -> bool:
v = os.environ.get(name)
if v is None:
return default
return v.strip().lower() in ("1", "true", "yes", "on", "ja")
def _default_id() -> str:
host = socket.gethostname() or "satellite"
slug = re.sub(r"[^a-zA-Z0-9_-]+", "-", host).strip("-").lower()
return slug or "satellite"
RVS_HOST = os.environ.get("RVS_HOST", "")
RVS_PORT = int(os.environ.get("RVS_PORT", "443") or "443")
RVS_TLS = _env_bool("RVS_TLS", True)
RVS_TOKEN = os.environ.get("RVS_TOKEN", "")
SATELLITE_ID = (os.environ.get("SATELLITE_ID") or _default_id()).strip()
SATELLITE_LOCATION = (os.environ.get("SATELLITE_LOCATION") or SATELLITE_ID).strip()
CONTROL_ENABLED = _env_bool("CONTROL_ENABLED", False)
CONTROL_ALLOWLIST = [
a.strip() for a in
os.environ.get("CONTROL_ALLOWLIST",
"dial.launch,wol,http.get,snmp.get,snmp.walk,snmp.printer,"
"snmp.ports,snmp.info,fritzbox.info,fritzbox.hosts").split(",")
if a.strip()
]
SCAN_INTERVAL_SEC = int(os.environ.get("SCAN_INTERVAL_SEC", "300") or "300")
DISCOVER_TIMEOUT_SEC = float(os.environ.get("DISCOVER_TIMEOUT_SEC", "6") or "6")
DEVICE_CACHE_TTL_SEC = int(os.environ.get("DEVICE_CACHE_TTL_SEC", "120") or "120")
# http.get/http.post: Body-Ausschnitt. Default grosszuegig (ganze Statusseiten
# passen), mit hartem Deckel gegen Riesen-Payloads durchs RVS. offset/max_chars
# pro Request ueberschreibbar; contains-Filter zieht nur relevante Zeilen.
HTTP_TIMEOUT_SEC = float(os.environ.get("HTTP_TIMEOUT_SEC", "10") or "10")
HTTP_MAX_CHARS = int(os.environ.get("HTTP_MAX_CHARS", "20000") or "20000")
HTTP_MAX_CHARS_HARD = int(os.environ.get("HTTP_MAX_CHARS_HARD", "200000") or "200000")
# SNMP (net-snmp-CLI): Default-Community/Version + Timeout. Drucker antworten
# i.d.R. auf community 'public', v2c.
SNMP_COMMUNITY = os.environ.get("SNMP_COMMUNITY", "public") or "public"
SNMP_VERSION = os.environ.get("SNMP_VERSION", "2c") or "2c"
SNMP_TIMEOUT_SEC = float(os.environ.get("SNMP_TIMEOUT_SEC", "5") or "5")
# Printer-MIB (RFC 3805) prtMarkerSuppliesEntry-Spalten (numerisch, ohne MIB-Files):
SNMP_SUPPLY_DESC = "1.3.6.1.2.1.43.11.1.1.6.1" # Beschreibung (z.B. "Black Ink")
SNMP_SUPPLY_MAX = "1.3.6.1.2.1.43.11.1.1.8.1" # Max-Kapazitaet
SNMP_SUPPLY_LVL = "1.3.6.1.2.1.43.11.1.1.9.1" # aktueller Fuellstand
# SNMP-Anreicherung bei der Discovery: jedes entdeckte Geraet mit IP wird kurz
# nach seiner System-Group (RFC 1213) gefragt. Switches/Router/APs/NAS geben so
# Name, Beschreibung, Standort & Uptime preis -> im Inventar (satellite_devices)
# sichtbar. Abschaltbar; kurzer Timeout + parallel, damit der Scan flott bleibt.
SNMP_DISCOVERY = _env_bool("SNMP_DISCOVERY", True)
SNMP_DISCOVERY_CONCURRENCY = int(os.environ.get("SNMP_DISCOVERY_CONCURRENCY", "16") or "16")
SNMP_DISCOVERY_TIMEOUT = float(os.environ.get("SNMP_DISCOVERY_TIMEOUT", "2") or "2")
# System-Group (RFC 1213) .0-Instanzen:
SNMP_SYS_OIDS = {
"descr": "1.3.6.1.2.1.1.1.0", # sysDescr
"objectid": "1.3.6.1.2.1.1.2.0", # sysObjectID
"uptime": "1.3.6.1.2.1.1.3.0", # sysUpTime
"contact": "1.3.6.1.2.1.1.4.0", # sysContact
"name": "1.3.6.1.2.1.1.5.0", # sysName
"location": "1.3.6.1.2.1.1.6.0", # sysLocation
}
# ─── Geraete-Credential-Store (verschluesselt, pro IP) ─────────────
# Diagnostic legt via sat_creds_set pro Geraet Zugangsdaten ab (SNMP-Community/
# v3, HTTP-Basic, FritzBox-Login). Der Satellit nutzt sie automatisch bei snmp.*/
# http/fritzbox. Persistiert verschluesselt (Fernet) in einem Bind-Volume.
CREDS_PATH = os.environ.get("CREDS_PATH", "/data/credentials.json.enc")
CREDS_KEY_PATH = os.environ.get("CREDS_KEY_PATH", "/data/creds.key")
_CREDS: dict = {} # {ip: {snmp:{...}, http:{...}, fritzbox:{...}}}
_creds_fernet = None # Fernet-Instanz (lazy)
def _creds_cipher():
"""Fernet-Instanz; Schluessel aus CREDS_KEY (env) oder Schluesseldatei im
Volume (wird beim ersten Start erzeugt, 0600)."""
global _creds_fernet
if _creds_fernet is not None:
return _creds_fernet
from cryptography.fernet import Fernet
key = os.environ.get("CREDS_KEY", "").strip().encode() or None
if not key:
kp = Path(CREDS_KEY_PATH)
if kp.exists():
key = kp.read_bytes().strip()
else:
key = Fernet.generate_key()
kp.parent.mkdir(parents=True, exist_ok=True)
kp.write_bytes(key)
try:
os.chmod(kp, 0o600)
except OSError:
pass
logger.info("[creds] neuer Verschluesselungs-Schluessel erzeugt: %s", CREDS_KEY_PATH)
_creds_fernet = Fernet(key)
return _creds_fernet
def _creds_load() -> None:
global _CREDS
p = Path(CREDS_PATH)
if not p.exists():
_CREDS = {}
return
try:
blob = p.read_bytes()
raw = _creds_cipher().decrypt(blob)
_CREDS = json.loads(raw.decode("utf-8")) or {}
logger.info("[creds] %d Geraete-Eintraege geladen", len(_CREDS))
except Exception as exc:
logger.warning("[creds] laden fehlgeschlagen (%s) — starte leer", exc)
_CREDS = {}
def _creds_save() -> None:
p = Path(CREDS_PATH)
p.parent.mkdir(parents=True, exist_ok=True)
blob = _creds_cipher().encrypt(json.dumps(_CREDS).encode("utf-8"))
p.write_bytes(blob)
try:
os.chmod(p, 0o600)
except OSError:
pass
def _creds_for(ip: str) -> dict:
return _CREDS.get((ip or "").strip(), {}) if ip else {}
def _creds_public_summary() -> list:
"""Fuer sat_creds_list: welche Geraete welche Cred-Typen haben — OHNE Secrets."""
out = []
for ip, entry in sorted(_CREDS.items()):
types = [t for t in ("snmp", "http", "fritzbox") if entry.get(t)]
out.append({"ip": ip, "types": types})
return out
HEARTBEAT_SEC = 25
# mDNS-Servicetypen, die fuer ARIA interessant sind.
MDNS_TYPES = [
"_googlecast._tcp.local.", # Chromecast / Google TV / Nest
"_airplay._tcp.local.", # Apple TV / AirPlay
"_raop._tcp.local.", # AirPlay-Audio
"_spotify-connect._tcp.local.", # Spotify-Geraete
"_sonos._tcp.local.", # Sonos
"_hap._tcp.local.", # HomeKit
"_printer._tcp.local.", # Drucker
"_ipp._tcp.local.", # Drucker (IPP)
"_smb._tcp.local.", # NAS / Fileshares
"_workstation._tcp.local.", # generische Hosts
"_http._tcp.local.", # Web-UIs (Router, NAS, IoT)
]
CAPABILITIES = ["discover"]
if CONTROL_ENABLED:
CAPABILITIES += CONTROL_ALLOWLIST
# ─── Netz-Kontext / Selbstdiagnose ──────────────────────────────────
def _in_docker_bridge(ip: str) -> bool:
# Docker-Default-Bridge-Range 172.16.0.0/12
try:
a, b = ip.split(".")[:2]
return a == "172" and 16 <= int(b) <= 31
except Exception:
return False
def _net_context() -> dict:
"""Ermittelt in welchem Netz der Satellit LAeUFT — und warnt, wenn das ein
Docker-/NAT-Netz ist (dann erreicht Discovery das echte LAN nicht)."""
ips: list[str] = []
primary = ""
try:
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
s.settimeout(1)
s.connect(("8.8.8.8", 80))
primary = s.getsockname()[0]
s.close()
except Exception:
pass
try:
for info in socket.getaddrinfo(socket.gethostname(), None, socket.AF_INET):
ip = info[4][0]
if ip and not ip.startswith("127.") and ip not in ips:
ips.append(ip)
except Exception:
pass
if primary and primary not in ips:
ips.insert(0, primary)
p = primary or (ips[0] if ips else "")
warning = ""
if p.startswith("192.168.65.") or _in_docker_bridge(p):
warning = (f"Satellit laeuft in einem Docker-/NAT-Netz ({p}), NICHT im echten LAN. "
"mDNS/SSDP erreichen die realen Geraete so nicht. Auf Docker Desktop "
"(Mac/Windows) geht LAN-Discovery nicht — den Satelliten NATIV (python "
"satellite.py) oder auf einem Linux-Host im Ziel-LAN betreiben.")
return {"primary_ip": p, "ips": ips, "warning": warning}
NET = _net_context()
# ─── Discovery ──────────────────────────────────────────────────────
def _discover_mdns(timeout: float) -> list[dict]:
"""Blockierend (im Executor): mDNS/Zeroconf-Sweep ueber MDNS_TYPES."""
out: dict[str, dict] = {}
try:
from zeroconf import Zeroconf, ServiceBrowser
except Exception as exc:
logger.warning("zeroconf nicht verfuegbar: %s", exc)
return []
class _Listener:
def add_service(self, zc, type_, name):
try:
info = zc.get_service_info(type_, name, timeout=2000)
except Exception:
info = None
if not info:
return
ips = []
try:
for addr in info.parsed_addresses():
ips.append(addr)
except Exception:
pass
props = {}
try:
for k, v in (info.properties or {}).items():
try:
props[k.decode("utf-8", "ignore")] = (
v.decode("utf-8", "ignore") if isinstance(v, (bytes, bytearray)) else v)
except Exception:
pass
except Exception:
pass
friendly = name.split("." + type_.split(".", 1)[0])[0].strip(".")
fn = props.get("fn") or props.get("friendlyName") or friendly
dev_id = _slug(f"{fn}-{ips[0] if ips else name}")
out[dev_id] = {
"id": dev_id,
"name": fn,
"type": _mdns_kind(type_),
"ip": ips[0] if ips else "",
"port": info.port,
"via": "mdns",
"service": type_,
"model": props.get("md") or props.get("model") or "",
}
def update_service(self, *a):
pass
def remove_service(self, *a):
pass
zc = None
try:
zc = Zeroconf()
listener = _Listener()
for t in MDNS_TYPES:
try:
ServiceBrowser(zc, t, listener)
except Exception:
pass
time.sleep(timeout)
except Exception as exc:
logger.warning("mDNS-Sweep-Fehler: %s", exc)
finally:
try:
if zc:
zc.close()
except Exception:
pass
return list(out.values())
def _mdns_kind(service_type: str) -> str:
m = {
"_googlecast": "cast", "_airplay": "airplay", "_raop": "airplay-audio",
"_spotify-connect": "spotify", "_sonos": "sonos", "_hap": "homekit",
"_printer": "printer", "_ipp": "printer", "_smb": "fileshare",
"_workstation": "host", "_http": "web",
}
for k, v in m.items():
if service_type.startswith(k):
return v
return "unknown"
def _discover_ssdp(timeout: float) -> list[dict]:
"""Blockierend: SSDP M-SEARCH (UPnP + DIAL). Liefert v.a. Smart-TVs / Fire
TV mit ihrer DIAL Application-URL (fuer App-Launch wie YouTube)."""
out: dict[str, dict] = {}
targets = [
"urn:dial-multiscreen-org:service:dial:1",
"ssdp:all",
]
for st in targets:
msg = (
"M-SEARCH * HTTP/1.1\r\n"
"HOST: 239.255.255.250:1900\r\n"
'MAN: "ssdp:discover"\r\n'
"MX: 2\r\n"
f"ST: {st}\r\n\r\n"
).encode("utf-8")
sock = None
try:
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM, socket.IPPROTO_UDP)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.setsockopt(socket.IPPROTO_IP, socket.IP_MULTICAST_TTL, 2)
sock.settimeout(timeout)
sock.sendto(msg, ("239.255.255.250", 1900))
deadline = time.time() + timeout
while time.time() < deadline:
try:
data, addr = sock.recvfrom(65507)
except socket.timeout:
break
except Exception:
break
headers = _parse_http_headers(data.decode("utf-8", "ignore"))
location = headers.get("location", "")
dial_app = headers.get("application-url", "")
ip = addr[0]
dev = _fetch_upnp_description(location) if location else {}
name = dev.get("name") or headers.get("server", "") or ip
dev_id = _slug(f"{name}-{ip}")
entry = out.get(dev_id, {
"id": dev_id, "name": name, "type": "media-renderer",
"ip": ip, "via": "ssdp",
})
if dev.get("name"):
entry["name"] = dev["name"]
if dev.get("model"):
entry["model"] = dev["model"]
if dev.get("manufacturer"):
entry["manufacturer"] = dev["manufacturer"]
if dial_app or dev.get("dialAppUrl"):
entry["dialAppUrl"] = dial_app or dev.get("dialAppUrl")
entry["type"] = "dial"
out[dev_id] = entry
except Exception as exc:
logger.debug("SSDP (%s) Fehler: %s", st, exc)
finally:
try:
if sock:
sock.close()
except Exception:
pass
return list(out.values())
def _fetch_upnp_description(location: str) -> dict:
try:
import requests
r = requests.get(location, timeout=3)
dial_app = r.headers.get("Application-URL", "")
xml = r.text
name = _xml_tag(xml, "friendlyName")
model = _xml_tag(xml, "modelName")
manuf = _xml_tag(xml, "manufacturer")
return {"name": name, "model": model, "manufacturer": manuf, "dialAppUrl": dial_app}
except Exception:
return {}
def _discover_arp() -> list[dict]:
"""Rohe Host-Liste aus der ARP-Tabelle (kein aktiver Scan)."""
out = []
try:
with open("/proc/net/arp", "r", encoding="utf-8") as f:
lines = f.read().splitlines()[1:]
for ln in lines:
parts = ln.split()
if len(parts) < 4:
continue
ip, _hw, _flags, mac = parts[0], parts[1], parts[2], parts[3]
if mac == "00:00:00:00:00:00":
continue
out.append({
"id": _slug(f"host-{ip}"), "name": ip, "type": "host",
"ip": ip, "mac": mac, "via": "arp",
})
except Exception:
pass
return out
def _merge_devices(*lists) -> list[dict]:
"""Fuehrt Geraetelisten zusammen, dedupt per IP (reichere Quelle gewinnt)."""
by_ip: dict[str, dict] = {}
loose: list[dict] = []
order = {"mdns": 3, "ssdp": 2, "arp": 1}
for lst in lists:
for d in lst:
ip = d.get("ip") or ""
if not ip:
loose.append(d)
continue
cur = by_ip.get(ip)
if not cur:
by_ip[ip] = d
else:
# bessere Quelle / mehr Felder → mergen
merged = {**d, **{k: v for k, v in cur.items() if v}}
if order.get(d.get("via"), 0) >= order.get(cur.get("via"), 0):
merged.update({k: v for k, v in d.items() if v})
# DIAL-URL / mac aus beiden behalten
for key in ("dialAppUrl", "mac", "model", "manufacturer"):
merged[key] = d.get(key) or cur.get(key) or merged.get(key)
by_ip[ip] = {k: v for k, v in merged.items() if v not in (None, "")}
return list(by_ip.values()) + loose
# ─── Control ────────────────────────────────────────────────────────
async def _control(action: str, params: dict, devices: list[dict]) -> dict:
"""Fuehrt eine Steuer-Aktion aus. Guards: CONTROL_ENABLED + Allowlist."""
if not CONTROL_ENABLED:
return {"ok": False, "error": "Steuerung ist an diesem Satelliten deaktiviert (CONTROL_ENABLED=false)."}
if action not in CONTROL_ALLOWLIST:
return {"ok": False, "error": f"Aktion '{action}' nicht erlaubt (Allowlist: {', '.join(CONTROL_ALLOWLIST)})."}
logger.info("[control] %s params=%s", action, {k: str(v)[:60] for k, v in (params or {}).items()})
loop = asyncio.get_event_loop()
try:
if action == "dial.launch":
return await loop.run_in_executor(None, _do_dial_launch, params, devices)
if action == "wol":
return await loop.run_in_executor(None, _do_wol, params)
if action in ("http.get", "http.post"):
return await loop.run_in_executor(None, _do_http, action, params)
if action in ("snmp.get", "snmp.walk"):
return await loop.run_in_executor(None, _do_snmp, action, params)
if action == "snmp.printer":
return await loop.run_in_executor(None, _do_snmp_printer, params)
if action == "snmp.ports":
return await loop.run_in_executor(None, _do_snmp_ports, params)
if action == "snmp.info":
return await loop.run_in_executor(None, _do_snmp_info, params)
if action in ("fritzbox.info", "fritzbox.hosts"):
return await loop.run_in_executor(None, _do_fritzbox, action, params)
return {"ok": False, "error": f"Aktion '{action}' nicht implementiert."}
except Exception as exc:
return {"ok": False, "error": f"{action} fehlgeschlagen: {exc}"}
def _find_device(devices: list[dict], ref: str) -> Optional[dict]:
ref = (ref or "").strip().lower()
if not ref:
return None
for d in devices:
if d.get("id", "").lower() == ref or d.get("ip", "") == ref:
return d
for d in devices:
if ref in (d.get("name", "").lower()):
return d
return None
def _do_dial_launch(params: dict, devices: list[dict]) -> dict:
"""DIAL-App-Launch, z.B. YouTube-Video auf Fire TV / Smart-TV.
params: {device, app='YouTube', v=<videoId> (oder beliebige app-params)}"""
import requests
ref = params.get("device") or ""
dev = _find_device(devices, ref)
app_url = (dev or {}).get("dialAppUrl") if dev else params.get("dialAppUrl")
if not app_url:
return {"ok": False, "error": f"Kein DIAL-Geraet fuer '{ref}' gefunden (oder keine Application-URL)."}
app = params.get("app") or "YouTube"
# app-Parameter (alles ausser device/app) als form-urlencoded Body.
body = {k: v for k, v in (params or {}).items() if k not in ("device", "app", "dialAppUrl")}
url = app_url.rstrip("/") + "/" + app
r = requests.post(url, data=body, timeout=5)
ok = r.status_code in (200, 201)
return {"ok": ok, "result": f"DIAL {app}{(dev or {}).get('name', ref)} (HTTP {r.status_code})"
if ok else None,
"error": None if ok else f"DIAL-Launch HTTP {r.status_code}: {r.text[:120]}"}
def _do_wol(params: dict) -> dict:
mac = (params.get("mac") or "").strip()
if not re.match(r"^([0-9A-Fa-f]{2}[:-]){5}[0-9A-Fa-f]{2}$", mac):
return {"ok": False, "error": f"Ungueltige MAC: {mac!r}"}
clean = re.sub(r"[:-]", "", mac)
packet = b"\xff" * 6 + bytes.fromhex(clean) * 16
s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
s.setsockopt(socket.SOL_SOCKET, socket.SO_BROADCAST, 1)
s.sendto(packet, ("255.255.255.255", 9))
s.close()
return {"ok": True, "result": f"Wake-on-LAN an {mac} gesendet."}
def _do_http(action: str, params: dict) -> dict:
"""HTTP-GET/POST vom Satelliten aus (lokale Webhooks, Geraete-Statusseiten …).
params:
url Pflicht (http/https).
body/headers optional (POST).
offset ab welchem Zeichen der Body zurueckgegeben wird (Default 0).
max_chars wie viele Zeichen max. (Default HTTP_MAX_CHARS, hart gedeckelt).
contains String oder Liste: nur Zeilen, die (case-insensitive) einen der
Begriffe enthalten, werden zurueckgegeben. Ideal um aus einer
grossen Statusseite nur die relevanten Werte (z.B. Tinte) zu
ziehen, ohne die ganze Seite zu paginieren.
Antwort enthaelt total_chars + truncated, damit der Aufrufer weiss, ob noch
mehr da ist."""
import requests
url = params.get("url") or ""
if not url.startswith(("http://", "https://")):
return {"ok": False, "error": "url (http/https) erforderlich."}
method = "GET" if action == "http.get" else "POST"
# HTTP-Basic-Auth: explizite params > gespeicherte http-Creds fuer den Host.
auth = None
hcreds = {}
try:
from urllib.parse import urlparse
host = urlparse(url).hostname or ""
hcreds = _creds_for(host).get("http", {})
except Exception:
pass
user = params.get("user") or hcreds.get("user")
pw = params.get("pass") or params.get("password") or hcreds.get("pass")
if user:
auth = (str(user), str(pw or ""))
r = requests.request(method, url, data=params.get("body"),
headers=params.get("headers"), auth=auth,
timeout=HTTP_TIMEOUT_SEC)
text = r.text
total = len(text)
contains = params.get("contains")
if contains:
terms = [contains] if isinstance(contains, str) else list(contains)
terms = [str(t).lower() for t in terms if str(t).strip()]
if terms:
lines = [ln for ln in text.splitlines()
if any(t in ln.lower() for t in terms)]
text = "\n".join(lines)
try:
offset = max(0, int(params.get("offset", 0)))
except (TypeError, ValueError):
offset = 0
try:
max_chars = int(params.get("max_chars", HTTP_MAX_CHARS))
except (TypeError, ValueError):
max_chars = HTTP_MAX_CHARS
max_chars = max(1, min(max_chars, HTTP_MAX_CHARS_HARD))
body = text[offset:offset + max_chars]
returned_end = offset + len(body)
truncated = returned_end < len(text)
return {"ok": True, "result": {
"status": r.status_code,
"body": body,
"total_chars": total, # Groesse der Roh-Antwort
"filtered": bool(contains), # contains-Filter aktiv?
"offset": offset,
"returned_chars": len(body),
"truncated": truncated, # noch mehr Text nach diesem Ausschnitt?
}}
def _snmp_run(args: list, timeout: float) -> tuple:
"""Fuehrt ein net-snmp-CLI-Tool aus. Gibt (ok, stdout|fehlertext)."""
import subprocess
try:
r = subprocess.run(args, capture_output=True, text=True, timeout=timeout)
except FileNotFoundError:
return False, "snmp-Tools fehlen im Container (Paket 'snmp' im Dockerfile)."
except subprocess.TimeoutExpired:
return False, "SNMP-Timeout — Geraet antwortet nicht (community/version/IP pruefen)."
if r.returncode != 0:
return False, (r.stderr or r.stdout or "SNMP-Fehler").strip()[:200]
return True, r.stdout
def _snmp_base_args(params: dict, ip: str = "") -> list:
"""Version/Community bzw. v3-Auth. Prioritaet: explizite params > gespeicherte
Creds fuer die IP > globale Defaults. OHNE -t/-r (haengt der Aufrufer an)."""
creds = _creds_for(ip).get("snmp", {}) if ip else {}
version = str(params.get("version") or creds.get("version") or SNMP_VERSION)
if version == "3":
v3 = creds.get("v3", {}) or {}
user = str(params.get("user") or v3.get("user") or "")
level = str(params.get("level") or v3.get("level") or "authPriv")
args = ["-v", "3", "-u", user, "-l", level]
ap = params.get("authProto") or v3.get("authProto")
ak = params.get("authKey") or v3.get("authKey")
pp = params.get("privProto") or v3.get("privProto")
pk = params.get("privKey") or v3.get("privKey")
if ap and ak:
args += ["-a", str(ap), "-A", str(ak)]
if pp and pk:
args += ["-x", str(pp), "-X", str(pk)]
return args
community = str(params.get("community") or creds.get("community") or SNMP_COMMUNITY)
return ["-v", version, "-c", community]
def _snmp_target(params: dict) -> str:
return (params.get("ip") or params.get("host") or params.get("device") or "").strip()
def _do_snmp(action: str, params: dict) -> dict:
"""Generisches snmp.get / snmp.walk.
params: {ip|host, oid, community?='public', version?='2c'}."""
ip = _snmp_target(params)
if not ip:
return {"ok": False, "error": "ip/host erforderlich."}
oid = str(params.get("oid") or "").strip()
if not oid:
return {"ok": False, "error": "oid erforderlich (z.B. 1.3.6.1.2.1.1.5.0 fuer sysName)."}
tool = "snmpwalk" if action == "snmp.walk" else "snmpget"
# -OQ: OID = Wert, ohne Typannotation; numerische OIDs brauchen keine MIB-Files.
args = [tool, "-OQ", *_snmp_base_args(params, ip), "-t", "2", "-r", "1", ip, oid]
ok, out = _snmp_run(args, SNMP_TIMEOUT_SEC)
if not ok:
return {"ok": False, "error": out}
lines = [ln.strip() for ln in out.splitlines() if ln.strip()]
return {"ok": True, "result": {"ip": ip, "oid": oid, "lines": lines[:200]}}
def _snmp_walk_values(ip: str, base: list, oid: str) -> list:
"""snmpwalk -Oqv (nur Werte, in OID-Index-Reihenfolge)."""
ok, out = _snmp_run(["snmpwalk", "-Oqv", *base, ip, oid], SNMP_TIMEOUT_SEC)
if not ok:
return []
return [ln.strip().strip('"') for ln in out.splitlines() if ln.strip()]
def _do_snmp_printer(params: dict) -> dict:
"""Komfort: liest die Verbrauchsmaterialien (Tinte/Toner) aus der Printer-MIB
und rechnet Fuellstaende in Prozent. params: {ip|host, community?, version?}."""
ip = _snmp_target(params)
if not ip:
return {"ok": False, "error": "ip/host erforderlich."}
base = [*_snmp_base_args(params, ip), "-t", "2", "-r", "1"]
descs = _snmp_walk_values(ip, base, SNMP_SUPPLY_DESC)
if not descs:
return {"ok": False, "error":
"Keine Printer-MIB-Daten (Geraet unterstuetzt kein SNMP, falsche "
"community/version, oder es ist kein Drucker)."}
lvls = _snmp_walk_values(ip, base, SNMP_SUPPLY_LVL)
maxs = _snmp_walk_values(ip, base, SNMP_SUPPLY_MAX)
supplies = []
for i, name in enumerate(descs):
lvl = _to_int(lvls[i]) if i < len(lvls) else None
mx = _to_int(maxs[i]) if i < len(maxs) else None
percent = None
if lvl is not None and mx and mx > 0 and lvl >= 0:
percent = round(lvl / mx * 100)
elif lvl == -3:
percent = "vorhanden (Stand unbekannt)" # RFC: some remaining
elif lvl in (-1, -2):
percent = "unbekannt"
supplies.append({"name": name, "level": lvl, "max": mx, "percent": percent})
return {"ok": True, "result": {"ip": ip, "supplies": supplies}}
# ifTable (RFC 1213) Spalten:
_IF_DESCR = "1.3.6.1.2.1.2.2.1.2"
_IF_TYPE = "1.3.6.1.2.1.2.2.1.3"
_IF_SPEED = "1.3.6.1.2.1.2.2.1.5"
_IF_ADMIN = "1.3.6.1.2.1.2.2.1.7" # up(1) down(2)
_IF_OPER = "1.3.6.1.2.1.2.2.1.8" # up(1) down(2) ...
_IF_ALIAS = "1.3.6.1.2.1.31.1.1.1.18" # ifAlias (ifXTable, optional)
def _do_snmp_ports(params: dict) -> dict:
"""Interface-Uebersicht eines Switches/Routers: welche Ports sind aktiv (Link),
welche frei. params: {ip|host, community?/v3?}. ethernetCsmacd(6)=echte Ports;
Loopback/VLAN etc. werden als 'other' markiert, nicht als freier Port gezaehlt."""
ip = _snmp_target(params)
if not ip:
return {"ok": False, "error": "ip/host erforderlich."}
base = [*_snmp_base_args(params, ip), "-t", "2", "-r", "1"]
descr = _snmp_walk_values(ip, base, _IF_DESCR)
if not descr:
return {"ok": False, "error":
"Keine Interface-Daten (kein SNMP / falsche Credentials / kein Switch)."}
types = _snmp_walk_values(ip, base, _IF_TYPE)
opers = _snmp_walk_values(ip, base, _IF_OPER)
admins = _snmp_walk_values(ip, base, _IF_ADMIN)
speeds = _snmp_walk_values(ip, base, _IF_SPEED)
aliases = _snmp_walk_values(ip, base, _IF_ALIAS)
ports = []
up = down_free = disabled = 0
for i, name in enumerate(descr):
itype = _to_int(types[i]) if i < len(types) else None
oper = _to_int(opers[i]) if i < len(opers) else None
admin = _to_int(admins[i]) if i < len(admins) else None
speed = _to_int(speeds[i]) if i < len(speeds) else None
is_eth = (itype == 6) # ethernetCsmacd
state = ("up" if oper == 1 else
"disabled" if admin == 2 else "down")
if is_eth:
if state == "up":
up += 1
elif state == "disabled":
disabled += 1
else:
down_free += 1
ports.append({
"name": name.strip('"'),
"alias": (aliases[i].strip('"') if i < len(aliases) else ""),
"physical": is_eth,
"state": state,
"speedMbps": round(speed / 1_000_000) if speed else None,
})
return {"ok": True, "result": {
"ip": ip,
"summary": {"physical_ports": up + down_free + disabled,
"up": up, "free": down_free, "disabled": disabled},
"ports": ports,
}}
# entPhysicalTable (RFC 4133) — Modell/Serie/Firmware:
_ENT_MODEL = "1.3.6.1.2.1.47.1.1.1.1.13" # entPhysicalModelName
_ENT_SERIAL = "1.3.6.1.2.1.47.1.1.1.1.11" # entPhysicalSerialNum
_ENT_SWREV = "1.3.6.1.2.1.47.1.1.1.1.10" # entPhysicalSoftwareRev
_ENT_FWREV = "1.3.6.1.2.1.47.1.1.1.1.9" # entPhysicalFirmwareRev
def _do_snmp_info(params: dict) -> dict:
"""Geraeteinfo: sysName/sysDescr + (falls vorhanden) Modell, Seriennummer,
Firmware-/Software-Version aus der Entity-MIB. Sagt die INSTALLIERTE Version —
ob ein Update existiert, weiss SNMP nicht (Hersteller-Sache)."""
ip = _snmp_target(params)
if not ip:
return {"ok": False, "error": "ip/host erforderlich."}
sysinfo = _snmp_system(ip, str(params.get("community") or ""), str(params.get("version") or ""))
base = [*_snmp_base_args(params, ip), "-t", "2", "-r", "1"]
def _first(oid):
vals = [v for v in _snmp_walk_values(ip, base, oid)
if v and "No Such" not in v]
return vals[0] if vals else None
result = {
"ip": ip,
"name": (sysinfo or {}).get("name"),
"descr": (sysinfo or {}).get("descr"),
"location": (sysinfo or {}).get("location"),
"uptime": (sysinfo or {}).get("uptime"),
"model": _first(_ENT_MODEL),
"serial": _first(_ENT_SERIAL),
"firmware": _first(_ENT_FWREV) or _first(_ENT_SWREV),
}
if not any(result[k] for k in ("name", "descr", "model", "firmware")):
return {"ok": False, "error": "Kein SNMP / keine verwertbaren Infos."}
return {"ok": True, "result": result}
# ─── FritzBox (TR-064) ─────────────────────────────────────────────
# TR-064 ist SOAP+Digest-Auth — zu fummelig fuer on-the-fly http.post, daher ein
# schlanker Reader. Braucht FritzBox-Login (Credential-Store, Typ 'fritzbox').
def _tr064(ip: str, user: str, pw: str, service: str, control: str,
action: str, args: Optional[dict] = None) -> dict:
"""Ein TR-064-SOAP-Call. Gibt {ok, fields|error}. fields = alle <NewX>-Tags."""
import requests
from requests.auth import HTTPDigestAuth
body = "".join(f"<{k}>{v}</{k}>" for k, v in (args or {}).items())
envelope = (
'<?xml version="1.0"?>'
'<s:Envelope xmlns:s="http://schemas.xmlsoap.org/soap/envelope/" '
's:encodingStyle="http://schemas.xmlsoap.org/soap/encoding/"><s:Body>'
f'<u:{action} xmlns:u="{service}">{body}</u:{action}>'
'</s:Body></s:Envelope>'
)
url = f"http://{ip}:49000{control}"
try:
r = requests.post(url, data=envelope.encode("utf-8"),
headers={"Content-Type": 'text/xml; charset="utf-8"',
"SOAPAction": f"{service}#{action}"},
auth=HTTPDigestAuth(user, pw), timeout=HTTP_TIMEOUT_SEC)
except Exception as exc:
return {"ok": False, "error": f"TR-064 nicht erreichbar: {exc}"}
if r.status_code == 401:
return {"ok": False, "error": "TR-064 Auth fehlgeschlagen (FritzBox-Login pruefen)."}
if r.status_code != 200:
return {"ok": False, "error": f"TR-064 HTTP {r.status_code}"}
fields = {m.group(1): m.group(2) for m in
re.finditer(r"<(New[^>/]+)>(.*?)</\1>", r.text, re.DOTALL)}
return {"ok": True, "fields": fields}
def _do_fritzbox(action: str, params: dict) -> dict:
"""fritzbox.info -> Modell/Firmware/Verbindung/externe IP/Datenrate.
fritzbox.hosts -> Liste der bekannten Geraete (Name/IP/MAC/aktiv)."""
ip = _snmp_target(params)
if not ip:
return {"ok": False, "error": "ip/host erforderlich."}
fb = _creds_for(ip).get("fritzbox", {})
user = str(params.get("user") or fb.get("user") or "")
pw = str(params.get("pass") or params.get("password") or fb.get("pass") or "")
if not pw:
return {"ok": False, "error":
"Kein FritzBox-Login hinterlegt. In der Geraeteliste Credentials "
"(Typ 'fritzbox') fuer diese IP setzen."}
if action == "fritzbox.hosts":
p = _tr064(ip, user, pw, "urn:dslforum-org:service:Hosts:1",
"/upnp/control/hosts", "X_AVM-DE_GetHostListPath")
if not p.get("ok"):
return p
path = p["fields"].get("NewX_AVM-DE_HostListPath", "")
if not path:
return {"ok": False, "error": "FritzBox lieferte keinen Host-Listen-Pfad."}
import requests
from requests.auth import HTTPDigestAuth
try:
r = requests.get(f"http://{ip}:49000{path}",
auth=HTTPDigestAuth(user, pw), timeout=HTTP_TIMEOUT_SEC)
except Exception as exc:
return {"ok": False, "error": f"Host-Liste nicht abrufbar: {exc}"}
hosts = []
for item in re.finditer(r"<Item>(.*?)</Item>", r.text, re.DOTALL):
blk = item.group(1)
def _t(tag):
m = re.search(rf"<{tag}>(.*?)</{tag}>", blk, re.DOTALL)
return m.group(1) if m else ""
hosts.append({"name": _t("HostName"), "ip": _t("IPAddress"),
"mac": _t("MACAddress"),
"active": _t("Active") in ("1", "true")})
return {"ok": True, "result": {"ip": ip, "count": len(hosts), "hosts": hosts}}
# fritzbox.info (Default): mehrere Services, Teil-Fehler tolerieren.
info = {"ip": ip}
dev = _tr064(ip, user, pw, "urn:dslforum-org:service:DeviceInfo:1",
"/upnp/control/deviceinfo", "GetInfo")
if dev.get("ok"):
f = dev["fields"]
info.update({"model": f.get("NewModelName"), "firmware": f.get("NewSoftwareVersion"),
"serial": f.get("NewSerialNumber"), "uptime_s": _to_int(f.get("NewUpTime"))})
st = _tr064(ip, user, pw, "urn:dslforum-org:service:WANIPConnection:1",
"/upnp/control/wanipconnection1", "GetStatusInfo")
if st.get("ok"):
info["connection"] = st["fields"].get("NewConnectionStatus")
info["connection_uptime_s"] = _to_int(st["fields"].get("NewUptime"))
ext = _tr064(ip, user, pw, "urn:dslforum-org:service:WANIPConnection:1",
"/upnp/control/wanipconnection1", "GetExternalIPAddress")
if ext.get("ok"):
info["external_ip"] = ext["fields"].get("NewExternalIPAddress")
link = _tr064(ip, user, pw, "urn:dslforum-org:service:WANCommonInterfaceConfig:1",
"/upnp/control/wancommonifconfig1", "GetCommonLinkProperties")
if link.get("ok"):
f = link["fields"]
dn = _to_int(f.get("NewLayer1DownstreamMaxBitRate"))
upr = _to_int(f.get("NewLayer1UpstreamMaxBitRate"))
info["downstream_mbit"] = round(dn / 1_000_000, 1) if dn else None
info["upstream_mbit"] = round(upr / 1_000_000, 1) if upr else None
info["physical_link"] = f.get("NewPhysicalLinkStatus")
if len(info) == 1:
return {"ok": False, "error":
"FritzBox antwortet nicht auf TR-064 (Login/Rechte pruefen; TR-064 in "
"der FritzBox unter Heimnetz > Netzwerkeinstellungen aktivieren)."}
return {"ok": True, "result": info}
def _to_int(s: str):
try:
return int(str(s).strip())
except (TypeError, ValueError):
return None
def _snmp_system(ip: str, community: str = "", version: str = "") -> Optional[dict]:
"""Fragt die SNMP-System-Group eines Hosts ab (ein snmpget, alle 6 OIDs).
Gibt {descr,name,contact,location,uptime,objectid} oder None (kein SNMP).
Nutzt gespeicherte Creds fuer die IP; kurzer Timeout, keine Retries ->
Nicht-SNMP-Hosts scheitern schnell."""
params = {}
if community:
params["community"] = community
if version:
params["version"] = version
base = [*_snmp_base_args(params, ip), "-t", "1", "-r", "0"]
keys = list(SNMP_SYS_OIDS.keys())
ok, out = _snmp_run(["snmpget", "-Oqv", *base, ip, *SNMP_SYS_OIDS.values()],
SNMP_DISCOVERY_TIMEOUT)
if not ok:
return None
vals = out.splitlines()
info = {}
for k, v in zip(keys, vals):
v = (v or "").strip().strip('"')
if v and "No Such" not in v and "No more" not in v:
info[k] = v
return info or None
def _snmp_kind(descr: str) -> str:
"""Grobe Geraeteklasse aus sysDescr (fuer type im Inventar)."""
d = (descr or "").lower()
if any(k in d for k in ("switch", "catalyst", "procurve", "aruba", "powerconnect")):
return "switch"
if any(k in d for k in ("router", "mikrotik", "routeros", "edgeos", "openwrt", "pfsense", "fritz!box")):
return "router"
if any(k in d for k in ("access point", "accesspoint", "unifi", "wifi", "wlan")):
return "access-point"
if any(k in d for k in ("printer", "laserjet", "officejet", "brother", "epson", "kyocera")):
return "printer"
if any(k in d for k in ("nas", "synology", "qnap", "truenas", "diskstation")):
return "nas"
if any(k in d for k in ("ups", "usv", "smart-ups")):
return "ups"
return ""
# ─── Helpers ────────────────────────────────────────────────────────
def _slug(s: str) -> str:
s = (s or "").strip().lower()
s = re.sub(r"[^a-z0-9]+", "-", s).strip("-")
return s or "dev"
def _parse_http_headers(text: str) -> dict:
headers = {}
for line in text.split("\r\n")[1:]:
if ":" in line:
k, _, v = line.partition(":")
headers[k.strip().lower()] = v.strip()
return headers
def _xml_tag(xml: str, tag: str) -> str:
m = re.search(rf"<{tag}>(.*?)</{tag}>", xml, re.IGNORECASE | re.DOTALL)
return m.group(1).strip() if m else ""
# ─── Satellit (RVS-Client) ──────────────────────────────────────────
class Satellite:
def __init__(self) -> None:
self.ws: Optional[websockets.WebSocketClientProtocol] = None
self._devices: list[dict] = []
self._devices_ts: float = 0.0
self._scanning = False
async def _scan(self, force: bool = False) -> list[dict]:
fresh = (time.time() - self._devices_ts) < DEVICE_CACHE_TTL_SEC
if self._devices and fresh and not force:
return self._devices
if self._scanning:
# Laufenden Scan abwarten (grob)
for _ in range(30):
await asyncio.sleep(0.2)
if not self._scanning:
break
return self._devices
self._scanning = True
try:
loop = asyncio.get_event_loop()
mdns = await loop.run_in_executor(None, _discover_mdns, DISCOVER_TIMEOUT_SEC)
ssdp = await loop.run_in_executor(None, _discover_ssdp, DISCOVER_TIMEOUT_SEC)
arp = await loop.run_in_executor(None, _discover_arp)
self._devices = _merge_devices(mdns, ssdp, arp)
if SNMP_DISCOVERY:
await self._enrich_snmp(self._devices)
self._devices_ts = time.time()
n_snmp = sum(1 for d in self._devices if d.get("snmpCapable"))
logger.info("[scan] %d Geraete (mdns=%d ssdp=%d arp=%d, snmp=%d)",
len(self._devices), len(mdns), len(ssdp), len(arp), n_snmp)
finally:
self._scanning = False
return self._devices
async def _enrich_snmp(self, devices: list) -> None:
"""Fragt jedes Geraet mit IP parallel per SNMP-System-Group ab und haengt
die Infos an. Verbessert Name/Typ, wenn bisher nur eine IP bekannt war."""
loop = asyncio.get_event_loop()
sem = asyncio.Semaphore(SNMP_DISCOVERY_CONCURRENCY)
async def _one(dev: dict) -> None:
ip = dev.get("ip") or ""
if not ip:
return
async with sem:
info = await loop.run_in_executor(None, _snmp_system, ip)
if not info:
return
dev["snmp"] = info
dev["snmpCapable"] = True
# Name aufwerten, wenn er bisher nur die IP/leer war.
if info.get("name") and dev.get("name", "") in ("", ip):
dev["name"] = info["name"]
# Typ aufwerten, wenn bisher generisch (host/leer).
kind = _snmp_kind(info.get("descr", ""))
if kind and dev.get("type", "") in ("", "host"):
dev["type"] = kind
await asyncio.gather(*(_one(d) for d in devices))
async def _send(self, message: dict) -> None:
if self.ws is None:
return
try:
await self.ws.send(json.dumps(message))
except Exception as exc:
logger.warning("send fehlgeschlagen: %s", exc)
async def _hello(self, log: bool = False) -> None:
await self._send({
"type": "sat_hello",
"payload": {
"id": SATELLITE_ID,
"location": SATELLITE_LOCATION,
"caps": CAPABILITIES,
"control": CONTROL_ENABLED,
"net": NET,
},
"timestamp": int(time.time() * 1000),
})
if log:
logger.info("sat_hello gesendet: id=%s location=%s caps=%s",
SATELLITE_ID, SATELLITE_LOCATION, CAPABILITIES)
def _for_me(self, payload: dict) -> bool:
target = (payload.get("satellite") or "").strip().lower()
if not target or target in ("all", "*"):
return True
return target in (SATELLITE_ID.lower(), SATELLITE_LOCATION.lower())
async def _handle(self, raw: str) -> None:
try:
msg = json.loads(raw)
except Exception:
return
mtype = msg.get("type", "")
payload = msg.get("payload", {}) or {}
if mtype == "sat_discover":
if not self._for_me(payload):
return
req_id = payload.get("requestId", "")
devices = await self._scan(force=bool(payload.get("force")))
await self._send({
"type": "sat_devices",
"payload": {"requestId": req_id, "satellite": SATELLITE_ID,
"location": SATELLITE_LOCATION, "devices": devices,
"net": NET},
"timestamp": int(time.time() * 1000),
})
elif mtype == "sat_command":
if not self._for_me(payload):
return
req_id = payload.get("requestId", "")
action = (payload.get("action") or "").strip()
params = payload.get("params") or {}
if payload.get("device") and "device" not in params:
params["device"] = payload.get("device")
# Geraeteliste nur scannen, wenn die Aktion sie wirklich braucht
# (dial.launch loest ein Geraet auf, oder es wurde ein device-Ref
# mitgegeben). http.get/http.post/wol arbeiten direkt mit url/mac —
# ein voller LAN-Scan davor kostete nur unnoetig viele Sekunden.
needs_devices = action == "dial.launch" or bool(params.get("device"))
devices = await self._scan() if needs_devices else self._devices
result = await _control(action, params, devices)
await self._send({
"type": "sat_result",
"payload": {"requestId": req_id, "satellite": SATELLITE_ID, **result},
"timestamp": int(time.time() * 1000),
})
elif mtype == "sat_creds_set":
if not self._for_me(payload):
return
ip = (payload.get("ip") or "").strip()
creds = payload.get("creds") or {}
ok = False
if ip and isinstance(creds, dict):
entry = _CREDS.setdefault(ip, {})
for t in ("snmp", "http", "fritzbox"):
if t in creds:
if creds[t]: # leeres Objekt = Typ loeschen
entry[t] = creds[t]
else:
entry.pop(t, None)
if not entry:
_CREDS.pop(ip, None)
try:
_creds_save()
ok = True
except Exception as exc:
logger.warning("[creds] speichern fehlgeschlagen: %s", exc)
await self._send({"type": "sat_creds_result",
"payload": {"requestId": payload.get("requestId", ""),
"satellite": SATELLITE_ID, "ok": ok, "ip": ip},
"timestamp": int(time.time() * 1000)})
elif mtype == "sat_creds_delete":
if not self._for_me(payload):
return
ip = (payload.get("ip") or "").strip()
ctype = (payload.get("type") or "").strip()
if ip in _CREDS:
if ctype:
_CREDS[ip].pop(ctype, None)
if not _CREDS[ip]:
_CREDS.pop(ip, None)
else:
_CREDS.pop(ip, None)
try:
_creds_save()
except Exception as exc:
logger.warning("[creds] speichern fehlgeschlagen: %s", exc)
await self._send({"type": "sat_creds_result",
"payload": {"requestId": payload.get("requestId", ""),
"satellite": SATELLITE_ID, "ok": True, "ip": ip},
"timestamp": int(time.time() * 1000)})
elif mtype == "sat_creds_list":
if not self._for_me(payload):
return
await self._send({"type": "sat_creds_list_result",
"payload": {"requestId": payload.get("requestId", ""),
"satellite": SATELLITE_ID,
"location": SATELLITE_LOCATION,
"items": _creds_public_summary()},
"timestamp": int(time.time() * 1000)})
async def _periodic_scan(self) -> None:
while True:
try:
await self._scan(force=True)
except Exception as exc:
logger.warning("periodischer Scan-Fehler: %s", exc)
await asyncio.sleep(SCAN_INTERVAL_SEC)
async def _heartbeat(self) -> None:
# Re-announce bei jedem Heartbeat: falls die Bridge NACH uns (neu)
# verbindet, lernt sie uns so innerhalb von HEARTBEAT_SEC — RVS replayt
# nichts. Haelt zugleich last_seen in der Bridge-Registry frisch.
while True:
await asyncio.sleep(HEARTBEAT_SEC)
await self._send({"type": "heartbeat", "timestamp": int(time.time() * 1000)})
await self._hello()
async def run(self) -> None:
if not RVS_HOST or not RVS_TOKEN:
logger.error("RVS_HOST und RVS_TOKEN sind Pflicht (siehe .env.example).")
return
_creds_load()
asyncio.create_task(self._periodic_scan())
backoff = 1
while True:
proto = "wss" if RVS_TLS else "ws"
url = f"{proto}://{RVS_HOST}:{RVS_PORT}?token={RVS_TOKEN}"
try:
logger.info("Verbinde mit RVS %s://%s:%s …", proto, RVS_HOST, RVS_PORT)
async with websockets.connect(url, max_size=8 * 1024 * 1024,
ping_interval=20, ping_timeout=20) as ws:
self.ws = ws
backoff = 1
await self._hello(log=True)
hb = asyncio.create_task(self._heartbeat())
try:
async for raw in ws:
await self._handle(raw)
finally:
hb.cancel()
except Exception as exc:
logger.warning("RVS-Verbindung verloren: %s", exc)
finally:
self.ws = None
await asyncio.sleep(backoff)
backoff = min(backoff * 2, 30)
def main() -> None:
logger.info("ARIA Satellit startet — id=%s location=%s control=%s",
SATELLITE_ID, SATELLITE_LOCATION, CONTROL_ENABLED)
logger.info("Netz: primary_ip=%s alle=%s", NET.get("primary_ip"), NET.get("ips"))
if NET.get("warning"):
logger.warning("⚠ %s", NET["warning"])
try:
asyncio.run(Satellite().run())
except KeyboardInterrupt:
pass
if __name__ == "__main__":
main()