Spaces:
Running
Running
Xiangyi Li
SkillsBench challenge; Terminal-Bench 2 out of the arena; gates check SkillsBench; practice off the board
77d8d06 Download mock_world.py from benchflow/posttrain-arena: direct link, hf CLI and curl.
- Browser
- Download file 42.5 kB
-
https://huggingface.co/spaces/benchflow/posttrain-arena/resolve/main/mock_world.py
- Command line
-
hf download hf://spaces/benchflow/posttrain-arena/mock_world.py
-
curl -L -o mock_world.py https://huggingface.co/spaces/benchflow/posttrain-arena/resolve/main/mock_world.py
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 ─────────────────────────────────────────────────────── | |
| 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) | |
| 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() | |