Spaces:
Running
Running
Download source/study/tests/test_controller.py from burtenshaw/beam-pi-programbench: direct link, hf CLI and curl.
- Browser
- Download file 16.5 kB
-
https://huggingface.co/spaces/burtenshaw/beam-pi-programbench/resolve/main/source/study/tests/test_controller.py
- Command line
-
hf download hf://spaces/burtenshaw/beam-pi-programbench/source/study/tests/test_controller.py
-
curl -L -o test_controller.py https://huggingface.co/spaces/burtenshaw/beam-pi-programbench/resolve/main/source/study/tests/test_controller.py
16.5 kB
| """Offline orchestration contracts; no containers, hosted APIs, or subprocess jobs.""" | |
| import copy | |
| import importlib.util | |
| import json | |
| from pathlib import Path | |
| import sys | |
| import tempfile | |
| import threading | |
| import time | |
| from types import SimpleNamespace | |
| import unittest | |
| from unittest import mock | |
| STUDY = Path(__file__).parents[1] | |
| sys.path.insert(0, str(STUDY)) | |
| spec = importlib.util.spec_from_file_location("study_controller", STUDY / "run_study.py") | |
| controller = importlib.util.module_from_spec(spec) | |
| spec.loader.exec_module(controller) | |
| class ControllerTests(unittest.TestCase): | |
| def setUp(self): | |
| self.temp = tempfile.TemporaryDirectory() | |
| self.root = Path(self.temp.name) | |
| self.config = json.loads((STUDY / "config.json").read_text()) | |
| self.config["snapshot_seconds"] = [30, 60, 120] | |
| self.config_path = self.root / "config.json" | |
| self.config_path.write_text(json.dumps(self.config)) | |
| self.study = controller.Study.__new__(controller.Study) | |
| self.study.root = self.root | |
| self.study.config = self.config | |
| self.study.manifest = {"tasks": self.config["tasks"]} | |
| self.study.args = SimpleNamespace(config=self.config_path, calibration=self.root / "calibration.json", | |
| manifest=self.root / "manifest.json", validate_only=False, | |
| evaluator_wheelhouse=self.root / "wheels") | |
| self.study.cancelled = False | |
| self.study.started = time.monotonic() | |
| self.study.events = self.root / "events.jsonl" | |
| self.study.event_lock = threading.Lock() | |
| self.study.gateway = SimpleNamespace(ledger=mock.Mock(), server_port=0) | |
| self.study.gateway.ledger.totals.return_value = { | |
| "api_tokens_charged": 0, "api_tokens_reserved": 0, "api_input_tokens": 0, | |
| "api_output_tokens": 0, "api_requests": 0, "unknown_usage_requests": 0} | |
| self.events = [] | |
| self.study.emit = mock.Mock(side_effect=lambda run_id, kind, **kwargs: self.events.append( | |
| {"run_id": run_id, "kind": kind, **kwargs})) | |
| self.study.validate = mock.Mock() | |
| self.study.prepare = mock.Mock(return_value={"container_id": "fake-container"}) | |
| self.study.cleanup_container = mock.Mock() | |
| self.study.runner = mock.Mock(side_effect=self.fake_runner) | |
| self.study.grade = mock.Mock(side_effect=self.fake_grade) | |
| calibration={"evaluator_dependencies_sha256":"test-dependencies","repetitions":2, | |
| "tasks":[{"instance_id":task["instance_id"],"repetitions":[ | |
| {"complete":True,"score":1.0,"evaluator_dependencies_sha256":"test-dependencies"} for _ in range(2) | |
| ]} for task in self.config["tasks"]]} | |
| self.calibration = mock.patch.object(controller.benchmark, "check_calibration", return_value=calibration) | |
| self.dependencies = mock.patch.object(controller.benchmark,"verify_wheelhouse",return_value="test-dependencies") | |
| self.process = mock.patch.object(controller.subprocess, "run", return_value=SimpleNamespace(returncode=0)) | |
| self.calibration.start() | |
| self.dependencies.start() | |
| self.process.start() | |
| def tearDown(self): | |
| self.calibration.stop() | |
| self.dependencies.stop() | |
| self.process.stop() | |
| self.temp.cleanup() | |
| def benchmarks(self): | |
| return [run for run in self.config["runs"] if run["category"] == "benchmark"] | |
| def fake_runner(self, run, prepared, directory, *args): | |
| result = {"stop_reason": "completed", "elapsed_seconds": 120, "snapshots": [], | |
| "snapshot_error": None, "agents": [{"id": "agent-0"}]} | |
| controller.save(directory / "runner" / "result.json", result) | |
| return result | |
| def fake_grade(self, run, directory, result): | |
| grades = [{"elapsed_seconds": 120, "valid": True, "score": 0.25}] | |
| controller.save(directory / "grades.json", grades) | |
| return grades | |
| def final_event(self): | |
| return [event for event in self.events if event["run_id"] == "study-summary" and event["kind"] == "run_finished"][-1] | |
| def test_plan_contains_exact_five_by_three_design_under_130m_cap(self): | |
| tasks = {task["instance_id"] for task in self.config["tasks"]} | |
| runs = self.benchmarks() | |
| self.assertEqual(len(tasks), 5) | |
| self.assertEqual(len(runs), 15) | |
| self.assertEqual({(run["task"], run["condition"]) for run in runs}, | |
| {(task, arm) for task in tasks for arm in ("single", "peers", "async")}) | |
| self.assertEqual(sum(run["api_token_cap"] for run in runs), 120_000_000) | |
| self.assertLessEqual(self.config["global_api_token_cap"], 130_000_000) | |
| def test_successful_controller_attempts_each_planned_run_once(self): | |
| self.study.run() | |
| seen = [call.args[0]["run_id"] for call in self.study.runner.call_args_list] | |
| self.assertCountEqual(seen, [run["run_id"] for run in self.benchmarks()]) | |
| self.assertEqual(len(seen), len(set(seen))) | |
| self.assertEqual(self.final_event()["status"], "completed") | |
| def test_failed_calibration_prevents_all_live_validation_and_benchmark_calls(self): | |
| controller.benchmark.check_calibration.side_effect = RuntimeError("reference failed") | |
| with self.assertRaises(RuntimeError): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| self.study.prepare.assert_not_called() | |
| def test_failed_live_validation_prevents_benchmark_calls(self): | |
| self.study.validate.side_effect = RuntimeError("live validation failed") | |
| with self.assertRaises(RuntimeError): | |
| self.study.run() | |
| self.study.prepare.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_evaluator_dependency_mismatch_prevents_all_model_calls(self): | |
| controller.benchmark.verify_wheelhouse.return_value="different-dependencies" | |
| with self.assertRaisesRegex(ValueError,"dependency cache differs"): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.prepare.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_missing_evaluator_dependency_cache_prevents_all_model_calls(self): | |
| self.study.args.evaluator_wheelhouse=None | |
| with self.assertRaisesRegex(ValueError,"evaluator-wheelhouse is required"): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_one_reference_repetition_is_not_enough_to_unlock_model_calls(self): | |
| controller.benchmark.check_calibration.return_value["repetitions"]=1 | |
| with self.assertRaisesRegex(ValueError,"Two complete reference"): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_failed_reference_score_is_not_hidden_by_top_level_pass_flag(self): | |
| reference=controller.benchmark.check_calibration.return_value["tasks"][0]["repetitions"][0] | |
| reference["score"]=0.89 | |
| with self.assertRaisesRegex(ValueError,"score at least 0.9"): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_each_reference_repetition_must_use_pinned_dependencies(self): | |
| reference=controller.benchmark.check_calibration.return_value["tasks"][0]["repetitions"][0] | |
| reference["evaluator_dependencies_sha256"]="old-dependencies" | |
| with self.assertRaisesRegex(ValueError,"repetition dependency cache differs"): | |
| self.study.run() | |
| self.study.validate.assert_not_called() | |
| self.study.runner.assert_not_called() | |
| def test_terminal_and_interrupted_existing_runs_never_get_paid_retry(self): | |
| prior_statuses = ["completed", "failed", "interrupted", "budget_exhausted", "running"] | |
| skipped = self.benchmarks()[:len(prior_statuses)] | |
| for run, status in zip(skipped, prior_statuses): | |
| controller.save(self.root / "episodes" / run["run_id"] / "status.json", {"status": status}) | |
| self.study.run() | |
| attempted = {call.args[0]["run_id"] for call in self.study.runner.call_args_list} | |
| self.assertFalse(attempted & {run["run_id"] for run in skipped}) | |
| self.assertEqual(len(attempted), 10) | |
| restarted = controller.load(self.root / "episodes" / skipped[-1]["run_id"] / "status.json") | |
| self.assertEqual(restarted["status"], "interrupted") | |
| def test_global_quota_stop_after_13_runs_is_not_completed_study(self): | |
| counter = 0 | |
| def quota_runner(run, prepared, directory, *args): | |
| nonlocal counter | |
| counter += 1 | |
| result = self.fake_runner(run, prepared, directory) | |
| if counter == 13: | |
| result["stop_reason"] = "daily_quota_headroom" | |
| return result | |
| self.study.runner.side_effect = quota_runner | |
| self.study.run() | |
| self.assertEqual(counter, 13) | |
| self.assertNotEqual(self.final_event()["status"], "completed") | |
| def test_live_validation_rejects_agent_error_even_when_artifact_tests_pass(self): | |
| self.study.config = copy.deepcopy(self.config) | |
| self.study.config["runs"] = [run for run in self.config["runs"] if run["run_id"] == "validation-single"] | |
| def failed_runner(run, prepared, directory, *args): | |
| result = self.fake_runner(run, prepared, directory) | |
| result["stop_reason"] = "agent_error" | |
| (directory / "runner" / "events.jsonl").write_text(json.dumps({"type": "model_request_end"}) + "\n") | |
| return result | |
| self.study.runner.side_effect = failed_runner | |
| with self.assertRaises(RuntimeError): | |
| controller.Study.validate(self.study) | |
| self.assertEqual(controller.load(self.root / "validation.json")["status"], "failed") | |
| def grade_fixture(self): | |
| run = self.benchmarks()[0] | |
| directory = self.root / "episodes" / run["run_id"] | |
| directory.mkdir(parents=True) | |
| first, second = directory / "first.tar.gz", directory / "second.tar.gz" | |
| first.write_bytes(b"first artifact") | |
| second.write_bytes(b"second artifact") | |
| result = {"elapsed_seconds": 120, "snapshots": [{"path": str(first), "elapsed_ms": 20_000}, | |
| {"path": str(second), "elapsed_ms": 40_000}]} | |
| return run, directory, result | |
| def test_grading_uses_latest_snapshot_before_cutoff_and_deduplicates_artifacts(self): | |
| run, directory, result = self.grade_fixture() | |
| self.study.adapter = mock.Mock(side_effect=[{"score": 0.75, "complete": True}, {"score": 0.25, "complete": True}]) | |
| grades = controller.Study.grade(self.study, run, directory, result) | |
| self.assertEqual(self.study.adapter.call_count, 2) | |
| self.assertEqual(next(grade for grade in grades if grade["elapsed_seconds"] == 30)["score"], 0.75) | |
| self.assertEqual(next(grade for grade in grades if grade["elapsed_seconds"] == 60)["score"], 0.25) | |
| def test_grader_infrastructure_failure_is_invalid_null_and_never_published_as_score(self): | |
| run, directory, result = self.grade_fixture() | |
| self.study.adapter = mock.Mock(side_effect=RuntimeError("unavailable evaluator")) | |
| grades = controller.Study.grade(self.study, run, directory, result) | |
| self.assertTrue(all(not grade["valid"] and grade["score"] is None for grade in grades)) | |
| for event in self.events: | |
| if event["kind"] == "score": | |
| self.assertNotIn("hidden_test_fraction", event["metrics"]) | |
| def test_classified_submission_compile_failure_is_valid_zero(self): | |
| run, directory, result = self.grade_fixture() | |
| self.study.adapter = mock.Mock(return_value={"score": 0.0, "complete": False, | |
| "valid": True, "analysis_score": 0.0, | |
| "scoring_status": "submission_failed", | |
| "error_code": "compile_failed", "passed": 0, | |
| "test_count": 100, "expected_test_count": 100}) | |
| grades = controller.Study.grade(self.study, run, directory, result) | |
| self.assertTrue(all(grade["valid"] and grade["score"] == 0.0 for grade in grades)) | |
| scores = [event["metrics"]["hidden_test_fraction"] for event in self.events if event["kind"] == "score"] | |
| self.assertEqual(scores, [0.0] * len(grades)) | |
| def test_cancelled_grade_does_not_dispatch_another_snapshot(self): | |
| run,directory,result=self.grade_fixture() | |
| def first_grade(*args,**kwargs): | |
| self.study.cancelled=True | |
| return {"score":0.25,"complete":True} | |
| self.study.adapter=mock.Mock(side_effect=first_grade) | |
| grades=controller.Study.grade(self.study,run,directory,result) | |
| self.assertEqual(self.study.adapter.call_count,1) | |
| self.assertEqual(len(grades),1) | |
| self.assertTrue((directory/"grades.json").exists()) | |
| def test_adapter_cancellation_terminates_process_group_and_reaps_parent(self): | |
| self.study.args.evaluator_python=sys.executable | |
| self.study.args.programbench_root=self.root | |
| self.study.args.resource_mode="slurm" | |
| self.study.current=None | |
| process=mock.Mock(pid=12345,returncode=None) | |
| process.poll.side_effect=lambda: process.returncode | |
| def kill_group(pid,sig): | |
| self.assertEqual(pid,12345) | |
| if sig==controller.signal.SIGINT: | |
| process.returncode=-int(sig) | |
| def wait(timeout): | |
| self.study.stop() | |
| return process.returncode | |
| process.wait.side_effect=wait | |
| with mock.patch.object(controller.subprocess,"Popen",return_value=process) as popen, \ | |
| mock.patch.object(controller.os,"killpg",side_effect=kill_group) as killpg: | |
| with self.assertRaises(InterruptedError): | |
| controller.Study.adapter(self.study,"grade",output_log=self.root/"adapter.log") | |
| self.assertTrue(popen.call_args.kwargs["start_new_session"]) | |
| self.assertEqual(killpg.call_args_list,[mock.call(12345,controller.signal.SIGINT),mock.call(12345,controller.signal.SIGKILL)]) | |
| self.assertIsNone(self.study.current) | |
| self.assertFalse(self.study.current_is_adapter) | |
| process.wait.assert_called() | |
| def test_adapter_is_not_started_after_cancellation(self): | |
| self.study.cancelled=True | |
| with mock.patch.object(controller.subprocess,"Popen") as popen: | |
| with self.assertRaises(InterruptedError): | |
| controller.Study.adapter(self.study,"grade",output_log=self.root/"adapter.log") | |
| popen.assert_not_called() | |
| def test_main_failure_and_interruption_return_nonzero_after_close(self): | |
| argv=["run_study.py","--programbench-root",str(self.root),"--calibration",str(self.root/"calibration.json"), | |
| "--output-dir",str(self.root/"output"),"--secrets-dir",str(self.root/"secrets")] | |
| for status,expected in [("failed",1),("interrupted",130),("completed",0)]: | |
| with self.subTest(status=status): | |
| study=mock.Mock(cancelled=False) | |
| study.run.return_value={"status":status} | |
| with mock.patch.object(controller,"Study",return_value=study), \ | |
| mock.patch.object(controller.signal,"signal"),mock.patch.object(sys,"argv",argv): | |
| self.assertEqual(controller.main(),expected) | |
| study.close.assert_called_once() | |
| def test_main_cancellation_exception_still_closes_and_returns_130(self): | |
| study=mock.Mock(cancelled=True) | |
| study.run.side_effect=InterruptedError("cancelled during validation") | |
| argv=["run_study.py","--programbench-root",str(self.root),"--calibration",str(self.root/"calibration.json"), | |
| "--output-dir",str(self.root/"output"),"--secrets-dir",str(self.root/"secrets")] | |
| with mock.patch.object(controller,"Study",return_value=study),mock.patch.object(controller.signal,"signal"),mock.patch.object(sys,"argv",argv): | |
| self.assertEqual(controller.main(),130) | |
| study.close.assert_called_once() | |
| if __name__ == "__main__": | |
| unittest.main() | |