diff --git a/monitoring/blktrace_monitoring.py b/monitoring/blktrace_monitoring.py index 0abb7de4..2a512cbf 100644 --- a/monitoring/blktrace_monitoring.py +++ b/monitoring/blktrace_monitoring.py @@ -21,23 +21,28 @@ def __init__(self, mconfig: dict[str, Any]) -> None: # Remove these casts once settings.cluster.get() is fully typed. self._osds_per_node = cast(int, settings.cluster.get("osds_per_node")) self._use_existing = cast(bool, settings.cluster.get("use_existing", True)) - self._user = cast(str, settings.cluster.get("user")) def start(self, directory: str) -> None: """Create the blktrace output directory and start tracing on each OSD device.""" + self._check_tool("blktrace") blktrace_dir = f"{directory}/blktrace" - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {blktrace_dir}").communicate() # type: ignore[no-untyped-call] + self._make_remote_dir(blktrace_dir) for device in range(self._osds_per_node): - common.pdsh( # type: ignore[no-untyped-call] - self._nodes, - f"cd {blktrace_dir};sudo blktrace -o device{device} -d /dev/disk/by-partlabel/osd-device-{device}-data", + cmd = ( + f"cd {blktrace_dir} && sudo blktrace" + f" -o device{device} -d /dev/disk/by-partlabel/osd-device-{device}-data" ) + common.pdsh(self._nodes, cmd) # type: ignore[no-untyped-call] + logger.info("Blktrace monitoring running in background for %d OSD devices per node.", self._osds_per_node) def stop(self, directory: Optional[str]) -> None: """Stop blktrace and optionally generate seekwatcher movies.""" common.pdsh(self._nodes, "sudo pkill -SIGINT -f blktrace").communicate() # type: ignore[no-untyped-call] + logger.info("Blktrace monitoring stopped.") if directory and not self._use_existing: + logger.info("Generating blktrace seekwatcher movies.") self._make_movies(directory) + logger.info("Blktrace seekwatcher movie generation complete.") def _make_movies(self, directory: str) -> None: """Generate an mpg movie for each OSD device using seekwatcher.""" @@ -46,5 +51,5 @@ def _make_movies(self, directory: str) -> None: for device in range(self._osds_per_node): common.pdsh( # type: ignore[no-untyped-call] self._nodes, - f"cd {blktrace_dir};{seekwatcher} -t device{device} -o device{device}.mpg --movie", + f"cd {blktrace_dir} && {seekwatcher} -t device{device} -o device{device}.mpg --movie", ).communicate() diff --git a/monitoring/collectl_monitoring.py b/monitoring/collectl_monitoring.py index a97a32b2..ee62591d 100644 --- a/monitoring/collectl_monitoring.py +++ b/monitoring/collectl_monitoring.py @@ -14,21 +14,24 @@ class CollectlMonitoring(Monitoring): "-s+mYZ -i 1:10 -F0 -f {collectl_dir} " r"--rawdskfilt \"+cciss/c\d+d\d+ |hd[ab] | sd[a-z]+ |dm-\d+ |" r"xvd[a-z]+ |fio[a-z]+ | vd[a-z]+ |emcpower[a-z]+ |psv\d+ |" - r"nvme[0-9]n[0-9]+p[0-9]+ \"" + r"nvme[0-9]n[0-9]+ |nvme[0-9]n[0-9]+p[0-9]+ \"" ) def __init__(self, mconfig: dict[str, Any]) -> None: """Initialize collectl monitoring configuration.""" super().__init__(mconfig) - self._args = mconfig.get("args", self.DEFAULT_ARGS) + args = mconfig.get("args", self.DEFAULT_ARGS) + self._args = args if isinstance(args, str) else self.DEFAULT_ARGS def start(self, directory: str) -> None: """Create the output directory and start collectl.""" + if not self._check_tool("collectl", fatal=False): + return collectl_dir = f"{directory}/collectl" - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {collectl_dir}").communicate() # type: ignore[no-untyped-call] - common.pdsh( - self._nodes, ["collectl", self._args.format(collectl_dir=collectl_dir)] - ) # type: ignore[no-untyped-call] + self._make_remote_dir(collectl_dir) + # collectl runs as a background daemon; do not call .communicate() here. + cmd = f"collectl {self._args.format(collectl_dir=collectl_dir)}" + common.pdsh(self._nodes, cmd) # type: ignore[no-untyped-call] def stop(self, directory: Optional[str]) -> None: """Stop running collectl processes.""" diff --git a/monitoring/monitoring.py b/monitoring/monitoring.py index 33e5814e..0230f692 100644 --- a/monitoring/monitoring.py +++ b/monitoring/monitoring.py @@ -1,10 +1,14 @@ """Base abstractions for monitoring backends.""" +import logging from abc import ABC, abstractmethod from typing import Any, ClassVar, Optional, cast +import common import settings +logger = logging.getLogger("cbt") + class Monitoring(ABC): """Abstract base class for monitoring backends.""" @@ -12,10 +16,47 @@ class Monitoring(ABC): DEFAULT_NODES: ClassVar[list[str]] def __init__(self, mconfig: dict[str, Any]) -> None: - """Resolve monitoring nodes from configuration or subclass defaults.""" + """Resolve monitoring nodes and common cluster settings from configuration.""" nodes_list = mconfig.get("nodes", self.DEFAULT_NODES) - # Remove this cast and ignore once settings.getnodes() is fully typed. + # Remove these casts once settings is fully typed. self._nodes = cast(str, settings.getnodes(*nodes_list)) # type: ignore[no-untyped-call] + self._user = cast(str, settings.cluster.get("user")) + + def _make_remote_dir(self, path: str) -> None: + """Create *path* on all monitored nodes, waiting for completion.""" + common.pdsh(self._nodes, f"mkdir -p -m0755 -- {path}").communicate() # type: ignore[no-untyped-call] + + def _check_tool(self, tool_name: str, *, fatal: bool = True) -> bool: + """Verify *tool_name* is present on all monitored nodes. + + Runs ``command -v `` on every node via pdsh. The behaviour on + failure depends on *fatal*: + + * ``fatal=True`` (default) — raises :exc:`RuntimeError`, aborting the + run. Use this for tools that are explicitly configured and required. + * ``fatal=False`` — logs a warning and returns ``False`` so the caller + can skip further work gracefully. Use this for tools with a built-in + default that may simply not be installed. + + Returns: + ``True`` if the tool was found on all nodes, ``False`` otherwise + (only reachable when *fatal* is ``False``). + + Raises: + RuntimeError: If the tool is not found and *fatal* is ``True``. + """ + try: + common.pdsh(self._nodes, f"command -v {tool_name}", continue_if_error=False).communicate() # type: ignore[no-untyped-call] + return True + except Exception as exc: + msg = ( + f"Monitoring tool '{tool_name}' not found on one or more nodes " + f"({self._nodes}). Install it before running CBT with this monitoring backend." + ) + if fatal: + raise RuntimeError(msg) from exc + logger.warning("%s — skipping %s monitoring.", msg, tool_name) + return False @abstractmethod def start(self, directory: str) -> None: diff --git a/monitoring/monitoring_factory.py b/monitoring/monitoring_factory.py index 96c7d090..aba6c656 100644 --- a/monitoring/monitoring_factory.py +++ b/monitoring/monitoring_factory.py @@ -9,8 +9,10 @@ from monitoring.blktrace_monitoring import BlktraceMonitoring from monitoring.collectl_monitoring import CollectlMonitoring from monitoring.monitoring import Monitoring -from monitoring.perf_monitoring import OsdPerfMonitoring, PerfMonitoring -from monitoring.top_monitoring import OsdTopMonitoring, TopMonitoring +from monitoring.osd_perf_monitoring import OsdPerfMonitoring +from monitoring.osd_top_monitoring import OsdTopMonitoring +from monitoring.perf_monitoring import PerfMonitoring +from monitoring.top_monitoring import TopMonitoring logger = logging.getLogger("cbt") diff --git a/monitoring/osd_perf_monitoring.py b/monitoring/osd_perf_monitoring.py new file mode 100644 index 00000000..e3d30722 --- /dev/null +++ b/monitoring/osd_perf_monitoring.py @@ -0,0 +1,22 @@ +"""OsdPerfMonitoring: PerfMonitoring specialised for Ceph OSD processes.""" + +from monitoring.osd_pid_monitoring import OsdPidMonitoring +from monitoring.perf_monitoring import PerfMonitoring + + +class OsdPerfMonitoring(OsdPidMonitoring, PerfMonitoring): + """PerfMonitoring specialised for Ceph OSD processes. + + Discovers the target PIDs by scanning PID files matching ``pid_glob`` + inside ``pid_dir`` (read from ``settings.cluster``), then launches a + separate ``perf`` invocation per OSD PID. + """ + + def start(self, directory: str) -> None: + """Create the perf output directory and start a perf instance per OSD PID.""" + self._check_tool(self._perf_cmd.split()[-1]) + perf_dir = f"{directory}/perf" + self._perf_dir = perf_dir + self._make_remote_dir(perf_dir) + cmd_template = f"{self._perf_cmd} {self._args_template}" + self._start_per_pid(perf_dir, "perf", cmd_template, self._perf_runners) diff --git a/monitoring/osd_pid_monitoring.py b/monitoring/osd_pid_monitoring.py new file mode 100644 index 00000000..812a9072 --- /dev/null +++ b/monitoring/osd_pid_monitoring.py @@ -0,0 +1,73 @@ +"""OsdPidMonitoring base class for per-OSD-PID monitoring backends.""" + +import glob as _glob +import logging +import os +from abc import ABC +from typing import Any, cast + +import common +import settings +from monitoring.monitoring import Monitoring + +logger = logging.getLogger("cbt") + + +class OsdPidMonitoring(Monitoring, ABC): + """Mixin base for monitoring backends that launch one process per OSD PID. + + Subclasses must supply :attr:`_pid_dir`, :attr:`_pid_glob`, and a + *runners* list attribute, then call :meth:`_start_per_pid` from their + ``start()`` implementation after creating the output directory. + """ + + def __init__(self, mconfig: dict[str, Any]) -> None: + """Read OSD PID discovery settings from *mconfig* and ``settings.cluster``.""" + super().__init__(mconfig) + self._pid_dir = cast(str, settings.cluster.get("pid_dir")) + self._pid_glob = mconfig.get("pid_glob", "osd.*.pid") + + # pylint: disable=too-many-locals + def _start_per_pid(self, output_dir: str, tool_name: str, cmd_template: str, runners: list[Any]) -> None: + """Launch *cmd_template* once per OSD PID, populating *runners*. + + *cmd_template* is a :meth:`str.format` template that accepts + ``{pid}`` and ``{output_dir}`` keyword arguments. + + On a local node each PID file under :attr:`_pid_dir` is read and the + command is spawned via :func:`common.sh`. On remote nodes a single + ``pdsh`` shell loop iterates over the PID files instead. + """ + pid_glob_path = f"{self._pid_dir}/{self._pid_glob}" + local_node = common.get_localnode(self._nodes) # type: ignore[no-untyped-call] + if local_node: + logger.debug("%s: local_node pid_dir=%s", type(self).__name__, self._pid_dir) + pid_paths = _glob.glob(os.path.join(self._pid_dir, self._pid_glob)) + if not pid_paths: + logger.warning( + "%s: no PID files matched %s in %s — no %s processes started", + type(self).__name__, + self._pid_glob, + self._pid_dir, + tool_name, + ) + for pid_path in pid_paths: + with open(pid_path, encoding="utf-8") as pidfile: + pid = pidfile.read().strip() + cmd = cmd_template.format(output_dir=output_dir, pid=pid) + runner = common.sh(local_node, cmd) # type: ignore[no-untyped-call] + runners.append(runner) + else: + logger.debug("%s: remote_node", type(self).__name__) + ls_runner = common.pdsh(self._nodes, f"ls {pid_glob_path} 2>/dev/null") # type: ignore[no-untyped-call] + stdout, _ = ls_runner.communicate() + if not stdout.strip(): + logger.warning( + "%s: no PID files matched %s on remote nodes — no %s processes started", + type(self).__name__, + pid_glob_path, + tool_name, + ) + cmd = cmd_template.format(output_dir=output_dir, pid='"$pid"') + loop_cmd = f'for f in {pid_glob_path}; do pid=$(cat "$f"); {cmd} & done' + runners.append(common.pdsh(self._nodes, loop_cmd)) # type: ignore[no-untyped-call] diff --git a/monitoring/osd_top_monitoring.py b/monitoring/osd_top_monitoring.py new file mode 100644 index 00000000..81974107 --- /dev/null +++ b/monitoring/osd_top_monitoring.py @@ -0,0 +1,30 @@ +"""OsdTopMonitoring: TopMonitoring specialised for Ceph OSD processes.""" + +from typing import Any + +from monitoring.osd_pid_monitoring import OsdPidMonitoring +from monitoring.top_monitoring import TopMonitoring + + +class OsdTopMonitoring(OsdPidMonitoring, TopMonitoring): + """TopMonitoring specialised for Ceph OSD processes. + + Discovers the target PIDs by scanning PID files matching ``pid_glob`` + inside ``pid_dir`` (read from ``settings.cluster``), then launches a + separate ``top`` invocation per OSD PID. + """ + + def __init__(self, mconfig: dict[str, Any]) -> None: + """Initialize OSD top monitoring configuration.""" + super().__init__(mconfig) + # Override default args to include per-pid placeholders. + # NOTE: see TopMonitoring.__init__ comment regarding Irix/Solaris mode. + self._args = mconfig.get("args", "-b -H -1 -p {pid} -n 30 > {output_dir}/{pid}_osd_top.out") + + def start(self, directory: str) -> None: + """Create the top output directory and start a top instance per OSD PID.""" + self._check_tool(self._top_cmd) + top_dir = f"{directory}/top" + self._make_remote_dir(top_dir) + cmd_template = f"{self._top_cmd} {self._args}" + self._start_per_pid(top_dir, "top", cmd_template, self._top_runners) diff --git a/monitoring/perf_monitoring.py b/monitoring/perf_monitoring.py index bad76ae0..1931392c 100644 --- a/monitoring/perf_monitoring.py +++ b/monitoring/perf_monitoring.py @@ -1,13 +1,13 @@ """Perf monitoring backend.""" -import glob +# pylint: disable=duplicate-code + import logging import os import re from typing import Any, ClassVar, Optional, cast import common -import settings from monitoring.monitoring import Monitoring logger = logging.getLogger("cbt") @@ -25,92 +25,69 @@ class PerfMonitoring(Monitoring): def __init__(self, mconfig: dict[str, Any]) -> None: """Initialize perf monitoring configuration.""" super().__init__(mconfig) - # Remove these casts once settings.cluster.get() is fully typed. - self._user = cast(str, settings.cluster.get("user")) + if "args" not in mconfig: + raise ValueError("PerfMonitoring requires 'args' in mconfig") self._perf_cmd = mconfig.get("perf_cmd", "sudo perf") - self._args_template = mconfig.get("args") + self._args_template: str = mconfig["args"] self._perf_runners: list[Any] = [] - self._perf_dir = "" + self._perf_dir: Optional[str] = None def start(self, directory: str) -> None: """Create the perf output directory and start perf collection.""" + self._check_tool(self._perf_cmd.split()[-1]) perf_dir = f"{directory}/perf" self._perf_dir = perf_dir - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {perf_dir}").communicate() # type: ignore[no-untyped-call] + self._make_remote_dir(perf_dir) - perf_cmd = f"{self._perf_cmd} {self._args_template} &".format(perf_dir=perf_dir) + perf_template = f"{self._perf_cmd} {self._args_template}" + perf_cmd = perf_template.format(perf_dir=perf_dir) local_node = common.get_localnode(self._nodes) # type: ignore[no-untyped-call] if local_node: runner = common.sh(local_node, perf_cmd) # type: ignore[no-untyped-call] self._perf_runners.append(runner) else: - common.pdsh(self._nodes, perf_cmd) # type: ignore[no-untyped-call] + runner = common.pdsh(self._nodes, perf_cmd) # type: ignore[no-untyped-call] + self._perf_runners.append(runner) + logger.info("Perf monitoring running in background (will be killed on stop).") def stop(self, directory: Optional[str]) -> None: """Stop perf collection and adjust file ownership when needed.""" - if self._perf_runners: - for runner in self._perf_runners: + common.pdsh(self._nodes, "sudo pkill -SIGINT -f 'perf '").communicate() # type: ignore[no-untyped-call] + for runner in self._perf_runners: + try: runner.kill() - else: - common.pdsh(self._nodes, r"sudo pkill -SIGINT -f perf\ ").communicate() # type: ignore[no-untyped-call] + except OSError: + pass if directory: common.pdsh( # type: ignore[no-untyped-call] - self._nodes, f"sudo chown {self._user}.{self._user} {directory}/perf/perf.data" - ) + self._nodes, f"sudo chown {self._user}:{self._user} {directory}/perf/perf.data" + ).communicate() common.pdsh( # type: ignore[no-untyped-call] - self._nodes, f"sudo chown {self._user}.{self._user} {directory}/perf/perf_stat.*" - ) + self._nodes, f"sudo chown {self._user}:{self._user} {directory}/perf/perf_stat.*" + ).communicate() + logger.info("Perf monitoring stopped.") def get_cpu_cycles(self, out_dir: str) -> Optional[int]: """Return total CPU cycles from perf stat output, if available.""" - total_cpu_cycles = 0 - perf_dir_name = str(glob.glob(out_dir + "/perf*")[0]) + perf_dir_name = os.path.join(out_dir, "perf") + if not os.path.isdir(perf_dir_name): + logger.warning("get_cpu_cycles: perf directory not found: %s", perf_dir_name) + return None perf_stat_fnames = os.listdir(perf_dir_name) + if not perf_stat_fnames: + logger.warning("get_cpu_cycles: no perf stat files found in %s", perf_dir_name) + return None + total_cpu_cycles = 0 for perf_out_fname in perf_stat_fnames: with open(f"{perf_dir_name}/{perf_out_fname}", encoding="utf-8") as perf_output_file: match = re.search(r"(.*) cycles(.*?) .*", perf_output_file.read(), re.M | re.I) if match: cpu_cycles = match.group(1).strip() else: + logger.warning( + "get_cpu_cycles: no cycles line found in %s — returning None", + perf_out_fname, + ) return None total_cpu_cycles = total_cpu_cycles + int(cpu_cycles.replace(",", "")) return cast(Optional[int], total_cpu_cycles) - - -class OsdPerfMonitoring(PerfMonitoring): - """PerfMonitoring specialised for Ceph OSD processes. - - Discovers the target PIDs by scanning PID files matching ``pid_glob`` - inside ``pid_dir`` (read from ``settings.cluster``), then launches a - separate ``perf`` invocation per OSD PID. - """ - - def __init__(self, mconfig: dict[str, Any]) -> None: - """Initialize OSD perf monitoring configuration.""" - super().__init__(mconfig) - self._pid_dir = cast(str, settings.cluster.get("pid_dir")) - self._pid_glob = mconfig.get("pid_glob", "osd.*.pid") - - def start(self, directory: str) -> None: - """Create the perf output directory and start a perf instance per OSD PID.""" - perf_dir = f"{directory}/perf" - self._perf_dir = perf_dir - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {perf_dir}").communicate() # type: ignore[no-untyped-call] - - perf_template = f"{self._perf_cmd} {self._args_template} &" - local_node = common.get_localnode(self._nodes) # type: ignore[no-untyped-call] - if local_node: - logger.debug("OsdPerfMonitoring: local_node pid_dir=%s", self._pid_dir) - for pid_path in glob.glob(os.path.join(self._pid_dir, self._pid_glob)): - with open(pid_path, encoding="utf-8") as pidfile: - pid = pidfile.read().strip() - perf_cmd = perf_template.format(perf_dir=perf_dir, pid=pid) - runner = common.sh(local_node, perf_cmd) # type: ignore[no-untyped-call] - self._perf_runners.append(runner) - else: - logger.debug("OsdPerfMonitoring: remote_node") - perf_cmd = perf_template.format(perf_dir=perf_dir, pid="${pid}") - common.pdsh( # type: ignore[no-untyped-call] - self._nodes, - [f"for pid in `cat {self._pid_dir}/{self._pid_glob}`;", "do", perf_cmd, ";", "done"], - ) diff --git a/monitoring/top_monitoring.py b/monitoring/top_monitoring.py index 657a2920..55ae9f22 100644 --- a/monitoring/top_monitoring.py +++ b/monitoring/top_monitoring.py @@ -1,13 +1,10 @@ """Top monitoring backend.""" -import glob import logging -import os import re -from typing import Any, ClassVar, Optional, cast +from typing import Any, ClassVar, Optional import common -import settings from monitoring.monitoring import Monitoring logger = logging.getLogger("cbt") @@ -18,6 +15,10 @@ def _estimate_top_duration(args: str) -> Optional[float]: Returns the estimated duration in seconds, or ``None`` if ``-n`` is not present (meaning top will run indefinitely and must be killed to stop). + + When ``-d`` is absent, this assumes the procps-ng compile-time default of + three seconds. The actual delay can be configured in ``~/.toprc``; pass + ``-d `` explicitly for an accurate estimate. """ n_match = re.search(r"-n\s+(\d+)", args) d_match = re.search(r"-d\s+([\d.]+)", args) @@ -41,7 +42,6 @@ class TopMonitoring(Monitoring): def __init__(self, mconfig: dict[str, Any]) -> None: """Initialize top monitoring configuration.""" super().__init__(mconfig) - self._user = cast(str, settings.cluster.get("user")) self._top_cmd = mconfig.get("top_cmd", "top") # NOTE: top's %CPU column behaviour depends on the Irix/Solaris mode # toggle ('I' key / Mode_irixps in ~/.toprc). The procps-ng build @@ -50,13 +50,17 @@ def __init__(self, mconfig: dict[str, Any]) -> None: # off in their ~/.toprc the resulting CPU figures will be incorrect. self._args = mconfig.get("args", "-b -H -1 -n 30 > {top_dir}/top.out") self._top_runners: list[Any] = [] + self._running_cmd: Optional[str] = None def start(self, directory: str) -> None: """Create the top output directory and start top collection.""" + self._check_tool(self._top_cmd) top_dir = f"{directory}/top" - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {top_dir}").communicate() # type: ignore[no-untyped-call] + self._make_remote_dir(top_dir) - top_cmd = f"{self._top_cmd} {self._args}".format(top_dir=top_dir) + top_template = f"{self._top_cmd} {self._args}" + top_cmd = top_template.format(top_dir=top_dir) + self._running_cmd = top_cmd # stored for use in stop() local_node = common.get_localnode(self._nodes) # type: ignore[no-untyped-call] if local_node: runner = common.sh(local_node, top_cmd) # type: ignore[no-untyped-call] @@ -65,91 +69,29 @@ def start(self, directory: str) -> None: duration = _estimate_top_duration(self._args) if duration is not None: logger.info( - "Top monitoring collecting %s samples (estimated ~%.0fs)...", + "Top monitoring running in background (%s samples, estimated ~%.0fs; " + "use -d for an explicit delay).", self._args.split("-n")[1].split()[0].strip(), duration, ) else: - logger.info("Top monitoring running (will be killed on stop)...") - common.pdsh(self._nodes, top_cmd).communicate() # type: ignore[no-untyped-call] - logger.info("Top monitoring collection complete.") + logger.info("Top monitoring running in background (will be killed on stop).") + runner = common.pdsh(self._nodes, top_cmd) # type: ignore[no-untyped-call] + self._top_runners.append(runner) def stop(self, directory: Optional[str]) -> None: """Stop top collection and adjust file ownership when needed.""" - if self._top_runners: - for runner in self._top_runners: - runner.kill() - else: - pkill_cmd = f"sudo pkill -SIGINT -f '{self._top_cmd} {self._args}'" + if self._running_cmd: + pkill_cmd = f"sudo pkill -SIGINT -f '{self._running_cmd}'" common.pdsh(self._nodes, pkill_cmd).communicate() # type: ignore[no-untyped-call] + for runner in self._top_runners: + try: + runner.kill() + except OSError: + pass if directory: common.pdsh( # type: ignore[no-untyped-call] self._nodes, - f"sudo chown {self._user}.{self._user} {directory}/top/*top.out", + f"sudo find {directory}/top -maxdepth 1 -name '*top.out' -exec chown {self._user}:{self._user} {{}} +", ) - - -class OsdTopMonitoring(TopMonitoring): - """TopMonitoring specialised for Ceph OSD processes. - - Discovers the target PIDs by scanning PID files matching ``pid_glob`` - inside ``pid_dir`` (read from ``settings.cluster``), then launches a - separate ``top`` invocation per OSD PID. - """ - - def __init__(self, mconfig: dict[str, Any]) -> None: - """Initialize OSD top monitoring configuration.""" - super().__init__(mconfig) - self._pid_dir = cast(str, settings.cluster.get("pid_dir")) - self._pid_glob = mconfig.get("pid_glob", "osd.*.pid") - # Override default args to include per-pid placeholders. - self._args = mconfig.get("args", "-b -H -1 -p {pid} -n 30 > {top_dir}/{pid}_osd_top.out") - # NOTE: see TopMonitoring.__init__ comment regarding Irix/Solaris mode. - - def start(self, directory: str) -> None: - """Create the top output directory and start a top instance per OSD PID.""" - top_dir = f"{directory}/top" - common.pdsh(self._nodes, f"mkdir -p -m0755 -- {top_dir}").communicate() # type: ignore[no-untyped-call] - - top_template = f"{self._top_cmd} {self._args}" - local_node = common.get_localnode(self._nodes) # type: ignore[no-untyped-call] - if local_node: - logger.debug("OsdTopMonitoring: local_node pid_dir=%s", self._pid_dir) - pid_paths = glob.glob(os.path.join(self._pid_dir, self._pid_glob)) - if not pid_paths: - logger.warning( - "OsdTopMonitoring: no PID files matched %s in %s — no top processes started", - self._pid_glob, - self._pid_dir, - ) - for pid_path in pid_paths: - with open(pid_path, encoding="utf-8") as pidfile: - pid = pidfile.read().strip() - top_cmd = top_template.format(top_dir=top_dir, pid=pid) - runner = common.sh(local_node, top_cmd) # type: ignore[no-untyped-call] - self._top_runners.append(runner) - else: - logger.debug("OsdTopMonitoring: remote_node") - pid_glob_path = f"{self._pid_dir}/{self._pid_glob}" - ls_runner = common.pdsh(self._nodes, f"ls {pid_glob_path} 2>/dev/null") # type: ignore[no-untyped-call] - stdout, _ = ls_runner.communicate() - if not stdout.strip(): - logger.warning( - "OsdTopMonitoring: no PID files matched %s on remote nodes — no top processes started", - pid_glob_path, - ) - duration = _estimate_top_duration(self._args) - if duration is not None: - logger.info( - "OSD top monitoring collecting %s samples per OSD (estimated ~%.0fs)...", - self._args.split("-n")[1].split()[0].strip(), - duration, - ) - else: - logger.info("OSD top monitoring running per OSD (will be killed on stop)...") - top_cmd = top_template.format(top_dir=top_dir, pid="${pid}") - common.pdsh( # type: ignore[no-untyped-call] - self._nodes, - [f"for pid in `cat {pid_glob_path}`;", "do", top_cmd, ";", "done"], - ).communicate() - logger.info("OSD top monitoring collection complete.") + logger.info("Top monitoring stopped.") diff --git a/tests/test_monitoring_blktrace.py b/tests/test_monitoring_blktrace.py index 5a53c8d9..d8ff87f6 100644 --- a/tests/test_monitoring_blktrace.py +++ b/tests/test_monitoring_blktrace.py @@ -5,6 +5,8 @@ from typing import Any, Optional from unittest.mock import MagicMock, call, patch +import pytest + from monitoring.blktrace_monitoring import BlktraceMonitoring @@ -23,6 +25,7 @@ def _make_monitor( patch("monitoring.blktrace_monitoring.common.pdsh"), ): mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.return_value = user mock_settings.cluster.get.side_effect = lambda key, default=None: { "osds_per_node": osds_per_node, "use_existing": use_existing, @@ -46,6 +49,8 @@ def test_init_stores_cluster_settings() -> None: def test_start_creates_directory_and_starts_traces() -> None: """start() calls pdsh for mkdir and once per device.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/sbin/blktrace\n", "") mkdir_runner = MagicMock() trace_runner = MagicMock() @@ -55,28 +60,51 @@ def test_start_creates_directory_and_starts_traces() -> None: patch("monitoring.blktrace_monitoring.common.pdsh") as mock_pdsh, ): mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.return_value = "ceph" mock_settings.cluster.get.side_effect = lambda key, default=None: { "osds_per_node": 2, "use_existing": True, "user": "ceph", }.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, trace_runner, trace_runner] + mock_pdsh.side_effect = [check_runner, mkdir_runner, trace_runner, trace_runner] monitor = BlktraceMonitoring({}) monitor.start("/tmp/output") + mock_pdsh.assert_any_call("resolved-nodes", "command -v blktrace", continue_if_error=False) mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/blktrace") mkdir_runner.communicate.assert_called_once_with() - mock_pdsh.assert_any_call( - "resolved-nodes", - "cd /tmp/output/blktrace;sudo blktrace -o device0 -d /dev/disk/by-partlabel/osd-device-0-data", - ) - mock_pdsh.assert_any_call( - "resolved-nodes", - "cd /tmp/output/blktrace;sudo blktrace -o device1 -d /dev/disk/by-partlabel/osd-device-1-data", + trace0 = "cd /tmp/output/blktrace && sudo blktrace -o device0 -d /dev/disk/by-partlabel/osd-device-0-data" + trace1 = "cd /tmp/output/blktrace && sudo blktrace -o device1 -d /dev/disk/by-partlabel/osd-device-1-data" + mock_pdsh.assert_any_call("resolved-nodes", trace0) + mock_pdsh.assert_any_call("resolved-nodes", trace1) + assert mock_pdsh.call_count == 4 # check + mkdir + 2 devices + + +def test_start_logs_background_collection() -> None: + """start() logs the number of OSD devices being traced per node.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/sbin/blktrace\n", "") + mkdir_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.blktrace_monitoring.settings") as mock_settings, + patch("monitoring.blktrace_monitoring.common.pdsh", side_effect=[check_runner, mkdir_runner, MagicMock()]), + patch("monitoring.blktrace_monitoring.logger") as mock_logger, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.return_value = "ceph" + mock_settings.cluster.get.side_effect = lambda key, default=None: { + "osds_per_node": 1, + "use_existing": True, + "user": "ceph", + }.get(key, default) + BlktraceMonitoring({}).start("/tmp/output") + + mock_logger.info.assert_called_once_with( + "Blktrace monitoring running in background for %d OSD devices per node.", 1 ) - assert mock_pdsh.call_count == 3 # mkdir + 2 devices def test_stop_issues_pkill() -> None: @@ -95,6 +123,19 @@ def test_stop_issues_pkill() -> None: mock_movies.assert_not_called() +def test_stop_logs_completion() -> None: + """stop() logs that blktrace monitoring has stopped.""" + monitor = _make_monitor() + + with ( + patch("monitoring.blktrace_monitoring.common.pdsh"), + patch("monitoring.blktrace_monitoring.logger") as mock_logger, + ): + monitor.stop(None) + + mock_logger.info.assert_called_once_with("Blktrace monitoring stopped.") + + def test_stop_calls_make_movies_when_not_use_existing() -> None: """stop() calls _make_movies when use_existing is False and directory is provided.""" pkill_runner = MagicMock() @@ -109,6 +150,26 @@ def test_stop_calls_make_movies_when_not_use_existing() -> None: mock_movies.assert_called_once_with("/tmp/output") +def test_stop_logs_movie_generation() -> None: + """stop() logs movie generation progress when rendering is requested.""" + monitor = _make_monitor(use_existing=False) + + with ( + patch("monitoring.blktrace_monitoring.common.pdsh"), + patch.object(monitor, "_make_movies"), + patch("monitoring.blktrace_monitoring.logger") as mock_logger, + ): + monitor.stop("/tmp/output") + + mock_logger.info.assert_has_calls( + [ + call("Blktrace monitoring stopped."), + call("Generating blktrace seekwatcher movies."), + call("Blktrace seekwatcher movie generation complete."), + ] + ) + + def test_stop_does_not_call_make_movies_when_use_existing() -> None: """stop() skips _make_movies when use_existing is True.""" pkill_runner = MagicMock() @@ -148,13 +209,54 @@ def test_make_movies_issues_seekwatcher_per_device() -> None: expected_calls = [ call( "resolved-nodes", - "cd /tmp/output/blktrace;/home/ceph/bin/seekwatcher -t device0 -o device0.mpg --movie", + "cd /tmp/output/blktrace && /home/ceph/bin/seekwatcher -t device0 -o device0.mpg --movie", ), call( "resolved-nodes", - "cd /tmp/output/blktrace;/home/ceph/bin/seekwatcher -t device1 -o device1.mpg --movie", + "cd /tmp/output/blktrace && /home/ceph/bin/seekwatcher -t device1 -o device1.mpg --movie", ), ] mock_pdsh.assert_has_calls(expected_calls) assert mock_pdsh.call_count == 2 assert movie_runner.communicate.call_count == 2 + + +# --------------------------------------------------------------------------- +# Tool availability checks +# --------------------------------------------------------------------------- + + +def test_blktrace_start_raises_when_blktrace_not_installed() -> None: + """BlktraceMonitoring.start() raises RuntimeError when blktrace is not on the nodes.""" + monitor = _make_monitor() + with patch("monitoring.blktrace_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="blktrace"): + monitor.start("/tmp/output") + + +def test_blktrace_start_proceeds_when_blktrace_is_installed() -> None: + """BlktraceMonitoring.start() continues normally when blktrace is found on the nodes.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/sbin/blktrace\n", "") + mkdir_runner = MagicMock() + trace_runner = MagicMock() + + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.blktrace_monitoring.settings") as mock_settings, + patch("monitoring.blktrace_monitoring.common.pdsh") as mock_bt_pdsh, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.return_value = "ceph" + mock_settings.cluster.get.side_effect = lambda key, default=None: { + "osds_per_node": 1, + "use_existing": True, + "user": "ceph", + }.get(key, default) + mock_bt_pdsh.side_effect = [check_runner, mkdir_runner, trace_runner] + + monitor = BlktraceMonitoring({}) + monitor.start("/tmp/output") + + mock_bt_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/blktrace") diff --git a/tests/test_monitoring_collectl.py b/tests/test_monitoring_collectl.py index e5aa1d81..330240a8 100644 --- a/tests/test_monitoring_collectl.py +++ b/tests/test_monitoring_collectl.py @@ -11,16 +11,29 @@ def test_init_sets_default_args() -> None: """Use the default collectl argument string when args are not configured.""" with patch("monitoring.monitoring.settings") as mock_settings: mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" monitor = CollectlMonitoring({}) assert monitor._args == CollectlMonitoring.DEFAULT_ARGS +def test_init_uses_default_args_when_args_is_dict() -> None: + """Fall back to default args when yaml sets args to an empty dict.""" + with patch("monitoring.monitoring.settings") as mock_settings: + mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" + + monitor = CollectlMonitoring({"args": {}}) + + assert monitor._args == CollectlMonitoring.DEFAULT_ARGS + + def test_init_uses_custom_args() -> None: """Use custom collectl args from monitoring config when provided.""" with patch("monitoring.monitoring.settings") as mock_settings: mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" monitor = CollectlMonitoring({"args": "--custom {collectl_dir}"}) @@ -29,22 +42,26 @@ def test_init_uses_custom_args() -> None: def test_start_creates_directory_and_starts_collectl() -> None: """Create the collectl directory and invoke collectl through pdsh.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/collectl\n", "") mkdir_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_settings, patch("monitoring.collectl_monitoring.common.pdsh") as mock_pdsh, ): mock_settings.getnodes.return_value = "resolved-nodes" - mock_pdsh.side_effect = [mkdir_runner, MagicMock()] + mock_settings.cluster.get.return_value = "ceph" + mock_pdsh.side_effect = [check_runner, mkdir_runner, MagicMock()] monitor = CollectlMonitoring({}) monitor.start("/tmp/output") + mock_pdsh.assert_any_call("resolved-nodes", "command -v collectl", continue_if_error=False) mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/collectl") mkdir_runner.communicate.assert_called_once_with() mock_pdsh.assert_any_call( "resolved-nodes", - ["collectl", monitor._args.format(collectl_dir="/tmp/output/collectl")], + f"collectl {monitor._args.format(collectl_dir='/tmp/output/collectl')}", ) @@ -56,6 +73,7 @@ def test_stop_calls_pdsh_with_collectl_pkill() -> None: patch("monitoring.collectl_monitoring.common.pdsh", return_value=stop_runner) as mock_pdsh, ): mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" monitor = CollectlMonitoring({}) monitor.stop(None) @@ -67,3 +85,61 @@ def test_stop_calls_pdsh_with_collectl_pkill() -> None: def test_default_nodes_matches_collectl_configuration() -> None: """Expose the expected default node groups for collectl monitoring.""" assert CollectlMonitoring.DEFAULT_NODES == ["clients", "osds", "mons", "rgws"] + + +# --------------------------------------------------------------------------- +# Tool availability checks +# --------------------------------------------------------------------------- + + +def test_collectl_start_skips_when_collectl_not_installed() -> None: + """CollectlMonitoring.start() skips all work and logs a warning when collectl is absent.""" + with patch("monitoring.monitoring.settings") as mock_settings: + mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" + monitor = CollectlMonitoring({}) + + with ( + patch("monitoring.collectl_monitoring.common.pdsh") as mock_collectl_pdsh, + patch("monitoring.monitoring.logger") as mock_logger, + ): + mock_collectl_pdsh.return_value.communicate.side_effect = Exception("exit 1") + monitor.start("/tmp/output") + + # mkdir and collectl commands must never have been called (only the check) + assert mock_collectl_pdsh.call_count == 1 + mock_collectl_pdsh.assert_called_once_with("resolved-nodes", "command -v collectl", continue_if_error=False) + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("collectl" in c for c in warning_calls) + + +def test_collectl_start_does_not_raise_when_collectl_not_installed() -> None: + """CollectlMonitoring.start() must not raise even when collectl is absent.""" + with patch("monitoring.monitoring.settings") as mock_settings: + mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" + monitor = CollectlMonitoring({}) + + with patch("monitoring.collectl_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + monitor.start("/tmp/output") # must not raise + + +def test_collectl_start_proceeds_when_collectl_is_installed() -> None: + """CollectlMonitoring.start() creates directory and starts collectl when present.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/collectl\n", "") + mkdir_runner = MagicMock() + collectl_runner = MagicMock() + + with ( + patch("monitoring.monitoring.settings") as mock_settings, + patch("monitoring.collectl_monitoring.common.pdsh") as mock_pdsh, + ): + mock_settings.getnodes.return_value = "resolved-nodes" + mock_settings.cluster.get.return_value = "ceph" + mock_pdsh.side_effect = [check_runner, mkdir_runner, collectl_runner] + monitor = CollectlMonitoring({}) + monitor.start("/tmp/output") + + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/collectl") diff --git a/tests/test_monitoring_factory.py b/tests/test_monitoring_factory.py index 4746bf86..07135eed 100644 --- a/tests/test_monitoring_factory.py +++ b/tests/test_monitoring_factory.py @@ -10,8 +10,10 @@ from monitoring.blktrace_monitoring import BlktraceMonitoring from monitoring.collectl_monitoring import CollectlMonitoring from monitoring.monitoring_factory import MonitoringFactory -from monitoring.perf_monitoring import OsdPerfMonitoring, PerfMonitoring -from monitoring.top_monitoring import OsdTopMonitoring, TopMonitoring +from monitoring.osd_perf_monitoring import OsdPerfMonitoring +from monitoring.osd_top_monitoring import OsdTopMonitoring +from monitoring.perf_monitoring import PerfMonitoring +from monitoring.top_monitoring import TopMonitoring # --------------------------------------------------------------------------- # Helpers @@ -46,13 +48,10 @@ def test_get_object_returns_collectl_monitoring() -> None: def test_get_object_returns_perf_monitoring() -> None: """get_object('perf') returns a PerfMonitoring instance.""" - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - ): + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "node1" - mock_settings.cluster.get.return_value = "dummy" - instance = MonitoringFactory.get_object("perf", {}) + mock_base_settings.cluster.get.return_value = "dummy" + instance = MonitoringFactory.get_object("perf", {"args": "record"}) assert isinstance(instance, PerfMonitoring) @@ -63,6 +62,7 @@ def test_get_object_returns_blktrace_monitoring() -> None: patch("monitoring.blktrace_monitoring.settings") as mock_settings, ): mock_base_settings.getnodes.return_value = "node1" + mock_base_settings.cluster.get.return_value = "dummy" mock_settings.cluster.get.return_value = "dummy" instance = MonitoringFactory.get_object("blktrace", {}) assert isinstance(instance, BlktraceMonitoring) @@ -70,37 +70,28 @@ def test_get_object_returns_blktrace_monitoring() -> None: def test_get_object_returns_top_monitoring() -> None: """get_object('top') returns a TopMonitoring instance.""" - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - ): + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "node1" - mock_settings.cluster.get.return_value = "dummy" + mock_base_settings.cluster.get.return_value = "dummy" instance = MonitoringFactory.get_object("top", {}) assert isinstance(instance, TopMonitoring) def test_get_object_returns_osd_top_monitoring() -> None: """get_object('osd_top') returns an OsdTopMonitoring instance.""" - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - ): + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "node1" - mock_settings.cluster.get.return_value = "dummy" + mock_base_settings.cluster.get.return_value = "dummy" instance = MonitoringFactory.get_object("osd_top", {}) assert isinstance(instance, OsdTopMonitoring) def test_get_object_returns_osd_perf_monitoring() -> None: """get_object('osd_perf') returns an OsdPerfMonitoring instance.""" - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - ): + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "node1" - mock_settings.cluster.get.return_value = "dummy" - instance = MonitoringFactory.get_object("osd_perf", {}) + mock_base_settings.cluster.get.return_value = "dummy" + instance = MonitoringFactory.get_object("osd_perf", {"args": "record"}) assert isinstance(instance, OsdPerfMonitoring) diff --git a/tests/test_monitoring_monitoring.py b/tests/test_monitoring_monitoring.py new file mode 100644 index 00000000..45a7fe35 --- /dev/null +++ b/tests/test_monitoring_monitoring.py @@ -0,0 +1,120 @@ +"""Tests for the Monitoring base class helpers.""" + +# pylint: disable=protected-access + +from typing import ClassVar +from unittest.mock import MagicMock, patch + +import pytest + +from monitoring.monitoring import Monitoring + +# --------------------------------------------------------------------------- +# Concrete stub to instantiate the abstract base +# --------------------------------------------------------------------------- + + +class _StubMonitoring(Monitoring): + """Minimal concrete subclass used only for testing base-class helpers.""" + + DEFAULT_NODES: ClassVar[list[str]] = ["osds"] + + def start(self, directory: str) -> None: # pragma: no cover + pass + + def stop(self, directory): # pragma: no cover + pass + + +def _make_stub(nodes: str = "node1") -> _StubMonitoring: + with patch("monitoring.monitoring.settings") as mock_settings: + mock_settings.getnodes.return_value = nodes + mock_settings.cluster.get.return_value = "ceph" + return _StubMonitoring({}) + + +# --------------------------------------------------------------------------- +# _check_tool — fatal=True (default) +# --------------------------------------------------------------------------- + + +def test_check_tool_returns_true_when_tool_found() -> None: + """_check_tool returns True when command -v succeeds on all nodes.""" + ok_runner = MagicMock() + ok_runner.communicate.return_value = ("/usr/bin/perf\n", "") + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh", return_value=ok_runner): + assert stub._check_tool("perf") is True + + +def test_check_tool_raises_runtime_error_when_tool_missing() -> None: + """_check_tool raises RuntimeError when command -v fails on any node.""" + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="perf"): + stub._check_tool("perf") + + +def test_check_tool_error_message_includes_nodes() -> None: + """_check_tool RuntimeError message contains the node list.""" + stub = _make_stub(nodes="osd-node1,osd-node2") + with patch("monitoring.monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="osd-node1,osd-node2"): + stub._check_tool("blktrace") + + +def test_check_tool_uses_continue_if_error_false() -> None: + """_check_tool passes continue_if_error=False so pdsh raises on non-zero exit.""" + ok_runner = MagicMock() + ok_runner.communicate.return_value = ("/usr/bin/top\n", "") + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh", return_value=ok_runner) as mock_pdsh: + stub._check_tool("top") + mock_pdsh.assert_called_once_with("node1", "command -v top", continue_if_error=False) + + +# --------------------------------------------------------------------------- +# _check_tool — fatal=False +# --------------------------------------------------------------------------- + + +def test_check_tool_fatal_false_returns_true_when_tool_found() -> None: + """_check_tool(fatal=False) returns True when the tool is present.""" + ok_runner = MagicMock() + ok_runner.communicate.return_value = ("/usr/bin/collectl\n", "") + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh", return_value=ok_runner): + assert stub._check_tool("collectl", fatal=False) is True + + +def test_check_tool_fatal_false_returns_false_when_tool_missing() -> None: + """_check_tool(fatal=False) returns False (does not raise) when tool is absent.""" + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + assert stub._check_tool("collectl", fatal=False) is False + + +def test_check_tool_fatal_false_logs_warning_when_tool_missing() -> None: + """_check_tool(fatal=False) logs a warning when the tool is absent.""" + stub = _make_stub() + with ( + patch("monitoring.monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.monitoring.logger") as mock_logger, + ): + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + stub._check_tool("collectl", fatal=False) + + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("collectl" in c for c in warning_calls) + + +def test_check_tool_fatal_false_does_not_raise_when_tool_missing() -> None: + """_check_tool(fatal=False) must not raise even when the tool is absent.""" + stub = _make_stub() + with patch("monitoring.monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + # should complete without raising + stub._check_tool("collectl", fatal=False) diff --git a/tests/test_monitoring_perf.py b/tests/test_monitoring_perf.py index 4df43d60..5fe0f629 100644 --- a/tests/test_monitoring_perf.py +++ b/tests/test_monitoring_perf.py @@ -1,17 +1,20 @@ """Tests for the perf monitoring backend.""" -# pylint: disable=protected-access +# pylint: disable=protected-access,duplicate-code from typing import Any, Optional from unittest.mock import MagicMock, mock_open, patch -from monitoring.perf_monitoring import OsdPerfMonitoring, PerfMonitoring +import pytest + +from monitoring.osd_perf_monitoring import OsdPerfMonitoring +from monitoring.perf_monitoring import PerfMonitoring # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- -_ARGS = "stat -p {pid} -o {perf_dir}/perf_stat.{pid}" +_ARGS = "stat -p {pid} -o {output_dir}/perf_stat.{pid}" def _make_perf_monitor( @@ -21,12 +24,10 @@ def _make_perf_monitor( """Construct a PerfMonitoring instance with mocked settings.""" if mconfig is None: mconfig = {"args": _ARGS} - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - ): + assert "args" in mconfig, "_make_perf_monitor: mconfig must include 'args'" + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": user}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": user}.get(key, default) return PerfMonitoring(mconfig) @@ -40,13 +41,15 @@ def _make_osd_perf_monitor( mconfig = {"args": _ARGS} with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": pid_dir, + mock_base_settings.cluster.get.side_effect = lambda key, default=None: { "user": user, }.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: { + "pid_dir": pid_dir, + }.get(key, default) return OsdPerfMonitoring(mconfig) @@ -76,50 +79,107 @@ def test_perf_has_no_pid_dir_or_pid_glob() -> None: def test_perf_start_local_node_runs_perf() -> None: """PerfMonitoring.start() runs a single perf command locally.""" + # _check_tool and _make_remote_dir both resolve through the same common.pdsh + # reference. Use side_effect to return distinct runners per call. + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") mkdir_runner = MagicMock() local_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - patch("monitoring.perf_monitoring.common.pdsh", return_value=mkdir_runner) as mock_pdsh, + patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh, patch("monitoring.perf_monitoring.common.get_localnode", return_value="node1"), patch("monitoring.perf_monitoring.common.sh", return_value=local_runner) as mock_sh, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner] monitor = PerfMonitoring({"args": "stat -o {perf_dir}/perf_stat.out"}) monitor.start("/tmp/output") - mock_pdsh.assert_called_once_with("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") + mock_pdsh.assert_any_call("resolved-nodes", "command -v perf", continue_if_error=False) + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") mkdir_runner.communicate.assert_called_once_with() - mock_sh.assert_called_once_with("node1", "sudo perf stat -o /tmp/output/perf/perf_stat.out &") + mock_sh.assert_called_once_with("node1", "sudo perf stat -o /tmp/output/perf/perf_stat.out") assert monitor._perf_runners == [local_runner] assert monitor._perf_dir == "/tmp/output/perf" +def test_perf_start_logs_background_collection() -> None: + """PerfMonitoring.start() logs that collection runs in the background.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") + mkdir_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.perf_monitoring.common.pdsh", side_effect=[check_runner, mkdir_runner, MagicMock()]), + patch("monitoring.perf_monitoring.common.get_localnode", return_value=None), + patch("monitoring.perf_monitoring.logger") as mock_logger, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + PerfMonitoring({"args": "stat -o {perf_dir}/perf_stat.out"}).start("/tmp/output") + + mock_logger.info.assert_called_once_with("Perf monitoring running in background (will be killed on stop).") + + def test_perf_start_remote_node_uses_pdsh() -> None: - """PerfMonitoring.start() dispatches a single pdsh command for remote nodes.""" + """PerfMonitoring.start() dispatches a background pdsh command for remote nodes.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") mkdir_runner = MagicMock() remote_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh, patch("monitoring.perf_monitoring.common.get_localnode", return_value=None), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, remote_runner] - PerfMonitoring({"args": "stat -o {perf_dir}/perf_stat.out"}).start("/tmp/output") + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner, remote_runner] + monitor = PerfMonitoring({"args": "stat -o {perf_dir}/perf_stat.out"}) + monitor.start("/tmp/output") + mock_pdsh.assert_any_call("resolved-nodes", "command -v perf", continue_if_error=False) mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") - mock_pdsh.assert_any_call("resolved-nodes", "sudo perf stat -o /tmp/output/perf/perf_stat.out &") + mock_pdsh.assert_any_call("resolved-nodes", "sudo perf stat -o /tmp/output/perf/perf_stat.out") + assert monitor._perf_runners == [remote_runner] + remote_runner.communicate.assert_not_called() def test_perf_stop_kills_local_runners() -> None: - """PerfMonitoring.stop() kills locally started runners when present.""" + """PerfMonitoring.stop() pkills remote perf before killing tracked runners.""" + runner = MagicMock() + stop_runner = MagicMock() + monitor = _make_perf_monitor() + monitor._perf_runners = [runner] + + with patch("monitoring.perf_monitoring.common.pdsh", return_value=stop_runner) as mock_pdsh: + monitor.stop(None) + + mock_pdsh.assert_called_once_with("resolved-nodes", "sudo pkill -SIGINT -f 'perf '") + stop_runner.communicate.assert_called_once_with() + runner.kill.assert_called_once_with() + + +def test_perf_stop_logs_completion() -> None: + """PerfMonitoring.stop() logs that monitoring has stopped.""" + monitor = _make_perf_monitor() + + with ( + patch("monitoring.perf_monitoring.common.pdsh"), + patch("monitoring.perf_monitoring.logger") as mock_logger, + ): + monitor.stop(None) + + mock_logger.info.assert_called_once_with("Perf monitoring stopped.") + + +def test_perf_stop_ignores_already_exited_runner() -> None: + """PerfMonitoring.stop() ignores an OSError from an exited runner.""" runner = MagicMock() + runner.kill.side_effect = OSError monitor = _make_perf_monitor() monitor._perf_runners = [runner] @@ -137,32 +197,33 @@ def test_perf_stop_uses_pdsh_when_no_local_runners() -> None: with patch("monitoring.perf_monitoring.common.pdsh", return_value=stop_runner) as mock_pdsh: monitor.stop(None) - mock_pdsh.assert_called_once_with("resolved-nodes", r"sudo pkill -SIGINT -f perf\ ") + mock_pdsh.assert_called_once_with("resolved-nodes", "sudo pkill -SIGINT -f 'perf '") stop_runner.communicate.assert_called_once_with() def test_perf_stop_chowns_output_files_when_directory_provided() -> None: - """PerfMonitoring.stop() adjusts ownership of generated perf files when a directory is given.""" + """PerfMonitoring.stop() adjusts ownership of generated perf files and awaits each chown.""" stop_runner = MagicMock() chown_data_runner = MagicMock() chown_stat_runner = MagicMock() - monitor = _make_perf_monitor() + monitor = _make_perf_monitor(user="ceph") with patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh: mock_pdsh.side_effect = [stop_runner, chown_data_runner, chown_stat_runner] monitor.stop("/tmp/output") - mock_pdsh.assert_any_call("resolved-nodes", r"sudo pkill -SIGINT -f perf\ ") - mock_pdsh.assert_any_call("resolved-nodes", "sudo chown ceph.ceph /tmp/output/perf/perf.data") - mock_pdsh.assert_any_call("resolved-nodes", "sudo chown ceph.ceph /tmp/output/perf/perf_stat.*") + mock_pdsh.assert_any_call("resolved-nodes", "sudo pkill -SIGINT -f 'perf '") + mock_pdsh.assert_any_call("resolved-nodes", "sudo chown ceph:ceph /tmp/output/perf/perf.data") + chown_data_runner.communicate.assert_called_once_with() + mock_pdsh.assert_any_call("resolved-nodes", "sudo chown ceph:ceph /tmp/output/perf/perf_stat.*") + chown_stat_runner.communicate.assert_called_once_with() def test_perf_get_cpu_cycles_returns_total_cycles() -> None: """PerfMonitoring.get_cpu_cycles() sums cycle counts from all perf stat output files.""" with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - patch("monitoring.perf_monitoring.glob.glob", return_value=["/tmp/output/perf"]), + patch("monitoring.perf_monitoring.os.path.isdir", return_value=True), patch("monitoring.perf_monitoring.os.listdir", return_value=["perf_stat.1", "perf_stat.2"]), patch( "builtins.open", @@ -173,7 +234,7 @@ def test_perf_get_cpu_cycles_returns_total_cycles() -> None: ), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) monitor = PerfMonitoring({"args": _ARGS}) assert monitor.get_cpu_cycles("/tmp/output") == 3500 @@ -183,18 +244,54 @@ def test_perf_get_cpu_cycles_returns_none_when_cycles_missing() -> None: """PerfMonitoring.get_cpu_cycles() returns None when perf output has no cycles line.""" with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - patch("monitoring.perf_monitoring.glob.glob", return_value=["/tmp/output/perf"]), + patch("monitoring.perf_monitoring.os.path.isdir", return_value=True), patch("monitoring.perf_monitoring.os.listdir", return_value=["perf_stat.1"]), patch("builtins.open", mock_open(read_data="nothing to match\n")), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) monitor = PerfMonitoring({"args": _ARGS}) assert monitor.get_cpu_cycles("/tmp/output") is None +def test_perf_get_cpu_cycles_returns_none_when_no_perf_dir() -> None: + """PerfMonitoring.get_cpu_cycles() returns None when the perf directory does not exist.""" + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.perf_monitoring.os.path.isdir", return_value=False), + patch("monitoring.perf_monitoring.logger") as mock_logger, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + monitor = PerfMonitoring({"args": _ARGS}) + + result = monitor.get_cpu_cycles("/tmp/output") + + assert result is None + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("perf directory not found" in c for c in warning_calls) + + +def test_perf_get_cpu_cycles_returns_none_when_perf_dir_is_empty() -> None: + """PerfMonitoring.get_cpu_cycles() returns None (not 0) when perf dir exists but has no files.""" + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.perf_monitoring.os.path.isdir", return_value=True), + patch("monitoring.perf_monitoring.os.listdir", return_value=[]), + patch("monitoring.perf_monitoring.logger") as mock_logger, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + monitor = PerfMonitoring({"args": _ARGS}) + + result = monitor.get_cpu_cycles("/tmp/output") + + assert result is None + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("no perf stat files" in c for c in warning_calls) + + # --------------------------------------------------------------------------- # OsdPerfMonitoring # --------------------------------------------------------------------------- @@ -225,64 +322,206 @@ def test_osd_perf_init_accepts_custom_pid_glob() -> None: def test_osd_perf_start_local_node_runs_perf_per_pid() -> None: """OsdPerfMonitoring.start() runs perf locally for each matching pid file.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") mkdir_runner = MagicMock() local_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, - patch("monitoring.perf_monitoring.common.pdsh", return_value=mkdir_runner) as mock_pdsh, - patch("monitoring.perf_monitoring.common.get_localnode", return_value="node1"), - patch("monitoring.perf_monitoring.common.sh", return_value=local_runner) as mock_sh, - patch("monitoring.perf_monitoring.glob.glob", return_value=["/var/run/ceph/osd.1.pid"]), + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.osd_pid_monitoring.common.sh", return_value=local_runner) as mock_sh, + patch("monitoring.osd_pid_monitoring._glob.glob", return_value=["/var/run/ceph/osd.1.pid"]), patch("builtins.open", mock_open(read_data="123\n")), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner] monitor = OsdPerfMonitoring({"args": _ARGS}) monitor.start("/tmp/output") - mock_pdsh.assert_called_once_with("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") + mock_pdsh.assert_any_call("resolved-nodes", "command -v perf", continue_if_error=False) + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") mkdir_runner.communicate.assert_called_once_with() - mock_sh.assert_called_once_with("node1", "sudo perf stat -p 123 -o /tmp/output/perf/perf_stat.123 &") + mock_sh.assert_called_once_with("node1", "sudo perf stat -p 123 -o /tmp/output/perf/perf_stat.123") assert monitor._perf_runners == [local_runner] assert monitor._perf_dir == "/tmp/output/perf" def test_osd_perf_start_remote_node_uses_pdsh_loop() -> None: """OsdPerfMonitoring.start() dispatches via pdsh for-loop when no local node is available.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") mkdir_runner = MagicMock() + ls_runner = MagicMock() + ls_runner.communicate.return_value = ("/var/run/ceph/osd.1.pid\n", "") remote_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.perf_monitoring.settings") as mock_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value=None), + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner, ls_runner, remote_runner] + monitor = OsdPerfMonitoring({"args": _ARGS}) + monitor.start("/tmp/output") + + mock_pdsh.assert_any_call("resolved-nodes", "command -v perf", continue_if_error=False) + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") + mock_pdsh.assert_any_call("resolved-nodes", "ls /var/run/ceph/osd.*.pid 2>/dev/null") + expected_loop = ( + 'for f in /var/run/ceph/osd.*.pid; do pid=$(cat "$f");' + ' sudo perf stat -p "$pid" -o /tmp/output/perf/perf_stat."$pid" & done' + ) # {output_dir} resolved + mock_pdsh.assert_any_call("resolved-nodes", expected_loop) + assert monitor._perf_runners == [remote_runner] + remote_runner.communicate.assert_not_called() + + +def test_osd_perf_start_local_node_warns_when_no_pid_files() -> None: + """OsdPerfMonitoring.start() logs a warning when no local PID files are found.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") + mkdir_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh, - patch("monitoring.perf_monitoring.common.get_localnode", return_value=None), + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.osd_pid_monitoring._glob.glob", return_value=[]), + patch("monitoring.osd_pid_monitoring.logger") as mock_logger, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, remote_runner] + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner] OsdPerfMonitoring({"args": _ARGS}).start("/tmp/output") - mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/perf") - mock_pdsh.assert_any_call( - "resolved-nodes", - [ - "for pid in `cat /var/run/ceph/osd.*.pid`;", - "do", - "sudo perf stat -p ${pid} -o /tmp/output/perf/perf_stat.${pid} &", - ";", - "done", - ], - ) + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("no PID files" in c for c in warning_calls) + + +def test_osd_perf_start_remote_node_warns_when_no_pid_files() -> None: + """OsdPerfMonitoring.start() logs a warning when the remote ls finds no PID files.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") + mkdir_runner = MagicMock() + ls_runner = MagicMock() + ls_runner.communicate.return_value = ("", "") + remote_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value=None), + patch("monitoring.osd_pid_monitoring.logger") as mock_logger, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner, ls_runner, remote_runner] + OsdPerfMonitoring({"args": _ARGS}).start("/tmp/output") + + warning_calls = [str(c) for c in mock_logger.warning.call_args_list] + assert any("no PID files" in c for c in warning_calls) def test_osd_perf_stop_inherited_from_perf_monitoring() -> None: """OsdPerfMonitoring inherits stop() from PerfMonitoring without override.""" assert OsdPerfMonitoring.stop is PerfMonitoring.stop + + +def test_perf_init_raises_when_args_missing() -> None: + """PerfMonitoring.__init__ raises ValueError when 'args' is absent from mconfig.""" + with patch("monitoring.monitoring.settings") as mock_base_settings: + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + with pytest.raises(ValueError, match="args"): + PerfMonitoring({}) + + +def test_osd_perf_init_raises_when_args_missing() -> None: + """OsdPerfMonitoring.__init__ raises ValueError when 'args' is absent from mconfig.""" + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + with pytest.raises(ValueError, match="args"): + OsdPerfMonitoring({}) + + +# --------------------------------------------------------------------------- +# Tool availability checks +# --------------------------------------------------------------------------- + + +def test_perf_start_raises_when_perf_not_installed() -> None: + """PerfMonitoring.start() raises RuntimeError when perf is not on the nodes.""" + with patch("monitoring.monitoring.settings") as mock_base_settings: + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + monitor = PerfMonitoring({"args": "stat -o {perf_dir}/out"}) + + with patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="perf"): + monitor.start("/tmp/output") + + +def test_perf_start_proceeds_when_perf_is_installed() -> None: + """PerfMonitoring.start() continues normally when perf is found on the nodes.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/perf\n", "") + mkdir_runner = MagicMock() + local_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.perf_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.perf_monitoring.common.sh", return_value=local_runner), + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner] + monitor = PerfMonitoring({"args": "stat -o {perf_dir}/out"}) + monitor.start("/tmp/output") + + assert monitor._perf_dir == "/tmp/output/perf" + + +def test_osd_perf_start_raises_when_perf_not_installed() -> None: + """OsdPerfMonitoring.start() raises RuntimeError when perf is not on the nodes.""" + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + monitor = OsdPerfMonitoring({"args": _ARGS}) + + with patch("monitoring.perf_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="perf"): + monitor.start("/tmp/output") diff --git a/tests/test_monitoring_top.py b/tests/test_monitoring_top.py index f6e198b9..17bff2e7 100644 --- a/tests/test_monitoring_top.py +++ b/tests/test_monitoring_top.py @@ -1,11 +1,14 @@ """Tests for the top monitoring backend.""" -# pylint: disable=protected-access +# pylint: disable=protected-access,duplicate-code from typing import Any, Optional from unittest.mock import MagicMock, mock_open, patch -from monitoring.top_monitoring import OsdTopMonitoring, TopMonitoring +import pytest + +from monitoring.osd_top_monitoring import OsdTopMonitoring +from monitoring.top_monitoring import TopMonitoring, _estimate_top_duration # --------------------------------------------------------------------------- # Helpers @@ -19,12 +22,9 @@ def _make_top_monitor( """Construct a TopMonitoring instance with mocked settings.""" if mconfig is None: mconfig = {} - with ( - patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - ): + with patch("monitoring.monitoring.settings") as mock_base_settings: mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": user}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": user}.get(key, default) return TopMonitoring(mconfig) @@ -38,13 +38,15 @@ def _make_osd_top_monitor( mconfig = {} with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": pid_dir, + mock_base_settings.cluster.get.side_effect = lambda key, default=None: { "user": user, }.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: { + "pid_dir": pid_dir, + }.get(key, default) return OsdTopMonitoring(mconfig) @@ -53,6 +55,19 @@ def _make_osd_top_monitor( # --------------------------------------------------------------------------- +@pytest.mark.parametrize( + ("args", "expected"), + [ + ("-b -n 30", 90.0), + ("-b -n 30 -d 1.5", 45.0), + ("-b", None), + ], +) +def test_estimate_top_duration(args: str, expected: Optional[float]) -> None: + """_estimate_top_duration() uses an explicit delay or the three-second default.""" + assert _estimate_top_duration(args) == expected + + def test_top_default_nodes_is_osds() -> None: """TopMonitoring.DEFAULT_NODES must be ['osds'].""" assert TopMonitoring.DEFAULT_NODES == ["osds"] @@ -83,22 +98,25 @@ def test_top_has_no_pid_dir_or_pid_glob() -> None: def test_top_start_local_node_runs_top() -> None: """TopMonitoring.start() runs top locally without per-pid iteration.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() local_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - patch("monitoring.top_monitoring.common.pdsh", return_value=mkdir_runner) as mock_pdsh, + patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, patch("monitoring.top_monitoring.common.get_localnode", return_value="node1"), patch("monitoring.top_monitoring.common.sh", return_value=local_runner) as mock_sh, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner] monitor = TopMonitoring({}) monitor.start("/tmp/output") - mock_pdsh.assert_called_once_with("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") + mock_pdsh.assert_any_call("resolved-nodes", "command -v top", continue_if_error=False) + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") mkdir_runner.communicate.assert_called_once_with() expected_cmd = "top -b -H -1 -n 30 > /tmp/output/top/top.out" mock_sh.assert_called_once_with("node1", expected_cmd) @@ -106,49 +124,92 @@ def test_top_start_local_node_runs_top() -> None: def test_top_start_remote_node_uses_pdsh() -> None: - """TopMonitoring.start() dispatches a single pdsh command for remote nodes.""" + """TopMonitoring.start() dispatches a background pdsh command for remote nodes.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() remote_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, patch("monitoring.top_monitoring.common.get_localnode", return_value=None), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, remote_runner] + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner, remote_runner] monitor = TopMonitoring({}) monitor.start("/tmp/output") + mock_pdsh.assert_any_call("resolved-nodes", "command -v top", continue_if_error=False) mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") mock_pdsh.assert_any_call("resolved-nodes", "top -b -H -1 -n 30 > /tmp/output/top/top.out") + assert monitor._top_runners == [remote_runner] + remote_runner.communicate.assert_not_called() def test_top_stop_kills_local_runners() -> None: - """TopMonitoring.stop() kills locally started runners when present.""" + """TopMonitoring.stop() pkills remote top before killing tracked runners.""" runner = MagicMock() + stop_runner = MagicMock() monitor = _make_top_monitor() + monitor._running_cmd = "top -b -n 1" monitor._top_runners = [runner] - with patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh: + with patch("monitoring.top_monitoring.common.pdsh", return_value=stop_runner) as mock_pdsh: monitor.stop(None) + mock_pdsh.assert_called_once_with("resolved-nodes", "sudo pkill -SIGINT -f 'top -b -n 1'") + stop_runner.communicate.assert_called_once_with() runner.kill.assert_called_once_with() - mock_pdsh.assert_not_called() -def test_top_stop_uses_pdsh_when_no_local_runners() -> None: - """TopMonitoring.stop() issues pkill via pdsh when no local runners are tracked.""" - stop_runner = MagicMock() +def test_top_stop_logs_completion() -> None: + """TopMonitoring.stop() logs that monitoring has stopped.""" monitor = _make_top_monitor() - expected_pkill = f"sudo pkill -SIGINT -f '{monitor._top_cmd} {monitor._args}'" - with patch("monitoring.top_monitoring.common.pdsh", return_value=stop_runner) as mock_pdsh: + with ( + patch("monitoring.top_monitoring.common.pdsh"), + patch("monitoring.top_monitoring.logger") as mock_logger, + ): + monitor.stop(None) + + mock_logger.info.assert_called_once_with("Top monitoring stopped.") + + +def test_top_stop_ignores_already_exited_runner() -> None: + """TopMonitoring.stop() ignores an OSError from an exited runner.""" + runner = MagicMock() + runner.kill.side_effect = OSError + monitor = _make_top_monitor() + monitor._top_runners = [runner] + + with patch("monitoring.top_monitoring.common.pdsh"): + monitor.stop(None) + + runner.kill.assert_called_once_with() + + +def test_top_stop_uses_pdsh_when_no_local_runners() -> None: + """TopMonitoring.stop() pkill uses the fully-formatted command stored at start() time.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") + stop_runner = MagicMock() + mkdir_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.top_monitoring.common.get_localnode", return_value=None), + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner, MagicMock(), stop_runner] + monitor = TopMonitoring({}) + monitor.start("/tmp/output") monitor.stop(None) - mock_pdsh.assert_called_once_with("resolved-nodes", expected_pkill) + expected_pkill = "sudo pkill -SIGINT -f 'top -b -H -1 -n 30 > /tmp/output/top/top.out'" + mock_pdsh.assert_any_call("resolved-nodes", expected_pkill) stop_runner.communicate.assert_called_once_with() @@ -162,7 +223,10 @@ def test_top_stop_chowns_output_files_when_directory_provided() -> None: mock_pdsh.side_effect = [stop_runner, chown_runner] monitor.stop("/tmp/output") - mock_pdsh.assert_any_call("resolved-nodes", "sudo chown ceph.ceph /tmp/output/top/*top.out") + mock_pdsh.assert_any_call( + "resolved-nodes", + "sudo find /tmp/output/top -maxdepth 1 -name '*top.out' -exec chown ceph:ceph {} +", + ) # --------------------------------------------------------------------------- @@ -188,10 +252,10 @@ def test_osd_top_init_stores_pid_dir_and_glob() -> None: def test_osd_top_init_args_include_pid_placeholder() -> None: - """OsdTopMonitoring default args template must contain {pid}.""" + """OsdTopMonitoring default args template must contain {pid} and {output_dir}.""" monitor = _make_osd_top_monitor() assert "{pid}" in monitor._args - assert "{top_dir}" in monitor._args + assert "{output_dir}" in monitor._args def test_osd_top_init_accepts_custom_pid_glob() -> None: @@ -202,27 +266,31 @@ def test_osd_top_init_accepts_custom_pid_glob() -> None: def test_osd_top_start_local_node_runs_top_per_pid() -> None: """OsdTopMonitoring.start() runs top locally for each matching pid file.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() local_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - patch("monitoring.top_monitoring.common.pdsh", return_value=mkdir_runner) as mock_pdsh, - patch("monitoring.top_monitoring.common.get_localnode", return_value="node1"), - patch("monitoring.top_monitoring.common.sh", return_value=local_runner) as mock_sh, - patch("monitoring.top_monitoring.glob.glob", return_value=["/var/run/ceph/osd.1.pid"]), + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.osd_pid_monitoring.common.sh", return_value=local_runner) as mock_sh, + patch("monitoring.osd_pid_monitoring._glob.glob", return_value=["/var/run/ceph/osd.1.pid"]), patch("builtins.open", mock_open(read_data="42\n")), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner] monitor = OsdTopMonitoring({}) monitor.start("/tmp/output") - mock_pdsh.assert_called_once_with("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") + mock_pdsh.assert_any_call("resolved-nodes", "command -v top", continue_if_error=False) + mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") mkdir_runner.communicate.assert_called_once_with() expected_cmd = "top -b -H -1 -p 42 -n 30 > /tmp/output/top/42_osd_top.out" mock_sh.assert_called_once_with("node1", expected_cmd) @@ -231,20 +299,23 @@ def test_osd_top_start_local_node_runs_top_per_pid() -> None: def test_osd_top_start_local_node_warns_when_no_pid_files() -> None: """OsdTopMonitoring.start() logs a warning when no local PID files are found.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - patch("monitoring.top_monitoring.common.pdsh", return_value=mkdir_runner), - patch("monitoring.top_monitoring.common.get_localnode", return_value="node1"), - patch("monitoring.top_monitoring.glob.glob", return_value=[]), - patch("monitoring.top_monitoring.logger") as mock_logger, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.osd_pid_monitoring._glob.glob", return_value=[]), + patch("monitoring.osd_pid_monitoring.logger") as mock_logger, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner] OsdTopMonitoring({}).start("/tmp/output") warning_calls = [str(c) for c in mock_logger.warning.call_args_list] @@ -253,57 +324,60 @@ def test_osd_top_start_local_node_warns_when_no_pid_files() -> None: def test_osd_top_start_remote_node_uses_pdsh_loop() -> None: """OsdTopMonitoring.start() dispatches via pdsh for-loop when no local node is available.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() ls_runner = MagicMock() ls_runner.communicate.return_value = ("/var/run/ceph/osd.1.pid\n", "") remote_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, - patch("monitoring.top_monitoring.common.get_localnode", return_value=None), + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value=None), ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, ls_runner, remote_runner] - OsdTopMonitoring({}).start("/tmp/output") + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner, ls_runner, remote_runner] + monitor = OsdTopMonitoring({}) + monitor.start("/tmp/output") + mock_pdsh.assert_any_call("resolved-nodes", "command -v top", continue_if_error=False) mock_pdsh.assert_any_call("resolved-nodes", "mkdir -p -m0755 -- /tmp/output/top") mock_pdsh.assert_any_call("resolved-nodes", "ls /var/run/ceph/osd.*.pid 2>/dev/null") - mock_pdsh.assert_any_call( - "resolved-nodes", - [ - "for pid in `cat /var/run/ceph/osd.*.pid`;", - "do", - "top -b -H -1 -p ${pid} -n 30 > /tmp/output/top/${pid}_osd_top.out", - ";", - "done", - ], + expected_loop = ( + 'for f in /var/run/ceph/osd.*.pid; do pid=$(cat "$f");' + ' top -b -H -1 -p "$pid" -n 30 > /tmp/output/top/"$pid"_osd_top.out & done' ) + mock_pdsh.assert_any_call("resolved-nodes", expected_loop) + assert monitor._top_runners == [remote_runner] + remote_runner.communicate.assert_not_called() def test_osd_top_start_remote_node_warns_when_no_pid_files() -> None: """OsdTopMonitoring.start() logs a warning when the remote ls finds no PID files.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") mkdir_runner = MagicMock() ls_runner = MagicMock() ls_runner.communicate.return_value = ("", "") remote_runner = MagicMock() with ( patch("monitoring.monitoring.settings") as mock_base_settings, - patch("monitoring.top_monitoring.settings") as mock_settings, - patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, - patch("monitoring.top_monitoring.common.get_localnode", return_value=None), - patch("monitoring.top_monitoring.logger") as mock_logger, + patch("monitoring.osd_pid_monitoring.settings") as mock_osd_settings, + patch("common.pdsh") as mock_pdsh, + patch("monitoring.osd_pid_monitoring.common.get_localnode", return_value=None), + patch("monitoring.osd_pid_monitoring.logger") as mock_logger, ): mock_base_settings.getnodes.return_value = "resolved-nodes" - mock_settings.cluster.get.side_effect = lambda key, default=None: { - "pid_dir": "/var/run/ceph", - "user": "ceph", - }.get(key, default) - mock_pdsh.side_effect = [mkdir_runner, ls_runner, remote_runner] + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_osd_settings.cluster.get.side_effect = lambda key, default=None: {"pid_dir": "/var/run/ceph"}.get( + key, default + ) + mock_pdsh.side_effect = [check_runner, mkdir_runner, ls_runner, remote_runner] OsdTopMonitoring({}).start("/tmp/output") warning_calls = [str(c) for c in mock_logger.warning.call_args_list] @@ -313,3 +387,47 @@ def test_osd_top_start_remote_node_warns_when_no_pid_files() -> None: def test_osd_top_stop_inherited_from_top_monitoring() -> None: """OsdTopMonitoring inherits stop() from TopMonitoring without override.""" assert OsdTopMonitoring.stop is TopMonitoring.stop + + +# --------------------------------------------------------------------------- +# Tool availability checks +# --------------------------------------------------------------------------- + + +def test_top_start_raises_when_top_not_installed() -> None: + """TopMonitoring.start() raises RuntimeError when top is not on the nodes.""" + monitor = _make_top_monitor() + with patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="top"): + monitor.start("/tmp/output") + + +def test_top_start_proceeds_when_top_is_installed() -> None: + """TopMonitoring.start() continues normally when top is found on the nodes.""" + check_runner = MagicMock() + check_runner.communicate.return_value = ("/usr/bin/top\n", "") + mkdir_runner = MagicMock() + local_runner = MagicMock() + with ( + patch("monitoring.monitoring.settings") as mock_base_settings, + patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh, + patch("monitoring.top_monitoring.common.get_localnode", return_value="node1"), + patch("monitoring.top_monitoring.common.sh", return_value=local_runner), + ): + mock_base_settings.getnodes.return_value = "resolved-nodes" + mock_base_settings.cluster.get.side_effect = lambda key, default=None: {"user": "ceph"}.get(key, default) + mock_pdsh.side_effect = [check_runner, mkdir_runner] + monitor = TopMonitoring({}) + monitor.start("/tmp/output") + + assert monitor._top_runners == [local_runner] + + +def test_osd_top_start_raises_when_top_not_installed() -> None: + """OsdTopMonitoring.start() raises RuntimeError when top is not on the nodes.""" + monitor = _make_osd_top_monitor() + with patch("monitoring.top_monitoring.common.pdsh") as mock_pdsh: + mock_pdsh.return_value.communicate.side_effect = Exception("exit 1") + with pytest.raises(RuntimeError, match="top"): + monitor.start("/tmp/output")