Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion benchmark/benchmark.py
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ def run(self):
with open(config_file, "w") as fd:
yaml.dump(config_dict, fd, default_flow_style=False)

def exists(self):
def exists(self) -> bool:
return False

def compare(self, baseline):
Expand Down
162 changes: 162 additions & 0 deletions benchmark/elbencho.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
"""Elbencho S3 benchmark, driven by the shared Workloads pipeline.
Comment thread
gitkenan marked this conversation as resolved.

The generated commands are fanned out to the client nodes through the pdsh-free
``RemoteExecutor``."""

import logging
Comment thread
gitkenan marked this conversation as resolved.
import os

import yaml

from monitoring.monitoring_factory import MonitoringFactory
import settings
from remote.async_ssh import AsyncSSHExecutor
from remote.remote_executor import RemoteExecutor

from .benchmark import Benchmark

logger = logging.getLogger("cbt")


class Elbencho(Benchmark):

def __init__(self, archive_dir: str, cluster, config: dict) -> None:
# auth comes in from the YAML as a nested dict; flatten it to
# strings now so the Workloads pipeline (which stringifies everything)
# doesn't mangle it.
self.auth = config.get("auth", {})
config["s3_auth_config"] = self.auth.get("config", "")
config["s3_session_token"] = self.auth.get("s3_session_token", "")
config.pop("auth", None)

super().__init__(archive_dir, cluster, config)

self.cmd_path = config.get("cmd_path", "/usr/local/bin/elbencho")

# RemoteExecutor is the ABC which allows us to easily
# swap out AsyncIO as the fan-out tool later if needed.
self._remote: RemoteExecutor = AsyncSSHExecutor()

self.base_run_dir = self.run_dir

workloads = config.get("workloads", {})
if not isinstance(workloads, dict):
raise ValueError(f"workloads must be a dict, got {type(workloads).__name__}")
self._validate_workloads(workloads)

for wl_name, wl_params in workloads.items():
logger.info("Elbencho workload '%s': %s", wl_name, wl_params)

# ------------------------------------------------------------------
# Lifecycle overrides (pdsh-free)
# ------------------------------------------------------------------

def exists(self) -> bool:
if os.path.exists(self.archive_dir):
logger.info("Skipping existing Elbencho results in %s.", self.archive_dir)
return True
return False

def initialize(self) -> None:
super().initialize()

logger.info("Verifying elbencho binary is executable on all client nodes: %s", self.cmd_path)
self._remote.run_command_with_error_checking(settings.getnodes('clients'), f"test -x {self.cmd_path}")

self.cleandir()

if not os.path.exists(self.archive_dir):
os.makedirs(self.archive_dir)

def cleandir(self) -> None:
clients = settings.getnodes('clients')
self._remote.clean_remote_dir(clients, self.run_dir)
self._remote.make_remote_dir(clients, self.run_dir)

def dropcaches(self) -> None:
nodes = settings.getnodes('clients', 'osds')
self._remote.run_command(nodes, 'sync', continue_if_error=False)
self._remote.run_command(
nodes,
'echo 3 | sudo tee /proc/sys/vm/drop_caches',
continue_if_error=False,
)

def run(self) -> None:
if self.osd_ra and self.osd_ra_changed:
logger.info('Setting OSD Read Ahead to: %s', self.osd_ra)
self.cluster.set_osd_param('read_ahead_kb', self.osd_ra)

config_file = os.path.join(self.archive_dir, 'benchmark_config.yaml')
if not os.path.exists(self.archive_dir):
os.makedirs(self.archive_dir)
if not os.path.exists(config_file):
config_dict = dict(cluster=self.config)
with open(config_file, 'w') as fd:
yaml.dump(config_dict, fd, default_flow_style=False)

if not self._workloads.exist():
logger.warning("Elbencho: no workloads defined — nothing to run.")
return

self.dropcaches()
# TODO: call super().run() once Benchmark.run() is executor-driven and
# is no longer using AsyncIO.
self._remote.make_remote_dir(settings.getnodes('clients'), self.run_dir)
self.cluster.dump_config(self.run_dir)

self._run_workloads()

self._remote.sync_files(settings.getnodes('clients'), self.run_dir, self.archive_dir)

def cleanup(self) -> None:
pass

# ------------------------------------------------------------------
# Validation
# ------------------------------------------------------------------

