beam-pi-programbench / source /study /tests /test_analysis.py
burtenshaw's picture
burtenshaw HF Staff
feat: publish beam pi study source
5741b22 verified
Raw History Blame Contribute Delete
11.4 kB
import importlib.util
import json
from pathlib import Path
import tempfile
import unittest
spec = importlib.util.spec_from_file_location("study_analysis", Path(__file__).parents[1] / "analyze.py")
analysis = importlib.util.module_from_spec(spec)
spec.loader.exec_module(analysis)
class CurveTests(unittest.TestCase):
def setUp(self):
self.temp = tempfile.TemporaryDirectory()
self.root = Path(self.temp.name)
self.config = {"tasks": [{"instance_id": f"task-{i}"} for i in range(5)],
"runs": [{"run_id": f"task-{i}-{condition}", "task": f"task-{i}",
"condition": condition, "category": "benchmark"}
for condition in analysis.CONDITIONS for i in range(5)]}
def tearDown(self):
self.temp.cleanup()
def write_episode(self, task, condition, grades, stop_reason="completed"):
directory = self.root / "episodes" / f"task-{task}-{condition}"
(directory / "runner").mkdir(parents=True)
(directory / "grades.json").write_text(json.dumps(grades))
(directory / "runner" / "result.json").write_text(json.dumps({"stop_reason": stop_reason, "elapsed_seconds": 120}))
(directory / "runner" / "events.jsonl").write_text(json.dumps({"type": "episode_start", "sequence": 1,
"utc": "2026-10-09T00:00:00Z", "elapsed_ms": 0}) + "\n")
def row(self, rows, instant):
return next(row for row in rows if row["axis"] == "wall_seconds" and row["condition"] == "single" and row["time_seconds"] == instant)
def test_regressions_are_preserved_and_terminal_scores_carried(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "score": 0.6, "valid": True},
{"elapsed_seconds": 60, "score": 0.2, "valid": True}])
summary, rows = analysis.analyze(self.root, self.config)
self.assertAlmostEqual(self.row(rows, 30)["mean_hidden_test_fraction"], 0.12)
self.assertAlmostEqual(self.row(rows, 60)["mean_hidden_test_fraction"], 0.04)
self.assertAlmostEqual(self.row(rows, 7200)["mean_hidden_test_fraction"], 0.04)
self.assertAlmostEqual(summary["final"]["single"]["mean_hidden_test_fraction"], 0.04)
self.assertEqual(summary["final"]["single"]["tasks_total"], 5)
def test_tasks_are_equal_weighted_not_pooled_by_test_count(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "score": 1.0, "valid": True,
"passed": 1, "test_count": 1}])
self.write_episode(1, "single", [{"elapsed_seconds": 30, "score": 0.5, "valid": True,
"passed": 500, "test_count": 1000}])
summary, rows = analysis.analyze(self.root, self.config)
self.assertAlmostEqual(summary["final"]["single"]["mean_hidden_test_fraction"], 0.3)
self.assertEqual(self.row(rows, 30)["tasks_with_valid_grade"], 2)
def test_invalid_grading_retains_last_valid_or_zero(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "score": 0.5, "valid": True},
{"elapsed_seconds": 60, "score": 0.9, "valid": False, "error": "grader_unavailable"}], "agent_error")
self.write_episode(1, "single", [{"elapsed_seconds": 30, "score": None, "valid": False}], "setup_error")
summary, _ = analysis.analyze(self.root, self.config)
self.assertAlmostEqual(summary["final"]["single"]["mean_hidden_test_fraction"], 0.1)
episode = next(row for row in summary["episodes"] if row["run_id"] == "task-0-single")
self.assertEqual(episode["status"], "failed")
self.assertEqual(episode["final_score"], 0.5)
self.assertEqual(episode["invalid_grade_count"], 1)
def test_explicit_valid_zero_replaces_previous_success(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "score": 0.8, "valid": True},
{"elapsed_seconds": 60, "score": 0.0, "valid": True}])
summary, rows = analysis.analyze(self.root, self.config)
self.assertEqual(summary["final"]["single"]["mean_hidden_test_fraction"], 0)
self.assertEqual(self.row(rows, 60)["tasks_with_valid_grade"], 1)
def test_grade_summary_fallback_uses_official_aggregate_only(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "sha256": "abc", "valid": True}])
directory = self.root / "episodes" / "task-0-single" / "evaluations" / "abc" / "task-0"
directory.mkdir(parents=True)
(directory / "summary.json").write_text(json.dumps({"complete": True, "score": 0.75, "test_count": 100}))
summary, _ = analysis.analyze(self.root, self.config)
self.assertEqual(summary["final"]["single"]["mean_hidden_test_fraction"], 0.15)
self.assertTrue((self.root / "curves.csv").exists())
self.assertTrue((self.root / "summary.json").exists())
def test_summary_fallback_respects_submission_and_infrastructure_classification(self):
fixtures = [{"complete": False, "valid": True, "analysis_score": 0.0,
"score": 0.0, "scoring_status": "submission_failed"},
{"complete": False, "valid": False, "analysis_score": None,
"score": 0.9, "scoring_status": "infrastructure_failed"}]
for task, summary in enumerate(fixtures):
self.write_episode(task, "single", [{"elapsed_seconds": 30, "sha256": "abc", "valid": True}])
directory = self.root / "episodes" / f"task-{task}-single" / "evaluations" / "abc" / f"task-{task}"
directory.mkdir(parents=True)
(directory / "summary.json").write_text(json.dumps(summary))
summary, _ = analysis.analyze(self.root, self.config)
episodes = {row["run_id"]: row for row in summary["episodes"]}
self.assertTrue(episodes["task-0-single"]["has_valid_grade"])
self.assertEqual(episodes["task-0-single"]["final_score"], 0.0)
self.assertFalse(episodes["task-1-single"]["has_valid_grade"])
self.assertEqual(episodes["task-1-single"]["invalid_grade_count"], 1)
def test_controller_grading_failure_overrides_completed_runner(self):
self.write_episode(0, "single", [{"elapsed_seconds": 30, "score": None, "valid": False}])
directory = self.root / "episodes" / "task-0-single"
(directory / "status.json").write_text(json.dumps({"status": "failed", "error": "grader failed"}))
summary, _ = analysis.analyze(self.root, self.config)
episode = next(row for row in summary["episodes"] if row["run_id"] == "task-0-single")
self.assertEqual(episode["status"], "failed")
self.assertEqual(episode["invalid_grade_count"], 1)
self.assertEqual(summary["final"]["single"]["episodes_completed"], 0)
def test_normal_wall_time_cutoff_is_a_completed_attempt(self):
self.write_episode(0, "single", [{"elapsed_seconds": 120, "score": 0.2, "valid": True}], "wall_time_limit")
summary, _ = analysis.analyze(self.root, self.config)
episode = next(row for row in summary["episodes"] if row["run_id"] == "task-0-single")
self.assertEqual(episode["status"], "completed")
self.assertEqual(episode["stop_reason"], "wall_time_limit")
directory = self.root / "episodes" / "task-0-single"
(directory / "status.json").write_text(json.dumps({"status": "interrupted", "stop_reason": "controller_restart"}))
summary, _ = analysis.analyze(self.root, self.config)
episode = next(row for row in summary["episodes"] if row["run_id"] == "task-0-single")
self.assertEqual(episode["status"], "interrupted")
self.assertEqual(episode["stop_reason"], "controller_restart")
class ClockTests(unittest.TestCase):
def event(self, type_, seconds, sequence, **kwargs):
return {"type": type_, "elapsed_ms": seconds * 1000, "sequence": sequence, **kwargs}
def request(self, ident, agent, prompt, context, output, started=0, ended=1):
return {"id": ident, "agent_id": agent, "client_request_id": ident, "status": "complete",
"prompt_tokens": prompt, "completion_tokens": output, "context_charged": context,
"started_at": f"2026-10-09T00:00:{started:02d}Z", "ended_at": f"2026-10-09T00:00:{ended:02d}Z"}
def test_parallel_clocks_message_handoff_and_child_inheritance(self):
events = [self.event("episode_start", 0, 1, utc="2026-10-09T00:00:00Z"),
self.event("agent_created", 0, 2, agent_id="a", parent_id=None),
self.event("agent_created", 0, 3, agent_id="b", parent_id=None),
self.event("tool_start", 1.1, 4, agent_id="a"),
self.event("tool_end", 4.1, 5, agent_id="a", causal_event_id=4, duration_ms=3000),
self.event("message_sent", 4.2, 6, agent_id="a"),
self.event("message_delivered", 4.3, 7, agent_id="b", causal_event_id=6),
self.event("tool_start", 4.4, 8, agent_id="b"),
self.event("tool_end", 6.4, 9, agent_id="b", causal_event_id=8, duration_ms=2000),
self.event("agent_created", 6.5, 10, agent_id="child", parent_id="b"),
self.event("tool_start", 6.6, 11, agent_id="child"),
self.event("tool_end", 7.6, 12, agent_id="child", causal_event_id=11, duration_ms=1000),
self.event("snapshot", 7.7, 13, path="snapshot.tar.gz")]
requests = [self.request("req-a", "a", 20_000, 10_265, 265),
self.request("req-b", "b", 0, 265, 265)]
clocks = analysis.clock_trace(events, requests)
self.assertTrue(clocks["valid"])
self.assertEqual(analysis.clocks_at(clocks, 1), (2, 3)) # max, not sum
self.assertEqual(analysis.clocks_at(clocks, 100, "snapshot.tar.gz"), (8, 9))
def test_gateway_only_compaction_is_counted(self):
events = [self.event("episode_start", 0, 1, utc="2026-10-09T00:00:00Z")]
requests = [self.request("ordinary", "a", 0, 265, 265),
self.request("compaction", "a", 0, 530, 530, started=2, ended=3)]
clocks = analysis.clock_trace(events, requests)
self.assertTrue(clocks["valid"])
self.assertEqual(analysis.clocks_at(clocks, 4), (3, 3))
def test_unknown_usage_invalidates_normalized_curve(self):
events = [self.event("episode_start", 0, 1, utc="2026-10-09T00:00:00Z")]
request = self.request("unknown", "a", 0, 0, 0)
request.update(completion_tokens=None, prompt_tokens=None, status="unknown_usage_conservatively_charged")
clocks = analysis.clock_trace(events, [request])
self.assertFalse(clocks["valid"])
episode = {"condition": "single", "task": "a", "clock_valid": False,
"points": [{"wall_seconds": 30, "reference_seconds": 2, "score": 0.5, "order": 0}]}
row = analysis.aggregate_at([episode], ["a"], "single", 5, "reference_seconds")
self.assertIsNone(row["mean_hidden_test_fraction"])
self.assertEqual(analysis.aggregate_at([episode], ["a"], "single", 30, "wall_seconds")["mean_hidden_test_fraction"], 0.5)
if __name__ == "__main__":
unittest.main()