Files
Athena-Deck/network/agent.py
T

305 lines
13 KiB
Python

"""Container-only privileged helper. Never run on the host.
Owns one wg interface, one policy-routing table, and two ingress listeners.
The web child runs without root; the helper accepts only fixed operations.
"""
import http.client
import json
import os
import secrets
import signal
import socket
import socketserver
import subprocess
import sys
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from auth import normalize, validate_record
from network.config import parse_config, wireguard_text
from network.policy import Policy
DATA = Path('/data')
RUN = Path('/run/deck')
INTERFACE = 'deckwg0'
PORT = 8110
LOCK = threading.RLock()
PROXY_TOKEN = secrets.token_hex(32)
def command(*args, input=None, check=True):
result = subprocess.run(args, input=input, text=True, capture_output=True, timeout=10)
if check and result.returncode:
raise ValueError('WireGuard-Netzwerkaktion fehlgeschlagen; Einstellungen prüfen.')
return result.stdout.strip()
def save(path, value, mode=0o600):
temporary = path.with_suffix('.tmp')
fd = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, mode)
with os.fdopen(fd, 'w') as stream:
json.dump(value, stream)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, path)
class DeviceServer(ThreadingHTTPServer):
daemon_threads = True
def __init__(self, address, device, ingress):
self.device, self.ingress = device, ingress
super().__init__((address, PORT), Proxy, bind_and_activate=False)
try:
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_BINDTODEVICE, device.encode()+b'\0')
self.server_bind()
self.server_activate()
except Exception:
self.server_close()
raise
class Proxy(BaseHTTPRequestHandler):
def log_message(self, *args):
pass
def do_GET(self):
self.forward()
def do_POST(self):
self.forward()
def forward(self):
self.connection.settimeout(10)
with LOCK:
allowed = AGENT.policy.allowed(self.server.ingress)
if not allowed:
self.send_error(403, 'This network access is disabled')
return
try:
allowed_hosts = {f'{self.server.server_address[0]}:{PORT}'}
if self.server.ingress == 'lan':
allowed_hosts.update({f'{AGENT.lan_ip}:{PORT}', f'127.0.0.1:{PORT}', f'localhost:{PORT}'})
if self.headers.get('Host') not in allowed_hosts:
self.send_error(403, 'Host rejected')
return
if self.headers.get('Transfer-Encoding'):
raise ValueError()
length = int(self.headers.get('Content-Length', '0'))
if not 0 <= length <= 32768 or len(self.path)>2048:
raise ValueError()
body = self.rfile.read(length) if length else None
headers = {name: self.headers[name] for name in ('Host', 'Origin', 'Content-Type', 'Cookie', 'Authorization', 'X-Athena-Deck') if name in self.headers}
headers.update({'X-Deck-Proxy': PROXY_TOKEN, 'X-Deck-Ingress': self.server.ingress})
conn = http.client.HTTPConnection('127.0.0.1', 8108, timeout=30)
try:
conn.request(self.command, self.path, body, headers)
response = conn.getresponse()
data = response.read(2*1024*1024)
self.send_response(response.status)
for name, value in response.getheaders():
if name.lower() not in ('connection', 'transfer-encoding', 'server', 'date', 'content-length'):
self.send_header(name, value)
self.send_header('Content-Length', str(len(data)))
self.end_headers()
self.wfile.write(data)
finally:
conn.close()
except (ValueError, OSError, http.client.HTTPException):
self.send_error(502, 'Deck unavailable')
class Agent:
def __init__(self):
mode_file = DATA/'mode.json'
mode = json.loads(mode_file.read_text()) if mode_file.exists() else 'lan'
self.policy = Policy(mode=mode, persist=lambda value: save(mode_file, value))
self.tunnel_server = None
self.address = None
self.error = None
self.lan_ip = os.environ.get('DECK_LAN_IP', '')
def connected(self):
try:
raw = command('wg', 'show', INTERFACE, 'latest-handshakes', check=False)
values = [int(line.split()[1]) for line in raw.splitlines()]
latest = max(values, default=0)
return latest > 0 and 0 <= time.time()-latest < 180, latest or None
except (ValueError, OSError, subprocess.SubprocessError):
return False, None
def status(self, ingress='management'):
connected, latest = self.connected()
config_file = DATA/'config.json'
config = parse_config(json.loads(config_file.read_text())) if config_file.exists() else None
return dict(installed=True, configured=config is not None, enabled=self.address is not None,
connected=connected, latest_handshake=latest, state='connected' if connected else ('connecting' if self.address else 'disabled'),
tunnel_url=f'http://{self.address}:{PORT}' if self.address else None,
lan_url=f'http://{self.lan_ip}:{PORT}' if self.lan_ip else None,
warnings=config['warnings'] if config else [], error=self.error,
ingress=ingress, **self.policy.status())
def disconnect(self):
if self.tunnel_server:
self.tunnel_server.shutdown()
self.tunnel_server.server_close()
self.tunnel_server = None
command('ip', 'link', 'delete', INTERFACE, check=False)
command('ip', '-4', 'rule', 'del', 'priority', '20000', check=False)
command('ip', '-4', 'route', 'flush', 'table', '51820', check=False)
self.address = None
def connect(self):
if self.address:
return
path = DATA/'config.json'
if not path.exists():
raise ValueError('Zuerst eine eigene WireGuard-Konfiguration importieren.')
config = parse_config(json.loads(path.read_text()))
# Avoid ambiguity with the LAN listener or loopback, even for unusual uploads.
eth = json.loads(command('ip', '-j', '-4', 'addr', 'show', 'dev', 'eth0'))
if any(a['local'] == config['address'] for i in eth for a in i.get('addr_info', [])):
raise ValueError('Tunneladresse kollidiert mit der Containeradresse.')
try:
command('ip', 'link', 'add', INTERFACE, 'type', 'wireguard')
secret = RUN/'wireguard.conf'
fd = os.open(secret, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
try:
with os.fdopen(fd, 'w') as stream:
stream.write(wireguard_text(config))
command('wg', 'setconf', INTERFACE, str(secret))
finally:
secret.unlink(missing_ok=True)
command('ip', '-4', 'addr', 'add', config['address']+'/32', 'dev', INTERFACE)
command('ip', 'link', 'set', INTERFACE, 'mtu', str(config['mtu']), 'up')
# Only packets sourced from the tunnel address use these routes.
# Never change the host or container main/default route.
for network in config['allowed_ips']:
command('ip', '-4', 'route', 'replace', network, 'dev', INTERFACE, 'table', '51820')
command('ip', '-4', 'rule', 'add', 'priority', '20000', 'from', config['address']+'/32', 'lookup', '51820')
self.tunnel_server = DeviceServer(config['address'], INTERFACE, 'tunnel')
threading.Thread(target=self.tunnel_server.serve_forever, daemon=True).start()
self.address = config['address']
self.error = None
except Exception:
self.disconnect()
raise ValueError('Tunnel konnte nicht eingerichtet werden. Keine Host-Konfiguration wurde geändert.') from None
def dispatch(self, data):
action = data.get('action')
ingress = data.get('ingress', 'management')
# Ingress may only be asserted by the web child using the private proxy token.
if not secrets.compare_digest(str(data.get('proxy_token', '')), PROXY_TOKEN):
ingress = 'management'
if action == 'credentials-read':
return {'credentials': validate_record(normalize(json.loads((DATA/'auth.json').read_text())))}
if action == 'credentials-write':
current = validate_record(normalize(json.loads((DATA/'auth.json').read_text())))
if current['revision'] != data.get('expected_revision'):
raise ValueError('Zugangsdaten wurden inzwischen geändert. Bitte erneut anmelden.')
record = validate_record(data.get('credentials'))
save(DATA/'auth.json', record)
return {'saved':True}
if action == 'status':
return self.status(ingress)
if action == 'import':
self.policy.expire()
if self.address or self.policy.mode != 'lan' or self.policy.pending:
raise ValueError('Import ist nur im bestätigten LAN-Modus bei deaktiviertem Tunnel möglich.')
config = parse_config(data.get('config'))
save(DATA/'config.json', data['config'])
self.error = None
elif action == 'connect':
self.connect()
save(DATA/'enabled.json', True)
elif action == 'disconnect':
self.disconnect()
save(DATA/'enabled.json', False)
elif action == 'mode':
self.policy.propose(data.get('mode'), self.connected()[0])
elif action == 'confirm':
self.policy.confirm(data.get('trial_id'), ingress)
elif action == 'cancel':
self.policy.cancel()
elif action == 'delete':
if self.address or self.policy.mode != 'lan' or self.policy.pending:
raise ValueError('Löschen ist nur bei deaktiviertem Tunnel im bestätigten LAN-Modus möglich.')
(DATA/'config.json').unlink(missing_ok=True)
else:
raise ValueError('Unbekannte Netzwerkaktion.')
return self.status(ingress)
class RPC(socketserver.StreamRequestHandler):
def handle(self):
self.connection.settimeout(30)
try:
raw = self.rfile.readline(32769)
if len(raw)>32768:
raise ValueError('Anfrage zu groß.')
data = json.loads(raw)
if not isinstance(data, dict):
raise ValueError('Ungültige Anfrage.')
with LOCK:
result = AGENT.dispatch(data)
except ValueError as exc:
# Parser errors are deliberately generic; never return submitted values.
result = {'error': str(exc) if not isinstance(exc, json.JSONDecodeError) else 'Ungültige Anfrage.'}
except Exception:
result = {'error': 'Netzwerkaktion fehlgeschlagen.'}
self.wfile.write(json.dumps(result).encode()+b'\n')
class RPCServer(socketserver.ThreadingUnixStreamServer):
daemon_threads = True
def main():
global AGENT
if not Path('/.dockerenv').exists() or os.environ.get('DECK_CONTAINER') != '1':
raise SystemExit('This helper must run in its isolated Deck container.')
os.umask(0o077)
DATA.mkdir(exist_ok=True)
DATA.chmod(0o700)
RUN.mkdir(exist_ok=True)
RUN.chmod(0o755)
path = RUN/'control.sock'
path.unlink(missing_ok=True)
AGENT = Agent()
rpc = RPCServer(str(path), RPC)
os.chown(path, 0, 65534)
path.chmod(0o660)
threading.Thread(target=rpc.serve_forever, daemon=True).start()
auth = DATA/'auth.json'
if not auth.exists():
raise SystemExit('Authentication must be provisioned before container startup.')
env = os.environ.copy()
env.update(DECK_PROXY_TOKEN=PROXY_TOKEN, DECK_AUTH_RPC='1', DECK_NETWORK_SOCKET='1', DECK_LOCAL_HARDWARE='1')
child = subprocess.Popen([sys.executable, '/app/server.py'], env=env, user=65534, group=65534, extra_groups=[], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
address = json.loads(command('ip', '-j', '-4', 'addr', 'show', 'dev', 'eth0'))[0]['addr_info'][0]['local']
lan = DeviceServer(address, 'eth0', 'lan')
threading.Thread(target=lan.serve_forever, daemon=True).start()
if (DATA/'enabled.json').exists() and json.loads((DATA/'enabled.json').read_text()):
try:
AGENT.connect()
except ValueError:
AGENT.error = 'Tunnel nach Neustart nicht verfügbar; gewählter Zugriffsmodus bleibt bestehen.'
done = threading.Event()
for sig in (signal.SIGTERM, signal.SIGINT):
signal.signal(sig, lambda *_: done.set())
try:
while not done.wait(1):
with LOCK:
AGENT.policy.expire()
if child.poll() is not None:
break
finally:
lan.shutdown()
rpc.shutdown()
AGENT.disconnect()
child.terminate()
try:
child.wait(timeout=5)
except subprocess.TimeoutExpired:
child.kill()
child.wait()
if __name__ == '__main__':
main()