diff --git a/doc/examples/nccl-dse/experiment.json b/doc/examples/nccl-dse/experiment.json new file mode 100644 index 000000000..1c2473bf4 --- /dev/null +++ b/doc/examples/nccl-dse/experiment.json @@ -0,0 +1,118 @@ +{ + "id": "nccl-dse_2026-09-21_14-55-45", + "name": "nccl-dse", + "system_name": "example-cluster", + "description": null, + "status": "completed", + "path": "/results/nccl-dse_2026-09-21_14-55-45", + "start": "2026-09-21T14:55:45.275325Z", + "finish": "2026-09-21T15:01:50.432934Z", + "duration": 365.157609, + "tests": [ + { + "id": "nccl", + "name": "nccl_base_test", + "description": "NCCL base test configuration", + "status": "completed", + "path": "/results/nccl-dse_2026-09-21_14-55-45/nccl", + "metrics": [ + { + "name": "Latency", + "value": 24.76, + "unit": "us", + "dimensions": [ + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + }, + { + "name": "Bandwidth", + "value": 74.12, + "unit": "GB/s", + "dimensions": [ + {"name": "Bandwidth basis", "value": "bus", "unit": "", "is_x": false}, + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + } + ], + "runs": [ + { + "path": "/results/nccl-dse_2026-09-21_14-55-45/nccl/0/1", + "jobid": "1001", + "status": "completed", + "metrics": [ + { + "name": "Latency", + "value": 24.76, + "unit": "us", + "dimensions": [ + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + }, + { + "name": "Bandwidth", + "value": 74.12, + "unit": "GB/s", + "dimensions": [ + {"name": "Bandwidth basis", "value": "bus", "unit": "", "is_x": false}, + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + } + ], + "start": "2026-09-21T14:55:46Z", + "finish": "2026-09-21T14:57:28Z", + "duration": 102.0, + "iteration": 0, + "step": 1 + }, + { + "path": "/results/nccl-dse_2026-09-21_14-55-45/nccl/0/2", + "jobid": "1002", + "status": "completed", + "metrics": [ + { + "name": "Latency", + "value": 37.36, + "unit": "us", + "dimensions": [ + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + }, + { + "name": "Bandwidth", + "value": 49.11, + "unit": "GB/s", + "dimensions": [ + {"name": "Bandwidth basis", "value": "bus", "unit": "", "is_x": false}, + {"name": "Operation", "value": "all_reduce", "unit": "", "is_x": false}, + {"name": "Placement", "value": "out_of_place", "unit": "", "is_x": false}, + {"name": "Size", "value": "1048576", "unit": "", "is_x": true} + ] + } + ], + "start": "2026-09-21T14:58:47Z", + "finish": "2026-09-21T15:00:28Z", + "duration": 101.0, + "iteration": 0, + "step": 2 + } + ], + "dse": { + "space": { + "extra_env_vars.NCCL_ALGO": ["Ring", "Tree"] + }, + "best_config": {"extra_env_vars.NCCL_ALGO": "Ring"}, + "best_step": 1 + } + } + ] +} diff --git a/doc/reporting.rst b/doc/reporting.rst index 0b907480a..386f6d72c 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -36,16 +36,36 @@ Unified experiment output CloudAI writes ``experiment.json`` in each scenario's results directory. -The file contains scenario details, the configured system name, test cases, status, timing, and result paths. Standalone -execution also records each run's process ID, status, iteration, timing, and workload metrics. +The file contains scenario details, the system name, and test cases under ``tests``. For Slurm runs, the system name is +the cluster reported by Slurm. Standalone and Slurm execution records appear in each test case's ``runs`` list. Each +record represents an iteration or DSE step and includes its number, process or Slurm job ID, status, timing, result path, +and workload metrics. -CloudAI updates the file when standalone runs start and finish, then finalizes it when scenario execution succeeds or -fails. Timestamps use UTC; durations use seconds. Unknown timestamps are ``null``. Dry runs also produce scenario and -test-case details without launching workloads. +Timestamps use UTC; durations use seconds. Unknown timestamps are ``null``. A final status of ``unknown`` means the outcome +could not be determined. Dry runs include scenario and test-case details without launching workloads. -Metrics come from ``TestDefinition.metric_observations()``, independently of reporter settings. A completed ordinary -test with one run also includes those metrics on the test case. Each file update replaces the previous snapshot -atomically. Metric extraction or write errors produce warnings without affecting execution. +Metrics come from ``TestDefinition.metric_observations()``, independently of reporter settings. + +When a test case executes once successfully, ``tests[].metrics`` contains that execution's metrics. For DSE, it contains +metrics from the successful step with the highest valid reward. The search space, selected step, and configuration appear +in ``tests[].dse``. + +NCCL DSE example +~~~~~~~~~~~~~~~~ + +This example comes from a Slurm NCCL all-reduce run on one node with eight H100 GPUs. DSE tried ``Ring`` and ``Tree`` +in two steps. ``Ring`` (step 1) won using inverse latency as the reward, so its measurements appear in ``tests[].metrics``. + +Identifiers and paths are anonymized. To keep the example small, each metrics list includes only out-of-place latency +and bus bandwidth for 1 MiB messages. Measurement values, timing, and DSE selection are unchanged. + +:download:`Download experiment.json `. + +.. literalinclude:: examples/nccl-dse/experiment.json + :language: json + +Each file update replaces the previous snapshot atomically. Metric extraction or write errors produce warnings without +affecting execution. .. _general-flow: diff --git a/src/cloudai/_core/base_runner.py b/src/cloudai/_core/base_runner.py index 17f85c2c9..0a75918c7 100644 --- a/src/cloudai/_core/base_runner.py +++ b/src/cloudai/_core/base_runner.py @@ -114,7 +114,9 @@ def update_run_output(self, job: BaseJob, result: JobStatusResult | None = None) self.experiment_output.write() def finish_output(self, successful: bool) -> None: - status = "completed" if successful else "failed" + status: cloudai.models.output.Status = "completed" if successful else "failed" + if self.shutting_down and not any(test.status == "failed" for test in self.experiment_output.experiment.tests): + status = "cancelled" self.experiment_output.finish(status=status, finish=datetime.datetime.now(datetime.timezone.utc)) def shutdown(self): diff --git a/src/cloudai/cli/handlers.py b/src/cloudai/cli/handlers.py index 313cbbed2..09402d619 100644 --- a/src/cloudai/cli/handlers.py +++ b/src/cloudai/cli/handlers.py @@ -173,7 +173,11 @@ def handle_dse_job(runner: Runner, args: argparse.Namespace) -> int: agent = agent_class(env, agent_config) logging.debug(f"Created agent {agent.__class__.__name__}.") - err |= agent.run() + env.update_output() + try: + err |= agent.run() + finally: + env.update_output() except Exception as exc: run_error = exc logging.exception("DSE job aborted by an unexpected error; generating reports before failing.") diff --git a/src/cloudai/configurator/cloudai_gym.py b/src/cloudai/configurator/cloudai_gym.py index cdafd5ce5..a472da9b5 100644 --- a/src/cloudai/configurator/cloudai_gym.py +++ b/src/cloudai/configurator/cloudai_gym.py @@ -16,6 +16,7 @@ import copy import logging +import math from pathlib import Path from typing import TYPE_CHECKING, Any, Dict, Optional, Tuple, cast @@ -62,6 +63,34 @@ def __init__( self.trajectory = Trajectory(iteration_dir=self.iteration_dir) super().__init__() + def _ranked_dse_candidates(self) -> list[tuple[int, dict[str, str | int | float]]]: + """Return valid trial steps and configurations in descending reward order.""" + tr = self.original_test_run + candidates = [ + row + for row in self.trajectory.dataframe.to_dict(orient="records") + if math.isfinite(row["reward"]) + and all( + math.isfinite(row[f"observation.{metric}"]) + and row[f"observation.{metric}"] != self.rewards.metric_failure + for metric in tr.test.agent_metrics + ) + ] + return [ + (int(row["step"]), {key: row[f"action.{key}"] for key in tr.param_space}) + for row in sorted(candidates, key=lambda row: row["reward"], reverse=True) + ] + + def update_output(self) -> None: + """Publish the DSE recommendation without interrupting execution on output errors.""" + try: + tr = self.original_test_run + output = self.runner.experiment_output + output.update_dse(str(tr.name), tr.param_space, self._ranked_dse_candidates()) + output.write() + except Exception as exc: + logging.warning("Cannot update DSE output for %s: %s", self.original_test_run.name, exc) + @property def upcoming_trial(self) -> int: """ diff --git a/src/cloudai/models/output.py b/src/cloudai/models/output.py index 71af6811f..76a41c576 100644 --- a/src/cloudai/models/output.py +++ b/src/cloudai/models/output.py @@ -76,7 +76,7 @@ class Test(pydantic.BaseModel): class Experiment(pydantic.BaseModel): - """Full snapshot spanning the entire scenario, including all DSE trials.""" + """Full scenario snapshot; changes to this model tree may require updating doc/reporting.rst.""" id: str name: str diff --git a/src/cloudai/output.py b/src/cloudai/output.py index 7d3e1aec0..67755d893 100644 --- a/src/cloudai/output.py +++ b/src/cloudai/output.py @@ -19,6 +19,7 @@ import pathlib import tempfile +import cloudai.metrics import cloudai.models.output @@ -29,6 +30,23 @@ def elapsed_seconds(start: datetime.datetime | None, finish: datetime.datetime | return max((finish - start).total_seconds(), 0.0) +def metric_output(observation: cloudai.metrics.MetricObservation) -> cloudai.models.output.Metric: + """Convert a canonical observation to an output metric.""" + return cloudai.models.output.Metric( + name=observation.metric.display_name, + value=observation.value, + unit=observation.metric.unit, + dimensions=[ + cloudai.models.output.Dimension( + name=cloudai.metrics.dimension_label(key), + value=str(value), + is_x=key == cloudai.metrics.SIZE_BYTES.key, + ) + for key, value in sorted(observation.dimensions.items()) + ], + ) + + class ExperimentOutput: """Collect experiment results and publish snapshots.""" @@ -63,6 +81,25 @@ def update_test(self, test: cloudai.models.output.Test) -> None: return self.experiment.tests.append(recorded) + def update_dse( + self, + test_id: str, + space: dict[str, list[str | int | float]], + candidates: list[tuple[int, dict[str, str | int | float]]], + ) -> None: + """Store the first ranked candidate with a successful run and its metrics.""" + test = next((test for test in self.experiment.tests if test.id == test_id), None) + if test is None: + raise KeyError(f"Unknown experiment test: {test_id}") + completed_runs = {run.step: run for run in test.runs if run.status == "completed"} + for step, config in candidates: + if step in completed_runs: + test.dse = cloudai.models.output.DSE(space=space, best_step=step, best_config=config) + test.metrics = [metric.model_copy(deep=True) for metric in completed_runs[step].metrics] + return + test.dse = cloudai.models.output.DSE(space=space) + test.metrics = [] + def snapshot(self) -> cloudai.models.output.Experiment: """Return an independent snapshot without finalizing the experiment.""" full = self.experiment.model_copy(deep=True) @@ -97,16 +134,25 @@ def write(self) -> None: logging.warning("Cannot remove temporary experiment output %s: %s", temporary_path, exc) def finish(self, status: cloudai.models.output.Status, finish: datetime.datetime | None) -> None: - self.experiment.status = status - self.experiment.finish = finish - self.experiment.duration = elapsed_seconds(self.experiment.start, finish) for test in self.experiment.tests: + for run in test.runs: + if run.status in ("pending", "running"): + run.status = "unknown" + self._update_test_status(test) if test.dse is None: test.metrics = self._single_run_metrics(test) if status == "completed": - for test in self.experiment.tests: - if test.status not in ("failed", "cancelled"): - test.status = "completed" + statuses = {test.status for test in self.experiment.tests} + for outcome in ("failed", "cancelled", "unknown"): + if outcome in statuses: + status = outcome + break + for test in self.experiment.tests: + if test.status in ("pending", "running"): + test.status = "completed" if status == "completed" else "unknown" + self.experiment.status = status + self.experiment.finish = finish + self.experiment.duration = elapsed_seconds(self.experiment.start, finish) self.write() @staticmethod @@ -128,5 +174,9 @@ def _update_test_status(test: cloudai.models.output.Test) -> None: test.status = "cancelled" elif "running" in statuses: test.status = "running" + elif "pending" in statuses: + test.status = "pending" + elif "unknown" in statuses: + test.status = "unknown" elif statuses == {"completed"}: test.status = "completed" diff --git a/src/cloudai/systems/slurm/single_sbatch_runner.py b/src/cloudai/systems/slurm/single_sbatch_runner.py index 5db16fc30..4aa0972de 100644 --- a/src/cloudai/systems/slurm/single_sbatch_runner.py +++ b/src/cloudai/systems/slurm/single_sbatch_runner.py @@ -23,7 +23,7 @@ from cloudai.configurator import CloudAIGymEnv from cloudai.configurator.env_params import EnvParams -from cloudai.core import BaseJob, Registry, System, TestRun, TestScenario +from cloudai.core import BaseJob, JobStatusResult, Registry, System, TestRun, TestScenario from cloudai.util import format_time_limit, parse_time_limit from .slurm_command_gen_strategy import SlurmCommandGenStrategy @@ -244,6 +244,9 @@ def handle_dse(self): def completed_test_runs(self, job: BaseJob) -> list[TestRun]: return list(self.all_trs) + def get_run_output(self, job: BaseJob, tr: TestRun, result: JobStatusResult | None = None) -> None: + return None + def _submit_test(self, tr: TestRun) -> SlurmJob: with open(self.scenario_root / "cloudai_sbatch_script.sh", "w") as f: f.write(self.gen_sbatch_content()) diff --git a/src/cloudai/systems/slurm/slurm_job.py b/src/cloudai/systems/slurm/slurm_job.py index c5633b765..2b903e375 100644 --- a/src/cloudai/systems/slurm/slurm_job.py +++ b/src/cloudai/systems/slurm/slurm_job.py @@ -18,9 +18,12 @@ from cloudai.core import BaseJob +from .slurm_metadata import SlurmJobMetadata + @dataclass class SlurmJob(BaseJob): """A job class for execution on a Slurm system.""" nodes: list[str] = field(default_factory=list, init=False) + metadata: SlurmJobMetadata | None = field(default=None, init=False) diff --git a/src/cloudai/systems/slurm/slurm_metadata.py b/src/cloudai/systems/slurm/slurm_metadata.py index 74a50e2be..ced74a8f8 100644 --- a/src/cloudai/systems/slurm/slurm_metadata.py +++ b/src/cloudai/systems/slurm/slurm_metadata.py @@ -42,6 +42,7 @@ class SlurmStepMetadata(_SlurmStepMetadataBase): step_id: str submit_line: str + cluster_name: str = Field(default="", exclude=True) @classmethod def from_sacct_output(cls, output: str, delimiter: str) -> list[SlurmStepMetadata]: @@ -54,10 +55,15 @@ def from_sacct_output(cls, output: str, delimiter: str) -> list[SlurmStepMetadat @classmethod def _from_sacct_single_line(cls, line: str, delimiter: str) -> SlurmStepMetadata | None: data = line.split(delimiter) + if data and not data[-1]: + data.pop() if len(data) < 8: return None job_id, step_id = data[0].split(".") if "." in data[0] else (data[0], "") + has_cluster_name = len(data) >= 9 + cluster_name = data[7] if has_cluster_name else "" + submit_line = delimiter.join(data[8:] if has_cluster_name else data[7:]) return cls( job_id=int(job_id), @@ -68,7 +74,8 @@ def _from_sacct_single_line(cls, line: str, delimiter: str) -> SlurmStepMetadata start_time=data[4], end_time=data[5], elapsed_time_sec=int(data[6]), - submit_line=data[7], + cluster_name=cluster_name, + submit_line=submit_line, ) diff --git a/src/cloudai/systems/slurm/slurm_rest_client.py b/src/cloudai/systems/slurm/slurm_rest_client.py index c42a27d26..d46033d51 100644 --- a/src/cloudai/systems/slurm/slurm_rest_client.py +++ b/src/cloudai/systems/slurm/slurm_rest_client.py @@ -362,10 +362,15 @@ def get_allocated_nodes(self) -> list[SlurmNode]: return nodes def _get_job(self, job_id: int, retry_threshold: int = 3) -> dict[str, Any] | None: - jobs = self._request("GET", "slurm", f"job/{job_id}", retry_threshold=retry_threshold).get("jobs") + response = self._request("GET", "slurm", f"job/{job_id}", retry_threshold=retry_threshold) + jobs = response.get("jobs") if not isinstance(jobs, list): raise RuntimeError("Slurm API returned an invalid jobs response.") - return cast(dict[str, Any], jobs[0]) if jobs else None + if not jobs: + return None + job = cast(dict[str, Any], jobs[0]).copy() + job["cluster_name"] = response.get("meta", {}).get("slurm", {}).get("cluster", "") + return job def get_job_state(self, job_id: int, retry_threshold: int = 3) -> str: """Return current job state from slurmctld.""" @@ -418,6 +423,7 @@ def get_job_status(self, job_id: int, retry_threshold: int = 3) -> list[SlurmSte end_time=end_time, elapsed_time_sec=elapsed_seconds, submit_line=job.get("command") or "", + cluster_name=job["cluster_name"], ) ] diff --git a/src/cloudai/systems/slurm/slurm_runner.py b/src/cloudai/systems/slurm/slurm_runner.py index 93a177962..67a44b159 100644 --- a/src/cloudai/systems/slurm/slurm_runner.py +++ b/src/cloudai/systems/slurm/slurm_runner.py @@ -14,13 +14,16 @@ # See the License for the specific language governing permissions and # limitations under the License. +import datetime import logging from pathlib import Path from typing import cast import toml -from cloudai.core import BaseJob, BaseRunner, System, TestRun, TestScenario +import cloudai.models.output +import cloudai.output +from cloudai.core import BaseJob, BaseRunner, JobStatusResult, System, TestRun, TestScenario from .slurm_command_gen_strategy import SlurmCommandGenStrategy from .slurm_job import SlurmJob @@ -36,6 +39,44 @@ def __init__(self, mode: str, system: System, test_scenario: TestScenario, outpu self.system = cast(SlurmSystem, system) self.pinned_nodes: dict[str, list[str]] = {} + def get_run_output( + self, job: BaseJob, tr: TestRun, result: JobStatusResult | None = None + ) -> cloudai.models.output.Run | None: + metadata = cast(SlurmJob, job).metadata + status: cloudai.models.output.Status = "pending" + metrics: list[cloudai.models.output.Metric] = [] + if result is not None: + status = "completed" if result.is_successful else "failed" + if job.terminated_by_dependency or (metadata is not None and metadata.state.startswith("CANCELLED")): + status = "cancelled" + if status == "completed": + try: + metrics = [ + cloudai.output.metric_output(observation) + for observation in tr.test.metric_observations(self.system, tr) + ] + except Exception as exc: + logging.warning("Cannot extract output metrics for Slurm job %s: %s", job.id, exc) + return cloudai.models.output.Run( + path=str(tr.output_path.absolute()), + jobid=str(job.id), + status=status, + metrics=metrics, + start=self._output_timestamp(metadata.start_time) if metadata is not None else None, + finish=self._output_timestamp(metadata.end_time) if metadata is not None else None, + duration=metadata.elapsed_time_sec if metadata is not None else None, + iteration=tr.current_iteration, + step=tr.step, + ) + + @staticmethod + def _output_timestamp(value: str) -> datetime.datetime | None: + try: + timestamp = datetime.datetime.fromisoformat(value.replace("Z", "+00:00")) + except ValueError: + return None + return timestamp.astimezone(datetime.timezone.utc) if timestamp.utcoffset() is not None else None + def submit_test(self, tr: TestRun) -> None: if tr.pin_nodes and tr.name in self.pinned_nodes: tr.nodes = self.pinned_nodes[tr.name].copy() @@ -119,7 +160,14 @@ def _get_job_metadata( def store_job_metadata(self, job: SlurmJob): system = cast(SlurmSystem, self.system) steps_metadata = [self._mock_job_metadata()] if self.mode == "dry-run" else system.get_job_status(job) + if not steps_metadata: + logging.warning("No Slurm accounting metadata available for job %s", job.id) + return + cluster_name = steps_metadata[0].cluster_name slurm_job_file, job_meta = self._get_job_metadata(job, steps_metadata) + job.metadata = job_meta + if cluster_name: + self.experiment_output.experiment.system_name = cluster_name logging.debug(f"Storing job metadata for job {job.id} to {slurm_job_file}") with slurm_job_file.open("w") as job_file: diff --git a/src/cloudai/systems/slurm/slurm_system.py b/src/cloudai/systems/slurm/slurm_system.py index 7b3d31010..d01f7bce7 100644 --- a/src/cloudai/systems/slurm/slurm_system.py +++ b/src/cloudai/systems/slurm/slurm_system.py @@ -481,7 +481,8 @@ def get_job_status(self, job: BaseJob, retry_threshold: int = 3) -> list[SlurmSt retry_count = 0 command = ( - f"sacct -j {job.id} --format=JobID,JobName,State,ExitCode,Start,End,ElapsedRAW,SubmitLine " + "TZ=UTC SLURM_TIME_FORMAT='%Y-%m-%dT%H:%M:%SZ' " + f"sacct -j {job.id} --format=JobID,JobName,State,ExitCode,Start,End,ElapsedRAW,Cluster,SubmitLine " "--delimiter='|' -p --noheader" ) diff --git a/src/cloudai/systems/standalone/standalone_runner.py b/src/cloudai/systems/standalone/standalone_runner.py index 6d31304e1..f8a0df2a7 100644 --- a/src/cloudai/systems/standalone/standalone_runner.py +++ b/src/cloudai/systems/standalone/standalone_runner.py @@ -19,7 +19,6 @@ from pathlib import Path from typing import cast -import cloudai.metrics import cloudai.models.output import cloudai.output from cloudai.core import BaseJob, BaseRunner, JobIdRetrievalError, JobStatusResult, System, TestRun, TestScenario @@ -53,7 +52,7 @@ def get_run_output( if result.is_successful: try: observations = tr.test.metric_observations(self.system, tr) - metrics = [self._metric_output(observation) for observation in observations] + metrics = [cloudai.output.metric_output(observation) for observation in observations] except Exception as exc: logging.warning("Cannot extract output metrics for standalone job %s: %s", job.id, exc) return cloudai.models.output.Run( @@ -72,22 +71,6 @@ def on_job_completion(self, job: BaseJob) -> None: standalone_job = cast(StandaloneJob, job) standalone_job.finish = datetime.datetime.now(datetime.timezone.utc) - @staticmethod - def _metric_output(observation: cloudai.metrics.MetricObservation) -> cloudai.models.output.Metric: - dimensions = [ - cloudai.models.output.Dimension( - name=cloudai.metrics.dimension_label(key), - value=str(value), - ) - for key, value in sorted(observation.dimensions.items()) - ] - return cloudai.models.output.Metric( - name=observation.metric.display_name, - value=observation.value, - unit=observation.metric.unit, - dimensions=dimensions, - ) - def _submit_test(self, tr: TestRun) -> StandaloneJob: logging.info(f"Running test: {tr.name}") tr.output_path = self.get_job_output_path(tr) diff --git a/tests/systems/slurm/test_runner.py b/tests/systems/slurm/test_runner.py new file mode 100644 index 000000000..d9d565616 --- /dev/null +++ b/tests/systems/slurm/test_runner.py @@ -0,0 +1,151 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import datetime +import pathlib +from unittest import mock + +import pytest + +import cloudai.core +import cloudai.metrics +from cloudai.systems.slurm import SlurmJob, SlurmRunner, SlurmSystem +from cloudai.systems.slurm.slurm_metadata import SlurmStepMetadata + + +@pytest.mark.parametrize( + "successful,state,status", + [(True, "FAILED", "completed"), (False, "COMPLETED", "failed"), (False, "CANCELLED", "cancelled")], +) +@pytest.mark.parametrize("step", [0, 2]) +def test_slurm_run_output( + tmp_path: pathlib.Path, + base_tr: cloudai.core.TestRun, + slurm_system: SlurmSystem, + successful: bool, + state: str, + status: str, + step: int, +) -> None: + runner = SlurmRunner("run", slurm_system, cloudai.core.TestScenario(name="scenario", test_runs=[base_tr]), tmp_path) + base_tr.output_path.mkdir(parents=True) + base_tr.step = step + job = SlurmJob(base_tr, id=123) + runner.update_run_output(job) + metadata = SlurmStepMetadata( + job_id=123, + step_id="", + name="job", + state=state, + exit_code="1:0", + elapsed_time_sec=3, + start_time="2026-01-02T03:04:05Z", + end_time="2026-01-02T03:04:08Z", + submit_line="sbatch run.sh", + cluster_name="actual-cluster", + ) + observation = cloudai.metrics.MetricObservation( + cloudai.metrics.BANDWIDTH, + 12.5, + {"size_bytes": 1024, "bandwidth_basis": "bus"}, + ) + with ( + mock.patch.object(SlurmSystem, "get_job_status", return_value=[metadata]) as get_metadata, + mock.patch.object( + runner, + "get_cmd_gen_strategy", + return_value=mock.Mock(gen_srun_command=lambda: "srun cmd", generate_test_command=lambda: ["cmd"]), + ), + mock.patch.object( + cloudai.core.TestDefinition, + "was_run_successful", + return_value=cloudai.core.JobStatusResult(is_successful=successful), + ), + mock.patch.object( + cloudai.core.TestDefinition, "metric_observations", return_value=[observation] + ) as get_metrics, + ): + runner.store_job_metadata(job) + runner.update_run_output(job, runner.get_job_status(job)) + get_metadata.assert_called_once_with(job) + assert get_metrics.call_count == int(successful) + + base_tr.step = 3 + runner.shutting_down = status == "cancelled" + runner.finish_output(successful=True) + experiment = runner.experiment_output.snapshot() + assert experiment.status == status + assert experiment.system_name == "actual-cluster" + test = experiment.tests[0] + assert test.metrics == (test.runs[0].metrics if successful and step == 0 else []) + start = datetime.datetime(2026, 1, 2, 3, 4, 5, tzinfo=datetime.timezone.utc) + assert [run.model_dump() for run in test.runs] == [ + { + "path": str(base_tr.output_path.absolute()), + "jobid": "123", + "status": status, + "metrics": [ + { + "name": "Bandwidth", + "value": 12.5, + "unit": "GB/s", + "dimensions": [ + {"name": "Bandwidth basis", "value": "bus", "unit": "", "is_x": False}, + {"name": "Size", "value": "1024", "unit": "", "is_x": True}, + ], + } + ] + if successful + else [], + "start": start, + "finish": start + datetime.timedelta(seconds=3), + "duration": 3, + "iteration": 0, + "step": step, + } + ] + + +@pytest.mark.parametrize("timestamp", ["Unknown", "", "2026-01-02T03:04:05"]) +def test_slurm_output_unknown_timing_and_metric_failure( + tmp_path: pathlib.Path, + base_tr: cloudai.core.TestRun, + slurm_system: SlurmSystem, + timestamp: str, + caplog: pytest.LogCaptureFixture, +) -> None: + runner = SlurmRunner("run", slurm_system, cloudai.core.TestScenario(name="scenario", test_runs=[base_tr]), tmp_path) + job = SlurmJob(base_tr, id=123) + with mock.patch.object(SlurmSystem, "get_job_status", return_value=[]): + runner.store_job_metadata(job) + with mock.patch.object( + cloudai.core.TestDefinition, "metric_observations", side_effect=ValueError("broken metrics") + ): + run = runner.get_run_output(job, base_tr, cloudai.core.JobStatusResult(is_successful=True)) + assert run is not None + assert run.model_dump() == { + "path": str(base_tr.output_path.absolute()), + "jobid": "123", + "status": "completed", + "metrics": [], + "start": None, + "finish": None, + "duration": None, + "iteration": 0, + "step": 0, + } + assert runner._output_timestamp(timestamp) is None + assert "broken metrics" in caplog.text diff --git a/tests/systems/slurm/test_system.py b/tests/systems/slurm/test_system.py index 561f742ef..d3774c4e7 100644 --- a/tests/systems/slurm/test_system.py +++ b/tests/systems/slurm/test_system.py @@ -139,6 +139,7 @@ def test_slurm_api_rejects_non_sbatch_launcher(rest_slurm_system: SlurmSystem): def test_slurm_api_job_lifecycle(rest_slurm_system: SlurmSystem): response = { + "meta": {"slurm": {"cluster": "rest-cluster"}}, "jobs": [ { "job_id": 42, @@ -149,7 +150,7 @@ def test_slurm_api_job_lifecycle(rest_slurm_system: SlurmSystem): "end_time": 120, "nodes": "node[01-02]", } - ] + ], } job = SlurmJob(test_run=Mock(), id=42) @@ -163,6 +164,7 @@ def test_slurm_api_job_lifecycle(rest_slurm_system: SlurmSystem): assert [(item.step_id, item.name, item.state) for item in metadata] == [("", "rest-test", "FAILED")] assert metadata[0].exit_code == "1:0" assert metadata[0].elapsed_time_sec == 20 + assert metadata[0].cluster_name == "rest-cluster" def test_slurm_api_nodes_cancel_and_validation(rest_slurm_system: SlurmSystem): @@ -902,30 +904,42 @@ def test_get_job_status(slurm_system: SlurmSystem, stdout: str, stderr: str, exp slurm_system.get_job_status(job) else: assert slurm_system.get_job_status(job) == expected + slurm_system.cmd_shell.execute.assert_called_with( + "TZ=UTC SLURM_TIME_FORMAT='%Y-%m-%dT%H:%M:%SZ' " + "sacct -j 1 --format=JobID,JobName,State,ExitCode,Start,End,ElapsedRAW,Cluster,SubmitLine " + "--delimiter='|' -p --noheader" + ) -sacct_output = """2623913,job,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,sbatch sbatch_script.sh, -2623913.batch,batch,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,, -2623913.extern,extern,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,, -2623913.0,bash,COMPLETED,0:0,2025-05-09T01:35:24,2025-05-09T01:35:58,34,srun --export=ALL --mpi=pmix ..., -2623913.1,bash,COMPLETED,0:0,2025-05-09T01:35:58,2025-05-09T01:36:16,18,srun --export=ALL --mpi=pmix ..., -2623913.2,all_reduce_perf_mpi,COMPLETED,0:0,2025-05-09T01:36:16,2025-05-09T01:37:02,46,srun -N2 ..., -2623913.3,all_reduce_perf_mpi,COMPLETED,0:0,2025-05-09T01:37:02,2025-05-09T01:37:59,57,srun -N2 ..., +sacct_output = ( + "2623913,job,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,cluster," + "sbatch sbatch_script.sh,\n" + """2623913.batch,batch,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,cluster,, +2623913.extern,extern,COMPLETED,0:0,2025-05-09T01:34:52,2025-05-09T01:59:27,1475,cluster,, +2623913.0,bash,COMPLETED,0:0,2025-05-09T01:35:24,2025-05-09T01:35:58,34,cluster,srun --export=ALL --mpi=pmix ..., +2623913.1,bash,COMPLETED,0:0,2025-05-09T01:35:58,2025-05-09T01:36:16,18,cluster,srun --export=ALL --mpi=pmix ..., +2623913.2,all_reduce_perf_mpi,COMPLETED,0:0,2025-05-09T01:36:16,2025-05-09T01:37:02,46,cluster,srun -N2 ..., +2623913.3,all_reduce_perf_mpi,COMPLETED,0:0,2025-05-09T01:37:02,2025-05-09T01:37:59,57,cluster,srun -N2 ..., """ -sacct_output2 = """2968718|job:run|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|sbatch run_submission.sh| -2968718.batch|batch|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|| -2968718.extern|extern|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|| -2968718.0|bash|COMPLETED|0:0|2025-06-16T07:40:54|2025-06-16T07:49:11|497|srun long cmd +) +sacct_output2 = ( + "2968718|job:run|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|cluster|" + "sbatch run_submission.sh|\n" + """2968718.batch|batch|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|cluster|| +2968718.extern|extern|COMPLETED|0:0|2025-06-16T07:40:16|2025-06-16T07:49:08|532|cluster|| +2968718.0|bash|COMPLETED|0:0|2025-06-16T07:40:54|2025-06-16T07:49:11|497|cluster|srun long cmd with multiple lines | """ +) @pytest.mark.parametrize("sacct_output,delimiter,expected_nsteps", [(sacct_output, ",", 7), (sacct_output2, "|", 4)]) def test_slurm_job_metadata_from_sacct_output(sacct_output: str, delimiter: str, expected_nsteps: int): job_metadata = SlurmStepMetadata.from_sacct_output(sacct_output, delimiter=delimiter) assert len(job_metadata) == expected_nsteps + assert {item.cluster_name for item in job_metadata} == {"cluster"} @pytest.mark.parametrize( diff --git a/tests/systems/standalone/test_runner.py b/tests/systems/standalone/test_runner.py index 8a3e437d7..91e408d16 100644 --- a/tests/systems/standalone/test_runner.py +++ b/tests/systems/standalone/test_runner.py @@ -76,7 +76,7 @@ def test_standalone_run_output_uses_workload_status_and_metrics( "name": "Bandwidth", "value": 12.5, "unit": "GB/s", - "dimensions": [{"name": "Size", "value": "1024", "unit": "", "is_x": False}], + "dimensions": [{"name": "Size", "value": "1024", "unit": "", "is_x": True}], } ], "start": start, diff --git a/tests/test_handlers.py b/tests/test_handlers.py index 560b2f409..326df1a96 100644 --- a/tests/test_handlers.py +++ b/tests/test_handlers.py @@ -32,6 +32,7 @@ verify_test_configs, verify_test_scenarios, ) +from cloudai.configurator import CloudAIGymEnv from cloudai.configurator.env_params import EnvParamSpec from cloudai.core import ( BaseAgent, @@ -387,14 +388,18 @@ def test_handle_dse_job_invokes_agent_run( slurm_system: SlurmSystem, dse_tr: TestRun, custom_run_agent_name: str, + monkeypatch: pytest.MonkeyPatch, ) -> None: """``handle_dse_job`` must delegate orchestration to ``agent.run()`` (polymorphism).""" dse_tr.test.agent = custom_run_agent_name test_scenario = TestScenario(name="test_scenario", test_runs=[dse_tr]) runner = Runner(mode="dry-run", system=slurm_system, test_scenario=test_scenario) + update_output = MagicMock() + monkeypatch.setattr(CloudAIGymEnv, "update_output", update_output) assert handle_dse_job(runner, argparse.Namespace(mode="dry-run")) == 0 assert CustomRunStubAgent.run_calls == 1 + assert update_output.call_count == 2 def test_handle_dse_job_propagates_agent_run_nonzero_rc( @@ -438,6 +443,7 @@ def test_handle_dse_job_propagates_agent_run_exception( slurm_system: SlurmSystem, dse_tr: TestRun, custom_run_agent_name: str, + monkeypatch: pytest.MonkeyPatch, ) -> None: """Hard failure: an exception out of ``agent.run()`` propagates instead of being swallowed. @@ -448,10 +454,13 @@ def test_handle_dse_job_propagates_agent_run_exception( dse_tr.test.agent = custom_run_agent_name test_scenario = TestScenario(name="test_scenario", test_runs=[dse_tr]) runner = Runner(mode="dry-run", system=slurm_system, test_scenario=test_scenario) + update_output = MagicMock() + monkeypatch.setattr(CloudAIGymEnv, "update_output", update_output) with pytest.raises(RuntimeError, match="agent blew up"): handle_dse_job(runner, argparse.Namespace(mode="dry-run")) assert CustomRunStubAgent.run_calls == 1 + assert update_output.call_count == 2 def test_handle_dse_job_hard_fail_aborts_remaining_runs( diff --git a/tests/test_output.py b/tests/test_output.py index dae5cce65..22bfd3c33 100644 --- a/tests/test_output.py +++ b/tests/test_output.py @@ -17,11 +17,16 @@ import datetime import pathlib +import pytest + import cloudai.models.output import cloudai.output -def test_experiment_output_preserves_runs_and_finalizes_failure(tmp_path: pathlib.Path) -> None: +@pytest.mark.parametrize("status", ["completed", "failed", "cancelled"]) +def test_experiment_output_preserves_runs_and_finalizes_failure( + tmp_path: pathlib.Path, status: cloudai.models.output.Status +) -> None: start = datetime.datetime(2026, 1, 2, 3, 4, 5, tzinfo=datetime.timezone.utc) experiment = cloudai.models.output.Experiment( id="experiment", @@ -30,35 +35,49 @@ def test_experiment_output_preserves_runs_and_finalizes_failure(tmp_path: pathli status="running", path=str(tmp_path), start=start, - tests=[cloudai.models.output.Test(id="case", name="workload", path=str(tmp_path / "case"))], + tests=[ + cloudai.models.output.Test(id=case, name="workload", path=str(tmp_path / case)) + for case in ("case", "interrupted", "not-started") + ], ) experiment_output = cloudai.output.ExperimentOutput(experiment, tmp_path) first_run = cloudai.models.output.Run( path=str(tmp_path / "case" / "0"), jobid="101", status="completed", + metrics=[cloudai.models.output.Metric(name="Bandwidth", value=12.5, unit="GB/s")], start=start, finish=start + datetime.timedelta(seconds=2), duration=2, iteration=0, - step=0, + step=1, ) second_run = cloudai.models.output.Run( path=str(tmp_path / "case" / "1"), jobid="102", - status="running", + status="pending", start=start + datetime.timedelta(seconds=2), iteration=1, - step=0, + step=2, ) experiment_output.update_run("case", first_run) experiment_output.update_run("case", second_run) + assert experiment_output.snapshot().tests[0].status == "pending" second_run.status = "failed" second_run.finish = start + datetime.timedelta(seconds=4) second_run.duration = 2 experiment_output.update_run("case", second_run) - experiment_output.finish("failed", start + datetime.timedelta(seconds=5)) + experiment_output.update_dse( + "case", + {"algorithm": ["first", "second"]}, + [(2, {"algorithm": "second"}), (1, {"algorithm": "first"})], + ) + experiment_output.update_run( + "interrupted", + cloudai.models.output.Run(path=str(tmp_path / "interrupted" / "0"), jobid="103", status="pending"), + ) + experiment_output.finish(status, start + datetime.timedelta(seconds=5)) stored = cloudai.models.output.Experiment.model_validate_json((tmp_path / "experiment.json").read_text()) assert stored.model_dump() == { @@ -66,7 +85,7 @@ def test_experiment_output_preserves_runs_and_finalizes_failure(tmp_path: pathli "name": "scenario", "system_name": "test-system", "description": None, - "status": "failed", + "status": "cancelled" if status == "cancelled" else "failed", "path": str(tmp_path), "start": start, "finish": start + datetime.timedelta(seconds=5), @@ -78,18 +97,18 @@ def test_experiment_output_preserves_runs_and_finalizes_failure(tmp_path: pathli "description": None, "status": "failed", "path": str(tmp_path / "case"), - "metrics": [], + "metrics": [{"name": "Bandwidth", "value": 12.5, "unit": "GB/s", "dimensions": []}], "runs": [ { "path": str(tmp_path / "case" / "0"), "jobid": "101", "status": "completed", - "metrics": [], + "metrics": [{"name": "Bandwidth", "value": 12.5, "unit": "GB/s", "dimensions": []}], "start": start, "finish": start + datetime.timedelta(seconds=2), "duration": 2, "iteration": 0, - "step": 0, + "step": 1, }, { "path": str(tmp_path / "case" / "1"), @@ -100,10 +119,46 @@ def test_experiment_output_preserves_runs_and_finalizes_failure(tmp_path: pathli "finish": start + datetime.timedelta(seconds=4), "duration": 2, "iteration": 1, - "step": 0, + "step": 2, }, ], + "dse": { + "space": {"algorithm": ["first", "second"]}, + "best_config": {"algorithm": "first"}, + "best_step": 1, + }, + }, + { + "id": "interrupted", + "name": "workload", + "description": None, + "status": "unknown", + "path": str(tmp_path / "interrupted"), + "metrics": [], + "runs": [ + { + "path": str(tmp_path / "interrupted" / "0"), + "jobid": "103", + "status": "unknown", + "metrics": [], + "start": None, + "finish": None, + "duration": None, + "iteration": None, + "step": None, + } + ], + "dse": None, + }, + { + "id": "not-started", + "name": "workload", + "description": None, + "status": "unknown", + "path": str(tmp_path / "not-started"), + "metrics": [], + "runs": [], "dse": None, - } + }, ], }