"""Bounded, synthetic model tuning; exclusive shared scheduler, no foreign process control.""" import copy,json,threading,time,uuid,socket from pathlib import Path from profiles import SCHEMAS,CHAT_GPU_DEFAULTS from image_test import probe from inference import InferenceError class AutoTests: def __init__(self,path,profiles,worker,scheduler): self.path=Path(path);self.profiles=profiles;self.worker=worker;self.scheduler=scheduler;self.lock=threading.RLock();self.cancel=threading.Event();self.conn=None;self.thread=None self.job=json.loads(self.path.read_text()) if self.path.exists() else None if self.job and self.job['state']=='running':self.job.update(state='interrupted',phase='Durch Neustart unterbrochen');self.persist() def persist(self): self.path.parent.mkdir(parents=True,exist_ok=True);p=self.path.with_suffix('.tmp');p.write_text(json.dumps(self.job));p.replace(self.path) def status(self): with self.lock:return {'job':copy.deepcopy(self.job)} def update(self,**kw): with self.lock:self.job.update(kw);self.persist() def start(self,data): if set(data)!={'model_id','max_context','mtp'} or type(data['max_context'])!=int or data['max_context'] not in (4096,8192,16384,32768,65536,131072,160000,192000,262144) or type(data['mtp'])!=bool:raise ValueError('Modell, Kontextgrenze und MTP auswählen.') model=self.worker.catalog.entry(data['model_id']) if model['kind']!='chat' or not model['file'].endswith('.gguf') or not model['profile_eligible']:raise ValueError('Ein heruntergeladenes Chat-GGUF auswählen.') devices=probe();primary=next((g for g in devices if '5080' in g['name']),None);secondary=next((g for g in devices if '3060' in g['name']),None) if not primary or not secondary:raise ValueError('Dieser Auto-Test benötigt RTX 5080 und RTX 3060. Keine stillschweigende andere GPU-Zuordnung.') with self.lock: if self.thread and self.thread.is_alive():raise ValueError('Ein Auto-Test läuft bereits.') self.cancel.clear();self.job=dict(id=uuid.uuid4().hex,state='running',phase='Wartet auf exklusive Modellreservierung',model_id=model['id'],model_file=model['file'],max_context=data['max_context'],mtp=data['mtp'],results=[],attempts=0,started_at=time.time(),error=None);self.persist() self.thread=threading.Thread(target=self.run,args=(primary['uuid'],secondary['uuid']),daemon=True);self.thread.start() return self.status() def stop(self): self.cancel.set() with self.lock:conn=self.conn if conn and conn.sock: try:conn.sock.shutdown(socket.SHUT_RDWR) except OSError:pass return {'cancellation_requested':True} def request(self,path,data): if self.cancel.is_set():raise InterruptedError() conn,key=self.worker.connect();conn.timeout=180 with self.lock:self.conn=conn try: conn.request('POST',path,json.dumps(data).encode(),{'Content-Type':'application/json','Authorization':'Bearer '+key}) response=conn.getresponse();raw=response.read(8*1024*1024) if response.status!=200:raise InferenceError('Synthetischer Test abgelehnt: Kontext oder Modell prüfen.') return json.loads(raw) finally: conn.close() with self.lock:self.conn=None def benchmark(self,context): # Tokenize synthetic text with this model; do not estimate token count from characters. target=int(context*.75);unit='alpha beta gamma delta epsilon zeta eta theta 0123456789. ' text=unit*(target//4+1);tokens=self.request('/tokenize',{'content':text,'add_special':True})['tokens'] if len(tokens)=48 or time.time()-self.job['started_at']>7200: self.update(state='complete',phase='Testbudget erreicht · Teilresultate verfügbar',finished_at=time.time());return params={**{k:v[2] for k,v in SCHEMAS['chat'].items()},**CHAT_GPU_DEFAULTS,'context':context,'slots':1,'batch':2048,'ubatch':micro,'gpu_devices':devices,'split_mode':split,'tensor_split':ratio,'gpu_offload':offload,'mtp':self.job['mtp']} p=dict(id='auto-'+self.job['id'],revision=self.job['attempts']+1,name='deck-auto-test',kind='chat',model_id=self.job['model_id'],parameters=params) self.update(attempts=p['revision'],phase=f'{context} Kontext · {tier} · Microbatch {micro}') row=dict(context=context,tier=tier,parameters=params,success=False) self.worker.stop() try: self.worker.plan(p,cancel=self.cancel.is_set) self.worker.ensure(p,cancel=self.cancel.is_set) # Small warm-up, then populated-context measurement, bounded output. self.request('/completion',dict(prompt='Synthetic warmup.',n_predict=16,temperature=0,cache_prompt=False)) row.update(success=True,timings=self.benchmark(context),memory_plan=copy.deepcopy(self.worker.memory_plan),tested_context_fraction=.75) successes.append(row) except Exception as exc: if self.cancel.is_set():raise InterruptedError() row['error']=str(exc) if isinstance(exc,ValueError) else 'Testprozess oder Verbindung fehlgeschlagen; kein eindeutiger OOM-Nachweis.' finally:self.worker.stop() with self.lock:self.job['results'].append(row);self.persist() # Smaller microbatch is useful after rejection; do not skip it. if successes:found=True;break if not found:break self.update(state='complete',phase='Testreihe abgeschlossen · Ergebnisse sind keine Garantie für beliebige Last',finished_at=time.time()) finally:self.worker.stop() except Exception as exc:self.update(state='cancelled' if self.cancel.is_set() else 'failed',phase='Abgebrochen' if self.cancel.is_set() else 'Test beendet',error=None if self.cancel.is_set() else str(exc),finished_at=time.time()) def save(self,data): with self.lock: if set(data)!={'job_id','index','name'} or not self.job or data['job_id']!=self.job['id'] or type(data['index'])!=int or not 0<=data['index']