diff --git a/doc/reporting.rst b/doc/reporting.rst index 386f6d72c..0d72b937c 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -8,6 +8,7 @@ This chapter describes the reporting system in CloudAI. In this chapter, we will - :ref:`Enabling, Disabling and Configuring Reports ` - :ref:`Reporting Registration ` - :ref:`Reporting Configuration Implementation ` +- :ref:`Uploading Results to Object Storage ` - :doc:`Reports ` .. toctree:: @@ -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 ``///``. + * - ``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 --tests-dir \ + --test-scenario --result-dir results/_ diff --git a/pyproject.toml b/pyproject.toml index cb643ee67..0df35be34 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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"] docs = [ "sphinx~=8.1", "nvidia-sphinx-theme~=0.0.8", diff --git a/src/cloudai/_core/registry.py b/src/cloudai/_core/registry.py index b0ccf5a75..5485733bd 100644 --- a/src/cloudai/_core/registry.py +++ b/src/cloudai/_core/registry.py @@ -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])) diff --git a/src/cloudai/core.py b/src/cloudai/core.py index 2ebd13c8c..2a6e1dc58 100644 --- a/src/cloudai/core.py +++ b/src/cloudai/core.py @@ -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", @@ -107,6 +109,7 @@ "MetricValue", "MissingTestError", "NsysConfiguration", + "ObjectStore", "ObsLeafDescriptor", "Parser", "PerTestReporter", @@ -118,6 +121,9 @@ "Reporter", "RewardOverrides", "Runner", + "S3ObjectStore", + "S3UploadConfig", + "S3UploadReporter", "StatusReporter", "StructuredObservationProducer", "System", @@ -131,6 +137,7 @@ "TestScenario", "TestScenarioParser", "TestScenarioParsingError", + "UploadStats", "case_name", "format_validation_error", ] diff --git a/src/cloudai/registration.py b/src/cloudai/registration.py index f471a5fcd..32961adce 100644 --- a/src/cloudai/registration.py +++ b/src/cloudai/registration.py @@ -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 @@ -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, diff --git a/src/cloudai/s3_reporter.py b/src/cloudai/s3_reporter.py new file mode 100644 index 000000000..8570ffb44 --- /dev/null +++ b/src/cloudai/s3_reporter.py @@ -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) + 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) diff --git a/src/cloudai/util/lazy_imports.py b/src/cloudai/util/lazy_imports.py index def790dfc..6b92ca391 100644 --- a/src/cloudai/util/lazy_imports.py +++ b/src/cloudai/util/lazy_imports.py @@ -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 @@ -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] diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py new file mode 100644 index 000000000..5a7706398 --- /dev/null +++ b/src/cloudai/util/object_store.py @@ -0,0 +1,192 @@ +# 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. + +from __future__ import annotations + +import fnmatch +import logging +import time +from abc import ABC, abstractmethod +from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass, field +from pathlib import Path +from typing import Optional + +from .lazy_imports import lazy + + +def join_key(*parts: str) -> str: + """Join object key parts with '/', dropping empties and collapsing separators.""" + cleaned = [p.strip("/") for p in parts if p and p.strip("/")] + return "/".join(cleaned) + + +@dataclass +class UploadStats: + """Summary of an upload operation.""" + + files_uploaded: int = 0 + bytes_uploaded: int = 0 + failures: list[tuple[Path, str]] = field(default_factory=list) + duration_seconds: float = 0.0 + + @property + def is_successful(self) -> bool: + return not self.failures + + +class ObjectStore(ABC): + """Minimal object storage interface used to publish CloudAI artifacts.""" + + @abstractmethod + def uri(self, key: str) -> str: + """Return a human-readable URI for the given key, for logging.""" + ... + + @abstractmethod + def upload_file(self, local_path: Path, key: str) -> None: + """Upload a single file to the given key.""" + ... + + def upload_directory( + self, + local_dir: Path, + key_prefix: str = "", + exclude: Optional[list[str]] = None, + max_workers: int = 8, + ) -> UploadStats: + """ + Upload every file under ``local_dir``, preserving relative paths. + + Args: + local_dir: Directory to walk. + key_prefix: Key prefix to place the tree under. + exclude: Glob patterns matched against paths relative to ``local_dir``. + Matching files are skipped. + max_workers: Number of files to upload concurrently. Bounded by the actual + number of files, so small trees don't spin up idle threads. + + Returns: + Stats describing what was uploaded, including wall-clock time spent + uploading (excluding the directory walk). Individual file failures are + collected rather than raised, so a single bad file does not abort the + whole upload. + """ + stats = UploadStats() + exclude = exclude or [] + + to_upload: list[tuple[Path, str]] = [] + for path in sorted(local_dir.rglob("*")): + if not path.is_file(): + continue + + relative = path.relative_to(local_dir) + if any(fnmatch.fnmatch(str(relative), pattern) for pattern in exclude): + logging.debug(f"Skipping excluded file {relative}") + continue + + to_upload.append((path, join_key(key_prefix, relative.as_posix()))) + + if not to_upload: + return stats + + start = time.perf_counter() + workers = min(max_workers, len(to_upload)) + with ThreadPoolExecutor(max_workers=workers) as executor: + futures = {executor.submit(self._upload_one, path, key): (path, key) for path, key in to_upload} + for future in as_completed(futures): + path, key = futures[future] + try: + size = future.result() + except Exception as e: + logging.debug(f"Failed to upload {path} to {self.uri(key)}: {e}", exc_info=True) + stats.failures.append((path, str(e))) + continue + + stats.files_uploaded += 1 + stats.bytes_uploaded += size + stats.duration_seconds = time.perf_counter() - start + + return stats + + def _upload_one(self, path: Path, key: str) -> int: + size = path.stat().st_size + self.upload_file(path, key) + return size + + +class S3ObjectStore(ObjectStore): + """ + S3-backed object store. + + Credentials are resolved by boto3's standard chain (``AWS_*`` environment variables, + ``~/.aws/credentials``, instance/IAM roles), so no secrets are read from CloudAI + configuration files. + """ + + def __init__( + self, + bucket: str, + endpoint_url: Optional[str] = None, + region: Optional[str] = None, + ) -> None: + self.bucket = bucket + self.endpoint_url = endpoint_url + self.region = region + self._client = None + + @property + def client(self): + if self._client is None: + self._client = lazy.boto3.client("s3", endpoint_url=self.endpoint_url, region_name=self.region) + return self._client + + def uri(self, key: str) -> str: + return f"s3://{join_key(self.bucket, key)}" + + def upload_file(self, local_path: Path, key: str) -> None: + self.client.upload_file(str(local_path), self.bucket, key) + + def exists(self, key: str) -> bool: + try: + self.client.head_object(Bucket=self.bucket, Key=key) + except self.client.exceptions.ClientError as e: + if e.response.get("ResponseMetadata", {}).get("HTTPStatusCode") == 404: + return False + raise + return True + + def bucket_exists(self) -> bool: + """ + Return False only when the configured bucket definitively does not exist (404). + + A 403 is not treated as missing: HeadBucket needs ``s3:ListBucket``, which + write-only credentials (``s3:PutObject`` only) commonly lack, even though + uploads would succeed. If access really is denied, the per-file uploads fail + and are reported. + """ + try: + self.client.head_bucket(Bucket=self.bucket) + except self.client.exceptions.ClientError as e: + status = e.response.get("ResponseMetadata", {}).get("HTTPStatusCode") + if status == 404: + logging.debug(f"Bucket '{self.bucket}' does not exist: {e}") + return False + if status == 403: + logging.debug(f"Cannot verify bucket '{self.bucket}' (access denied on HeadBucket), proceeding: {e}") + return True + raise + return True diff --git a/tests/test_init.py b/tests/test_init.py index 2943b4676..da198ab80 100644 --- a/tests/test_init.py +++ b/tests/test_init.py @@ -18,6 +18,7 @@ from cloudai.core import Registry from cloudai.report_generator.training import TrainingReporter from cloudai.reporter import DSEReporter, JUnitReporter, PerTestReporter, StatusReporter, TarballReporter +from cloudai.s3_reporter import S3UploadReporter from cloudai.systems.kubernetes import KubernetesInstaller, KubernetesSystem from cloudai.systems.lsf import LSFInstaller, LSFSystem from cloudai.systems.runai import RunAISystem @@ -287,6 +288,7 @@ def test_scenario_reports(): "junit", "dse", "tarball", + "s3", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -303,6 +305,7 @@ def test_scenario_reports(): JUnitReporter, DSEReporter, TarballReporter, + S3UploadReporter, NIXLBenchComparisonReport, NixlEPComparisonReport, NcclComparisonReport, @@ -323,6 +326,7 @@ def test_report_configs(): "junit", "dse", "tarball", + "s3", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -331,7 +335,7 @@ def test_report_configs(): "vllm_comparison", "sglang_comparison", ] - disabled_configs = {"junit"} + disabled_configs = {"junit", "s3"} # uploading off-box must be opt-in for name, rep_config in configs.items(): expected_be_enabled = name not in disabled_configs assert rep_config.enable is expected_be_enabled, f"Report {name} has an unexpected default state" diff --git a/tests/test_reporter.py b/tests/test_reporter.py index b6e3a0f65..d3bdad65e 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -30,7 +30,14 @@ from cloudai.handlers import generate_reports from cloudai.models.scenario import ReportConfig, TestRunDetails from cloudai.report_generator.dse_report import build_dse_summaries -from cloudai.reporter import DSEReporter, JUnitReporter, PerTestReporter, ReportItem, StatusReporter, TarballReporter +from cloudai.reporter import ( + DSEReporter, + JUnitReporter, + PerTestReporter, + ReportItem, + StatusReporter, + TarballReporter, +) from cloudai.systems.slurm.slurm_metadata import ( MetadataCUDA, MetadataMPI, @@ -440,9 +447,10 @@ def test_scenario_report_escapes_error_message( def test_report_order() -> None: reports = Registry().ordered_scenario_reports() assert reports[0][0] == "per_test" - assert reports[-3][0] == "status" - assert reports[-2][0] == "dse" - assert reports[-1][0] == "tarball" + assert reports[-4][0] == "status" + assert reports[-3][0] == "dse" + assert reports[-2][0] == "tarball" + assert reports[-1][0] == "s3" def test_junit_reporter_generates_testcases_with_status_logs_and_duration( diff --git a/tests/test_s3_reporter.py b/tests/test_s3_reporter.py new file mode 100644 index 000000000..223b21d74 --- /dev/null +++ b/tests/test_s3_reporter.py @@ -0,0 +1,195 @@ +# 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 os +import tarfile +from pathlib import Path +from typing import Any +from unittest.mock import patch + +import pytest +from pydantic import ValidationError + +from cloudai import TestScenario +from cloudai.core import Registry +from cloudai.s3_reporter import S3UploadConfig, S3UploadReporter +from cloudai.systems.slurm.slurm_system import SlurmSystem +from cloudai.util.object_store import UploadStats + + +class TestS3UploadReporter: + """Tests for uploading a results directory to object storage.""" + + @pytest.fixture + def results_dir(self, tmp_path: Path) -> Path: + results_dir = tmp_path / "nccl-test_2025-04-16_14-27-45" + (results_dir / "nccl" / "0").mkdir(parents=True) + (results_dir / "nccl" / "0" / "stdout.txt").write_text("out") + (results_dir / "report.html").write_text("") + return results_dir + + def reporter(self, slurm_system: SlurmSystem, results_dir: Path, **kwargs: Any) -> S3UploadReporter: + return S3UploadReporter( + slurm_system, + TestScenario(name="dummy", test_runs=[]), + results_dir, + S3UploadConfig(enable=True, **kwargs), + ) + + def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + store.upload_directory.return_value = UploadStats(files_uploaded=2, bytes_uploaded=16) + + self.reporter(slurm_system, results_dir, bucket="my-bucket", prefix="cloudai").generate() + + mock_store_cls.assert_called_once_with(bucket="my-bucket", endpoint_url=None, region=None) + store.upload_directory.assert_called_once_with( + results_dir, "cloudai/test_system/nccl-test_2025-04-16_14-27-45", max_workers=8 + ) + + def test_upload_concurrency_is_configurable(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + store.upload_directory.return_value = UploadStats() + + self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_concurrency=16).generate() + + _, kwargs = store.upload_directory.call_args + assert kwargs["max_workers"] == 16 + + def test_no_bucket_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + self.reporter(slurm_system, results_dir).generate() + + mock_store_cls.assert_not_called() + + def test_missing_results_dir_uploads_nothing(self, slurm_system: SlurmSystem, tmp_path: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + self.reporter(slurm_system, tmp_path / "nope", bucket="my-bucket").generate() + + mock_store_cls.assert_not_called() + + def test_missing_bucket_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + store.bucket_exists.return_value = False + + self.reporter(slurm_system, results_dir, bucket="my-bucket").generate() + + store.upload_directory.assert_not_called() + + def test_upload_tree_disabled(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + + self.reporter( + slurm_system, results_dir, bucket="my-bucket", upload_tree=False, upload_tarball=True + ).generate() + + store.upload_directory.assert_not_called() + store.upload_file.assert_called_once() + + def test_tarball_created_when_absent(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + tarball_path = Path(str(results_dir) + ".tgz") + assert not tarball_path.exists() + + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + store.upload_directory.return_value = UploadStats() + + self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_tarball=True).generate() + + assert tarball_path.exists(), "TarballReporter only tarballs on failure, so it must be created here" + store.upload_file.assert_called_once_with( + tarball_path, "test_system/nccl-test_2025-04-16_14-27-45/nccl-test_2025-04-16_14-27-45.tgz" + ) + + def test_fresh_tarball_is_reused(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + tarball_path = Path(str(results_dir) + ".tgz") + tarball_path.write_bytes(b"pre-existing") + self._set_tarball_newer_than_contents(results_dir, tarball_path) + + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + mock_store_cls.return_value.upload_directory.return_value = UploadStats() + + self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_tarball=True).generate() + + assert tarball_path.read_bytes() == b"pre-existing" + mock_store_cls.return_value.upload_file.assert_called_once() + + def test_stale_tarball_is_regenerated(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + """A tarball left by an earlier run must not be uploaded in place of the regenerated reports.""" + tarball_path = Path(str(results_dir) + ".tgz") + tarball_path.write_bytes(b"stale") + self._set_tarball_newer_than_contents(results_dir, tarball_path) + (results_dir / "report.html").write_text("regenerated") + os.utime(tarball_path, (1, 1)) # tarball predates the regenerated report + + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + mock_store_cls.return_value.upload_directory.return_value = UploadStats() + + self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_tarball=True).generate() + + assert tarball_path.read_bytes() != b"stale" + with tarfile.open(tarball_path) as tar: + member = tar.extractfile(f"{results_dir.name}/report.html") + assert member is not None + assert member.read() == b"regenerated" + + @staticmethod + def _set_tarball_newer_than_contents(results_dir: Path, tarball_path: Path) -> None: + newest = max(p.stat().st_mtime for p in [results_dir, *results_dir.rglob("*")]) + os.utime(tarball_path, (newest + 10, newest + 10)) + + def test_env_var_fallback(self, slurm_system: SlurmSystem, results_dir: Path, monkeypatch) -> None: + monkeypatch.setenv("CLOUDAI_S3_BUCKET", "env-bucket") + monkeypatch.setenv("CLOUDAI_S3_PREFIX", "env-prefix") + monkeypatch.setenv("CLOUDAI_S3_ENDPOINT_URL", "http://localhost:9000") + + config = S3UploadConfig(enable=True) + + assert config.bucket == "env-bucket" + assert config.prefix == "env-prefix" + assert config.endpoint_url == "http://localhost:9000" + + def test_toml_overrides_env_var(self, monkeypatch) -> None: + monkeypatch.setenv("CLOUDAI_S3_BUCKET", "env-bucket") + + assert S3UploadConfig(enable=True, bucket="toml-bucket").bucket == "toml-bucket" + + def test_upload_failure_does_not_raise(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.s3_reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + store.upload_directory.return_value = UploadStats(failures=[(results_dir / "report.html", "denied")]) + + self.reporter(slurm_system, results_dir, bucket="my-bucket").generate() + + def test_requires_at_least_one_upload_mode(self) -> None: + with pytest.raises(ValidationError, match="upload_tree"): + S3UploadConfig(enable=True, bucket="my-bucket", upload_tree=False, upload_tarball=False) + + @pytest.mark.parametrize("concurrency", [0, -1]) + def test_upload_concurrency_must_be_positive(self, concurrency: int) -> None: + with pytest.raises(ValidationError): + S3UploadConfig(enable=True, bucket="my-bucket", upload_concurrency=concurrency) + + +def test_s3_upload_runs_after_tarball() -> None: + order = [name for name, _ in Registry().ordered_scenario_reports()] + + assert order.index("s3") > order.index("tarball") + assert order[-1] == "s3" diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py new file mode 100644 index 000000000..7b7072fd4 --- /dev/null +++ b/tests/util/test_object_store.py @@ -0,0 +1,245 @@ +# 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. + +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pytest +from botocore.exceptions import ClientError + +from cloudai.util.object_store import ObjectStore, S3ObjectStore, UploadStats, join_key + + +def _client_error(status: int, code: str = "") -> ClientError: + return ClientError( + {"Error": {"Code": code or str(status)}, "ResponseMetadata": {"HTTPStatusCode": status}}, + "HeadObject", + ) + + +@pytest.fixture +def tree(tmp_path: Path) -> Path: + root = tmp_path / "scenario_2025-04-16_14-27-45" + (root / "test-a" / "0").mkdir(parents=True) + (root / "test-a" / "0" / "stdout.txt").write_text("out") + (root / "test-a" / "0" / "stderr.txt").write_text("err!!") + (root / "report.html").write_text("") + return root + + +class RecordingStore(ObjectStore): + """In-memory ObjectStore that records uploads, to exercise upload_directory().""" + + def __init__(self, fail_on: str = "") -> None: + self.uploads: list[tuple[Path, str]] = [] + self.fail_on = fail_on + + def uri(self, key: str) -> str: + return f"mem://{key}" + + def upload_file(self, local_path: Path, key: str) -> None: + if self.fail_on and self.fail_on in key: + raise RuntimeError("boom") + self.uploads.append((local_path, key)) + + +@pytest.mark.parametrize( + "parts,expected", + [ + (("a", "b"), "a/b"), + (("", "b"), "b"), + (("a/", "/b"), "a/b"), + (("", ""), ""), + (("bucket", "", "key"), "bucket/key"), + ], +) +def test_join_key(parts: tuple[str, ...], expected: str) -> None: + assert join_key(*parts) == expected + + +def test_upload_directory_preserves_relative_paths(tree: Path) -> None: + store = RecordingStore() + stats = store.upload_directory(tree, "runs") + + assert sorted(key for _, key in store.uploads) == [ + "runs/report.html", + "runs/test-a/0/stderr.txt", + "runs/test-a/0/stdout.txt", + ] + assert stats.files_uploaded == 3 + assert stats.bytes_uploaded == len("") + len("err!!") + len("out") + assert stats.is_successful + + +def test_upload_directory_records_duration(tree: Path) -> None: + with patch("cloudai.util.object_store.time.perf_counter", side_effect=[100.0, 102.5]): + stats = RecordingStore().upload_directory(tree, "runs") + + assert stats.duration_seconds == 2.5 + + +def test_upload_directory_empty_dir_has_zero_duration(tmp_path: Path) -> None: + root = tmp_path / "empty" + root.mkdir() + + stats = RecordingStore().upload_directory(root, "runs") + + assert stats.duration_seconds == 0.0 + + +def test_upload_directory_without_prefix(tree: Path) -> None: + store = RecordingStore() + store.upload_directory(tree) + + assert "report.html" in [key for _, key in store.uploads] + + +def test_upload_directory_honours_exclude(tree: Path) -> None: + store = RecordingStore() + stats = store.upload_directory(tree, "runs", exclude=["*.txt"]) + + assert [key for _, key in store.uploads] == ["runs/report.html"] + assert stats.files_uploaded == 1 + + +def test_upload_directory_collects_failures_and_continues(tree: Path) -> None: + store = RecordingStore(fail_on="stdout.txt") + stats = store.upload_directory(tree, "runs") + + assert stats.files_uploaded == 2 + assert not stats.is_successful + assert len(stats.failures) == 1 + failed_path, message = stats.failures[0] + assert failed_path.name == "stdout.txt" + assert "boom" in message + + +def test_upload_directory_skips_empty_dirs(tmp_path: Path) -> None: + root = tmp_path / "empty" + (root / "nested").mkdir(parents=True) + + stats = RecordingStore().upload_directory(root, "runs") + + assert stats == UploadStats() + + +def test_upload_directory_concurrent_stats_are_accurate(tmp_path: Path) -> None: + root = tmp_path / "many" + root.mkdir() + for i in range(40): + name = f"fail_{i}.txt" if i % 3 == 0 else f"ok_{i}.txt" + (root / name).write_text("x") + + store = RecordingStore(fail_on="fail_") + stats = store.upload_directory(root, "runs", max_workers=8) + + expected_failures = sum(1 for i in range(40) if i % 3 == 0) + assert stats.files_uploaded == 40 - expected_failures + assert len(stats.failures) == expected_failures + assert len(store.uploads) == 40 - expected_failures + assert stats.files_uploaded + len(stats.failures) == 40 + + +def test_upload_directory_max_workers_bounded_by_file_count(tree: Path) -> None: + store = RecordingStore() + stats = store.upload_directory(tree, "runs", max_workers=100) + + assert stats.files_uploaded == 3 + + +def test_s3_object_store_uri() -> None: + store = S3ObjectStore(bucket="my-bucket") + assert store.uri("runs/report.html") == "s3://my-bucket/runs/report.html" + + +def test_s3_object_store_upload_file_and_client_reuse(tmp_path: Path) -> None: + local = tmp_path / "report.html" + local.write_text("") + + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket", endpoint_url="http://localhost:9000", region="us-east-1") + store.upload_file(local, "runs/report.html") + store.upload_file(local, "runs/again.html") + + mock_lazy.boto3.client.assert_called_once_with( + "s3", endpoint_url="http://localhost:9000", region_name="us-east-1" + ) + client.upload_file.assert_any_call(str(local), "my-bucket", "runs/report.html") + assert client.upload_file.call_count == 2 + + +def test_s3_object_store_exists() -> None: + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + client.exceptions.ClientError = ClientError + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket") + assert store.exists("present") is True + + client.head_object.side_effect = _client_error(404) + assert store.exists("missing") is False + + +def test_s3_object_store_exists_reraises_non_404_errors() -> None: + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + client.exceptions.ClientError = ClientError + client.head_object.side_effect = _client_error(403) + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket") + with pytest.raises(ClientError): + store.exists("forbidden") + + +def test_s3_object_store_bucket_exists() -> None: + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + client.exceptions.ClientError = ClientError + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket") + assert store.bucket_exists() is True + + client.head_bucket.side_effect = _client_error(404) + assert store.bucket_exists() is False + + +def test_s3_object_store_bucket_exists_tolerates_403() -> None: + """Write-only credentials lack s3:ListBucket, so HeadBucket returns 403 even though uploads work.""" + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + client.exceptions.ClientError = ClientError + client.head_bucket.side_effect = _client_error(403) + mock_lazy.boto3.client.return_value = client + + assert S3ObjectStore(bucket="my-bucket").bucket_exists() is True + + +def test_s3_object_store_bucket_exists_reraises_other_errors() -> None: + with patch("cloudai.util.object_store.lazy") as mock_lazy: + client = MagicMock() + client.exceptions.ClientError = ClientError + client.head_bucket.side_effect = _client_error(500) + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket") + with pytest.raises(ClientError): + store.bucket_exists() diff --git a/uv.lock b/uv.lock index ba191e1b2..4317be098 100644 --- a/uv.lock +++ b/uv.lock @@ -130,6 +130,34 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/f6/a8/877f306720bc114c612579c5af36bcb359026b83d051226945499b306b1a/bokeh-3.8.2-py3-none-any.whl", hash = "sha256:5e2c0d84f75acb25d60efb9e4d2f434a791c4639b47d685534194c4e07bd0111", size = 7207131, upload-time = "2026-01-06T00:20:04.917Z" }, ] +[[package]] +name = "boto3" +version = "1.43.109" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "botocore" }, + { name = "jmespath" }, + { name = "s3transfer" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/2a/80/0430e16c302b0d1a4ca3c8c69ffee9e4f8f2170d56162d2761b34fd512ba/boto3-1.43.109.tar.gz", hash = "sha256:c829bc3352e1e3922d8fe1db10b697182ac411b64ce5fc3acb0a06a29837417e", size = 112661, upload-time = "2026-10-07T15:24:51.293Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/a1/56/c83a8fdadf06f631d1a0c4d82fc8daa9e1a40f6354264fba20c186363f49/boto3-1.43.109-py3-none-any.whl", hash = "sha256:703a81baed186d31e27cf50065d25a9faeb704f009e7c56b0f002ca52a328896", size = 140044, upload-time = "2026-10-07T15:24:49.793Z" }, +] + +[[package]] +name = "botocore" +version = "1.43.109" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "jmespath" }, + { name = "python-dateutil" }, + { name = "urllib3" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/3c/da/424aaef4727876f4095bd2c10e683bb8531ac3bea56b0b6efd29eeb57f91/botocore-1.43.109.tar.gz", hash = "sha256:46b15bea4d942aaea6d7ab7c4b8a3a7290305064a3b41a6cb271cde9da3ba371", size = 16301678, upload-time = "2026-10-07T15:24:46.458Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/7c/d7/bf9e5988baf1245f3ac373e09f3204f641e7a2db170f512f236aa90961d6/botocore-1.43.109-py3-none-any.whl", hash = "sha256:cd00edff465264fb448e36acc9ee29a828f738fff3746270b3f685a19194a5a6", size = 15998671, upload-time = "2026-10-07T15:24:43.76Z" }, +] + [[package]] name = "build" version = "1.4.0" @@ -289,6 +317,7 @@ dependencies = [ [package.optional-dependencies] dev = [ + { name = "boto3" }, { name = "build" }, { name = "gymnasium" }, { name = "import-linter" }, @@ -324,12 +353,17 @@ docs-cms = [ rl = [ { name = "gymnasium" }, ] +s3 = [ + { name = "boto3" }, +] [package.metadata] requires-dist = [ { name = "autodoc-pydantic", marker = "extra == 'docs'", specifier = "~=2.2" }, { name = "autodoc-pydantic", marker = "extra == 'docs-cms'", specifier = "~=2.2" }, { name = "bokeh", specifier = "~=3.8" }, + { name = "boto3", marker = "extra == 'dev'", specifier = "~=1.40" }, + { name = "boto3", marker = "extra == 's3'", specifier = "~=1.40" }, { name = "build", marker = "extra == 'dev'", specifier = "~=1.4" }, { name = "click", specifier = "~=8.3" }, { name = "gymnasium", marker = "extra == 'dev'", specifier = "~=1.2" }, @@ -369,7 +403,7 @@ requires-dist = [ { name = "vulture", marker = "extra == 'dev'", specifier = "==2.14" }, { name = "websockets", specifier = "~=16.0" }, ] -provides-extras = ["dev", "rl", "docs", "docs-cms"] +provides-extras = ["dev", "rl", "s3", "docs", "docs-cms"] [[package]] name = "cloudpickle" @@ -1105,6 +1139,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/62/a1/3d680cbfd5f4b8f15abc1d571870c5fc3e594bb582bc3b64ea099db13e56/jinja2-3.1.6-py3-none-any.whl", hash = "sha256:85ece4451f492d0c13c5dd7c13a64681a86afae63a5f347908daf103ce6d2f67", size = 134899, upload-time = "2025-03-05T20:05:00.369Z" }, ] +[[package]] +name = "jmespath" +version = "1.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/d3/59/322338183ecda247fb5d1763a6cbe46eff7222eaeebafd9fa65d4bf5cb11/jmespath-1.1.0.tar.gz", hash = "sha256:472c87d80f36026ae83c6ddd0f1d05d4e510134ed462851fd5f754c8c3cbb88d", size = 27377, upload-time = "2026-01-22T16:35:26.279Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/14/2f/967ba146e6d58cf6a652da73885f52fc68001525b4197effc174321d70b4/jmespath-1.1.0-py3-none-any.whl", hash = "sha256:a5663118de4908c91729bea0acadca56526eb2698e83de10cd116ae0f4e97c64", size = 20419, upload-time = "2026-01-22T16:35:24.919Z" }, +] + [[package]] name = "kubernetes" version = "35.0.0" @@ -2147,6 +2190,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/f6/b0/2d823f6e77ebe560f4e397d078487e8d52c1516b331e3521bc75db4272ca/ruff-0.15.0-py3-none-win_arm64.whl", hash = "sha256:c480d632cc0ca3f0727acac8b7d053542d9e114a462a145d0b00e7cd658c515a", size = 10865753, upload-time = "2026-02-03T17:53:03.014Z" }, ] +[[package]] +name = "s3transfer" +version = "0.19.2" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "botocore" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/76/43/35e4d8aa320bffe8287fe8f65f578fa2d2db0a64212f0e710dce58267854/s3transfer-0.19.2.tar.gz", hash = "sha256:ba0309fd86be3c27dbf78cdd813c13c5e1df16e5874b99d2535ebbdfb9892993", size = 165592, upload-time = "2026-07-22T19:30:44.432Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/bc/e7/5c595c75e9f41a44f30e526eda465ea0b4eec93470e074e4a111b253f13a/s3transfer-0.19.2-py3-none-any.whl", hash = "sha256:d8168eccca828cbb2cd573675333f3bddd254313a9c42494b84c76b539e8ba25", size = 90216, upload-time = "2026-07-22T19:30:43.251Z" }, +] + [[package]] name = "setuptools" version = "83.0.0"