Files
Athena-Deck/video_comfy.py
T

159 lines
12 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 secrets
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,music_profiles=getattr(self,'music_profiles',lambda:[])(),separator=getattr(self,'separator_status',lambda:{})(),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','music','separator'):raise ValueError('Modus muss llm, video, music oder separator sein.')
with self.lock:
if self.state=='switching':raise ValueError('Moduswechsel läuft bereits.')
if mode=='music' and not getattr(self,'music_profiles',lambda:[])():raise ValueError('Zuerst ein ausführbares Musikprofil am API-Endpunkt freigeben.')
if mode=='separator' and not any(m['installed'] for m in getattr(self,'separator_status',lambda:{'models':[]})().get('models',[])):raise ValueError('Zuerst Audio-Separator-Laufzeit und ein Modellpaket einrichten.')
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:
self.stop_owned();self._stop();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
if mode in ('llm','music','separator'):
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 in ('video','music','separator') else 'idle'
if mode=='separator':self.switch_phase='Audio-Trennung bereit · Gewichte laden bei Auftrag'
if mode=='music':self.switch_phase='Musik bereit · Gewichte laden bei Anfrage'
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')
swarm_nodes=self.root/'swarm-comfy-nodes'
config={'deck':{'base_path':str(work),'custom_nodes':'custom_nodes',**{key:'models/'+key for key in paths}}};
if swarm_nodes.is_dir():config['swarm']={'base_path':str(self.root),'custom_nodes':'swarm-comfy-nodes'}
(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))
if (self.root/'swarm-python').is_dir():env['PYTHONPATH']=str(self.root/'swarm-python')+os.pathsep+env.get('PYTHONPATH','')
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 service_authenticated(self,header):
if self.scheduler.gpu_mode!='video' or not header.startswith('Bearer '):return False
path=self.root/'comfy-client-token'
try:return secrets.compare_digest(header[7:].encode(),path.read_bytes().strip())
except OSError:return False
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()