def _validate_workloads(self, workloads: dict) -> None:
"""Fail early with a precise error rather than mid-run inside the pipeline."""
for name, params in workloads.items():
if not isinstance(params, dict):
raise ValueError(f"workload '{name}' must be a dict")
for field in ("mode", "s3_bucket"):
if field not in params:
raise ValueError(f"workload '{name}' missing required key '{field}'")
for field in ("threads", "iodepth"):
values = params.get(field)
if values is None:
continue
for value in (values if isinstance(values, list) else [values]):
try:
int(value)
except (TypeError, ValueError):
raise ValueError(
f"workload '{name}': {field} value {value!r} is not an integer"
)

# ------------------------------------------------------------------
# Run loop
# ------------------------------------------------------------------

def _run_workloads(self) -> None:
clients = settings.getnodes("clients")

self._workloads.set_benchmark_type("elbencho")
self._workloads.set_executable(self.cmd_path)

for output_directory, commands in self._workloads.command_groups():
live_commands = [cmd for cmd in commands if cmd]
if not live_commands:
continue

self._remote.make_remote_dir(settings.getnodes('clients'), output_directory)
logger.info("Elbencho: running %d command(s) → %s", len(live_commands), output_directory)
MonitoringFactory.start(output_directory)
for cmd in live_commands:
logger.debug("Elbencho cmd: %s", cmd)
self._remote.run_command(clients, cmd, continue_if_error=False)
MonitoringFactory.stop()

logger.info("Elbencho: all workloads complete.")
4 changes: 3 additions & 1 deletion benchmarkfactory.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import settings
from common import all_configs
from benchmark.radosbench import Radosbench
from benchmark.elbencho import Elbencho
from benchmark.fio import Fio
from benchmark.hsbench import Hsbench
from benchmark.rbdfio import RbdFio
Expand Down Expand Up @@ -32,7 +33,8 @@ def get_object(archive, cluster, benchmark, bconfig):
'librbdfio': LibrbdFio,
'cosbench': Cosbench,
'cephtestrados': CephTestRados,
'getput': Getput}
'getput': Getput,
'elbencho': Elbencho}
try:
return benchmarks[benchmark](archive, cluster, bconfig)
except KeyError:
Expand Down
180 changes: 180 additions & 0 deletions command/elbencho_command.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
"""Builds the elbencho command line for a single S3 workload instance.

It returns the full executable string that can be used to run a cli command.
It is instantiated by ``Workload._create_command_class`` as part of the
shared Workloads pipeline."""

import os
import re
import shlex
from logging import Logger, getLogger

from cli_options import CliOptions
from command.command import Command

log: Logger = getLogger("cbt")

_BS_SUFFIXES = {"k": 1024, "m": 1024 ** 2, "g": 1024 ** 3}


class ElbenchoCommand(Command):
"""A single elbencho S3 command line for one run cell."""

_MODE_FLAGS = {
"write": ["--write"],
"read": ["--read"],
"readwrite": ["--write", "--read"],
"stat": ["--stat"],
"list": ["--s3listobjpar"],
}

_MODES_NO_BLOCKSIZE = {"stat", "list"}

def __init__(self, options: dict[str, str], workload_output_directory: str) -> None:
# Must be set before super().__init__() because the base constructor calls _parse_options().
self._workload_output_directory: str = workload_output_directory
super().__init__(options)

@classmethod
def mode_is_supported(cls, mode: str) -> bool:
"""Return True if mode is currently supported (stat/list are not yet)."""
return mode not in cls._MODES_NO_BLOCKSIZE

@staticmethod
def parse_blocksize_to_bytes(blocksize: str) -> int:
s = str(blocksize).strip().lower()
m = re.fullmatch(r"(\d+(?:\.\d+)?)([kmg]?)", s)
if not m:
raise ValueError(f"Unrecognised blocksize format: {blocksize!r}")
value, suffix = m.group(1), m.group(2)
return int(float(value) * _BS_SUFFIXES.get(suffix, 1))

@staticmethod
def build_auth_flags(auth: dict[str, str]) -> list[str]:
"""Return S3 auth flags for whatever credentials are present; empty auth yields no flags."""
flags: list[str] = []

