236 lines
13 KiB
Python
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()
|