Files
Athena-Deck/image_test.py
T

167 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.
"""Owned, serial ComfyUI test jobs. Never talks to the production image router."""
import json
import os
from pathlib import Path
import secrets
import shutil
import signal
import socket
import subprocess
import threading
import time
import urllib.request
import uuid
from profiles import QWEN_REPO
PYTHON=Path('/opt/deck-image-python/bin/python')
COMFY=Path('/opt/deck-comfy')
GIB=1024**3
def probe():
rows=subprocess.run(['nvidia-smi','--query-gpu=uuid,name,memory.total,memory.free','--format=csv,noheader,nounits'],capture_output=True,text=True,check=True,timeout=5).stdout
apps=subprocess.run(['nvidia-smi','--query-compute-apps=gpu_uuid,pid','--format=csv,noheader,nounits'],capture_output=True,text=True,check=True,timeout=5).stdout
counts={}
for line in apps.splitlines():
gpu,_=line.split(',',1);counts[gpu.strip()]=counts.get(gpu.strip(),0)+1
return [dict(uuid=u.strip(),name=n.strip(),total_mib=float(t),free_mib=float(f),processes=counts.get(u.strip(),0)) for u,n,t,f in (line.split(',') for line in rows.splitlines())]
def workflow(prompt,params,seed):
return {
'1':{'class_type':'UnetLoaderGGUF','inputs':{'unet_name':'model.gguf'}},
'2':{'class_type':'CLIPLoader','inputs':{'clip_name':'encoder.safetensors','type':'qwen_image','device':'cpu'}},
'3':{'class_type':'VAELoader','inputs':{'vae_name':'vae.safetensors'}},
'4':{'class_type':'TextEncodeQwenImage21','inputs':{'clip':['2',0],'prompt':prompt,'negative_prompt':'','resolution':max(params['width'],params['height'])}},
'5':{'class_type':'EmptyLatentImage','inputs':{'width':params['width'],'height':params['height'],'batch_size':1}},
'6':{'class_type':'KSampler','inputs':{'model':['1',0],'positive':['4',0],'negative':['4',1],'latent_image':['5',0],'seed':seed,'steps':params['steps'],'cfg':params['guidance'],'sampler_name':'euler','scheduler':'simple','denoise':1}},
'7':{'class_type':'VAEDecode','inputs':{'samples':['6',0],'vae':['3',0]}},
'8':{'class_type':'SaveImage','inputs':{'filename_prefix':'result','images':['7',0]}}
}
class ImageTests:
def __init__(self,root,profiles):
self.root=Path(root);self.profiles=profiles;self.lock=threading.RLock();self.process=None;self.cancel=threading.Event();self.job=None
if (self.root/'status.json').exists():
self.job=json.loads((self.root/'status.json').read_text())
if self.job.get('state')=='running':self.job.update(state='interrupted',phase='Deck wurde neu gestartet; Test erneut starten.')
def _save(self):
self.root.mkdir(parents=True,exist_ok=True,mode=0o700)
p=self.root/'status.tmp';p.write_text(json.dumps(self.job));p.replace(self.root/'status.json')
def _phase(self,text):
with self.lock:
self.job['phase']=text;self._save()
def status(self):
with self.lock:return dict(job=dict(self.job) if self.job else None,runtime_installed=PYTHON.is_file() and (COMFY/'main.py').is_file())
def stop(self):
self.cancel.set()
with self.lock:process=self.process
if process and process.poll() is None:
try:
os.killpg(process.pid,signal.SIGTERM)
try:process.wait(timeout=3)
except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait()
except ProcessLookupError:pass
return {'cancellation_requested':True}
def model_path(self,item):
return (self.profiles.catalog.root/item['id']/('model'+Path(item['file']).suffix)).resolve()
def start(self,profile_id,prompt):
if not isinstance(prompt,str) or not 1<=len(prompt.strip())<=4000:raise ValueError('Bitte einen Prompt mit 1–4000 Zeichen eingeben.')
with self.lock:
if self.job and self.job['state']=='running':raise ValueError('Ein Bildtest läuft bereits.')
if not self.status()['runtime_installed']:raise ValueError('Die eigene ComfyUI-Bildlaufzeit ist noch nicht installiert.')
profile=next((p for p in self.profiles.status()['profiles'] if p['id']==profile_id),None)
if not profile or profile['kind']!='image' or not profile['model']:raise ValueError('Bildprofil nicht verfügbar.')
model=profile['model']
if model['repo']!=QWEN_REPO or not model['file'].endswith('.gguf'):raise ValueError('Der Bildtest unterstützt zunächst Qwen-Image-2.1 GGUF aus dem hinterlegten Rezept.')
params=profile['parameters']
if params['width']>1024 or params['height']>1024:raise ValueError('Der isolierte Test unterstützt maximal 1024 × 1024 Pixel.')
encoder=self.profiles._component(model,'text_encoder',profile.get('components',{}).get('text_encoder'))
vae=self.profiles._component(model,'vae',profile.get('components',{}).get('vae'))
mem={line.split(':')[0]:int(line.split()[1])*1024 for line in Path('/proc/meminfo').read_text().splitlines() if line.startswith(('MemAvailable:','MemTotal:'))}
required_ram=max(20*GIB,encoder['size']*2+4*GIB)
try:
raw=Path('/sys/fs/cgroup/memory.max').read_text().strip()
current=int(Path('/sys/fs/cgroup/memory.current').read_text())
if raw!='max' and int(raw)-current<required_ram:raise ValueError('Deck-RAM-Limit reicht für Textencoder und Arbeitsdaten nicht aus.')
except FileNotFoundError:pass
if mem.get('MemAvailable',0)<required_ram+4*GIB:raise ValueError('Aktuell zu wenig freier System-RAM; produktive Dienste bleiben unverändert.')
needed=(model['size']+vae['size']+3*GIB)/1024**2
gpus=[g for g in probe() if g['processes']==0 and g['free_mib']>needed]
if not gpus:raise ValueError('Keine unbelegte GPU mit ausreichend freiem VRAM. Deck stoppt keine produktiven Modelle.')
gpu=max(gpus,key=lambda g:g['free_mib']);self.root.mkdir(parents=True,exist_ok=True,mode=0o700)
if shutil.disk_usage(self.root).free<10*GIB:raise ValueError('Weniger als 10 GiB freier Plattenspeicher.')
job_id=uuid.uuid4().hex;seed=params['seed'] if params['seed']>=0 else secrets.randbelow(2147483648)
self.cancel.clear();self.job=dict(id=job_id,state='running',phase='Bildlaufzeit startet',profile_id=profile_id,profile_name=profile['name'],started_at=time.time(),gpu=gpu['name'],seed=seed,parameters=params)
self._save();threading.Thread(target=self._run,args=(job_id,prompt,params,seed,gpu,model,encoder,vae),daemon=True).start()
return dict(self.job)
def _run(self,job_id,prompt,params,seed,gpu,model,encoder,vae):
directory=self.root/job_id;process=None
try:
directory.mkdir(mode=0o700)
for role,item,filename in [('unet',model,'model.gguf'),('clip',encoder,'encoder.safetensors'),('vae',vae,'vae.safetensors')]:
dest=directory/'models'/role;dest.mkdir(parents=True);(dest/filename).symlink_to(self.model_path(item))
for folder in ('output','temp','user','input'):(directory/folder).mkdir()
config={'deck':{'base_path':str(directory/'models'),'unet':'unet','clip':'clip','vae':'vae'}}
(directory/'paths.json').write_text(json.dumps(config)) # JSON is valid YAML.
with socket.socket() as sock:sock.bind(('127.0.0.1',0));port=sock.getsockname()[1]
env=dict(os.environ,CUDA_VISIBLE_DEVICES=gpu['uuid'],OMP_NUM_THREADS='2',MKL_NUM_THREADS='2',HOME=str(directory),HF_HUB_OFFLINE='1',TRANSFORMERS_OFFLINE='1',PYTHONDONTWRITEBYTECODE='1')
args=[str(PYTHON),str(COMFY/'main.py'),'--listen','127.0.0.1','--port',str(port),'--disable-auto-launch','--disable-metadata','--lowvram','--reserve-vram','1.5','--extra-model-paths-config',str(directory/'paths.json'),'--output-directory',str(directory/'output'),'--temp-directory',str(directory/'temp'),'--user-directory',str(directory/'user'),'--input-directory',str(directory/'input')]
with self.lock:
if self.cancel.is_set():raise InterruptedError()
process=subprocess.Popen(args,cwd=COMFY,env=env,stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL,start_new_session=True);self.process=process
def request(path,data=None):
req=urllib.request.Request(f'http://127.0.0.1:{port}'+path,data=json.dumps(data).encode() if data is not None else None,headers={'Content-Type':'application/json'})
with urllib.request.urlopen(req,timeout=5) as r:
raw=r.read(4*1024*1024+1)
if len(raw)>4*1024*1024:raise ValueError('Bildlaufzeit-Antwort zu groß.')
return json.loads(raw)
def guard():
if self.cancel.is_set():raise InterruptedError()
if process.poll() is not None:raise ValueError('Bildlaufzeit wurde beendet (Speicherlimit oder Startfehler).')
current=next((g for g in probe() if g['uuid']==gpu['uuid']),None)
if not current or current['processes']>1:raise ValueError('GPU wird inzwischen von einem weiteren Prozess verwendet. Deck-Test beendet.')
deadline=time.monotonic()+180
while True:
guard()
try:request('/system_stats');break
except (OSError,ValueError):
if time.monotonic()>deadline:raise ValueError('Bildlaufzeit wurde nicht rechtzeitig bereit.')
time.sleep(1)
nodes=request('/object_info')
required={'UnetLoaderGGUF','CLIPLoader','VAELoader','TextEncodeQwenImage21','EmptyLatentImage','KSampler','VAEDecode','SaveImage'}
if not required.issubset(nodes):raise ValueError('Der installierten Bildlaufzeit fehlen erforderliche Qwen/GGUF-Nodes.')
self._phase('Modell und Textencoder laden · Bild wird erzeugt')
response=request('/prompt',{'prompt':workflow(prompt,params,seed),'client_id':'deck-'+job_id});prompt_id=response['prompt_id']
deadline=time.monotonic()+1800
while time.monotonic()<deadline:
guard();history=request('/history/'+prompt_id)
if prompt_id in history:
result=history[prompt_id]
if result.get('status',{}).get('status_str')=='error':
errors=[m[1] for m in result.get('status',{}).get('messages',[]) if m[0]=='execution_error']
error=errors[0] if errors else {}
# Do not include prompts, tensor dumps or upstream exception messages.
raise ValueError('Generierung fehlgeschlagen: '+str(error.get('node_type','Worker'))+' / '+str(error.get('exception_type','Fehler')))
images=[i for value in result.get('outputs',{}).values() for i in value.get('images',[])]
if images:
image=images[0];path=(directory/'output'/image.get('subfolder','')/image['filename']).resolve()
if not path.is_relative_to((directory/'output').resolve()) or path.stat().st_size>32*1024**2:raise ValueError('Ungültiges Ergebnisbild.')
with path.open('rb') as f:
if f.read(8)!=b'\x89PNG\r\n\x1a\n':raise ValueError('Ergebnis ist kein PNG.')
shutil.copyfile(path,directory/'result.png');break
time.sleep(2)
else:raise ValueError('Zeitlimit der Bildgenerierung erreicht.')
final_state='complete';phase='Bild fertig · Modell wurde entladen'
except InterruptedError:final_state='cancelled';phase='Bildtest abgebrochen'
except Exception as exc:final_state='failed';phase=str(exc) if isinstance(exc,ValueError) else 'Bildlaufzeit nicht erreichbar oder nicht bereit. Komponenten und Installation prüfen.'
finally:
if process and process.poll() is None:
try:os.killpg(process.pid,signal.SIGTERM);process.wait(timeout=10)
except subprocess.TimeoutExpired:os.killpg(process.pid,signal.SIGKILL);process.wait()
except ProcessLookupError:pass
with self.lock:
self.process=None;self.job.update(state=final_state,phase=phase,finished_at=time.time());self._save()
def image(self,job_id):
with self.lock:
if not self.job or self.job['id']!=job_id or self.job['state']!='complete':raise ValueError('Ergebnisbild nicht verfügbar.')
return (self.root/job_id/'result.png').read_bytes()