304 lines
20 KiB
Python
304 lines
20 KiB
Python
"""Owned llama.cpp worker and fair GPU leases. No production service control."""
|
|
import contextlib
|
|
import http.client
|
|
import json
|
|
import os
|
|
from pathlib import Path
|
|
import secrets
|
|
import shlex
|
|
import signal
|
|
import socket
|
|
import subprocess
|
|
import threading
|
|
import time
|
|
from image_test import probe, GIB, cgroup_headroom
|
|
|
|
class InferenceError(ValueError):
|
|
pass
|
|
|
|
class Scheduler:
|
|
"""FIFO admission; same-profile requests share configured slots, switches drain."""
|
|
def __init__(self, worker):
|
|
self.worker=worker;self.evict_tts=lambda:None;self.cv=threading.Condition();self.queue=[];self.key=None;self.active=0;self.tts_active=0;self.transition=False;self.gpu_mode="llm"
|
|
def status(self):
|
|
with self.cv:return dict(active_requests=self.active+self.tts_active,waiting_requests=len(self.queue),switching=self.transition)
|
|
@contextlib.contextmanager
|
|
def lease(self, key, slots=1, prepare=None, timeout=600, allowed=lambda:True, modes=("llm",)):
|
|
ticket=object();deadline=time.monotonic()+timeout;claimed=False
|
|
with self.cv:
|
|
if len(self.queue)>=16:raise InferenceError('Warteschlange voll. Bitte später erneut versuchen.')
|
|
self.queue.append(ticket)
|
|
try:
|
|
while True:
|
|
if self.gpu_mode not in modes:raise InferenceError("Video-Modus aktiv oder Moduswechsel läuft; GPU-Aufträge sind gesperrt.")
|
|
if not allowed():raise InferenceError('Endpunkt wird gestoppt oder Profil ist nicht mehr aktiviert.')
|
|
if self.queue[0] is ticket and not self.transition and (not self.active or (self.key==key and self.active<slots)) and (not self.tts_active or self.key==key and key!=('image',)):
|
|
self.queue.pop(0);self.active+=1;claimed=True
|
|
switch=self.key!=key;self.key=key
|
|
# Also verifies/restarts a crashed worker when there are no other leases.
|
|
self.transition=self.active==1
|
|
break
|
|
if time.monotonic()>=deadline:raise InferenceError('Ressourcen noch belegt. Anfrage erneut versuchen.')
|
|
self.cv.wait(min(1,max(.01,deadline-time.monotonic())))
|
|
finally:
|
|
if ticket in self.queue:self.queue.remove(ticket)
|
|
self.cv.notify_all()
|
|
try:
|
|
if self.transition:
|
|
try:
|
|
if switch:self.worker.stop()
|
|
if key==('image',) or switch:self.evict_tts()
|
|
if prepare:prepare()
|
|
except Exception:
|
|
self.worker.stop()
|
|
with self.cv:self.key=None
|
|
raise
|
|
finally:
|
|
with self.cv:self.transition=False;self.cv.notify_all()
|
|
yield
|
|
finally:
|
|
if claimed:
|
|
with self.cv:self.active-=1;self.cv.notify_all()
|
|
def image_reservation(self, wait=False,kind='image'):
|
|
lease=self.lease((kind,),timeout=600 if wait else 0,modes=("music",) if kind=="music" else ("llm",))
|
|
lease.__enter__()
|
|
return lambda:lease.__exit__(None,None,None)
|
|
def tts_reservation(self,wait=False):
|
|
deadline=time.monotonic()+(600 if wait else 0)
|
|
with self.cv:
|
|
while True:
|
|
if self.gpu_mode!='llm':raise InferenceError('Video-Modus aktiv; TTS ist nicht verfügbar.')
|
|
if not self.transition and not (self.key==('image',) and self.active) and not self.queue:
|
|
self.tts_active+=1;break
|
|
if time.monotonic()>=deadline:raise InferenceError('GPU-Auftrag läuft; TTS später erneut versuchen.')
|
|
self.cv.wait(min(1,max(.01,deadline-time.monotonic())))
|
|
def release():
|
|
with self.cv:self.tts_active-=1;self.cv.notify_all()
|
|
return release
|
|
def unload_idle(self):
|
|
with self.cv:
|
|
if self.active or self.tts_active:return False
|
|
self.transition=True
|
|
try:self.worker.stop();self.evict_tts();self.key=None
|
|
finally:self.transition=False;self.cv.notify_all()
|
|
return True
|
|
|
|
class LlamaWorker:
|
|
def __init__(self,root,catalog,runtime):
|
|
self.root=Path(root);self.catalog=catalog;self.runtime=runtime;self.lock=threading.RLock()
|
|
self.process=None;self.fit_process=None;self.generation=0;self.port=None;self.profile=None;self.state='stopped';self.error=None;self.key=None;self.devices=[];self.memory_plan=None;self.phase='Gestoppt';self.oom_before=0
|
|
def build(self):
|
|
state=self.runtime.status()
|
|
build=next((b for b in state['builds'] if b['id']==state['active']),None)
|
|
if not build or build['backend']!='CUDA':raise InferenceError('Kein aktiver CUDA-Build. Unter Laufzeiten → llama.cpp einen Build auswählen.')
|
|
directory=(self.runtime.root/build['id']).resolve()
|
|
allowed=(self.runtime.root.resolve(),(self.runtime.root.parent/'runtime-verification').resolve())
|
|
if not any(directory.is_relative_to(root) for root in allowed):raise InferenceError('Ungültiger Build-Pfad.')
|
|
if not (directory/'build/bin/llama-server').is_file() or not build.get('fit_tool'):raise InferenceError('llama-server oder Fit-Werkzeug fehlen.')
|
|
return directory
|
|
def blockers(self,p):
|
|
try:
|
|
self.build()
|
|
if p['parameters'].get('vision_projector'):
|
|
projector=self.catalog.entry(p['parameters']['vision_projector'])
|
|
if projector.get('role')!='vision_projector':raise InferenceError('Vision-Projektor fehlt oder ist ungültig.')
|
|
if not p.get('model') or not p['model']['file'].endswith('.gguf'):raise InferenceError('Ein lokales GGUF-Modell wird benötigt.')
|
|
return []
|
|
except ValueError as exc:return [str(exc)]
|
|
@staticmethod
|
|
def oom_count():
|
|
try:return int(dict(line.split() for line in Path('/sys/fs/cgroup/memory.events').read_text().splitlines()).get('oom_kill',0))
|
|
except (OSError,ValueError):return 0
|
|
def crash_message(self):
|
|
if self.oom_count()>self.oom_before:return 'OOM-Kill im Deck-RAM-Bereich erkannt. Das RAM-Limit wurde überschritten; Kontext oder Modellgröße reduzieren.'
|
|
return 'Modellprozess beendet. GPU-OOM oder Modell-/Buildfehler möglich; Ursache nicht eindeutig bestätigt. Kontext oder Microbatch reduzieren und erneut testen.'
|
|
def status(self):
|
|
with self.lock:
|
|
if self.process and self.process.poll() is not None and self.state=='ready':self.state='failed';self.error=self.crash_message()
|
|
return dict(state=self.state,profile_id=self.profile['id'] if self.profile else None,profile_name=self.profile['name'] if self.profile else None,port=self.port,error=self.error,gpus=list(self.devices),memory_plan=self.memory_plan,phase=self.phase)
|
|
def stop(self):
|
|
with self.lock:
|
|
self.generation+=1
|
|
if self.fit_process and self.fit_process.poll() is None:
|
|
self.fit_process.terminate()
|
|
try:self.fit_process.wait(timeout=3)
|
|
except subprocess.TimeoutExpired:self.fit_process.kill();self.fit_process.wait()
|
|
p=self.process
|
|
if p and p.poll() is None:
|
|
try:
|
|
os.killpg(p.pid,signal.SIGTERM)
|
|
try:p.wait(timeout=10)
|
|
except subprocess.TimeoutExpired:os.killpg(p.pid,signal.SIGKILL);p.wait(timeout=5)
|
|
except ProcessLookupError:pass
|
|
self.process=None;self.profile=None;self.port=None;self.state='stopped';self.phase='Gestoppt';self.devices=[];self.key=None
|
|
(self.root/'worker.key').unlink(missing_ok=True)
|
|
def ensure(self,profile,cancel=lambda:False):
|
|
with self.lock:
|
|
if self.process and self.process.poll() is None and self.profile and (self.profile['id'],self.profile['revision'])==(profile['id'],profile['revision']):return
|
|
self.stop();self.state='loading';self.error=None;self.memory_plan=None;self.phase='Speicherprüfung';self.oom_before=self.oom_count();generation=self.generation
|
|
try:self._start(profile,generation,cancel)
|
|
except Exception as exc:
|
|
self.stop()
|
|
with self.lock:self.state='failed';self.phase='Modellstart fehlgeschlagen';self.error=str(exc) if isinstance(exc,ValueError) else 'llama.cpp konnte nicht gestartet werden.'
|
|
raise InferenceError(self.error) from None
|
|
def plan(self,profile,cancel=lambda:False):
|
|
self._start(profile,self.generation,cancel,plan_only=True)
|
|
return self.memory_plan
|
|
def _start(self,profile,generation,cancel=lambda:False,plan_only=False):
|
|
directory=self.build();params=profile['parameters'];entry=self.catalog.entry(profile['model_id'])
|
|
model=(self.catalog.root/entry['id']/('model'+Path(entry['file']).suffix)).resolve()
|
|
mtp=params.get('mtp',False)
|
|
if mtp and not (directory/'build/bin/deck-mtp-fit-v1').is_file():raise InferenceError('MTP benötigt einen neu gebauten llama.cpp-Build mit Deck-MTP-Speicherprüfung. Unter Laufzeiten erneut bauen.')
|
|
devices=probe();ids=params['gpu_devices']
|
|
if ids:
|
|
selected=[next((g for g in devices if g['uuid']==ident),None) for ident in ids]
|
|
if any(g is None for g in selected):raise InferenceError('Eine ausgewählte GPU ist nicht verfügbar.')
|
|
else:
|
|
selected=sorted([g for g in devices if g['processes']==0],key=lambda g:g['free_mib'],reverse=True)
|
|
if params['split_mode']=='none':selected=selected[:1]
|
|
if not selected or any(g['processes'] or g['free_mib']<2048 for g in selected):raise InferenceError('Benötigte GPU ist durch einen anderen Dienst belegt. Deck stoppt keine fremden Prozesse.')
|
|
mem={line.split(':')[0]:int(line.split()[1])*1024 for line in Path('/proc/meminfo').read_text().splitlines() if line.startswith('MemAvailable:')}
|
|
projector=None;vision_args=[];vision_bytes=0;vision_device=params.get('vision_device','cpu')
|
|
if params.get('vision_projector'):
|
|
projector=self.catalog.entry(params['vision_projector'])
|
|
if projector.get('role')!='vision_projector':raise InferenceError('Kein gültiger Vision-Projektor.')
|
|
vision_bytes=projector['size']+2*GIB
|
|
vision_args=['--mmproj',str((self.catalog.root/projector['id']/'model.gguf').resolve())]
|
|
if vision_device=='cpu':vision_args+=['--no-mmproj-offload']
|
|
else:
|
|
if vision_device not in [g['uuid'] for g in selected]:raise InferenceError('Projektor-GPU ist nicht ausgewählt.')
|
|
vision_args+=['--mmproj-device','CUDA'+str([g['uuid'] for g in selected].index(vision_device)),'--mmproj-offload']
|
|
staging=entry['size']+4*GIB+vision_bytes
|
|
if mem.get('MemAvailable',0)<staging+4*GIB:raise InferenceError('Zu wenig freier RAM für Modellladen und Host-Reserve.')
|
|
headroom=cgroup_headroom()
|
|
if headroom is not None and headroom<staging:raise InferenceError('Deck-RAM-Limit reicht für das Laden dieses Modells nicht aus.')
|
|
# Never inherit LLAMA_ARG_* or user HF credentials into the model server.
|
|
env={k:v for k,v in os.environ.items() if not k.startswith(('LLAMA_','HF_','DECK_'))}
|
|
env.update(CUDA_VISIBLE_DEVICES=','.join(g['uuid'] for g in selected),HF_HUB_OFFLINE='1',OMP_NUM_THREADS=str(params['threads']))
|
|
reserve_mode=params.get('gpu_reserve_mode','auto')
|
|
margins=[0 if reserve_mode=='none' else params.get('gpu_reserve_mib',{}).get(g['uuid'],512) if reserve_mode=='manual' else max(512,round(g['total_mib']*.025)) for g in selected]
|
|
common=['--model',str(model),'--ctx-size',str(params['context']),'--parallel',str(params['slots']),'--batch-size',str(params['batch']),'--ubatch-size',str(params['ubatch']),'--cache-type-k',params.get('cache_type_k','q4_0'),'--cache-type-v',params.get('cache_type_v','q4_0'),'--flash-attn','on','--split-mode',params['split_mode'],'--fit-target',','.join(map(str,margins))]
|
|
if params['tensor_split']:common+=['--tensor-split',','.join(map(str,params['tensor_split']))]
|
|
tool=str(directory/'build/bin/llama-fit-params')
|
|
mtp_fit=['--deck-mtp'] if mtp else []
|
|
def estimate(layers,extra=()):
|
|
if cancel():raise InferenceError('Modellstart abgebrochen.')
|
|
output=self._fit_command([tool]+mtp_fit+common+['--gpu-layers',str(layers),'--fit-print','on']+list(extra),env,generation,cancel)
|
|
rows={}
|
|
for line in output.splitlines():
|
|
parts=line.split()
|
|
if len(parts)==4 and all(x.isdigit() for x in parts[1:]):rows[parts[0]]=sum(map(int,parts[1:]))
|
|
if 'Host' not in rows or any('CUDA'+str(i) not in rows for i in range(len(selected))):raise InferenceError('Speicherprognose unvollständig; Start abgebrochen.')
|
|
if projector:
|
|
target='Host' if vision_device=='cpu' else 'CUDA'+str([g['uuid'] for g in selected].index(vision_device))
|
|
rows[target]+=round(vision_bytes/1024**2)
|
|
gpu_fits=all(rows['CUDA'+str(i)]+margins[i]<=g['free_mib'] for i,g in enumerate(selected))
|
|
need=max(staging,rows['Host']*1024**2+4*GIB)
|
|
host_fits=mem.get('MemAvailable',0)>=need+4*GIB and (headroom is None or headroom>=need)
|
|
with self.lock:self.memory_plan=dict(profile_name=profile['name'],estimated=True,reserve_mode=reserve_mode,gpu_layers=layers,full_gpu=layers==999,context=params['context'],slots=params['slots'],fits=gpu_fits and host_fits,host_required_mib=round(need/1024**2),host_available_mib=round(min(mem.get('MemAvailable',0)-4*GIB,headroom if headroom is not None else mem.get('MemAvailable',0))/1024**2),gpus=[dict(name=g.get('name',g['uuid']),uuid=g['uuid'],free_mib=g['free_mib'],required_mib=rows['CUDA'+str(i)],reserve_mib=margins[i]) for i,g in enumerate(selected)])
|
|
return gpu_fits and host_fits,rows
|
|
extra=[]
|
|
if params.get('gpu_offload')=='full':
|
|
layers=999;fits,memory=estimate(layers)
|
|
if not fits:raise InferenceError('Vollständige GPU-Ausführung passt einschließlich Reserve nicht. Kontext reduzieren oder automatische CPU-Auslagerung wählen.')
|
|
elif params['tensor_split']:
|
|
# Upstream automatic fitting refuses user tensor_split. Predict the
|
|
# fixed split instead; bounded layer search preserves ratio/context.
|
|
layers=999;fits,memory=estimate(layers)
|
|
if not fits:
|
|
low,high,best=0,998,None
|
|
for _ in range(10):
|
|
if low>high:break
|
|
middle=(low+high)//2;ok,usage=estimate(middle)
|
|
if ok:best=(middle,usage);low=middle+1
|
|
else:high=middle-1
|
|
if best is None:raise InferenceError('Kontext und fester GPU-Split passen nicht sicher in GPU/RAM.')
|
|
layers,memory=best
|
|
estimate(layers)
|
|
else:
|
|
output=self._fit_command([tool]+mtp_fit+common,env,generation,cancel)
|
|
flags=shlex.split(output.strip());fit={}
|
|
if len(flags)%2:raise InferenceError('Fit-Werkzeug lieferte ungültige Parameter.')
|
|
for i in range(0,len(flags),2):
|
|
if flags[i] not in ('-c','-ngl','-ts'):raise InferenceError('Unbekannte Fit-Option.')
|
|
fit[flags[i]]=flags[i+1]
|
|
if fit.get('-c')!=str(params['context']) or not fit.get('-ngl','').isdigit():raise InferenceError('Gewünschter Kontext konnte nicht unverändert eingepasst werden.')
|
|
layers=int(fit['-ngl'])
|
|
if '-ts' in fit:
|
|
parts=fit['-ts'].split(',')
|
|
if len(parts)!=len(selected) or any(not x.replace('.','',1).isdigit() for x in parts):raise InferenceError('Ungültige automatische GPU-Verteilung.')
|
|
extra=['--tensor-split',fit['-ts']]
|
|
fits,memory=estimate(layers,extra)
|
|
if not fits:raise InferenceError('Speicherprognose überschreitet die GPU-/RAM-Reserve.')
|
|
if cancel():raise InferenceError('Modellstart abgebrochen.')
|
|
if plan_only:return
|
|
with self.lock:self.phase='Modell wird geladen'
|
|
launch=common+extra+vision_args+['--gpu-layers',str(layers),'--fit','off','--kv-unified','--threads',str(params['threads']),'--load-mode','none','--host','127.0.0.1','--alias',profile['name'],'--metrics','--no-webui','--log-disable']
|
|
if mtp:launch+=['--spec-type','draft-mtp','--spec-draft-n-max',str(params.get('mtp_tokens',2)),'--spec-draft-p-min',str(params.get('mtp_min_p',.05)),'--spec-draft-type-k','f16','--spec-draft-type-v','f16']
|
|
self.root.mkdir(parents=True,exist_ok=True,mode=0o700)
|
|
with socket.socket() as sock:sock.bind(('127.0.0.1',0));port=sock.getsockname()[1]
|
|
key=secrets.token_urlsafe(32);keypath=self.root/'worker.key'
|
|
fd=os.open(keypath,os.O_WRONLY|os.O_CREAT|os.O_TRUNC,0o600)
|
|
with os.fdopen(fd,'w') as f:f.write(key+'\n')
|
|
launch+=['--port',str(port),'--api-key-file',str(keypath)]
|
|
with self.lock:
|
|
if self.generation!=generation:raise InferenceError('Modellstart abgebrochen.')
|
|
self.process=subprocess.Popen([str(directory/'build/bin/llama-server')]+launch,env=env,cwd=directory,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True)
|
|
self.profile=profile;self.port=port;self.key=key;self.devices=[g['uuid'] for g in selected]
|
|
deadline=time.monotonic()+300
|
|
while time.monotonic()<deadline:
|
|
if cancel():raise InferenceError('Modellstart abgebrochen.')
|
|
with self.lock:
|
|
if not self.process or self.process.poll() is not None:raise InferenceError(self.crash_message())
|
|
current={g['uuid']:g for g in probe()}
|
|
if any(current.get(g['uuid'],{}).get('processes',2)>1 for g in selected):raise InferenceError('Ein weiterer Prozess verwendet die reservierte GPU. Deck bricht seinen Start ab.')
|
|
conn=http.client.HTTPConnection('127.0.0.1',port,timeout=2)
|
|
try:
|
|
conn.request('GET','/health',headers={'Authorization':'Bearer '+key});r=conn.getresponse()
|
|
if r.status==200:
|
|
with self.lock:self.state='ready';self.phase='Modell bereit'
|
|
threading.Thread(target=self._monitor,args=(generation,),daemon=True).start()
|
|
return
|
|
except OSError:pass
|
|
finally:conn.close()
|
|
time.sleep(1)
|
|
raise InferenceError('llama.cpp wurde nicht rechtzeitig bereit.')
|
|
def _fit_command(self,args,env,generation,cancel=lambda:False):
|
|
with self.lock:
|
|
if self.generation!=generation:raise InferenceError('Modellstart abgebrochen.')
|
|
process=subprocess.Popen(args,env=env,stdout=subprocess.PIPE,stderr=subprocess.DEVNULL,text=True)
|
|
self.fit_process=process
|
|
try:
|
|
deadline=time.monotonic()+180
|
|
while True:
|
|
if cancel():
|
|
process.terminate();process.communicate(timeout=3);raise InferenceError('Modellstart abgebrochen.')
|
|
try:output,_=process.communicate(timeout=1);break
|
|
except subprocess.TimeoutExpired:
|
|
if time.monotonic()>=deadline:raise
|
|
except subprocess.TimeoutExpired:
|
|
process.kill();process.communicate();raise InferenceError('Zeitlimit der Speicher-Einpassung überschritten.') from None
|
|
finally:
|
|
with self.lock:self.fit_process=None
|
|
if process.returncode:raise InferenceError('Speicher-Einpassung fehlgeschlagen. Kontext, Split oder Modellgröße reduzieren.')
|
|
return output
|
|
def _monitor(self,generation):
|
|
while True:
|
|
time.sleep(3)
|
|
with self.lock:
|
|
if self.generation!=generation or not self.process or self.process.poll() is not None:return
|
|
selected=list(self.devices)
|
|
try:
|
|
current={g['uuid']:g for g in probe()}
|
|
conflict=any(current.get(ident,{}).get('processes',2)>1 for ident in selected)
|
|
except Exception:conflict=True
|
|
if conflict:
|
|
with self.lock:
|
|
if self.generation!=generation:return
|
|
self.stop();self.state='failed';self.error='GPU-Konflikt oder Telemetrie ausgefallen. Nur der eigene Modellprozess wurde beendet.'
|
|
return
|
|
def connect(self):
|
|
with self.lock:
|
|
if not self.process or self.process.poll() is not None:raise InferenceError('Modellprozess nicht verfügbar.')
|
|
return http.client.HTTPConnection('127.0.0.1',self.port,timeout=120),self.key
|