"""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()