Files
Athena-Deck/network/native.py
T

236 lines
13 KiB
Python

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