From ff3f8aedd7432ed237474a95187a52d65e938db0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jaroslav=20=C5=A0karvada?= Date: Sat, 4 Jul 2015 11:23:32 +0200 Subject: [PATCH] plugin_scheduler: added support for runtime tuning through perf MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit It can now tune newly created processes. This functionality is by default on and can be disabled in tuned profile by setting: [scheduler] runtime = 0 This is temporaly workaround. In the future this will be handled by per-plugin dynamic_tuning configuration and not by "runtime" option. resolves: rhbz#1148546 Signed-off-by: Jaroslav Škarvada --- tuned.spec | 2 +- tuned/plugins/plugin_scheduler.py | 130 ++++++++++++++++++++++++++---- tuned/utils/commands.py | 14 +++- 3 files changed, 129 insertions(+), 17 deletions(-) diff --git a/tuned.spec b/tuned.spec index 05fa6c3..750bb25 100644 --- a/tuned.spec +++ b/tuned.spec @@ -22,7 +22,7 @@ Requires(preun): systemd Requires(postun): systemd Requires: python-decorator, dbus-python, pygobject3-base, python-pyudev Requires: virt-what, python-configobj, ethtool, gawk, kernel-tools, hdparm -Requires: util-linux +Requires: util-linux, python-perf %description The tuned package contains a daemon that tunes system settings dynamically. diff --git a/tuned/plugins/plugin_scheduler.py b/tuned/plugins/plugin_scheduler.py index 2ecb559..8c0b00d 100644 --- a/tuned/plugins/plugin_scheduler.py +++ b/tuned/plugins/plugin_scheduler.py @@ -3,6 +3,10 @@ from decorators import * import tuned.logs import re from subprocess import * +import threading +import perf +import select +import tuned.consts as consts from tuned.utils.commands import commands log = tuned.logs.get() @@ -16,9 +20,14 @@ class SchedulerPlugin(base.Plugin): _dict_sched2param = {"SCHED_FIFO":"f", "SCHED_BATCH":"b", "SCHED_RR":"r", "SCHED_OTHER":"o", "SCHED_IDLE":"i"} - def __init__(self, *args, **kwargs): - super(self.__class__, self).__init__(*args, **kwargs) + def __init__(self, monitors_repository, storage_factory, hardware_inventory, device_matcher, instance_factory, global_cfg, variables): + super(self.__class__, self).__init__(monitors_repository, storage_factory, hardware_inventory, device_matcher, instance_factory, global_cfg, variables) self._has_dynamic_options = True + self._daemon = consts.CFG_DEF_DAEMON + self._sleep_interval = int(consts.CFG_DEF_SLEEP_INTERVAL) + if global_cfg is not None: + self._daemon = global_cfg.get_bool(consts.CFG_DAEMON, consts.CFG_DEF_DAEMON) + self._sleep_interval = int(global_cfg.get(consts.CFG_SLEEP_INTERVAL, consts.CFG_DEF_SLEEP_INTERVAL)) self._cmd = commands() def _scheduler_storage_key(self, instance): @@ -27,6 +36,10 @@ class SchedulerPlugin(base.Plugin): def _instance_init(self, instance): instance._has_dynamic_tuning = False instance._has_static_tuning = True + # this is hack, runtime_tuning should be covered by dynamic_tuning configuration + # TODO: add per plugin dynamic tuning configuration and use dynamic_tuning configuration + # instead of runtime_tuning + instance._runtime_tuning = True # FIXME: do we want to do this here? # recover original values in case of crash @@ -38,10 +51,40 @@ class SchedulerPlugin(base.Plugin): self._storage.unset(self._scheduler_storage_key(instance)) instance._scheduler = instance.options + for k in instance._scheduler: + instance._scheduler[k] = self._variables.expand(instance._scheduler[k]) + if self._cmd.get_bool(instance._scheduler.get("runtime", 1)) == "0": + instance._runtime_tuning = False + instance._terminate = threading.Event() + if self._daemon and instance._runtime_tuning: + try: + instance._cpus = perf.cpu_map() + instance._threads = perf.thread_map() + evsel = perf.evsel(task = 1, comm = 1, mmap = 0, + wakeup_events = 1, watermark = 1, + sample_type = perf.SAMPLE_TID | perf.SAMPLE_CPU) + evsel.open(cpus = instance._cpus, threads = instance._threads) + instance._evlist = perf.evlist(instance._cpus, instance._threads) + instance._evlist.add(evsel) + instance._evlist.mmap() + # no perf + except: + instance._runtime_tuning = False def _instance_cleanup(self, instance): pass + def get_process(self, pid): + cmd = self._cmd.read_file("/proc/" + pid + "/comm", no_error = True) + if cmd == "": + return "" + cmd = cmd.strip() + cmdline = self._cmd.read_file("/proc/" + pid + "/cmdline", no_error = True) + if cmdline == "": + return "[" + cmd + "]" + else: + return cmdline.replace("\0", " ").strip() + def get_processes(self): (rc, out) = self._cmd.execute(["ps", "-eopid,cmd", "--no-headers"]) if rc != 0 or len(out) <= 0: @@ -81,6 +124,7 @@ class SchedulerPlugin(base.Plugin): def _schedcfg2param(self, sched): if sched in ["f", "b", "r", "o"]: return "-" + sched + # including '*' else: return "" @@ -107,38 +151,57 @@ class SchedulerPlugin(base.Plugin): log.debug("setting affinity to '%s' for PID '%s'" % (affinity, pid)) self._cmd.execute(["taskset", "-p", str(affinity), str(pid)]) + #tune process and store previous values + def _tune_process(self, instance, pid, cmd, sched, prio, affinity): + #rt[0] - prev_sched, rt[1] - prev_prio + rt = self._get_rt(pid) + prev_affinity = self._get_affinity(pid) + if prev_affinity is not None and rt is not None and len(rt) == 2 and rt[0] is not None and rt[1] is not None: + instance._scheduler_original[pid] = (cmd, rt[0], rt[1], prev_affinity) + self._set_rt(pid, self._schedcfg2param(sched), prio) + if affinity != "*": + self._set_affinity(pid, affinity) + def _instance_apply_static(self, instance): ps = self.get_processes() if ps is None: log.error("error applying tuning, cannot get information about running processes") return - for k in instance._scheduler: - instance._scheduler[k] = self._variables.expand(instance._scheduler[k]) - sched_cfg = map(lambda (option, value): (option, value.split(":", 4)), instance._scheduler.items()) - buf = filter(lambda (option, vals): re.match(r"group\.", option) and len(vals) == 5, sched_cfg) - sched_cfg = sorted(buf, key=lambda (option, vals): vals[0]) + instance._sched_cfg = map(lambda (option, value): (option, value.split(":", 4)), instance._scheduler.items()) + buf = filter(lambda (option, vals): re.match(r"group\.", option) and len(vals) == 5, instance._sched_cfg) + instance._sched_cfg = sorted(buf, key=lambda (option, vals): vals[0]) sched_all = dict() - for option, vals in sched_cfg: + # for runtime tunning + instance._sched_lookup = {} + for option, vals in instance._sched_cfg: try: r = re.compile(vals[4]) except re.error as e: log.error("error compiling regular expression: '%s'" % str(vals[4])) continue processes = filter(lambda (pid, cmd): re.search(r, cmd) is not None, ps.items()) + #cmd - process name, option - group name, vals[0] - rule prio, vals[1] - sched, vals[2] - prio, + #vals[3] - affinity, vals[4] - regex sched = dict(map(lambda (pid, cmd): (pid, (cmd, option, vals[1], vals[2], vals[3], vals[4])), processes)) sched_all.update(sched) + v4 = str(vals[4]).replace("(", r"\(") + v4 = v4.replace(")", r"\)") + instance._sched_lookup[v4] = [vals[1], vals[2], vals[3]] for pid, vals in sched_all.items(): - (sched, prio) = self._get_rt(pid) - affinity = self._get_affinity(pid) - if affinity is not None and sched is not None and prio is not None: - instance._scheduler_original[pid] = (vals[0], sched, prio, affinity) - self._set_rt(pid, self._schedcfg2param(vals[2]), vals[3]) - if vals[4] != "*": - self._set_affinity(pid, vals[4]) + #vals[0] - process name, vals[1] - rule prio, vals[2] - sched, vals[3] - prio, vals[4] - affinity, + #vals[5] - regex + self._tune_process(instance, pid, vals[0], vals[2], vals[3], vals[4]) self._storage.set("options", instance._scheduler_original) + if self._daemon and instance._runtime_tuning: + instance._thread = threading.Thread(target = self._thread_code, args = [instance]) + instance._thread.start() def _instance_unapply_static(self, instance, profile_switch = False): ps = self.get_processes() + if self._daemon and instance._runtime_tuning: + instance._terminate.set() + instance._thread.join() + for pid, vals in instance._scheduler_original.iteritems(): # if command line for the pid didn't change, it's very probably the same process try: @@ -147,3 +210,40 @@ class SchedulerPlugin(base.Plugin): self._set_affinity(pid, vals[3]) except KeyError as e: pass + + def _add_pid(self, instance, pid): + cmd = self.get_process(pid) + # check to filter short living process + if cmd == "": + return + v = self._cmd.re_lookup(instance._sched_lookup, cmd) + if v is not None and not pid in instance._scheduler_original: + log.debug("tuning new process '%s' with pid '%s' by '%s'" % (cmd, pid, str(v))) + #v[0] - sched, v[1] - prio, v[2] - affinity + self._tune_process(instance, pid, cmd, v[0], v[1], v[2]) + self._storage.set("options", instance._scheduler_original) + + def _remove_pid(self, instance, pid): + if pid in instance._scheduler_original: + del instance._scheduler_original[pid] + log.debug("removed PID %s from the rollback database" % pid) + self._storage.set("options", instance._scheduler_original) + + def _thread_code(self, instance): + poll = select.poll() + for fd in instance._evlist.get_pollfd(): + poll.register(fd) + while not instance._terminate.is_set(): + # timeout to poll in milliseconds + if len(poll.poll(self._sleep_interval * 1000)) > 0 and not instance._terminate.is_set(): + read_events = True + while read_events: + read_events = False + for cpu in instance._cpus: + event = instance._evlist.read_on_cpu(cpu) + if event: + read_events = True + if event.type == perf.RECORD_COMM: + self._add_pid(instance, str(event.pid)) + elif event.type == perf.RECORD_EXIT: + self._remove_pid(instance, str(event.pid)) diff --git a/tuned/utils/commands.py b/tuned/utils/commands.py index 9c94572..1a31912 100644 --- a/tuned/utils/commands.py +++ b/tuned/utils/commands.py @@ -41,13 +41,25 @@ class commands: return l # Do multiple regex replaces in 's' according to lookup table described by - # dictionary 'd', e.g.: d = {"re1": "replace1", "re2": "replace2"} + # dictionary 'd', e.g.: d = {"re1": "replace1", "re2": "replace2", ...} def multiple_re_replace(self, d, s): if len(d) == 0 or s is None: return s r = re.compile("(%s)" % ")|(".join(d.keys())) return r.sub(lambda mo: d.values()[mo.lastindex - 1], s) + # Do regex lookup on 's' according to lookup table described by + # dictionary 'd' and return corresponding value from the dictionary, + # e.g.: d = {"re1": val1, "re2": val2, ...} + def re_lookup(self, d, s): + if len(d) == 0 or s is None: + return None + r = re.compile("(%s)" % ")|(".join(d.keys())) + mo = r.search(s) + if mo: + return d.values()[mo.lastindex - 1] + return None + def write_to_file(self, f, data): self._debug("Writing to file: %s < %s" % (f, data)) try: