From a6a35f8aa3fe9737bfcca4c4114fdea315c08dcf Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Thu, 10 Sep 2026 12:33:30 -0600 Subject: [PATCH 01/15] Add generic ObjectStore interface with an S3 implementation Introduce cloudai.util.object_store, a domain-free transport layer for publishing artifacts to object storage: - ObjectStore ABC with upload_file/uri plus a shared upload_directory() that walks a tree, preserves relative paths, honours exclude globs, and collects per-file failures into UploadStats instead of aborting. - S3ObjectStore, backed by boto3. Credentials come from boto3's standard chain (AWS_* env, ~/.aws/credentials, IAM role), so no secrets are ever read from CloudAI configuration. boto3 is an optional 'cloudai[s3]' extra reached through the existing LazyImports pattern, matching how gymnasium is handled. This keeps it out of module-level imports, which the ruff banned-module-level-imports rule and the filterwarnings=["error"] pytest setting both require. The module lives under cloudai.util so the import-linter leaf-dependency contract keeps it free of cloudai domain concepts. --- pyproject.toml | 2 + src/cloudai/util/lazy_imports.py | 15 ++++ src/cloudai/util/object_store.py | 143 +++++++++++++++++++++++++++++++ uv.lock | 57 +++++++++++- 4 files changed, 216 insertions(+), 1 deletion(-) create mode 100644 src/cloudai/util/object_store.py 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/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..00fd7f032 --- /dev/null +++ b/src/cloudai/util/object_store.py @@ -0,0 +1,143 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2025-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. + +"""Generic object storage interface and an S3 implementation.""" + +from __future__ import annotations + +import fnmatch +import logging +from abc import ABC, abstractmethod +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) + + @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 + ) -> 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. + + Returns: + Stats describing what was uploaded. Individual file failures are collected + rather than raised, so a single bad file does not abort the whole upload. + """ + stats = UploadStats() + exclude = exclude or [] + + 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 + + key = join_key(key_prefix, relative.as_posix()) + try: + size = path.stat().st_size + self.upload_file(path, key) + 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 + + return stats + + +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 Exception: + return False + return True diff --git a/uv.lock b/uv.lock index ba191e1b2..9228af30f 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.91" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "botocore" }, + { name = "jmespath" }, + { name = "s3transfer" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/42/76/e5c5fb601b09273599bf9c1a678ea2bc645385e758fbc66a36c9614b6324/boto3-1.43.91.tar.gz", hash = "sha256:98643e500883bf6fcd13d04bb19da983bf97e5f1c600fc11470ce45437fa96b4", size = 112693, upload-time = "2026-09-09T19:24:38.927Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/24/3f/0288c2383455f1c94741fd8de819fb6d7360065274d905711a62fc6cd1f2/boto3-1.43.91-py3-none-any.whl", hash = "sha256:5ff948cfac8bff72227930e0894a8542731e1998634d1ffd8a074c20a90af919", size = 140023, upload-time = "2026-09-09T19:24:37.588Z" }, +] + +[[package]] +name = "botocore" +version = "1.43.91" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "jmespath" }, + { name = "python-dateutil" }, + { name = "urllib3" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/9a/1e/c6d26cb1124799e2bee54854813b02aa8a813c57335f97b0b8ed50f902ff/botocore-1.43.91.tar.gz", hash = "sha256:0f12bceb8c5d90a0c60f326e883bf674ded8c7d7c18ee0b559843b3a6c52f2ad", size = 16089153, upload-time = "2026-09-09T19:24:34.426Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/8f/1a/f672cc408d808035397c5d63a1b31a6a56ee95417a9d60d120d96349d182/botocore-1.43.91-py3-none-any.whl", hash = "sha256:f96363d4caf50bce45fe2fdc2d4d86b3bea7f5cd805169d53c1371c2dd4772b6", size = 15782250, upload-time = "2026-09-09T19:24:31.222Z" }, +] + [[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" From aa152fa2419a9c425211616466746bef36b0a5aa Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Thu, 10 Sep 2026 12:34:19 -0600 Subject: [PATCH 02/15] Add ResultsUploadReporter to publish results to object storage ResultsUploadConfig subclasses ReportConfig with the destination fields (bucket, prefix, endpoint_url, region) plus upload_tree/upload_tarball/ exclude toggles. Each destination field falls back to a CLOUDAI_S3_* environment variable via default_factory, so TOML wins when set and the environment supplies the default otherwise. ReportConfig forbids extra keys, so every field has to be declared explicitly. ResultsUploadReporter uploads self.results_root, which the Reporter base class already hands it. Two deliberate choices: - It does not call load_test_runs(). This reporter treats the directory as raw bytes, so coupling the upload to workload parsing would only add failure modes. - upload_tarball() creates the tarball when absent. TarballReporter only writes .tgz when a test run failed, so its presence cannot be assumed; the archiving logic is reused rather than duplicated. A missing bucket warns loudly instead of silently succeeding, so a misconfigured upload is not mistaken for an empty one. --- src/cloudai/reporter.py | 84 ++++++++++++++++++++++++++++++++++++++++- 1 file changed, 83 insertions(+), 1 deletion(-) diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 918fe26fb..d6d3b96bc 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -16,6 +16,7 @@ import contextlib import logging +import os import tarfile import xml.etree.ElementTree as ET from dataclasses import dataclass @@ -24,6 +25,7 @@ import jinja2 import toml +from pydantic import Field from rich import box from rich.console import Console from rich.table import Table @@ -32,7 +34,8 @@ from cloudai.report_generator.util import load_system_metadata from .core import CommandGenStrategy, Reporter, TestRun, case_name -from .models.scenario import TestRunDetails +from .models.scenario import ReportConfig, TestRunDetails +from .util.object_store import S3ObjectStore, join_key @dataclass @@ -313,3 +316,82 @@ def create_tarball(self, directory: Path) -> None: with tarfile.open(tarball_path, "w:gz") as tar: tar.add(directory, arcname=directory.name) logging.info(f"Created tarball at {tarball_path}") + + +class ResultsUploadConfig(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 + exclude: list[str] = Field(default_factory=list) + + +class ResultsUploadReporter(Reporter): + """Uploads the scenario results directory to object storage.""" + + def generate(self) -> None: + config = self.config + if not isinstance(config, ResultsUploadConfig): + logging.warning(f"Expected ResultsUploadConfig, got {type(config).__name__}, skipping results upload.") + return + + if not config.bucket: + logging.warning( + "Results 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) + key_prefix = join_key(config.prefix, self.results_root.name) + + if config.upload_tree: + stats = store.upload_directory(self.results_root, key_prefix, config.exclude) + logging.info( + f"Uploaded {stats.files_uploaded} file(s), {stats.bytes_uploaded} byte(s) to {store.uri(key_prefix)}" + ) + 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, creating it if it does not exist. + + TarballReporter only produces a tarball when a test run failed, so we cannot + assume one is already present. + """ + tarball_path = Path(str(self.results_root) + ".tgz") + if not tarball_path.exists(): + 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)}") From 41b74dd8eb05be8cfb467a699421e71d1c15a311 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Thu, 10 Sep 2026 12:35:05 -0600 Subject: [PATCH 03/15] Register results_upload and order it after every other report Registry.ordered_scenario_reports() sorts by a hardcoded map in which any unlisted name falls through to priority 1. Left at that default the uploader would run before StatusReporter, DSEReporter and TarballReporter, and would therefore publish an incomplete results directory. Give "results_upload" priority 5 so it always observes the finished tree. Registered with enable=False: shipping results off-box has to be opt-in. Registering here also surfaces it in `cloudai list reports` for free. --- src/cloudai/_core/registry.py | 3 ++- src/cloudai/core.py | 15 ++++++++++++++- src/cloudai/registration.py | 11 ++++++++++- 3 files changed, 26 insertions(+), 3 deletions(-) diff --git a/src/cloudai/_core/registry.py b/src/cloudai/_core/registry.py index b0ccf5a75..71dc0ee45 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, + "results_upload": 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..eb2b37fd4 100644 --- a/src/cloudai/core.py +++ b/src/cloudai/core.py @@ -70,9 +70,17 @@ from .configurator.gymnasium_adapter import GymnasiumAdapter from .models.workload import CmdArgs, NsysConfiguration, PredictorConfig, TestDefinition from .parser import Parser -from .reporter import JUnitReporter, PerTestReporter, StatusReporter, TarballReporter +from .reporter import ( + JUnitReporter, + PerTestReporter, + ResultsUploadConfig, + ResultsUploadReporter, + StatusReporter, + TarballReporter, +) from .test_parser import TestParser from .test_scenario_parser import TestScenarioParser +from .util.object_store import ObjectStore, S3ObjectStore, UploadStats __all__ = [ "METRIC_ERROR", @@ -107,6 +115,7 @@ "MetricValue", "MissingTestError", "NsysConfiguration", + "ObjectStore", "ObsLeafDescriptor", "Parser", "PerTestReporter", @@ -116,8 +125,11 @@ "Registry", "ReportGenerationStrategy", "Reporter", + "ResultsUploadConfig", + "ResultsUploadReporter", "RewardOverrides", "Runner", + "S3ObjectStore", "StatusReporter", "StructuredObservationProducer", "System", @@ -131,6 +143,7 @@ "TestScenario", "TestScenarioParser", "TestScenarioParsingError", + "UploadStats", "case_name", "format_validation_error", ] diff --git a/src/cloudai/registration.py b/src/cloudai/registration.py index f471a5fcd..27443b0a5 100644 --- a/src/cloudai/registration.py +++ b/src/cloudai/registration.py @@ -47,7 +47,15 @@ def register_all(): from cloudai.core import Registry from cloudai.models.scenario import ReportConfig from cloudai.report_generator.training import TrainingReporter - from cloudai.reporter import DSEReporter, JUnitReporter, PerTestReporter, StatusReporter, TarballReporter + from cloudai.reporter import ( + DSEReporter, + JUnitReporter, + PerTestReporter, + ResultsUploadConfig, + ResultsUploadReporter, + StatusReporter, + TarballReporter, + ) # Import systems from cloudai.systems.kubernetes import KubernetesInstaller, KubernetesRunner, KubernetesSystem @@ -343,6 +351,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("results_upload", ResultsUploadReporter, ResultsUploadConfig(enable=False)) Registry().add_scenario_report( "nixl_bench_summary", NIXLBenchComparisonReport, From c8027f3c5880feb84384fbbfa849c5bfbf133022 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Thu, 10 Sep 2026 12:37:40 -0600 Subject: [PATCH 04/15] Test the object store and results upload reporter tests/util/test_object_store.py drives upload_directory() through a small in-memory ObjectStore rather than mocking the ABC, so the shared walk logic - relative-path preservation, exclude globs, and failure collection - is tested independently of boto3. The S3 tests patch cloudai.util.object_store.lazy, matching the repo convention of patching at the consumer's import site. tests/test_reporter.py covers the reporter: env-var fallback and TOML precedence, the missing-bucket and missing-directory guards, and the tarball being created when absent but reused when present. Two existing tests needed updating for the new registration: - test_report_order asserted positions by negative index. - test_init hardcodes the registry contents and asserted that every report except junit defaults to enabled; results_upload joins junit as opt-in. --- tests/test_init.py | 14 ++- tests/test_reporter.py | 126 ++++++++++++++++++++++++++- tests/util/test_object_store.py | 148 ++++++++++++++++++++++++++++++++ 3 files changed, 282 insertions(+), 6 deletions(-) create mode 100644 tests/util/test_object_store.py diff --git a/tests/test_init.py b/tests/test_init.py index 2943b4676..5932464d6 100644 --- a/tests/test_init.py +++ b/tests/test_init.py @@ -17,7 +17,14 @@ from cloudai.core import Registry from cloudai.report_generator.training import TrainingReporter -from cloudai.reporter import DSEReporter, JUnitReporter, PerTestReporter, StatusReporter, TarballReporter +from cloudai.reporter import ( + DSEReporter, + JUnitReporter, + PerTestReporter, + ResultsUploadReporter, + StatusReporter, + TarballReporter, +) from cloudai.systems.kubernetes import KubernetesInstaller, KubernetesSystem from cloudai.systems.lsf import LSFInstaller, LSFSystem from cloudai.systems.runai import RunAISystem @@ -287,6 +294,7 @@ def test_scenario_reports(): "junit", "dse", "tarball", + "results_upload", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -303,6 +311,7 @@ def test_scenario_reports(): JUnitReporter, DSEReporter, TarballReporter, + ResultsUploadReporter, NIXLBenchComparisonReport, NixlEPComparisonReport, NcclComparisonReport, @@ -323,6 +332,7 @@ def test_report_configs(): "junit", "dse", "tarball", + "results_upload", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -331,7 +341,7 @@ def test_report_configs(): "vllm_comparison", "sglang_comparison", ] - disabled_configs = {"junit"} + disabled_configs = {"junit", "results_upload"} # 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..81d827355 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -21,6 +21,7 @@ from dataclasses import asdict from pathlib import Path from typing import Any +from unittest.mock import patch import pytest import toml @@ -30,7 +31,16 @@ 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, + ResultsUploadConfig, + ResultsUploadReporter, + StatusReporter, + TarballReporter, +) from cloudai.systems.slurm.slurm_metadata import ( MetadataCUDA, MetadataMPI, @@ -44,6 +54,7 @@ ) from cloudai.systems.slurm.slurm_system import SlurmSystem from cloudai.systems.standalone.standalone_system import StandaloneSystem +from cloudai.util.object_store import UploadStats from cloudai.workloads.nccl_test import NCCLCmdArgs, NCCLTestDefinition @@ -440,9 +451,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] == "results_upload" def test_junit_reporter_generates_testcases_with_status_logs_and_duration( @@ -680,3 +692,109 @@ def test_dse_reporter( assert (slurm_system.output_path / "single-dse-scenario-dse-report.html").exists() assert (slurm_system.output_path / dse_case.name / "0" / f"{dse_case.name}.toml").exists() + + +class TestResultsUploadReporter: + """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) -> ResultsUploadReporter: + return ResultsUploadReporter( + slurm_system, + TestScenario(name="dummy", test_runs=[]), + results_dir, + ResultsUploadConfig(enable=True, **kwargs), + ) + + def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.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/nccl-test_2025-04-16_14-27-45", []) + + def test_no_bucket_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.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.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_upload_tree_disabled(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.reporter.S3ObjectStore") as mock_store_cls: + store = mock_store_cls.return_value + + self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_tree=False).generate() + + store.upload_directory.assert_not_called() + + 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.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, "nccl-test_2025-04-16_14-27-45/nccl-test_2025-04-16_14-27-45.tgz" + ) + + def test_existing_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") + + with patch("cloudai.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" + + 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 = ResultsUploadConfig(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 ResultsUploadConfig(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.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_results_upload_runs_after_tarball() -> None: + order = [name for name, _ in Registry().ordered_scenario_reports()] + + assert order.index("results_upload") > order.index("tarball") + assert order[-1] == "results_upload" diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py new file mode 100644 index 000000000..3da2ba208 --- /dev/null +++ b/tests/util/test_object_store.py @@ -0,0 +1,148 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2025-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 cloudai.util.object_store import ObjectStore, S3ObjectStore, UploadStats, join_key + + +@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_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_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() + mock_lazy.boto3.client.return_value = client + + store = S3ObjectStore(bucket="my-bucket") + assert store.exists("present") is True + + client.head_object.side_effect = RuntimeError("404") + assert store.exists("missing") is False From 6b80da2781dc076df2540a084aac6caaff5f21f1 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Thu, 10 Sep 2026 12:38:06 -0600 Subject: [PATCH 05/15] Document the results_upload report Covers the cloudai[s3] extra, every config option with its environment variable fallback, the fact that credentials come from boto3's standard chain rather than CloudAI config, that an upload failure does not change the run's exit status, and how to upload an earlier run's directory with generate-report. --- doc/reporting.rst | 81 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 81 insertions(+) diff --git a/doc/reporting.rst b/doc/reporting.rst index 386f6d72c..38e42718b 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,83 @@ 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 ``results_upload`` 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] + results_upload = { 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] + results_upload = { 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, creating it if absent. + * - ``exclude`` + - ``[]`` + - Glob patterns matched against paths relative to the results directory. + +Destination fields fall back to the environment variable shown above when not set in +TOML, so a cluster-wide default can come from the environment while an individual +scenario can still override it. + +**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/_ From 4b9feca500524c250d2b74201b69bf4691e920f0 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Mon, 21 Sep 2026 10:06:10 -0600 Subject: [PATCH 06/15] Drop exclude from ResultsUploadConfig It had no real consumer yet and TarballReporter.create_tarball() never honored it, so enabling upload_tarball alongside exclude would silently ship files the config claimed to filter out. ObjectStore. upload_directory() still supports exclude for a future caller that actually needs it. Co-Authored-By: Claude Sonnet 5 --- doc/reporting.rst | 3 --- src/cloudai/reporter.py | 3 +-- tests/test_reporter.py | 2 +- 3 files changed, 2 insertions(+), 6 deletions(-) diff --git a/doc/reporting.rst b/doc/reporting.rst index 38e42718b..b850433fe 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -243,9 +243,6 @@ Configuration options: * - ``upload_tarball`` - ``false`` - Also upload a ``.tgz`` of the whole directory, creating it if absent. - * - ``exclude`` - - ``[]`` - - Glob patterns matched against paths relative to the results directory. Destination fields fall back to the environment variable shown above when not set in TOML, so a cluster-wide default can come from the environment while an individual diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index d6d3b96bc..2e3f843ce 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -334,7 +334,6 @@ class ResultsUploadConfig(ReportConfig): region: Optional[str] = None upload_tree: bool = True upload_tarball: bool = False - exclude: list[str] = Field(default_factory=list) class ResultsUploadReporter(Reporter): @@ -361,7 +360,7 @@ def generate(self) -> None: key_prefix = join_key(config.prefix, self.results_root.name) if config.upload_tree: - stats = store.upload_directory(self.results_root, key_prefix, config.exclude) + stats = store.upload_directory(self.results_root, key_prefix) logging.info( f"Uploaded {stats.files_uploaded} file(s), {stats.bytes_uploaded} byte(s) to {store.uri(key_prefix)}" ) diff --git a/tests/test_reporter.py b/tests/test_reporter.py index 81d827355..d95903506 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -721,7 +721,7 @@ def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> Non 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/nccl-test_2025-04-16_14-27-45", []) + store.upload_directory.assert_called_once_with(results_dir, "cloudai/nccl-test_2025-04-16_14-27-45") def test_no_bucket_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: with patch("cloudai.reporter.S3ObjectStore") as mock_store_cls: From 273a83642b029dfd74181694e9ce5f02d2f65486 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Mon, 21 Sep 2026 16:24:14 -0600 Subject: [PATCH 07/15] Validate bucket accessibility before upload, tighten exists() error handling S3ObjectStore.exists() previously swallowed every exception (permission errors, network failures, throttling) into a bare False, indistinguishable from a genuinely missing object. It now catches only ClientError and treats a 404 status as "does not exist," re-raising anything else. Add S3ObjectStore.bucket_exists(), using head_bucket rather than head_object, since head_object's 404 can't distinguish a missing bucket from a missing key. ResultsUploadReporter now checks bucket_exists() right after constructing the store, so a misconfigured or inaccessible bucket fails fast with one clear warning instead of every individual file upload failing separately with the same root cause. Co-Authored-By: Claude Sonnet 5 --- src/cloudai/reporter.py | 6 ++++ src/cloudai/util/object_store.py | 18 +++++++++-- tests/test_reporter.py | 9 ++++++ tests/util/test_object_store.py | 51 +++++++++++++++++++++++++++++++- 4 files changed, 81 insertions(+), 3 deletions(-) diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 2e3f843ce..302941a4f 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -357,6 +357,12 @@ def generate(self) -> None: 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 or is not accessible, skipping results upload." + ) + return + key_prefix = join_key(config.prefix, self.results_root.name) if config.upload_tree: diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py index 00fd7f032..6ba745166 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -138,6 +138,20 @@ def upload_file(self, local_path: Path, key: str) -> None: def exists(self, key: str) -> bool: try: self.client.head_object(Bucket=self.bucket, Key=key) - except Exception: - return False + 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 whether the configured bucket exists and is accessible.""" + try: + self.client.head_bucket(Bucket=self.bucket) + except self.client.exceptions.ClientError as e: + status = e.response.get("ResponseMetadata", {}).get("HTTPStatusCode") + if status in (404, 403): + logging.debug(f"Bucket '{self.bucket}' is not accessible: {e}") + return False + raise return True diff --git a/tests/test_reporter.py b/tests/test_reporter.py index d95903506..154f000e7 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -735,6 +735,15 @@ def test_missing_results_dir_uploads_nothing(self, slurm_system: SlurmSystem, tm mock_store_cls.assert_not_called() + def test_bucket_not_accessible_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + with patch("cloudai.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.reporter.S3ObjectStore") as mock_store_cls: store = mock_store_cls.return_value diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py index 3da2ba208..21f7c9d1a 100644 --- a/tests/util/test_object_store.py +++ b/tests/util/test_object_store.py @@ -18,10 +18,18 @@ 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" @@ -139,10 +147,51 @@ def test_s3_object_store_upload_file_and_client_reuse(tmp_path: Path) -> None: 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 = RuntimeError("404") + 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 + + client.head_bucket.side_effect = _client_error(403) + assert store.bucket_exists() is False + + +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() From fc4f9ea2197e30a88e97a7f02869983a83943e7d Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Mon, 21 Sep 2026 16:24:34 -0600 Subject: [PATCH 08/15] Scope uploaded S3 keys under the system name Without it, two clusters sharing one bucket (or one prefix) have no structural way to tell their results apart, and per-cluster IAM scoping by key prefix isn't enforceable since nothing guarantees the prefix is actually cluster-specific. system.name is already a required field on every system config, so this needs no new configuration. Co-Authored-By: Claude Sonnet 5 --- doc/reporting.rst | 2 +- src/cloudai/reporter.py | 2 +- tests/test_reporter.py | 6 ++++-- 3 files changed, 6 insertions(+), 4 deletions(-) diff --git a/doc/reporting.rst b/doc/reporting.rst index b850433fe..93aa41791 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -230,7 +230,7 @@ Configuration options: - Destination bucket. Required; the upload is skipped with a warning if unset. * - ``prefix`` - ``$CLOUDAI_S3_PREFIX`` - - Key prefix. Objects are written under ``//``. + - Key prefix. Objects are written under ``///``. * - ``endpoint_url`` - ``$CLOUDAI_S3_ENDPOINT_URL`` - Custom endpoint, for MinIO or other S3-compatible stores. diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 302941a4f..489ac0981 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -363,7 +363,7 @@ def generate(self) -> None: ) return - key_prefix = join_key(config.prefix, self.results_root.name) + 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) diff --git a/tests/test_reporter.py b/tests/test_reporter.py index 154f000e7..c1efbcaf1 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -721,7 +721,9 @@ def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> Non 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/nccl-test_2025-04-16_14-27-45") + store.upload_directory.assert_called_once_with( + results_dir, "cloudai/test_system/nccl-test_2025-04-16_14-27-45" + ) def test_no_bucket_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: with patch("cloudai.reporter.S3ObjectStore") as mock_store_cls: @@ -764,7 +766,7 @@ def test_tarball_created_when_absent(self, slurm_system: SlurmSystem, results_di assert tarball_path.exists(), "TarballReporter only tarballs on failure, so it must be created here" store.upload_file.assert_called_once_with( - tarball_path, "nccl-test_2025-04-16_14-27-45/nccl-test_2025-04-16_14-27-45.tgz" + tarball_path, "test_system/nccl-test_2025-04-16_14-27-45/nccl-test_2025-04-16_14-27-45.tgz" ) def test_existing_tarball_is_reused(self, slurm_system: SlurmSystem, results_dir: Path) -> None: From 85fc94e8683101f73201de05ac88fad577e5de69 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Tue, 22 Sep 2026 09:28:04 -0600 Subject: [PATCH 09/15] Upload tree files concurrently instead of one at a time upload_directory() previously uploaded files serially, one blocking S3 call at a time, which dominates wall-clock time for result trees with many small files. It now uploads concurrently via a ThreadPoolExecutor (max_workers, default 8, bounded by the actual file count), aggregating UploadStats back on the main thread via as_completed() so concurrent updates can't race. Exposed as upload_concurrency on ResultsUploadConfig. Measured against a real S3 bucket, uploading a 114-file / 3.68MB results tree (cloudai-db-test, same code path, only max_workers changed): serial (max_workers=1): 17.95s, 114 files, 0 failures threaded (max_workers=8): 3.89s, 114 files, 0 failures ~4.6x speedup, identical output in both cases. Co-Authored-By: Claude Sonnet 5 --- doc/reporting.rst | 3 +++ src/cloudai/reporter.py | 3 ++- src/cloudai/util/object_store.py | 43 ++++++++++++++++++++++++-------- tests/test_reporter.py | 12 ++++++++- tests/util/test_object_store.py | 24 ++++++++++++++++++ 5 files changed, 72 insertions(+), 13 deletions(-) diff --git a/doc/reporting.rst b/doc/reporting.rst index 93aa41791..d2f034f38 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -243,6 +243,9 @@ Configuration options: * - ``upload_tarball`` - ``false`` - Also upload a ``.tgz`` of the whole directory, creating it if absent. + * - ``upload_concurrency`` + - ``8`` + - Number of files uploaded concurrently when ``upload_tree`` is enabled. Destination fields fall back to the environment variable shown above when not set in TOML, so a cluster-wide default can come from the environment while an individual diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 489ac0981..88a57d6e8 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -334,6 +334,7 @@ class ResultsUploadConfig(ReportConfig): region: Optional[str] = None upload_tree: bool = True upload_tarball: bool = False + upload_concurrency: int = 8 class ResultsUploadReporter(Reporter): @@ -366,7 +367,7 @@ def generate(self) -> None: 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) + 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)}" ) diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py index 6ba745166..b369b9026 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -21,6 +21,7 @@ import fnmatch import logging 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 @@ -61,7 +62,11 @@ def upload_file(self, local_path: Path, key: str) -> None: ... def upload_directory( - self, local_dir: Path, key_prefix: str = "", exclude: Optional[list[str]] = None + 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. @@ -71,6 +76,8 @@ def upload_directory( 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. Individual file failures are collected @@ -79,6 +86,7 @@ def upload_directory( stats = UploadStats() exclude = exclude or [] + to_upload: list[tuple[Path, str]] = [] for path in sorted(local_dir.rglob("*")): if not path.is_file(): continue @@ -88,20 +96,33 @@ def upload_directory( logging.debug(f"Skipping excluded file {relative}") continue - key = join_key(key_prefix, relative.as_posix()) - try: - size = path.stat().st_size - self.upload_file(path, key) - 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 + to_upload.append((path, join_key(key_prefix, relative.as_posix()))) + + if not to_upload: + return stats - stats.files_uploaded += 1 - stats.bytes_uploaded += size + 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 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): """ diff --git a/tests/test_reporter.py b/tests/test_reporter.py index c1efbcaf1..32d317943 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -722,9 +722,19 @@ def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> Non 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" + 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.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.reporter.S3ObjectStore") as mock_store_cls: self.reporter(slurm_system, results_dir).generate() diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py index 21f7c9d1a..4cd070344 100644 --- a/tests/util/test_object_store.py +++ b/tests/util/test_object_store.py @@ -120,6 +120,30 @@ def test_upload_directory_skips_empty_dirs(tmp_path: Path) -> None: 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" From 3ec947b78169e7253b719700e91d5f17c2c058b0 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Mon, 28 Sep 2026 15:02:57 -0600 Subject: [PATCH 10/15] add upload time to upload stats --- src/cloudai/reporter.py | 3 ++- src/cloudai/util/object_store.py | 10 ++++++++-- tests/util/test_object_store.py | 16 ++++++++++++++++ 3 files changed, 26 insertions(+), 3 deletions(-) diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index 88a57d6e8..e8b1aa939 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -369,7 +369,8 @@ def generate(self) -> None: 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"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( diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py index b369b9026..f58c1a2b6 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -20,6 +20,7 @@ import fnmatch import logging +import time from abc import ABC, abstractmethod from concurrent.futures import ThreadPoolExecutor, as_completed from dataclasses import dataclass, field @@ -42,6 +43,7 @@ class UploadStats: 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: @@ -80,8 +82,10 @@ def upload_directory( number of files, so small trees don't spin up idle threads. Returns: - Stats describing what was uploaded. Individual file failures are collected - rather than raised, so a single bad file does not abort the whole upload. + 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 [] @@ -101,6 +105,7 @@ def upload_directory( 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} @@ -115,6 +120,7 @@ def upload_directory( stats.files_uploaded += 1 stats.bytes_uploaded += size + stats.duration_seconds = time.perf_counter() - start return stats diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py index 4cd070344..0c8e7c338 100644 --- a/tests/util/test_object_store.py +++ b/tests/util/test_object_store.py @@ -84,6 +84,22 @@ def test_upload_directory_preserves_relative_paths(tree: Path) -> None: 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) From 9c4797ed18ad9361d11e14e5b99b86eb91d1b429 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Wed, 7 Oct 2026 10:23:29 -0600 Subject: [PATCH 11/15] address comments: rename results_upload to s3, move s3 config and reporter class into separate module and validate one of `upload_tree` or `upload_tarball` are set --- doc/reporting.rst | 10 +- src/cloudai/_core/registry.py | 2 +- src/cloudai/core.py | 14 +-- src/cloudai/registration.py | 13 +-- src/cloudai/reporter.py | 91 +----------------- src/cloudai/s3_reporter.py | 118 ++++++++++++++++++++++++ tests/test_init.py | 18 ++-- tests/test_reporter.py | 133 +-------------------------- tests/test_s3_reporter.py | 167 ++++++++++++++++++++++++++++++++++ 9 files changed, 307 insertions(+), 259 deletions(-) create mode 100644 src/cloudai/s3_reporter.py create mode 100644 tests/test_s3_reporter.py diff --git a/doc/reporting.rst b/doc/reporting.rst index d2f034f38..68740abb2 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -190,7 +190,7 @@ And it can be used in a test scenario as follows: Uploading Results to Object Storage ------------------------------------ -The ``results_upload`` scenario report publishes the scenario results directory to an +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. @@ -208,14 +208,14 @@ Then enable it in a test scenario: .. code-block:: toml [reports] - results_upload = { enable = true, bucket = "my-bucket", prefix = "cloudai/runs", upload_tarball = true } + 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] - results_upload = { enable = true, bucket = "my-bucket" } + s3 = { enable = true, bucket = "my-bucket" } Configuration options: @@ -245,7 +245,9 @@ Configuration options: - Also upload a ``.tgz`` of the whole directory, creating it if absent. * - ``upload_concurrency`` - ``8`` - - Number of files uploaded concurrently when ``upload_tree`` is enabled. + - 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. Destination fields fall back to the environment variable shown above when not set in TOML, so a cluster-wide default can come from the environment while an individual diff --git a/src/cloudai/_core/registry.py b/src/cloudai/_core/registry.py index 71dc0ee45..5485733bd 100644 --- a/src/cloudai/_core/registry.py +++ b/src/cloudai/_core/registry.py @@ -279,7 +279,7 @@ def report_order(k: str) -> int: "status": 2, "dse": 3, "tarball": 4, - "results_upload": 5, # last, must observe every other report's output + "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 eb2b37fd4..2a6e1dc58 100644 --- a/src/cloudai/core.py +++ b/src/cloudai/core.py @@ -70,14 +70,8 @@ from .configurator.gymnasium_adapter import GymnasiumAdapter from .models.workload import CmdArgs, NsysConfiguration, PredictorConfig, TestDefinition from .parser import Parser -from .reporter import ( - JUnitReporter, - PerTestReporter, - ResultsUploadConfig, - ResultsUploadReporter, - StatusReporter, - TarballReporter, -) +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 @@ -125,11 +119,11 @@ "Registry", "ReportGenerationStrategy", "Reporter", - "ResultsUploadConfig", - "ResultsUploadReporter", "RewardOverrides", "Runner", "S3ObjectStore", + "S3UploadConfig", + "S3UploadReporter", "StatusReporter", "StructuredObservationProducer", "System", diff --git a/src/cloudai/registration.py b/src/cloudai/registration.py index 27443b0a5..32961adce 100644 --- a/src/cloudai/registration.py +++ b/src/cloudai/registration.py @@ -47,15 +47,8 @@ def register_all(): from cloudai.core import Registry from cloudai.models.scenario import ReportConfig from cloudai.report_generator.training import TrainingReporter - from cloudai.reporter import ( - DSEReporter, - JUnitReporter, - PerTestReporter, - ResultsUploadConfig, - ResultsUploadReporter, - StatusReporter, - TarballReporter, - ) + 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 @@ -351,7 +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("results_upload", ResultsUploadReporter, ResultsUploadConfig(enable=False)) + Registry().add_scenario_report("s3", S3UploadReporter, S3UploadConfig(enable=False)) Registry().add_scenario_report( "nixl_bench_summary", NIXLBenchComparisonReport, diff --git a/src/cloudai/reporter.py b/src/cloudai/reporter.py index e8b1aa939..918fe26fb 100644 --- a/src/cloudai/reporter.py +++ b/src/cloudai/reporter.py @@ -16,7 +16,6 @@ import contextlib import logging -import os import tarfile import xml.etree.ElementTree as ET from dataclasses import dataclass @@ -25,7 +24,6 @@ import jinja2 import toml -from pydantic import Field from rich import box from rich.console import Console from rich.table import Table @@ -34,8 +32,7 @@ from cloudai.report_generator.util import load_system_metadata from .core import CommandGenStrategy, Reporter, TestRun, case_name -from .models.scenario import ReportConfig, TestRunDetails -from .util.object_store import S3ObjectStore, join_key +from .models.scenario import TestRunDetails @dataclass @@ -316,89 +313,3 @@ def create_tarball(self, directory: Path) -> None: with tarfile.open(tarball_path, "w:gz") as tar: tar.add(directory, arcname=directory.name) logging.info(f"Created tarball at {tarball_path}") - - -class ResultsUploadConfig(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 = 8 - - -class ResultsUploadReporter(Reporter): - """Uploads the scenario results directory to object storage.""" - - def generate(self) -> None: - config = self.config - if not isinstance(config, ResultsUploadConfig): - logging.warning(f"Expected ResultsUploadConfig, got {type(config).__name__}, skipping results upload.") - return - - if not config.bucket: - logging.warning( - "Results 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 or is not accessible, 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, creating it if it does not exist. - - TarballReporter only produces a tarball when a test run failed, so we cannot - assume one is already present. - """ - tarball_path = Path(str(self.results_root) + ".tgz") - if not tarball_path.exists(): - 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)}") diff --git a/src/cloudai/s3_reporter.py b/src/cloudai/s3_reporter.py new file mode 100644 index 000000000..419fb9cb3 --- /dev/null +++ b/src/cloudai/s3_reporter.py @@ -0,0 +1,118 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2025-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 or is not accessible, 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, creating it if it does not exist. + + TarballReporter only produces a tarball when a test run failed, so we cannot + assume one is already present. + """ + tarball_path = Path(str(self.results_root) + ".tgz") + if not tarball_path.exists(): + 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)}") diff --git a/tests/test_init.py b/tests/test_init.py index 5932464d6..da198ab80 100644 --- a/tests/test_init.py +++ b/tests/test_init.py @@ -17,14 +17,8 @@ from cloudai.core import Registry from cloudai.report_generator.training import TrainingReporter -from cloudai.reporter import ( - DSEReporter, - JUnitReporter, - PerTestReporter, - ResultsUploadReporter, - StatusReporter, - TarballReporter, -) +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 @@ -294,7 +288,7 @@ def test_scenario_reports(): "junit", "dse", "tarball", - "results_upload", + "s3", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -311,7 +305,7 @@ def test_scenario_reports(): JUnitReporter, DSEReporter, TarballReporter, - ResultsUploadReporter, + S3UploadReporter, NIXLBenchComparisonReport, NixlEPComparisonReport, NcclComparisonReport, @@ -332,7 +326,7 @@ def test_report_configs(): "junit", "dse", "tarball", - "results_upload", + "s3", "nixl_bench_summary", "nixl_ep_comparison", "nccl_comparison", @@ -341,7 +335,7 @@ def test_report_configs(): "vllm_comparison", "sglang_comparison", ] - disabled_configs = {"junit", "results_upload"} # uploading off-box must be opt-in + 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 32d317943..d3bdad65e 100644 --- a/tests/test_reporter.py +++ b/tests/test_reporter.py @@ -21,7 +21,6 @@ from dataclasses import asdict from pathlib import Path from typing import Any -from unittest.mock import patch import pytest import toml @@ -36,8 +35,6 @@ JUnitReporter, PerTestReporter, ReportItem, - ResultsUploadConfig, - ResultsUploadReporter, StatusReporter, TarballReporter, ) @@ -54,7 +51,6 @@ ) from cloudai.systems.slurm.slurm_system import SlurmSystem from cloudai.systems.standalone.standalone_system import StandaloneSystem -from cloudai.util.object_store import UploadStats from cloudai.workloads.nccl_test import NCCLCmdArgs, NCCLTestDefinition @@ -454,7 +450,7 @@ def test_report_order() -> None: assert reports[-4][0] == "status" assert reports[-3][0] == "dse" assert reports[-2][0] == "tarball" - assert reports[-1][0] == "results_upload" + assert reports[-1][0] == "s3" def test_junit_reporter_generates_testcases_with_status_logs_and_duration( @@ -692,130 +688,3 @@ def test_dse_reporter( assert (slurm_system.output_path / "single-dse-scenario-dse-report.html").exists() assert (slurm_system.output_path / dse_case.name / "0" / f"{dse_case.name}.toml").exists() - - -class TestResultsUploadReporter: - """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) -> ResultsUploadReporter: - return ResultsUploadReporter( - slurm_system, - TestScenario(name="dummy", test_runs=[]), - results_dir, - ResultsUploadConfig(enable=True, **kwargs), - ) - - def test_uploads_tree(self, slurm_system: SlurmSystem, results_dir: Path) -> None: - with patch("cloudai.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.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.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.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_bucket_not_accessible_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: - with patch("cloudai.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.reporter.S3ObjectStore") as mock_store_cls: - store = mock_store_cls.return_value - - self.reporter(slurm_system, results_dir, bucket="my-bucket", upload_tree=False).generate() - - store.upload_directory.assert_not_called() - - 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.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_existing_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") - - with patch("cloudai.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" - - 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 = ResultsUploadConfig(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 ResultsUploadConfig(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.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_results_upload_runs_after_tarball() -> None: - order = [name for name, _ in Registry().ordered_scenario_reports()] - - assert order.index("results_upload") > order.index("tarball") - assert order[-1] == "results_upload" diff --git a/tests/test_s3_reporter.py b/tests/test_s3_reporter.py new file mode 100644 index 000000000..2b947b6af --- /dev/null +++ b/tests/test_s3_reporter.py @@ -0,0 +1,167 @@ +# SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES +# Copyright (c) 2025-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 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_bucket_not_accessible_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_existing_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") + + 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" + + 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" From bc785cfee05913ecf29221408551fc48ccbc7d40 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Wed, 7 Oct 2026 10:30:14 -0600 Subject: [PATCH 12/15] clean up uv.lock, handle 403 and 404 differently and handle stale tarball correctly in s3uploader --- src/cloudai/s3_reporter.py | 15 +++-- src/cloudai/util/object_store.py | 16 +++++- tests/test_s3_reporter.py | 32 ++++++++++- tests/util/test_object_store.py | 10 +++- uv.lock | 94 ++++++++++++++++---------------- 5 files changed, 110 insertions(+), 57 deletions(-) diff --git a/src/cloudai/s3_reporter.py b/src/cloudai/s3_reporter.py index 419fb9cb3..8382ded64 100644 --- a/src/cloudai/s3_reporter.py +++ b/src/cloudai/s3_reporter.py @@ -75,7 +75,7 @@ def generate(self) -> None: 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 or is not accessible, skipping results upload.") + 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) @@ -97,13 +97,14 @@ def generate(self) -> None: def upload_tarball(self, store: S3ObjectStore, key_prefix: str) -> None: """ - Upload a tarball of the results directory, creating it if it does not exist. + 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. + 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(): + 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 ) @@ -116,3 +117,9 @@ def upload_tarball(self, store: S3ObjectStore, key_prefix: str) -> None: 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/object_store.py b/src/cloudai/util/object_store.py index f58c1a2b6..8549b2004 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -172,13 +172,23 @@ def exists(self, key: str) -> bool: return True def bucket_exists(self) -> bool: - """Return whether the configured bucket exists and is accessible.""" + """ + 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 in (404, 403): - logging.debug(f"Bucket '{self.bucket}' is not accessible: {e}") + 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_s3_reporter.py b/tests/test_s3_reporter.py index 2b947b6af..c3887bf3f 100644 --- a/tests/test_s3_reporter.py +++ b/tests/test_s3_reporter.py @@ -14,6 +14,8 @@ # 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 @@ -81,7 +83,7 @@ def test_missing_results_dir_uploads_nothing(self, slurm_system: SlurmSystem, tm mock_store_cls.assert_not_called() - def test_bucket_not_accessible_uploads_nothing(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + 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 @@ -116,9 +118,10 @@ def test_tarball_created_when_absent(self, slurm_system: SlurmSystem, results_di tarball_path, "test_system/nccl-test_2025-04-16_14-27-45/nccl-test_2025-04-16_14-27-45.tgz" ) - def test_existing_tarball_is_reused(self, slurm_system: SlurmSystem, results_dir: Path) -> None: + 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() @@ -126,6 +129,31 @@ def test_existing_tarball_is_reused(self, slurm_system: SlurmSystem, results_dir 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") diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py index 0c8e7c338..d8defd94d 100644 --- a/tests/util/test_object_store.py +++ b/tests/util/test_object_store.py @@ -221,8 +221,16 @@ def test_s3_object_store_bucket_exists() -> None: 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) - assert store.bucket_exists() is False + 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: diff --git a/uv.lock b/uv.lock index 9228af30f..20163599a 100644 --- a/uv.lock +++ b/uv.lock @@ -132,30 +132,30 @@ wheels = [ [[package]] name = "boto3" -version = "1.43.91" +version = "1.43.109" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "botocore" }, { name = "jmespath" }, { name = "s3transfer" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/42/76/e5c5fb601b09273599bf9c1a678ea2bc645385e758fbc66a36c9614b6324/boto3-1.43.91.tar.gz", hash = "sha256:98643e500883bf6fcd13d04bb19da983bf97e5f1c600fc11470ce45437fa96b4", size = 112693, upload-time = "2026-09-09T19:24:38.927Z" } +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/24/3f/0288c2383455f1c94741fd8de819fb6d7360065274d905711a62fc6cd1f2/boto3-1.43.91-py3-none-any.whl", hash = "sha256:5ff948cfac8bff72227930e0894a8542731e1998634d1ffd8a074c20a90af919", size = 140023, upload-time = "2026-09-09T19:24:37.588Z" }, + { 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.91" +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/9a/1e/c6d26cb1124799e2bee54854813b02aa8a813c57335f97b0b8ed50f902ff/botocore-1.43.91.tar.gz", hash = "sha256:0f12bceb8c5d90a0c60f326e883bf674ded8c7d7c18ee0b559843b3a6c52f2ad", size = 16089153, upload-time = "2026-09-09T19:24:34.426Z" } +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/8f/1a/f672cc408d808035397c5d63a1b31a6a56ee95417a9d60d120d96349d182/botocore-1.43.91-py3-none-any.whl", hash = "sha256:f96363d4caf50bce45fe2fdc2d4d86b3bea7f5cd805169d53c1371c2dd4772b6", size = 15782250, upload-time = "2026-09-09T19:24:31.222Z" }, + { 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]] @@ -431,7 +431,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/66/54/eb9bfc647b19f2009dd5c7f5ec51c4e6ca831725f1aea7a993034f483147/contourpy-1.3.2.tar.gz", hash = "sha256:b6945942715a034c671b7fc54f9588126b0b8bf23db2696e3ca8328f3ff0ab54", size = 13466130, upload-time = "2025-04-15T17:47:53.79Z" } wheels = [ @@ -503,7 +503,7 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/58/01/1253e6698a07380cd31a736d248a3f2a50a7c88779a1813da27503cadc2a/contourpy-1.3.3.tar.gz", hash = "sha256:083e12155b210502d0bca491432bb04d56dc3432f95a979b429f2848c3dbe880", size = 13466174, upload-time = "2025-07-26T12:03:12.549Z" } @@ -731,7 +731,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions" }, + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -1111,7 +1111,7 @@ name = "importlib-metadata" version = "8.7.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "zipp" }, + { name = "zipp", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/f3/49/3b30cad09e7771a4982d9975a8cbf64f00d4a1ececb53297f1d9a7be1b10/importlib_metadata-8.7.1.tar.gz", hash = "sha256:49fef1ae6440c182052f407c8d34a68f72efc36db9ca90dc0113398f2fdde8bb", size = 57107, upload-time = "2025-12-21T10:00:19.278Z" } wheels = [ @@ -2158,7 +2158,7 @@ name = "roman-numerals-py" version = "4.1.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "roman-numerals" }, + { name = "roman-numerals", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/cb/b5/de96fca640f4f656eb79bbee0e79aeec52e3e0e359f8a3e6a0d366378b64/roman_numerals_py-4.1.0.tar.gz", hash = "sha256:f5d7b2b4ca52dd855ef7ab8eb3590f428c0b1ea480736ce32b01fef2a5f8daf9", size = 4274, upload-time = "2025-12-17T18:25:41.153Z" } wheels = [ @@ -2255,23 +2255,23 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "alabaster" }, - { name = "babel" }, - { name = "colorama", marker = "sys_platform == 'win32'" }, - { name = "docutils" }, - { name = "imagesize" }, - { name = "jinja2" }, - { name = "packaging" }, - { name = "pygments" }, - { name = "requests" }, - { name = "snowballstemmer" }, - { name = "sphinxcontrib-applehelp" }, - { name = "sphinxcontrib-devhelp" }, - { name = "sphinxcontrib-htmlhelp" }, - { name = "sphinxcontrib-jsmath" }, - { name = "sphinxcontrib-qthelp" }, - { name = "sphinxcontrib-serializinghtml" }, - { name = "tomli" }, + { name = "alabaster", marker = "python_full_version < '3.11'" }, + { name = "babel", marker = "python_full_version < '3.11'" }, + { name = "colorama", marker = "python_full_version < '3.11' and sys_platform == 'win32'" }, + { name = "docutils", marker = "python_full_version < '3.11'" }, + { name = "imagesize", marker = "python_full_version < '3.11'" }, + { name = "jinja2", marker = "python_full_version < '3.11'" }, + { name = "packaging", marker = "python_full_version < '3.11'" }, + { name = "pygments", marker = "python_full_version < '3.11'" }, + { name = "requests", marker = "python_full_version < '3.11'" }, + { name = "snowballstemmer", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-applehelp", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-devhelp", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-htmlhelp", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-jsmath", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-qthelp", marker = "python_full_version < '3.11'" }, + { name = "sphinxcontrib-serializinghtml", marker = "python_full_version < '3.11'" }, + { name = "tomli", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/6f/6d/be0b61178fe2cdcb67e2a92fc9ebb488e3c51c4f74a36a7824c0adf23425/sphinx-8.1.3.tar.gz", hash = "sha256:43c1911eecb0d3e161ad78611bc905d1ad0e523e4ddc202a58a821773dc4c927", size = 8184611, upload-time = "2024-10-13T20:27:13.93Z" } wheels = [ @@ -2288,23 +2288,23 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "alabaster" }, - { name = "babel" }, - { name = "colorama", marker = "sys_platform == 'win32'" }, - { name = "docutils" }, - { name = "imagesize" }, - { name = "jinja2" }, - { name = "packaging" }, - { name = "pygments" }, - { name = "requests" }, - { name = "roman-numerals-py" }, - { name = "snowballstemmer" }, - { name = "sphinxcontrib-applehelp" }, - { name = "sphinxcontrib-devhelp" }, - { name = "sphinxcontrib-htmlhelp" }, - { name = "sphinxcontrib-jsmath" }, - { name = "sphinxcontrib-qthelp" }, - { name = "sphinxcontrib-serializinghtml" }, + { name = "alabaster", marker = "python_full_version >= '3.11'" }, + { name = "babel", marker = "python_full_version >= '3.11'" }, + { name = "colorama", marker = "python_full_version >= '3.11' and sys_platform == 'win32'" }, + { name = "docutils", marker = "python_full_version >= '3.11'" }, + { name = "imagesize", marker = "python_full_version >= '3.11'" }, + { name = "jinja2", marker = "python_full_version >= '3.11'" }, + { name = "packaging", marker = "python_full_version >= '3.11'" }, + { name = "pygments", marker = "python_full_version >= '3.11'" }, + { name = "requests", marker = "python_full_version >= '3.11'" }, + { name = "roman-numerals-py", marker = "python_full_version >= '3.11'" }, + { name = "snowballstemmer", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-applehelp", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-devhelp", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-htmlhelp", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-jsmath", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-qthelp", marker = "python_full_version >= '3.11'" }, + { name = "sphinxcontrib-serializinghtml", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/38/ad/4360e50ed56cb483667b8e6dadf2d3fda62359593faabbe749a27c4eaca6/sphinx-8.2.3.tar.gz", hash = "sha256:398ad29dee7f63a75888314e9424d40f52ce5a6a87ae88e7071e80af296ec348", size = 8321876, upload-time = "2025-03-02T22:31:59.658Z" } wheels = [ @@ -2332,7 +2332,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "sphinx", version = "8.1.3", source = { registry = "https://pypi.org/simple" } }, + { name = "sphinx", version = "8.1.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/2b/69/b34e0cb5336f09c6866d53b4a19d76c227cdec1bbc7ac4de63ca7d58c9c7/sphinx_design-0.6.1.tar.gz", hash = "sha256:b44eea3719386d04d765c1a8257caca2b3e6f8421d7b3a5e742c0fd45f84e632", size = 2193689, upload-time = "2024-08-02T13:48:44.277Z" } wheels = [ @@ -2349,7 +2349,7 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "sphinx", version = "8.2.3", source = { registry = "https://pypi.org/simple" } }, + { name = "sphinx", version = "8.2.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/13/7b/804f311da4663a4aecc6cf7abd83443f3d4ded970826d0c958edc77d4527/sphinx_design-0.7.0.tar.gz", hash = "sha256:d2a3f5b19c24b916adb52f97c5f00efab4009ca337812001109084a740ec9b7a", size = 2203582, upload-time = "2026-01-19T13:12:53.297Z" } wheels = [ From 299e97d198e1054b95b9792c0a86cecf3351b3b1 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Wed, 7 Oct 2026 12:04:42 -0600 Subject: [PATCH 13/15] Fix copyright year in new S3 upload files The copyright header check derives the year from git history; these files were first committed in 2026, so the header must be 2026 rather than 2025-2026. Co-Authored-By: Claude Sonnet 5.5 --- src/cloudai/s3_reporter.py | 2 +- src/cloudai/util/object_store.py | 2 +- tests/test_s3_reporter.py | 2 +- tests/util/test_object_store.py | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/src/cloudai/s3_reporter.py b/src/cloudai/s3_reporter.py index 8382ded64..8570ffb44 100644 --- a/src/cloudai/s3_reporter.py +++ b/src/cloudai/s3_reporter.py @@ -1,5 +1,5 @@ # SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES -# Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# 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"); diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py index 8549b2004..6e1048f42 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -1,5 +1,5 @@ # SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES -# Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# 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"); diff --git a/tests/test_s3_reporter.py b/tests/test_s3_reporter.py index c3887bf3f..223b21d74 100644 --- a/tests/test_s3_reporter.py +++ b/tests/test_s3_reporter.py @@ -1,5 +1,5 @@ # SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES -# Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# 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"); diff --git a/tests/util/test_object_store.py b/tests/util/test_object_store.py index d8defd94d..7b7072fd4 100644 --- a/tests/util/test_object_store.py +++ b/tests/util/test_object_store.py @@ -1,5 +1,5 @@ # SPDX-FileCopyrightText: NVIDIA CORPORATION & AFFILIATES -# Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# 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"); From c4a5b5872d9f307d57e3a692e9758d2de94766f7 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Wed, 7 Oct 2026 12:18:42 -0600 Subject: [PATCH 14/15] clean uv.lock --- uv.lock | 82 ++++++++++++++++++++++++++++----------------------------- 1 file changed, 41 insertions(+), 41 deletions(-) diff --git a/uv.lock b/uv.lock index 20163599a..4317be098 100644 --- a/uv.lock +++ b/uv.lock @@ -431,7 +431,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, ] sdist = { url = "https://files.pythonhosted.org/packages/66/54/eb9bfc647b19f2009dd5c7f5ec51c4e6ca831725f1aea7a993034f483147/contourpy-1.3.2.tar.gz", hash = "sha256:b6945942715a034c671b7fc54f9588126b0b8bf23db2696e3ca8328f3ff0ab54", size = 13466130, upload-time = "2025-04-15T17:47:53.79Z" } wheels = [ @@ -503,7 +503,7 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, { name = "numpy", version = "2.5.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/58/01/1253e6698a07380cd31a736d248a3f2a50a7c88779a1813da27503cadc2a/contourpy-1.3.3.tar.gz", hash = "sha256:083e12155b210502d0bca491432bb04d56dc3432f95a979b429f2848c3dbe880", size = 13466174, upload-time = "2025-07-26T12:03:12.549Z" } @@ -731,7 +731,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions", marker = "python_full_version < '3.11'" }, + { name = "typing-extensions" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -1111,7 +1111,7 @@ name = "importlib-metadata" version = "8.7.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "zipp", marker = "python_full_version < '3.11'" }, + { name = "zipp" }, ] sdist = { url = "https://files.pythonhosted.org/packages/f3/49/3b30cad09e7771a4982d9975a8cbf64f00d4a1ececb53297f1d9a7be1b10/importlib_metadata-8.7.1.tar.gz", hash = "sha256:49fef1ae6440c182052f407c8d34a68f72efc36db9ca90dc0113398f2fdde8bb", size = 57107, upload-time = "2025-12-21T10:00:19.278Z" } wheels = [ @@ -2158,7 +2158,7 @@ name = "roman-numerals-py" version = "4.1.0" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "roman-numerals", marker = "python_full_version >= '3.11'" }, + { name = "roman-numerals" }, ] sdist = { url = "https://files.pythonhosted.org/packages/cb/b5/de96fca640f4f656eb79bbee0e79aeec52e3e0e359f8a3e6a0d366378b64/roman_numerals_py-4.1.0.tar.gz", hash = "sha256:f5d7b2b4ca52dd855ef7ab8eb3590f428c0b1ea480736ce32b01fef2a5f8daf9", size = 4274, upload-time = "2025-12-17T18:25:41.153Z" } wheels = [ @@ -2255,23 +2255,23 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "alabaster", marker = "python_full_version < '3.11'" }, - { name = "babel", marker = "python_full_version < '3.11'" }, - { name = "colorama", marker = "python_full_version < '3.11' and sys_platform == 'win32'" }, - { name = "docutils", marker = "python_full_version < '3.11'" }, - { name = "imagesize", marker = "python_full_version < '3.11'" }, - { name = "jinja2", marker = "python_full_version < '3.11'" }, - { name = "packaging", marker = "python_full_version < '3.11'" }, - { name = "pygments", marker = "python_full_version < '3.11'" }, - { name = "requests", marker = "python_full_version < '3.11'" }, - { name = "snowballstemmer", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-applehelp", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-devhelp", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-htmlhelp", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-jsmath", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-qthelp", marker = "python_full_version < '3.11'" }, - { name = "sphinxcontrib-serializinghtml", marker = "python_full_version < '3.11'" }, - { name = "tomli", marker = "python_full_version < '3.11'" }, + { name = "alabaster" }, + { name = "babel" }, + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "docutils" }, + { name = "imagesize" }, + { name = "jinja2" }, + { name = "packaging" }, + { name = "pygments" }, + { name = "requests" }, + { name = "snowballstemmer" }, + { name = "sphinxcontrib-applehelp" }, + { name = "sphinxcontrib-devhelp" }, + { name = "sphinxcontrib-htmlhelp" }, + { name = "sphinxcontrib-jsmath" }, + { name = "sphinxcontrib-qthelp" }, + { name = "sphinxcontrib-serializinghtml" }, + { name = "tomli" }, ] sdist = { url = "https://files.pythonhosted.org/packages/6f/6d/be0b61178fe2cdcb67e2a92fc9ebb488e3c51c4f74a36a7824c0adf23425/sphinx-8.1.3.tar.gz", hash = "sha256:43c1911eecb0d3e161ad78611bc905d1ad0e523e4ddc202a58a821773dc4c927", size = 8184611, upload-time = "2024-10-13T20:27:13.93Z" } wheels = [ @@ -2288,23 +2288,23 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "alabaster", marker = "python_full_version >= '3.11'" }, - { name = "babel", marker = "python_full_version >= '3.11'" }, - { name = "colorama", marker = "python_full_version >= '3.11' and sys_platform == 'win32'" }, - { name = "docutils", marker = "python_full_version >= '3.11'" }, - { name = "imagesize", marker = "python_full_version >= '3.11'" }, - { name = "jinja2", marker = "python_full_version >= '3.11'" }, - { name = "packaging", marker = "python_full_version >= '3.11'" }, - { name = "pygments", marker = "python_full_version >= '3.11'" }, - { name = "requests", marker = "python_full_version >= '3.11'" }, - { name = "roman-numerals-py", marker = "python_full_version >= '3.11'" }, - { name = "snowballstemmer", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-applehelp", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-devhelp", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-htmlhelp", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-jsmath", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-qthelp", marker = "python_full_version >= '3.11'" }, - { name = "sphinxcontrib-serializinghtml", marker = "python_full_version >= '3.11'" }, + { name = "alabaster" }, + { name = "babel" }, + { name = "colorama", marker = "sys_platform == 'win32'" }, + { name = "docutils" }, + { name = "imagesize" }, + { name = "jinja2" }, + { name = "packaging" }, + { name = "pygments" }, + { name = "requests" }, + { name = "roman-numerals-py" }, + { name = "snowballstemmer" }, + { name = "sphinxcontrib-applehelp" }, + { name = "sphinxcontrib-devhelp" }, + { name = "sphinxcontrib-htmlhelp" }, + { name = "sphinxcontrib-jsmath" }, + { name = "sphinxcontrib-qthelp" }, + { name = "sphinxcontrib-serializinghtml" }, ] sdist = { url = "https://files.pythonhosted.org/packages/38/ad/4360e50ed56cb483667b8e6dadf2d3fda62359593faabbe749a27c4eaca6/sphinx-8.2.3.tar.gz", hash = "sha256:398ad29dee7f63a75888314e9424d40f52ce5a6a87ae88e7071e80af296ec348", size = 8321876, upload-time = "2025-03-02T22:31:59.658Z" } wheels = [ @@ -2332,7 +2332,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "sphinx", version = "8.1.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "sphinx", version = "8.1.3", source = { registry = "https://pypi.org/simple" } }, ] sdist = { url = "https://files.pythonhosted.org/packages/2b/69/b34e0cb5336f09c6866d53b4a19d76c227cdec1bbc7ac4de63ca7d58c9c7/sphinx_design-0.6.1.tar.gz", hash = "sha256:b44eea3719386d04d765c1a8257caca2b3e6f8421d7b3a5e742c0fd45f84e632", size = 2193689, upload-time = "2024-08-02T13:48:44.277Z" } wheels = [ @@ -2349,7 +2349,7 @@ resolution-markers = [ "python_full_version == '3.11.*'", ] dependencies = [ - { name = "sphinx", version = "8.2.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11'" }, + { name = "sphinx", version = "8.2.3", source = { registry = "https://pypi.org/simple" } }, ] sdist = { url = "https://files.pythonhosted.org/packages/13/7b/804f311da4663a4aecc6cf7abd83443f3d4ded970826d0c958edc77d4527/sphinx_design-0.7.0.tar.gz", hash = "sha256:d2a3f5b19c24b916adb52f97c5f00efab4009ca337812001109084a740ec9b7a", size = 2203582, upload-time = "2026-01-19T13:12:53.297Z" } wheels = [ From ab99744a50d52e0c22996c285f97dfd7469aa359 Mon Sep 17 00:00:00 2001 From: Jakob Johnson Date: Wed, 7 Oct 2026 12:33:51 -0600 Subject: [PATCH 15/15] update docs with env var info and remove module level comment --- doc/reporting.rst | 38 ++++++++++++++++++++++++++++---- src/cloudai/util/object_store.py | 2 -- 2 files changed, 34 insertions(+), 6 deletions(-) diff --git a/doc/reporting.rst b/doc/reporting.rst index 68740abb2..0d72b937c 100644 --- a/doc/reporting.rst +++ b/doc/reporting.rst @@ -242,16 +242,46 @@ Configuration options: - Upload each file individually, preserving relative paths. * - ``upload_tarball`` - ``false`` - - Also upload a ``.tgz`` of the whole directory, creating it if absent. + - 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. -Destination fields fall back to the environment variable shown above when not set in -TOML, so a cluster-wide default can come from the environment while an individual -scenario can still override it. +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``, diff --git a/src/cloudai/util/object_store.py b/src/cloudai/util/object_store.py index 6e1048f42..5a7706398 100644 --- a/src/cloudai/util/object_store.py +++ b/src/cloudai/util/object_store.py @@ -14,8 +14,6 @@ # See the License for the specific language governing permissions and # limitations under the License. -"""Generic object storage interface and an S3 implementation.""" - from __future__ import annotations import fnmatch