Replace WireGuard prototype with native host service and scoped access
This commit is contained in:
1 parent
37b5394df1
commit
11f16ef4a3
16 files changed
+667
-136
No files matched your search
@@ -0,0 +1,235 @@
|
||||
"""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()
|
||||
Reference in new issue
Block a user