Replace Deck video generation with selected external service lifecycle
This commit is contained in:
@@ -1,190 +1,71 @@
|
||||
"""Exclusive GPU mode and owned persistent LTX subprocess with asynchronous jobs."""
|
||||
"""Deck controls the original Videodienst service lifecycle, never video generation."""
|
||||
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 __init__(self,scheduler,stop_owned,helper,root):
|
||||
self.scheduler=scheduler;self.stop_owned=stop_owned;self.helper=helper;self.lock=threading.RLock();self.thread=None;self.state='idle';self.error=None;self.target_mode=None;self.switch_phase=None;self.switch_started_at=None
|
||||
self.path=Path(root)/'active-service.json';self.selected=json.loads(self.path.read_text()).get('id') if self.path.exists() else None
|
||||
self.services=[];self.service={};self.refresh()
|
||||
def refresh(self):
|
||||
try:self.services=self.helper.call('video-status')['services'];unknown=False
|
||||
except (OSError,ValueError,KeyError):self.services=[];unknown=True
|
||||
running=[s for s in self.services if s.get('running')]
|
||||
if self.selected is None and len(running)==1:self.selected=running[0]['id']
|
||||
if self.selected is None and len(self.services)==1:self.selected=self.services[0]['id']
|
||||
self.service=next((s for s in self.services if s['id']==self.selected),{})
|
||||
if self.state!='switching':
|
||||
blocked=any(s.get('running') is not False for s in self.services) or (unknown and self.selected is not None)
|
||||
with self.scheduler.cv:self.scheduler.gpu_mode='video' if blocked else 'llm';self.scheduler.cv.notify_all()
|
||||
self.state='ready' if self.service.get('ready') else 'starting' if self.service.get('running') else 'failed' if blocked else 'idle'
|
||||
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()
|
||||
self.refresh()
|
||||
if self.state=='switching' or self.scheduler.gpu_mode!='llm':raise ValueError('Zuerst auf LLM wechseln und den laufenden Videodienst beenden.')
|
||||
if not any(s['id']==ident for s in self.services):raise ValueError('Videodienst nicht eingerichtet.')
|
||||
self.path.parent.mkdir(parents=True,exist_ok=True);temp=self.path.with_suffix('.tmp');temp.write_text(json.dumps({'id':ident}));temp.replace(self.path);self.selected=ident
|
||||
return self.status()
|
||||
def status(self):
|
||||
with self.lock:
|
||||
self.refresh()
|
||||
return dict(mode=self.scheduler.gpu_mode,state=self.state,error=self.error,target_mode=self.target_mode,switch_phase=self.switch_phase,switch_started_at=self.switch_started_at,service=dict(self.service),services=list(self.services),selected=self.selected,memory_policy='Deck steuert nur den aktiven Videodienst. Modelle, Eingaben und Generierung werden über die Original-API des Dienstes verwaltet.')
|
||||
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.
|
||||
self.refresh()
|
||||
if not self.service:raise ValueError('Zuerst einen eingerichteten Videodienst unter Weitere Dienste aktivieren.')
|
||||
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()
|
||||
self.state='switching';self.error=None;self.target_mode=mode;self.switch_started_at=time.time();self.switch_phase='Deck-GPU-Aufträge beenden' if mode=='video' else 'Videodienst beenden und GPU-Speicher freigeben'
|
||||
self.thread=threading.Thread(target=self._switch,args=(mode,),daemon=True);self.thread.start()
|
||||
return self.status()
|
||||
def _switch(self,mode,profile):
|
||||
def _switch(self,mode):
|
||||
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()
|
||||
self.stop_owned()
|
||||
deadline=time.monotonic()+60
|
||||
with self.scheduler.cv:
|
||||
while self.scheduler.active or self.scheduler.transition:
|
||||
if time.monotonic()>deadline:raise ValueError('Deck-Aufträge konnten noch nicht beendet werden.')
|
||||
self.scheduler.cv.wait(.2)
|
||||
self.scheduler.key=None
|
||||
self.switch_phase='Freie GPUs prüfen und Videodienst-Dienst starten'
|
||||
self.helper.call('video-start',self.selected);self.switch_phase='Auf Original-API des Dienstes warten'
|
||||
deadline=time.monotonic()+240
|
||||
while True:
|
||||
service=next(s for s in self.helper.call('video-status')['services'] if s['id']==self.selected)
|
||||
if service.get('ready'):break
|
||||
if not service.get('running'):raise ValueError('Videodienst-Dienst wurde beendet. Installation prüfen.')
|
||||
if time.monotonic()>deadline:raise ValueError('Videodienst läuft, seine API ist noch nicht bereit. Erneut prüfen oder auf LLM zurückschalten.')
|
||||
time.sleep(2)
|
||||
else:
|
||||
for service in self.helper.call('video-status')['services']:
|
||||
if service.get('running'):self.helper.call('video-stop',service['id'])
|
||||
except Exception as exc:self.error=str(exc) if isinstance(exc,ValueError) else 'Videodienst-Systemhelfer nicht erreichbar.'
|
||||
finally:
|
||||
with self.lock:self.state='idle';self.refresh()
|
||||
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')
|
||||
# Do not terminate an external service on a UI restart. Restore its gate at startup.
|
||||
if self.thread and self.thread.is_alive():self.thread.join(5)
|
||||
|
||||
Reference in New Issue
Block a user