"""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.cv=threading.Condition();self.queue=[];self.key=None;self.active=0;self.transition=False def status(self): with self.cv:return dict(active_requests=self.active,waiting_requests=len(self.queue),switching=self.transition) @contextlib.contextmanager def lease(self, key, slots=1, prepare=None, timeout=600, allowed=lambda:True): 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 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=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 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): lease=self.lease(('image',),timeout=600 if wait else 0) lease.__enter__() return lambda:lease.__exit__(None,None,None) def unload_idle(self): with self.cv: if self.active:return False self.transition=True try:self.worker.stop();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 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 _start(self,profile,generation,cancel=lambda: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:')} staging=entry['size']+4*GIB if 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,gpu_layers=layers,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['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.') with self.lock:self.phase='Modell wird geladen' launch=common+extra+['--gpu-layers',str(layers),'--fit','off','--kv-unified','--threads',str(params['threads']),'--load-mode','none','--host','127.0.0.1','--alias',profile['name'],'--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()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