Add isolated WireGuard module with guarded access modes and server setup
This commit is contained in:
@@ -0,0 +1,294 @@
|
||||
"""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 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', '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 == '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_JSON=auth.read_text(), 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()
|
||||
Reference in New Issue
Block a user