diff --git a/doc/source/lib/galaxy.job_metrics.collectl.rst b/doc/source/lib/galaxy.job_metrics.collectl.rst deleted file mode 100644 index d302c9f9f5b..00000000000 --- a/doc/source/lib/galaxy.job_metrics.collectl.rst +++ /dev/null @@ -1,42 +0,0 @@ -galaxy.job\_metrics.collectl package -==================================== - -.. automodule:: galaxy.job_metrics.collectl - :members: - :undoc-members: - :show-inheritance: - -Submodules ----------- - -galaxy.job\_metrics.collectl.cli module ---------------------------------------- - -.. automodule:: galaxy.job_metrics.collectl.cli - :members: - :undoc-members: - :show-inheritance: - -galaxy.job\_metrics.collectl.processes module ---------------------------------------------- - -.. automodule:: galaxy.job_metrics.collectl.processes - :members: - :undoc-members: - :show-inheritance: - -galaxy.job\_metrics.collectl.stats module ------------------------------------------ - -.. automodule:: galaxy.job_metrics.collectl.stats - :members: - :undoc-members: - :show-inheritance: - -galaxy.job\_metrics.collectl.subsystems module ----------------------------------------------- - -.. automodule:: galaxy.job_metrics.collectl.subsystems - :members: - :undoc-members: - :show-inheritance: diff --git a/doc/source/lib/galaxy.job_metrics.rst b/doc/source/lib/galaxy.job_metrics.rst index c51c202b505..f8c01bd1cb5 100644 --- a/doc/source/lib/galaxy.job_metrics.rst +++ b/doc/source/lib/galaxy.job_metrics.rst @@ -12,7 +12,6 @@ Subpackages .. toctree:: :maxdepth: 4 - galaxy.job_metrics.collectl galaxy.job_metrics.instrumenters Submodules diff --git a/lib/galaxy/job_metrics/__init__.py b/lib/galaxy/job_metrics/__init__.py index 50a79af9745..0a9a39f0b2b 100644 --- a/lib/galaxy/job_metrics/__init__.py +++ b/lib/galaxy/job_metrics/__init__.py @@ -13,9 +13,24 @@ collect the output of these from a job directory. import collections import logging import os +from abc import ( + ABCMeta, + abstractmethod, +) +from typing import ( + Any, + Dict, + List, + NamedTuple, + Optional, +) from galaxy import util from galaxy.util import plugin_config +from .safety import ( + DEFAULT_SAFETY, + Safety, +) from ..job_metrics import formatting log = logging.getLogger(__name__) @@ -24,6 +39,32 @@ log = logging.getLogger(__name__) DEFAULT_FORMATTER = formatting.JobMetricFormatter() +class DictifiableMetric(NamedTuple): + """The full context of a metric that is ready to be exposed to an external client.""" + + title: str + value: str + raw_value: str + name: str + plugin: str + safety: Safety = Safety.POTENTIALLY_SENSITVE + + def dict(self) -> Dict[str, str]: + return dict( + title=self.title, + value=self.value, + plugin=self.plugin, + name=self.name, + raw_value=self.raw_value, + ) + + +class RawMetric(NamedTuple): + metric_name: str + metric_value: Any + metric_plugin: str + + class JobMetrics: """Load and store a collection of :class:`JobInstrumenter` objects.""" @@ -33,7 +74,7 @@ class JobMetrics: self.default_job_instrumenter = JobInstrumenter.from_file(self.plugin_classes, conf_file, **kwargs) self.job_instrumenters = collections.defaultdict(lambda: self.default_job_instrumenter) - def format(self, plugin, key, value): + def format(self, plugin: str, key: str, value: Any) -> formatting.FormattedMetric: """Find :class:`formatting.JobMetricFormatter` corresponding to instrumented plugin value.""" if plugin in self.plugin_classes: plugin_class = self.plugin_classes[plugin] @@ -42,6 +83,30 @@ class JobMetrics: formatter = DEFAULT_FORMATTER return formatter.format(key, value) + def dictifiable_metrics(self, raw_metrics: List[RawMetric], allowed_safety: Safety) -> List[DictifiableMetric]: + def raw_to_dictifiable(raw_metric: RawMetric) -> DictifiableMetric: + metric_name, metric_value, metric_plugin = raw_metric + title, value = self.format(metric_plugin, metric_name, metric_value) + configured_plugin = self.default_job_instrumenter.get_configured_plugin(metric_plugin) + if configured_plugin is not None: + safety = configured_plugin.safety(metric_name) + elif metric_plugin in self.plugin_classes: + plugin_class = self.plugin_classes[metric_plugin] + safety = plugin_class.default_safety + else: + safety = DEFAULT_SAFETY + return DictifiableMetric( + title, + value, + str(metric_value), + metric_name, + metric_plugin, + safety, + ) + + metrics = map(raw_to_dictifiable, raw_metrics) + return [m for m in metrics if m.safety.value >= allowed_safety.value] + def set_destination_conf_file(self, destination_id, conf_file): instrumenter = JobInstrumenter.from_file(self.plugin_classes, conf_file) self.set_destination_instrumenter(destination_id, instrumenter) @@ -70,7 +135,25 @@ class JobMetrics: return plugin_config.plugins_dict(galaxy.job_metrics.instrumenters, "plugin_type") -class NullJobInstrumenter: +class JobInstrumenterI(metaclass=ABCMeta): + @abstractmethod + def pre_execute_commands(self, job_directory: str) -> Optional[str]: + return None + + @abstractmethod + def post_execute_commands(self, job_directory: str) -> Optional[str]: + return None + + @abstractmethod + def collect_properties(self, job_id, job_directory: str) -> Dict[str, Any]: + return {} + + @abstractmethod + def get_configured_plugin(self, plugin_type: str): + return None + + +class NullJobInstrumenter(JobInstrumenterI): def pre_execute_commands(self, job_directory): return None @@ -80,16 +163,25 @@ class NullJobInstrumenter: def collect_properties(self, job_id, job_directory): return {} + def get_configured_plugin(self, plugin_type: str): + return None + NULL_JOB_INSTRUMENTER = NullJobInstrumenter() -class JobInstrumenter: +class JobInstrumenter(JobInstrumenterI): def __init__(self, plugin_classes, plugins_source, **kwargs): self.extra_kwargs = kwargs self.plugin_classes = plugin_classes self.plugins = self.__plugins_from_source(plugins_source) + def get_configured_plugin(self, plugin_type: str): + for plugin in self.plugins: + if plugin.plugin_type == plugin_type: + return plugin + return None + def pre_execute_commands(self, job_directory): commands = [] for plugin in self.plugins: @@ -127,8 +219,14 @@ class JobInstrumenter: return plugin_config.load_plugins(self.plugin_classes, plugins_source, self.extra_kwargs) @staticmethod - def from_file(plugin_classes, conf_file, **kwargs): + def from_file(plugin_classes, conf_file, **kwargs) -> "JobInstrumenterI": if not conf_file or not os.path.exists(conf_file): return NULL_JOB_INSTRUMENTER plugins_source = plugin_config.plugin_source_from_path(conf_file) return JobInstrumenter(plugin_classes, plugins_source, **kwargs) + + +__all__ = ( + "JobInstrumenter", + "Safety", +) diff --git a/lib/galaxy/job_metrics/collectl/__init__.py b/lib/galaxy/job_metrics/collectl/__init__.py deleted file mode 100644 index c3e8815049d..00000000000 --- a/lib/galaxy/job_metrics/collectl/__init__.py +++ /dev/null @@ -1,4 +0,0 @@ -"""Helper functions and data structures for interacting with collectl & data. - -More information on collectl can be found at: http://collectl.sourceforge.net/. -""" diff --git a/lib/galaxy/job_metrics/collectl/cli.py b/lib/galaxy/job_metrics/collectl/cli.py deleted file mode 100644 index 1c370fe1858..00000000000 --- a/lib/galaxy/job_metrics/collectl/cli.py +++ /dev/null @@ -1,141 +0,0 @@ -"""This module describes :class:`CollectlCli` - an abstraction for building collectl command lines.""" -import logging -import subprocess -from string import Template - -log = logging.getLogger(__name__) - -COMMAND_LINE_TEMPLATE = Template( - "$collectl_path $destination_arg $mode_arg $subsystems_arg $interval_arg $procfilt_arg $flush_arg $sep_arg" -) -MODE_RECORD = "record" -MODE_PLAYBACK = "playback" - - -class CollectlCli: - """ - Abstraction over (some of) the command-line arguments of collectl. - Ideally this will be useful for building up command line arguments for - remote execution as well as runnning directly on local host. - - This is meant to be a fairly generic utility - for interfacing with - collectl CLI - logic more directly related to the Galaxy job metric plugin - plugin should be placed in other modules. - - **Keyword Arguments:** - - ``collectl_path`` - Path to collectl executable (defaults to collectl - i.e. - search the PATH). - - ``playback_path`` (defaults to ``None``) - If this is ``None``, collectl will run in - record mode, else it will playback specified file. - - **Playback Mode Options:** - - ``sep`` - Separator used in playback mode (set to 9 to produce tsv) - (defaults to None). - - **Record Mode Options** (some of these may work in playback mode also) - - ``destination_path`` - Location of path files to write to (defaults to None - and collectl will just use cwd). Really this is just to prefix - - collectl will append hostname and datetime to file. - ``interval`` - Setup polling interval (secs) for most subsystems (defaults - to None and when unspecified collectl will use default of 1 second). - ``interval2`` - Setup polling interval (secs) for process information - (defaults to None and when unspecified collectl will use default to - 60 seconds). - ``interval3`` - Setup polling interval (secs) for environment information - (defaults to None and when unspecified collectl will use default to - 300 seconds). - ``procfilt`` - Optional argument to procfilt. (defaults to None). - ``flush`` - Optional flush interval (defaults to None). - """ - - def __init__(self, **kwargs): - command_args = {} - command_args["collectl_path"] = kwargs.get("collectl_path", "collectl") - playback_path = kwargs.get("playback_path", None) - self.mode = MODE_RECORD if not playback_path else MODE_PLAYBACK - if self.mode == MODE_RECORD: - mode_arg = "" - elif self.mode == MODE_PLAYBACK: - mode_arg = f"-P -p '{playback_path}'" - else: - raise Exception(f"Invalid mode supplied to CollectlCli - {self.mode}") - command_args["mode_arg"] = mode_arg - command_args["interval_arg"] = self.__interval_arg(kwargs) - destination = kwargs.get("destination_path", None) - if destination: - destination_arg = f"-f '{destination}'" - else: - destination_arg = "" - command_args["destination_arg"] = destination_arg - procfilt = kwargs.get("procfilt", None) - command_args["procfilt_arg"] = "" if not procfilt else f"--procfilt {procfilt}" - command_args["subsystems_arg"] = self.__subsystems_arg(kwargs.get("subsystems", [])) - flush = kwargs.get("flush", None) - command_args["flush_arg"] = f"--flush {flush}" if flush else "" - sep = kwargs.get("sep", None) - command_args["sep_arg"] = f"--sep={sep}" if sep else "" - - self.command_args = command_args - - def __subsystems_arg(self, subsystems): - if subsystems: - return f"-s{''.join(s.command_line_arg for s in subsystems)}" - else: - return "" - - def __interval_arg(self, kwargs): - if self.mode != MODE_RECORD: - return "" - - interval = kwargs.get("interval", None) - if not interval: - return "" - - self.__validate_interval_arg(interval) - interval_arg = f"-i {interval}" - interval2 = kwargs.get("interval2", None) - if not interval2: - return interval_arg - self.__validate_interval_arg(interval2, multiple_of=int(interval)) - interval_arg = f"{interval_arg}:{interval2}" - - interval3 = kwargs.get("interval3", None) - if not interval3: - return interval_arg - self.__validate_interval_arg(interval3, multiple_of=int(interval)) - interval_arg = f"{interval_arg}:{interval3}" - return interval_arg - - def __validate_interval_arg(self, value, multiple_of=None): - if value and not str(value).isdigit(): - raise Exception(f"Invalid interval argument supplied, must be integer {value}") - if multiple_of: - if int(value) % multiple_of != 0: - raise Exception(f"Invalid interval argument supplied, must multiple of {multiple_of}") - - def build_command_line(self): - return COMMAND_LINE_TEMPLATE.substitute(**self.command_args) - - def run(self, stdout=subprocess.PIPE, stderr=subprocess.PIPE): - command_line = self.build_command_line() - log.info(f"Executing {command_line}") - proc = subprocess.Popen(command_line, shell=True, stdout=stdout, stderr=stderr) - return_code = proc.wait() - if return_code: - raise Exception("Problem running collectl command.") - - -__all__ = ("CollectlCli",) diff --git a/lib/galaxy/job_metrics/collectl/processes.py b/lib/galaxy/job_metrics/collectl/processes.py deleted file mode 100644 index a665e76b82a..00000000000 --- a/lib/galaxy/job_metrics/collectl/processes.py +++ /dev/null @@ -1,249 +0,0 @@ -""" Modules will run collectl in playback mode and collect various process -statistics for a given pid's process and process ancestors. -""" -import collections -import csv -import logging -import tempfile - -from galaxy import util -from ..collectl import stats - -log = logging.getLogger(__name__) - -# Collectl process information cheat sheet: -# -# Record process information for current user. -# % collectl -sZ -f./__instrument_collectl -i 10:10 --procfilt U$USER -# -# TSV Replay of processing information in plottable mode... -# -# % collectl -sZ -P --sep=9 -p __instrument_collectl-jlaptop13-20140322-120919.raw.gz -# -# Has following columns: -# Date Time PID User PR PPID THRD S VmSize VmLck VmRSS VmData VmStk VmExe VmLib CPU SysT UsrT PCT AccumT RKB WKB RKBC WKBC RSYS WSYS CNCL MajF MinF Command -# - -# Process data dumped one row per process per interval. -# http://collectl.sourceforge.net/Data-detail.html -PROCESS_COLUMNS = [ - "#Date", # Date of interval - e.g. 20140322 - "Time", # Time of interval - 12:18:58 - "PID", # Process pid. - "User", # Process user. - "PR", # Priority of process. - "PPID", # Parent PID of process. - "THRD", # Thread??? - "S", # Process state - S - Sleeping, D - Uninterruptable Sleep, R - Running, Z - Zombie or T - Stopped/Traced - # Memory options - http://ewx.livejournal.com/579283.html - "VmSize", - "VmLck", - "VmRSS", - "VmData", - "VmStk", - "VmExe", - "VmLib", - "CPU", # CPU number of process - "SysT", # Amount of system time consumed during interval - "UsrT", # Amount user time consumed during interval - "PCT", # Percentage of current interval consumed by task - "AccumT", # Total accumulated System and User time since the process began execution - # kilobytes read/written - requires I/O level monitoring to be enabled in kernel. - "RKB", # kilobytes read by process - requires I/O monitoring in kernel - "WKB", - "RKBC", - "WKBC", - "RSYS", # Number of read system calls - "WSYS", # Number of write system calls - "CNCL", - "MajF", # Number of major page faults - "MinF", # Number of minor page faults - "Command", # Command executed -] - -# Types of statistics this module can summarize -STATISTIC_TYPES = ["max", "min", "sum", "count", "avg"] - -COLUMN_INDICES = {col: i for i, col in enumerate(PROCESS_COLUMNS)} -PID_INDEX = COLUMN_INDICES["PID"] -PARENT_PID_INDEX = COLUMN_INDICES["PPID"] - -DEFAULT_STATISTICS = [ - ("max", "VmSize"), - ("avg", "VmSize"), - ("max", "VmRSS"), - ("avg", "VmRSS"), - ("sum", "SysT"), - ("sum", "UsrT"), - ("max", "PCT"), - ("avg", "PCT"), - ("max", "AccumT"), - ("sum", "RSYS"), - ("sum", "WSYS"), -] - - -def parse_process_statistics(statistics): - """Turn string or list of strings into list of tuples in format ( stat, - resource ) where stat is a value from STATISTIC_TYPES and resource is a - value from PROCESS_COLUMNS. - """ - if statistics is None: - statistics = DEFAULT_STATISTICS - - statistics = util.listify(statistics) - statistics = [_tuplize_statistic(_) for _ in statistics] - # Check for validity... - for statistic in statistics: - if statistic[0] not in STATISTIC_TYPES: - raise Exception(f"Unknown statistic type encountered {statistic[0]}") - if statistic[1] not in PROCESS_COLUMNS: - raise Exception(f"Unknown process column encountered {statistic[1]}") - return statistics - - -def generate_process_statistics(collectl_playback_cli, pid, statistics=DEFAULT_STATISTICS): - """Playback collectl file and generate summary statistics.""" - with tempfile.NamedTemporaryFile() as tmp_tsv: - collectl_playback_cli.run(stdout=tmp_tsv) - with open(tmp_tsv.name) as tsv_file: - return _read_process_statistics(tsv_file, pid, statistics) - - -def _read_process_statistics(tsv_file, pid, statistics): - process_summarizer = CollectlProcessSummarizer(pid, statistics) - current_interval = None - - for row in csv.reader(tsv_file, dialect="excel-tab"): - if current_interval is None: - for header, expected_header in zip(row, PROCESS_COLUMNS): - if header.lower() != expected_header.lower(): - raise Exception(f"Unknown header value encountered while processing collectl playback - {header}") - - # First row, check contains correct header. - current_interval = CollectlProcessInterval() - continue - - if current_interval.row_is_in(row): - current_interval.add_row(row) - else: - process_summarizer.handle_interval(current_interval) - current_interval = CollectlProcessInterval() - - # Do we have unsummarized rows... - if current_interval and current_interval.rows: - process_summarizer.handle_interval(current_interval) - - return process_summarizer.get_statistics() - - -class CollectlProcessSummarizer: - def __init__(self, pid, statistics): - self.pid = pid - self.statistics = statistics - self.columns_of_interest = {s[1] for s in statistics} - self.tree_statistics = collections.defaultdict(stats.StatisticsTracker) - self.process_accum_statistics = collections.defaultdict(stats.StatisticsTracker) - self.interval_count = 0 - - def handle_interval(self, interval): - self.interval_count += 1 - rows = self.__rows_for_process(interval.rows, self.pid) - for column_name in self.columns_of_interest: - column_index = COLUMN_INDICES[column_name] - - if column_name == "AccumT": - # Should not sum this across pids each interval, sum max at end... - for r in rows: - pid_seconds = self.__time_to_seconds(r[column_index]) - self.process_accum_statistics[r[PID_INDEX]].track(pid_seconds) - else: - # All other stastics should be summed across whole process tree - # at each interval I guess. - if column_name in ["SysT", "UsrT", "PCT"]: - to_num = float - else: - to_num = int - - interval_stat = sum(to_num(r[column_index]) for r in rows) - self.tree_statistics[column_name].track(interval_stat) - - def get_statistics(self): - if self.interval_count == 0: - return [] - - computed_statistics = [] - for statistic in self.statistics: - statistic_type, column = statistic - if column == "AccumT": - # Only thing that makes sense is sum - if statistic_type != "max": - log.warning("Only statistic max makes sense for AccumT") - continue - - value = sum(v.max for v in self.process_accum_statistics.values()) - else: - statistics_tracker = self.tree_statistics[column] - value = getattr(statistics_tracker, statistic_type) - - computed_statistic = (statistic, value) - computed_statistics.append(computed_statistic) - - return computed_statistics - - def __rows_for_process(self, rows, pid): - process_rows = [] - pids = self.__all_child_pids(rows, pid) - for row in rows: - if row[PID_INDEX] in pids: - process_rows.append(row) - return process_rows - - def __all_child_pids(self, rows, pid): - pids_in_process_tree = {str(self.pid)} - added = True - while added: - added = False - for row in rows: - pid = row[PID_INDEX] - parent_pid = row[PARENT_PID_INDEX] - if parent_pid in pids_in_process_tree and pid not in pids_in_process_tree: - pids_in_process_tree.add(pid) - added = True - return pids_in_process_tree - - def __time_to_seconds(self, minutes_str): - parts = minutes_str.split(":") - seconds = 0.0 - for i, val in enumerate(parts): - seconds += float(val) * (60 ** (len(parts) - (i + 1))) - return seconds - - -class CollectlProcessInterval: - """Represent all rows in collectl playback file for given time slice with - ability to filter out just rows corresponding to the process tree - corresponding to a given pid. - """ - - def __init__(self): - self.rows = [] - - def row_is_in(self, row): - if not self.rows: # No rows, this row defines interval. - return True - first_row = self.rows[0] - return first_row[0] == row[0] and first_row[1] == row[1] - - def add_row(self, row): - self.rows.append(row) - - -def _tuplize_statistic(statistic): - if not isinstance(statistic, tuple): - statistic_split = statistic.split("_", 1) - statistic = (statistic_split[0].lower(), statistic_split[1]) - return statistic - - -__all__ = ("generate_process_statistics",) diff --git a/lib/galaxy/job_metrics/collectl/stats.py b/lib/galaxy/job_metrics/collectl/stats.py deleted file mode 100644 index a5685dd4288..00000000000 --- a/lib/galaxy/job_metrics/collectl/stats.py +++ /dev/null @@ -1,26 +0,0 @@ -""" Primitive module for tracking running statistics without storing values in -memory. -""" - - -class StatisticsTracker: - def __init__(self): - self.min = None - self.max = None - self.count = 0 - self.sum = 0 - - def track(self, value): - if self.min is None or value < self.min: - self.min = value - if self.max is None or value > self.max: - self.max = value - self.count += 1 - self.sum += value - - @property - def avg(self): - if self.count > 0: - return self.sum / self.count - else: - return None diff --git a/lib/galaxy/job_metrics/collectl/subsystems.py b/lib/galaxy/job_metrics/collectl/subsystems.py deleted file mode 100644 index d4816593c20..00000000000 --- a/lib/galaxy/job_metrics/collectl/subsystems.py +++ /dev/null @@ -1,75 +0,0 @@ -"""Abstractions describing collectl subsystems (specified with the collectl ``-s`` parameter). - -Subsystems are essentially monitoring plugins available within collectl. -""" -from abc import ( - ABCMeta, - abstractmethod, -) - - -class CollectlSubsystem(metaclass=ABCMeta): - """Class providing an abstraction of collectl subsytems.""" - - @property - @abstractmethod - def command_line_arg(self): - """Return single letter command-line argument used by collectl CLI.""" - - @property - @abstractmethod - def name(self): - """High-level name for subsystem as consumed by this module.""" - - -class ProcessesSubsystem(CollectlSubsystem): - command_line_arg = "Z" - name = "process" - - -class CpuSubsystem(CollectlSubsystem): - command_line_arg = "C" - name = "cpu" - - -class DiskSubsystem(CollectlSubsystem): - command_line_arg = "D" - name = "disk" - - -class NetworkSubsystem(CollectlSubsystem): - command_line_arg = "N" - name = "network" - - -class EnvironmentSubsystem(CollectlSubsystem): - command_line_arg = "E" - name = "environment" - - -class MemorySubsystem(CollectlSubsystem): - command_line_arg = "M" - name = "memory" - - -SUBSYSTEMS = [ - ProcessesSubsystem(), - CpuSubsystem(), - DiskSubsystem(), - NetworkSubsystem(), - EnvironmentSubsystem(), - MemorySubsystem(), -] -SUBSYSTEM_DICT = {s.name: s for s in SUBSYSTEMS} - - -def get_subsystem(name): - """ - - >>> get_subsystem( "process" ).command_line_arg == "Z" - True - """ - return SUBSYSTEM_DICT[name] - - -__all__ = ("get_subsystem",) diff --git a/lib/galaxy/job_metrics/formatting.py b/lib/galaxy/job_metrics/formatting.py index 6b5c4b81106..1b89548ffa1 100644 --- a/lib/galaxy/job_metrics/formatting.py +++ b/lib/galaxy/job_metrics/formatting.py @@ -1,14 +1,23 @@ """Utilities related to formatting job metrics for human consumption.""" +from typing import ( + Any, + NamedTuple, +) + + +class FormattedMetric(NamedTuple): + title: str + value: str class JobMetricFormatter: """Format job metric key-value pairs for human consumption in Web UI.""" - def format(self, key, value): - return (str(key), str(value)) + def format(self, key: Any, value: Any) -> FormattedMetric: + return FormattedMetric(str(key), str(value)) -def seconds_to_str(value): +def seconds_to_str(value: int) -> str: """Convert seconds to a simple simple string describing the amount of time.""" mins, secs = divmod(value, 60) hours, mins = divmod(mins, 60) diff --git a/lib/galaxy/job_metrics/instrumenters/__init__.py b/lib/galaxy/job_metrics/instrumenters/__init__.py index cc15408da54..cac587d60b4 100644 --- a/lib/galaxy/job_metrics/instrumenters/__init__.py +++ b/lib/galaxy/job_metrics/instrumenters/__init__.py @@ -7,30 +7,43 @@ from abc import ( ABCMeta, abstractmethod, ) +from typing import ( + Any, + Dict, + List, + Optional, + Union, +) from .. import formatting +from ..safety import ( + DEFAULT_SAFETY, + Safety, +) INSTRUMENT_FILE_PREFIX = "__instrument" +InstrumentableT = Optional[Union[str, List[str]]] class InstrumentPlugin(metaclass=ABCMeta): """Describes how to instrument job scripts and retrieve collected metrics.""" formatter = formatting.JobMetricFormatter() + default_safety = DEFAULT_SAFETY @property @abstractmethod def plugin_type(self): """Short string providing labelling this plugin""" - def pre_execute_instrument(self, job_directory): + def pre_execute_instrument(self, job_directory: str) -> InstrumentableT: """Optionally return one or more commands to instrument job. These commands will be executed on the compute server prior to the job running. """ return None - def post_execute_instrument(self, job_directory): + def post_execute_instrument(self, job_directory: str) -> InstrumentableT: """Optionally return one or more commands to instrument job. These commands will be executed on the compute server after the tool defined command is ran. @@ -38,18 +51,26 @@ class InstrumentPlugin(metaclass=ABCMeta): return None @abstractmethod - def job_properties(self, job_id, job_directory): + def job_properties(self, job_id, job_directory: str) -> Dict[str, Any]: """Collect properties for this plugin from specified job directory. This method will run on the Galaxy server and can assume files created in job_directory with pre_execute_instrument and post_execute_instrument are available. """ - def _instrument_file_name(self, name): + def safety(self, metric_name: str) -> Safety: + """Return safety level of metric.""" + # None of the plugins override this to dispatch on metric_name but on next + # iteration it would make sense to allow admins to expose particular env vars + # or to have cgroup keys we know are about runtime or memeory to be exposed + # at a safer level. + return self.default_safety + + def _instrument_file_name(self, name: str) -> str: """Provide a common pattern for naming files used by instrumentation plugins - to ease their staging out of remote job directories. """ return f"{INSTRUMENT_FILE_PREFIX}_{self.plugin_type}_{name}" - def _instrument_file_path(self, job_directory, name): + def _instrument_file_path(self, job_directory: str, name: str) -> str: return os.path.join(job_directory, self._instrument_file_name(name)) diff --git a/lib/galaxy/job_metrics/instrumenters/cgroup.py b/lib/galaxy/job_metrics/instrumenters/cgroup.py index 48f07c1499a..d4739019032 100644 --- a/lib/galaxy/job_metrics/instrumenters/cgroup.py +++ b/lib/galaxy/job_metrics/instrumenters/cgroup.py @@ -2,6 +2,11 @@ import logging import numbers from collections import namedtuple +from typing import ( + Any, + Dict, + List, +) from galaxy.util import ( asbool, @@ -71,7 +76,7 @@ class CgroupPluginFormatter(formatting.JobMetricFormatter): return title, nice_size(value) except ValueError: pass - elif isinstance(value, numbers.Number) and value == int(value): + elif isinstance(value, (numbers.Integral, numbers.Real)) and value == int(value): value = int(value) return title, value @@ -89,26 +94,26 @@ class CgroupPlugin(InstrumentPlugin): if params_str: params = [v.strip() for v in params_str.split(",")] else: - params = TITLES.keys() + params = list(TITLES.keys()) self.params = params - def post_execute_instrument(self, job_directory): - commands = [] + def post_execute_instrument(self, job_directory: str) -> List[str]: + commands: List[str] = [] commands.append(self.__record_cgroup_cpu_usage(job_directory)) commands.append(self.__record_cgroup_memory_usage(job_directory)) return commands - def job_properties(self, job_id, job_directory): + def job_properties(self, job_id, job_directory: str) -> Dict[str, Any]: metrics = self.__read_metrics(self.__cgroup_metrics_file(job_directory)) return metrics - def __record_cgroup_cpu_usage(self, job_directory): + def __record_cgroup_cpu_usage(self, job_directory: str) -> str: # comounted cgroups (which cpu and cpuacct are on the supported Linux distros) can appear in any order (cpu,cpuacct or cpuacct,cpu) return CPU_USAGE_TEMPLATE.format( metrics=self.__cgroup_metrics_file(job_directory), cgroup_mount=self.cgroup_mount ) - def __record_cgroup_memory_usage(self, job_directory): + def __record_cgroup_memory_usage(self, job_directory: str) -> str: return MEMORY_USAGE_TEMPLATE.format( metrics=self.__cgroup_metrics_file(job_directory), cgroup_mount=self.cgroup_mount ) @@ -117,7 +122,7 @@ class CgroupPlugin(InstrumentPlugin): return self._instrument_file_path(job_directory, "_metrics") def __read_metrics(self, path): - metrics = {} + metrics: Dict[str, str] = {} key = None with open(path) as infile: for line in infile: diff --git a/lib/galaxy/job_metrics/instrumenters/collectl.py b/lib/galaxy/job_metrics/instrumenters/collectl.py deleted file mode 100644 index ee740c274d9..00000000000 --- a/lib/galaxy/job_metrics/instrumenters/collectl.py +++ /dev/null @@ -1,215 +0,0 @@ -"""The module describes the ``collectl`` job metrics plugin.""" -import logging -import os -import shutil - -from galaxy import util -from . import InstrumentPlugin -from .. import formatting -from ..collectl import ( - cli, - processes, - subsystems, -) - -log = logging.getLogger(__name__) - -# By default, only grab statistics for user processes (as identified by -# username). -DEFAULT_PROCFILT_ON = "username" -DEFAULT_SUBSYSTEMS = "process" -# Set to zero to flush every collection. -DEFAULT_FLUSH_INTERVAL = "0" - -FORMATTED_RESOURCE_TITLES = { - "PCT": "Percent CPU Usage", - "RSYS": "Disk Reads", - "WSYS": "Disk Writes", -} - -EMPTY_COLLECTL_FILE_MESSAGE = ( - "Skipping process summary due to empty file... job probably did not run long enough for collectl to gather data." -) - - -class CollectlFormatter(formatting.JobMetricFormatter): - def format(self, key, value): - if key == "pid": - return ("Process ID", int(value)) - elif key == "raw_log_path": - return ("Relative Path of Full Collectl Log", value) - elif key == "process_max_AccumT": - return ("Job Runtime (System+User)", formatting.seconds_to_str(float(value))) - else: - _, stat_type, resource_type = key.split("_", 2) - if resource_type.startswith("Vm"): - value_str = f"{int(value)} KB" - elif resource_type in ["RSYS", "WSYS"] and stat_type in ["count", "max", "sum"]: - value_str = "%d (# system calls)" % int(value) - else: - value_str = str(value) - resource_title = FORMATTED_RESOURCE_TITLES.get(resource_type, resource_type) - return (f"{resource_title} ({stat_type})", value_str) - - -class CollectlPlugin(InstrumentPlugin): - """Run collectl along with job to capture system and/or process data - according to specified collectl subsystems. - """ - - plugin_type = "collectl" - formatter = CollectlFormatter() - - def __init__(self, **kwargs): - self.__configure_paths(kwargs) - self.__configure_subsystems(kwargs) - saved_logs_path = kwargs.get("saved_logs_path", "") - if "app" in kwargs: - log.debug(f"Found path for saved logs: {saved_logs_path}") - saved_logs_path = kwargs["app"].config.resolve_path(saved_logs_path) - self.saved_logs_path = saved_logs_path - self.__configure_collectl_recorder_args(kwargs) - self.summarize_process_data = util.asbool(kwargs.get("summarize_process_data", True)) - self.log_collectl_program_output = util.asbool(kwargs.get("log_collectl_program_output", False)) - if self.summarize_process_data: - if subsystems.get_subsystem("process") not in self.subsystems: - raise Exception( - "Collectl plugin misconfigured - cannot summarize_process_data without process subsystem being enabled." - ) - - process_statistics = kwargs.get("process_statistics", None) - # None will let processes module use default set of statistics - # defined there. - self.process_statistics = processes.parse_process_statistics(process_statistics) - - def pre_execute_instrument(self, job_directory): - commands = [] - # Capture PID of process so we can walk its ancestors when building - # statistics for the whole job. - commands.append(f"""echo "$$" > '{self.__pid_file(job_directory)}' """) - # Run collectl in record mode to capture process and system level - # statistics according to supplied subsystems. - commands.append(self.__collectl_record_command(job_directory)) - return commands - - def post_execute_instrument(self, job_directory): - commands = [] - # collectl dies when job script completes, perhaps capture pid of - # collectl above and check if it is still alive to allow tracking if - # collectl ran successfully through the whole job. - return commands - - def job_properties(self, job_id, job_directory): - pid = open(self.__pid_file(job_directory)).read().strip() - contents = os.listdir(job_directory) - try: - rel_path = next(iter(filter(self._is_instrumented_collectl_log, contents))) - path = os.path.join(job_directory, rel_path) - except IndexError: - message = f"Failed to find collectl log in directory {job_directory}, files were {contents}" - raise Exception(message) - - properties = dict( - pid=int(pid), - ) - - if self.saved_logs_path: - destination_rel_dir = os.path.join(*util.directory_hash_id(job_id)) - destination_rel_path = os.path.join(destination_rel_dir, rel_path) - destination_path = os.path.join(self.saved_logs_path, destination_rel_path) - destination_dir = os.path.dirname(destination_path) - if not os.path.isdir(destination_dir): - os.makedirs(destination_dir) - shutil.copyfile(path, destination_path) - properties["raw_log_path"] = destination_rel_path - - if self.summarize_process_data: - # Run collectl in playback and generate statistics of interest - summary_statistics = self.__summarize_process_data(pid, path) - for statistic, value in summary_statistics: - properties[f"process_{'_'.join(statistic)}"] = value - - return properties - - def __configure_paths(self, kwargs): - # 95% of time I would expect collectl to just be installed with apt or - # yum, but if it is manually installed on not on path, allow - # configuration of explicit path - and allow path to be different - # between galaxy job handler (local_collectl_path) and compute node - # (remote_collectl_path). - collectl_path = kwargs.get("collectl_path", "collectl") - self.remote_collectl_path = kwargs.get("remote_collectl_path", collectl_path) - self.local_collectl_path = kwargs.get("local_collectl_path", collectl_path) - - def __configure_subsystems(self, kwargs): - raw_subsystems_str = kwargs.get("subsystems", DEFAULT_SUBSYSTEMS) - raw_subsystems = util.listify(raw_subsystems_str, do_strip=True) - self.subsystems = [subsystems.get_subsystem(_) for _ in raw_subsystems] - - def __configure_collectl_recorder_args(self, kwargs): - collectl_recorder_args = kwargs.copy() - - # Allow deployer to configure separate system and process intervals, - # but if they specify just one - use it for both. Thinking here is this - # plugin's most useful feature is the process level information so - # this is likely what the deployer is attempting to configure. - if "interval" in kwargs and "interval2" not in kwargs: - collectl_recorder_args["interval2"] = kwargs["interval"] - - if "flush" not in kwargs: - collectl_recorder_args["flush"] = DEFAULT_FLUSH_INTERVAL - - procfilt_on = kwargs.get("procfilt_on", DEFAULT_PROCFILT_ON).lower() - # Calculate explicit arguments, rest can just be passed through from - # constructor arguments. - explicit_args = dict( - collectl_path=self.remote_collectl_path, - procfilt=procfilt_argument(procfilt_on), - subsystems=self.subsystems, - ) - collectl_recorder_args.update(explicit_args) - self.collectl_recorder_args = collectl_recorder_args - - def __summarize_process_data(self, pid, collectl_log_path): - playback_cli_args = dict(collectl_path=self.local_collectl_path, playback_path=collectl_log_path, sep="9") - if not os.stat(collectl_log_path).st_size: - log.debug(EMPTY_COLLECTL_FILE_MESSAGE) - return [] - - playback_cli = cli.CollectlCli(**playback_cli_args) - return processes.generate_process_statistics(playback_cli, pid, self.process_statistics) - - def __collectl_recorder_cli(self, job_directory): - cli_args = self.collectl_recorder_args.copy() - cli_args["destination_path"] = self._instrument_file_path(job_directory, "log") - return cli.CollectlCli(**cli_args) - - def __collectl_record_command(self, job_directory): - collectl_cli = self.__collectl_recorder_cli(job_directory) - if self.log_collectl_program_output: - redirect_to = self._instrument_file_path(job_directory, "program_output") - else: - redirect_to = "/dev/null" - return f"{collectl_cli.build_command_line()} > {redirect_to} 2>&1 &" - - def __pid_file(self, job_directory): - return self._instrument_file_path(job_directory, "pid") - - def _is_instrumented_collectl_log(self, filename): - prefix = self._instrument_file_name("log") - return filename.startswith(prefix) and filename.endswith(".raw.gz") - - -def procfilt_argument(procfilt_on): - if procfilt_on == "username": - return "U$USER" - elif procfilt_on == "uid": - return "u$UID" - else: - # Ensure it is empty of None - if procfilt_on or procfilt_on.lower() != "none": - raise Exception("Invalid procfilt_on argument encountered") - return "" - - -__all__ = ("CollectlPlugin",) diff --git a/lib/galaxy/job_metrics/instrumenters/core.py b/lib/galaxy/job_metrics/instrumenters/core.py index 756524979cc..f449c5b29a9 100644 --- a/lib/galaxy/job_metrics/instrumenters/core.py +++ b/lib/galaxy/job_metrics/instrumenters/core.py @@ -1,9 +1,19 @@ """The module describes the ``core`` job metrics plugin.""" import logging import time +from typing import ( + Any, + Dict, + List, +) from . import InstrumentPlugin -from .. import formatting +from ..formatting import ( + FormattedMetric, + JobMetricFormatter, + seconds_to_str, +) +from ..safety import Safety log = logging.getLogger(__name__) @@ -14,19 +24,19 @@ END_EPOCH_KEY = "end_epoch" RUNTIME_SECONDS_KEY = "runtime_seconds" -class CorePluginFormatter(formatting.JobMetricFormatter): - def format(self, key, value): +class CorePluginFormatter(JobMetricFormatter): + def format(self, key: str, value: Any) -> FormattedMetric: value = int(value) if key == GALAXY_SLOTS_KEY: - return ("Cores Allocated", "%d" % value) + return FormattedMetric("Cores Allocated", "%d" % value) elif key == GALAXY_MEMORY_MB_KEY: - return ("Memory Allocated (MB)", "%d" % value) + return FormattedMetric("Memory Allocated (MB)", "%d" % value) elif key == RUNTIME_SECONDS_KEY: - return ("Job Runtime (Wall Clock)", formatting.seconds_to_str(value)) + return FormattedMetric("Job Runtime (Wall Clock)", seconds_to_str(value)) else: # TODO: Use localized version of this from galaxy.ini title = "Job Start Time" if key == START_EPOCH_KEY else "Job End Time" - return (title, time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(value))) + return FormattedMetric(title, time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(value))) class CorePlugin(InstrumentPlugin): @@ -36,23 +46,24 @@ class CorePlugin(InstrumentPlugin): plugin_type = "core" formatter = CorePluginFormatter() + default_safety = Safety.SAFE def __init__(self, **kwargs): pass - def pre_execute_instrument(self, job_directory): + def pre_execute_instrument(self, job_directory: str) -> List[str]: commands = [] commands.append(self.__record_galaxy_slots_command(job_directory)) commands.append(self.__record_galaxy_memory_mb_command(job_directory)) commands.append(self.__record_seconds_since_epoch_to_file(job_directory, "start")) return commands - def post_execute_instrument(self, job_directory): + def post_execute_instrument(self, job_directory: str) -> List[str]: commands = [] commands.append(self.__record_seconds_since_epoch_to_file(job_directory, "end")) return commands - def job_properties(self, job_id, job_directory): + def job_properties(self, job_id, job_directory: str) -> Dict[str, Any]: galaxy_slots_file = self.__galaxy_slots_file(job_directory) galaxy_memory_mb_file = self.__galaxy_memory_mb_file(job_directory) diff --git a/lib/galaxy/job_metrics/instrumenters/cpuinfo.py b/lib/galaxy/job_metrics/instrumenters/cpuinfo.py index 4c1672d3fbf..5d806a4e989 100644 --- a/lib/galaxy/job_metrics/instrumenters/cpuinfo.py +++ b/lib/galaxy/job_metrics/instrumenters/cpuinfo.py @@ -1,22 +1,26 @@ """The module describes the ``cpuinfo`` job metrics plugin.""" import logging import re +from typing import Any from galaxy import util from . import InstrumentPlugin -from .. import formatting +from ..formatting import ( + FormattedMetric, + JobMetricFormatter, +) log = logging.getLogger(__name__) PROCESSOR_LINE = re.compile(r"processor\s*\:\s*(\d+)") -class CpuInfoFormatter(formatting.JobMetricFormatter): - def format(self, key, value): +class CpuInfoFormatter(JobMetricFormatter): + def format(self, key: str, value: Any) -> FormattedMetric: if key == "processor_count": - return "Processor Count", f"{int(value)}" + return FormattedMetric("Processor Count", f"{int(value)}") else: - return key, value + return super().format(key, value) class CpuInfoPlugin(InstrumentPlugin): diff --git a/lib/galaxy/job_metrics/instrumenters/env.py b/lib/galaxy/job_metrics/instrumenters/env.py index b46d4a68796..fd28e89b9e5 100644 --- a/lib/galaxy/job_metrics/instrumenters/env.py +++ b/lib/galaxy/job_metrics/instrumenters/env.py @@ -1,14 +1,19 @@ """The module describes the ``env`` job metrics plugin.""" import logging import re +from typing import ( + List, + Optional, +) from . import InstrumentPlugin -from .. import formatting +from ..formatting import JobMetricFormatter +from ..safety import Safety log = logging.getLogger(__name__) -class EnvFormatter(formatting.JobMetricFormatter): +class EnvFormatter(JobMetricFormatter): pass @@ -19,16 +24,17 @@ class EnvPlugin(InstrumentPlugin): plugin_type = "env" formatter = EnvFormatter() + variables: Optional[List[str]] + default_safety = Safety.UNSAFE def __init__(self, **kwargs): variables_str = kwargs.get("variables", None) if variables_str: - variables = [v.strip() for v in variables_str.split(",")] + self.variables = [v.strip() for v in variables_str.split(",")] else: - variables = None - self.variables = variables + self.variables = None - def pre_execute_instrument(self, job_directory): + def pre_execute_instrument(self, job_directory: str): """Use env to dump all environment variables to a file.""" return f"env > '{self.__env_file(job_directory)}'" diff --git a/lib/galaxy/job_metrics/instrumenters/meminfo.py b/lib/galaxy/job_metrics/instrumenters/meminfo.py index 9c10a381f02..d6da041b0a8 100644 --- a/lib/galaxy/job_metrics/instrumenters/meminfo.py +++ b/lib/galaxy/job_metrics/instrumenters/meminfo.py @@ -4,6 +4,7 @@ import re from galaxy import util from . import InstrumentPlugin from .. import formatting +from ..safety import Safety MEMINFO_LINE = re.compile(r"(\w+)\s*\:\s*(\d+) kB") @@ -24,6 +25,7 @@ class MemInfoPlugin(InstrumentPlugin): plugin_type = "meminfo" formatter = MemInfoFormatter() + default_safety = Safety.SAFE def __init__(self, **kwargs): self.verbose = util.asbool(kwargs.get("verbose", False)) diff --git a/lib/galaxy/job_metrics/safety.py b/lib/galaxy/job_metrics/safety.py new file mode 100644 index 00000000000..99a7ff7a607 --- /dev/null +++ b/lib/galaxy/job_metrics/safety.py @@ -0,0 +1,16 @@ +from enum import Enum + + +class Safety(Enum): + UNSAFE = 2 + POTENTIALLY_SENSITVE = 4 + SAFE = 6 + + +DEFAULT_SAFETY = Safety.POTENTIALLY_SENSITVE + + +__all__ = ( + "DEFAULT_SAFETY", + "Safety", +) diff --git a/lib/galaxy/managers/jobs.py b/lib/galaxy/managers/jobs.py index 4c7cca0ed4c..2124f0a74e2 100644 --- a/lib/galaxy/managers/jobs.py +++ b/lib/galaxy/managers/jobs.py @@ -22,6 +22,10 @@ from galaxy.exceptions import ( ObjectNotFound, RequestParameterInvalidException, ) +from galaxy.job_metrics import ( + RawMetric, + Safety, +) from galaxy.managers.collections import DatasetCollectionManager from galaxy.managers.datasets import DatasetManager from galaxy.managers.hdas import HDAManager @@ -658,24 +662,21 @@ def summarize_job_metrics(trans, job): Precondition: the caller has verified the job is accessible to the user represented by the trans parameter. """ - if not trans.user_is_admin and not trans.app.config.expose_potentially_sensitive_job_metrics: - return [] - - def metric_to_dict(metric): - metric_name = metric.metric_name - metric_value = metric.metric_value - metric_plugin = metric.plugin - title, value = trans.app.job_metrics.format(metric_plugin, metric_name, metric_value) - return dict( - title=title, - value=value, - plugin=metric_plugin, - name=metric_name, - raw_value=str(metric_value), + safety_level = Safety.SAFE + if trans.user_is_admin: + safety_level = Safety.UNSAFE + elif trans.app.config.expose_potentially_sensitive_job_metrics: + safety_level = Safety.POTENTIALLY_SENSITVE + raw_metrics = [ + RawMetric( + m.metric_name, + m.metric_value, + m.plugin, ) - - metrics = [m for m in job.metrics if m.plugin != "env" or trans.user_is_admin] - return list(map(metric_to_dict, metrics)) + for m in job.metrics + ] + dictifiable_metrics = trans.app.job_metrics.dictifiable_metrics(raw_metrics, safety_level) + return [d.dict() for d in dictifiable_metrics] def summarize_destination_params(trans, job): diff --git a/packages/job_metrics/setup.py b/packages/job_metrics/setup.py index 25a3f672873..3ee4855fc0d 100644 --- a/packages/job_metrics/setup.py +++ b/packages/job_metrics/setup.py @@ -36,7 +36,6 @@ PACKAGES = [ "galaxy", "galaxy.job_metrics", "galaxy.job_metrics.instrumenters", - "galaxy.job_metrics.collectl", ] ENTRY_POINTS = """ [console_scripts] diff --git a/setup.cfg b/setup.cfg index 1746400441b..f09ce93db68 100644 --- a/setup.cfg +++ b/setup.cfg @@ -202,8 +202,6 @@ check_untyped_defs = False check_untyped_defs = False [mypy-galaxy.model.custom_types] check_untyped_defs = False -[mypy-galaxy.job_metrics.collectl.processes] -check_untyped_defs = False [mypy-galaxy.datatypes.util.gff_util] check_untyped_defs = False [mypy-galaxy.datatypes.dataproviders.line] @@ -245,8 +243,6 @@ check_untyped_defs = False check_untyped_defs = False [mypy-galaxy.jobs.runners.util.cli.shell.rsh] check_untyped_defs = False -[mypy-galaxy.job_metrics] -check_untyped_defs = False [mypy-galaxy.datatypes.display_applications.parameters] check_untyped_defs = False [mypy-tool_shed.webapp.search.tool_search] @@ -303,12 +299,6 @@ check_untyped_defs = False check_untyped_defs = False [mypy-galaxy.model.dataset_collections.matching] check_untyped_defs = False -[mypy-galaxy.job_metrics.instrumenters.env] -check_untyped_defs = False -[mypy-galaxy.job_metrics.instrumenters.collectl] -check_untyped_defs = False -[mypy-galaxy.job_metrics.instrumenters.cgroup] -check_untyped_defs = False [mypy-galaxy.datatypes.sniff] check_untyped_defs = False [mypy-galaxy.containers] diff --git a/test/unit/job_metrics/test_cgroups.py b/test/unit/job_metrics/test_cgroups.py new file mode 100644 index 00000000000..863ea89c558 --- /dev/null +++ b/test/unit/job_metrics/test_cgroups.py @@ -0,0 +1,133 @@ +from galaxy.job_metrics.instrumenters.cgroup import CgroupPlugin + +CGROUP_PRODUCTION_EXAMPLE_2201 = """__cpu.cfs_period_us__ +100000 +__cpu.cfs_quota_us__ +-1 +__cpu.rt_period_us__ +1000000 +__cpu.rt_runtime_us__ +0 +__cpu.shares__ +1024 +__cpu.stat__ +nr_periods 0 +nr_throttled 0 +throttled_time 0 +__cpuacct.stat__ +user 616 +system 98 +__cpuacct.usage__ +7265342042 +__cpuacct.usage_percpu__ +0 7269963877 0 0 0 0 0 0 0 0 +__memory.failcnt__ +0 +__memory.force_empty__ +__memory.kmem.failcnt__ +0 +__memory.kmem.limit_in_bytes__ +9223372036854771712 +__memory.kmem.max_usage_in_bytes__ +0 +__memory.kmem.slabinfo__ +__memory.kmem.tcp.failcnt__ +0 +__memory.kmem.tcp.limit_in_bytes__ +9223372036854771712 +__memory.kmem.tcp.max_usage_in_bytes__ +0 +__memory.kmem.tcp.usage_in_bytes__ +0 +__memory.kmem.usage_in_bytes__ +0 +__memory.limit_in_bytes__ +9223372036854771712 +__memory.max_usage_in_bytes__ +264409088 +__memory.memsw.failcnt__ +0 +__memory.memsw.limit_in_bytes__ +9223372036854771712 +__memory.memsw.max_usage_in_bytes__ +264409088 +__memory.memsw.usage_in_bytes__ +103460864 +__memory.move_charge_at_immigrate__ +0 +__memory.numa_stat__ +total=25226 N0=25226 +file=25071 N0=25071 +anon=155 N0=155 +unevictable=0 N0=0 +__memory.oom_control__ +oom_kill_disable 0 +under_oom 0 +__memory.pressure_level__ +__memory.soft_limit_in_bytes__ +9223372036854771712 +__memory.stat__ +cache 102690816 +rss 655360 +rss_huge 0 +mapped_file 0 +swap 0 +pgpgin 136493 +pgpgout 111262 +pgfault 170611 +pgmajfault 961 +inactive_anon 0 +active_anon 634880 +inactive_file 101662720 +active_file 1028096 +unevictable 0 +hierarchical_memory_limit 7853834240 +hierarchical_memsw_limit 7853834240 +total_cache 102690816 +total_rss 655360 +total_rss_huge 0 +total_mapped_file 0 +total_swap 0 +total_pgpgin 136493 +total_pgpgout 111262 +total_pgfault 170611 +total_pgmajfault 961 +total_inactive_anon 0 +total_active_anon 634880 +total_inactive_file 101662720 +total_active_file 1028096 +total_unevictable 0 +__memory.swappiness__ +30 +__memory.usage_in_bytes__ +103391232 +__memory.use_hierarchy__ +1 +""" + + +def test_cgroup_collection(tmpdir): + plugin = CgroupPlugin() + job_dir = tmpdir.mkdir("job") + job_dir.join("__instrument_cgroup__metrics").write(CGROUP_PRODUCTION_EXAMPLE_2201) + properties = plugin.job_properties(1, job_dir) + assert "cpuacct.usage" in properties + assert properties["cpuacct.usage"] == 7265342042 + assert "memory.limit_in_bytes" in properties + assert properties["memory.limit_in_bytes"] == 9223372036854771712 + + +def test_instrumentation(tmpdir): + # don't actually run the instrumentation but at least exercise the code the and make + # sure templating includes cgroup_mount override. + mock_cgroup_mount = "/proc/sys/made/up/cgroup/mount" + plugin = CgroupPlugin( + cgroup_mount=mock_cgroup_mount, + ) + assert plugin.pre_execute_instrument(tmpdir) is None + collection_commands = plugin.post_execute_instrument(tmpdir) + assert len(collection_commands) == 2 + cpu_command = collection_commands[0] + assert mock_cgroup_mount in cpu_command + memory_command = collection_commands[1] + assert mock_cgroup_mount in memory_command diff --git a/test/unit/job_metrics/test_core.py b/test/unit/job_metrics/test_core.py new file mode 100644 index 00000000000..b333c8b2bd8 --- /dev/null +++ b/test/unit/job_metrics/test_core.py @@ -0,0 +1,31 @@ +import subprocess + +from galaxy.job_metrics.instrumenters.core import ( + CorePlugin, + GALAXY_MEMORY_MB_KEY, + GALAXY_SLOTS_KEY, +) +from galaxy.util import listify + + +def test_core_instrumentation(tmpdir): + core_plugin = CorePlugin() + env = {"GALAXY_SLOTS": "4", "GALAXY_MEMORY_MB": "1024"} + _run_plugin(core_plugin, tmpdir, env) + properties = core_plugin.job_properties(1, tmpdir) + assert properties[GALAXY_SLOTS_KEY] == 4 + assert properties[GALAXY_MEMORY_MB_KEY] == 1024 + + +def _run_plugin(plugin, work_dir, env=None): + setup_commands = plugin.pre_execute_instrument(work_dir) + teardown_commands = plugin.post_execute_instrument(work_dir) + if setup_commands is not None: + _run(setup_commands, work_dir, env) + if teardown_commands is not None: + _run(teardown_commands, work_dir, env) + + +def _run(commands, work_dir, env): + command_str = "\n".join(listify(commands)) + return subprocess.run(command_str, shell=True, cwd=work_dir, env=env) diff --git a/test/unit/job_metrics/test_cpuinfo.py b/test/unit/job_metrics/test_cpuinfo.py new file mode 100644 index 00000000000..04f7c4b0355 --- /dev/null +++ b/test/unit/job_metrics/test_cpuinfo.py @@ -0,0 +1,261 @@ +from galaxy.job_metrics.instrumenters.cpuinfo import CpuInfoPlugin + +CPUINFO_PRODUCTION_EXAMPLE_2201 = """processor : 0 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 0 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 0 +initial apicid : 0 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 1 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 1 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 1 +initial apicid : 1 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 2 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 2 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 2 +initial apicid : 2 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 3 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 3 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 3 +initial apicid : 3 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 4 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 4 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 4 +initial apicid : 4 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 5 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 5 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 5 +initial apicid : 5 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 6 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 6 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 6 +initial apicid : 6 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 7 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 7 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 7 +initial apicid : 7 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 8 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 8 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 8 +initial apicid : 8 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +processor : 9 +vendor_id : GenuineIntel +cpu family : 6 +model : 63 +model name : Intel(R) Xeon(R) CPU E5-2680 v3 @ 2.50GHz +stepping : 2 +microcode : 0x1 +cpu MHz : 2494.222 +cache size : 16384 KB +physical id : 9 +siblings : 1 +core id : 0 +cpu cores : 1 +apicid : 9 +initial apicid : 9 +fpu : yes +fpu_exception : yes +cpuid level : 13 +wp : yes +flags : fpu vme de pse tsc msr pae mce cx8 apic sep mtrr pge mca cmov pat pse36 clflush mmx fxsr sse sse2 ss syscall nx pdpe1gb rdtscp lm constant_tsc arch_perfmon rep_good nopl xtopology eagerfpu pni pclmulqdq vmx ssse3 fma cx16 pcid sse4_1 sse4_2 x2apic movbe popcnt tsc_deadline_timer aes xsave avx f16c rdrand hypervisor lahf_lm abm invpcid_single ssbd ibrs ibpb stibp tpr_shadow vnmi flexpriority ept vpid fsgsbase tsc_adjust bmi1 avx2 smep bmi2 erms invpcid xsaveopt arat md_clear spec_ctrl intel_stibp +bogomips : 4988.44 +clflush size : 64 +cache_alignment : 64 +address sizes : 46 bits physical, 48 bits virtual +power management: +""" + + +def test_cpuinfo_collection(tmpdir): + plugin = CpuInfoPlugin() + job_dir = tmpdir.mkdir("job") + job_dir.join("__instrument_cpuinfo_cpuinfo").write(CPUINFO_PRODUCTION_EXAMPLE_2201) + properties = plugin.job_properties(1, job_dir) + assert properties["processor_count"] == 10 diff --git a/test/unit/job_metrics/test_env.py b/test/unit/job_metrics/test_env.py new file mode 100644 index 00000000000..8d2758128ee --- /dev/null +++ b/test/unit/job_metrics/test_env.py @@ -0,0 +1,27 @@ +from galaxy.job_metrics.instrumenters.env import EnvPlugin +from .test_core import _run_plugin + +TEST_ENV = { + "FOOBAR": "moocow", + "DOGCAT": "catdog", +} + + +def test_env_instrumentation(tmpdir): + env_plugin = EnvPlugin() + _run_plugin(env_plugin, tmpdir, TEST_ENV) + props = env_plugin.job_properties(1, tmpdir) + assert props + assert props["FOOBAR"] == "moocow" + assert props["DOGCAT"] == "catdog" + + +def test_env_restrict_variables(tmpdir): + env_plugin = EnvPlugin( + variables="FOOBAR,FOOBAR2", + ) + _run_plugin(env_plugin, tmpdir, TEST_ENV) + props = env_plugin.job_properties(1, tmpdir) + assert props + assert props["FOOBAR"] == "moocow" + assert "DOGCAT" not in props diff --git a/test/unit/job_metrics/test_job_metrics.py b/test/unit/job_metrics/test_job_metrics.py index 7c5acd5bb21..07930c0684e 100644 --- a/test/unit/job_metrics/test_job_metrics.py +++ b/test/unit/job_metrics/test_job_metrics.py @@ -1,13 +1,80 @@ +from typing import ( + Any, + Optional, +) + from galaxy.job_metrics import ( formatting, JobMetrics, + RawMetric, ) +from galaxy.job_metrics.safety import Safety + +TEST_JOBS_METRICS = JobMetrics() +TEST_LINUX_OS = "Ubuntu Linux 21.10" -def test_job_metrics_load(): - # Just construct the manager to make sure all the plugin classes load fine - # and package is configured properly. - JobMetrics() +def test_job_metrics_format_env(): + _assert_format( + "env", + "moo", + "cow", + assert_title="moo", + assert_value="cow", + ) + + +def test_job_metrics_format_core(): + _assert_format( + "core", + "galaxy_slots", + "4", + assert_title="Cores Allocated", + assert_value="4", + ) + + +def test_job_metrics_format_cgroup(): + _assert_format( + "cgroup", + "cpuacct.usage", + 7265342042, + assert_title="CPU Time", + assert_value="7.265342042 seconds", + ) + _assert_format( + "cgroup", + "memory.limit_in_bytes", + 9223372036854771712, + assert_title="Memory limit on cgroup (MEM)", + assert_value="8.0 EB", + ) + + +def test_job_metrics_uname(): + _assert_format( + "uname", + "moo", + TEST_LINUX_OS, + assert_title="Operating System", + assert_value=TEST_LINUX_OS, + ) + + +def test_metrics_dictifiable(): + test_metrics = [ + RawMetric("galaxy_slots", "4", "core"), + RawMetric("uname", TEST_LINUX_OS, "uname"), + RawMetric("SSH_AUTH_SOCK", "/private/tmp/com.apple.launchd.Nw6gC2VOCr/Listeners", "env"), + ] + dictifiable_metrics = TEST_JOBS_METRICS.dictifiable_metrics(test_metrics, Safety.POTENTIALLY_SENSITVE) + _assert_metrics_of_type(dictifiable_metrics, ["core", "uname"]) + + dictifiable_metrics = TEST_JOBS_METRICS.dictifiable_metrics(test_metrics, Safety.SAFE) + _assert_metrics_of_type(dictifiable_metrics, ["core"]) + + dictifiable_metrics = TEST_JOBS_METRICS.dictifiable_metrics(test_metrics, Safety.UNSAFE) + _assert_metrics_of_type(dictifiable_metrics, ["core", "uname", "env"]) def test_job_metric_formatting(): @@ -30,3 +97,19 @@ def test_job_metric_formatting(): assert formatting.seconds_to_str(7260) == "2 hours and 1 minute" assert formatting.seconds_to_str(7320) == "2 hours and 2 minutes" assert formatting.seconds_to_str(36181) == "10 hours and 3 minutes" + + +def _assert_metrics_of_type(metric_list, expected_types): + assert len(metric_list) == len(expected_types) + for dictifiable_metric, expected_type in zip(metric_list, expected_types): + assert dictifiable_metric.plugin == expected_type + + +def _assert_format( + plugin: str, key: str, value: Any, assert_title: Optional[str] = None, assert_value: Optional[str] = None +): + result = TEST_JOBS_METRICS.format(plugin, key, value) + if assert_title is not None: + assert result[0] == assert_title + if assert_value is not None: + assert result[1] == assert_value