config_str = auth.get("config", "")
if config_str:
pairs = dict(
kv.split("=", 1)
for kv in config_str.split(";")
if "=" in kv
)
if "url" in pairs:
flags += ["--s3endpoints", pairs["url"]]
if "access_key" in pairs:
flags += ["--s3key", pairs["access_key"]]
if "secret_key" in pairs:
flags += ["--s3secret", pairs["secret_key"]]

token = auth.get("s3_session_token", "")
if token:
flags += ["--s3authtoken", token]

return flags

# ------------------------------------------------------------------
# Command ABC implementation
# ------------------------------------------------------------------

def _parse_options(self, options: dict[str, str]) -> CliOptions:
# Populate CliOptions with parsed options and defaults.
parsed_options: CliOptions = CliOptions()

parsed_options["mode"] = options.get("mode")
parsed_options["s3_bucket"] = options.get("s3_bucket")
parsed_options["s3_region"] = options.get("s3_region", "default")
parsed_options["threads"] = str(options.get("threads", 1))
parsed_options["blocksize"] = str(options.get("blocksize", "4k"))
parsed_options["iodepth"] = str(options.get("iodepth", 1))

# Optional value flags
for key in ("size", "num_objects", "num_dirs", "duration", "hosts"):
value = options.get(key)
parsed_options[key] = str(value) if value is not None else None

# Boolean presence flags: "true" when truthy, None otherwise.
for key in ("deldirs", "s3nompcheck", "mkdirs"):
parsed_options[key] = "true" if options.get(key) else None

parsed_options["s3_auth_config"] = options.get("s3_auth_config", "")
parsed_options["s3_session_token"] = options.get("s3_session_token", "")

return parsed_options

def _parse_global_options(self, options: dict[str, str]) -> CliOptions:
return CliOptions(options)

def _generate_output_directory_path(self) -> str:
# {base}/elbencho/{mode}_{blocksize_bytes}/threads-{NNN}/iodepth-{NNN}
options = self._options
mode = str(options["mode"])
blocksize = str(options["blocksize"])
threads = int(str(options["threads"]))
iodepth = int(str(options["iodepth"]))
return os.path.join(
self._workload_output_directory,
self.benchmark,
f"{mode}_{self.parse_blocksize_to_bytes(blocksize)}",
f"threads-{threads:03d}",
f"iodepth-{iodepth:03d}",
)

@property
def benchmark(self) -> str:
return "elbencho"

def _generate_full_command(self) -> str:
if self._executable is None:
return ""

options = self._options
mode = str(options["mode"])

if not self.mode_is_supported(mode):
log.warning(
"Elbencho: mode '%s' is not yet supported by the formatter. "
"Skipping run (blocksize=%s, threads=%s, iodepth=%s).",
mode, options["blocksize"], options["threads"], options["iodepth"],
)
return ""

mode_flags = self._MODE_FLAGS.get(mode)
if mode_flags is None:
raise ValueError(f"Unknown elbencho mode: {mode!r}")

cmd_parts: list[str] = [self._executable]
cmd_parts += mode_flags
cmd_parts += ["--threads", str(options["threads"])]
cmd_parts += ["--block", str(options["blocksize"])]
cmd_parts += ["--iodepth", str(options["iodepth"])]

if options["size"] is not None:
cmd_parts += ["--size", str(options["size"])]
if options["num_objects"] is not None:
cmd_parts += ["--files", str(options["num_objects"])]
if options["num_dirs"] is not None:
cmd_parts += ["--dirs", str(options["num_dirs"])]
if options["duration"] is not None:
cmd_parts += ["--timelimit", str(options["duration"])]
if options["deldirs"]:
cmd_parts += ["--deldirs"]
if options["s3nompcheck"]:
cmd_parts += ["--s3nompcheck"]
if options["hosts"] is not None:
cmd_parts += ["--hosts", str(options["hosts"])]

cmd_parts += self.build_auth_flags({
"config": options["s3_auth_config"] or "",
"s3_session_token": options["s3_session_token"] or "",
})
cmd_parts += ["--s3region", str(options["s3_region"])]
cmd_parts += ["--resfile", os.path.join(self._generate_output_directory_path(), "result.csv")]

if options["mkdirs"]:
cmd_parts += ["--mkdirs"]
cmd_parts += [f"s3://{options['s3_bucket']}"]

# Safely quote arguments for remote shell execution.
return shlex.join(cmd_parts)
Loading