140 lines
10 KiB
Python
140 lines
10 KiB
Python
"""Deck-owned LTX process, sharing the installed image ComfyUI environment."""
|
|
import json
|
|
import os
|
|
from pathlib import Path, PurePosixPath
|
|
import shutil
|
|
import signal
|
|
import socket
|
|
import subprocess
|
|
import threading
|
|
import time
|
|
import urllib.request
|
|
from profiles import video_recipe
|
|
|
|
class VideoComfy:
|
|
def __init__(self,scheduler,stop_owned,helper,root,catalog,profiles,runtime):
|
|
self.scheduler=scheduler;self.stop_owned=stop_owned;self.helper=helper;self.root=Path(root);self.catalog=catalog;self.profiles=profiles;self.runtime=runtime
|
|
self.lock=threading.RLock();self.process=None;self.thread=None;self.port=None;self.state='idle';self.error=None;self.target_mode=None;self.switch_phase=None;self.switch_started_at=None
|
|
self.path=self.root/'comfy-selection.json';self.selection=json.loads(self.path.read_text()) if self.path.exists() else {}
|
|
self.foreign_blocked=False
|
|
try:self.foreign_blocked=any(x.get('running') is not False for x in self._foreign_services())
|
|
except (OSError,ValueError,KeyError,TypeError):self.foreign_blocked=True
|
|
if self.foreign_blocked:
|
|
self.scheduler.gpu_mode='video';self.state='failed';self.error='Fremder Videodienst läuft oder ist nicht prüfbar; GPU-Aufträge gesperrt.'
|
|
def _foreign_services(self):
|
|
path=Path(os.environ.get('DECK_DOCKER_HELPER_SOCKET','/run/athena-deck-docker/control.sock'))
|
|
if not path.is_socket() and not os.environ.get('DECK_DOCKER_HELPER_SOCKET'):return []
|
|
return self.helper.call('video-status')['services']
|
|
def models(self):
|
|
entries=self.catalog.status()['entries'];out=[]
|
|
for model in entries:
|
|
if not video_recipe(model):continue
|
|
components={};blockers=[]
|
|
try:self.catalog.entry(model['id'])
|
|
except ValueError as exc:blockers.append(str(exc))
|
|
for role,info in video_recipe(model).items():
|
|
item=next((x for x in entries if x['repo']==model['repo'] and x['revision']==model['revision'] and x['file'] in info['files']),None)
|
|
try:components[role]=self.profiles._component(model,role,item['id'] if item else None)['id']
|
|
except ValueError:blockers.append(info['label']+' fehlt oder ist unvollständig.')
|
|
python,comfy=self.runtime.paths()
|
|
if not python.is_file() or not (comfy/'comfy/ldm/lightricks').is_dir():blockers.append('Passende ComfyUI-LTX-Anbindung fehlt in der vorhandenen Laufzeit.')
|
|
out.append(dict(model=model,components=components,blockers=blockers,runnable=not blockers,state='ready' if not blockers else 'configured',label='Laufzeit bereit · lädt bei Anfrage' if not blockers else 'Voraussetzungen fehlen'))
|
|
return out
|
|
def select(self,ident):
|
|
with self.lock:
|
|
if self.scheduler.gpu_mode!='llm' or self.state=='switching':raise ValueError('Zuerst auf LLM wechseln.')
|
|
row=next((x for x in self.models() if x['model']['id']==ident),None)
|
|
if ident and not row:raise ValueError('Unterstütztes LTX-Modell aus der Bibliothek auswählen.')
|
|
self.root.mkdir(parents=True,exist_ok=True);self.selection={'model_id':ident or None};temp=self.path.with_suffix('.tmp');temp.write_text(json.dumps(self.selection));temp.replace(self.path)
|
|
return self.status()
|
|
def status(self):
|
|
with self.lock:
|
|
running=bool(self.process and self.process.poll() is None)
|
|
if self.process and not running and self.state!='switching':
|
|
self.state='failed';self.error='ComfyUI wurde unerwartet beendet. Auf LLM wechseln und Laufzeit erneut starten.'
|
|
rows=self.models();ident=self.selection.get('model_id');row=next((x for x in rows if x['model']['id']==ident),None)
|
|
service=dict(id=ident,name=row['model']['file'] if row else 'Kein Videomodell ausgewählt',running=running,ready=running and self.state=='ready',api_url='/ (ComfyUI auf dem gemeinsamen API-Port)',message='5080: Diffusion/VAE · 3060: Textencoder · große BF16-Gewichte nutzen CPU-Auslagerung.') if row else {}
|
|
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=service,services=[dict(id=x['model']['id'],name=x['model']['file'],api_url=x['label'],**{k:x[k] for k in ('blockers','runnable')}) for x in rows],models=rows,selected=ident,memory_policy='Gemeinsame ComfyUI-Installation. Video startet einen eigenen Prozess; Gewichte laden bei Workflow-Anfrage. Der Wechsel auf LLM beendet den Prozess.',runtime=self.runtime.status())
|
|
def switch(self,mode):
|
|
if mode not in ('llm','video'):raise ValueError('Modus muss llm oder video sein.')
|
|
with self.lock:
|
|
if self.state=='switching':raise ValueError('Moduswechsel läuft bereits.')
|
|
if mode=='video':
|
|
row=next((x for x in self.models() if x['model']['id']==self.selection.get('model_id')),None)
|
|
if not row:raise ValueError('Zuerst Videomodell aktivieren.')
|
|
if row['blockers']:raise ValueError(' '.join(row['blockers']))
|
|
self.state='switching';self.error=None;self.target_mode=mode;self.switch_phase='Deck-GPU-Aufträge beenden' if mode=='video' else 'Video entladen';self.switch_started_at=time.time()
|
|
with self.scheduler.cv:self.scheduler.gpu_mode='switching';self.scheduler.cv.notify_all()
|
|
self.thread=threading.Thread(target=self._switch,args=(mode,),daemon=True);self.thread.start()
|
|
return self.status()
|
|
def _switch(self,mode):
|
|
try:
|
|
if mode=='llm':
|
|
self._stop()
|
|
if any(x.get('running') is not False for x in self._foreign_services()):raise ValueError('Fremder Videodienst blockiert LLM-Modus.')
|
|
self.foreign_blocked=False
|
|
else:
|
|
foreign=self._foreign_services()
|
|
self.foreign_blocked=any(x.get('running') is not False for x in foreign)
|
|
if self.foreign_blocked:raise ValueError('Fremder Videodienst läuft oder ist nicht prüfbar. Er wird nicht von Deck beendet.')
|
|
self.stop_owned();deadline=time.monotonic()+60
|
|
with self.scheduler.cv:
|
|
while self.scheduler.active or self.scheduler.tts_active or self.scheduler.transition:
|
|
if time.monotonic()>deadline:raise ValueError('Deck-Aufträge noch aktiv.')
|
|
self.scheduler.cv.wait(.2)
|
|
self.scheduler.key=None
|
|
self.switch_phase='Freie GPUs prüfen';self._start()
|
|
with self.scheduler.cv:self.scheduler.gpu_mode=mode;self.scheduler.cv.notify_all()
|
|
self.state='ready' if mode=='video' else 'idle'
|
|
except Exception as exc:
|
|
self._stop();self.error=str(exc);self.state='failed'
|
|
with self.scheduler.cv:self.scheduler.gpu_mode='video' if self.foreign_blocked else 'llm';self.scheduler.cv.notify_all()
|
|
def _start(self):
|
|
raw=subprocess.check_output(['nvidia-smi','--query-gpu=uuid,name,memory.used','--format=csv,noheader,nounits'],text=True)
|
|
gpus=[line.split(',') for line in raw.strip().splitlines()]
|
|
ordered=[next((g for g in gpus if name in g[1]),None) for name in ('5080','3060')]
|
|
if any(g is None for g in ordered):raise ValueError('RTX 5080 und RTX 3060 benötigt.')
|
|
if any(float(g[2])>256 for g in ordered):raise ValueError('GPUs sind nicht frei; fremde Prozesse werden nicht beendet.')
|
|
row=next(x for x in self.models() if x['model']['id']==self.selection['model_id'])
|
|
work=self.root/'comfy-work';work.mkdir(parents=True,exist_ok=True)
|
|
paths={'diffusion_models':row['model']['id'],'text_encoders':row['components']['text_encoder'],'vae':None,'latent_upscale_models':row['components']['spatial_upsampler']}
|
|
for key in paths:(work/'models'/key).mkdir(parents=True,exist_ok=True)
|
|
for folder,ident in [('diffusion_models',row['model']['id']),('text_encoders',row['components']['text_encoder']),('vae',row['components']['video_vae']),('vae',row['components']['audio_vae']),('latent_upscale_models',row['components']['spatial_upsampler'])]:
|
|
item=self.catalog.entry(ident);source=self.catalog.root/ident/('model'+PurePosixPath(item['file']).suffix);link=work/'models'/folder/PurePosixPath(item['file']).name
|
|
for old in link.parent.iterdir():
|
|
if old.is_symlink() and old.name!=link.name and folder!='vae':old.unlink()
|
|
if link.is_symlink():link.unlink()
|
|
link.symlink_to(source)
|
|
node=work/'custom_nodes/deck_ltx';node.mkdir(parents=True,exist_ok=True);shutil.copyfile(Path(__file__).with_name('video_comfy_node.py'),node/'__init__.py')
|
|
config={'deck':{'base_path':str(work),'custom_nodes':'custom_nodes',**{key:'models/'+key for key in paths}}};(work/'paths.json').write_text(json.dumps(config))
|
|
for directory in ('input','output','temp','user'):(work/directory).mkdir(exist_ok=True)
|
|
with socket.socket() as sock:sock.bind(('127.0.0.1',0));self.port=sock.getsockname()[1]
|
|
python,comfy=self.runtime.paths();env=dict(os.environ,CUDA_VISIBLE_DEVICES=','.join(g[0].strip() for g in ordered))
|
|
args=[str(python),str(comfy/'main.py'),'--listen','127.0.0.1','--port',str(self.port),'--disable-auto-launch','--disable-metadata','--lowvram','--reserve-vram','1.5','--extra-model-paths-config',str(work/'paths.json')]
|
|
for directory in ('input','output','temp','user'):args+=['--'+directory+'-directory',str(work/directory)]
|
|
self.switch_phase='ComfyUI starten · Gewichte laden erst bei Anfrage';self.process=subprocess.Popen(args,env=env,cwd=comfy,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True)
|
|
deadline=time.monotonic()+180
|
|
while time.monotonic()<deadline:
|
|
if self.process.poll() is not None:raise ValueError('ComfyUI-Start fehlgeschlagen; keine Gewichte geladen.')
|
|
try:
|
|
with urllib.request.urlopen(f'http://127.0.0.1:{self.port}/object_info',timeout=2) as response:info=json.load(response)
|
|
required={'DeckLTXTextEncoderLoader','LTXVEmptyLatentAudio','LTXVConcatAVLatent','LTXVAudioVAEDecode'}
|
|
if not required<=info.keys():raise ValueError('ComfyUI fehlen erforderliche LTX-Nodes.')
|
|
return
|
|
except (OSError,TimeoutError):time.sleep(.5)
|
|
raise ValueError('ComfyUI ist noch nicht bereit.')
|
|
def _stop(self):
|
|
process=self.process
|
|
if process and process.poll() is None:
|
|
os.killpg(process.pid,signal.SIGTERM)
|
|
try:process.wait(timeout=15)
|
|
except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait(timeout=10)
|
|
self.process=None;self.port=None
|
|
def relay(self,handler):
|
|
from video_comfy_proxy import relay
|
|
if not self.port or not self.process or self.process.poll() is not None:raise ValueError('ComfyUI nicht bereit.')
|
|
relay(handler,self.port)
|
|
def close(self):
|
|
if self.thread and self.thread.is_alive():self.thread.join(190)
|
|
self._stop()
|