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
113 changes: 113 additions & 0 deletions doc/reporting.rst
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ This chapter describes the reporting system in CloudAI. In this chapter, we will
- :ref:`Enabling, Disabling and Configuring Reports <enabling-disabling-and-configuring-reports>`
- :ref:`Reporting Registration <reporting-registration>`
- :ref:`Reporting Configuration Implementation <reporting-configuration-implementation>`
- :ref:`Uploading Results to Object Storage <uploading-results-to-object-storage>`
- :doc:`Reports <reports>`

.. toctree::
Expand Down Expand Up @@ -183,3 +184,115 @@ And it can be used in a test scenario as follows:

[reports]
custom = { enable = true, greeting = "Hello, world!" }

.. _uploading-results-to-object-storage:

Uploading Results to Object Storage
------------------------------------

The ``s3`` scenario report publishes the scenario results directory to an
S3-compatible bucket. It is disabled by default, because shipping results off-box should
be a deliberate choice.

It is registered last, after ``tarball``, so it always observes the complete results
directory including every other report's output.

Install the optional dependency first:

.. code-block:: bash

pip install 'cloudai[s3]'

Then enable it in a test scenario:

.. code-block:: toml

[reports]
s3 = { enable = true, bucket = "my-bucket", prefix = "cloudai/runs", upload_tarball = true }

Or, for Slurm systems, once per cluster in the system config:

.. code-block:: toml

[reports]
s3 = { enable = true, bucket = "my-bucket" }

Configuration options:

.. list-table::
:header-rows: 1

* - Option
- Default
- Description
* - ``bucket``
- ``$CLOUDAI_S3_BUCKET``
- Destination bucket. Required; the upload is skipped with a warning if unset.
* - ``prefix``
- ``$CLOUDAI_S3_PREFIX``
- Key prefix. Objects are written under ``<prefix>/<system_name>/<results_dir_name>/``.
* - ``endpoint_url``
- ``$CLOUDAI_S3_ENDPOINT_URL``
- Custom endpoint, for MinIO or other S3-compatible stores.
* - ``region``
- unset
- AWS region. When unset, boto3 resolves it (e.g. ``AWS_DEFAULT_REGION``).
* - ``upload_tree``
- ``true``
- Upload each file individually, preserving relative paths.
* - ``upload_tarball``
- ``false``
- Also upload a ``.tgz`` of the whole directory. It is created if absent and rebuilt if older than the results.
* - ``upload_concurrency``
- ``8``
- Number of files uploaded concurrently when ``upload_tree`` is enabled. Must be at least 1.

At least one of ``upload_tree`` and ``upload_tarball`` must be enabled; the config is rejected otherwise.

Environment variables
~~~~~~~~~~~~~~~~~~~~~

The destination can be supplied through three environment variables, so a cluster-wide
default does not have to be repeated in every scenario:

.. list-table::
:header-rows: 1

* - Variable
- Sets option
- Notes
* - ``CLOUDAI_S3_BUCKET``
- ``bucket``
- Destination bucket.
* - ``CLOUDAI_S3_PREFIX``
- ``prefix``
- Key prefix. Empty by default.
* - ``CLOUDAI_S3_ENDPOINT_URL``
- ``endpoint_url``
- Custom endpoint, for MinIO or other S3-compatible stores. An empty value is treated as unset.

A value set in TOML takes precedence over the environment variable. The variables are read by
the ``cloudai`` process when the report configuration is loaded, so export them where
``cloudai`` runs:

.. code-block:: bash

export CLOUDAI_S3_BUCKET=my-bucket
export CLOUDAI_S3_PREFIX=cloudai/runs
export CLOUDAI_S3_ENDPOINT_URL=http://localhost:9000 # only for MinIO or other S3-compatible stores

With these set, enabling the report only needs ``s3 = { enable = true }``.

**Credentials are never read from CloudAI configuration.** They are resolved by boto3's
standard chain: ``AWS_ACCESS_KEY_ID``/``AWS_SECRET_ACCESS_KEY``, ``~/.aws/credentials``,
or an instance/IAM role.

Because reports run inside a ``try``/``except``, an upload failure logs a warning and
leaves the run's exit status unchanged.

To upload a results directory from an earlier run, re-run the reports against it:

.. code-block:: bash

