@cryptotaxi247 / netdata-1 / commits / 2175673e2

python.d.plugin: use separate process for initial module checking (#5552)

##### Summary This PR adds (major) changes only to `python.d.plugin` file. Fixes: #5525 `pyhton.d.plugin` imports a lot of additional packages during initial module initialization/job creating/checking and there is no way to unimport them, even if they arn't needed. It consumes relatively a lot of ram. ___ Memory utilization comparing before/after the PR (one job `example` module, py3.7.2): > 21.1 => 8.8 MiB ![screenshot_20190305_111837](https://user-images.githubusercontent.com/22274335/53791147-c27a6e00-3f39-11e9-8eaf-8ac3809a3b6e.png) ##### Component Name [`collectors/python.d.plugin`](https://github.com/netdata/netdata/blob/master/collectors/python.d.plugin/python.d.plugin.in) ##### Additional Information This PR adds separate process for initial module checking. Logic: - main process spawns checker process - checker process loads every module, loads module config, creates jobs and runs job.check() for every job, if check success it adds the job to the list. - checker process returns list of modules and jobs. - main process loads only active modules, etc.

Ilya Mashchenko committed Mar 7, 2019 at 13:49 UTC 2175673e29f8ee35d2e48719c6cf0169ddd31405
5 files changed +648 -356
collectors/python.d.plugin/powerdns/powerdns.chart.py
+1
@@ -125,6 +125,7 @@ class Service(UrlService):
125 UrlService.__init__(self, configuration=configuration, name=name)
126 self.order = ORDER
127 self.definitions = CHARTS
128 + self.url = configuration.get('url', 'http://127.0.0.1:8081/api/v1/servers/localhost/statistics')
129
130 def check(self):
131 self._manager = self._build_manager()
collectors/python.d.plugin/python.d.plugin.in
+623 -350
@@ -8,426 +8,699 @@ echo "ERROR python IS NOT AVAILABLE IN THIS SYSTEM")" "$0" "$@" # '''
8 # Author: Ilya Mashchenko (l2isbad)
9 # SPDX-License-Identifier: GPL-3.0-or-later
10
11 +
12 +import collections
13 +import copy
14 import gc
15 +import multiprocessing
16 import os
17 +import re
18 import sys
19 +import time
20 import threading
21
16 -from re import sub
17 -from sys import version_info, argv, stdout
18 -from time import sleep
22 +ENV_NETDATA_USER_CONFIG_DIR = 'NETDATA_USER_CONFIG_DIR'
23 +ENV_NETDATA_STOCK_CONFIG_DIR = 'NETDATA_STOCK_CONFIG_DIR'
24 +ENV_NETDATA_PLUGINS_DIR = 'NETDATA_PLUGINS_DIR'
25 +ENV_NETDATA_UPDATE_EVERY = 'NETDATA_UPDATE_EVERY'
26 +
27 +
28 +def dirs():
29 + user_config = os.getenv(
30 + ENV_NETDATA_USER_CONFIG_DIR,
31 + '@configdir_POST@',
32 + )
33 + stock_config = os.getenv(
34 + ENV_NETDATA_STOCK_CONFIG_DIR,
35 + '@libconfigdir_POST@',
36 + )
37 + modules_user_config = os.path.join(user_config, 'python.d')
38 + modules_stock_config = os.path.join(stock_config, 'python.d')
39 +
40 + modules = os.path.abspath(
41 + os.getenv(
42 + ENV_NETDATA_PLUGINS_DIR,
43 + os.path.dirname(__file__),
44 + ) + '/../python.d'
45 + )
46 + pythond_packages = os.path.join(modules, 'python_modules')
47 +
48 + return collections.namedtuple(
49 + 'Dirs',
50 + [
51 + 'user_config',
52 + 'stock_config',
53 + 'modules_user_config',
54 + 'modules_stock_config',
55 + 'modules',
56 + 'pythond_packages',
57 + ]
58 + )(
59 + user_config,
60 + stock_config,
61 + modules_user_config,
62 + modules_stock_config,
63 + modules,
64 + pythond_packages,
65 + )
66 +
67 +
68 +DIRS = dirs()
69 +
70 +sys.path.append(DIRS.pythond_packages)
71 +
72 +
73 +from bases.collection import safe_print
74 +from bases.loggers import PythonDLogger
75 +from bases.loaders import load_config
76 +from bases.loaders import load_module as _load_module
77
20 -GC_RUN = True
21 -GC_COLLECT_EVERY = 300
78 +try:
79 + from collections import OrderedDict
80 +except ImportError:
81 + from third_party.ordereddict import OrderedDict
82
23 -PY_VERSION = version_info[:2]
83
25 -USER_CONFIG_DIR = os.getenv('NETDATA_USER_CONFIG_DIR', '@configdir_POST@')
26 -STOCK_CONFIG_DIR = os.getenv('NETDATA_STOCK_CONFIG_DIR', '@libconfigdir_POST@')
84 +END_TASK_MARKER = None
85
28 -PLUGINS_USER_CONFIG_DIR = os.path.join(USER_CONFIG_DIR, 'python.d')
29 -PLUGINS_STOCK_CONFIG_DIR = os.path.join(STOCK_CONFIG_DIR, 'python.d')
86 +IS_ATTY = sys.stdout.isatty()
87
88 +PLUGIN_CONF_FILE = 'python.d.conf'
89
32 -PLUGINS_DIR = os.path.abspath(os.getenv(
33 - 'NETDATA_PLUGINS_DIR',
34 - os.path.dirname(__file__)) + '/../python.d')
90 +MODULE_SUFFIX = '.chart.py'
91
92 +OBSOLETED_MODULES = (
93 + 'apache_cache', # replaced by web_log
94 + 'gunicorn_log', # replaced by web_log
95 + 'nginx_log', # replaced by web_log
96 + 'cpufreq', # rewritten in C
97 + 'cpuidle', # rewritten in C
98 + 'mdstat', # rewritten in C
99 + 'linux_power_supply', # rewritten in C
100 +)
101
37 -PYTHON_MODULES_DIR = os.path.join(PLUGINS_DIR, 'python_modules')
102
39 -sys.path.append(PYTHON_MODULES_DIR)
103 +AVAILABLE_MODULES = [
104 + m[:-len(MODULE_SUFFIX)] for m in sorted(os.listdir(DIRS.modules))
105 + if m.endswith(MODULE_SUFFIX) and m[:-len(MODULE_SUFFIX)] not in OBSOLETED_MODULES
106 +]
107
41 -from bases.loaders import ModuleAndConfigLoader # noqa: E402
42 -from bases.loggers import PythonDLogger # noqa: E402
43 -from bases.collection import setdefault_values, run_and_exit, safe_print # noqa: E402
108 +PLUGIN_BASE_CONF = {
109 + 'enabled': True,
110 + 'default_run': True,
111 + 'gc_run': True,
112 + 'gc_interval': 300,
113 +}
114
45 -try:
46 - from collections import OrderedDict
47 -except ImportError:
48 - from third_party.ordereddict import OrderedDict
115 +JOB_BASE_CONF = {
116 + 'update_every': os.getenv(ENV_NETDATA_UPDATE_EVERY, 1),
117 + 'priority': 60000,
118 + 'autodetection_retry': 0,
119 + 'chart_cleanup': 10,
120 + 'penalty': True,
121 + 'name': str(),
122 +}
123
50 -IS_ATTY = stdout.isatty()
124
52 -BASE_CONFIG = {'update_every': os.getenv('NETDATA_UPDATE_EVERY', 1),
53 - 'priority': 60000,
54 - 'autodetection_retry': 0,
55 - 'chart_cleanup': 10,
56 - 'penalty': True,
57 - 'name': str()}
125 +class HeartBeat(threading.Thread):
126 + def __init__(self, every):
127 + threading.Thread.__init__(self)
128 + self.daemon = True
129 + self.every = every
130
131 + def run(self):
132 + while True:
133 + time.sleep(self.every)
134 + if IS_ATTY:
135 + continue
136 + safe_print('\n')
137 +
138 +
139 +def load_module(name):
140 + abs_path = os.path.join(DIRS.modules, '{0}{1}'.format(name, MODULE_SUFFIX))
141 + return _load_module(name, abs_path)
142 +
143 +
144 +def multi_path_find(name, paths):
145 + for path in paths:
146 + abs_name = os.path.join(path, name)
147 + if os.path.isfile(abs_name):
148 + return abs_name
149 + return ''
150 +
151 +
152 +Task = collections.namedtuple(
153 + 'Task',
154 + [
155 + 'module_name',
156 + 'explicitly_enabled',
157 + ],
158 +)
159 +
160 +Result = collections.namedtuple(
161 + 'Result',
162 + [
163 + 'module_name',
164 + 'jobs_configs',
165 + ],
166 +)
167 +
168 +
169 +class ModuleChecker(multiprocessing.Process):
170 + def __init__(
171 + self,
172 + task_queue,
173 + result_queue,
174 + ):
175 + multiprocessing.Process.__init__(self)
176 + self.log = PythonDLogger()
177 + self.log.job_name = 'checker'
178 + self.task_queue = task_queue
179 + self.result_queue = result_queue
180 +
181 + def run(self):
182 + self.log.info('starting...')
183 + HeartBeat(1).start()
184 + while self.run_once():
185 + pass
186 + self.log.info('terminating...')
187 +
188 + def run_once(self):
189 + task = self.task_queue.get()
190 +
191 + if task is END_TASK_MARKER:
192 + self.task_queue.task_done()
193 + self.result_queue.put(END_TASK_MARKER)
194 + return False
195
60 -MODULE_EXTENSION = '.chart.py'
61 -OBSOLETE_MODULES = ['apache_cache', 'gunicorn_log', 'nginx_log', 'cpufreq', 'cpuidle', 'mdstat', 'linux_power_supply']
196 + result = self.do_task(task)
197 + if result:
198 + self.result_queue.put(result)
199 + self.task_queue.task_done()
200
201 + return True
202
64 -def module_ok(m):
65 - return m.endswith(MODULE_EXTENSION) and m[:-len(MODULE_EXTENSION)] not in OBSOLETE_MODULES
203 + def do_task(self, task):
204 + self.log.info("{0} : checking".format(task.module_name))
205
206 + # LOAD SOURCE
207 + module = Module(task.module_name)
208 + try:
209 + module.load_source()
210 + except Exception as error:
211 + self.log.warning("{0} : error on loading source : {1}, skipping module".format(
212 + task.module_name,
213 + error,
214 + ))
215 + return None
216 + else:
217 + self.log.info("{0} : source successfully loaded".format(task.module_name))
218
68 -ALL_MODULES = [m for m in sorted(os.listdir(PLUGINS_DIR)) if module_ok(m)]
219 + if module.is_disabled_by_default() and not task.explicitly_enabled:
220 + self.log.info("{0} : disabled by default".format(task.module_name))
221 + return None
222
223 + # LOAD CONFIG
224 + paths = [
225 + DIRS.modules_user_config,
226 + DIRS.modules_stock_config,
227 + ]
228
71 -def parse_cmd():
72 - debug = 'debug' in argv[1:]
73 - trace = 'trace' in argv[1:]
74 - override_update_every = next((arg for arg in argv[1:] if arg.isdigit() and int(arg) > 1), False)
75 - modules = [''.join([m, MODULE_EXTENSION]) for m in argv[1:] if ''.join([m, MODULE_EXTENSION]) in ALL_MODULES]
76 - return debug, trace, override_update_every, modules or ALL_MODULES
229 + conf_abs_path = multi_path_find(
230 + name='{0}.conf'.format(task.module_name),
231 + paths=paths,
232 + )
233
234 + if conf_abs_path:
235 + self.log.info("{0} : found config file '{1}'".format(task.module_name, conf_abs_path))
236 + try:
237 + module.load_config(conf_abs_path)
238 + except Exception as error:
239 + self.log.warning("{0} : error on loading config : {1}, skipping module".format(
240 + task.module_name, error))
241 + return None
242 + else:
243 + self.log.info("{0} : config was not found in '{1}', using default 1 job config".format(
244 + task.module_name, paths))
245 +
246 + # CHECK JOBS
247 + jobs = module.create_jobs()
248 + self.log.info("{0} : created {1} job(s) from the config".format(task.module_name, len(jobs)))
249 +
250 + successful_jobs_configs = list()
251 + for job in jobs:
252 + if job.autodetection_retry() > 0:
253 + successful_jobs_configs.append(job.config)
254 + self.log.info("{0}[{1}]: autodetection job, will be checked in main".format(task.module_name, job.name))
255 + continue
256
79 -def multi_job_check(config):
80 - return next((True for key in config if isinstance(config[key], dict)), False)
257 + try:
258 + job.init()
259 + except Exception as error:
260 + self.log.warning("{0}[{1}] : unhandled exception on init : {2}, skipping the job)".format(
261 + task.module_name, job.name, error))
262 + continue
263
264 + try:
265 + ok = job.check()
266 + except Exception as error:
267 + self.log.warning("{0}[{1}] : unhandled exception on check : {2}, skipping the job".format(
268 + task.module_name, job.name, error))
269 + continue
270
83 -class RawModule:
84 - def __init__(self, name, path, explicitly_enabled=True):
85 - self.name = name
86 - self.path = path
87 - self.explicitly_enabled = explicitly_enabled
88 -
89 -
90 -class Job(object):
91 - def __init__(self, initialized_job, job_id):
92 - """
93 - :param initialized_job: instance of <Class Service>
94 - :param job_id: <str>
95 - """
96 - self.job = initialized_job
97 - self.id = job_id # key in Modules.jobs()
98 - self.module_name = self.job.__module__ # used in Plugin.delete_job()
99 - self.recheck_every = self.job.configuration.pop('autodetection_retry')
100 - self.checked = False # used in Plugin.check_job()
101 - self.created = False # used in Plugin.create_job_charts()
102 - if self.job.update_every < int(OVERRIDE_UPDATE_EVERY):
103 - self.job.update_every = int(OVERRIDE_UPDATE_EVERY)
104 -
105 - def __getattr__(self, item):
106 - return getattr(self.job, item)
107 -
108 - def __repr__(self):
109 - return self.job.__repr__()
110 -
111 - def is_dead(self):
112 - return bool(self.ident) and not self.is_alive()
113 -
114 - def not_launched(self):
115 - return not bool(self.ident)
116 -
117 - def is_autodetect(self):
118 - return self.recheck_every
119 -
120 -
121 -class Module(object):
122 - def __init__(self, service, config):
123 - """
124 - :param service: <Module>
125 - :param config: <dict>
126 - """
271 + if not ok:
272 + self.log.info("{0}[{1}] : check failed, skipping the job".format(task.module_name, job.name))
273 + continue
274 +
275 + self.log.info("{0}[{1}] : check successful".format(task.module_name, job.name))
276 +
277 + job.config['autodetection_retry'] = job.config['update_every']
278 + successful_jobs_configs.append(job.config)
279 +
280 + if not successful_jobs_configs:
281 + self.log.info("{0} : all jobs failed, skipping module".format(task.module_name))
282 + return None
283 +
284 + return Result(module.source.__name__, successful_jobs_configs)
285 +
286 +
287 +class JobConf(OrderedDict):
288 + def __init__(self, *args):
289 + OrderedDict.__init__(self, *args)
290 +
291 + def set_defaults_from_module(self, module):
292 + for k in [k for k in JOB_BASE_CONF if hasattr(module, k)]:
293 + self[k] = getattr(module, k)
294 +
295 + def set_defaults_from_config(self, module_config):
296 + for k in [k for k in JOB_BASE_CONF if k in module_config]:
297 + self[k] = module_config[k]
298 +
299 + def set_job_name(self, name):
300 + self['job_name'] = re.sub(r'\s+', '_', name)
301 +
302 + def set_override_name(self, name):
303 + self['override_name'] = re.sub(r'\s+', '_', name)
304 +
305 + def as_dict(self):
306 + return copy.deepcopy(OrderedDict(self))
307 +
308 +
309 +class Job:
310 + def __init__(
311 + self,
312 + service,
313 + module_name,
314 + config,
315 + ):
316 self.service = service
128 - self.name = service.__name__
129 - self.config = self.jobs_configurations_builder(config)
130 - self.jobs = OrderedDict()
131 - self.counter = 1
317 + self.config = config
318 + self.module_name = module_name
319 + self.name = config['job_name']
320 + self.override_name = config['override_name']
321 + self.wrapped = None
322
133 - self.initialize_jobs()
323 + def init(self):
324 + self.wrapped = self.service(configuration=self.config.as_dict())
325
135 - def __repr__(self):
136 - return "<Class Module '{name}'>".format(name=self.name)
326 + def check(self):
327 + return self.wrapped.check()
328
138 - def __iter__(self):
139 - return iter(OrderedDict(self.jobs).values())
329 + def post_check(self, min_update_every):
330 + if self.wrapped.update_every < min_update_every:
331 + self.wrapped.update_every = min_update_every
332
141 - def __getitem__(self, item):
142 - return self.jobs[item]
333 + def create(self):
334 + return self.wrapped.create()
335
144 - def __delitem__(self, key):
145 - del self.jobs[key]
336 + def autodetection_retry(self):
337 + return self.config['autodetection_retry']
338
147 - def __len__(self):
148 - return len(self.jobs)
339 + def run(self):
340 + self.wrapped.run()
341
150 - def __bool__(self):
151 - return bool(self.jobs)
342
153 - def __nonzero__(self):
154 - return self.__bool__()
343 +class Module:
344 + def __init__(self, name):
345 + self.name = name
346 + self.source = None
347 + self.config = dict()
348
156 - def jobs_configurations_builder(self, config):
157 - """
158 - :param config: <dict>
159 - :return:
160 - """
161 - counter = 0
162 - job_base_config = dict()
349 + def is_disabled_by_default(self):
350 + return bool(getattr(self.source, 'disabled_by_default', False))
351
164 - for attr in BASE_CONFIG:
165 - job_base_config[attr] = config.pop(attr, getattr(self.service, attr, BASE_CONFIG[attr]))
352 + def load_source(self):
353 + self.source = load_module(self.name)
354
167 - if not config:
168 - config = {str(): dict()}
169 - elif not multi_job_check(config):
170 - config = {str(): config}
355 + def load_config(self, abs_path):
356 + self.config = load_config(abs_path) or dict()
357
172 - for job_name in config:
173 - if not isinstance(config[job_name], dict):
174 - continue
358 + def gather_jobs_configs(self):
359 + job_names = [v for v in self.config if isinstance(self.config[v], dict)]
360
176 - job_config = setdefault_values(config[job_name], base_dict=job_base_config)
177 - job_name = sub(r'\s+', '_', job_name)
178 - config[job_name]['name'] = sub(r'\s+', '_', config[job_name]['name'])
179 - counter += 1
180 - job_id = 'job' + str(counter).zfill(3)
361 + if len(job_names) == 0:
362 + job_conf = JobConf(JOB_BASE_CONF)
363 + job_conf.set_defaults_from_module(self.source)
364 + job_conf.update(self.config)
365 + job_conf.set_job_name(self.name)
366 + job_conf.set_override_name(job_conf.pop('name'))
367 + return [job_conf]
368
182 - yield job_id, job_name, job_config
369 + configs = list()
370 + for job_name in job_names:
371 + raw_job_conf = self.config[job_name]
372 + job_conf = JobConf(JOB_BASE_CONF)
373 + job_conf.set_defaults_from_module(self.source)
374 + job_conf.set_defaults_from_config(self.config)
375 + job_conf.update(raw_job_conf)
376 + job_conf.set_job_name(job_name)
377 + job_conf.set_override_name(job_conf.pop('name'))
378 + configs.append(job_conf)
379
184 - def initialize_jobs(self):
185 - """
186 - :return:
187 - """
188 - for job_id, job_name, job_config in self.config:
189 - job_config['job_name'] = job_name
190 - job_config['override_name'] = job_config.pop('name')
380 + return configs
381
192 - try:
193 - initialized_job = self.service.Service(configuration=job_config)
194 - except Exception as error:
195 - Logger.error("job initialization: '{module_name} {job_name}' "
196 - "=> ['FAILED'] ({error})".format(module_name=self.name,
197 - job_name=job_name,
198 - error=error))
199 - continue
200 - else:
201 - Logger.debug("job initialization: '{module_name} {job_name}' "
202 - "=> ['OK']".format(module_name=self.name,
203 - job_name=job_name or self.name))
204 - self.jobs[job_id] = Job(initialized_job=initialized_job,
205 - job_id=job_id)
206 - del self.config
207 - del self.service
208 -
209 -
210 -class Plugin(object):
211 - def __init__(self):
212 - self.loader = ModuleAndConfigLoader()
213 - self.modules = OrderedDict()
214 - self.sleep_time = 1
215 - self.runs_counter = 0
216 -
217 - user_config = os.path.join(USER_CONFIG_DIR, 'python.d.conf')
218 - stock_config = os.path.join(STOCK_CONFIG_DIR, 'python.d.conf')
219 -
220 - Logger.debug("loading '{0}'".format(user_config))
221 - self.config, error = self.loader.load_config_from_file(user_config)
222 -
223 - if error:
224 - Logger.error("cannot load '{0}': {1}. Will try stock version.".format(user_config, error))
225 - Logger.debug("loading '{0}'".format(stock_config))
226 - self.config, error = self.loader.load_config_from_file(stock_config)
227 - if error:
228 - Logger.error("cannot load '{0}': {1}".format(stock_config, error))
229 -
230 - self.do_gc = self.config.get("gc_run", GC_RUN)
231 - self.gc_interval = self.config.get("gc_interval", GC_COLLECT_EVERY)
232 -
233 - if not self.config.get('enabled', True):
234 - run_and_exit(Logger.info)('DISABLED in configuration file.')
235 -
236 - self.load_and_initialize_modules()
237 - if not self.modules:
238 - run_and_exit(Logger.info)('No modules to run. Exit...')
239 -
240 - def __iter__(self):
241 - return iter(OrderedDict(self.modules).values())
242 -
243 - @property
244 - def jobs(self):
245 - return (job for mod in self for job in mod)
246 -
247 - @property
248 - def dead_jobs(self):
249 - return (job for job in self.jobs if job.is_dead())
250 -
251 - @property
252 - def autodetect_jobs(self):
253 - return [job for job in self.jobs if job.not_launched()]
254 -
255 - def enabled_modules(self):
256 - for mod in MODULES_TO_RUN:
257 - mod_name = mod[:-len(MODULE_EXTENSION)]
258 - mod_path = os.path.join(PLUGINS_DIR, mod)
259 - if any(
260 - [
261 - self.config.get('default_run', True) and self.config.get(mod_name, True),
262 - (not self.config.get('default_run')) and self.config.get(mod_name),
263 - ]
264 - ):
265 - yield RawModule(
266 - name=mod_name,
267 - path=mod_path,
268 - explicitly_enabled=self.config.get(mod_name),
269 - )
270 -
271 - def load_and_initialize_modules(self):
272 - for mod in self.enabled_modules():
273 -
274 - # Load module from file ------------------------------------------------------------
275 - loaded_module, error = self.loader.load_module_from_file(mod.name, mod.path)
276 - log = Logger.error if error else Logger.debug
277 - log("module load source: '{module_name}' => [{status}]".format(status='FAILED' if error else 'OK',
278 - module_name=mod.name))
279 - if error:
280 - Logger.error("load source error : {0}".format(error))
281 - continue
382 + def create_jobs(self, jobs_conf=None):
383 + return [Job(self.source.Service, self.name, conf) for conf in jobs_conf or self.gather_jobs_configs()]
384
283 - # Load module config from file ------------------------------------------------------
284 - user_config = os.path.join(PLUGINS_USER_CONFIG_DIR, mod.name + '.conf')
285 - stock_config = os.path.join(PLUGINS_STOCK_CONFIG_DIR, mod.name + '.conf')
385
287 - Logger.debug("loading '{0}'".format(user_config))
288 - loaded_config, error = self.loader.load_config_from_file(user_config)
289 - if error:
290 - Logger.error("cannot load '{0}' : {1}. Will try stock version.".format(user_config, error))
291 - Logger.debug("loading '{0}'".format(stock_config))
292 - loaded_config, error = self.loader.load_config_from_file(stock_config)
386 +class JobRunner(threading.Thread):
387 + def __init__(self, job):
388 + threading.Thread.__init__(self)
389 + self.daemon = True
390 + self.wrapped = job
391
294 - if error:
295 - Logger.error("cannot load '{0}': {1}".format(stock_config, error))
392 + def run(self):
393 + self.wrapped.run()
394
297 - # Skip disabled modules
298 - if getattr(loaded_module, 'disabled_by_default', False) and not mod.explicitly_enabled:
299 - Logger.info("module '{0}' disabled by default".format(loaded_module.__name__))
300 - continue
395
302 - # Module initialization ---------------------------------------------------
396 +class PluginConf(dict):
397 + def __init__(self, *args):
398 + dict.__init__(self, *args)
399
304 - initialized_module = Module(service=loaded_module, config=loaded_config)
305 - Logger.debug("module status: '{module_name}' => [{status}] "
306 - "(jobs: {jobs_number})".format(status='OK' if initialized_module else 'FAILED',
307 - module_name=initialized_module.name,
308 - jobs_number=len(initialized_module)))
309 - if initialized_module:
310 - self.modules[initialized_module.name] = initialized_module
400 + def is_module_enabled(self, module_name, explicit):
401 + if module_name in self:
402 + return self[module_name]
403 + if explicit:
404 + return False
405 + return self['default_run']
406 +
407 +
408 +class Plugin:
409 + def __init__(
410 + self,
411 + min_update_every=1,
412 + modules_to_run=tuple(AVAILABLE_MODULES),
413 + ):
414 + self.log = PythonDLogger()
415 + self.config = PluginConf(PLUGIN_BASE_CONF)
416 + self.task_queue = multiprocessing.JoinableQueue()
417 + self.result_queue = multiprocessing.JoinableQueue()
418 + self.min_update_every = min_update_every
419 + self.modules_to_run = modules_to_run
420 + self.auto_detection_jobs = list()
421 + self.tasks = list()
422 + self.results = list()
423 + self.checked_jobs = collections.defaultdict(list)
424 + self.runs = 0
425
426 @staticmethod
313 - def check_job(job):
314 - """
315 - :param job: <Job>
316 - :return:
317 - """
318 - try:
319 - check_ok = bool(job.check())
320 - except Exception as error:
321 - job.error('check() unhandled exception: {error}'.format(error=error))
322 - return None
323 - else:
324 - return check_ok
427 + def shutdown():
428 + safe_print('DISABLE')
429 + exit(0)
430
326 - @staticmethod
327 - def create_job_charts(job):
328 - """
329 - :param job: <Job>
330 - :return:
331 - """
431 + def run(self):
432 + jobs = self.create_jobs()
433 + if not jobs:
434 + return
435 +
436 + for job in self.prepare_jobs(jobs):
437 + self.log.info('{0}[{1}] : started in thread'.format(job.module_name, job.name))
438 + JobRunner(job).start()
439 +
440 + self.serve()
441 +
442 + def enqueue_tasks(self):
443 + for task in self.tasks:
444 + self.task_queue.put(task)
445 + self.task_queue.put(END_TASK_MARKER)
446 +
447 + def dequeue_results(self):
448 + while True:
449 + result = self.result_queue.get()
450 + self.result_queue.task_done()
451 + if result is END_TASK_MARKER:
452 + break
453 + self.results.append(result)
454 +
455 + def load_config(self):
456 + paths = [
457 + DIRS.user_config,
458 + DIRS.stock_config,
459 + ]
460 +
461 + self.log.info("checking for config in {0}".format(paths))
462 + abs_path = multi_path_find(name=PLUGIN_CONF_FILE, paths=paths)
463 + if not abs_path:
464 + self.log.warning('config was not found, using defaults')
465 + return True
466 +
467 + self.log.info("config found, loading config '{0}'".format(abs_path))
468 try:
333 - create_ok = job.create()
469 + config = load_config(abs_path) or dict()
470 except Exception as error:
335 - job.error('create() unhandled exception: {error}'.format(error=error))
471 + self.log.error('error on loading config : {0}'.format(error))
472 return False
337 - else:
338 - return create_ok
339 -
340 - def delete_job(self, job):
341 - """
342 - :param job: <Job>
343 - :return:
344 - """
345 - del self.modules[job.module_name][job.id]
346 -
347 - def run_check(self):
348 - checked = list()
349 - for job in self.jobs:
350 - if job.name in checked:
351 - job.info('check() => [DROPPED] (already served by another job)')
352 - self.delete_job(job)
473 +
474 + self.log.info('config successfully loaded')
475 + self.config.update(config)
476 + return True
477 +
478 + def setup(self):
479 + self.log.info('starting setup')
480 + if not self.load_config():
481 + return False
482 +
483 + if not self.config['enabled']:
484 + self.log.info('disabled in configuration file')
485 + return False
486 +
487 + for mod in self.modules_to_run:
488 + if self.config.is_module_enabled(mod, False):
489 + task = Task(mod, self.config.is_module_enabled(mod, True))
490 + self.tasks.append(task)
491 + else:
492 + self.log.info("{0} : disabled in configuration file".format(mod))
493 +
494 + if not self.tasks:
495 + self.log.info('no modules to run')
496 + return False
497 +
498 + worker = ModuleChecker(self.task_queue, self.result_queue)
499 + self.log.info('starting checker process ({0} module(s) to check)'.format(len(self.tasks)))
500 + worker.start()
501 +
502 + # TODO: timeouts?
503 + self.enqueue_tasks()
504 + self.task_queue.join()
505 + self.dequeue_results()
506 + self.result_queue.join()
507 + self.log.info('stopping checker process')
508 + worker.join()
509 +
510 + if not self.results:
511 + self.log.info('no modules to run')
512 + return False
513 +
514 + self.log.info("setup complete, {0} active module(s) : '{1}'".format(
515 + len(self.results),
516 + [v.module_name for v in self.results])
517 + )
518 +
519 + return True
520 +
521 + def create_jobs(self):
522 + jobs = list()
523 + for result in self.results:
524 + module = Module(result.module_name)
525 + try:
526 + module.load_source()
527 + except Exception as error:
528 + self.log.warning("{0} : error on loading module source : {1}, skipping module".format(
529 + result.module_name, error))
530 continue
354 - ok = self.check_job(job)
355 - if ok:
356 - job.info('check() => [OK]')
357 - checked.append(job.name)
358 - job.checked = True
531 +
532 + module_jobs = module.create_jobs(result.jobs_configs)
533 + self.log.info("{0} : created {1} job(s)".format(module.name, len(module_jobs)))
534 + jobs.extend(module_jobs)
535 +
536 + return jobs
537 +
538 + def prepare_jobs(self, jobs):
539 + prepared = list()
540 +
541 + for job in jobs:
542 + check_name = job.override_name or job.name
543 + if check_name in self.checked_jobs[job.module_name]:
544 + self.log.info('{0}[{1}] : already served by another job, skipping the job'.format(
545 + job.module_name, job.name))
546 continue
360 - if not job.is_autodetect() or ok is None:
361 - job.info('check() => [FAILED]')
362 - self.delete_job(job)
363 - else:
364 - job.info('check() => [RECHECK] (autodetection_retry: {0})'.format(job.recheck_every))
547
366 - def run_create(self):
367 - for job in self.jobs:
368 - if not job.checked:
369 - # skip autodetection_retry jobs
548 + try:
549 + job.init()
550 + except Exception as error:
551 + self.log.warning("{0}[{1}] : unhandled exception on init : {2}, skipping the job".format(
552 + job.module_name, job.name, error))
553 + continue
554 +
555 + self.log.info("{0}[{1}] : init successful".format(job.module_name, job.name))
556 +
557 + try:
558 + ok = job.check()
559 + except Exception as error:
560 + self.log.warning("{0}[{1}] : unhandled exception on check : {2}, skipping the job".format(
561 + job.module_name, job.name, error))
562 continue
371 - ok = self.create_job_charts(job)
372 - if ok:
373 - job.debug('create() => [OK] (charts: {0})'.format(len(job.charts)))
374 - job.created = True
563 +
564 + if not ok:
565 + self.log.info('{0}[{1}] : check failed'.format(job.module_name, job.name))
566 + if job.autodetection_retry() > 0:
567 + self.log.info('{0}[{1}] : will recheck every {2} second(s)'.format(
568 + job.module_name, job.name, job.autodetection_retry()))
569 + self.auto_detection_jobs.append(job)
570 continue
376 - job.error('create() => [FAILED] (charts: {0})'.format(len(job.charts)))
377 - self.delete_job(job)
571
379 - def start(self):
380 - self.run_check()
381 - self.run_create()
382 - for job in self.jobs:
383 - if job.created:
384 - job.start()
572 + self.log.info('{0}[{1}] : check successful'.format(job.module_name, job.name))
573 +
574 + job.post_check(int(self.min_update_every))
575 +
576 + if not job.create():
577 + self.log.info('{0}[{1}] : create failed'.format(job.module_name, job.name))
578 +
579 + self.checked_jobs[job.module_name].append(check_name)
580 + prepared.append(job)
581 +
582 + return prepared
583 +
584 + def serve(self):
585 + gc_run = self.config['gc_run']
586 + gc_interval = self.config['gc_interval']
587
588 while True:
387 - if threading.active_count() <= 1 and not self.autodetect_jobs:
388 - run_and_exit(Logger.info)('FINISHED')
589 + self.runs += 1
590
390 - sleep(self.sleep_time)
391 - self.cleanup()
392 - self.autodetect_retry()
591 + if threading.active_count() <= 3 and not self.auto_detection_jobs:
592 + return
593
394 - # FIXME: https://github.com/netdata/netdata/issues/3817
395 - if self.do_gc and self.runs_counter % self.gc_interval == 0:
594 + time.sleep(1)
595 +
596 + if gc_run and self.runs % gc_interval == 0:
597 v = gc.collect()
397 - Logger.debug("GC full collection run result: {0}".format(v))
398 -
399 - # for exiting on SIGPIPE
400 - if not IS_ATTY:
401 - safe_print('\n')
402 -
403 - def cleanup(self):
404 - for job in self.dead_jobs:
405 - self.delete_job(job)
406 - for mod in self:
407 - if not mod:
408 - del self.modules[mod.name]
409 -
410 - def autodetect_retry(self):
411 - self.runs_counter += self.sleep_time
412 - for job in self.autodetect_jobs:
413 - if self.runs_counter % job.recheck_every == 0:
414 - checked = self.check_job(job)
415 - if checked:
416 - created = self.create_job_charts(job)
417 - if not created:
418 - self.delete_job(job)
419 - continue
420 - job.start()
598 + self.log.debug('GC collection run result: {0}'.format(v))
599 +
600 + self.auto_detection_jobs = [job for job in self.auto_detection_jobs if not self.retry_job(job)]
601 +
602 + def retry_job(self, job):
603 + stop_retrying = True
604 + retry_later = False
605 +
606 + if self.runs % job.autodetection_retry() != 0:
607 + return retry_later
608 +
609 + check_name = job.override_name or job.name
610 + if check_name in self.checked_jobs[job.module_name]:
611 + self.log.info("{0}[{1}]: already served by another job, give up on retrying".format(
612 + job.module_name, job.name))
613 + return stop_retrying
614 +
615 + try:
616 + ok = job.check()
617 + except Exception as error:
618 + self.log.warning("{0}[{1}] : unhandled exception on recheck : {2}, give up on retrying".format(
619 + job.module_name, job.name, error))
620 + return stop_retrying
621 +
622 + if not ok:
623 + self.log.info('{0}[{1}] : recheck failed, will retry in {2} second(s)'.format(
624 + job.module_name, job.name, job.autodetection_retry()))
625 + return retry_later
626 + self.log.info('{0}[{1}] : recheck successful'.format(job.module_name, job.name))
627 +
628 + if not job.create():
629 + return stop_retrying
630 +
631 + job.post_check(int(self.min_update_every))
632 + self.checked_jobs[job.module_name].append(check_name)
633 + return stop_retrying
634 +
635 +
636 +def parse_cmd():
637 + opts = sys.argv[:][1:]
638 + debug = False
639 + trace = False
640 + update_every = 1
641 + modules_to_run = list()
642 +
643 + v = next((opt for opt in opts if opt.isdigit() and int(opt) >= 1), None)
644 + if v:
645 + update_every = v
646 + opts.remove(v)
647 + if 'debug' in opts:
648 + debug = True
649 + opts.remove('debug')
650 + if 'trace' in opts:
651 + trace = True
652 + opts.remove('trace')
653 + if opts:
654 + modules_to_run = list(opts)
655 +
656 + return collections.namedtuple(
657 + 'CMD',
658 + [
659 + 'update_every',
660 + 'debug',
661 + 'trace',
662 + 'modules_to_run',
663 + ],
664 + )(
665 + update_every,
666 + debug,
667 + trace,
668 + modules_to_run,
669 + )
670 +
671 +
672 +def main():
673 + cmd = parse_cmd()
674 + logger = PythonDLogger()
675 +
676 + if cmd.debug:
677 + logger.logger.severity = 'DEBUG'
678 + if cmd.trace:
679 + logger.log_traceback = True
680 +
681 + logger.info('using python v{0}'.format(sys.version_info[:2][0]))
682 +
683 + unknown_modules = set(cmd.modules_to_run) - set(AVAILABLE_MODULES)
684 + if unknown_modules:
685 + logger.error('unknown modules : {0}'.format(sorted(list(unknown_modules))))
686 + safe_print('DISABLE')
687 + return
688 +
689 + plugin = Plugin(
690 + cmd.update_every,
691 + cmd.modules_to_run or AVAILABLE_MODULES,
692 + )
693 +
694 + HeartBeat(1).start()
695 +
696 + if not plugin.setup():
697 + safe_print('DISABLE')
698 + return
699 +
700 + plugin.run()
701 + logger.info('exiting from main...')
702 + plugin.shutdown()
703
704
705 if __name__ == '__main__':
424 - DEBUG, TRACE, OVERRIDE_UPDATE_EVERY, MODULES_TO_RUN = parse_cmd()
425 - Logger = PythonDLogger()
426 - if DEBUG:
427 - Logger.logger.severity = 'DEBUG'
428 - if TRACE:
429 - Logger.log_traceback = True
430 - Logger.info('Using python {version}'.format(version=PY_VERSION[0]))
431 -
432 - plugin = Plugin()
433 - plugin.start()
706 + main()
collectors/python.d.plugin/python_modules/bases/FrameworkServices/SimpleService.py
+3 -5
@@ -4,7 +4,7 @@
4 # Author: Ilya Mashchenko (l2isbad)
5 # SPDX-License-Identifier: GPL-3.0-or-later
6
7 -from threading import Thread
7 +
8 from time import sleep, time
9
10 from third_party.monotonic import monotonic
@@ -55,7 +55,7 @@ class RuntimeCounters:
55 self.penalty = round(min(self.retries * self.update_every / 2, MAX_PENALTY))
56
57
58 -class SimpleService(Thread, PythonDLimitedLogger, OldVersionCompatibility, object):
58 +class SimpleService(PythonDLimitedLogger, OldVersionCompatibility, object):
59 """
60 Prototype of Service class.
61 Implemented basic functionality to run jobs by `python.d.plugin`
@@ -65,8 +65,6 @@ class SimpleService(Thread, PythonDLimitedLogger, OldVersionCompatibility, objec
65 :param configuration: <dict>
66 :param name: <str>
67 """
68 - Thread.__init__(self)
69 - self.daemon = True
68 PythonDLimitedLogger.__init__(self)
69 OldVersionCompatibility.__init__(self)
70 self.configuration = configuration
@@ -91,7 +89,7 @@ class SimpleService(Thread, PythonDLimitedLogger, OldVersionCompatibility, objec
89
90 @property
91 def name(self):
94 - if self.job_name:
92 + if self.job_name and self.job_name != self.module_name:
93 return '_'.join([self.module_name, self.override_name or self.job_name])
94 return self.module_name
95
collectors/python.d.plugin/python_modules/bases/loaders.py
+20
@@ -81,3 +81,23 @@ class SourceLoader:
81
82 class ModuleAndConfigLoader(YamlOrderedLoader, SourceLoader):
83 pass
84 +
85 +
86 +def load_module(name, path):
87 + module = SourceFileLoader(name, path)
88 + if isinstance(module, types.ModuleType):
89 + return module
90 + return module.load_module()
91 +
92 +
93 +def load_yaml(stream):
94 + loader = YamlSafeLoader(stream)
95 + try:
96 + return loader.get_single_data()
97 + finally:
98 + loader.dispose()
99 +
100 +
101 +def load_config(file_name):
102 + with open(file_name, 'r') as stream:
103 + return load_yaml(stream)
collectors/python.d.plugin/python_modules/bases/loggers.py
+1 -1
@@ -26,7 +26,7 @@ LOGGING_LEVELS = {'CRITICAL': 50,
26 DEFAULT_LOG_LINE_FORMAT = '%(asctime)s: %(name)s %(levelname)s : %(message)s'
27 DEFAULT_LOG_TIME_FORMAT = '%Y-%m-%d %H:%M:%S'
28
29 -PYTHON_D_LOG_LINE_FORMAT = '%(asctime)s: %(name)s %(levelname)s: %(module_name)s: %(job_name)s: %(message)s'
29 +PYTHON_D_LOG_LINE_FORMAT = '%(asctime)s: %(name)s %(levelname)s: %(module_name)s[%(job_name)s] : %(message)s'
30 PYTHON_D_LOG_NAME = 'python.d'
31
32