Merge pull request #13634 from jmchilton/job_metrics_revision

Overhaul job metrics.
This commit is contained in:
Björn Grüning
2022-04-04 22:27:21 +02:00
committed by GitHub
25 changed files with 770 additions and 826 deletions
@@ -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:
-1
View File
@@ -12,7 +12,6 @@ Subpackages
.. toctree::
:maxdepth: 4
galaxy.job_metrics.collectl
galaxy.job_metrics.instrumenters
Submodules
+102 -4
View File
@@ -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",
)
@@ -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/.
"""
-141
View File
@@ -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",)
@@ -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",)
-26
View File
@@ -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
@@ -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",)
+12 -3
View File
@@ -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)
@@ -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))
+13 -8
View File
@@ -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:
@@ -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",)
+21 -10
View File
@@ -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)
@@ -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):
+12 -6
View File
@@ -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)}'"
@@ -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))
+16
View File
@@ -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",
)
+18 -17
View File
@@ -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):
-1
View File
@@ -36,7 +36,6 @@ PACKAGES = [
"galaxy",
"galaxy.job_metrics",
"galaxy.job_metrics.instrumenters",
"galaxy.job_metrics.collectl",
]
ENTRY_POINTS = """
[console_scripts]
-10
View File
@@ -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]
+133
View File
@@ -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
+31
View File
@@ -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)
+261
View File
@@ -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
+27
View File
@@ -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
+87 -4
View File
@@ -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