posttrain-arena / mock_world.py
Xiangyi Li
SkillsBench challenge; Terminal-Bench 2 out of the arena; gates check SkillsBench; practice off the board
77d8d06
Raw History Blame Contribute Delete
42.5 kB
"""A mock arena that follows the real protocol, built through the same code as live data.
Nothing here describes a view. It produces the raw records the arena keeps, in their real formats, and then asks the real
code to build every view from them:
- collections: synthetic task packages written to disk and checked by the real static gates
(validation_gates.static_report + compact), stored as registry records shaped like POST /api/environments writes them;
- runs: a competition simulated under the open challenge's rules, read from its config: one active arena run at a
time (compute.concurrent_runs), one run per submission per day that counts (compute.runs_per_submission_per_day), the
project cap (arena_jobs.CAP) with a reservation of flavor price x job timeout per run, the job timeout itself, and the
recipe (gate_task_count, max_steps, num_generations, one held-out trial on the held-out suite). While the open
challenge's own recipe values and compute are not final (SkillsBench: runs paused, recipe values owner-set, Nebius
planned), the simulation runs it under SIMULATED: the arena's earlier, fully specified two-step recipe on HF a100x8, on
the challenge's real held-out suite (simulated_row);
- each run's HF job log uses the pipeline's line formats ([posttrainarena] markers, [PASS]/[FAIL]/[ERR] verdicts,
"Job: N tasks", "Job complete: k/N ...", grpo_rollout_<step>_<index> markers, TRL step dicts, TRAINER_EXIT=); failures
use the platform faults in configs/known_issues.toml;
- each finished run uploads its gate and training attempts in the pipeline's file layout (result.json, timing.json,
verifier/test-stdout.txt, trajectory/acp_trajectory.jsonl, rollout_error.json), which traces.py reads;
- each finished run has a score report in the pipeline's schema, a result written the way POST .../collect writes it,
and an organizer review the way POST .../review records it.
Planned challenges get no runs, because the protocol refuses them. build() patches the storage reads of challenges,
environments and arena_jobs to point at this world, calls challenges.formula(), compute_metrics() and leaderboard(), and
restores them. Run it in its own process (store.py does), since the real modules cache what they read.
With PTA_MOCK_OFFLINE set (the tests set it), build() reads nothing from the network: it skips the sealed suites'
decontamination check, and the stored files the real code reads that this world doesn't simulate come from
fixture/mock-world/<dataset>/<path> instead of the HF datasets (offline_files()).
"""
import hashlib, json, math, os, random, sys, tempfile
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from pathlib import Path
NOW = datetime(2026, 10, 19, 0, 0, tzinfo=timezone.utc) # the moment this simulation shows: two weeks after the challenge opens
# How the simulation runs a challenge whose recipe and compute are not final: changes to its challenge file, applied
# before challenge_row builds the row (None removes a key). The recipe is grpo-v1 (2 optimizer steps on one group of 8
# rollouts, one held-out trial), the arena's last recipe that ran end to end, on the model fragment's HF a100x8 layout.
SIMULATED = {'binding': {'method': 'grpo-v1'}, 'metric': {'trials_per_run': 1},
'compute': {'provider': 'huggingface', 'provider_status': None, 'layout': None, 'eval_trials': None, 'sandbox_concurrency': None},
'recipe': {'note': 'Simulated: the real recipe values are not final, so this world runs grpo-v1 (2 optimizer steps on one group of 8 rollouts).',
'serving_note': 'Simulated: one A100 (device 4) serves the policy; the trainer uses GPUs 0-3.'}}
# What an offline build (PTA_MOCK_OFFLINE) reads for the HF datasets' files, by dataset and path: the submission tracks'
# catalog and the challenge's reference baseline as the Space serves them publicly (/api/v2/challenges, /api/app/meta on
# Sept 29, 2026), and no organizer notices or board messages. A file that isn't here reads as absent from its dataset.
OFFLINE = Path(__file__).with_name('fixture') / 'mock-world'
PRICE_PER_HOUR = 20.0
FAULT_RATE = 0.3 # each run's chance to stop on one of the platform faults below # a100x8 on HF Jobs: $159.9998 per 8 h reservation, as quoted live
iso = lambda d: d.isoformat().replace('+00:00', 'Z')
rng = random.Random(20260926)
TRACES = None # where build() writes the attempts' files (traces.py reads them for source=mock)
TEAMS = ['shellsmiths', 'tty-labs', 'cron-collective', 'kernel-panic', 'pipefitters', 'grep-gang', 'sandbox-sailors', 'posix-pals',
'stack-tracers', 'root-cause', 'devnull-club', 'regex-rangers']
THEMES = [('Log forensics', ['system-administration', 'security']), ('CI repair', ['software-engineering', 'debugging']),
('Data wrangling', ['data-processing', 'file-operations']), ('Build systems', ['software-engineering', 'system-administration']),
('Package hell', ['system-administration', 'software-engineering']), ('Network ops', ['system-administration', 'security']),
('Scientific pipelines', ['scientific-computing', 'data-processing']), ('Puzzle shells', ['games', 'mathematics']),
('Crypto CTF-lite', ['security', 'mathematics']), ('Filesystem surgery', ['file-operations', 'system-administration']),
('Perf tuning', ['optimization', 'debugging']), ('Tool-use drills', ['tool-use', 'file-operations']),
('Database migrations', ['data-processing', 'software-engineering']), ('ML ops', ['model-training', 'machine-learning'])]
VERBS = ['Fix', 'Repair', 'Diagnose', 'Speed up', 'Migrate', 'Recover', 'Harden', 'Refactor', 'Reproduce', 'Automate', 'Audit', 'Untangle']
OBJECTS = ['a systemd unit that restarts in a loop', 'a flaky integration test', 'a CSV export with mixed encodings', 'a slow SQL report',
'a broken apt state', 'a leaking file descriptor', 'a Makefile build', 'a cron job across DST', 'a corrupted tar archive',
'a log rotation that drops lines', 'a config loader', 'a race between two workers', 'a JSONL file with bad rows',
'an nginx reverse proxy', 'a git history with a leaked secret', 'a Python package that will not install', 'a disk that keeps filling',
'a checksum mismatch in a backup', 'a shell pipeline that loses data', 'a TLS certificate chain', 'a memory spike after reload']
PLACES = ['on Debian 12', 'on Alpine', 'in a Python 3.12 repo', 'in a Rust workspace', 'with only busybox', 'on a read-only root', 'in a Go monorepo',
'on Ubuntu 24.04', 'inside a container', 'with no network', 'in a Node 22 project', 'on Fedora 40']
# ── collections: synthetic packages, checked by the real static gates ─────────────────────────────
def write_task(root: Path, name: str, category: str, team: str, defect: str | None):
"""One task package in the layout the validator reads, with at most one planted defect."""
d = root / name; (d / 'environment').mkdir(parents=True); (d / 'verifier').mkdir()
title = f'{rng.choice(VERBS)} {rng.choice(OBJECTS)} {rng.choice(PLACES)}'
credit = '' if defect == 'no-credit' else (f'metadata:\n author_name: {team}\n author_email: {team}@example.org\n license: MIT\n'
f' category: {category}\n origin: original\n')
(d / 'task.md').write_text(f'---\nname: {name}\n{credit}---\n{title}. Leave the result in /app/output and do not change the tests.\n')
solve = '#!/bin/bash\nset -euo pipefail\ncd /app\n' + ''.join(f'python3 tools/step_{i}.py --input data/in_{i}.txt --out output/part_{i}.txt\n' for i in range(4))
if defect == 'stub-oracle': solve = '#!/bin/bash\nexit 0\n'
if defect != 'no-oracle': (d / 'solution').mkdir(); (d / 'solution' / 'solve.sh').write_text(solve)
docker = 'FROM python:3.12-slim\nWORKDIR /app\nCOPY data/ /app/data/\n'
(d / 'environment' / 'data').mkdir(); (d / 'environment' / 'data' / 'in_0.txt').write_text('\n'.join(str(rng.random()) for _ in range(20)))
if defect == 'oracle-in-image': # the reference solution travels into the agent's image
(d / 'environment' / 'data' / 'helper.sh').write_text(solve)
if defect == 'answer-file':
(d / 'environment' / 'data' / 'expected_output.txt').write_text('42\n')
(d / 'environment' / 'Dockerfile').write_text(docker)
checks = [' assert os.path.exists("/app/output/part_0.txt")\n'] if defect == 'existence-only' else \
[' assert Path("/app/output/part_0.txt").read_text().strip() == EXPECTED\n', ' assert len(Path("/app/output/part_1.txt").read_text().splitlines()) == 20\n']
(d / 'verifier' / 'test_outputs.py').write_text('import os\nfrom pathlib import Path\nEXPECTED = "ok"\n\ndef test_output():\n' + ''.join(checks))
test_sh = '#!/bin/bash\necho 1 > /logs/verifier/reward.txt\n' if defect == 'always-reward' else \
'#!/bin/bash\nset -e\nif python3 -m pytest -q /verifier/test_outputs.py; then echo 1 > /logs/verifier/reward.txt; else echo 0 > /logs/verifier/reward.txt; fi\n'
(d / 'verifier' / 'test.sh').write_text(test_sh)
return title
TITLES = {} # task name -> the sentence its task.md asks for (the prompt an attempt's transcript starts with)
DEFECTS = [(None, 0.52), ('no-oracle', 0.14), ('existence-only', 0.07), ('oracle-in-image', 0.06), ('answer-file', 0.05),
('always-reward', 0.04), ('stub-oracle', 0.04), ('no-credit', 0.08)]
def simulated_row(row):
"""The open challenge as this world runs it: its challenge file with SIMULATED applied, built by challenge_row."""
import challenges, tomllib
spec = tomllib.loads((challenges.CHALLENGE_DIR / f"{row['id']}.toml").read_text())
for table, changes in SIMULATED.items():
for key, value in changes.items():
if value is None: spec[table].pop(key, None)
else: spec[table][key] = value
spec['runs_paused'] = spec['status_note'] = None
return challenges.challenge_row(spec)
def collections(row, root: Path, suites):
import environments as env, validation_gates as gates
track, _ = env.resolve_challenge(row['id'])
out, opens = [], datetime.fromisoformat(row['opens']).replace(tzinfo=timezone.utc)
for i in range(16):
team = TEAMS[i % len(TEAMS)]; theme, cats = THEMES[i % len(THEMES)]
n = rng.choice([8, 12, 16, 24, 32, 40, 48, 64])
created = opens + timedelta(hours=rng.uniform(2, (NOW - opens).total_seconds() / 3600 - 2))
base = root / f'c{i:02d}'; packages, titles = {}, {}
for k in range(n):
name = f'{theme.lower().replace(" ", "-")}-{k + 1:03d}'
defect = rng.choices([d for d, _ in DEFECTS], weights=[w for _, w in DEFECTS])[0]
titles[name] = TITLES[name] = write_task(base, name, rng.choice(cats), team, defect)
packages[name] = (base / name, gates.local_manifest(base / name))
report = gates.static_report(packages, suites, require_oracle=gates.REQUIRE_ORACLE, defaults={'license': 'MIT', 'origin': 'original'})
errors, warnings = gates.submission_messages(report) if hasattr(gates, 'submission_messages') else ([], [])
rev = hashlib.sha1(f'{team}{theme}{i}'.encode()).hexdigest()
record = {'agent_id': None, 'challenge_id': track, 'repo_type': 'dataset', 'repo_id': f'{team}/{theme.lower().replace(" ", "-")}', 'revision': rev,
'entry_path': 'tasks', 'title': theme + (' II' if i >= len(THEMES) else ''), 'notes': f'{theme}: {n} terminal tasks written by {team}.',
**({'target_challenge_id': row['id']} if track != row['id'] else {}), 'id': 'env-' + hashlib.sha1(f'mock{i}'.encode()).hexdigest()[:12],
'author': team, 'source_url': None, 'task_count': n, 'status': 'Validated', 'created_at': iso(created),
'warnings': list(warnings)[:30], 'quality_gates': {'static': gates.compact(report)}}
out.append(record)
return out
# ── runs: the competition under the challenge's own rules ───────────────────────────────────────
def fault_catalog():
"""Platform faults the arena has actually seen and explains in configs/known_issues.toml: (stage, the note on the tasks
that errored, or for setup and snapshot the line the job prints)."""
return [('setup', 'VLLM_FAILED'), ('snapshot', 'RuntimeError: bench tasks snapshot-hf failed with exit code 1'),
('baseline', 'LiteLLM proxy failed to start for model…'), ('gate', 'ACP initialize timed out after 120s')]
def tolerated(n):
"""The pipeline's bound on errored tasks per evaluation: ceil(max_infra_error_fraction x tasks), fraction 0.1."""
return int(math.ceil(0.1 * n))
def solve_rate(e, task):
"""The untrained model's chance to solve one task: about 45% of tasks it never solves, the rest somewhere in 15-95%.
Fixed per task, so the gate and the rollouts on a task agree."""
r = random.Random(f'{e["id"]}/{task}'); return 0.0 if r.random() < 0.45 else r.uniform(0.15, 0.95)
def grpo_config(row):
"""The [grpo] table of the challenge's recipe template (rollout_attempts and friends are not in the challenge row)."""
import challenges, tomllib
try: return tomllib.loads((challenges.ROOT / row['recipe']['template']).read_text()).get('grpo', {})
except Exception: return {}
# How a failed attempt that did not time out ends under the pinned pipeline (configs/known_issues.toml): once an attempt
# fills the model's context (16,384 tokens in training, 24,576 in evaluation) the bridge cuts its reply off mid tool call.
ENDINGS = {'training': (0.6, 0.1), 'gate': (0.1, 0.07)} # (cut off, malformed tool call)
COMMANDS = [('ls -la /app', 'total 24\ndrwxr-xr-x 1 root root 4096 .\ndrwxr-xr-x 1 root root 4096 data\ndrwxr-xr-x 1 root root 4096 tools'),
('head -n 5 /app/data/in_0.txt', '0.8444218515250481\n0.7579544029403025\n0.420571580830845\n0.25891675029296335\n0.5112747213686085'),
('python3 tools/step_0.py --input data/in_0.txt --out output/part_0.txt', 'Traceback (most recent call last):\n File "tools/step_0.py", line 12\nFileNotFoundError: output/part_0.txt'),
('mkdir -p /app/output && ls /app/output', ''), ('python3 -m pytest -q /app/tests 2>&1 | tail -3', '1 failed, 1 passed in 0.05s'),
('wc -l /app/output/part_1.txt', '20 /app/output/part_1.txt'), ('cat /app/output/part_0.txt', 'ok')]
THOUGHTS = ['Let me look at what is in /app first.', 'The input has one number per line. I will run the first step.', 'The output directory is missing; create it and rerun.',
'Check the result against what the task asks for.', 'That step failed; read the traceback and fix the path.']
TESTS = ['test_output', 'test_row_count', 'test_no_extra_files', 'test_idempotent', 'test_encoding', 'test_exit_code']
def write_attempts(root: Path, run_id, attempts, start):
"""The files the pipeline uploads for a run's gate and training attempts, in its layout (see traces.py)."""
import shutil
base = root / 'runs' / run_id / 'jobs'; shutil.rmtree(base, ignore_errors=True)
stamp = lambda t: t.strftime('%Y-%m-%d__%H-%M-%S')
for i, a in enumerate(attempts):
r = random.Random(f'{run_id}/{i}'); tid = f'{a["task"]}__{hashlib.sha1(f"{run_id}{i}".encode()).hexdigest()[:8]}'
if a['stage'] == 'gate': d = base / 'grpo-gate' / stamp(start) / tid
else:
top = base / 'grpo-train' / 'step-000000' / 'rank-00' / f'rollout-{a["rollout"]:06d}' / f'attempt-{a["attempt"]:02d}'; d = top / 'jobs' / stamp(a['at']) / tid
top.mkdir(parents=True, exist_ok=True)
if a.get('retried'): (top / 'rollout_error.json').write_text(json.dumps({'attempt': a['attempt'], 'error': a['retried'], 'error_type': 'RuntimeError'}, indent=2))
(d / 'verifier').mkdir(parents=True); (d / 'trajectory').mkdir()
cut, bad = ENDINGS[a['stage']]; roll = r.random()
end = 'errored' if a['error'] else 'timed out' if a['timeout'] else 'finished' if a['ok'] else 'cut off' if roll < cut else 'malformed tool call' if roll < cut + bad else 'finished'
ev = [{'type': 'user_message', 'text': TITLES.get(a['task'], a['task']) + '. Leave the result in /app/output and do not change the tests.'}]
for k in range(a['tools']):
if k % 2 == 0: ev.append({'type': 'agent_message', 'text': r.choice(THOUGHTS) + '\n</think>'})
cmd, out = r.choice(COMMANDS)
ev.append({'type': 'tool_call', 'tool_call_id': f'call_{k:04d}', 'kind': 'execute', 'title': 'bash', 'status': 'completed', 'content': [{'type': 'content', 'content': {'type': 'text', 'text': f'$ {cmd}\n{out}'}}]})
ev.append({'finished': {'type': 'agent_message', 'text': 'The output is in /app/output and the checks I can run pass.'},
'cut off': {'type': 'agent_message', 'text': 'Now write the fixed script.\n</think>\n\n<tool_call>\n<function=write>\n<parameter=content>\n#!/usr/bin/env python3\nimport sys\nfor line in open(sys.argv[1]):\n '},
'malformed tool call': {'type': 'agent_message', 'text': 'Check the output.\n</think>\n\n<tool_call>\n<function=bash>\n<parameter=command>\nls /app/output\n</parameter>\n</function>\n</function>\n</tool_call>'},
'timed out': {'type': 'agent_timeout'}, 'errored': {'type': 'agent_message', 'text': ''}}[end])
(d / 'trajectory' / 'acp_trajectory.jsonl').write_text(''.join(json.dumps(x) + '\n' for x in ev))
n = r.randint(2, 6); passed = n if a['ok'] else r.randint(0, n - 1); fails = TESTS[passed:n]
summary = ', '.join(x for x in [f'{len(fails)} failed' if fails else '', f'{passed} passed' if passed else ''] if x)
if not a['error']:
(d / 'verifier' / 'test-stdout.txt').write_text(''.join(f'FAILED ../root::{t} - AssertionError: /app/output does not match the expected result\n' for t in fails) + f'========== {summary} in 0.04s ==========\n')
(d / 'verifier' / 'reward.txt').write_text('1' if a['ok'] else '0')
agent_s = a['secs']; setup = r.uniform(10, 140)
(d / 'timing.json').write_text(json.dumps({'environment_setup': round(setup, 1), 'agent_setup': round(r.uniform(15, 25), 1), 'agent_execution': round(agent_s, 1), 'verifier': round(r.uniform(2, 12), 1), 'total': round(setup + agent_s + 30, 1)}))
tin = a['tools'] * r.randint(6000, 11000)
(d / 'result.json').write_text(json.dumps({'task_name': a['task'], 'rollout_name': tid, 'rewards': None if a['error'] else {'reward': 1.0 if a['ok'] else 0.0}, 'agent': 'opencode', 'n_tool_calls': a['tools'],
'agent_result': {'n_tool_calls': a['tools'], 'n_input_tokens': tin, 'n_output_tokens': a['tools'] * r.randint(150, 400)},
'error': a['error'] or (f'Agent prompt exceeded wall-clock budget 900s' if a['timeout'] else None)}))
def simulate(row, collections_, sealed):
"""Arena runs in time order under the challenge's rules. Returns ledger records, per-job logs, statuses, costs,
score reports and collected results."""
import arena_jobs as jobs
c, rec = row['compute'], row['recipe']
reserve = round(PRICE_PER_HOUR * c['timeout_seconds'] / 3600, 4) - 0.0002 # flavor price x timeout, as quote() computes it
per_day, cap, prior = c['runs_per_submission_per_day'], jobs.CAP, jobs.PRIOR_ALLOWANCE
eligible = [e for e in collections_ if e['quality_gates']['static']['summary']['eligible']]
requests = sorted((datetime.fromisoformat(e['created_at'].replace('Z', '+00:00')) + timedelta(hours=rng.uniform(0.5, 30)), e) for e in eligible)
requests += sorted((t + timedelta(days=rng.uniform(1.0, 2.5)), e) for t, e in requests if rng.random() < 0.5) # some teams rerun
requests.sort(key=lambda q: q[0])
runs, logs, statuses, reports, results, costs = [], {}, {}, {}, [], {}
free_at, spent = datetime.fromisoformat(row['opens']).replace(tzinfo=timezone.utc), prior
last_run = {}
faults = fault_catalog(); ref = 0.15 # invented: the untrained model's pass rate on the held-out suite (no baseline is measured yet)
for want, e in requests:
start = max(want, free_at, last_run.get(e['id'], want - timedelta(days=2)) + timedelta(days=1 / per_day))
start += timedelta(minutes=rng.uniform(2, 40))
if start > NOW: continue
if spent + reserve > cap: break # the budget check refuses every further run
run_id = 'challenge-' + hashlib.sha1(f'{e["id"]}{start}'.encode()).hexdigest()[:12]
job_id = hashlib.sha1(run_id.encode()).hexdigest()[:24]
fault = rng.choice(faults) if rng.random() < FAULT_RATE else None
events, end, summary, attempts = run_log(row, e, run_id, start, sealed, fault, ref)
running = end is None or end > NOW
if running: end = None; events = [(t, l) for t, l in events if datetime.fromisoformat(t.replace('Z', '+00:00')) <= NOW]
status = 'RUNNING' if running else ('COMPLETED' if summary else 'ERROR')
hours = ((end or NOW) - start).total_seconds() / 3600 + 0.05
record = {'run_id': run_id, 'request_key': f'{e["author"]}-{start:%Y%m%d%H%M}', 'kind': 'challenge-run', 'author': e['author'], 'status': status,
'config': {'challenge_id': row['id'], 'environment_id': e['id'], 'environment_revision': e['revision'], 'agent_id': None, 'recipe_id': rec['id'],
'model': row['base_model']['repo_id'], 'model_revision': row['base_model']['revision'], 'eval_repo': row['eval_suite']['repo_id'],
'eval_revision': row['eval_suite']['revision'], 'train_task_count': e['quality_gates']['static']['summary']['eligible'],
'excluded_task_count': e['quality_gates']['static']['summary']['tasks'] - e['quality_gates']['static']['summary']['eligible'],
'eval_task_count': row['eval_suite']['task_count'], 'flavor': c['flavor'], 'timeout_seconds': c['timeout_seconds']},
'max_compute_usd': reserve, 'created_at': iso(start - timedelta(minutes=2)), 'job_id': job_id, 'job_url': None}
if not running:
record.update(settled_usd=round(PRICE_PER_HOUR * hours, 2), settled_at=iso(end + timedelta(minutes=20)))
spent += record['settled_usd']
else:
spent += reserve
costs[job_id] = record.get('settled_usd')
if not running and TRACES: write_attempts(TRACES, run_id, attempts, start) # the pipeline uploads a run's attempts when its job exits
runs.append(record); logs[job_id] = events; statuses[job_id] = status
if summary and not running:
reports[run_id] = summary
collected = end + timedelta(hours=rng.uniform(0.3, 8))
if collected < NOW:
results.append(collect(row, record, summary, collected))
if not summary and fault is None and running: pass
last_run[e['id']] = start if (summary or running) else last_run.get(e['id'], start - timedelta(days=2)) # failed runs do not count toward the daily limit
free_at = (end or NOW + timedelta(hours=8)) + timedelta(minutes=rng.uniform(5, 45))
for r in results: # organizers review the evidence within a day
reviewed = datetime.fromisoformat(r['collected_at'].replace('Z', '+00:00')) + timedelta(hours=rng.uniform(3, 20))
if reviewed < NOW:
r.update(verification='valid', verification_note='Evidence reviewed: per-task results, the GRPO update and the isolation report check out.',
reviewed_by='organizer', reviewed_at=iso(reviewed))
return runs, logs, statuses, reports, results, costs
def run_log(row, e, run_id, start, sealed, fault, ref):
"""The HF job log of one run in the pipeline's line formats, the time it ended (None: still going) and its score report."""
rec, suite = {**grpo_config(row), **row['recipe']}, row['eval_suite']['task_count']; h = rec['harness']
ev = []; clock = [start]; attempts = []
def at(t, line): ev.append((iso(t), line)); clock[0] = max(clock[0], t)
def say(line, minutes=0.0): at(clock[0] + timedelta(minutes=minutes), line)
def stop(line=None):
if line: say(line, 0.1)
say('TRAINER_EXIT=1', 0.2); return ev, clock[0], None, attempts
def evaluate(outcome, faulty, stage=None):
"""[PASS]/[FAIL]/[ERR] lines in completion order: tasks start on free slots of the harness's concurrency and a task
that hits the per-task limit fails at the limit. A faulty stage has more errored tasks than it tolerates. With a
stage, each attempt is also kept for the files the pipeline uploads (sealed stages are not kept)."""
slots = [clock[0]] * h['concurrency']; done = []
errs = set(rng.sample(list(outcome), min(len(outcome), tolerated(len(outcome)) + rng.randint(1, 2)))) if faulty else set()
for name in outcome:
k = min(range(len(slots)), key=lambda i: slots[i]); begin = slots[k]
timeout = name not in errs and not outcome[name] and rng.random() < 0.35
secs = rng.uniform(40, 120) if name in errs else h['agent_timeout_sec'] + rng.uniform(5, 40) if timeout else rng.uniform(90, h['agent_timeout_sec'] - 30)
slots[k] = begin + timedelta(seconds=rng.uniform(20, 60) + secs); tools = 0 if name in errs else rng.randint(0, 60)
line = f'[ERR] {name} (tools=0) ({faulty})' if name in errs else f'[{"PASS" if outcome[name] else "FAIL"}] {name} (tools={tools})' + (f' (Agent prompt exceeded wall-clock budget {h["agent_timeout_sec"]}s)' if timeout else '')
done.append((slots[k], line))
if stage: attempts.append({'stage': stage, 'task': name, 'ok': outcome[name], 'tools': tools, 'secs': secs, 'timeout': timeout, 'error': faulty if name in errs else None, 'at': slots[k]})
for t, line in sorted(done): at(t, line)
return len(errs)
def complete(outcome, errors, began):
n = len(outcome); p = sum(v for k, v in outcome.items()); say(f'Job complete: {p}/{n} ({100 * p / max(1, n):.1f}%), errors={errors}, idle_timeouts=0, time={(clock[0] - began).total_seconds() / 60:.1f}min', 0.2)
if errors > tolerated(n): return stop(f'RuntimeError: OpenCode evaluation contains agent or verifier errors ({errors} > {tolerated(n)} tolerated)')
kind, note = fault if fault else (None, None)
say(f'pipeline ref {rec["pipeline"]["ref"]} at {rec["pipeline"]["ref"][:7]}')
if kind == 'setup': say(note, rng.uniform(18, 21)); return stop()
say('vllm up', rng.uniform(4.5, 6)); say('bridge up', 0.2); say(f'relay reachable via https://benchflow-posttrain-arena.hf.space/relay/{run_id}/v1', 0.1)
say(f'[posttrainarena] snapshot_train_tasks: bench tasks snapshot-hf benchflow/posttrain-runs --path bundle/submissions/{e["id"]}/{e["revision"]}', 0.1)
if kind == 'snapshot': return stop(note)
say(f'[posttrainarena] snapshot_eval_tasks: bench tasks snapshot-hf {row["eval_suite"]["repo_id"]} --revision {row["eval_suite"]["revision"]}', 0.3)
say('[posttrainarena] validate_task_content_isolation: posttrainarena isolation --train data/train --eval data/eval', 0.3)
# held-out before: the held-out suite, one trial
before = {task: rng.random() < ref + rng.uniform(-0.03, 0.03) for task in sealed}
say(f'[posttrainarena] baseline_eval: bench eval run --tasks-dir data/eval --agent {h["agent"]} --model vllm/{row["base_model"]["repo_id"]} --sandbox {rec["sandbox"]} --concurrency {h["concurrency"]}', 0.4)
say(f'Job: {suite} tasks, 0 done, {suite} to run (concurrency={h["concurrency"]})', 0.1); began = clock[0]
errors = evaluate(before, note if kind == 'baseline' else None)
failed = complete(before, errors, began)
if failed: return failed
passed = sum(before.values())
# base-model gate on up to gate_task_count of the collection's eligible tasks (run policy "always": it never stops the run)
static = e['quality_gates']['static']; names = sorted(n for n in static['task_identity'] if n not in static['excluded'])
gate = names[:rec['gate_task_count']]
gated = {task: rng.random() < solve_rate(e, task) for task in gate}
say(f'[posttrainarena] grpo_gate_eval: bench eval run --tasks-dir data/train --agent {h["agent"]} --model vllm/{row["base_model"]["repo_id"]} --concurrency {h["concurrency"]}', 0.3)
say(f'Job: {len(gate)} tasks, 0 done, {len(gate)} to run (concurrency={h["concurrency"]})', 0.1); began = clock[0]
errors = evaluate(gated, note if kind == 'gate' else None, 'gate')
failed = complete(gated, errors, began)
if failed: return failed
g = sum(gated.values())
# GRPO: TRL draws the step's task from all training tasks with a fixed seed (every run of one revision draws the same
# task) and runs num_generations rollouts on it at once; a rollout that times out has incomplete telemetry and is
# retried, up to rollout_attempts. Both optimizer steps train on that group, and a group whose rewards are all equal
# stops the run (require_reward_variance).
say('[posttrainarena] sync_grpo_endpoint: posttrainarena sync --endpoint student', 1.0)
task = names[random.Random(f'{e["revision"]}/seed-0').randrange(len(names))]; p = solve_rate(e, task); rewards = []; done = []; began = clock[0]
for k in range(rec['num_generations']):
t = began + timedelta(seconds=rng.uniform(1, 20))
for attempt in range(1, rec.get('rollout_attempts', 2) + 1):
ok = rng.random() < p; timeout = not ok and rng.random() < 0.15; tools = rng.randint(3, 60)
secs = h['agent_timeout_sec'] + rng.uniform(5, 40) if timeout else rng.uniform(120, h['agent_timeout_sec'] - 30); end = t + timedelta(seconds=secs + rng.uniform(20, 60))
done += [(t, f'[posttrainarena] grpo_rollout_{0:06d}_{k:06d}: bench eval run --tasks-dir data/train --task-id {task} --model vllm/student'),
(end, f'[{"PASS" if ok else "FAIL"}] {task} (tools={tools})' + (f' (Agent prompt exceeded wall-clock budget {h["agent_timeout_sec"]}s)' if timeout else ''))]
attempts.append({'stage': 'training', 'task': task, 'rollout': k, 'attempt': attempt, 'ok': ok, 'tools': tools, 'secs': secs, 'timeout': timeout, 'error': None, 'at': end,
'retried': 'OpenCode evaluation telemetry coverage is incomplete' if timeout else None})
t = end
if not timeout: rewards.append(1.0 if ok else 0.0); break
for t, line in sorted(done): at(t, line)
if len(rewards) < rec['num_generations']:
return stop(f"RuntimeError: OpenCode GRPO rollout failed for {task} after {rec.get('rollout_attempts', 2)} attempts")
mean = sum(rewards) / len(rewards); std = math.sqrt(sum((r - mean) ** 2 for r in rewards) / len(rewards))
for step in range(1, rec['max_steps'] + 1):
say(str({'loss': round(rng.gauss(0, 0.002), 5) if std else 0.0, 'grad_norm': round(abs(rng.gauss(0.2, 0.05)), 4) if std else 0.0, 'learning_rate': rec['learning_rate'], 'num_tokens': float(rng.randint(60000, 120000)),
'completions/mean_length': round(rng.uniform(6000, 9000), 1), 'rewards/opencode_reward/mean': round(mean, 4), 'reward': round(mean, 4), 'reward_std': round(std, 4),
'frac_reward_zero_std': 0.0 if std else 1.0, 'kl': round(abs(rng.gauss(0.0004 * step, 0.0002)), 6), 'entropy': round(rng.uniform(0.7, 0.9), 4), 'epoch': round(step / rec['max_steps'], 2)}), rng.uniform(2, 6))
if not std: return stop('RuntimeError: GRPO produced zero within-group reward variance; increase runtime.num_generations or improve reward shaping')
# held-out after: the same held-out suite; two optimizer steps move almost nothing
after = {task: (rng.random() < 0.88) if p else (rng.random() < 0.03) for task, p in before.items()}
say(f'[posttrainarena] posttrain_eval: bench eval run --tasks-dir data/eval --agent {h["agent"]} --model vllm/student', 0.5)
say(f'Job: {suite} tasks, 0 done, {suite} to run (concurrency={h["concurrency"]})', 0.1); began = clock[0]
evaluate(after, None); a = sum(after.values()); complete(after, 0, began)
say('[posttrainarena] compare_eval_lift: posttrainarena compare --baseline jobs/baseline --trained jobs/posttrain', 0.2)
say('TRAINER_EXIT=0', 1.0)
b = passed / suite; f = a / suite
summary = {'schema_version': 1, 'run_name': run_id, 'model': row['base_model']['repo_id'], 'model_revision': row['base_model']['revision'], 'final_model': 'grpo_merged',
'train_task_ids': names, 'eval_task_ids': sorted(sealed), 'baseline_score': b, 'sft_score': None, 'grpo_gate_score': g / max(1, len(gate)),
'score_after_posttrain': f, 'delta_score': f - b, 'grpo_planned': True, 'grpo_ran': True, 'grpo_effective_update': bool(std),
'grpo_threshold': 0.0, 'grpo_run_policy': 'always',
'eval_dataset': {'repo_id': row['eval_suite']['repo_id'], 'revision': row['eval_suite']['revision']}, '_before': before, '_after': after}
return ev, clock[0], summary, attempts
def collect(row, record, summary, when):
"""The result POST .../collect writes: pass rates recomputed from per-task outcomes, one trial on the held-out suite."""
before, after = summary['_before'], summary['_after']; n = len(before)
b = sum(before.values()) / n; f = sum(after.values()) / n
se = lambda p: math.sqrt(p * (1 - p) / n)
return {'run_id': record['run_id'], 'challenge_id': row['id'], 'environment_id': record['config']['environment_id'], 'environment_revision': record['config']['environment_revision'],
'author': record['author'], 'agent_id': None, 'baseline_pass_rate': round(b, 6), 'after_pass_rate': round(f, 6), 'delta_pp': round(100 * (f - b), 4),
'stderr_pp': round(100 * math.sqrt(se(b) ** 2 + se(f) ** 2), 4), 'n_tasks': n, 'errors': {'baseline': 0, 'posttrain': 0}, 'grpo_ran': True,
'grpo_effective_update': True, 'trials': 1, 'report_url': None, 'job_url': None, 'artifact_revision': None, 'verification': 'pending',
'verification_note': 'Recomputed from per-task results; awaiting organizer evidence review.', 'collected_by': record['author'], 'collected_at': iso(when)}
# ── the world, read through the real code ───────────────────────────────────────────────────────
@contextmanager
def patched(world):
import arena_jobs as jobs, challenges, environments as env
saved = []
def put(module, name, value):
saved.append((module, name, getattr(module, name))); setattr(module, name, value)
registry, results = world['collections'], world['results']
put(jobs, 'read', lambda head=None: {'prior_allowance_usd': jobs.PRIOR_ALLOWANCE, 'runs': world['ledger']})
put(jobs, 'stage', lambda record: world['statuses'].get(record.get('job_id'), record.get('status')))
put(jobs, 'recorded', lambda: {'jobs': [{'id': j, 'stage': s, 'cost_usd': world['costs'].get(j)} for j, s in world['statuses'].items()],
'spent_usd': jobs.PRIOR_ALLOWANCE + sum(c or 0 for c in world['costs'].values())})
put(challenges, 'job_log', lambda job_id: world['logs'].get(job_id))
put(challenges, 'uploaded_log', lambda run_id: None)
put(challenges, 'report', lambda run_id: {k: v for k, v in world['reports'][run_id].items() if not k.startswith('_')} if run_id in world['reports'] else None)
put(challenges, 'isolation', lambda run_id: 0 if run_id in world['reports'] else None)
put(challenges, 'phase2_rows', lambda: [])
put(challenges, 'relay_stats', lambda run_id: None)
put(challenges, 'collection_rows_cached', lambda: registry)
put(challenges, 'hardware', lambda: [{'name': 'a100x8', 'unitLabel': 'minute', 'unitCostUSD': PRICE_PER_HOUR / 60}])
import compute_jobs
put(compute_jobs, 'compute_jobs', lambda: world['jobs'])
put(env, 'environments', lambda challenge_id=None: registry)
real_read = env.read
put(env, 'read', lambda path=env.PATH, head=None, default=None: results if path == challenges.RESULTS else registry if path == env.PATH else real_read(path, head, default))
put(challenges, 'CHALLENGES', [world['row'] if c['id'] == world['row']['id'] else c for c in challenges.CHALLENGES]) # the challenge as this world runs it
for row in challenges.CHALLENGES: # the organizers' status line describes the real arena, not this one
saved.append((row, 'status_note', row.get('status_note'))); row['status_note'] = None
saved.append((row, 'runs_paused', row.get('runs_paused'))); row['runs_paused'] = None # nor does a pause of the real arena
for cache in ('_health', '_finished', '_scored', '_metrics'):
if isinstance(getattr(challenges, cache, None), dict): getattr(challenges, cache).clear()
try: yield
finally:
for target, name, value in reversed(saved):
if isinstance(target, dict): target[name] = value
else: setattr(target, name, value)
@contextmanager
def offline_files():
"""With PTA_MOCK_OFFLINE set, the HF datasets' files the real code downloads while it builds this world (the tracks'
catalog, notices, board messages, the reference baseline) come from OFFLINE; without it, nothing changes."""
if not os.environ.get('PTA_MOCK_OFFLINE'):
yield
return
import challenges, environments as env
from huggingface_hub.errors import EntryNotFoundError
def download(repo_id, filename, **kwargs):
path = OFFLINE / repo_id / filename
if not path.is_file():
raise EntryNotFoundError(f'{repo_id}: {filename} is not among the offline mock world\'s files')
return str(path)
saved = [(module, module.hf_hub_download) for module in (env, challenges)]
for module, _ in saved: module.hf_hub_download = download
try: yield
finally:
for module, real in saved: module.hf_hub_download = real
def build():
"""The mock payload in the shape store.load reads: {'formula', 'jobs', 'metrics', 'boards'}, all built by the real code."""
with offline_files():
return _build()
def _build():
import arena_jobs as jobs, challenges, validation_gates as gates
import store, shutil
global TRACES
rng.seed(20260926); TITLES.clear(); TRACES = store.DIR / 'mock-traces'; shutil.rmtree(TRACES, ignore_errors=True)
row = simulated_row(next(c for c in challenges.CHALLENGES if c['status'] == 'open'))
sealed = challenges.suite_task_ids(row)
try: suites = [] if os.environ.get('PTA_MOCK_OFFLINE') else gates.heldout_suites() # decontamination needs the sealed suites; offline it runs without them
except Exception: suites = []
with tempfile.TemporaryDirectory() as tmp:
cols = collections(row, Path(tmp), suites)
ledger, logs, statuses, reports, results, costs = simulate(row, cols, sealed)
job_rows = [{'id': r['job_id'], 'name': r['run_id'], 'kind': 'challenge run', 'purpose': f'{r["author"]}: a run of {r["config"]["environment_id"]}', 'run_id': r['run_id'],
'challenge': r['config']['challenge_id'], 'flavor': r['config']['flavor'], 'provider': 'huggingface', 'stage': statuses[r['job_id']], 'created_at': r['created_at'],
'started_at': r['created_at'], 'finished_at': r.get('settled_at'), 'seconds': None, 'cost_usd': r.get('settled_usd'), 'url': None} for r in ledger]
world = {'row': row, 'collections': cols, 'ledger': ledger, 'logs': logs, 'statuses': statuses, 'reports': reports, 'results': results, 'costs': costs, 'jobs': {'jobs': job_rows}}
import store
with patched(world):
payload = store.live_payload() # the same assembly as live data, reading this world
payload['jobs']['budget'] = jobs.budget(jobs.read(), jobs.recorded())
payload['formula']['expected'] = {'as_of': iso(NOW), 'basis': BASIS, 'band_attempts': None}
return payload
BASIS = [
'The real challenge’s runs are paused and its recipe values are not final, so this world runs it under the arena’s earlier two-step recipe on its real held-out suite: one arena run at a time, one counted run per submission per day, the $800 project cap with a $160 reservation per run (a100x8 price × the 8 h job timeout), 2 GRPO steps on one group of 8 rollouts, and one trial on the 87 SkillsBench tasks, so every Δ is a multiple of 1/87 (about 1.15 pp).',
'Collections are synthetic task packages checked by the real static gates, the same code that checks a real submission; about half carry one planted defect (no reference solution, a leaked solution, an existence-only verifier, and so on).',
'Run logs use the pipeline’s exact line formats and are read by the real log parser: evaluations run 8 tasks at a time with the 900 s per-task limit, and each run has a 30% chance to stop on a platform fault the arena has actually hit (a failed evaluation has more errored tasks than the pipeline tolerates, ceil(10% of tasks)). Results are recomputed and reviewed the way collect and review do it.',
'Training follows the recipe: one task drawn from the collection’s training tasks with a fixed seed, 8 attempts on it at once, a timed-out attempt retried once, and the run stops when all 8 score the same, because GRPO has nothing to learn from them. The untrained model never solves about 45% of tasks, so many runs stop there, as the arena’s own TMax run did.',
'Each finished run uploads its gate and training attempts in the pipeline’s file layout (transcript, verifier output, timings); failed attempts end the way they do under the pinned pipeline, including replies cut off once an attempt fills the model’s context.',
'Teams, collections and outcomes are invented, and so is the untrained model’s pass rate (about 15%): no SkillsBench baseline is measured yet.',
]
def main():
payload = build()
json.dump(payload, sys.stdout, default=str)
if __name__ == '__main__':
main()