diff --git a/.github/workflows/regression-whisk.yml b/.github/workflows/regression-whisk.yml index c8250027..7d6ed992 100644 --- a/.github/workflows/regression-whisk.yml +++ b/.github/workflows/regression-whisk.yml @@ -132,9 +132,7 @@ jobs: - name: Configure OpenWhisk run: | WORKER_IP=$(kubectl get node kind-worker -o jsonpath='{.status.addresses[0].address}') - HOST_IP=$(hostname -I | awk '{print $1}') echo "WORKER_IP=${WORKER_IP}" >> $GITHUB_ENV - echo "HOST_IP=${HOST_IP}" >> $GITHUB_ENV git clone --depth 1 https://github.com/apache/openwhisk-deploy-kube.git /tmp/ow @@ -159,25 +157,16 @@ jobs: - name: Create OpenWhisk regression config run: | - jq \ - --arg host_ip "${HOST_IP}" \ - ' - .object.minio.address = ($host_ip + ":" + (.object.minio.mapped_port | tostring)) - | .nosql.scylladb.address = ($host_ip + ":" + (.nosql.scylladb.mapped_port | tostring)) - ' storage.json > storage-openwhisk.json - jq \ --arg language "${LANGUAGE}" \ --arg version "${LANGUAGE_VERSION}" \ --arg architecture "${ARCHITECTURE}" \ --arg registry "localhost:${REGISTRY_PORT}" \ - --slurpfile storage storage-openwhisk.json \ ' .experiments.architecture = $architecture | .experiments.runtime.language = $language | .experiments.runtime.version = $version | .deployment.openwhisk.docker_registry.registry = $registry - | .deployment.openwhisk.storage = $storage[0] ' configs/openwhisk.json > openwhisk-regression.json - name: Run OpenWhisk regression @@ -188,6 +177,7 @@ jobs: set -o pipefail uv run sebs benchmark regression test \ --config openwhisk-regression.json \ + --storage-configuration storage.json \ --deployment openwhisk \ --language ${LANGUAGE} \ --language-version ${LANGUAGE_VERSION} \ @@ -267,7 +257,6 @@ jobs: diagnostics/ regression-cache/ storage.json - storage-openwhisk.json openwhisk-regression.json regression_*.json if-no-files-found: ignore diff --git a/CHANGELOG.md b/CHANGELOG.md index 5881fa87..a1933961 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,16 +2,24 @@ ### Features +* Support for [RustFS](https://github.com/rustfs/rustfs) as an S3-compatible object storage for local and OpenWhisk deployments, next to Minio. The self-hosted object storage implementation is now shared between backends. + ### Bug Fixes * Change to input of 120.uploader benchmark to conform with new Wikipedia policies (#308) +* Pull Minio images from quay.io, as the `minio/minio` repository was removed from Docker Hub. + +* Restore cached ScyllaDB configuration independently of Minio; make the OpenWhisk `shutdownStorage` option work again. + ### Improvements * Support for multiple variants of the same programming language (#286). * Support for new versions of Python and Java on AWS (#311). +* Self-hosted storage (Minio, ScyllaDB) advertises an externally reachable address to functions, detected automatically or set with `--external-address` when starting the storage, and overridable with `--storage-address` when running benchmarks. This removes the manual editing of storage addresses for OpenWhisk (#229). + ### Deprecations ### Contributors diff --git a/configs/example.json b/configs/example.json index a3f3f9b2..799de59f 100644 --- a/configs/example.json +++ b/configs/example.json @@ -89,14 +89,19 @@ }, "local": { "storage": { - "address": "", - "mapped_port": -1, - "access_key": "", - "secret_key": "", - "instance_id": "", - "input_buckets": [], - "output_buckets": [], - "type": "minio" + "object": { + "type": "minio", + "minio": { + "address": "", + "external_address": "", + "mapped_port": -1, + "access_key": "", + "secret_key": "", + "instance_id": "", + "input_buckets": [], + "output_buckets": [] + } + } } }, "openwhisk": { @@ -112,14 +117,19 @@ "password": "" }, "storage": { - "address": "", - "mapped_port": -1, - "access_key": "", - "secret_key": "", - "instance_id": "", - "input_buckets": [], - "output_buckets": [], - "type": "minio" + "object": { + "type": "minio", + "minio": { + "address": "", + "external_address": "", + "mapped_port": -1, + "access_key": "", + "secret_key": "", + "instance_id": "", + "input_buckets": [], + "output_buckets": [] + } + } } } } diff --git a/configs/openwhisk.json b/configs/openwhisk.json index e02f3553..c774a7a2 100644 --- a/configs/openwhisk.json +++ b/configs/openwhisk.json @@ -35,7 +35,8 @@ "secret_key": "", "instance_id": "", "input_buckets": [], - "output_buckets": [] + "output_buckets": [], + "external_address": "" } } } diff --git a/configs/storage-rustfs.json b/configs/storage-rustfs.json new file mode 100644 index 00000000..917f0b13 --- /dev/null +++ b/configs/storage-rustfs.json @@ -0,0 +1,20 @@ +{ + "object": { + "type": "rustfs", + "rustfs": { + "mapped_port": 9011, + "version": "1.0.0-rc.6", + "data_volume": "rustfs-volume" + } + }, + "nosql": { + "type": "scylladb", + "scylladb": { + "mapped_port": 9012, + "version": "6.0", + "cpus": 1, + "memory": "750", + "data_volume": "scylladb-volume" + } + } +} diff --git a/docs/platforms.md b/docs/platforms.md index b01937c4..efd0698a 100644 --- a/docs/platforms.md +++ b/docs/platforms.md @@ -461,5 +461,7 @@ To use that feature in SeBS, set the `experimentalManifest` flag to true. ### Storage +Start the storage with `sebs storage start` and pass the generated configuration to benchmark commands with `--storage-configuration`; see the [storage documentation](storage.md) for details, including how to override the address advertised to functions with `--storage-address`. + OpenWhisk has a `shutdownStorage` switch that controls the behavior of SeBS. -When set to true, SeBS will remove the Minio instance after finishing all work. +When set to true, SeBS will stop the Minio and ScyllaDB instances after finishing all work. diff --git a/docs/storage.md b/docs/storage.md index 73645a57..76876362 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -7,12 +7,20 @@ SeBS will automatically allocate resources and configure them. With open-source platforms like OpenWhisk or local deployment, SeBS needs a self-hosted storage instance. In this document, we explain how to deploy and configure storage systems for benchmarking with SeBS. -We use [Minio](https://github.com/minio/minio), a high-performance and S3-compatible object storage, and [ScyllaDB](https://github.com/scylladb/scylladb) -with an adapter that provides a DynamoDB-compatible interface. +For object storage, we support two S3-compatible systems: [Minio](https://github.com/minio/minio) and [RustFS](https://github.com/rustfs/rustfs). +For NoSQL storage, we use [ScyllaDB](https://github.com/scylladb/scylladb) with an adapter that provides a DynamoDB-compatible interface. The storage instance is deployed as a Docker container and can be retained across multiple experiments. While we provide a default configuration that automatically deploys each storage instance, you can deploy them on any cloud resource and adapt the configuration to fit your needs. +## Object Storage Backends + +Benchmark functions access object storage through the S3 API, so both backends are interchangeable and no benchmark code changes when switching between them. +Select the backend with the `type` field of the object storage configuration; the default configuration files are `configs/storage.json` for Minio and `configs/storage-rustfs.json` for RustFS. + +* **Minio** is the established default. Its community edition is no longer maintained and its images were removed from Docker Hub; SeBS pulls the pinned version from `quay.io/minio/minio`. +* **RustFS** is an actively developed, Apache-2.0 licensed alternative. Its data is kept in a named Docker volume, since the container runs as a fixed unprivileged user. At the time of writing, RustFS has not published a stable release yet, so we pin a release candidate. + ## Starting Storage Services You can start the necessary storage services using the `storage` command in SeBS: @@ -37,6 +45,7 @@ This file contains all the necessary information to connect to the storage servi "type": "minio", "minio": { "address": "172.17.0.2:9000", + "external_address": "10.10.1.15:9011", "mapped_port": 9011, "access_key": "XXX", "secret_key": "XXX", @@ -52,6 +61,7 @@ This file contains all the necessary information to connect to the storage servi "type": "scylladb", "scylladb": { "address": "172.17.0.3:8000", + "external_address": "10.10.1.15:9012", "mapped_port": 9012, "alternator_port": 8000, "access_key": "None", @@ -67,65 +77,51 @@ This file contains all the necessary information to connect to the storage servi } ``` -As we can see, the Minio container is running on the default Docker bridge network with address `172.17.0.2` and uses port `9000`. -The default configuration maps the container's port to the host, making the storage instance available directly without referring to the container's IP address. Minio is mapped to port 9011, and ScyllaDB is mapped to port 9012. - -## Network Configuration - -The storage instance must be accessible from the host network, and in some cases, from external networks. -For example, the storage can be deployed on a separate virtual machine or container. -Furthermore, even on a local machine, it's necessary to configure the network address, as OpenWhisk functions -are running isolated from the host network and won't be able to reach other containers running on the Docker bridge. +Each storage instance has two addresses: -When using Minio with cloud-hosted FaaS platforms like OpenWhisk or for local deployment, you need to ensure that the functions can reach the storage instance. -By default, the container runs on the Docker bridge network with an address (e.g., `172.17.0.2`) that is not accessible from outside the host. -Even when deploying both OpenWhisk and storage on the same system, the local bridge network is not accessible from the Kubernetes cluster. -To make it accessible, functions need to use the public IP address of the machine hosting the container instance and the mapped port. -You can typically find an externally accessible address via `ip addr`, and then replace the storage's address with the external address of the machine and the mapped port. +* `address` is used by SeBS itself, e.g., to upload benchmark inputs. On Linux, this is the container's address on the default Docker bridge network (`172.17.0.2`) and the container's port (`9000`). Functions of the local deployment run on the same bridge network and use this address as well. +* `external_address` is advertised to benchmark functions that run outside of the Docker bridge network, e.g., in a Kubernetes cluster hosting OpenWhisk. It combines the IP address of the machine with the port mapped on the host: Minio is mapped to port 9011, and ScyllaDB to port 9012. -For example, for an external address `10.10.1.15` (a LAN-local address on CloudLab) and mapped port `9011`, set the SeBS configuration as follows: +The external address is detected automatically as the IP address of the host's default network interface, and SeBS verifies that the storage answers on it. To use a different interface or a hostname, pass the `--external-address` flag when starting the storage: ```bash -# For a LAN-local address (e.g., on CloudLab) -jq --slurpfile file1 storage.json '.deployment.openwhisk.storage = $file1[0] | .deployment.openwhisk.storage.object.minio.address = "10.10.1.15:9011"' configs/example.json > configs/openwhisk.json +sebs storage start all configs/storage.json --output-json storage.json --external-address 10.10.1.15 ``` -You can validate the configuration of Minio with an HTTP request by using `curl`: +> [!WARNING] +> The mapped ports are bound on all interfaces of the host. On a machine with a public IP address, restrict access to these ports with a firewall or use a private address. + +## Network Configuration + +To use the deployed storage with a benchmark, pass the generated configuration file with the `--storage-configuration` flag. +The storage configuration is merged into the deployment section of the SeBS configuration, so no manual editing of JSON files is needed: ```bash -$ curl -i 10.10.1.15:9011/minio/health/live -HTTP/1.1 200 OK -Accept-Ranges: bytes -Content-Length: 0 -Content-Security-Policy: block-all-mixed-content -Server: MinIO -Strict-Transport-Security: max-age=31536000; includeSubDomains -Vary: Origin -X-Amz-Request-Id: 16F3D9B9FDFFA340 -X-Content-Type-Options: nosniff -X-Xss-Protection: 1; mode=block -Date: Mon, 30 May 2022 10:01:21 GMT +sebs benchmark invoke 210.thumbnailer test --config configs/openwhisk.json --storage-configuration storage.json ``` -If you use benchmarks relying on NoSQL storage (ScyllaDB), then you need to apply the same change to reflect the different address as well. -Here, we again assume the external IP address of the system is `10.10.1.15`, and the mapped port changes to `9012`. +Functions running in OpenWhisk or another Kubernetes-based platform cannot reach the Docker bridge network of the host, even when the cluster runs on the same machine. +They connect to the storage through the external address, which is detected when starting the storage. +If the detected address is not reachable from the functions, e.g., because the machine has multiple network interfaces or the storage runs on a different host, override it without changing any files: ```bash -# For a LAN-local address (e.g., on CloudLab) -jq '.deployment.openwhisk.storage.nosql.scylladb.address = "10.10.1.15:9012"' configs/openwhisk.json | sponge configs/openwhisk.json +sebs benchmark invoke 210.thumbnailer test --config configs/openwhisk.json --storage-configuration storage.json --storage-address 10.10.1.15 ``` -You can validate the configuration of ScyllaDB with an HTTP request by using `curl`: +The override applies to all storage instances, each with its own mapped port. Alternatively, provide the address once when starting the storage with `--external-address`. + +You can validate that the storage is reachable with an HTTP request to Minio's health endpoint and ScyllaDB's root endpoint: ```bash -curl -i 10.10.1.15:9012 +$ curl -i 10.10.1.15:9011/minio/health/live HTTP/1.1 200 OK -Content-Length: 26 -Content-Type: text/plain -Date: Sun, 07 Dec 2025 14:07:29 GMT -Server: Seastar httpd +... +Server: MinIO -healthy: 192.168.0.20:9012 +$ curl -i 10.10.1.15:9012 +HTTP/1.1 200 OK +... +healthy: 10.10.1.15:9012 ``` ## Lifecycle Management @@ -147,4 +143,4 @@ sebs storage stop all storage.json Each storage service uses a Docker volume to persist data. The name of the volume is included in the storage configuration file under the `data_volume` field. In Minio, the volume is mapped to a physical location on the filesystem, and the directory can be removed once the experiments are finished. -For ScyllaDB, we use named Docker volumes that can be removed using Docker commands: `docker volume rm scylladb-volume`. +For RustFS and ScyllaDB, we use named Docker volumes that can be removed using Docker commands: `docker volume rm rustfs-volume scylladb-volume`. diff --git a/install.py b/install.py index 879ac16e..53a04c38 100755 --- a/install.py +++ b/install.py @@ -73,7 +73,7 @@ def execute(cmd, cwd=None): print("Install Python dependencies for local") execute(". {}/bin/activate && pip3 install -r requirements.local.txt".format(env_dir)) print("Initialize Docker image for local storage.") - execute("docker pull minio/minio:latest") + execute("docker pull quay.io/minio/minio:latest") # One of the installed dependencies causes a downgrade, which in turns breaks static typing. print("Update typing-extensions (resolving bug with mypy)") diff --git a/sebs/cli.py b/sebs/cli.py index d99daf94..1f051046 100755 --- a/sebs/cli.py +++ b/sebs/cli.py @@ -13,7 +13,7 @@ import os import sys import traceback -from typing import cast, List, Optional +from typing import cast, Dict, List, Optional import click import docker @@ -144,6 +144,62 @@ def wrapper(*args, **kwargs): return wrapper +def storage_params(func): + """Decorator that adds CLI parameters for user-deployed storage.""" + + @click.option( + "--storage-configuration", + type=str, + multiple=True, + help="JSON configuration of deployed storage, as written by 'sebs storage start'.", + ) + @click.option( + "--storage-address", + default=None, + type=str, + help="Override the address (IP or hostname, optional port) that benchmark functions " + "use to reach self-hosted storage. Applied to all storage types.", + ) + @functools.wraps(func) + def wrapper(*args, **kwargs): + """Internal Click wrapper.""" + return func(*args, **kwargs) + + return wrapper + + +def override_storage_address( + config_obj: dict, deployment: str, storage_address: str +) -> Dict[str, str]: + """Override the externally reachable address of all self-hosted storage instances. + + Each storage instance receives the given host combined with its own mapped + port, unless the user already provided a port. + + Args: + config_obj: Full SeBS configuration + deployment: Name of the selected deployment + storage_address: IP address or hostname, optionally with a port + + Returns: + Dict[str, str]: Applied address per storage type + """ + from sebs.utils import resolve_external_address + + applied: Dict[str, str] = {} + storage_cfg = config_obj.get("deployment", {}).get(deployment, {}).get("storage", {}) + for storage_type, type_cfg in storage_cfg.items(): + impl = type_cfg.get("type") + if impl is None or impl not in type_cfg: + continue + impl_cfg = type_cfg[impl] + impl_cfg["external_address"] = resolve_external_address( + storage_address, impl_cfg.get("mapped_port", -1) + ) + applied[storage_type] = impl_cfg["external_address"] + return applied + + def parse_common_params( config, output_dir, @@ -163,6 +219,7 @@ def parse_common_params( initialize_deployment: bool = True, ignore_cache: bool = False, storage_configuration: Optional[List[str]] = None, + storage_address: Optional[str] = None, ): """Parse and process common CLI parameters, initialize SeBS and deployment clients. @@ -207,7 +264,12 @@ def parse_common_params( sebs_client.logging.info(f"Loading storage configuration from {cfg_f}") cfg = json.load(open(cfg_f, "r")) - append_nested_dict(config_obj, ["deployment", deployment, "storage"], cfg) + append_nested_dict(config_obj, ["deployment", selected_deployment, "storage"], cfg) + + if storage_address is not None: + overrides = override_storage_address(config_obj, selected_deployment, storage_address) + for storage_type, address in overrides.items(): + sebs_client.logging.info(f"Using storage address {address} for {storage_type} storage") if initialize_deployment: deployment_client = sebs_client.get_deployment( @@ -273,12 +335,7 @@ def benchmark(): type=str, help="Attach prefix to generated Docker image tag.", ) -@click.option( - "--storage-configuration", - type=str, - multiple=True, - help="JSON configuration of deployed storage.", -) +@storage_params @click.option( "--validate/--no-validate", default=False, @@ -483,12 +540,7 @@ def package( type=str, help="Run only the selected benchmark.", ) -@click.option( - "--storage-configuration", - type=str, - multiple=True, - help="JSON configuration of deployed storage.", -) +@storage_params @click.option( "--selected-architecture/--all-architectures", type=bool, @@ -506,6 +558,7 @@ def regression( benchmark_input_size, benchmark_name, storage_configuration, + storage_address, selected_architecture, filter_output, **kwargs, @@ -523,6 +576,7 @@ def regression( (config, output_dir, logging_filename, sebs_client, _) = parse_common_params( initialize_deployment=False, storage_configuration=storage_configuration, + storage_address=storage_address, **kwargs, ) architecture = config["experiments"]["architecture"] if selected_architecture else None @@ -563,8 +617,21 @@ def storage(): default=True, help="Remove containers after stopping.", ) -def storage_start(storage, config, output_json, remove_containers): - """Start local storage instances (object storage, NoSQL, or both).""" +@click.option( + "--external-address", + default=None, + type=str, + help="Address (IP or hostname) advertised to benchmark functions. Each storage instance " + "appends its mapped port. Defaults to the IP of the host's default network interface.", +) +def storage_start(storage, config, output_json, remove_containers, external_address): + """Start local storage instances (object storage, NoSQL, or both). + + The written configuration contains two addresses per storage instance: the + address used by SeBS on this host, and the external address advertised to + benchmark functions, which is required for functions running outside of the + Docker bridge network, e.g., in a Kubernetes cluster. + """ import docker sebs.utils.global_logging() @@ -578,11 +645,14 @@ def storage_start(storage, config, output_json, remove_containers): storage_config = sebs.SeBS.get_storage_config_implementation(storage_type_enum) config = storage_config.deserialize(user_storage_config["object"][storage_type_name]) config.remove_containers = remove_containers + if external_address is not None: + config.external_address = external_address storage_instance = storage_type(docker.from_env(), None, None, True) storage_instance.config = config storage_instance.start() + storage_instance.check_external_address() user_storage_config["object"][storage_type_name] = storage_instance.serialize() else: @@ -596,10 +666,13 @@ def storage_start(storage, config, output_json, remove_containers): storage_config = sebs.SeBS.get_nosql_config_implementation(storage_type_enum) config = storage_config.deserialize(user_storage_config["nosql"][storage_type_name]) config.remove_containers = remove_containers + if external_address is not None: + config.external_address = external_address storage_instance = storage_type(docker.from_env(), None, config) storage_instance.start() + storage_instance.check_external_address() key, value = storage_instance.serialize() user_storage_config["nosql"][key] = value @@ -662,12 +735,7 @@ def local(): @click.argument("benchmark-input-size", type=click.Choice(["test", "small", "large"])) @click.argument("output", type=str) @click.option("--deployments", default=1, type=int, help="Number of deployed containers.") -@click.option( - "--storage-configuration", - type=str, - multiple=True, - help="JSON configuration of deployed storage.", -) +@storage_params @click.option( "--measure-interval", type=int, @@ -692,6 +760,7 @@ def start( output, deployments, storage_configuration, + storage_address, measure_interval, remove_containers, architecture, @@ -706,6 +775,7 @@ def start( update_storage=False, deployment="local", storage_configuration=storage_configuration, + storage_address=storage_address, system_variant="package", architecture=architecture, **kwargs, diff --git a/sebs/faas/nosql.py b/sebs/faas/nosql.py index 502b904d..81f3163b 100644 --- a/sebs/faas/nosql.py +++ b/sebs/faas/nosql.py @@ -128,13 +128,18 @@ def update_cache(self, benchmark: str): """ pass - def envs(self) -> dict: + def envs(self, external: bool = True) -> dict: """ Return a dictionary of environment variables that are required by functions to access this NoSQL storage (e.g., connection strings, table names). Default implementation returns an empty dictionary. Subclasses should override if they need to expose environment variables. + Args: + external: For self-hosted storage, advertise the externally reachable + address. Functions running on the same Docker bridge as the storage + (local deployment) must use the internal address instead. + Returns: dict: Dictionary of environment variables """ diff --git a/sebs/local/deployment.py b/sebs/local/deployment.py index b06fa787..09c31598 100644 --- a/sebs/local/deployment.py +++ b/sebs/local/deployment.py @@ -18,7 +18,8 @@ from sebs.cache import Cache from sebs.local.function import LocalFunction from sebs.local.config import LocalResources -from sebs.storage.minio import Minio, MinioConfig +from sebs.storage.resources import OBJECT_STORAGE_IMPLEMENTATIONS +from sebs.storage.s3compatible import S3CompatibleStorage from sebs.utils import serialize, LoggingBase @@ -27,7 +28,7 @@ class Deployment(LoggingBase): Attributes: _functions: List of deployed local functions - _storage: Optional Minio storage instance + _storage: Optional S3-compatible storage instance _inputs: List of function input configurations _memory_measurement_pids: PIDs of memory measurement processes _measurement_file: Path to memory measurement output file @@ -55,7 +56,7 @@ def __init__(self): """Initialize a new deployment.""" super().__init__() self._functions: List[LocalFunction] = [] - self._storage: Optional[Minio] + self._storage: Optional[S3CompatibleStorage] self._inputs: List[dict] = [] self._memory_measurement_pids: List[int] = [] self._measurement_file: Optional[str] = None @@ -80,11 +81,11 @@ def add_input(self, func_input: dict) -> None: """ self._inputs.append(func_input) - def set_storage(self, storage: Minio) -> None: + def set_storage(self, storage: S3CompatibleStorage) -> None: """Set the storage instance for the deployment. Args: - storage: Minio storage instance to use + storage: Storage instance to use """ self._storage = storage @@ -135,8 +136,11 @@ def deserialize(path: str, cache_client: Cache) -> "Deployment": if "memory_measurements" in input_data: deployment._memory_measurement_pids = input_data["memory_measurements"]["pids"] deployment._measurement_file = input_data["memory_measurements"]["file"] - deployment._storage = Minio.deserialize( - MinioConfig.deserialize(input_data["storage"]), cache_client, LocalResources() + storage_type = OBJECT_STORAGE_IMPLEMENTATIONS[input_data["storage"]["type"]] + deployment._storage = storage_type.deserialize( + storage_type.CONFIG_TYPE.deserialize(input_data["storage"]), + cache_client, + LocalResources(), ) return deployment diff --git a/sebs/local/local.py b/sebs/local/local.py index 99c03c7a..19ef80d3 100644 --- a/sebs/local/local.py +++ b/sebs/local/local.py @@ -270,12 +270,15 @@ def _start_container( "CONTAINER_GID": str(os.getgid()), "CONTAINER_USER": self._system_config.username(self.name(), code_package.language_name), } + # Function containers share the Docker bridge with the storage containers, + # so they use the internal addresses. if self.config.resources.storage_config: - environment = {**self.config.resources.storage_config.envs(), **environment} + storage_envs = self.config.resources.storage_config.envs(external=False) + environment = {**storage_envs, **environment} if code_package.uses_nosql: nosql_storage = self.system_resources.get_nosql_storage() - environment = {**environment, **nosql_storage.envs()} + environment = {**environment, **nosql_storage.envs(external=False)} for original_name, actual_name in nosql_storage.get_tables( code_package.benchmark diff --git a/sebs/openwhisk/function.py b/sebs/openwhisk/function.py index 79cc776a..17ca2032 100644 --- a/sebs/openwhisk/function.py +++ b/sebs/openwhisk/function.py @@ -13,7 +13,8 @@ from sebs.benchmark import Benchmark from sebs.faas.function import Function, FunctionConfig, Runtime -from sebs.storage.config import MinioConfig, ScyllaDBConfig +from sebs.storage.config import S3CompatibleConfig, ScyllaDBConfig +from sebs.storage.resources import OBJECT_STORAGE_IMPLEMENTATIONS @dataclass @@ -28,7 +29,7 @@ class OpenWhiskFunctionConfig(FunctionConfig): Attributes: docker_image: Docker image URI used for the function deployment namespace: OpenWhisk namespace (default: "_" for default namespace) - object_storage: Minio object storage configuration if required + object_storage: S3-compatible object storage configuration if required nosql_storage: ScyllaDB NoSQL storage configuration if required Note: @@ -39,7 +40,7 @@ class OpenWhiskFunctionConfig(FunctionConfig): docker_image: str = "" namespace: str = "_" - object_storage: Optional[MinioConfig] = None + object_storage: Optional[S3CompatibleConfig] = None nosql_storage: Optional[ScyllaDBConfig] = None @staticmethod @@ -57,7 +58,8 @@ def deserialize(data: Dict[str, Any]) -> OpenWhiskFunctionConfig: data = {k: v for k, v in data.items() if k in keys} data["runtime"] = Runtime.deserialize(data["runtime"]) if data["object_storage"] is not None: - data["object_storage"] = MinioConfig.deserialize(data["object_storage"]) + storage_type = OBJECT_STORAGE_IMPLEMENTATIONS[data["object_storage"]["type"]] + data["object_storage"] = storage_type.CONFIG_TYPE.deserialize(data["object_storage"]) if data["nosql_storage"] is not None: data["nosql_storage"] = ScyllaDBConfig.deserialize(data["nosql_storage"]) return OpenWhiskFunctionConfig(**data) diff --git a/sebs/openwhisk/openwhisk.py b/sebs/openwhisk/openwhisk.py index 12a5abd2..b204baf6 100644 --- a/sebs/openwhisk/openwhisk.py +++ b/sebs/openwhisk/openwhisk.py @@ -22,7 +22,7 @@ from sebs.openwhisk.container import OpenWhiskContainer from sebs.openwhisk.triggers import LibraryTrigger, HTTPTrigger from sebs.storage.resources import SelfHostedSystemResources -from sebs.storage.minio import Minio +from sebs.storage.s3compatible import S3CompatibleStorage from sebs.storage.scylladb import ScyllaDB from sebs.utils import LoggingHandlers from sebs.faas.config import Resources @@ -147,8 +147,11 @@ def shutdown(self) -> None: This method stops storage services if configured and optionally removes the OpenWhisk cluster based on configuration settings. """ - if hasattr(self, "storage") and self.config.shutdownStorage: - self.storage.stop() + if self.config.shutdownStorage: + if self.config.resources.storage_config: + cast(S3CompatibleStorage, self.system_resources.get_storage()).stop() + if self.config.resources.nosql_storage_config: + cast(ScyllaDB, self.system_resources.get_nosql_storage()).stop() if self.config.removeCluster: from tools.openwhisk_preparation import delete_cluster # type: ignore @@ -413,7 +416,7 @@ def create_function( function_cfg = OpenWhiskFunctionConfig.from_benchmark(code_package) if code_package.uses_storage: function_cfg.object_storage = cast( - Minio, self.system_resources.get_storage() + S3CompatibleStorage, self.system_resources.get_storage() ).config if code_package.uses_nosql: function_cfg.nosql_storage = cast( @@ -639,7 +642,7 @@ def is_configuration_changed(self, cached_function: Function, benchmark: Benchma changed = super().is_configuration_changed(cached_function, benchmark) if benchmark.uses_storage: - storage = cast(Minio, self.system_resources.get_storage()) + storage = cast(S3CompatibleStorage, self.system_resources.get_storage()) function = cast(OpenWhiskFunction, cached_function) # check if now we're using a new storage if function.config.object_storage != storage.config: diff --git a/sebs/sebs.py b/sebs/sebs.py index d99aac4c..c31369f3 100644 --- a/sebs/sebs.py +++ b/sebs/sebs.py @@ -26,7 +26,7 @@ from sebs.faas.storage import PersistentStorage from sebs.faas.nosql import NoSQLStorage from sebs.faas.config import Config -from sebs.storage import minio, config, scylladb +from sebs.storage import minio, rustfs, config, scylladb from sebs.utils import has_platform, LoggingHandlers, LoggingBase from sebs.experiments.config import Config as ExperimentConfig @@ -389,7 +389,10 @@ def get_storage_implementation(storage_type: types.Storage) -> Type[PersistentSt Raises: AssertionError: If the requested storage type is not supported """ - _storage_implementations = {types.Storage.MINIO: minio.Minio} + _storage_implementations = { + types.Storage.MINIO: minio.Minio, + types.Storage.RUSTFS: rustfs.RustFS, + } impl = _storage_implementations.get(storage_type) assert impl, f"Storage type {storage_type} not supported" return impl @@ -431,7 +434,10 @@ def get_storage_config_implementation(storage_type: types.Storage): Raises: AssertionError: If the requested storage type is not supported """ - _storage_implementations = {types.Storage.MINIO: config.MinioConfig} + _storage_implementations = { + types.Storage.MINIO: config.MinioConfig, + types.Storage.RUSTFS: config.RustFSConfig, + } impl = _storage_implementations.get(storage_type) assert impl, f"Storage configuration for type {storage_type} not supported" return impl diff --git a/sebs/sebs_types.py b/sebs/sebs_types.py index ba6d129c..b5ea4894 100644 --- a/sebs/sebs_types.py +++ b/sebs/sebs_types.py @@ -53,12 +53,14 @@ class Storage(str, Enum): - AZURE_BLOB_STORAGE: Microsoft Azure Blob Storage - GCP_STORAGE: Google Cloud Storage - MINIO: MinIO object storage (local or self-hosted) + - RUSTFS: RustFS object storage (local or self-hosted) """ AWS_S3 = "aws-s3" AZURE_BLOB_STORAGE = "azure-blob-storage" GCP_STORAGE = "google-cloud-storage" MINIO = "minio" + RUSTFS = "rustfs" class Language(str, Enum): diff --git a/sebs/storage/__init__.py b/sebs/storage/__init__.py index 21b4a2f3..94f46eff 100644 --- a/sebs/storage/__init__.py +++ b/sebs/storage/__init__.py @@ -4,7 +4,7 @@ It includes: - Configuration classes for different storage backends -- MinIO implementation for local S3-compatible object storage +- MinIO and RustFS implementations for local S3-compatible object storage - ScyllaDB implementation for local DynamoDB-compatible NoSQL storage - Resource management classes for self-hosted storage deployments @@ -15,7 +15,9 @@ Key Components: - config: Configuration dataclasses for storage backends + - s3compatible: shared implementation of self-hosted S3-compatible object storage - minio: MinIO-based object storage implementation + - rustfs: RustFS-based object storage implementation - scylladb: ScyllaDB-based NoSQL storage implementation - resources: Resource management for self-hosted storage deployments diff --git a/sebs/storage/config.py b/sebs/storage/config.py index 883cea42..7af50f9b 100644 --- a/sebs/storage/config.py +++ b/sebs/storage/config.py @@ -7,7 +7,7 @@ from abc import ABC, abstractmethod from dataclasses import dataclass, field -from typing import Any, Dict, List +from typing import Any, Dict, List, Type, TypeVar from sebs.cache import Cache @@ -37,9 +37,15 @@ def serialize(self) -> Dict[str, Any]: pass @abstractmethod - def envs(self) -> Dict[str, str]: + def envs(self, external: bool = True) -> Dict[str, str]: """Generate environment variables for the storage configuration. + Args: + external: Advertise the externally reachable address. Functions running + on the same Docker bridge as the storage (local deployment) must use + the internal address instead, as Docker does not route traffic from + the bridge to ports published on the host. + Returns: Dict[str, str]: Environment variables to be set in benchmark runtime """ @@ -47,27 +53,33 @@ def envs(self) -> Dict[str, str]: @dataclass -class MinioConfig(PersistentStorageConfig): - """Configuration for MinIO object storage. +class S3CompatibleConfig(PersistentStorageConfig): + """Configuration for self-hosted, S3-compatible object storage. - MinIO provides a local S3-compatible object storage service that runs in - a Docker container. This configuration class stores all the necessary - parameters for deploying and connecting to a MinIO instance. + The storage runs in a Docker container; see MinioConfig and RustFSConfig + for the supported implementations. This configuration class stores all the + necessary parameters for deploying and connecting to the instance. Attributes: - address: Network address where MinIO is accessible (auto-detected) - mapped_port: Host port mapped to MinIO's internal port 9000 - access_key: Access key for MinIO authentication (auto-generated) - secret_key: Secret key for MinIO authentication (auto-generated) - instance_id: Docker container ID of the running MinIO instance + address: Network address used by SeBS itself to reach the storage (auto-detected). + On Linux this is the container's bridge IP and internal port. + external_address: Network address advertised to benchmark functions, + e.g., the host's IP and the mapped port. Functions running outside + the Docker bridge network (OpenWhisk, Kubernetes) need this address. + Falls back to `address` when empty. + mapped_port: Host port mapped to the container's S3 port 9000 + access_key: Access key for authentication (auto-generated) + secret_key: Secret key for authentication (auto-generated) + instance_id: Docker container ID of the running instance output_buckets: List of bucket names used for benchmark output input_buckets: List of bucket names used for benchmark input - version: MinIO Docker image version to use - data_volume: Host directory path for persistent data storage - type: Storage type identifier (always "minio") + version: Docker image version to use + data_volume: Host directory or named Docker volume for persistent data storage + type: Storage type identifier, set by the subclasses """ address: str = "" + external_address: str = "" mapped_port: int = -1 access_key: str = "" secret_key: str = "" @@ -76,7 +88,7 @@ class MinioConfig(PersistentStorageConfig): input_buckets: List[str] = field(default_factory=lambda: []) version: str = "" data_volume: str = "" - type: str = "minio" + type: str = "" remove_containers: bool = False def update_cache(self, path: List[str], cache: Cache) -> None: @@ -90,16 +102,16 @@ def update_cache(self, path: List[str], cache: Cache) -> None: path: Cache key path prefix for this configuration cache: Cache instance to store configuration in """ - for key in MinioConfig.__dataclass_fields__.keys(): - if key == "resources": - continue + for key in self.__dataclass_fields__.keys(): cache.update_config(val=getattr(self, key), keys=[*path, key]) - @staticmethod - def deserialize(data: Dict[str, Any]) -> "MinioConfig": + T = TypeVar("T", bound="S3CompatibleConfig") + + @classmethod + def deserialize(cls: Type[T], data: Dict[str, Any]) -> T: """Deserialize configuration from a dictionary. - Creates a new MinioConfig instance from dictionary data, typically + Creates a new configuration instance from dictionary data, typically loaded from cache or configuration files. Only known configuration fields are used, unknown fields are ignored. @@ -107,14 +119,12 @@ def deserialize(data: Dict[str, Any]) -> "MinioConfig": data: Dictionary containing configuration data Returns: - MinioConfig: New configuration instance + T: New configuration instance of the calling class """ - keys = list(MinioConfig.__dataclass_fields__.keys()) + keys = list(cls.__dataclass_fields__.keys()) data = {k: v for k, v in data.items() if k in keys} - cfg = MinioConfig(**data) - - return cfg + return cls(**data) def serialize(self) -> Dict[str, Any]: """Serialize the configuration to a dictionary. @@ -124,22 +134,41 @@ def serialize(self) -> Dict[str, Any]: """ return self.__dict__ - def envs(self) -> Dict[str, str]: - """Generate environment variables for MinIO configuration. + def envs(self, external: bool = True) -> Dict[str, str]: + """Generate environment variables for the storage configuration. Creates environment variables that can be used by benchmark functions - to connect to the MinIO storage instance. + to connect to the storage instance. The variable names are shared by + all S3-compatible implementations, so function code stays unchanged. + + Args: + external: Advertise the externally reachable address instead of the + internal one; see PersistentStorageConfig.envs. Returns: - Dict[str, str]: Environment variables for MinIO connection + Dict[str, str]: Environment variables for the storage connection """ return { - "MINIO_ADDRESS": self.address, + "MINIO_ADDRESS": (self.external_address or self.address) if external else self.address, "MINIO_ACCESS_KEY": self.access_key, "MINIO_SECRET_KEY": self.secret_key, } +@dataclass +class MinioConfig(S3CompatibleConfig): + """Configuration for MinIO object storage.""" + + type: str = "minio" + + +@dataclass +class RustFSConfig(S3CompatibleConfig): + """Configuration for RustFS object storage.""" + + type: str = "rustfs" + + @dataclass class NoSQLStorageConfig(ABC): """Abstract base class for NoSQL database storage configuration. @@ -174,7 +203,10 @@ class ScyllaDBConfig(NoSQLStorageConfig): the necessary parameters for deploying and connecting to a ScyllaDB instance. Attributes: - address: Network address where ScyllaDB is accessible (auto-detected) + address: Network address used by SeBS itself to reach ScyllaDB (auto-detected). + On Linux this is the container's bridge IP and the Alternator port. + external_address: Network address advertised to benchmark functions, + e.g., the host's IP and the mapped port. Falls back to `address` when empty. mapped_port: Host port mapped to ScyllaDB's Alternator port alternator_port: Internal port for DynamoDB-compatible API (default: 8000) access_key: Access key for DynamoDB API (placeholder value) @@ -188,6 +220,7 @@ class ScyllaDBConfig(NoSQLStorageConfig): """ address: str = "" + external_address: str = "" mapped_port: int = -1 alternator_port: int = 8000 access_key: str = "None" diff --git a/sebs/storage/minio.py b/sebs/storage/minio.py index 37d4bb14..3915bdb0 100644 --- a/sebs/storage/minio.py +++ b/sebs/storage/minio.py @@ -1,49 +1,25 @@ # Copyright 2020-2025 ETH Zurich and the SeBS authors. All rights reserved. -""" -Module for MinIO S3-compatible storage in the Serverless Benchmarking Suite. +"""MinIO implementation of self-hosted, S3-compatible object storage.""" -MinIO runs in a Docker container and provides persistent -storage for benchmark data and results. It is primarily used for local -testing and on cloud platforms with no object storage, e.g., OpenWhisk. -""" - -import copy -import json -import os -import secrets -import uuid -from typing import Any, Dict, List, Optional, Type, TypeVar - -import docker -import minio - -from sebs.cache import Cache -from sebs.faas.config import Resources -from sebs.faas.storage import PersistentStorage from sebs.storage.config import MinioConfig -from sebs.utils import is_linux +from sebs.storage.s3compatible import S3CompatibleStorage -class Minio(PersistentStorage): - """ - This class manages a self-hosted MinIO storage instance running - in a Docker container. It handles bucket creation, file uploads/downloads, - and container lifecycle management. +class Minio(S3CompatibleStorage): + """Self-hosted MinIO storage instance running in a Docker container. - Attributes: - config: MinIO configuration settings - connection: MinIO client connection + The data volume is a host directory and the container runs as the host user, + so the directory stays writable and removable without elevated privileges. """ - @staticmethod - def typename() -> str: - """ - Get the qualified type name of this class. - - Returns: - str: Full type name including deployment name - """ - return f"{Minio.deployment_name()}.Minio" + # Docker Hub no longer serves minio/minio; the same images are published on quay.io + IMAGE = "quay.io/minio/minio" + COMMAND = "server /data" + ACCESS_KEY_ENV = "MINIO_ACCESS_KEY" + SECRET_KEY_ENV = "MINIO_SECRET_KEY" + HEALTH_PATH = "/minio/health/live" + BIND_MOUNT = True + CONFIG_TYPE = MinioConfig @staticmethod def deployment_name() -> str: @@ -54,505 +30,3 @@ def deployment_name() -> str: str: Deployment name ('minio') """ return "minio" - - # The region setting is required by S3 API but not used for local MinIO - MINIO_REGION = "us-east-1" - - def __init__( - self, - docker_client: docker.DockerClient, - cache_client: Cache, - resources: Resources, - replace_existing: bool, - ): - """ - Initialize a MinIO storage instance. - - Args: - docker_client: Docker client for managing the MinIO container - cache_client: Cache client for storing storage configuration - resources: Resources configuration - replace_existing: Whether to replace existing buckets - """ - super().__init__(self.MINIO_REGION, cache_client, resources, replace_existing) - self._docker_client: docker.DockerClient = docker_client - self._storage_container: Optional[docker.models.containers.Container] = None - self._cfg = MinioConfig() - - @property - def config(self) -> MinioConfig: - """ - Get the MinIO configuration. - - Returns: - MinioConfig: The configuration object - """ - return self._cfg - - @config.setter - def config(self, config: MinioConfig): - """ - Set the MinIO configuration. - - Args: - config: New configuration object - """ - self._cfg = config - - @staticmethod - def _define_http_client() -> Any: - """ - Configure HTTP client for MinIO with appropriate timeouts and retries. - - MinIO does not provide a direct way to configure connection timeouts, so - we need to create a custom HTTP client with proper timeout settings. - The rest of configuration follows MinIO's default client settings. - - Returns: - urllib3.PoolManager: Configured HTTP client for MinIO - """ - import urllib3 - from datetime import timedelta - - timeout = timedelta(seconds=1).seconds - - return urllib3.PoolManager( - timeout=urllib3.util.Timeout(connect=timeout, read=timeout), - maxsize=10, - retries=urllib3.Retry( - total=5, backoff_factor=0.2, status_forcelist=[500, 502, 503, 504] - ), - ) - - def start(self) -> None: - """ - Start a MinIO storage container. - - Creates and runs a Docker container with MinIO, configuring it with - random credentials and mounting a volume for persistent storage. - The container runs in detached mode and is accessible via the - configured port. - - Raises: - RuntimeError: If starting the MinIO container fails - """ - # Set up data volume location - if self._cfg.data_volume == "": - minio_volume = os.path.join(os.getcwd(), "minio-volume") - self._cfg.data_volume = minio_volume - else: - minio_volume = self._cfg.data_volume - minio_volume = os.path.abspath(minio_volume) - - # Create volume directory if it doesn't exist - os.makedirs(minio_volume, exist_ok=True) - volumes = { - minio_volume: { - "bind": "/data", - "mode": "rw", - } - } - - # Generate random credentials for security - self._cfg.access_key = secrets.token_urlsafe(32) - self._cfg.secret_key = secrets.token_hex(32) - self._cfg.address = "" - self.logging.info("Minio storage ACCESS_KEY={}".format(self._cfg.access_key)) - self.logging.info("Minio storage SECRET_KEY={}".format(self._cfg.secret_key)) - - try: - self.logging.info(f"Starting storage Minio on port {self._cfg.mapped_port}") - # Run the MinIO container - self._storage_container = self._docker_client.containers.run( - f"minio/minio:{self._cfg.version}", - command="server /data", - network_mode="bridge", - user=os.getuid(), - ports={"9000": self._cfg.mapped_port}, - environment={ - "MINIO_ACCESS_KEY": self._cfg.access_key, - "MINIO_SECRET_KEY": self._cfg.secret_key, - }, - volumes=volumes, - remove=self.config.remove_containers, - stdout=True, - stderr=True, - detach=True, - ) - assert self._storage_container.id is not None - self._cfg.instance_id = self._storage_container.id - self.configure_connection() - except docker.errors.APIError as e: - self.logging.error("Starting Minio storage failed! Reason: {}".format(e)) - raise RuntimeError("Starting Minio storage unsuccessful") - except Exception as e: - self.logging.error("Starting Minio storage failed! Unknown error: {}".format(e)) - raise RuntimeError("Starting Minio storage unsuccessful") - - def configure_connection(self) -> None: - """ - Configure the connection to the MinIO container. - - Determines the appropriate address to connect to the MinIO container - based on the host platform. For Linux, it uses the container's - bridge IP address, hile for Windows, macOS, or WSL it uses - localhost with the mapped port. - - Raises: - RuntimeError: If the MinIO container is not available or if the IP address - cannot be detected - """ - # Only configure if the address is not already set - if self._cfg.address == "": - # Verify container existence - if self._storage_container is None: - raise RuntimeError( - "Minio container is not available! Make sure that you deployed " - "the Minio storage and provided configuration!" - ) - - # Reload to ensure we have the latest container attributes - self._storage_container.reload() - - # Platform-specific address configuration - if is_linux(): - # On native Linux, use the container's bridge network IP - networks = self._storage_container.attrs["NetworkSettings"]["Networks"] - self._cfg.address = "{IPAddress}:{Port}".format( - IPAddress=networks["bridge"]["IPAddress"], Port=9000 - ) - else: - # On Windows, macOS, or WSL, use localhost with the mapped port - self._cfg.address = f"localhost:{self._cfg.mapped_port}" - - # Verify address was successfully determined - if not self._cfg.address: - self.logging.error( - f"Couldn't read the IP address of container from attributes " - f"{json.dumps(self._storage_container.attrs, indent=2)}" - ) - raise RuntimeError( - f"Incorrect detection of IP address for container with id " - f"{self._cfg.instance_id}" - ) - self.logging.info("Starting minio instance at {}".format(self._cfg.address)) - - # Create the connection using the configured address - self.connection = self.get_connection() - - def stop(self) -> None: - """ - Stop the MinIO container. - - Gracefully stops the running MinIO container if it exists. - Logs an error if the container is not known. - """ - if self._storage_container is not None: - self.logging.info(f"Stopping minio container at {self._cfg.address}.") - self._storage_container.stop() - self.logging.info(f"Stopped minio container at {self._cfg.address}.") - else: - self.logging.error("Stopping minio was not successful, storage container not known!") - - def get_connection(self) -> minio.Minio: - """ - Create a new MinIO client connection. - - Creates a connection to the MinIO server using the configured address, - credentials, and HTTP client settings. - - Returns: - minio.Minio: Configured MinIO client - """ - return minio.Minio( - self._cfg.address, - access_key=self._cfg.access_key, - secret_key=self._cfg.secret_key, - secure=False, # Local MinIO doesn't use HTTPS - http_client=Minio._define_http_client(), - ) - - def _create_bucket( - self, - name: str, - buckets: Optional[List[str]] = None, - randomize_name: bool = False, - ) -> str: - """ - Create a new bucket if it doesn't already exist. - - Checks if a bucket with the given name already exists in the list of buckets. - If not, creates a new bucket with either the exact name or a randomized name. - - Args: - name: Base name for the bucket - buckets: List of existing bucket names to check against - randomize_name: Whether to append a random UUID to the bucket name - - Returns: - str: Name of the existing or newly created bucket - - Raises: - minio.error.ResponseError: If bucket creation fails - """ - - if buckets is None: - buckets = [] - - # Check if bucket already exists - for bucket_name in buckets: - if name in bucket_name: - self.logging.info( - "Bucket {} for {} already exists, skipping.".format(bucket_name, name) - ) - return bucket_name - - # MinIO has limit of bucket name to 16 characters - if randomize_name: - bucket_name = "{}-{}".format(name, str(uuid.uuid4())[0:16]) - else: - bucket_name = name - - try: - self.connection.make_bucket(bucket_name, location=self.MINIO_REGION) - self.logging.info("Created bucket {}".format(bucket_name)) - return bucket_name - except ( - minio.error.BucketAlreadyOwnedByYou, - minio.error.BucketAlreadyExists, - minio.error.ResponseError, - ) as err: - self.logging.error("Bucket creation failed!") - # Rethrow the error for handling by the caller - raise err - - def uploader_func(self, path_idx: int, file: str, filepath: str) -> None: - """ - Upload a file to the MinIO storage. - - Uploads a file to the specified input prefix in the benchmarks bucket. - This function is passed to benchmarks for uploading their input data. - - Args: - path_idx: Index of the input prefix to use - file: Name of the file within the bucket - filepath: Local path to the file to upload - - Raises: - minio.error.ResponseError: If the upload fails - """ - try: - key = os.path.join(self.input_prefixes[path_idx], file) - bucket_name = self.get_bucket(Resources.StorageBucketType.BENCHMARKS) - self.logging.info("Upload {} to {}".format(filepath, bucket_name)) - self.connection.fput_object(bucket_name, key, filepath) - except minio.error.ResponseError as err: - self.logging.error("Upload failed!") - raise err - - def clean_bucket(self, bucket_name: str) -> None: - """ - Remove all objects from a bucket. - - Deletes all objects within the specified bucket but keeps the bucket itself. - Logs any errors that occur during object deletion. - - Args: - bucket: Name of the bucket to clean - """ - delete_object_list = map( - lambda x: minio.DeleteObject(x.object_name), - self.connection.list_objects(bucket_name=bucket_name), - ) - errors = self.connection.remove_objects(bucket_name, delete_object_list) - for error in errors: - self.logging.error(f"Error when deleting object from bucket {bucket_name}: {error}!") - - def remove_bucket(self, bucket: str) -> None: - """ - Delete a bucket completely. - - Removes the specified bucket from the MinIO storage. - The bucket must be empty before it can be deleted. - - Args: - bucket: Name of the bucket to remove - """ - self.connection.remove_bucket(Bucket=bucket) - - def correct_name(self, name: str) -> str: - """ - Format a bucket name to comply with MinIO naming requirements. - - For MinIO, no name correction is needed (unlike some cloud providers - that enforce additional restrictions). - - Args: - name: Original bucket name - - Returns: - str: Bucket name (unchanged for MinIO) - """ - return name - - def download(self, bucket_name: str, key: str, filepath: str) -> None: - """ - Download an object from a bucket to a local file. - - Args: - bucket_name: Name of the source bucket - key: Object key/path in the bucket - filepath: Local destination path - - Raises: - RuntimeError: If the bucket does not exist - minio.error.ResponseError: If the download fails - """ - if not self.exists_bucket(bucket_name): - raise RuntimeError(f"Attempting to download from a non-existing bucket {bucket_name}!") - try: - self.connection.fget_object(bucket_name, key, filepath) - except minio.error.ResponseError as err: - raise err - - def exists_bucket(self, bucket_name: str) -> bool: - """ - Check if a bucket exists. - - Args: - bucket_name: Name of the bucket to check - - Returns: - bool: True if the bucket exists, False otherwise - """ - return self.connection.bucket_exists(bucket_name) - - def list_bucket(self, bucket_name: str, prefix: str = "") -> List[str]: - """ - List all objects in a bucket with an optional prefix filter. - - Args: - bucket_name: Name of the bucket to list - prefix: Optional prefix to filter objects - - Returns: - List[str]: List of object names in the bucket - - Raises: - RuntimeError: If the bucket does not exist - """ - try: - objects_list = self.connection.list_objects(bucket_name) - return [obj.object_name for obj in objects_list if prefix in obj.object_name] - except minio.error.NoSuchBucket: - raise RuntimeError( - f"Attempting to access a non-existing bucket {bucket_name}!" - ) from None - - def list_buckets(self, bucket_name: Optional[str] = None) -> List[str]: - """ - List all buckets, optionally filtered by name. - - Args: - bucket_name: Optional filter for bucket names - - Returns: - List[str]: List of bucket names - """ - buckets = self.connection.list_buckets() - if bucket_name is not None: - return [bucket.name for bucket in buckets if bucket_name in bucket.name] - else: - return [bucket.name for bucket in buckets] - - def upload(self, bucket_name: str, filepath: str, key: str) -> None: - """ - Upload a file to a bucket. - - Not implemented for this class. Use fput_object directly or uploader_func. - - Raises: - NotImplementedError: This method is not implemented - """ - raise NotImplementedError() - - def serialize(self) -> Dict[str, Any]: - """ - Serialize MinIO configuration to a dictionary. - - Returns: - dict: Serialized configuration data - """ - return self._cfg.serialize() - - T = TypeVar("T", bound="Minio") - - @staticmethod - def _deserialize( - cached_config: MinioConfig, - cache_client: Cache, - resources: Resources, - obj_type: Type[T], - ) -> T: - """ - Deserialize a MinIO instance from cached configuration with custom type. - - Creates a new instance of the specified class type from cached configuration - data. This allows platform-specific versions to be deserialized correctly - while sharing the core implementation. When overriding the implementation in - Local/OpenWhisk/..., we call the _deserialize method and provide an - alternative implementation type. - - FIXME: is this still needed? It looks like we stopped using - platform-specific implementations. - - Args: - cached_config: Cached MinIO configuration - cache_client: Cache client - resources: Resources configuration - obj_type: Type of object to create (a Minio subclass) - - Returns: - T: Deserialized instance of the specified type - - Raises: - RuntimeError: If the storage container does not exist - """ - docker_client = docker.from_env() - obj = obj_type(docker_client, cache_client, resources, False) - obj._cfg = cached_config - - # Try to reconnect to existing container if ID is available - if cached_config.instance_id: - instance_id = cached_config.instance_id - try: - obj._storage_container = docker_client.containers.get(instance_id) - except docker.errors.NotFound: - raise RuntimeError(f"Storage container {instance_id} does not exist!") - else: - obj._storage_container = None - - # Copy bucket information - obj._input_prefixes = copy.copy(cached_config.input_buckets) - obj._output_prefixes = copy.copy(cached_config.output_buckets) - - # Set up connection - obj.configure_connection() - return obj - - @staticmethod - def deserialize(cached_config: MinioConfig, cache_client: Cache, res: Resources) -> "Minio": - """ - Deserialize a MinIO instance from cached configuration. - - Creates a new Minio instance from cached configuration data. - - Args: - cached_config: Cached MinIO configuration - cache_client: Cache client - res: Resources configuration - - Returns: - Minio: Deserialized Minio instance - """ - return Minio._deserialize(cached_config, cache_client, res, Minio) diff --git a/sebs/storage/resources.py b/sebs/storage/resources.py index 0f190ac6..f59b10d9 100644 --- a/sebs/storage/resources.py +++ b/sebs/storage/resources.py @@ -10,7 +10,7 @@ """ import docker -from typing import cast, Dict, Optional, Tuple, Any +from typing import cast, Dict, Optional, Tuple, Type, Any from sebs.cache import Cache from sebs.faas.config import Config, Resources @@ -18,15 +18,23 @@ from sebs.faas.storage import PersistentStorage from sebs.faas.nosql import NoSQLStorage from sebs.storage.minio import Minio +from sebs.storage.rustfs import RustFS +from sebs.storage.s3compatible import S3CompatibleStorage from sebs.storage.scylladb import ScyllaDB from sebs.storage.config import ( NoSQLStorageConfig, PersistentStorageConfig, + S3CompatibleConfig, ScyllaDBConfig, - MinioConfig, ) from sebs.utils import LoggingHandlers +# Self-hosted object storage implementations, keyed by storage type +OBJECT_STORAGE_IMPLEMENTATIONS: Dict[str, Type[S3CompatibleStorage]] = { + Minio.deployment_name(): Minio, + RustFS.deployment_name(): RustFS, +} + class SelfHostedResources(Resources): """Resource configuration for self-hosted storage deployments. @@ -98,7 +106,7 @@ def update_cache(self, cache: Cache) -> None: """ super().update_cache(cache) if self._object_storage is not None: - cast(MinioConfig, self._object_storage).update_cache( + cast(S3CompatibleConfig, self._object_storage).update_cache( [self._name, "resources", "storage"], cache ) if self._nosql_storage is not None: @@ -139,7 +147,7 @@ def _deserialize_storage( cached_config is not None and "resources" in cached_config and "storage" in cached_config["resources"] - and "object" in cached_config["resources"]["storage"] + and storage_type in cached_config["resources"]["storage"] ): storage_impl = cached_config["resources"]["storage"][storage_type]["type"] storage_config = cached_config["resources"]["storage"][storage_type][storage_impl] @@ -168,9 +176,10 @@ def _deserialize( config, cached_config, "object" ) - if obj_storage_impl == "minio": - ret._object_storage = MinioConfig.deserialize(obj_storage_cfg) - ret.logging.info("Deserializing access data to Minio storage") + if obj_storage_impl in OBJECT_STORAGE_IMPLEMENTATIONS: + config_type = OBJECT_STORAGE_IMPLEMENTATIONS[obj_storage_impl].CONFIG_TYPE + ret._object_storage = config_type.deserialize(obj_storage_cfg) + ret.logging.info(f"Deserializing access data to {obj_storage_impl} storage") elif obj_storage_impl != "": ret.logging.warning(f"Unknown object storage type: {obj_storage_impl}") else: @@ -226,7 +235,7 @@ def __init__( def get_storage(self, replace_existing: Optional[bool] = None) -> PersistentStorage: """Get or create a persistent storage instance. - Creates a MinIO storage instance if one doesn't exist, or returns the + Creates a storage instance if one doesn't exist, or returns the existing instance. The storage is deserialized from a serialized config of an existing storage deployment. @@ -234,7 +243,7 @@ def get_storage(self, replace_existing: Optional[bool] = None) -> PersistentStor replace_existing: Whether to replace existing buckets (optional) Returns: - PersistentStorage: MinIO storage instance + PersistentStorage: S3-compatible storage instance (Minio, RustFS) Raises: RuntimeError: If storage configuration is missing or unsupported @@ -244,12 +253,15 @@ def get_storage(self, replace_existing: Optional[bool] = None) -> PersistentStor if storage_config is None: self.logging.error( f"The {self._name} deployment is missing the " - "configuration of pre-allocated storage!" + "configuration of pre-allocated storage! Start the storage with " + "'sebs storage start object configs/storage.json --output-json storage.json' " + "and pass the result with '--storage-configuration storage.json'." ) raise RuntimeError(f"Cannot run {self._name} deployment without any object storage") - if isinstance(storage_config, MinioConfig): - self._storage = Minio.deserialize( + if isinstance(storage_config, S3CompatibleConfig): + impl = OBJECT_STORAGE_IMPLEMENTATIONS[storage_config.type] + self._storage = impl.deserialize( storage_config, self._cache_client, self._config.resources, @@ -285,7 +297,9 @@ def get_nosql_storage(self) -> NoSQLStorage: if storage_config is None: self.logging.error( f"The {self._name} deployment is missing the configuration " - "of pre-allocated NoSQL storage!" + "of pre-allocated NoSQL storage! Start the storage with " + "'sebs storage start nosql configs/storage.json --output-json storage.json' " + "and pass the result with '--storage-configuration storage.json'." ) raise RuntimeError("Cannot allocate NoSQL storage!") diff --git a/sebs/storage/rustfs.py b/sebs/storage/rustfs.py new file mode 100644 index 00000000..e0ccac6f --- /dev/null +++ b/sebs/storage/rustfs.py @@ -0,0 +1,31 @@ +# Copyright 2020-2025 ETH Zurich and the SeBS authors. All rights reserved. +"""RustFS implementation of self-hosted, S3-compatible object storage.""" + +from sebs.storage.config import RustFSConfig +from sebs.storage.s3compatible import S3CompatibleStorage + + +class RustFS(S3CompatibleStorage): + """Self-hosted RustFS storage instance running in a Docker container. + + RustFS runs as a fixed, unprivileged user inside the container, so the data + is kept in a named Docker volume instead of a host directory. + """ + + IMAGE = "rustfs/rustfs" + COMMAND = None + ACCESS_KEY_ENV = "RUSTFS_ACCESS_KEY" + SECRET_KEY_ENV = "RUSTFS_SECRET_KEY" + HEALTH_PATH = "/health" + BIND_MOUNT = False + CONFIG_TYPE = RustFSConfig + + @staticmethod + def deployment_name() -> str: + """ + Get the deployment platform name. + + Returns: + str: Deployment name ('rustfs') + """ + return "rustfs" diff --git a/sebs/storage/s3compatible.py b/sebs/storage/s3compatible.py new file mode 100644 index 00000000..22c865fc --- /dev/null +++ b/sebs/storage/s3compatible.py @@ -0,0 +1,609 @@ +# Copyright 2020-2025 ETH Zurich and the SeBS authors. All rights reserved. +""" +Module for self-hosted, S3-compatible object storage in the Serverless Benchmarking Suite. + +The storage runs in a Docker container and provides persistent storage for +benchmark data and results. It is primarily used for local testing and on +platforms with no object storage of their own, e.g., OpenWhisk. Concrete +implementations (MinIO, RustFS) only define the container image and its +configuration; the S3 client side is shared. +""" + +import copy +import json +import os +import secrets +import uuid +from typing import Any, Dict, List, Optional, Type, TypeVar + +import docker +import minio + +from sebs.cache import Cache +from sebs.faas.config import Resources +from sebs.faas.storage import PersistentStorage +from sebs.storage.config import S3CompatibleConfig +from sebs.utils import is_linux, probe_http, resolve_external_address + + +class S3CompatibleStorage(PersistentStorage): + """ + This class manages a self-hosted, S3-compatible storage instance running + in a Docker container. It handles bucket creation, file uploads/downloads, + and container lifecycle management. + + Subclasses define the container image and how it is configured: + + Attributes: + IMAGE: Docker image without a tag; the tag comes from the configuration's version + COMMAND: Command passed to the container, or None for the image default + ACCESS_KEY_ENV: Environment variable receiving the generated access key + SECRET_KEY_ENV: Environment variable receiving the generated secret key + HEALTH_PATH: HTTP path of the liveness endpoint + BIND_MOUNT: Mount a host directory as the data volume and run the container + as the host user. Otherwise, use a named Docker volume and the image's user. + CONFIG_TYPE: Configuration class of the implementation + config: Storage configuration settings + connection: S3 client connection + """ + + IMAGE: str + COMMAND: Optional[str] + ACCESS_KEY_ENV: str + SECRET_KEY_ENV: str + HEALTH_PATH: str + BIND_MOUNT: bool + CONFIG_TYPE: Type[S3CompatibleConfig] + + @classmethod + def typename(cls) -> str: + """ + Get the qualified type name of this class. + + Returns: + str: Full type name including deployment name + """ + return f"{cls.deployment_name()}.{cls.__name__}" + + # The region setting is required by S3 API but not used for self-hosted storage + S3_REGION = "us-east-1" + + def __init__( + self, + docker_client: docker.DockerClient, + cache_client: Cache, + resources: Resources, + replace_existing: bool, + ): + """ + Initialize a storage instance. + + Args: + docker_client: Docker client for managing the storage container + cache_client: Cache client for storing storage configuration + resources: Resources configuration + replace_existing: Whether to replace existing buckets + """ + super().__init__(self.S3_REGION, cache_client, resources, replace_existing) + self._docker_client: docker.DockerClient = docker_client + self._storage_container: Optional[docker.models.containers.Container] = None + self._cfg: S3CompatibleConfig = self.CONFIG_TYPE() + + @property + def config(self) -> S3CompatibleConfig: + """ + Get the storage configuration. + + Returns: + S3CompatibleConfig: The configuration object + """ + return self._cfg + + @config.setter + def config(self, config: S3CompatibleConfig): + """ + Set the storage configuration. + + Args: + config: New configuration object + """ + self._cfg = config + + @staticmethod + def _define_http_client() -> Any: + """ + Configure HTTP client for the S3 client with appropriate timeouts and retries. + + The MinIO SDK does not provide a direct way to configure connection timeouts, so + we need to create a custom HTTP client with proper timeout settings. + The rest of configuration follows the SDK's default client settings. + + Returns: + urllib3.PoolManager: Configured HTTP client + """ + import urllib3 + from datetime import timedelta + + timeout = timedelta(seconds=1).seconds + + return urllib3.PoolManager( + timeout=urllib3.util.Timeout(connect=timeout, read=timeout), + maxsize=10, + retries=urllib3.Retry( + total=5, backoff_factor=0.2, status_forcelist=[500, 502, 503, 504] + ), + ) + + def start(self) -> None: + """ + Start the storage container. + + Creates and runs the Docker container, configuring it with + random credentials and mounting a volume for persistent storage. + The container runs in detached mode and is accessible via the + configured port. + + Raises: + RuntimeError: If starting the container fails + """ + name = self.deployment_name() + # Set up data volume: a host directory, or a named Docker volume + if self._cfg.data_volume == "": + self._cfg.data_volume = f"{name}-volume" + if self.BIND_MOUNT: + self._cfg.data_volume = os.path.abspath(self._cfg.data_volume) + os.makedirs(self._cfg.data_volume, exist_ok=True) + volumes = {self._cfg.data_volume: {"bind": "/data", "mode": "rw"}} + + # Generate random credentials for security + self._cfg.access_key = secrets.token_urlsafe(32) + self._cfg.secret_key = secrets.token_hex(32) + self._cfg.address = "" + self.logging.info(f"{name} storage ACCESS_KEY={self._cfg.access_key}") + self.logging.info(f"{name} storage SECRET_KEY={self._cfg.secret_key}") + + try: + self.logging.info(f"Starting storage {name} on port {self._cfg.mapped_port}") + self._storage_container = self._docker_client.containers.run( + f"{self.IMAGE}:{self._cfg.version}", + command=self.COMMAND, + network_mode="bridge", + user=os.getuid() if self.BIND_MOUNT else None, + ports={"9000": self._cfg.mapped_port}, + environment={ + self.ACCESS_KEY_ENV: self._cfg.access_key, + self.SECRET_KEY_ENV: self._cfg.secret_key, + }, + volumes=volumes, + remove=self.config.remove_containers, + stdout=True, + stderr=True, + detach=True, + ) + assert self._storage_container.id is not None + self._cfg.instance_id = self._storage_container.id + self.configure_connection() + except docker.errors.APIError as e: + self.logging.error(f"Starting {name} storage failed! Reason: {e}") + raise RuntimeError(f"Starting {name} storage unsuccessful") + except Exception as e: + self.logging.error(f"Starting {name} storage failed! Unknown error: {e}") + raise RuntimeError(f"Starting {name} storage unsuccessful") + + def configure_connection(self) -> None: + """ + Configure the connection to the storage container. + + Determines the appropriate address to connect to the container + based on the host platform. For Linux, it uses the container's + bridge IP address, while for Windows, macOS, or WSL it uses + localhost with the mapped port. + + Additionally, it determines the address advertised to benchmark + functions: the user-provided external address, or the host's + default-route IP combined with the mapped port. This address is + reachable from outside the Docker bridge network, e.g., from + Kubernetes pods. + + Raises: + RuntimeError: If the container is not available or if the IP address + cannot be detected + """ + # Only configure if the address is not already set + if self._cfg.address == "": + # Verify container existence + if self._storage_container is None: + raise RuntimeError( + "Storage container is not available! Make sure that you deployed " + "the Minio storage and provided configuration!" + ) + + # Reload to ensure we have the latest container attributes + self._storage_container.reload() + + # Platform-specific address configuration + if is_linux(): + # On native Linux, use the container's bridge network IP + networks = self._storage_container.attrs["NetworkSettings"]["Networks"] + self._cfg.address = "{IPAddress}:{Port}".format( + IPAddress=networks["bridge"]["IPAddress"], Port=9000 + ) + else: + # On Windows, macOS, or WSL, use localhost with the mapped port + self._cfg.address = f"localhost:{self._cfg.mapped_port}" + + # Verify address was successfully determined + if not self._cfg.address: + self.logging.error( + f"Couldn't read the IP address of container from attributes " + f"{json.dumps(self._storage_container.attrs, indent=2)}" + ) + raise RuntimeError( + f"Incorrect detection of IP address for container with id " + f"{self._cfg.instance_id}" + ) + self.logging.info(f"Starting {self.deployment_name()} instance at {self._cfg.address}") + self.configure_external_address() + + # Create the connection using the configured address + self.connection = self.get_connection() + + def configure_external_address(self) -> None: + """ + Determine the address advertised to benchmark functions. + + Uses the user-provided external address, or the host's default-route IP, + combined with the mapped port. + """ + name = self.deployment_name() + self._cfg.external_address = resolve_external_address( + self._cfg.external_address, self._cfg.mapped_port + ) + if self._cfg.external_address: + self.logging.info(f"{name} advertised to functions at {self._cfg.external_address}") + else: + self.logging.warning( + "Could not detect the host's IP address. Functions running outside of the " + f"Docker bridge network will not reach {name}; provide --external-address." + ) + + def check_external_address(self) -> bool: + """ + Verify that the storage is reachable through the address advertised to functions. + + Failures are reported as warnings, since the host running SeBS is not + always able to reach the same network as benchmark functions. + + Returns: + bool: True if the probe succeeded + """ + if not self._cfg.external_address: + return False + name = self.deployment_name() + url = f"http://{self._cfg.external_address}{self.HEALTH_PATH}" + error = probe_http(url, timeout_seconds=15) + if error is None: + self.logging.info(f"{name} is reachable at {url}") + return True + self.logging.warning( + f"{name} is not reachable at {url}: {error}. Benchmark functions might not be " + f"able to reach the storage. Verify the address with: curl -i {url}" + ) + return False + + def stop(self) -> None: + """ + Stop the storage container. + + Gracefully stops the running MinIO container if it exists. + Logs an error if the container is not known. + """ + if self._storage_container is not None: + self.logging.info(f"Stopping storage container at {self._cfg.address}.") + self._storage_container.stop() + self.logging.info(f"Stopped storage container at {self._cfg.address}.") + else: + self.logging.error("Stopping storage was not successful, container not known!") + + def get_connection(self) -> minio.Minio: + """ + Create a new S3 client connection. + + Creates a connection to the storage server using the configured address, + credentials, and HTTP client settings. + + Returns: + minio.Minio: Configured S3 client from the MinIO SDK + """ + return minio.Minio( + self._cfg.address, + access_key=self._cfg.access_key, + secret_key=self._cfg.secret_key, + secure=False, # Self-hosted storage doesn't use HTTPS + http_client=self._define_http_client(), + ) + + def _create_bucket( + self, + name: str, + buckets: Optional[List[str]] = None, + randomize_name: bool = False, + ) -> str: + """ + Create a new bucket if it doesn't already exist. + + Checks if a bucket with the given name already exists in the list of buckets. + If not, creates a new bucket with either the exact name or a randomized name. + + Args: + name: Base name for the bucket + buckets: List of existing bucket names to check against + randomize_name: Whether to append a random UUID to the bucket name + + Returns: + str: Name of the existing or newly created bucket + + Raises: + minio.error.ResponseError: If bucket creation fails + """ + + if buckets is None: + buckets = [] + + # Check if bucket already exists + for bucket_name in buckets: + if name in bucket_name: + self.logging.info( + "Bucket {} for {} already exists, skipping.".format(bucket_name, name) + ) + return bucket_name + + # Bucket names are limited to 16 characters + if randomize_name: + bucket_name = "{}-{}".format(name, str(uuid.uuid4())[0:16]) + else: + bucket_name = name + + try: + self.connection.make_bucket(bucket_name, location=self.S3_REGION) + self.logging.info("Created bucket {}".format(bucket_name)) + return bucket_name + except ( + minio.error.BucketAlreadyOwnedByYou, + minio.error.BucketAlreadyExists, + minio.error.ResponseError, + ) as err: + self.logging.error("Bucket creation failed!") + # Rethrow the error for handling by the caller + raise err + + def uploader_func(self, path_idx: int, file: str, filepath: str) -> None: + """ + Upload a file to the storage. + + Uploads a file to the specified input prefix in the benchmarks bucket. + This function is passed to benchmarks for uploading their input data. + + Args: + path_idx: Index of the input prefix to use + file: Name of the file within the bucket + filepath: Local path to the file to upload + + Raises: + minio.error.ResponseError: If the upload fails + """ + try: + key = os.path.join(self.input_prefixes[path_idx], file) + bucket_name = self.get_bucket(Resources.StorageBucketType.BENCHMARKS) + self.logging.info("Upload {} to {}".format(filepath, bucket_name)) + self.connection.fput_object(bucket_name, key, filepath) + except minio.error.ResponseError as err: + self.logging.error("Upload failed!") + raise err + + def clean_bucket(self, bucket_name: str) -> None: + """ + Remove all objects from a bucket. + + Deletes all objects within the specified bucket but keeps the bucket itself. + Logs any errors that occur during object deletion. + + Args: + bucket: Name of the bucket to clean + """ + delete_object_list = map( + lambda x: minio.DeleteObject(x.object_name), + self.connection.list_objects(bucket_name=bucket_name), + ) + errors = self.connection.remove_objects(bucket_name, delete_object_list) + for error in errors: + self.logging.error(f"Error when deleting object from bucket {bucket_name}: {error}!") + + def remove_bucket(self, bucket: str) -> None: + """ + Delete a bucket completely. + + Removes the specified bucket from the storage. + The bucket must be empty before it can be deleted. + + Args: + bucket: Name of the bucket to remove + """ + self.connection.remove_bucket(Bucket=bucket) + + def correct_name(self, name: str) -> str: + """ + Format a bucket name to comply with naming requirements. + + No name correction is needed (unlike some cloud providers + that enforce additional restrictions). + + Args: + name: Original bucket name + + Returns: + str: Bucket name (unchanged) + """ + return name + + def download(self, bucket_name: str, key: str, filepath: str) -> None: + """ + Download an object from a bucket to a local file. + + Args: + bucket_name: Name of the source bucket + key: Object key/path in the bucket + filepath: Local destination path + + Raises: + RuntimeError: If the bucket does not exist + minio.error.ResponseError: If the download fails + """ + if not self.exists_bucket(bucket_name): + raise RuntimeError(f"Attempting to download from a non-existing bucket {bucket_name}!") + try: + self.connection.fget_object(bucket_name, key, filepath) + except minio.error.ResponseError as err: + raise err + + def exists_bucket(self, bucket_name: str) -> bool: + """ + Check if a bucket exists. + + Args: + bucket_name: Name of the bucket to check + + Returns: + bool: True if the bucket exists, False otherwise + """ + return self.connection.bucket_exists(bucket_name) + + def list_bucket(self, bucket_name: str, prefix: str = "") -> List[str]: + """ + List all objects in a bucket with an optional prefix filter. + + Args: + bucket_name: Name of the bucket to list + prefix: Optional prefix to filter objects + + Returns: + List[str]: List of object names in the bucket + + Raises: + RuntimeError: If the bucket does not exist + """ + try: + objects_list = self.connection.list_objects(bucket_name) + return [obj.object_name for obj in objects_list if prefix in obj.object_name] + except minio.error.NoSuchBucket: + raise RuntimeError( + f"Attempting to access a non-existing bucket {bucket_name}!" + ) from None + + def list_buckets(self, bucket_name: Optional[str] = None) -> List[str]: + """ + List all buckets, optionally filtered by name. + + Args: + bucket_name: Optional filter for bucket names + + Returns: + List[str]: List of bucket names + """ + buckets = self.connection.list_buckets() + if bucket_name is not None: + return [bucket.name for bucket in buckets if bucket_name in bucket.name] + else: + return [bucket.name for bucket in buckets] + + def upload(self, bucket_name: str, filepath: str, key: str) -> None: + """ + Upload a file to a bucket. + + Not implemented for this class. Use fput_object directly or uploader_func. + + Raises: + NotImplementedError: This method is not implemented + """ + raise NotImplementedError() + + def serialize(self) -> Dict[str, Any]: + """ + Serialize the storage configuration to a dictionary. + + Returns: + dict: Serialized configuration data + """ + return self._cfg.serialize() + + T = TypeVar("T", bound="S3CompatibleStorage") + + @staticmethod + def _deserialize( + cached_config: S3CompatibleConfig, + cache_client: Cache, + resources: Resources, + obj_type: Type[T], + ) -> T: + """ + Deserialize a storage instance from cached configuration with custom type. + + Creates a new instance of the specified class type from cached configuration + data. This allows platform-specific versions to be deserialized correctly + while sharing the core implementation. When overriding the implementation in + Local/OpenWhisk/..., we call the _deserialize method and provide an + alternative implementation type. + + FIXME: is this still needed? It looks like we stopped using + platform-specific implementations. + + Args: + cached_config: Cached storage configuration + cache_client: Cache client + resources: Resources configuration + obj_type: Type of object to create (a S3CompatibleStorage subclass) + + Returns: + T: Deserialized instance of the specified type + + Raises: + RuntimeError: If the storage container does not exist + """ + docker_client = docker.from_env() + obj = obj_type(docker_client, cache_client, resources, False) + obj._cfg = cached_config + + # Try to reconnect to existing container if ID is available + if cached_config.instance_id: + instance_id = cached_config.instance_id + try: + obj._storage_container = docker_client.containers.get(instance_id) + except docker.errors.NotFound: + raise RuntimeError(f"Storage container {instance_id} does not exist!") + else: + obj._storage_container = None + + # Copy bucket information + obj._input_prefixes = copy.copy(cached_config.input_buckets) + obj._output_prefixes = copy.copy(cached_config.output_buckets) + + # Set up connection + obj.configure_connection() + return obj + + @classmethod + def deserialize( + cls: Type[T], cached_config: S3CompatibleConfig, cache_client: Cache, res: Resources + ) -> T: + """ + Deserialize a storage instance from cached configuration. + + Args: + cached_config: Cached storage configuration + cache_client: Cache client + res: Resources configuration + + Returns: + T: Deserialized instance of the calling class + """ + return cls._deserialize(cached_config, cache_client, res, cls) diff --git a/sebs/storage/scylladb.py b/sebs/storage/scylladb.py index 6c532686..f2a0cbe2 100644 --- a/sebs/storage/scylladb.py +++ b/sebs/storage/scylladb.py @@ -23,6 +23,7 @@ from sebs.faas.nosql import NoSQLStorage from sebs.sebs_types import NoSQLStorage as StorageType from sebs.storage.config import ScyllaDBConfig +from sebs.utils import probe_http, resolve_external_address class ScyllaDB(NoSQLStorage): @@ -208,6 +209,10 @@ def configure_connection(self) -> None: based on the host platform. For Linux, it uses the container's IP address, while for Windows, macOS, or WSL it uses localhost with the mapped port. + Additionally, it determines the address advertised to benchmark + functions: the user-provided external address, or the host's + default-route IP combined with the mapped port. + Creates a boto3 DynamoDB client configured to connect to ScyllaDB's Alternator interface. @@ -244,6 +249,7 @@ def configure_connection(self) -> None: f"{self._cfg.instance_id}" ) self.logging.info("Starting ScyllaDB instance at {}".format(self._cfg.address)) + self.configure_external_address() # Create the DynamoDB client for ScyllaDB's Alternator interface self.client = boto3.client( @@ -254,6 +260,45 @@ def configure_connection(self) -> None: endpoint_url=f"http://{self._cfg.address}", ) + def configure_external_address(self) -> None: + """Determine the address advertised to benchmark functions. + + Uses the user-provided external address, or the host's default-route IP, + combined with the mapped port. + """ + self._cfg.external_address = resolve_external_address( + self._cfg.external_address, self._cfg.mapped_port + ) + if self._cfg.external_address: + self.logging.info(f"ScyllaDB advertised to functions at {self._cfg.external_address}") + else: + self.logging.warning( + "Could not detect the host's IP address. Functions running outside of the " + "Docker bridge network will not reach ScyllaDB; provide --external-address." + ) + + def check_external_address(self) -> bool: + """Verify that ScyllaDB is reachable through the address advertised to functions. + + Failures are reported as warnings, since the host running SeBS is not + always able to reach the same network as benchmark functions. + + Returns: + bool: True if the probe succeeded + """ + if not self._cfg.external_address: + return False + url = f"http://{self._cfg.external_address}/" + error = probe_http(url, timeout_seconds=15) + if error is None: + self.logging.info(f"ScyllaDB is reachable at {url}") + return True + self.logging.warning( + f"ScyllaDB is not reachable at {url}: {error}. Benchmark functions might not be " + f"able to reach the storage. Verify the address with: curl -i {url}" + ) + return False + def stop(self) -> None: """Stop the ScyllaDB container. @@ -266,16 +311,25 @@ def stop(self) -> None: else: self.logging.error("Stopping ScyllaDB was not successful, storage container not known!") - def envs(self) -> Dict[str, str]: + def envs(self, external: bool = True) -> Dict[str, str]: """Generate environment variables for ScyllaDB configuration. Creates environment variables that can be used by benchmark functions to connect to the ScyllaDB storage instance. + Args: + external: Advertise the externally reachable address instead of the + internal one; see NoSQLStorage.envs. + Returns: Dict[str, str]: Environment variables for ScyllaDB connection """ - return {"NOSQL_STORAGE_TYPE": "scylladb", "NOSQL_STORAGE_ENDPOINT": self._cfg.address} + return { + "NOSQL_STORAGE_TYPE": "scylladb", + "NOSQL_STORAGE_ENDPOINT": ( + (self._cfg.external_address or self._cfg.address) if external else self._cfg.address + ), + } def serialize(self) -> Tuple[StorageType, Dict[str, Any]]: """Serialize ScyllaDB configuration to a tuple. diff --git a/sebs/utils.py b/sebs/utils.py index 4ebdcb82..1567b033 100644 --- a/sebs/utils.py +++ b/sebs/utils.py @@ -19,7 +19,9 @@ import click import datetime import platform +import socket import threading +import time import re from pathlib import Path @@ -191,7 +193,7 @@ def append_nested_dict(cfg: dict, keys: List[str], value: Optional[dict]) -> Non # make sure parent keys exist for key in keys[:-1]: cfg = cfg.setdefault(key, {}) - cfg[keys[-1]] = {**cfg[keys[-1]], **value} + cfg[keys[-1]] = {**cfg.get(keys[-1], {}), **value} def find(name: str, path: str) -> Optional[str]: @@ -708,6 +710,75 @@ def is_linux() -> bool: return platform.system() == "Linux" and "microsoft" not in platform.release().lower() +def detect_external_address() -> str: + """ + Detect the IP address of the host on its default-route network interface. + + No packet is sent: connecting a UDP socket only selects the outgoing interface. + + Returns: + str: IPv4 address of the default-route interface, or an empty string + if the detection fails, e.g., on a host without a default route. + """ + try: + with socket.socket(socket.AF_INET, socket.SOCK_DGRAM) as sock: + sock.connect(("8.8.8.8", 80)) + return sock.getsockname()[0] + except OSError: + return "" + + +def resolve_external_address(address: str, port: int) -> str: + """ + Determine the address advertised to benchmark functions for self-hosted storage. + + Functions running outside the Docker bridge network, e.g., in a Kubernetes + cluster, reach the storage through the host's IP and the port mapped on the host. + + Args: + address: User-provided IP or hostname, optionally with a port. When empty, + the IP of the host's default-route interface is used. + port: Port mapped on the host, appended when the address has no port. + + Returns: + str: Address in the form "host:port", or an empty string if no address was + given and the detection failed. + """ + host = address if address else detect_external_address() + if not host: + return "" + has_port = "]:" in host if host.startswith("[") else ":" in host + return host if has_port else f"{host}:{port}" + + +def probe_http(url: str, timeout_seconds: int) -> Optional[str]: + """ + Repeatedly query an HTTP endpoint until it answers with status 200. + + Args: + url: Endpoint to query + timeout_seconds: How long to keep retrying, e.g., while a server starts up + + Returns: + Optional[str]: None on success, otherwise a description of the last failure + """ + import urllib3 + + http = urllib3.PoolManager(timeout=urllib3.util.Timeout(connect=2, read=2)) + last_error = "timeout" + deadline = time.monotonic() + timeout_seconds + while time.monotonic() < deadline: + try: + resp = http.request("GET", url, retries=False) + if resp.status == 200: + return None + last_error = f"status {resp.status}" + except Exception as e: + last_error = str(e) + time.sleep(0.5) + return last_error + + def catch_interrupt() -> None: """ Set up a signal handler to catch interrupt signals (Ctrl+C).