ProCreations's picture
Pin matched Agnes native validation suite
63f6cfd verified
Raw History Blame
18.9 kB
"""Paired BF16/NVFP4 generation tests in a disposable, credential-free HF Job.
The private quantized repository is mounted read-only by HF Jobs. This program
does not receive a Hub token. Results are returned through authenticated job logs.
These are diagnostic subsets, not official full-benchmark leaderboard scores.
"""
import base64
import concurrent.futures
from decimal import Decimal
import hashlib
import io
import importlib.metadata
import itertools
import json
import os
import random
import re
import subprocess
import sys
import time
from pathlib import Path
import requests
from datasets import load_dataset
from huggingface_hub import HfApi, snapshot_download, hf_hub_download
from stage_checkpoint import stage_checkpoint
ROOT = Path(os.environ.get('AGNES_EVAL_WORKSPACE','/workspace/native-eval'))
ROOT.mkdir(parents=True, exist_ok=True)
SOURCE = 'Agnes-AI/Agnes-3.0-Flash'
REV = '24f712ce59379b54c4a141d2708c35daf5ff613b'
SEED = 20260908
PORTS = {'bf16': 31000, 'nvfp4': 31001}
MODE = os.environ.get('AGNES_EVAL_MODE','paired')
assert MODE in ['paired','bf16','nvfp4']
if MODE != 'paired':
PORTS = {MODE:PORTS[MODE]}
DATA_REVISIONS = {}
CHECKPOINT_STAGING = None
PLAN = json.loads((Path(__file__).resolve().parent/'native_plan.json').read_text())
RECHECK = (json.loads((Path(__file__).resolve().parent/'recheck_cases.json').read_text())
if os.environ.get('AGNES_RECHECK')=='1' else None)
def data(repo, split, config=None):
rev = PLAN['dataset_revisions'][repo]
DATA_REVISIONS[repo] = rev
return load_dataset(repo, name=config, split=split, revision=rev, token=False)
def norm(s):
return re.sub(r'\s+', ' ', s.strip().lower()).strip(' .,:;!"\'')
def request(model, task):
body = dict(model=model, messages=task['messages'], temperature=0,
max_tokens=task.get('max_tokens', 8192), chat_template_kwargs=dict(reasoning_effort='xhigh', enable_thinking=task.get('effort') != 'none'),
seed=SEED, top_p=1.0)
if 'tools' in task:
body['tools'] = task['tools']
body['tool_choice'] = 'auto'
start = time.monotonic()
response = requests.post(f'http://127.0.0.1:{PORTS[model]}/v1/chat/completions', json=body, timeout=900)
response.raise_for_status()
result = response.json()
choice = result['choices'][0]
message = choice['message']
content = message.get('content') or ''
# Preserve full outputs locally in the disposable runner; publish compact evidence.
(ROOT / f"{model}-{task['id'].replace('/', '_')}.json").write_text(json.dumps(result))
score = score_task(task, message)
row = dict(model=model, task=task['kind'], id=task['id'], score=score,
finish_reason=choice.get('finish_reason'), seconds=time.monotonic()-start,
usage=result.get('usage'), output_sha256=hashlib.sha256(content.encode()).hexdigest())
print('RESULT_ITEM ' + json.dumps(row), flush=True)
if score == 0:
print('RESULT_FAILURE ' + json.dumps(dict(model=model,id=task['id'],content=content[-4000:],
tool_calls=message.get('tool_calls'))),flush=True)
return row
def score_task(task, message):
answer = message.get('content') or ''
if task['kind'] == 'gsm8k_96':
values = re.findall(r'[-+]?\d[\d,]*(?:\.\d+)?', answer)
return float(bool(values) and Decimal(values[-1].replace(',', '')) == Decimal(task['answer'].replace(',', '')))
if task['kind'] in ['mmlu_pro_84','longbench_v2_24']:
boxed = re.findall(r'\\boxed\{([A-J])\}', answer)
letters = re.findall(r'(?<![A-Za-z])([A-J])(?![A-Za-z])', answer)
prediction = (boxed or letters or [''])[-1]
return float(prediction == task['answer'])
if task['kind'] == 'document_vqa_64':
return float(any(norm(answer) == norm(x) for x in task['answers']))
if task['kind'] in ['structured_multilingual_32', 'long_context_12']:
try:
value = json.loads(re.sub(r'^```(?:json)?\s*|\s*```$', '', answer.strip()))
return float(value == task['answer'])
except (ValueError, TypeError):
return 0.0
if task['kind'] == 'tool_calling_32':
calls = message.get('tool_calls') or []
if len(calls) != 1:
return 0.0
function = calls[0]['function']
try:
args = json.loads(function['arguments']) if isinstance(function['arguments'], str) else function['arguments']
except (ValueError, TypeError):
return 0.0
return float(function['name'] == 'lookup_inventory' and args == task['answer'])
if task['kind'] == 'humaneval_plus_32':
from evalplus_compat import score
return score(task['problem'], None, answer)
raise ValueError(task['kind'])
def tasks():
rng = random.Random(SEED)
output = []
gsm = data('openai/gsm8k', 'test', 'main')
for i in rng.sample(range(len(gsm)), 96):
row = gsm[i]
output.append(dict(id=f'gsm8k-{i}', kind='gsm8k_96', answer=row['answer'].split('####')[-1].strip(),
messages=[dict(role='user', content=row['question']+'\nEnd with only the numerical answer on the final line.')]))
mmlu = data('TIGER-Lab/MMLU-Pro', 'test')
categories = sorted(set(mmlu['category']))
for category in categories:
ids = [i for i, c in enumerate(mmlu['category']) if c == category]
for i in rng.sample(ids, 6):
row = mmlu[i]
prompt = row['question'] + '\n' + '\n'.join(f'{chr(65+j)}. {o}' for j, o in enumerate(row['options']))
prompt += '\nPut the final answer letter inside \\boxed{}.'
output.append(dict(id=f'mmlu-{i}', kind='mmlu_pro_84', answer=row['answer'],
messages=[dict(role='user', content=prompt)]))
repo = 'HuggingFaceM4/DocumentVQA'
revision = PLAN['dataset_revisions'][repo]
DATA_REVISIONS[repo] = revision
images = load_dataset(repo, split='validation', revision=revision, streaming=True, token=False)
for i, row in enumerate(itertools.islice(images, 256, 320), start=256):
image = row['image'].convert('RGB')
image.thumbnail((1600, 1600))
buffer = io.BytesIO()
image.save(buffer, format='PNG')
url = 'data:image/png;base64,' + base64.b64encode(buffer.getvalue()).decode()
output.append(dict(id=f'docvqa-{i}', kind='document_vqa_64', answers=row['answers'], effort='none', max_tokens=128,
messages=[dict(role='user', content=[dict(type='image_url', image_url={'url':url}),
dict(type='text', text=row['question']+'\nReturn only the short answer from the document.')])]))
languages = [
'Return only JSON with keys "sum" (the sum) and "product" (the product) for the two numbers',
'Gib nur JSON mit den Schlüsseln "sum" (Summe) und "product" (Produkt) für diese Zahlen zurück',
'Devuelve solo JSON con las claves "sum" (suma) y "product" (producto) para estos números',
'次の二つの数について、和を "sum"、積を "product" とするJSONのみを返してください',
'请只返回JSON,键 "sum" 表示两数的和,键 "product" 表示两数的乘积,数字是',
'Retourne uniquement du JSON avec "sum" pour la somme et "product" pour le produit de ces nombres',
'أعد JSON فقط بالمفتاح "sum" لمجموع العددين و"product" لحاصل ضربهما',
'Верни только JSON с ключами "sum" (сумма) и "product" (произведение) для чисел',
]
for i in range(32):
a, b = rng.randint(12, 99), rng.randint(12, 99)
output.append(dict(id=f'multilingual-{i}', kind='structured_multilingual_32', effort='none', max_tokens=128,
answer=dict(sum=a+b, product=a*b), messages=[dict(role='user', content=f'{languages[i%8]}: {a}, {b}.')]))
sku = f'ITEM-{rng.randint(1000,9999)}'
warehouse = ['north','south','east','west'][i%4]
schema = dict(type='function', function=dict(name='lookup_inventory', description='Look up stock for an item in a warehouse.',
parameters=dict(type='object', properties=dict(sku=dict(type='string'), warehouse=dict(type='string', enum=['north','south','east','west'])), required=['sku','warehouse'])))
output.append(dict(id=f'tool-{i}', kind='tool_calling_32', effort='none', max_tokens=256,
answer=dict(sku=sku, warehouse=warehouse), tools=[schema],
messages=[dict(role='user', content=f'Use the inventory tool to check {sku} in the {warehouse} warehouse.')]))
for i in range(12):
lines = [f'Record {j:05d}: {rng.randrange(100000,999999)}' for j in range([1500,3000,6000,9000,14000,16000][i//2])]
location = [50, len(lines)//2, len(lines)-50][i%3]
key = rng.randrange(10000000, 99999999)
lines[location] = f'The unique access code for project Project-{i} is {key}.'
prompt = '\n'.join(lines) + f'\nReturn only JSON {{"code": number}} with the access code for Project-{i}.'
output.append(dict(id=f'long-{i}', kind='long_context_12', effort='none', max_tokens=128,
answer=dict(code=key), messages=[dict(role='user', content=prompt)]))
from select_longbench import prompt_for
selection=json.loads((Path(__file__).resolve().parent/'longbench_selection.json').read_text())
longbench=json.loads(Path(hf_hub_download(selection['dataset'],'data.json',repo_type='dataset',
revision=selection['revision'],token=False)).read_text())
longbench={row['_id']:row for row in longbench}
DATA_REVISIONS[selection['dataset']]=selection['revision']
for selected in selection['records']:
row=longbench[selected['id']]
prompt=prompt_for(row)
assert hashlib.sha256(prompt.encode()).hexdigest()==selected['prompt_sha256']
output.append(dict(id='longbench-'+row['_id'],kind='longbench_v2_24',answer=row['answer'],
messages=[dict(role='user',content=prompt)]))
from evalplus.data import get_human_eval_plus, get_human_eval_plus_hash
problems = get_human_eval_plus()
selected = {k: problems[k] for k in rng.sample(sorted(problems), 32)}
DATA_REVISIONS['evalplus/humaneval+'] = get_human_eval_plus_hash()
for key, problem in selected.items():
output.append(dict(id=key, kind='humaneval_plus_32', problem=problem,
messages=[dict(role='user', content='Implement the following Python function. Return the complete function and needed imports in one Python code block.\n'+problem['prompt'])]))
return output
def start_server(model, path, gpu):
env = dict(os.environ, CUDA_VISIBLE_DEVICES=str(gpu))
assert not any(env.get(k) for k in ['HF_TOKEN','HUGGING_FACE_HUB_TOKEN']), 'Native runner must not receive Hub credentials'
log_path = ROOT / (model+'-server.log')
log = log_path.open('w')
command = ['bash',str(Path(path)/'serve.sh'),
'--served-model-name',model,'--host','127.0.0.1','--port',str(PORTS[model]),
'--tp-size','1','--context-length','262144','--max-running-requests','4',
'--chunked-prefill-size','4096',
'--mem-fraction-static','0.88','--cuda-graph-max-bs','4',
'--reasoning-parser','qwen3','--tool-call-parser','qwen3_coder',
'--chat-template',str(Path(path)/'chat_template.jinja'),
'--mamba-scheduler-strategy','extra_buffer']
if model=='nvfp4':
# RTX PRO 6000 is SM120. Select its native FP4 backend explicitly;
# the ModelOpt MoE auto fallback can select the SM100 TRTLLM path.
command += ['--fp4-gemm-backend','flashinfer_cutlass']
proc = subprocess.Popen(command, env=env, stdout=log, stderr=subprocess.STDOUT)
return proc, log_path, command
def prepare_quant_model():
global CHECKPOINT_STAGING
destination=ROOT/'quant-model'
CHECKPOINT_STAGING=stage_checkpoint('/models/quant',destination,
Path(__file__).resolve().parent/'native_quant_files.json')
return destination
def get_checkpoint_staging_report():
return CHECKPOINT_STAGING
def report_startup_progress(model, log, last_report):
now=time.monotonic()
if now-last_report.get(model,0)>=60:
last_report[model]=now
tail=log.read_text(errors='replace')[-2000:]
if tail:
print('SERVER_STARTUP_PROGRESS',model,tail.replace('\r','\n'),flush=True)
def main():
assert not os.environ.get('HF_TOKEN')
paths = {'nvfp4':prepare_quant_model()} if 'nvfp4' in PORTS else {}
if 'bf16' in PORTS:
paths['bf16']=snapshot_download(SOURCE, revision=REV, token=False, max_workers=8)
servers = {}
try:
for gpu, model in enumerate(PORTS):
servers[model] = start_server(model, paths[model], gpu)
deadline = time.monotonic()+1200
ready = set()
last_startup_report = {}
while len(ready) != len(PORTS):
for model, (proc, log, command) in servers.items():
if proc.poll() is not None:
raise RuntimeError(f'{model} server failed: '+log.read_text()[-12000:])
if model not in ready:
report_startup_progress(model,log,last_startup_report)
try:
response = requests.get(f'http://127.0.0.1:{PORTS[model]}/health',timeout=3)
if response.ok:
ready.add(model)
except requests.RequestException:
pass
if time.monotonic()>deadline:
raise RuntimeError('Server startup timeout: '+ '\n'.join(s[1].read_text()[-6000:] for s in servers.values()))
time.sleep(5)
print('NATIVE_SERVERS_READY',list(PORTS),flush=True)
suite = tasks()
if RECHECK:
suite=[task for task in suite if task['id'] in RECHECK['ids']]
for task in suite:
task['max_tokens']=(RECHECK['protocol']['short_format_output_cap']
if task.get('effort')=='none' else RECHECK['protocol']['broad_reasoning_output_cap'])
from evalplus_compat import preflight, PROTOCOL
preflight(suite)
print('NATIVE_METADATA '+json.dumps(dict(source_revision=REV,mode=MODE,
evaluation_code_revision=os.environ.get('AGNES_CODE_REVISION'),
versions={p:importlib.metadata.version(p) for p in ['torch','transformers','sglang','evalplus','datasets','huggingface_hub']},
dataset_revisions=DATA_REVISIONS,code_scoring_protocol=PROTOCOL,
server_commands={k:v[2] for k,v in servers.items()})),flush=True)
results=[]
# Four simultaneous requests per server; both models see identical tasks.
with concurrent.futures.ThreadPoolExecutor(max_workers=4) as pool:
pending = [pool.submit(request, model, task) for task in suite for model in PORTS]
for future in concurrent.futures.as_completed(pending):
results.append(future.result())
from agent_eval import run_suite, GORILLA_REV
results.extend(run_suite(Path('/workspace/native-eval/gorilla/berkeley-function-call-leaderboard'), PORTS, ROOT,
selected_ids=RECHECK['ids'] if RECHECK else None,
per_step_output_cap=RECHECK['protocol']['bfcl_per_step_output_cap'] if RECHECK else 4096))
DATA_REVISIONS['BFCL_v4_source_and_data'] = GORILLA_REV
from long_agent_eval import scenarios as long_scenarios, oracle_check, run_case
long_suite=long_scenarios(); oracle_check(long_suite)
for case in long_suite:
for model in PORTS:
results.append(run_case(model,case))
summary={}
for kind in sorted({r['task'] for r in results}):
summary[kind]={}
for model in PORTS:
rows=[r for r in results if r['task']==kind and r['model']==model]
summary[kind][model]=dict(n=len(rows), score=sum(r['score'] for r in rows)/len(rows),
length_terminated=sum(r['finish_reason']=='length' for r in rows))
if rows and 'state_response_score' in rows[0]:
for component in ['state_response_score','irrelevance_score']:
summary[kind][model][component]=sum(r[component] for r in rows)/len(rows)
report=dict(seed=SEED, source_revision=REV, quant_revision=os.environ.get('AGNES_QUANT_REVISION') if MODE!='bf16' else None,
checkpoint_staging=CHECKPOINT_STAGING,
evaluation_code_revision=os.environ.get('AGNES_CODE_REVISION'),mode=MODE,
code_scoring_protocol=PROTOCOL,
output_budget_recheck=RECHECK, examples=results,
versions={p:importlib.metadata.version(p) for p in ['torch','transformers','sglang','evalplus','datasets','huggingface_hub']},
dataset_revisions=DATA_REVISIONS, summary=summary,
server_commands={k:v[2] for k,v in servers.items()},
caveats=['Diagnostic subsets; not official complete benchmark scores.',
('Expanded output-budget recheck; caps and deterministic case-selection recorded in output_budget_recheck.' if RECHECK else
'Temperature 0, source xhigh thinking; max 8192 tokens for broad reasoning tasks and long workflows; 4096 per BFCL step. Thinking disabled for short exact-format tasks. Truncations reported.'),
'BFCL v4 subsets use upstream simulator and state/response/irrelevance checks; max 21 generation steps per user turn.',
'LongBench v2 uses 24 untruncated cases, four per domain; see longbench_selection.json for fixed IDs, lengths and filtering.',
'Document images resized to fit 1600x1600; exact-match scoring.',
'Twelve synthetic multi-step tool workflows include approximately 32K, 98K and 224K-token archives; not an official benchmark.', 'No HF credentials were passed to this evaluator.'])
print('NATIVE_REPORT '+json.dumps(report),flush=True)
finally:
for proc, log, command in servers.values():
proc.terminate()
for proc, log, command in servers.values():
try:
proc.wait(timeout=30)
except subprocess.TimeoutExpired:
proc.kill()
if __name__ == '__main__':
main()