Files
Athena-Deck/video.py
T

191 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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,video_parameters
def video_devices(profile,gpus):
if not gpus:raise ValueError('Keine Video-GPU verfügbar.')
params=video_parameters(profile.get('parameters',{}))
def selected(value):
found=next((g for g in gpus if g['uuid']==value),None)
if not found:raise ValueError('Im Videoprofil ausgewählte GPU nicht verfügbar: '+value)
return found
gpu=(next((g for g in gpus if '5080' in g.get('name','')),None) or max(gpus,key=lambda g:g['total_mib'])) if params['video_device']=='auto' else selected(params['video_device'])
encoder=gpu if params['text_encoder_device']=='same' else selected(params['text_encoder_device'])
devices=[gpu] if encoder['uuid']==gpu['uuid'] else [gpu,encoder]
return gpu,encoder,devices
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.device_names={}
self.switch_phase=None;self.switch_started_at=None;self.target_mode=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(devices=dict(self.device_names),target_mode=self.target_mode,switch_phase=self.switch_phase,switch_started_at=self.switch_started_at,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 den im Profil gewählten GPUs; 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.target_mode=mode;self.switch_started_at=time.time();self.switch_phase='Laufende Deck-Aufträge beenden'
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()
with self.lock:self.switch_phase='Auf Freigabe des GPU-Speichers warten'
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':
with self.lock:self.switch_phase='GPUs prüfen und Videoprofil vorbereiten'
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,encoder,devices=video_devices(profile,gpus)
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']))
config['text_encoder_device']='cuda:0' if len(devices)==1 else 'cuda:1'
config['device_names']={'video':gpu['name'],'text_encoder':encoder['name']}
env={k:v for k,v in os.environ.items() if not k.startswith(('HF_','DECK_','LLAMA_'))}
env.update(CUDA_VISIBLE_DEVICES=','.join(g['uuid'] for g in devices),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'];self.device_names=dict(config['device_names'])
if self.closing:self.stop_process();raise ValueError('Deck wird beendet.')
with self.lock:self.switch_phase='Video-Worker laden und Komponenten prüfen'
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')