Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion doc/reporting.rst
Original file line number Diff line number Diff line change
Expand Up @@ -98,7 +98,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.

Expand Down
23 changes: 22 additions & 1 deletion src/cloudai/_core/base_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
from pathlib import Path
from typing import Dict, List

import toml

import cloudai.models.output
import cloudai.output

Expand Down Expand Up @@ -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
Expand Down
4 changes: 3 additions & 1 deletion src/cloudai/configurator/cloudai_gym.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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}")

Expand Down
2 changes: 2 additions & 0 deletions src/cloudai/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
from ._core.exceptions import (
JobFailureError,
JobIdRetrievalError,
JobSubmissionError,
MissingTestError,
SystemConfigParsingError,
TestConfigParsingError,
Expand Down Expand Up @@ -102,6 +103,7 @@
"JobFailureError",
"JobIdRetrievalError",
"JobStatusResult",
"JobSubmissionError",
"JsonGenStrategy",
"MetricErrorSentinel",
"MetricValue",
Expand Down
8 changes: 8 additions & 0 deletions src/cloudai/models/scenario.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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:
Expand Down
29 changes: 22 additions & 7 deletions src/cloudai/reporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -148,29 +148,44 @@ 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:
suite_attributes["time"] = self._format_duration(sum(durations))

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)
Expand Down
27 changes: 27 additions & 0 deletions tests/test_base_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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"]
18 changes: 17 additions & 1 deletion tests/test_cloudaigym.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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()
12 changes: 11 additions & 1 deletion tests/test_reporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -467,6 +467,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]),
Expand All @@ -481,7 +486,7 @@ def was_run_successful(tr: TestRun):
"name": "test-scenario",
"tests": "3",
"failures": "1",
"errors": "0",
"errors": "1",
"skipped": "0",
"time": "6.000",
}
Expand All @@ -497,6 +502,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:
Expand Down
Loading