cloudai generate-report --system-config <system.toml> --tests-dir <tests/> \
--test-scenario <scenario.toml> --result-dir results/<scenario>_<timestamp>
2 changes: 2 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -61,8 +61,10 @@ classifiers = [
"pytest-deadfixtures~=3.1",
"taplo~=0.9.3",
"gymnasium~=1.2",
"boto3~=1.40",
]
rl = ["gymnasium~=1.2"]
s3 = ["boto3~=1.40"]
Comment thread
jj10306 marked this conversation as resolved.
docs = [
"sphinx~=8.1",
"nvidia-sphinx-theme~=0.0.8",
Expand Down
3 changes: 2 additions & 1 deletion src/cloudai/_core/registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -278,7 +278,8 @@ def report_order(k: str) -> int:
"per_test": 0, # first
"status": 2,
"dse": 3,
"tarball": 4, # last
"tarball": 4,
"s3": 5, # last, must observe every other report's output
}.get(k, 1)

return sorted(self.scenario_reports.items(), key=lambda kv: report_order(kv[0]))
Expand Down
7 changes: 7 additions & 0 deletions src/cloudai/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -71,8 +71,10 @@
from .models.workload import CmdArgs, NsysConfiguration, PredictorConfig, TestDefinition
from .parser import Parser
from .reporter import JUnitReporter, PerTestReporter, StatusReporter, TarballReporter
from .s3_reporter import S3UploadConfig, S3UploadReporter
from .test_parser import TestParser
from .test_scenario_parser import TestScenarioParser
from .util.object_store import ObjectStore, S3ObjectStore, UploadStats

__all__ = [
"METRIC_ERROR",
Expand Down Expand Up @@ -107,6 +109,7 @@
"MetricValue",
"MissingTestError",
"NsysConfiguration",
"ObjectStore",
"ObsLeafDescriptor",
"Parser",
"PerTestReporter",
Expand All @@ -118,6 +121,9 @@
"Reporter",
"RewardOverrides",
"Runner",
"S3ObjectStore",
"S3UploadConfig",
"S3UploadReporter",
"StatusReporter",
"StructuredObservationProducer",
"System",
Expand All @@ -131,6 +137,7 @@
"TestScenario",
"TestScenarioParser",
"TestScenarioParsingError",
"UploadStats",
"case_name",
"format_validation_error",
]
2 changes: 2 additions & 0 deletions src/cloudai/registration.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ def register_all():
from cloudai.models.scenario import ReportConfig
from cloudai.report_generator.training import TrainingReporter
from cloudai.reporter import DSEReporter, JUnitReporter, PerTestReporter, StatusReporter, TarballReporter
from cloudai.s3_reporter import S3UploadConfig, S3UploadReporter

