diff --git a/doc/reporting.rst b/doc/reporting.rst index 30ad69743..0bfdba85e 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -104,7 +104,7 @@ Enabling or disabling a report needs to be done in the system configuration: junit = { enable = true } The ``junit`` scenario reporter is disabled by default. When enabled, it writes ``junit.xml`` in the scenario results -directory. It emits one test case for every regular test iteration and every DSE step, including pass/fail status, +directory. It emits one test case for each existing regular test iteration and DSE step, including pass/fail status, failure details, scheduler duration when available, and the contents of ``stdout.txt`` and ``stderr.txt``. The artifact can be consumed directly by Jenkins, GitLab, GitHub Actions, and other CI systems that support JUnit XML. diff --git a/src/cloudai/_core/base_runner.py b/src/cloudai/_core/base_runner.py index 0a75918c7..43ded3286 100644 --- a/src/cloudai/_core/base_runner.py +++ b/src/cloudai/_core/base_runner.py @@ -21,6 +21,8 @@ from pathlib import Path from typing import Dict, List +import toml + import cloudai.models.output import cloudai.output @@ -162,7 +164,26 @@ def submit_test(self, tr: TestRun): self.update_run_output(job) except JobSubmissionError as e: logging.error(e) - exit(1) + self.record_execution_error(tr, e) + raise + + def record_execution_error(self, tr: TestRun, error: JobSubmissionError) -> None: + """Preserve a CloudAI failure in the run dump for report regeneration.""" + from cloudai.models.scenario import ExecutionError, TestRunDetails + + path = tr.output_path / CommandGenStrategy.TEST_RUN_DUMP_FILE_NAME + try: + if path.is_file(): + details = toml.load(path) + else: + details = TestRunDetails.from_test_run(tr, test_cmd="", full_cmd=error.command).model_dump( + exclude_none=True + ) + details["execution_error"] = ExecutionError(type=type(error).__name__, message=str(error)).model_dump() + with path.open("w") as stream: + toml.dump(details, stream) + except (OSError, toml.TomlDecodeError): + logging.exception("Failed to persist execution error for %s", tr.name) def on_job_submit(self, tr: TestRun) -> None: return diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index a472da9b5..f818278e8 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -20,7 +20,7 @@ from pathlib import Path from typing import TYPE_CHECKING, Any, Dict, Optional, Tuple, cast -from cloudai.core import METRIC_ERROR, BaseRunner, Registry, TestRun +from cloudai.core import METRIC_ERROR, BaseRunner, JobSubmissionError, Registry, TestRun from cloudai.util.lazy_imports import lazy from .base_agent import RewardOverrides @@ -195,6 +195,8 @@ def step(self, action: Any) -> Tuple[list, float, bool, dict]: try: self.runner.run() + except JobSubmissionError: + raise except Exception as e: logging.error(f"Error running step {self.test_run.step}: {e}") diff --git a/src/cloudai/core.py b/src/cloudai/core.py index 2a6e1dc58..55c251afe 100644 --- a/src/cloudai/core.py +++ b/src/cloudai/core.py @@ -25,6 +25,7 @@ from ._core.exceptions import ( JobFailureError, JobIdRetrievalError, + JobSubmissionError, MissingTestError, SystemConfigParsingError, TestConfigParsingError, @@ -104,6 +105,7 @@ "JobFailureError", "JobIdRetrievalError", "JobStatusResult", + "JobSubmissionError", "JsonGenStrategy", "MetricErrorSentinel", "MetricValue", diff --git a/src/cloudai/models/scenario.py b/src/cloudai/models/scenario.py index 543a0cf3c..249ed3ae4 100644 --- a/src/cloudai/models/scenario.py +++ b/src/cloudai/models/scenario.py @@ -265,6 +265,13 @@ def parse_reports(cls, value: dict[str, Any] | None) -> dict[str, ReportConfig] return parse_reports_spec(value) +class ExecutionError(BaseModel): + """Unrecoverable CloudAI execution error, distinct from a workload failure.""" + + type: str + message: str + + class TestRunDetails(BaseModel): """ Model for test run dump. @@ -285,6 +292,7 @@ class TestRunDetails(BaseModel): test_cmd: str full_cmd: str test_definition: Any + execution_error: ExecutionError | None = None @field_serializer("output_path") def _path_serializer(self, v: Path) -> str: diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 918fe26fb..9f217f46f 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -32,7 +32,7 @@ from cloudai.report_generator.util import load_system_metadata from .core import CommandGenStrategy, Reporter, TestRun, case_name -from .models.scenario import TestRunDetails +from .models.scenario import ExecutionError, TestRunDetails @dataclass @@ -148,15 +148,24 @@ class JUnitReporter(Reporter): def generate(self) -> None: self.load_test_runs() - results = [(tr, tr.test.was_run_successful(tr), self._duration(tr.output_path)) for tr in self.trs] - failures = sum(not status.is_successful for _, status, _ in results) - durations = [duration for _, _, duration in results if duration is not None] + results = [] + for tr in self.trs: + dump_path = tr.output_path / CommandGenStrategy.TEST_RUN_DUMP_FILE_NAME + details = toml.load(dump_path) if dump_path.is_file() else {} + error = ( + ExecutionError.model_validate(details["execution_error"]) if details.get("execution_error") else None + ) + status = tr.test.was_run_successful(tr) if error is None else None + results.append((tr, status, self._duration(tr.output_path), error)) + failures = sum(status is not None and not status.is_successful for _, status, _, _ in results) + errors = sum(error is not None for _, _, _, error in results) + durations = [duration for _, _, duration, _ in results if duration is not None] suite_attributes = { "name": self.test_scenario.name, "tests": str(len(results)), "failures": str(failures), - "errors": "0", + "errors": str(errors), "skipped": "0", } if durations: @@ -164,13 +173,19 @@ def generate(self) -> None: root = ET.Element("testsuites", suite_attributes) suite = ET.SubElement(root, "testsuite", suite_attributes) - for tr, status, duration in results: + for tr, status, duration, error in results: attributes = {"name": case_name(tr), "classname": self.test_scenario.name} if duration is not None: attributes["time"] = self._format_duration(duration) testcase = ET.SubElement(suite, "testcase", attributes) - if not status.is_successful: + if error is not None: + message = self._xml_text(f"{error.type}: {error.message}") + error_element = ET.SubElement( + testcase, "error", {"type": self._xml_text(error.type), "message": message} + ) + error_element.text = message + elif status is not None and not status.is_successful: message = status.error_message or "Test run failed" failure = ET.SubElement(testcase, "failure", {"message": self._xml_text(message)}) failure.text = self._xml_text(message) diff --git a/tests/test_base_runner.py b/tests/test_base_runner.py index 2a671a851..77ec71309 100644 --- a/tests/test_base_runner.py +++ b/tests/test_base_runner.py @@ -19,11 +19,13 @@ from typing import cast import pytest +import toml from pydantic import ConfigDict from cloudai.core import ( BaseJob, BaseRunner, + JobIdRetrievalError, JobStatusResult, System, TestDefinition, @@ -159,3 +161,28 @@ def test_end_post_comp(self, runner: MyRunner, tr_main: TestRun): runner.handle_dependencies(BaseJob(tr_main, 0)) assert len(runner.killed_by_dependency) == 1 assert runner.killed_by_dependency[0].test_run == tr_dep + + +@pytest.mark.parametrize("existing_dump", [False, True]) +def test_submission_error_is_preserved(runner: MyRunner, monkeypatch: pytest.MonkeyPatch, existing_dump: bool) -> None: + error = JobIdRetrievalError("tr-name", "sbatch test.sh", "", "submission timeout", "No job ID") + + def submit(tr: TestRun) -> BaseJob: + raise error + + monkeypatch.setattr(runner, "_submit_test", submit) + tr = runner.test_scenario.test_runs[0] + if existing_dump: + tr.output_path = runner.get_job_output_path(tr) + with (tr.output_path / "test-run.toml").open("w") as stream: + toml.dump({"name": tr.name, "full_cmd": "original command"}, stream) + with pytest.raises(JobIdRetrievalError) as caught: + runner.submit_test(tr) + assert caught.value is error + assert runner.jobs == [] + details = toml.load(tr.output_path / "test-run.toml") + assert details["name"] == tr.name + assert details["full_cmd"] == ("original command" if existing_dump else error.command) + assert details["execution_error"]["type"] == "JobIdRetrievalError" + assert "submission timeout" in details["execution_error"]["message"] + assert "sbatch test.sh" in details["execution_error"]["message"] diff --git a/tests/test_cloudaigym.py b/tests/test_cloudaigym.py index a7aa60720..164af22b0 100644 --- a/tests/test_cloudaigym.py +++ b/tests/test_cloudaigym.py @@ -27,7 +27,7 @@ Trajectory, ) from cloudai.configurator.env_params import EnvParamSpec, ObsLeafDescriptor -from cloudai.core import BaseRunner, RewardOverrides, Runner, TestRun, TestScenario +from cloudai.core import BaseRunner, JobIdRetrievalError, RewardOverrides, Runner, TestRun, TestScenario from cloudai.systems.slurm import SlurmRunner, SlurmSystem from cloudai.util import flatten_dict from cloudai.workloads.nemo_run import ( @@ -934,3 +934,19 @@ def test_reset_reports_the_regime_step_will_apply_on_info(self, tmp_path: Path) assert obs == env.define_observation_space(), "reset's flat obs stays the metrics placeholder" assert info["env_params"] == {"ball_speed": upcoming}, "reset peeks step+1 and reports the regime" assert env.encode_env_params(info["env_params"]) == {"ball_speed": [1, 2, 3].index(upcoming)} + + +def test_step_propagates_submission_error(setup_env: tuple[TestRun, BaseRunner], monkeypatch: pytest.MonkeyPatch): + test_run, runner = setup_env + test_run.test.cmd_args.data.global_batch_size = 8 + env = CloudAIGymEnv(test_run=test_run, runner=runner, rewards=RewardOverrides()) + agent = GridSearchAgent(env, GridSearchAgent.get_config_class()()) + _, action = agent.select_action() + error = JobIdRetrievalError(test_run.name, "sbatch test.sh", "", "timeout", "No job ID") + monkeypatch.setattr(runner, "run", MagicMock(side_effect=error)) + observation = MagicMock() + monkeypatch.setattr(env, "get_observation", observation) + with pytest.raises(JobIdRetrievalError) as caught: + env.step(action) + assert caught.value is error + observation.assert_not_called() diff --git a/tests/test_reporter.py b/tests/test_reporter.py index d3bdad65e..3d63d4489 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -475,6 +475,11 @@ def was_run_successful(tr: TestRun): return JobStatusResult(successful, message) monkeypatch.setattr(type(benchmark_tr.test), "was_run_successful", lambda self, tr: was_run_successful(tr)) + details = TestRunDetails.from_test_run(benchmark_tr, test_cmd="benchmark", full_cmd="sbatch test.sh") + dump = details.model_dump(exclude_none=True) + dump["execution_error"] = {"type": "JobIdRetrievalError", "message": "submission timeout"} + with (run_dirs[2] / "test-run.toml").open("w") as stream: + toml.dump(dump, stream) reporter = JUnitReporter( slurm_system, TestScenario(name="test-scenario", test_runs=[benchmark_tr]), @@ -489,7 +494,7 @@ def was_run_successful(tr: TestRun): "name": "test-scenario", "tests": "3", "failures": "1", - "errors": "0", + "errors": "1", "skipped": "0", "time": "6.000", } @@ -505,6 +510,11 @@ def was_run_successful(tr: TestRun): assert failure is not None assert failure.attrib["message"] == "benchmark failed" assert cases[1].findtext("system-err") == "stderr 1\n" + error = cases[2].find("error") + assert error is not None + assert error.attrib == {"type": "JobIdRetrievalError", "message": "JobIdRetrievalError: submission timeout"} + assert error.text == "JobIdRetrievalError: submission timeout" + assert cases[2].find("failure") is None def _write_slurm_job(step_dir: Path, elapsed_time_sec: int) -> None: