174 lines
11 KiB
Python
174 lines
11 KiB
Python
"""Exclusive GPU mode and owned persistent LTX subprocess with asynchronous jobs."""
|
||
import json
|
||
import os
|
||
from pathlib import Path
|
||
import re
|
||
import selectors
|
||
import signal
|
||
import shutil
|
||
import subprocess
|
||
import threading
|
||
import time
|
||
import uuid
|
||
from image_test import probe,cgroup_headroom,GIB
|
||
from profiles import video_recipe
|
||
|
||
class Video:
|
||
def __init__(self,root,profiles,runtime,scheduler,stop_owned):
|
||
self.root=Path(root);self.profiles=profiles;self.runtime=runtime;self.scheduler=scheduler;self.stop_owned=stop_owned
|
||
self.lock=threading.RLock();self.mode='llm';self.state='idle';self.error=None;self.process=None;self.profile=None;self.job=None;self.thread=None;self.gpu=None
|
||
self.closing=False;self.selected=None
|
||
if (self.root/'selection.json').exists():self.selected=json.loads((self.root/'selection.json').read_text()).get('profile_id')
|
||
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='Video-Worker beendet; RAM-OOM oder Laufzeitfehler möglich.'
|
||
return dict(mode=self.mode,state=self.state,error=self.error,profile_name=self.profile['name'] if self.profile else None,profile_id=self.selected,loaded_profile_id=self.profile['id'] if self.profile else None,gpu=self.gpu,job=dict(self.job) if self.job else None,runtime=self.runtime.status(),memory_policy='Disk-Streaming auf RTX 5080; beide GPUs exklusiv reserviert. Gewichte werden bedarfsgerecht geladen.')
|
||
def blockers(self,p):
|
||
if not video_recipe(p.get('model') or {}):return ['Für diese Videovariante ist kein Worker angebunden.']
|
||
errors=[]
|
||
if not self.runtime.status()['installed']:errors.append('Video-Laufzeit unter Einstellungen → Laufzeiten → Video installieren.')
|
||
for role,info in video_recipe(p['model']).items():
|
||
try:self.profiles._component(p['model'],role,p.get('components',{}).get(role))
|
||
except ValueError:errors.append(info['label']+' fehlt oder ist noch nicht zugeordnet.')
|
||
return errors
|
||
def select(self,ident):
|
||
with self.lock:
|
||
if self.mode!='llm' or self.state=='switching':raise ValueError('Videoprofil nur im LLM-Modus wechseln.')
|
||
p=next((p for p in self.profiles.status()['profiles'] if p['id']==ident and p['kind']=='video'),None)
|
||
if not p:raise ValueError('Videoprofil nicht gefunden.')
|
||
self.root.mkdir(parents=True,exist_ok=True,mode=0o700)
|
||
tmp=self.root/'selection.tmp';tmp.write_text(json.dumps({'profile_id':ident}));tmp.replace(self.root/'selection.json');self.selected=ident
|
||
return self.status()
|
||
def switch(self,mode):
|
||
if mode not in ('llm','video'):raise ValueError('Modus muss llm oder video sein.')
|
||
with self.lock:
|
||
if self.closing:raise ValueError('Deck wird beendet.')
|
||
if self.state=='switching':raise ValueError('Moduswechsel läuft bereits.')
|
||
if mode==self.mode and self.state in ('idle','ready'):return self.status()
|
||
profile=None
|
||
if mode=='video':
|
||
profile=next((p for p in self.profiles.status()['profiles'] if p['id']==self.selected),None)
|
||
if not profile:raise ValueError('Zuerst ein aktives Videoprofil auswählen.')
|
||
errors=self.blockers(profile)
|
||
if errors:raise ValueError(' '.join(errors))
|
||
# Close admission before stopping any process. Queued leases recheck this flag.
|
||
with self.scheduler.cv:self.scheduler.gpu_mode='switching';self.scheduler.cv.notify_all()
|
||
self.state='switching';self.error=None
|
||
self.thread=threading.Thread(target=self._switch,args=(mode,profile),daemon=True);self.thread.start()
|
||
return self.status()
|
||
def _switch(self,mode,profile):
|
||
try:
|
||
self.stop_owned();self.stop_process()
|
||
deadline=time.monotonic()+60
|
||
with self.scheduler.cv:
|
||
while self.scheduler.active or self.scheduler.transition:
|
||
if time.monotonic()>deadline:raise ValueError('Laufende Deck-Aufträge konnten noch nicht beendet werden. Erneut umschalten.')
|
||
self.scheduler.cv.wait(.2)
|
||
self.scheduler.key=None
|
||
if self.closing:raise ValueError('Deck wird beendet.')
|
||
if mode=='video':self.load(profile)
|
||
if self.closing:raise ValueError('Deck wird beendet.')
|
||
with self.lock:self.mode=mode;self.state='ready' if mode=='video' else 'idle'
|
||
with self.scheduler.cv:self.scheduler.gpu_mode=mode;self.scheduler.cv.notify_all()
|
||
except Exception as exc:
|
||
self.stop_process()
|
||
with self.lock:self.mode='llm';self.state='failed';self.error=str(exc) if isinstance(exc,ValueError) else 'Video-Moduswechsel fehlgeschlagen.'
|
||
with self.scheduler.cv:self.scheduler.gpu_mode='llm';self.scheduler.cv.notify_all()
|
||
def stop_process(self):
|
||
with self.lock:
|
||
p=self.process;self.process=None;self.profile=None;self.gpu=None
|
||
if p:
|
||
if p.poll() is None:
|
||
try:
|
||
os.killpg(p.pid,signal.SIGTERM)
|
||
try:p.wait(timeout=5)
|
||
except subprocess.TimeoutExpired:os.killpg(p.pid,signal.SIGKILL);p.wait(timeout=5)
|
||
except ProcessLookupError:pass
|
||
for stream in (p.stdin,p.stdout):
|
||
if stream:stream.close()
|
||
def close(self):
|
||
self.closing=True
|
||
self.stop_process()
|
||
if self.thread and self.thread is not threading.current_thread():self.thread.join(timeout=7)
|
||
def read(self,p,timeout):
|
||
deadline=time.monotonic()+timeout
|
||
with selectors.DefaultSelector() as selector:
|
||
selector.register(p.stdout,selectors.EVENT_READ)
|
||
while time.monotonic()<deadline:
|
||
if self.process is not p:raise ValueError('Video-Auftrag durch Moduswechsel oder Abbruch beendet.')
|
||
buffered=getattr(p,'_deck_buffer',b'')
|
||
if b'\n' not in buffered:
|
||
if not selector.select(.5):continue
|
||
chunk=os.read(p.stdout.fileno(),8192)
|
||
if not chunk:raise ValueError('Video-Worker beendet; RAM-OOM oder Laufzeitfehler möglich.')
|
||
buffered+=chunk
|
||
if len(buffered)>65536:raise ValueError('Ungültige Worker-Antwort.')
|
||
p._deck_buffer=buffered
|
||
if b'\n' not in buffered:continue
|
||
line,p._deck_buffer=buffered.split(b'\n',1)
|
||
result=json.loads(line)
|
||
if result.get('state')=='progress':
|
||
with self.lock:
|
||
if self.job:self.job['phase']=result.get('phase','Verarbeitung')
|
||
continue
|
||
return result
|
||
raise ValueError('Zeitlimit des Video-Workers erreicht.')
|
||
def load(self,profile):
|
||
gpus=probe()
|
||
if not gpus or any(g['processes'] for g in gpus):raise ValueError('GPU durch fremden Dienst belegt. Fremden Dienst zuerst freigeben; Deck beendet ihn nicht.')
|
||
gpu=next((g for g in gpus if '5080' in g.get('name','')),None)
|
||
if gpu is None:gpu=max(gpus,key=lambda g:g['total_mib'])
|
||
if gpu['free_mib']<12000:raise ValueError('Mindestens 12 GiB freier GPU-Speicher für den Video-Test erforderlich.')
|
||
available=next((int(line.split()[1])*1024 for line in Path('/proc/meminfo').read_text().splitlines() if line.startswith('MemAvailable:')),0)
|
||
if available<10*GIB:raise ValueError('Mindestens 10 GiB verfügbarer Host-RAM erforderlich.')
|
||
headroom=cgroup_headroom()
|
||
if headroom is not None and headroom<8*GIB:raise ValueError('Mindestens 8 GiB freier Deck-RAM für Disk-Streaming erforderlich.')
|
||
def path(entry):return str((self.profiles.catalog.root/entry['id']/('model'+Path(entry['file']).suffix)).resolve())
|
||
components={role:self.profiles._component(profile['model'],role,profile['components'][role]) for role in video_recipe(profile['model'])}
|
||
config=dict(paths=dict(transformer_path=path(profile['model']),text_encoder_path=path(components['text_encoder']),video_vae_path=path(components['video_vae']),audio_vae_path=path(components['audio_vae'])),spatial_upsampler=path(components['spatial_upsampler']))
|
||
env={k:v for k,v in os.environ.items() if not k.startswith(('HF_','DECK_','LLAMA_'))}
|
||
env.update(CUDA_VISIBLE_DEVICES=gpu['uuid'],HF_HUB_OFFLINE='1',TRANSFORMERS_OFFLINE='1',OMP_NUM_THREADS='2',TOKENIZERS_PARALLELISM='false')
|
||
self.root.mkdir(parents=True,exist_ok=True,mode=0o700);env.update(HOME=str(self.root),TMPDIR=str(self.root))
|
||
python,_=self.runtime.paths()
|
||
p=subprocess.Popen([str(python),str(Path(__file__).with_name('video_worker.py'))],stdin=subprocess.PIPE,stdout=subprocess.PIPE,stderr=subprocess.DEVNULL,bufsize=0,start_new_session=True,env=env)
|
||
with self.lock:self.process=p;self.gpu=gpu['uuid']
|
||
if self.closing:self.stop_process();raise ValueError('Deck wird beendet.')
|
||
p.stdin.write((json.dumps(config)+'\n').encode());p.stdin.flush();result=self.read(p,180)
|
||
if result.get('state')!='ready':raise ValueError(result.get('error','Video-Profil nicht bereit.'))
|
||
with self.lock:self.profile=profile
|
||
def generate(self,data):
|
||
with self.lock:
|
||
if self.mode!='video' or self.state!='ready' or not self.profile:raise ValueError('Video-Modus mit geladenem Profil erforderlich.')
|
||
if self.job and self.job['state']=='running':raise ValueError('Ein Video-Auftrag läuft bereits.')
|
||
if set(data)-{'prompt','width','height','fps','frames','seed','model'}:raise ValueError('Unbekannte Video-Parameter.')
|
||
if data.get('model',self.profile['name']) not in (self.profile['name'],'athena-video'):raise ValueError('Angefragtes Videoprofil ist nicht aktiv.')
|
||
prompt=data.get('prompt')
|
||
if not isinstance(prompt,str) or not 1<=len(prompt.strip())<=8000:raise ValueError('Prompt mit 1–8000 Zeichen erforderlich.')
|
||
params={k:data.get(k,self.profile['parameters'][k]) for k in ('width','height','fps','frames','seed')}
|
||
for k,lo,hi in [('width',256,1920),('height',256,1088),('fps',1,60),('frames',9,241),('seed',-1,2147483647)]:
|
||
if type(params[k]) is not int or not lo<=params[k]<=hi:raise ValueError('Ungültiger Video-Parameter: '+k)
|
||
if params['width']%64 or params['height']%64 or params['frames']%8!=1:raise ValueError('Breite/Höhe müssen durch 64 teilbar sein; Bildanzahl muss 8n+1 sein.')
|
||
if params['seed']==-1:params['seed']=int.from_bytes(os.urandom(4),'big')%2147483648
|
||
if shutil.disk_usage(self.root).free<10*GIB:raise ValueError('Mindestens 10 GiB freier Plattenspeicher erforderlich.')
|
||
ident=uuid.uuid4().hex;out=self.root/(ident+'.mp4')
|
||
self.job=dict(id=ident,state='running',phase='Videoauftrag gestartet',created_at=time.time(),parameters=params,error=None)
|
||
process=self.process
|
||
threading.Thread(target=self._generate,args=(process,dict(params,prompt=prompt,output=str(out)),ident),daemon=True).start()
|
||
return dict(self.job)
|
||
def _generate(self,p,data,ident):
|
||
try:
|
||
p.stdin.write((json.dumps(data)+'\n').encode());p.stdin.flush();result=self.read(p,3600)
|
||
if result.get('state')!='complete':raise ValueError(result.get('error','Video-Auftrag fehlgeschlagen.'))
|
||
with self.lock:
|
||
if self.job['id']==ident:self.job.update(state='complete',finished_at=time.time())
|
||
except Exception as exc:
|
||
with self.lock:
|
||
if self.job['id']==ident:self.job.update(state='failed',error=str(exc) if isinstance(exc,ValueError) else 'Video-Worker nicht erreichbar.')
|
||
if self.process is p:self.state='failed';self.error=self.job['error']
|
||
if self.process is p:self.stop_process()
|
||
(self.root/(ident+'.mp4')).unlink(missing_ok=True)
|
||
def result(self,ident):
|
||
with self.lock:
|
||
if not re.fullmatch('[a-f0-9]{32}',ident) or not self.job or self.job['id']!=ident or self.job['state']!='complete':raise ValueError('Video noch nicht verfügbar.')
|
||
return self.root/(ident+'.mp4')
|