"""Native host WireGuard, fixed Unix RPC and address-bound Deck access gateways. No wg-quick hooks, default route replacement, DNS changes or Docker operations. """ import argparse import http.client from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer import ipaddress import json import os from pathlib import Path import secrets import selectors import signal import socket import socketserver import struct import subprocess import threading import time from network.config import parse_config, wireguard_text from network.policy import Policy INTERFACE='wgdeck0' TABLE='51826' PRIORITIES={4:'12126',6:'12127'} ENV={'PATH':'/usr/sbin:/usr/bin:/sbin:/bin','LANG':'C.UTF-8'} def run(args,check=True): r=subprocess.run(args,capture_output=True,text=True,timeout=15,env=ENV) if check and r.returncode:raise ValueError('WireGuard-Systemaktion fehlgeschlagen ('+' '.join(args)+'); Verbindung und Hostzustand prüfen.') return r.stdout.strip() def save(path,value): temp=path.with_suffix('.tmp');fd=os.open(temp,os.O_WRONLY|os.O_CREAT|os.O_TRUNC,0o600) with os.fdopen(fd,'w') as f:json.dump(value,f);f.flush();os.fsync(f.fileno()) temp.replace(path) class Manager: def __init__(self,state,policy): self.root=Path(state);self.root.mkdir(mode=0o700,exist_ok=True);self.settings=policy self.lock=threading.RLock();self.listeners=[];self.enabled=False;self.error=None self.policy=Policy(mode=self.read('mode.json','lan'),persist=lambda v:save(self.root/'mode.json',v)) token=self.root/'proxy-token' if not token.exists():token.write_text(secrets.token_hex(32));token.chmod(0o600) self.token=token.read_text().strip() self.reconcile() def read(self,name,default=None): p=self.root/name return json.loads(p.read_text()) if p.exists() else default def config(self): text=self.read('config.json') return parse_config(text,native=True) if text else None def connected(self): if not self.enabled:return False,0 raw=run(['wg','show',INTERFACE,'latest-handshakes'],False) latest=max((int(l.split()[1]) for l in raw.splitlines() if len(l.split())==2),default=0) return latest>0 and 0<=time.time()-latest<180,latest def status(self,ingress='management'): connected,latest=self.connected();c=self.config() return dict(installed=True,install_supported=False,backend='native-systemd',interface=INTERFACE,enabled=self.enabled,configured=bool(c),state='connected' if connected else 'connecting' if self.enabled else 'disabled',latest_handshake=latest,ingress=ingress,lan_url='http://'+self.settings['lan_address']+':'+str(self.settings['gui_port']),tunnel_url='http://'+c['address']+':'+str(self.settings['gui_port']) if c else None,warnings=c['warnings'] if c else [],error=self.error,**self.policy.status()) def connect(self): if self.enabled:return c=self.config() if not c:raise ValueError('Zuerst eine WireGuard-Konfiguration importieren.') existing=json.loads(run(['ip','-j','address','show'])) if any(ipaddress.ip_interface(a).ip==ipaddress.ip_address(v['local']) for a in c['addresses'] for d in existing for v in d.get('addr_info',[])): raise ValueError('Tunnel-IP ist bereits auf diesem Host aktiv. Alten Tunnel zuerst kontrolliert deaktivieren.') if run(['ip','link','show',INTERFACE],False):raise ValueError('Schnittstelle wgdeck0 ist bereits belegt.') for version in (4,6): if run(['ip','-'+str(version),'route','show','table',TABLE],False) or any(line.split(':',1)[0].strip()==PRIORITIES[version] for line in run(['ip','-'+str(version),'rule','show']).splitlines()):raise ValueError('Reservierte Routingtabelle oder Regelpriorität ist bereits belegt.') try: run(['ip','link','add',INTERFACE,'type','wireguard']);self.enabled=True temp=self.root/'setconf.conf';fd=os.open(temp,os.O_CREAT|os.O_EXCL|os.O_WRONLY,0o600) try: with os.fdopen(fd,'w') as f:f.write(wireguard_text(c)) run(['wg','setconf',INTERFACE,str(temp)]) finally:temp.unlink(missing_ok=True) run(['ip','link','set',INTERFACE,'mtu',str(c['mtu']),'up']) for address in c['addresses']: a=ipaddress.ip_interface(address);version=str(a.version) run(['ip','-'+version,'address','add',str(a.ip)+('/32' if a.version==4 else '/128'),'dev',INTERFACE]) for network in c['allowed_ips']: if ipaddress.ip_network(network).version==a.version:run(['ip','-'+version,'route','add',network,'dev',INTERFACE,'table',TABLE]) run(['ip','-'+version,'rule','add','priority',PRIORITIES[a.version],'from',str(a.ip)+('/32' if a.version==4 else '/128'),'lookup',TABLE]) self.reconcile();save(self.root/'enabled.json',True);self.error=None except Exception: self.disconnect();raise def disconnect(self,persist=True): if self.enabled: self.enabled=False;self.reconcile() for version in (4,6): run(['ip','-'+str(version),'rule','del','priority',PRIORITIES[version],'lookup',TABLE],False) run(['ip','-'+str(version),'route','flush','table',TABLE],False) run(['ip','link','del',INTERFACE],False) if persist:save(self.root/'enabled.json',False) def reconcile(self): # LAN never uses 0.0.0.0; API and applications keep their loopback upstream. desired={(self.settings['lan_address'],p,'lan') for p in self.settings['ports']} if self.enabled: c=self.config();desired|={(str(ipaddress.ip_interface(a).ip),p,'tunnel') for a in c['addresses'] for p in self.settings['ports']} present={(s.bind_address,s.port,s.ingress) for s in self.listeners} added=[] try: for address,port,ingress in sorted(desired-present): cls=GUI if port==self.settings['gui_port'] else Relay server=Listener(address,port,ingress,cls,self) threading.Thread(target=server.serve_forever,daemon=True).start();added.append(server) except Exception: for s in added:s.shutdown();s.server_close() raise ValueError('Netzwerkport ist belegt; bestehende Listener bleiben unverändert.') from None self.listeners+=added for s in self.listeners[:]: if (s.bind_address,s.port,s.ingress) not in desired: s.shutdown();s.server_close();self.listeners.remove(s) def dispatch(self,data): action=data.get('action');ingress='management' if secrets.compare_digest(str(data.get('proxy_token','')),self.token) and data.get('ingress') in ('lan','tunnel'):ingress=data['ingress'] with self.lock: if action=='status':return self.status(ingress) if action=='import': if self.enabled or self.policy.mode!='lan' or self.policy.pending:raise ValueError('Import nur bei deaktiviertem Tunnel im bestätigten lokalen Modus.') text=data.get('config');parse_config(text,native=True);save(self.root/'config.json',text) elif action=='connect':self.connect() elif action=='disconnect':self.disconnect() 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.enabled or self.policy.mode!='lan' or self.policy.pending:raise ValueError('Tunnel zuerst deaktivieren und lokalen Modus bestätigen.') (self.root/'config.json').unlink(missing_ok=True) elif action=='backup-export':return dict(config=self.read('config.json'),enabled=self.enabled,mode=self.policy.mode) elif action=='backup-restore': b=data.get('backup',{}) if set(b)!={'config','enabled','mode'} or type(b['enabled']) is not bool or b['mode'] not in Policy.MODES:raise ValueError('Ungültiger Netzwerkstand.') if self.enabled or self.policy.mode!='lan':raise ValueError('Restore benötigt deaktivierten Tunnel im lokalen Modus.') if b['config']:parse_config(b['config'],native=True);save(self.root/'config.json',b['config']) if b['enabled']:self.connect() if b['mode']!='lan':self.policy.propose(b['mode'],self.connected()[0]) else:raise ValueError('Unbekannte Netzwerkaktion.') return self.status(ingress) class Listener(socketserver.ThreadingTCPServer): allow_reuse_address=True;daemon_threads=True def __init__(self,address,port,ingress,handler,manager): self.address_family=socket.AF_INET6 if ':' in address else socket.AF_INET self.bind_address=address;self.port=port;self.ingress=ingress;self.manager=manager super().__init__((address,port),handler,bind_and_activate=False) try: device=INTERFACE if ingress=='tunnel' else manager.settings.get('lan_interface') if device: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 def allowed(self): with self.manager.lock:return self.manager.policy.allowed(self.ingress) class Relay(socketserver.BaseRequestHandler): def handle(self): if not self.server.allowed():return try: with socket.create_connection(('127.0.0.1',self.server.port),timeout=5) as upstream,selectors.DefaultSelector() as sel: self.request.setblocking(False);upstream.setblocking(False) sel.register(self.request,selectors.EVENT_READ,upstream);sel.register(upstream,selectors.EVENT_READ,self.request) while self.server.allowed(): for key,_ in sel.select(1): data=key.fileobj.recv(65536) if not data:return # Bounded blocking send: preserve streaming and WebSocket bytes unchanged. dst=key.data;dst.settimeout(10);dst.sendall(data);dst.setblocking(False) except (OSError,ValueError):pass class GUI(BaseHTTPRequestHandler): protocol_version='HTTP/1.1' def log_message(self,*args):pass def do_GET(self):self.forward() def do_POST(self):self.forward() def do_HEAD(self):self.forward() def forward(self): if not self.server.allowed():self.send_error(403,'Network access disabled');return expected=('['+self.server.bind_address+']' if ':' in self.server.bind_address else self.server.bind_address)+':'+str(self.server.port) if self.headers.get('Host')!=expected:self.send_error(403,'Host rejected');return try: if self.headers.get('Transfer-Encoding'):raise ValueError() length=int(self.headers.get('Content-Length','0')) if not 0<=length<=64*1024*1024:raise ValueError() self.connection.settimeout(30) body=self.rfile.read(length) if length else None if body is not None and len(body)!=length:raise ValueError() headers={k:v for k,v in self.headers.items() if k.lower() not in ('connection','x-deck-proxy','x-deck-ingress','transfer-encoding')} headers.update({'X-Deck-Proxy':self.server.manager.token,'X-Deck-Ingress':self.server.ingress,'Connection':'close'}) conn=http.client.HTTPConnection('127.0.0.1',self.server.port,timeout=120) try: conn.request(self.command,self.path,body,headers);response=conn.getresponse();self.send_response(response.status) for k,v in response.getheaders(): if k.lower() not in ('connection','transfer-encoding','server','date'):self.send_header(k,v) self.send_header('Connection','close');self.end_headers();self.close_connection=True if self.command!='HEAD': while chunk:=response.read(65536):self.wfile.write(chunk) finally:conn.close() except (OSError,ValueError,http.client.HTTPException):self.close_connection=True class RPC(socketserver.StreamRequestHandler): def handle(self): self.connection.settimeout(30) try: _,uid,_=struct.unpack('3i',self.connection.getsockopt(socket.SOL_SOCKET,socket.SO_PEERCRED,12)) if uid not in (0,self.server.client_uid):raise ValueError('Client nicht freigegeben.') 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.') result=self.server.manager.dispatch(data) except Exception as e:result={'error':str(e) if isinstance(e,ValueError) else 'Nativer Netzwerkdienst nicht verfügbar.'} self.wfile.write(json.dumps(result).encode()+b'\n') class Server(socketserver.ThreadingUnixStreamServer):daemon_threads=True def main(): p=argparse.ArgumentParser();p.add_argument('--client-uid',type=int,required=True);p.add_argument('--client-gid',type=int,required=True);p.add_argument('--state',default='/var/lib/athena-deck-network');args=p.parse_args() root=Path(args.state);policy=json.loads((root/'policy.json').read_text()) manager=Manager(root,policy) runtime=Path('/run/athena-deck-network');runtime.mkdir(exist_ok=True) token=runtime/'proxy-token';token.write_text(manager.token);os.chown(token,0,args.client_gid);token.chmod(0o640) sock=runtime/'control.sock';sock.unlink(missing_ok=True) with Server(str(sock),RPC) as server: os.chown(sock,0,args.client_gid);sock.chmod(0o660);server.client_uid=args.client_uid;server.manager=manager # enabled is persisted only on explicit successful activation; native service owns recovery. if manager.read('enabled.json',False): try:manager.connect() except Exception:manager.error='Gespeicherter Tunnel konnte nicht aktiviert werden; Hostzustand prüfen.' signal.signal(signal.SIGTERM,lambda *_:threading.Thread(target=server.shutdown,daemon=True).start()) try:server.serve_forever() finally: manager.disconnect(persist=False) for listener in manager.listeners:listener.shutdown();listener.server_close() if __name__=='__main__':main()