"""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 ident is not None and 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()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')