Skip to content
Merged
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
118 changes: 118 additions & 0 deletions doc/examples/nccl-dse/experiment.json
Original file line number Diff line number Diff line change
@@ -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
}
}
]
}
36 changes: 28 additions & 8 deletions doc/reporting.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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 <examples/nccl-dse/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:
Expand Down
4 changes: 3 additions & 1 deletion src/cloudai/_core/base_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Comment thread
podkidyshev marked this conversation as resolved.
status = "cancelled"
self.experiment_output.finish(status=status, finish=datetime.datetime.now(datetime.timezone.utc))

def shutdown(self):
Expand Down
6 changes: 5 additions & 1 deletion src/cloudai/cli/handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.")
Expand Down
29 changes: 29 additions & 0 deletions src/cloudai/configurator/cloudai_gym.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@

import copy
import logging
import math
from pathlib import Path
from typing import TYPE_CHECKING, Any, Dict, Optional, Tuple, cast

Expand Down Expand Up @@ -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:
"""
Expand Down
2 changes: 1 addition & 1 deletion src/cloudai/models/output.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
62 changes: 56 additions & 6 deletions src/cloudai/output.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import pathlib
import tempfile

import cloudai.metrics
import cloudai.models.output


Expand All @@ -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,
Comment thread
podkidyshev marked this conversation as resolved.
)
for key, value in sorted(observation.dimensions.items())
],
)


class ExperimentOutput:
"""Collect experiment results and publish snapshots."""

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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"):
Comment thread
podkidyshev marked this conversation as resolved.
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
Expand All @@ -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"
5 changes: 4 additions & 1 deletion src/cloudai/systems/slurm/single_sbatch_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Comment thread
podkidyshev marked this conversation as resolved.

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())
Expand Down
3 changes: 3 additions & 0 deletions src/cloudai/systems/slurm/slurm_job.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Loading
Loading