187 lines
12 KiB
Python
187 lines
12 KiB
Python
"""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
|
|
|
|
CONTEXT_STEPS=(2048,4096,8192,16384,32768,65536,131072,160000,192000,262144)
|
|
MAX_CANDIDATES=300
|
|
MAX_SECONDS=7200
|
|
GPU_SPLITS=([98,2],[95,5],[90,10],[85,15],[75,25],[50,50])
|
|
|
|
def model_layers(path,mtp):
|
|
"""Read bounded GGUF metadata only; never load tensors or tokenizer strings."""
|
|
import struct
|
|
formats={0:'B',1:'b',2:'H',3:'h',4:'I',5:'i',6:'f',7:'?',10:'Q',11:'q',12:'d'}
|
|
with open(path,'rb') as f:
|
|
size=Path(path).stat().st_size
|
|
def read(n):
|
|
b=f.read(n)
|
|
if len(b)!=n:raise ValueError('Unvollständige GGUF-Metadaten.')
|
|
return b
|
|
def num(fmt):return struct.unpack('<'+fmt,read(struct.calcsize('<'+fmt)))[0]
|
|
def string(keep=False):
|
|
n=num('Q')
|
|
if n>size-f.tell() or (keep and n>65536):raise ValueError('Ungültige GGUF-Zeichenfolge.')
|
|
if keep:return read(n).decode('utf-8')
|
|
f.seek(n,1)
|
|
def value(t,keep=False,depth=0):
|
|
if t in formats:return num(formats[t])
|
|
if t==8:return string(keep)
|
|
if t==9 and depth<2:
|
|
subtype=num('I');count=num('Q')
|
|
if count>10000000:raise ValueError('GGUF-Array zu groß.')
|
|
if subtype in formats:
|
|
length=struct.calcsize('<'+formats[subtype])*count
|
|
if length>size-f.tell():raise ValueError('Ungültiges GGUF-Array.')
|
|
f.seek(length,1)
|
|
else:
|
|
for _ in range(count):value(subtype,False,depth+1)
|
|
return None
|
|
raise ValueError('Unbekannter GGUF-Metadatentyp.')
|
|
if read(4)!=b'GGUF' or num('I') not in (2,3):raise ValueError('GGUF-Version nicht unterstützt.')
|
|
num('Q');count=num('Q');fields={}
|
|
if count>100000:raise ValueError('Zu viele GGUF-Felder.')
|
|
for _ in range(count):
|
|
key=string(True);t=num('I');wanted=key=='general.architecture' or key.endswith(('.block_count','.nextn_predict_layers'))
|
|
v=value(t,wanted)
|
|
if wanted:fields[key]=v
|
|
arch=fields.get('general.architecture');blocks=fields.get(str(arch)+'.block_count');draft=fields.get(str(arch)+'.nextn_predict_layers',0)
|
|
if type(blocks)!=int or type(draft)!=int or not 1<=blocks<=4096 or not 0<=draft<blocks:raise ValueError('Keine verlässliche Layerzahl im GGUF.')
|
|
return blocks-(0 if mtp else draft)+1 # output layer also participates in layer split
|
|
|
|
def fine_splits(total,ratio,upper):
|
|
"""llama layer split uses ceil(fraction * total) layers on the first GPU."""
|
|
import math
|
|
first=math.ceil(total*ratio[0]/sum(ratio))
|
|
stop=min(total,math.ceil(total*upper/100))
|
|
return [(k,[100*(k-.5)/total,100-100*(k-.5)/total]) for k in range(first+1,stop)]
|
|
|
|
class TestBudget(Exception):
|
|
pass
|
|
|
|
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):
|
|
fields=set(data);base={'model_id','max_context','mtp'}
|
|
if fields-base-{'start_context','cache_type_k','cache_type_v'} or not base<=fields or (('cache_type_k' in fields)!=('cache_type_v' in fields)):raise ValueError('Modell, Start- und Endkontext, MTP und beide KV-Cache-Typen auswählen.')
|
|
start_context=data.get('start_context',2048)
|
|
if type(start_context)!=int or start_context not in CONTEXT_STEPS or type(data['max_context'])!=int or data['max_context'] not in CONTEXT_STEPS or start_context>data['max_context'] or type(data['mtp'])!=bool:raise ValueError('Start- und Endkontext müssen gültige Stufen in aufsteigender Reihenfolge sein.')
|
|
if data.get('cache_type_k','q4_0') not in ('f16','q8_0','q4_0') or data.get('cache_type_v','q4_0') not in ('f16','q8_0','q4_0'):raise ValueError('KV-Cache: nur f16, q8_0 oder q4_0 unterstützt.')
|
|
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'],start_context=start_context,max_context=data['max_context'],mtp=data['mtp'],cache_type_k=data.get('cache_type_k','q4_0'),cache_type_v=data.get('cache_type_v','q4_0'),results=[],attempts=0,max_candidates=MAX_CANDIDATES,completed_contexts=[],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.
|
|
# Keep room for 64 output tokens and template overhead, but verify nearly the full context.
|
|
target=context-max(128,context//100);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)<target:raise InferenceError('Synthetischer Kontext konnte nicht gefüllt werden.')
|
|
result=self.request('/completion',dict(prompt=tokens[:target],n_predict=64,temperature=0,seed=1234,cache_prompt=False,stream=False))
|
|
timing=result.get('timings',{})
|
|
if not timing.get('predicted_n') or timing.get('prompt_n',0)<target*.99:raise InferenceError('Keine vollständige Kontext-/Ausgabemessung zurückgegeben.')
|
|
return {k:timing[k] for k in ('prompt_n','prompt_ms','prompt_per_second','predicted_n','predicted_ms','predicted_per_second') if k in timing}
|
|
def test_candidate(self,context,tier,devices,split,ratio,offload,layer_count=None):
|
|
successes=[]
|
|
for micro in (128,64,256):
|
|
if self.cancel.is_set():raise InterruptedError()
|
|
if self.job['attempts']>=self.job.get('max_candidates',MAX_CANDIDATES):raise TestBudget('Kandidatenlimit erreicht')
|
|
if time.time()-self.job['started_at']>MAX_SECONDS:raise TestBudget('Zeitlimit erreicht')
|
|
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'],'cache_type_k':self.job.get('cache_type_k','q4_0'),'cache_type_v':self.job.get('cache_type_v','q4_0')}
|
|
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,primary_gpu_layers=layer_count)
|
|
self.worker.stop()
|
|
stage='prediction'
|
|
try:
|
|
self.worker.plan(p,cancel=self.cancel.is_set)
|
|
stage='model_start'
|
|
self.worker.ensure(p,cancel=self.cancel.is_set)
|
|
# Small warm-up, then populated-context measurement, bounded output.
|
|
stage='measurement'
|
|
self.request('/completion',dict(prompt='Synthetic warmup.',n_predict=16,temperature=0,cache_prompt=False))
|
|
timings=self.benchmark(context)
|
|
row.update(success=True,timings=timings,memory_plan=copy.deepcopy(self.worker.memory_plan),tested_context_fraction=timings.get('prompt_n',0)/context)
|
|
successes.append(row)
|
|
except Exception as exc:
|
|
if self.cancel.is_set():raise InterruptedError()
|
|
row['failure_stage']=stage
|
|
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.
|
|
return successes
|
|
def run(self,primary,secondary):
|
|
try:
|
|
with self.scheduler.lease(('auto-test',self.job['id']),timeout=0,allowed=lambda:not self.cancel.is_set()):
|
|
try:
|
|
contexts=[x for x in CONTEXT_STEPS if self.job.get('start_context',2048)<=x<=self.job['max_context']]
|
|
fine_search=[]
|
|
for context in contexts:
|
|
self.update(phase=f'{context} Kontext · Grundverteilungen',current_context=context)
|
|
primary_success=self.test_candidate(context,'5080',[primary],'none',[],'full')
|
|
dual_success=[]
|
|
for ratio in GPU_SPLITS:
|
|
rows=self.test_candidate(context,'5080 + 3060',[primary,secondary],'layer',ratio,'full')
|
|
if rows:dual_success.append(ratio)
|
|
if not primary_success and not dual_success:
|
|
self.test_candidate(context,'5080 + 3060 + CPU',[primary,secondary],'layer',[85,15],'auto')
|
|
self.update(completed_contexts=self.job.get('completed_contexts',[])+[context])
|
|
if dual_success and not primary_success:fine_search.append((context,dual_success[0]))
|
|
# Run one-layer refinement only after every requested context has its
|
|
# coarse 5080/3060 splits, so tuning cannot starve later context levels.
|
|
for context,ratio in reversed(fine_search):
|
|
try:
|
|
entry=self.worker.catalog.entry(self.job['model_id'])
|
|
total=model_layers(self.worker.catalog.root/entry['id']/('model'+Path(entry['file']).suffix),self.job['mtp'])
|
|
except (OSError,ValueError,TypeError) as exc:
|
|
self.update(fine_search_note='Layer-Feinsuche nicht verfügbar: '+str(exc));continue
|
|
for layers,finer in fine_splits(total,ratio,100):
|
|
result=self.test_candidate(context,'5080 + 3060 · Layer-Feinsuche',[primary,secondary],'layer',finer,'full',layers)
|
|
if not result: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 TestBudget as exc:self.update(state='partial',phase=f'{exc} · Testreihe unvollständig, Teilresultate verfügbar',finished_at=time.time())
|
|
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']<len(self.job['results']):raise ValueError('Testergebnis nicht gefunden.')
|
|
row=copy.deepcopy(self.job['results'][data['index']]);model_id=self.job['model_id']
|
|
if not row['success']:raise ValueError('Nur erfolgreich getestete Kandidaten speichern.')
|
|
return self.profiles.save(dict(id=None,revision=0,name=data['name'],kind='chat',model_id=model_id,parameters=row['parameters']))
|