Refine dual GPU allocation in individual layer steps
This commit is contained in:
+92
-25
@@ -5,6 +5,57 @@ from profiles import SCHEMAS,CHAT_GPU_DEFAULTS
|
||||
from image_test import probe
|
||||
from inference import InferenceError
|
||||
|
||||
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
|
||||
@@ -55,6 +106,31 @@ class AutoTests:
|
||||
timing=result.get('timings',{})
|
||||
if not timing.get('predicted_n') or timing.get('prompt_n',0)<target*.95: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']>=96 or time.time()-self.job['started_at']>7200:
|
||||
raise TestBudget()
|
||||
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,primary_gpu_layers=layer_count)
|
||||
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.
|
||||
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()):
|
||||
@@ -62,35 +138,26 @@ class AutoTests:
|
||||
contexts=[x for x in (2048,4096,8192,16384,32768,65536,131072,160000,192000,262144) if x<=self.job['max_context']]
|
||||
for context in contexts:
|
||||
candidates=[('5080',[primary],'none',[],'full')]+[('5080 + 3060',[primary,secondary],'layer',s,'full') for s in ([95,5],[90,10],[85,15],[75,25],[50,50])]+[('5080 + 3060 + CPU',[primary,secondary],'layer',[85,15],'auto')]
|
||||
found=False
|
||||
found=False;upper=100
|
||||
for tier,devices,split,ratio,offload in candidates:
|
||||
successes=[]
|
||||
for micro in (128,64,256):
|
||||
if self.cancel.is_set():raise InterruptedError()
|
||||
if self.job['attempts']>=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
|
||||
successes=self.test_candidate(context,tier,devices,split,ratio,offload)
|
||||
if successes:
|
||||
found=True
|
||||
if tier=='5080 + 3060':
|
||||
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));break
|
||||
for layers,finer in fine_splits(total,ratio,upper):
|
||||
result=self.test_candidate(context,'5080 + 3060 · Layer-Feinsuche',devices,split,finer,offload,layers)
|
||||
if not result:break
|
||||
break
|
||||
if tier=='5080 + 3060':upper=ratio[0]
|
||||
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 TestBudget:self.update(state='complete',phase='Testbudget erreicht · 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:
|
||||
|
||||
Reference in New Issue
Block a user