Files

237 lines
16 KiB
Python

#!/usr/bin/env python3
"""Root-owned local helper. Only inventory and fresh Docker install; no generic RPC."""
import argparse
import http.client
from urllib.parse import urlsplit
import re
import urllib.request
import urllib.error
import json
import os
from pathlib import Path
import shutil
import socket
import socketserver
import struct
import subprocess
import threading
import time
MANAGED_LABEL='io.athena-deck.managed=true'
ROLE_LABEL='io.athena-deck.role=application'
PACKAGES=['docker-ce','docker-ce-cli','containerd.io','docker-buildx-plugin','docker-compose-plugin']
CONFLICTS=['docker.io','docker-compose','docker-doc','docker-buildx','podman-docker','containerd','runc']
ENV={'PATH':'/usr/sbin:/usr/bin:/sbin:/bin','HOME':'/var/lib/athena-deck-docker','LANG':'C.UTF-8','DEBIAN_FRONTEND':'noninteractive','NEEDRESTART_MODE':'l'}
def run(args,timeout=15):
return subprocess.run(args,capture_output=True,text=True,timeout=timeout,env=ENV)
def release():
return dict(line.split('=',1) for line in Path('/etc/os-release').read_text().replace('"','').splitlines() if '=' in line)
class Manager:
def __init__(self,state,allow_install=False):
self.state=Path(state);self.allow_install=allow_install;self.lock=threading.RLock();self.job=None
if (self.state/'job.json').exists():
self.job=json.loads((self.state/'job.json').read_text())
if self.job['state']=='running':self.job.update(state='interrupted',phase='Installation unterbrochen. Hostzustand prüfen; vorhandene Pakete werden nicht automatisch ersetzt.')
def inventory(self):
# Filter at the daemon before retrieving rows. No inspect, env, mounts or logs.
result=run(['/usr/bin/docker','ps','-a','--no-trunc','--filter','label='+MANAGED_LABEL,'--filter','label='+ROLE_LABEL,'--format','{{.ID}}\t{{.Names}}\t{{.Image}}\t{{.State}}\t{{.Status}}'])
if result.returncode:raise ValueError('Docker-Daemon nicht erreichbar.')
rows=[]
for line in result.stdout.splitlines()[:500]:
fields=line.split('\t')
if len(fields)==5:rows.append(dict(zip(('id','name','image','state','status'),fields)))
return rows
def preflight(self):
info=release()
if info.get('ID')!='debian' or info.get('VERSION_ID') not in ('12','13'):return False,'Unterstützt werden Debian 12 und 13.'
if Path('/usr/bin/docker').exists() or Path('/var/run/docker.sock').exists():return False,'Docker ist vorhanden. Keine automatische Neuinstallation oder Aktualisierung.'
if not self.allow_install:return False,'Docker-Erstinstallation ist für diesen Helfer nicht freigegeben.'
if any(Path(p).exists() for p in ('/var/lib/docker','/var/lib/containerd','/etc/docker')):return False,'Vorhandene Containerdaten oder Konfiguration erkannt. Manuelle Prüfung erforderlich.'
for package in CONFLICTS+PACKAGES:
value=run(['/usr/bin/dpkg-query','-W','-f=${Status}',package])
if value.returncode==0 and 'install ok installed' in value.stdout:return False,'Vorhandene Containerpakete erkannt. Keine automatische Entfernung oder Ersetzung.'
if shutil.disk_usage('/var').free<3*1024**3:return False,'Mindestens 3 GiB freier Speicher erforderlich.'
return True,'Docker Engine, Compose und Buildx aus dem offiziellen Docker-Repository installieren. Docker richtet einen Systemdienst sowie eigene Netzwerk-/Firewallregeln ein. Keine NVIDIA-Treiberinstallation und kein Host-Neustart.'
def status(self):
installed=Path('/usr/bin/docker').exists();reachable=False;version=None;services=[]
if installed:
try:
version_result=run(['/usr/bin/docker','version','--format','{{.Server.Version}}'])
if version_result.returncode==0:version=version_result.stdout.strip();services=self.inventory();reachable=True
except (OSError,ValueError,subprocess.SubprocessError):pass
supported,message=self.preflight()
with self.lock:job=dict(self.job) if self.job else None
return dict(installed=installed,reachable=reachable,version=version,services=services,install_supported=supported,message=message,job=job,labels=[MANAGED_LABEL,ROLE_LABEL])
def _phase(self,text,state='running'):
with self.lock:
self.job.update(state=state,phase=text);self.state.mkdir(parents=True,exist_ok=True,mode=0o700)
temp=self.state/'job.tmp';temp.write_text(json.dumps(self.job));temp.replace(self.state/'job.json')
def install(self):
with self.lock:
if self.job and self.job['state']=='running':raise ValueError('Docker-Installation läuft bereits.')
allowed,message=self.preflight()
if not allowed:raise ValueError(message)
self.job=dict(state='running',phase='Docker-Erstinstallation wird vorbereitet',started_at=time.time())
self._phase(self.job['phase']);threading.Thread(target=self._install,daemon=True).start()
return {'started':True}
def _checked(self,args):
result=run(args,timeout=1200)
if result.returncode:raise ValueError('Installationsschritt fehlgeschlagen. Paketverwaltung und Internetverbindung auf dem Host prüfen.')
def _install(self):
try:
info=release();codename={'12':'bookworm','13':'trixie'}[info['VERSION_ID']]
arch=run(['/usr/bin/dpkg','--print-architecture']).stdout.strip()
if arch not in ('amd64','arm64'):raise ValueError('Diese Architektur wird vom Deck-Installer noch nicht unterstützt.')
self._phase('Paketquellen prüfen und Download-Werkzeuge installieren')
self._checked(['/usr/bin/apt-get','update'])
self._checked(['/usr/bin/apt-get','install','-y','--no-remove','ca-certificates','curl'])
key=Path('/etc/apt/keyrings/athena-deck-docker.asc');source=Path('/etc/apt/sources.list.d/athena-deck-docker.sources')
key.parent.mkdir(mode=0o755,exist_ok=True)
self._checked(['/usr/bin/curl','--fail','--silent','--show-error','--proto','=https','--tlsv1.2','--max-time','60','https://download.docker.com/linux/debian/gpg','-o',str(key)])
key.chmod(0o644)
source.write_text(f'Types: deb\nURIs: https://download.docker.com/linux/debian\nSuites: {codename}\nComponents: stable\nArchitectures: {arch}\nSigned-By: {key}\n');source.chmod(0o644)
self._phase('Docker Engine, Compose und Buildx installieren')
self._checked(['/usr/bin/apt-get','update'])
self._checked(['/usr/bin/apt-get','install','-y','--no-remove',*PACKAGES])
self._phase('Docker-Systemdienst starten und Erreichbarkeit prüfen')
self._checked(['/usr/bin/systemctl','enable','--now','docker.service'])
self._checked(['/usr/bin/docker','version','--format','{{.Server.Version}}'])
self._checked(['/usr/bin/docker','compose','version','--short'])
self._phase('Docker bereit. Keine Anwendungen oder Modelle installiert.','complete')
except Exception as exc:self._phase(str(exc) if isinstance(exc,ValueError) else 'Docker-Installation fehlgeschlagen. Hostzustand vor erneutem Versuch prüfen.','failed')
def video_configs(self):
path=self.state/'video-services.json'
if not path.exists():return []
if path.is_symlink() or path.stat().st_uid!=0 or path.stat().st_mode&0o022:raise ValueError('Videodienste müssen durch root eingerichtet werden.')
configs=json.loads(path.read_text())
if not isinstance(configs,list) or len(configs)>20:raise ValueError('Ungültige Dienstliste.')
ids=set()
for config in configs:
if not re.fullmatch('[a-z0-9-]{1,40}',config.get('id','')) or config['id'] in ids or not re.fullmatch('[a-f0-9]{64}',config.get('container_id','')):raise ValueError('Ungültige Dienstzuordnung.')
ids.add(config['id'])
return configs
def video_status(self,config):
result=run(['/usr/bin/docker','inspect','--format','{{.State.Running}}',config['container_id']])
if result.returncode:return dict(id=config['id'],name=config['name'],running=None,ready=False,api_url=config['api_url'],message='Zugeordneter Container fehlt.')
running=result.stdout.strip()=='true';ready=False
if running:
class NoRedirect(urllib.request.HTTPRedirectHandler):
def redirect_request(self,*args,**kwargs):return None
try:
headers={}
if config.get('token_file'):headers['Authorization']='Bearer '+Path(config['token_file']).read_text().strip()
request=urllib.request.Request(config['health_url'],headers=headers)
with urllib.request.build_opener(NoRedirect).open(request,timeout=2) as response:ready=response.status==200
except (OSError,ValueError,urllib.error.URLError):pass
return dict(id=config['id'],name=config['name'],running=running,ready=ready,api_url=config['api_url'],message=config.get('message','Original-API des Dienstes.'))
def video_inventory(self):return {'services':[self.video_status(config) for config in self.video_configs()]}
def video_action(self,ident,start):
with self.lock:
configs=self.video_configs();config=next((c for c in configs if c['id']==ident),None)
if config is None:raise ValueError('Videodienst nicht freigegeben.')
status=self.video_status(config)
if start and status['running']:return status
if start:
if any(self.video_status(c)['running'] is not False for c in configs if c['id']!=ident):raise ValueError('Ein anderer Videodienst ist noch aktiv oder nicht prüfbar.')
probe=run(['/usr/bin/nvidia-smi','--query-compute-apps=pid','--format=csv,noheader,nounits'])
if probe.returncode or probe.stdout.strip():raise ValueError('GPU durch einen anderen Dienst belegt. Deck beendet keine fremden Modelle. Alten Router zuerst freigeben.')
result=run(['/usr/bin/docker','start',config['container_id']] if start else ['/usr/bin/docker','stop','--time','30',config['container_id']],timeout=45)
if result.returncode:raise ValueError('Videodienst konnte nicht gestartet oder beendet werden.')
return self.video_status(config)
def container_action(self,ident,action):
if action not in ('start','stop','delete') or not isinstance(ident,str) or not re.fullmatch('[a-f0-9]{64}',ident):raise ValueError('Ungültige Containeraktion.')
with self.lock:
row=next((c for c in self.inventory() if c['id']==ident),None)
if row is None:raise ValueError('Container nicht für Deck freigegeben oder nicht mehr vorhanden.')
if any(c['container_id']==ident for c in self.video_configs()):raise ValueError('Registrierte Video-Laufzeiten über die Steuerung verwalten.')
if action=='delete' and row['state'] not in ('exited','created','dead'):raise ValueError('Container vor dem Löschen stoppen.')
args=['/usr/bin/docker']+(['stop','--time','30',ident] if action=='stop' else ['rm',ident] if action=='delete' else ['start',ident])
result=run(args,timeout=45)
if result.returncode:raise ValueError('Containeraktion fehlgeschlagen. Status erneut prüfen.')
return dict(changed=True,action=action,id=ident)
def dispatch(self,data):
if not isinstance(data,dict):raise ValueError('Ungültige Helferanfrage.')
if set(data)=={'action','service'} and data['action'] in ('container-start','container-stop','container-delete'):
return self.container_action(data['service'],data['action'].removeprefix('container-'))
if data.get('action')=='backup-export' and set(data)=={'action'}:
import docker_backup
return docker_backup.export(self)
if data.get('action')=='backup-restore' and set(data)=={'action','container'}:
import docker_backup
return docker_backup.restore(self,data['container'])
if set(data)=={'action','service'} and data['action'] in ('video-start','video-stop'):
return self.video_action(data['service'],data['action']=='video-start')
if set(data)!={'action'}:raise ValueError('Ungültige Helferanfrage.')
if data['action']=='video-status':return self.video_inventory()
if data['action']=='status':return self.status()
if data['action']=='install':return self.install()
raise ValueError('Aktion nicht erlaubt.')
def video_http(manager,data,source,target):
"""Fixed registered upstream only; no model translation or secret returned to Deck."""
sent=False;connection=None
try:
if set(data)!={'action','service','method','path','length','headers'}:raise ValueError('Ungültige Proxy-Anfrage.')
method=data['method'];path=data['path'];length=data['length'];headers=data['headers']
if method not in ('GET','POST','PUT','PATCH','DELETE','HEAD','OPTIONS') or not isinstance(path,str) or not path.startswith('/') or path.startswith('//') or any(ord(c)<32 for c in path):raise ValueError('Ungültige HTTP-Anfrage.')
if type(length) is not int or not 0<=length<=256*1024*1024 or not isinstance(headers,dict):raise ValueError('Ungültige Anfragegröße.')
allowed={'content-type','accept','range','if-range','if-none-match','if-modified-since'}
if any(not isinstance(k,str) or k.lower() not in allowed or not isinstance(v,str) or '\r' in v or '\n' in v for k,v in headers.items()):raise ValueError('Ungültige Header.')
config=next((c for c in manager.video_configs() if c['id']==data['service']),None)
if config is None:raise ValueError('Videodienst nicht freigegeben.')
if manager.video_status(config)['running'] is not True:raise ValueError('Videodienst nicht gestartet.')
url=urlsplit(config['api_url'])
if url.scheme not in ('http','https') or not url.hostname or url.username or url.password:raise ValueError('Ungültiges Backend-Ziel.')
connection=(http.client.HTTPSConnection if url.scheme=='https' else http.client.HTTPConnection)(url.hostname,url.port,timeout=3600)
connection.putrequest(method,path,skip_accept_encoding=True)
for key,value in headers.items():connection.putheader(key,value)
if config.get('token_file'):connection.putheader('Authorization','Bearer '+Path(config['token_file']).read_text().strip())
connection.putheader('Content-Length',str(length));connection.putheader('Connection','close');connection.endheaders()
remaining=length
while remaining:
chunk=source.read(min(1024*1024,remaining))
if not chunk:raise ValueError('Unvollständige Anfrage.')
connection.send(chunk);remaining-=len(chunk)
response=connection.getresponse();sent=True
target.write(f'HTTP/1.1 {response.status} {response.reason}\r\n'.encode())
hop={'connection','keep-alive','proxy-authenticate','proxy-authorization','te','trailer','transfer-encoding','upgrade'}
for key,value in response.getheaders():
if key.lower() not in hop:target.write(f'{key}: {value}\r\n'.encode('latin-1'))
target.write(b'Connection: close\r\n\r\n');target.flush()
if method!='HEAD':
while chunk:=response.read(1024*1024):target.write(chunk);target.flush()
except Exception:
if not sent:
body=b'{"error":"Video-API nicht erreichbar oder Anfrage nicht erlaubt."}'
target.write(b'HTTP/1.1 502 Bad Gateway\r\nContent-Type: application/json\r\nContent-Length: '+str(len(body)).encode()+b'\r\nConnection: close\r\n\r\n'+body)
finally:
if connection:connection.close()
class Handler(socketserver.StreamRequestHandler):
def handle(self):
self.connection.settimeout(20)
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 erlaubt.')
raw=self.rfile.readline(67108865)
if len(raw)>67108864:raise ValueError('Anfrage zu groß.')
data=json.loads(raw)
if isinstance(data,dict) and data.get('action')=='video-http':
self.connection.settimeout(3600);video_http(self.server.manager,data,self.rfile,self.wfile);return
result=self.server.manager.dispatch(data)
except Exception as exc:result={'error':str(exc) if isinstance(exc,ValueError) else 'Docker-Systemhelfer nicht verfügbar.'}
self.wfile.write(json.dumps(result).encode()+b'\n')
class Server(socketserver.ThreadingUnixStreamServer):
daemon_threads=True
def main():
parser=argparse.ArgumentParser();parser.add_argument('--client-uid',type=int,required=True);parser.add_argument('--client-gid',type=int,required=True);parser.add_argument('--allow-install',action='store_true');parser.add_argument('--state',default='/var/lib/athena-deck-docker');args=parser.parse_args()
socket_path=Path('/run/athena-deck-docker/control.sock');socket_path.parent.mkdir(mode=0o755,exist_ok=True);socket_path.unlink(missing_ok=True)
with Server(str(socket_path),Handler) as server:
os.chown(socket_path,0,args.client_gid);socket_path.chmod(0o660);server.client_uid=args.client_uid;server.manager=Manager(args.state,args.allow_install);server.serve_forever()
if __name__=='__main__':main()