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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 11 additions & 6 deletions monitoring/blktrace_monitoring.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand All @@ -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()
15 changes: 9 additions & 6 deletions monitoring/collectl_monitoring.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
45 changes: 43 additions & 2 deletions monitoring/monitoring.py
Original file line number Diff line number Diff line change
@@ -1,21 +1,62 @@
"""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."""

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 <tool>`` 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:
Expand Down
6 changes: 4 additions & 2 deletions monitoring/monitoring_factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")

Expand Down
22 changes: 22 additions & 0 deletions monitoring/osd_perf_monitoring.py
Original file line number Diff line number Diff line change
@@ -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)
73 changes: 73 additions & 0 deletions monitoring/osd_pid_monitoring.py
Original file line number Diff line number Diff line change
@@ -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]
30 changes: 30 additions & 0 deletions monitoring/osd_top_monitoring.py
Original file line number Diff line number Diff line change
@@ -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)
Loading