# Import systems
from cloudai.systems.kubernetes import KubernetesInstaller, KubernetesRunner, KubernetesSystem
Expand Down Expand Up @@ -343,6 +344,7 @@ def register_all():
Registry().add_scenario_report("junit", JUnitReporter, ReportConfig(enable=False))
Registry().add_scenario_report("dse", DSEReporter, ReportConfig(enable=True))
Registry().add_scenario_report("tarball", TarballReporter, ReportConfig(enable=True))
Registry().add_scenario_report("s3", S3UploadReporter, S3UploadConfig(enable=False))
Registry().add_scenario_report(
"nixl_bench_summary",
NIXLBenchComparisonReport,
Expand Down
125 changes: 125 additions & 0 deletions src/cloudai/s3_reporter.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
# 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 logging
import os
from pathlib import Path
from typing import Optional

from pydantic import Field, model_validator
from typing_extensions import Self

from .core import Reporter
from .models.scenario import ReportConfig
from .reporter import TarballReporter
from .util.object_store import S3ObjectStore, join_key


class S3UploadConfig(ReportConfig):
"""
Configuration for uploading a scenario results directory to object storage.

Destination fields fall back to environment variables when not set in TOML, so a
cluster-wide destination can be supplied by the environment while a scenario can
still override it. Credentials are never read from here; boto3 resolves them from
its standard chain.
"""

bucket: str = Field(default_factory=lambda: os.getenv("CLOUDAI_S3_BUCKET", ""))
prefix: str = Field(default_factory=lambda: os.getenv("CLOUDAI_S3_PREFIX", ""))
endpoint_url: Optional[str] = Field(default_factory=lambda: os.getenv("CLOUDAI_S3_ENDPOINT_URL") or None)
Comment thread
jj10306 marked this conversation as resolved.
region: Optional[str] = None
upload_tree: bool = True
upload_tarball: bool = False
upload_concurrency: int = Field(default=8, ge=1)

@model_validator(mode="after")
def at_least_one_upload_mode(self) -> Self:
if not (self.upload_tree or self.upload_tarball):
raise ValueError("At least one of 'upload_tree' or 'upload_tarball' must be enabled.")
return self


class S3UploadReporter(Reporter):
"""Uploads the scenario results directory to object storage."""

def generate(self) -> None:
config = self.config
if not isinstance(config, S3UploadConfig):
logging.warning(f"Expected S3UploadConfig, got {type(config).__name__}, skipping results upload.")
return

if not config.bucket:
logging.warning(
"S3 upload is enabled but no bucket is configured. "
"Set 'bucket' in the report config or the CLOUDAI_S3_BUCKET environment variable."
)
return

if not self.results_root.exists():
logging.warning(f"Results directory {self.results_root} does not exist, skipping results upload.")
return

store = S3ObjectStore(bucket=config.bucket, endpoint_url=config.endpoint_url, region=config.region)
if not store.bucket_exists():
logging.warning(f"Bucket '{config.bucket}' does not exist, skipping results upload.")
return

key_prefix = join_key(config.prefix, self.system.name, self.results_root.name)

if config.upload_tree:
stats = store.upload_directory(self.results_root, key_prefix, max_workers=config.upload_concurrency)
logging.info(
f"Uploaded {stats.files_uploaded} file(s), {stats.bytes_uploaded} byte(s) to {store.uri(key_prefix)} "
f"in {stats.duration_seconds:.2f}s"
)
if stats.failures:
logging.warning(
f"Failed to upload {len(stats.failures)} file(s) to {store.uri(key_prefix)}, "
"see debug log for details"
)

if config.upload_tarball:
self.upload_tarball(store, key_prefix)

def upload_tarball(self, store: S3ObjectStore, key_prefix: str) -> None:
"""
Upload a tarball of the results directory, (re)creating it if it is missing or stale.

TarballReporter only produces a tarball when a test run failed, so we cannot
assume one is already present. A leftover tarball from an earlier run may predate
regenerated reports, so it is reused only if nothing in the directory is newer.
"""
tarball_path = Path(str(self.results_root) + ".tgz")
if not tarball_path.exists() or self._tarball_is_stale(tarball_path):
TarballReporter(self.system, self.test_scenario, self.results_root, self.config).create_tarball(
self.results_root
)

key = join_key(key_prefix, tarball_path.name)
try:
store.upload_file(tarball_path, key)
except Exception as e:
logging.warning(f"Failed to upload tarball to {store.uri(key)}, see debug log for details")
logging.debug(e, exc_info=True)
return
logging.info(f"Uploaded tarball to {store.uri(key)}")

def _tarball_is_stale(self, tarball_path: Path) -> bool:
"""Whether anything in the results directory was modified after the tarball was written."""
tarball_mtime = tarball_path.stat().st_mtime
entries = [self.results_root, *self.results_root.rglob("*")]
return any(entry.stat().st_mtime > tarball_mtime for entry in entries)
15 changes: 15 additions & 0 deletions src/cloudai/util/lazy_imports.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import bokeh.palettes as bokeh_pallettes
import bokeh.plotting as bokeh_plotting
import bokeh.transform as bokeh_transform
import boto3
import gymnasium
import kubernetes as k8s
import numpy as np
Expand All @@ -48,6 +49,20 @@ def __init__(self):
self._bokeh_transform: ModuleType | None = None
self._bokeh_pallettes: ModuleType | None = None
self._bokeh_embed: ModuleType | None = None
self._boto3: ModuleType | None = None

@property
def boto3(self) -> boto3: # type: ignore[no-any-return]
"""Lazy import of boto3 (optional ``cloudai[s3]`` extra)."""
if self._boto3 is None:
try:
import boto3
except ImportError as exc:
raise ImportError(
"boto3 is required for S3 object storage. Install it with: pip install 'cloudai[s3]'"
) from exc
self._boto3 = boto3
return cast("boto3", self._boto3)

@property
def np(self) -> np: # type: ignore[no-any-return]
Expand Down
Loading
Loading