@cryptotaxi247 / netdata-1 / commits / 35ff8fe45

pythond respect prev running jobs and refactor (#6661)

##### Summary Fixes: #6498 **Respect previously running jobs** - mantain active job statuses - periodically dump active job statuses to a file (in `NETDATA_LIB_DIR`) - read the file from start, if job was active on previous run then set (`recovering`): - autodetection_retry seconds: 30 - autodetection_retry retries: 10

Ilya Mashchenko committed Sep 11, 2019 at 17:06 UTC 35ff8fe456e163e1e075bd83c8e90465a9e243a4
1 file changed +527 -491
collectors/python.d.plugin/python.d.plugin.in
+527 -491
@@ -1,8 +1,5 @@
1 #!/usr/bin/env bash
2 '''':;
3 -if [[ "$OSTYPE" == "darwin"* ]]; then
4 - export OBJC_DISABLE_INITIALIZE_FORK_SAFETY=YES
5 -fi
3 exec "$(command -v python || command -v python3 || command -v python2 ||
4 echo "ERROR python IS NOT AVAILABLE IN THIS SYSTEM")" "$0" "$@" # '''
5
@@ -12,19 +9,25 @@ echo "ERROR python IS NOT AVAILABLE IN THIS SYSTEM")" "$0" "$@" # '''
9 # Author: Ilya Mashchenko (l2isbad)
10 # SPDX-License-Identifier: GPL-3.0-or-later
11
15 -
12 import collections
13 import copy
14 import gc
19 -import multiprocessing
15 +import json
16 import os
17 +import pprint
18 import re
19 import sys
20 import time
21 import threading
22 import types
23
27 -PY_VERSION = sys.version_info[:2]
24 +try:
25 + from queue import Queue
26 +except ImportError:
27 + from Queue import Queue
28 +
29 +PY_VERSION = sys.version_info[:2] # (major=3, minor=7, micro=3, releaselevel='final', serial=0)
30 +
31
32 if PY_VERSION > (3, 1):
33 from importlib.machinery import SourceFileLoader
@@ -35,98 +38,101 @@ else:
38 ENV_NETDATA_USER_CONFIG_DIR = 'NETDATA_USER_CONFIG_DIR'
39 ENV_NETDATA_STOCK_CONFIG_DIR = 'NETDATA_STOCK_CONFIG_DIR'
40 ENV_NETDATA_PLUGINS_DIR = 'NETDATA_PLUGINS_DIR'
41 +ENV_NETDATA_LIB_DIR = 'NETDATA_LIB_DIR'
42 ENV_NETDATA_UPDATE_EVERY = 'NETDATA_UPDATE_EVERY'
43
44
45 +def add_pythond_packages():
46 + pluginsd = os.getenv(ENV_NETDATA_PLUGINS_DIR, os.path.dirname(__file__))
47 + pythond = os.path.abspath(pluginsd + '/../python.d')
48 + packages = os.path.join(pythond, 'python_modules')
49 + sys.path.append(packages)
50 +
51 +
52 +add_pythond_packages()
53 +
54 +
55 +from bases.collection import safe_print
56 +from bases.loggers import PythonDLogger
57 +from bases.loaders import load_config
58 +
59 +try:
60 + from collections import OrderedDict
61 +except ImportError:
62 + from third_party.ordereddict import OrderedDict
63 +
64 +
65 def dirs():
42 - user_config = os.getenv(
66 + var_lib = os.getenv(
67 + ENV_NETDATA_LIB_DIR,
68 + '@varlibdir_POST@',
69 + )
70 + plugin_user_config = os.getenv(
71 ENV_NETDATA_USER_CONFIG_DIR,
72 '@configdir_POST@',
73 )
46 - stock_config = os.getenv(
74 + plugin_stock_config = os.getenv(
75 ENV_NETDATA_STOCK_CONFIG_DIR,
76 '@libconfigdir_POST@',
77 )
50 - modules_user_config = os.path.join(user_config, 'python.d')
51 - modules_stock_config = os.path.join(stock_config, 'python.d')
52 -
53 - modules = os.path.abspath(
54 - os.getenv(
55 - ENV_NETDATA_PLUGINS_DIR,
56 - os.path.dirname(__file__),
57 - ) + '/../python.d'
78 + pluginsd = os.getenv(
79 + ENV_NETDATA_PLUGINS_DIR,
80 + os.path.dirname(__file__),
81 )
59 - pythond_packages = os.path.join(modules, 'python_modules')
82 + modules_user_config = os.path.join(plugin_user_config, 'python.d')
83 + modules_stock_config = os.path.join(plugin_stock_config, 'python.d')
84 + modules = os.path.abspath(pluginsd + '/../python.d')
85
61 - return collections.namedtuple(
86 + Dirs = collections.namedtuple(
87 'Dirs',
88 [
64 - 'user_config',
65 - 'stock_config',
89 + 'plugin_user_config',
90 + 'plugin_stock_config',
91 'modules_user_config',
92 'modules_stock_config',
93 'modules',
69 - 'pythond_packages',
94 + 'var_lib',
95 ]
71 - )(
72 - user_config,
73 - stock_config,
96 + )
97 + return Dirs(
98 + plugin_user_config,
99 + plugin_stock_config,
100 modules_user_config,
101 modules_stock_config,
102 modules,
77 - pythond_packages,
103 + var_lib,
104 )
105
106
107 DIRS = dirs()
108
83 -sys.path.append(DIRS.pythond_packages)
84 -
85 -
86 -from bases.collection import safe_print
87 -from bases.loggers import PythonDLogger
88 -from bases.loaders import load_config
89 -
90 -try:
91 - from collections import OrderedDict
92 -except ImportError:
93 - from third_party.ordereddict import OrderedDict
94 -
95 -
96 -END_TASK_MARKER = None
97 -
109 IS_ATTY = sys.stdout.isatty()
110
100 -PLUGIN_CONF_FILE = 'python.d.conf'
101 -
111 MODULE_SUFFIX = '.chart.py'
112
104 -OBSOLETED_MODULES = (
105 - 'apache_cache', # replaced by web_log
106 - 'cpuidle', # rewritten in C
107 - 'cpufreq', # rewritten in C
108 - 'gunicorn_log', # replaced by web_log
109 - 'linux_power_supply', # rewritten in C
110 - 'nginx_log', # replaced by web_log
111 - 'mdstat', # rewritten in C
112 - 'sslcheck', # memory leak bug https://github.com/netdata/netdata/issues/5624
113 -)
113
114 +def available_modules():
115 + obsolete = (
116 + 'apache_cache', # replaced by web_log
117 + 'cpuidle', # rewritten in C
118 + 'cpufreq', # rewritten in C
119 + 'gunicorn_log', # replaced by web_log
120 + 'linux_power_supply', # rewritten in C
121 + 'nginx_log', # replaced by web_log
122 + 'mdstat', # rewritten in C
123 + 'sslcheck', # rewritten in Go, memory leak bug https://github.com/netdata/netdata/issues/5624
124 + )
125
116 -AVAILABLE_MODULES = [
117 - m[:-len(MODULE_SUFFIX)] for m in sorted(os.listdir(DIRS.modules))
118 - if m.endswith(MODULE_SUFFIX) and m[:-len(MODULE_SUFFIX)] not in OBSOLETED_MODULES
119 -]
126 + files = sorted(os.listdir(DIRS.modules))
127 + modules = [m[:-len(MODULE_SUFFIX)] for m in files if m.endswith(MODULE_SUFFIX)]
128 + avail = [m for m in modules if m not in obsolete]
129 + return tuple(avail)
130
121 -PLUGIN_BASE_CONF = {
122 - 'enabled': True,
123 - 'default_run': True,
124 - 'gc_run': True,
125 - 'gc_interval': 300,
126 -}
131 +
132 +AVAILABLE_MODULES = available_modules()
133
134 JOB_BASE_CONF = {
129 - 'update_every': os.getenv(ENV_NETDATA_UPDATE_EVERY, 1),
135 + 'update_every': int(os.getenv(ENV_NETDATA_UPDATE_EVERY, 1)),
136 'priority': 60000,
137 'autodetection_retry': 0,
138 'chart_cleanup': 10,
@@ -134,23 +140,20 @@ JOB_BASE_CONF = {
140 'name': str(),
141 }
142
143 +PLUGIN_BASE_CONF = {
144 + 'enabled': True,
145 + 'default_run': True,
146 + 'gc_run': True,
147 + 'gc_interval': 300,
148 +}
149
138 -def heartbeat():
139 - if IS_ATTY:
140 - return
141 - safe_print('\n')
142 -
143 -
144 -class HeartBeat(threading.Thread):
145 - def __init__(self, every):
146 - threading.Thread.__init__(self)
147 - self.daemon = True
148 - self.every = every
150
150 - def run(self):
151 - while True:
152 - time.sleep(self.every)
153 - heartbeat()
151 +def multi_path_find(name, *paths):
152 + for path in paths:
153 + abs_name = os.path.join(path, name)
154 + if os.path.isfile(abs_name):
155 + return abs_name
156 + return str()
157
158
159 def load_module(name):
@@ -161,265 +164,268 @@ def load_module(name):
164 return module.load_module()
165
166
164 -def multi_path_find(name, paths):
165 - for path in paths:
166 - abs_name = os.path.join(path, name)
167 - if os.path.isfile(abs_name):
168 - return abs_name
169 - return ''
170 -
171 -
172 -Task = collections.namedtuple(
173 - 'Task',
174 - [
175 - 'module_name',
176 - 'explicitly_enabled',
177 - ],
178 -)
179 -
180 -Result = collections.namedtuple(
181 - 'Result',
182 - [
183 - 'module_name',
184 - 'jobs_configs',
185 - ],
186 -)
187 -
188 -
189 -class ModuleChecker(multiprocessing.Process):
190 - def __init__(
191 - self,
192 - task_queue,
193 - result_queue,
194 - ):
195 - multiprocessing.Process.__init__(self)
196 - self.log = PythonDLogger()
197 - self.log.job_name = 'checker'
198 - self.task_queue = task_queue
199 - self.result_queue = result_queue
200 -
201 - def run(self):
202 - self.log.info('starting...')
203 - HeartBeat(1).start()
204 - while self.run_once():
205 - pass
206 - self.log.info('terminating...')
167 +class ModuleConfig:
168 + def __init__(self, name, config=None):
169 + self.name = name
170 + self.config = config or OrderedDict()
171
208 - def run_once(self):
209 - task = self.task_queue.get()
172 + def load(self, abs_path):
173 + self.config.update(load_config(abs_path) or dict())
174
211 - if task is END_TASK_MARKER:
212 - # TODO: find better solution, understand why heartbeat thread doesn't work
213 - heartbeat()
214 - self.task_queue.task_done()
215 - self.result_queue.put(END_TASK_MARKER)
216 - return False
175 + def defaults(self):
176 + keys = (
177 + 'update_every',
178 + 'priority',
179 + 'autodetection_retry',
180 + 'chart_cleanup',
181 + 'penalty',
182 + )
183 + return dict((k, self.config[k]) for k in keys if k in self.config)
184
218 - result = self.do_task(task)
219 - if result:
220 - self.result_queue.put(result)
221 - self.task_queue.task_done()
185 + def create_job(self, job_name, job_config=None):
186 + job_config = job_config or dict()
187
223 - return True
188 + config = OrderedDict()
189 + config.update(job_config)
190 + config['job_name'] = job_name
191 + for k, v in self.defaults().items():
192 + config.setdefault(k, v)
193
225 - def do_task(self, task):
226 - self.log.info("{0} : checking".format(task.module_name))
194 + return config
195
228 - # LOAD SOURCE
229 - module = Module(task.module_name)
230 - try:
231 - module.load_source()
232 - except Exception as error:
233 - self.log.warning("{0} : error on loading source : {1}, skipping module".format(
234 - task.module_name,
235 - error,
236 - ))
237 - return None
238 - else:
239 - self.log.info("{0} : source successfully loaded".format(task.module_name))
196 + def job_names(self):
197 + return [v for v in self.config if isinstance(self.config.get(v), dict)]
198
241 - if module.is_disabled_by_default() and not task.explicitly_enabled:
242 - self.log.info("{0} : disabled by default".format(task.module_name))
243 - return None
199 + def single_job(self):
200 + return [self.create_job(self.name)]
201
245 - # LOAD CONFIG
246 - paths = [
247 - DIRS.modules_user_config,
248 - DIRS.modules_stock_config,
249 - ]
202 + def multi_job(self):
203 + return [self.create_job(n, self.config[n]) for n in self.job_names()]
204
251 - conf_abs_path = multi_path_find(
252 - name='{0}.conf'.format(task.module_name),
253 - paths=paths,
254 - )
205 + def create_jobs(self):
206 + return self.multi_job() or self.single_job()
207
256 - if conf_abs_path:
257 - self.log.info("{0} : found config file '{1}'".format(task.module_name, conf_abs_path))
258 - try:
259 - module.load_config(conf_abs_path)
260 - except Exception as error:
261 - self.log.warning("{0} : error on loading config : {1}, skipping module".format(
262 - task.module_name, error))
263 - return None
264 - else:
265 - self.log.info("{0} : config was not found in '{1}', using default 1 job config".format(
266 - task.module_name, paths))
267 -
268 - # CHECK JOBS
269 - jobs = module.create_jobs()
270 - self.log.info("{0} : created {1} job(s) from the config".format(task.module_name, len(jobs)))
271 -
272 - successful_jobs_configs = list()
273 - for job in jobs:
274 - if job.autodetection_retry() > 0:
275 - successful_jobs_configs.append(job.config)
276 - self.log.info("{0}[{1}]: autodetection job, will be checked in main".format(task.module_name, job.name))
277 - continue
208
279 - try:
280 - job.init()
281 - except Exception as error:
282 - self.log.warning("{0}[{1}] : unhandled exception on init : {2}, skipping the job)".format(
283 - task.module_name, job.name, error))
284 - continue
209 +class JobsConfigsBuilder:
210 + def __init__(self, config_dirs):
211 + self.config_dirs = config_dirs
212 + self.log = PythonDLogger()
213 + self.job_defaults = None
214 + self.module_defaults = None
215 + self.min_update_every = None
216
286 - try:
287 - ok = job.check()
288 - except Exception as error:
289 - self.log.warning("{0}[{1}] : unhandled exception on check : {2}, skipping the job".format(
290 - task.module_name, job.name, error))
291 - continue
217 + def load_module_config(self, module_name):
218 + name = '{0}.conf'.format(module_name)
219 + self.log.debug("[{0}] looking for '{1}' in {2}".format(module_name, name, self.config_dirs))
220 + config = ModuleConfig(module_name)
221
293 - if not ok:
294 - self.log.info("{0}[{1}] : check failed, skipping the job".format(task.module_name, job.name))
295 - continue
222 + abs_path = multi_path_find(name, *self.config_dirs)
223 + if not abs_path:
224 + self.log.warning("[{0}] '{1}' was not found".format(module_name, name))
225 + return config
226
297 - self.log.info("{0}[{1}] : check successful".format(task.module_name, job.name))
227 + self.log.debug("[{0}] loading '{1}'".format(module_name, abs_path))
228 + try:
229 + config.load(abs_path)
230 + except Exception as error:
231 + self.log.error("[{0}] error on loading '{1}' : {2}".format(module_name, abs_path, repr(error)))
232 + return None
233
299 - job.config['autodetection_retry'] = job.config['update_every']
300 - successful_jobs_configs.append(job.config)
234 + self.log.debug("[{0}] '{1}' is loaded".format(module_name, abs_path))
235 + return config
236
302 - if not successful_jobs_configs:
303 - self.log.info("{0} : all jobs failed, skipping module".format(task.module_name))
304 - return None
237 + @staticmethod
238 + def apply_defaults(jobs, defaults):
239 + if defaults is None:
240 + return
241 + for k, v in defaults.items():
242 + for job in jobs:
243 + job.setdefault(k, v)
244
306 - return Result(module.source.__name__, successful_jobs_configs)
245 + def set_min_update_every(self, jobs, min_update_every):
246 + if min_update_every is None:
247 + return
248 + for job in jobs:
249 + if 'update_every' in job and job['update_every'] < self.min_update_every:
250 + job['update_every'] = self.min_update_every
251
252 + def build(self, module_name):
253 + config = self.load_module_config(module_name)
254 + if config is None:
255 + return None
256
309 -class JobConf(OrderedDict):
310 - def __init__(self, *args):
311 - OrderedDict.__init__(self, *args)
257 + configs = config.create_jobs()
258 + self.log.info("[{0}] built {1} job(s) configs".format(module_name, len(configs)))
259
313 - def set_defaults_from_module(self, module):
314 - for k in [k for k in JOB_BASE_CONF if hasattr(module, k)]:
315 - self[k] = getattr(module, k)
260 + self.apply_defaults(configs, self.module_defaults)
261 + self.apply_defaults(configs, self.job_defaults)
262 + self.set_min_update_every(configs, self.min_update_every)
263
317 - def set_defaults_from_config(self, module_config):
318 - for k in [k for k in JOB_BASE_CONF if k in module_config]:
319 - self[k] = module_config[k]
264 + return configs
265
321 - def set_job_name(self, name):
322 - self['job_name'] = re.sub(r'\s+', '_', name)
266
324 - def set_override_name(self, name):
325 - self['override_name'] = re.sub(r'\s+', '_', name)
267 +JOB_STATUS_ACTIVE = 'active'
268 +JOB_STATUS_RECOVERING = 'recovering'
269 +JOB_STATUS_DROPPED = 'dropped'
270 +JOB_STATUS_INIT = 'initial'
271
327 - def as_dict(self):
328 - return copy.deepcopy(OrderedDict(self))
272
273 +class Job(threading.Thread):
274 + inf = -1
275
331 -class Job:
332 - def __init__(
333 - self,
334 - service,
335 - module_name,
336 - config,
337 - ):
276 + def __init__(self, service, module_name, config):
277 + threading.Thread.__init__(self)
278 + self.daemon = True
279 self.service = service
339 - self.config = config
280 self.module_name = module_name
341 - self.name = config['job_name']
342 - self.override_name = config['override_name']
343 - self.wrapped = None
281 + self.config = config
282 + self.real_name = config['job_name']
283 + self.actual_name = config['override_name'] or self.real_name
284 + self.autodetection_retry = config['autodetection_retry']
285 + self.checks = self.inf
286 + self.job = None
287 + self.status = JOB_STATUS_INIT
288 +
289 + def is_inited(self):
290 + return self.job is not None
291
292 def init(self):
346 - self.wrapped = self.service(configuration=self.config.as_dict())
293 + self.job = self.service(configuration=copy.deepcopy(self.config))
294
295 def check(self):
349 - return self.wrapped.check()
350 -
351 - def post_check(self, min_update_every):
352 - if self.wrapped.update_every < min_update_every:
353 - self.wrapped.update_every = min_update_every
296 + ok = self.job.check()
297 + self.checks -= self.checks != self.inf and not ok
298 + return ok
299
300 def create(self):
356 - return self.wrapped.create()
301 + self.job.create()
302
358 - def autodetection_retry(self):
359 - return self.config['autodetection_retry']
303 + def need_to_recheck(self):
304 + return self.autodetection_retry != 0 and self.checks != 0
305
306 def run(self):
362 - self.wrapped.run()
307 + self.job.run()
308
309
365 -class Module:
310 +class ModuleSrc:
311 def __init__(self, name):
312 self.name = name
368 - self.source = None
369 - self.config = dict()
313 + self.src = None
314 +
315 + def load(self):
316 + self.src = load_module(self.name)
317 +
318 + def get(self, key):
319 + return getattr(self.src, key, None)
320 +
321 + def service(self):
322 + return self.get('Service')
323 +
324 + def defaults(self):
325 + keys = (
326 + 'update_every',
327 + 'priority',
328 + 'autodetection_retry',
329 + 'chart_cleanup',
330 + 'penalty',
331 + )
332 + return dict((k, self.get(k)) for k in keys if self.get(k) is not None)
333
334 def is_disabled_by_default(self):
372 - return bool(getattr(self.source, 'disabled_by_default', False))
373 -
374 - def load_source(self):
375 - self.source = load_module(self.name)
376 -
377 - def load_config(self, abs_path):
378 - self.config = load_config(abs_path) or dict()
379 -
380 - def gather_jobs_configs(self):
381 - job_names = [v for v in self.config if isinstance(self.config[v], dict)]
382 -
383 - if len(job_names) == 0:
384 - job_conf = JobConf(JOB_BASE_CONF)
385 - job_conf.set_defaults_from_module(self.source)
386 - job_conf.update(self.config)
387 - job_conf.set_job_name(self.name)
388 - job_conf.set_override_name(job_conf.pop('name'))
389 - return [job_conf]
390 -
391 - configs = list()
392 - for job_name in job_names:
393 - raw_job_conf = self.config[job_name]
394 - job_conf = JobConf(JOB_BASE_CONF)
395 - job_conf.set_defaults_from_module(self.source)
396 - job_conf.set_defaults_from_config(self.config)
397 - job_conf.update(raw_job_conf)
398 - job_conf.set_job_name(job_name)
399 - job_conf.set_override_name(job_conf.pop('name'))
400 - configs.append(job_conf)
335 + return bool(self.get('disabled_by_default'))
336
402 - return configs
337
404 - def create_jobs(self, jobs_conf=None):
405 - return [Job(self.source.Service, self.name, conf) for conf in jobs_conf or self.gather_jobs_configs()]
338 +class JobsStatuses:
339 + def __init__(self):
340 + self.items = OrderedDict()
341
342 + def dump(self):
343 + return json.dumps(self.items, indent=2)
344
408 -class JobRunner(threading.Thread):
409 - def __init__(self, job):
410 - threading.Thread.__init__(self)
411 - self.daemon = True
412 - self.wrapped = job
345 + def get(self, module_name, job_name):
346 + if module_name not in self.items:
347 + return None
348 + return self.items[module_name].get(job_name)
349
414 - def run(self):
415 - self.wrapped.run()
350 + def has(self, module_name, job_name):
351 + return self.get(module_name, job_name) is not None
352
353 + def from_file(self, path):
354 + with open(path) as f:
355 + data = json.load(f)
356 + return self.from_json(data)
357
418 -class PluginConf(dict):
358 + @staticmethod
359 + def from_json(items):
360 + if not isinstance(items, dict):
361 + raise Exception('items obj has wrong type : {0}'.format(type(items)))
362 + if not items:
363 + return JobsStatuses()
364 +
365 + v = OrderedDict()
366 + for mod_name in sorted(items):
367 + if not items[mod_name]:
368 + continue
369 + v[mod_name] = OrderedDict()
370 + for job_name in sorted(items[mod_name]):
371 + v[mod_name][job_name] = items[mod_name][job_name]
372 +
373 + rv = JobsStatuses()
374 + rv.items = v
375 + return rv
376 +
377 + @staticmethod
378 + def from_jobs(jobs):
379 + v = OrderedDict()
380 + for job in jobs:
381 + status = job.status
382 + if status not in (JOB_STATUS_ACTIVE, JOB_STATUS_RECOVERING):
383 + continue
384 + if job.module_name not in v:
385 + v[job.module_name] = OrderedDict()
386 + v[job.module_name][job.real_name] = status
387 +
388 + rv = JobsStatuses()
389 + rv.items = v
390 + return rv
391 +
392 +
393 +class StdoutSaver:
394 + @staticmethod
395 + def save(dump):
396 + print(dump)
397 +
398 +
399 +class CachedFileSaver:
400 + def __init__(self, path):
401 + self.last_save_success = False
402 + self.last_saved_dump = str()
403 + self.path = path
404 +
405 + def save(self, dump):
406 + if self.last_save_success and self.last_saved_dump == dump:
407 + return
408 + try:
409 + with open(self.path, 'w') as out:
410 + out.write(dump)
411 + except Exception:
412 + self.last_save_success = False
413 + raise
414 + self.last_saved_dump = dump
415 + self.last_save_success = True
416 +
417 +
418 +class PluginConfig(dict):
419 def __init__(self, *args):
420 dict.__init__(self, *args)
421
422 - def is_module_enabled(self, module_name, explicit):
422 + def is_module_explicitly_enabled(self, module_name):
423 + return self._is_module_enabled(module_name, True)
424 +
425 + def is_module_enabled(self, module_name):
426 + return self._is_module_enabled(module_name, False)
427 +
428 + def _is_module_enabled(self, module_name, explicit):
429 if module_name in self:
430 return self[module_name]
431 if explicit:
@@ -428,249 +434,253 @@ class PluginConf(dict):
434
435
436 class Plugin:
431 - def __init__(
432 - self,
433 - min_update_every=1,
434 - modules_to_run=tuple(AVAILABLE_MODULES),
435 - ):
436 - self.log = PythonDLogger()
437 - self.config = PluginConf(PLUGIN_BASE_CONF)
438 - self.task_queue = multiprocessing.JoinableQueue()
439 - self.result_queue = multiprocessing.JoinableQueue()
440 - self.min_update_every = min_update_every
437 + config_name = 'python.d.conf'
438 + jobs_status_dump_name = 'pythond-jobs-statuses.json'
439 +
440 + def __init__(self, modules_to_run, min_update_every):
441 self.modules_to_run = modules_to_run
442 - self.auto_detection_jobs = list()
443 - self.tasks = list()
444 - self.results = list()
445 - self.checked_jobs = collections.defaultdict(list)
442 + self.min_update_every = min_update_every
443 + self.config = PluginConfig(PLUGIN_BASE_CONF)
444 + self.log = PythonDLogger()
445 + self.started_jobs = collections.defaultdict(dict)
446 + self.jobs = list()
447 + self.saver = None
448 self.runs = 0
449
448 - @staticmethod
449 - def shutdown():
450 - safe_print('DISABLE')
451 - exit(0)
452 -
453 - def run(self):
454 - jobs = self.create_jobs()
455 - if not jobs:
456 - return
457 -
458 - for job in self.prepare_jobs(jobs):
459 - self.log.info('{0}[{1}] : started in thread'.format(job.module_name, job.name))
460 - JobRunner(job).start()
461 -
462 - self.serve()
463 -
464 - def enqueue_tasks(self):
465 - for task in self.tasks:
466 - self.task_queue.put(task)
467 - self.task_queue.put(END_TASK_MARKER)
468 -
469 - def dequeue_results(self):
470 - while True:
471 - result = self.result_queue.get()
472 - self.result_queue.task_done()
473 - if result is END_TASK_MARKER:
474 - break
475 - self.results.append(result)
476 -
450 def load_config(self):
451 paths = [
479 - DIRS.user_config,
480 - DIRS.stock_config,
452 + DIRS.plugin_user_config,
453 + DIRS.plugin_stock_config,
454 ]
482 -
483 - self.log.info("checking for config in {0}".format(paths))
484 - abs_path = multi_path_find(name=PLUGIN_CONF_FILE, paths=paths)
455 + self.log.debug("looking for '{0}' in {1}".format(self.config_name, paths))
456 + abs_path = multi_path_find(self.config_name, *paths)
457 if not abs_path:
486 - self.log.warning('config was not found, using defaults')
458 + self.log.warning("'{0}' was not found, using defaults".format(self.config_name))
459 return True
460
489 - self.log.info("config found, loading config '{0}'".format(abs_path))
461 + self.log.debug("loading '{0}'".format(abs_path))
462 try:
491 - config = load_config(abs_path) or dict()
463 + config = load_config(abs_path)
464 except Exception as error:
493 - self.log.error('error on loading config : {0}'.format(error))
465 + self.log.error("error on loading '{0}' : {1}".format(abs_path, repr(error)))
466 return False
467
496 - self.log.info('config successfully loaded')
468 + self.log.debug("'{0}' is loaded".format(abs_path))
469 self.config.update(config)
470 return True
471
500 - def setup(self):
501 - self.log.info('starting setup')
502 - if not self.load_config():
503 - return False
504 -
505 - if not self.config['enabled']:
506 - self.log.info('disabled in configuration file')
507 - return False
508 -
509 - for mod in self.modules_to_run:
510 - if self.config.is_module_enabled(mod, False):
511 - task = Task(mod, self.config.is_module_enabled(mod, True))
512 - self.tasks.append(task)
513 - else:
514 - self.log.info("{0} : disabled in configuration file".format(mod))
515 -
516 - if not self.tasks:
517 - self.log.info('no modules to run')
518 - return False
472 + def load_job_statuses(self):
473 + self.log.debug("looking for '{0}' in {1}".format(self.jobs_status_dump_name, DIRS.var_lib))
474 + abs_path = multi_path_find(self.jobs_status_dump_name, DIRS.var_lib)
475 + if not abs_path:
476 + self.log.warning("'{0}' was not found".format(self.jobs_status_dump_name))
477 + return
478
520 - worker = ModuleChecker(self.task_queue, self.result_queue)
521 - self.log.info('starting checker process ({0} module(s) to check)'.format(len(self.tasks)))
522 - worker.start()
523 -
524 - # TODO: timeouts?
525 - self.enqueue_tasks()
526 - self.task_queue.join()
527 - self.dequeue_results()
528 - self.result_queue.join()
529 - self.task_queue.close()
530 - self.result_queue.close()
531 - self.log.info('stopping checker process')
532 - worker.join()
533 -
534 - if not self.results:
535 - self.log.info('no modules to run')
536 - return False
479 + self.log.debug("loading '{0}'".format(abs_path))
480 + try:
481 + statuses = JobsStatuses().from_file(abs_path)
482 + except Exception as error:
483 + self.log.warning("error on loading '{0}' : {1}".format(abs_path, repr(error)))
484 + return None
485 + self.log.debug("'{0}' is loaded".format(abs_path))
486 + return statuses
487
538 - self.log.info("setup complete, {0} active module(s) : '{1}'".format(
539 - len(self.results),
540 - [v.module_name for v in self.results])
541 - )
488 + def create_jobs(self, job_statuses=None):
489 + paths = [
490 + DIRS.modules_user_config,
491 + DIRS.modules_stock_config,
492 + ]
493
543 - return True
494 + builder = JobsConfigsBuilder(paths)
495 + builder.job_defaults = JOB_BASE_CONF
496 + builder.min_update_every = self.min_update_every
497
545 - def create_jobs(self):
498 jobs = list()
547 - for result in self.results:
548 - module = Module(result.module_name)
499 + for mod_name in self.modules_to_run:
500 + if not self.config.is_module_enabled(mod_name):
501 + self.log.info("[{0}] is disabled in the configuration file, skipping it".format(mod_name))
502 + continue
503 +
504 + src = ModuleSrc(mod_name)
505 try:
550 - module.load_source()
506 + src.load()
507 except Exception as error:
552 - self.log.warning("{0} : error on loading module source : {1}, skipping module".format(
553 - result.module_name, error))
508 + self.log.warning("[{0}] error on loading source : {1}, skipping it".format(mod_name, repr(error)))
509 continue
510
556 - module_jobs = module.create_jobs(result.jobs_configs)
557 - self.log.info("{0} : created {1} job(s)".format(module.name, len(module_jobs)))
558 - jobs.extend(module_jobs)
559 -
560 - return jobs
561 -
562 - def prepare_jobs(self, jobs):
563 - prepared = list()
511 + if not (src.service() and callable(src.service())):
512 + self.log.warning("[{0}] has no callable Service object, skipping it".format(mod_name))
513 + continue
514
565 - for job in jobs:
566 - check_name = job.override_name or job.name
567 - if check_name in self.checked_jobs[job.module_name]:
568 - self.log.info('{0}[{1}] : already served by another job, skipping the job'.format(
569 - job.module_name, job.name))
515 + if src.is_disabled_by_default() and not self.config.is_module_explicitly_enabled(mod_name):
516 + self.log.info("[{0}] is disabled by default, skipping it".format(mod_name))
517 continue
518
572 - try:
573 - job.init()
574 - except Exception as error:
575 - self.log.warning("{0}[{1}] : unhandled exception on init : {2}, skipping the job".format(
576 - job.module_name, job.name, error))
519 + builder.module_defaults = src.defaults()
520 + configs = builder.build(mod_name)
521 + if not configs:
522 + self.log.info("[{0}] has no job configs, skipping it".format(mod_name))
523 continue
524
579 - self.log.info("{0}[{1}] : init successful".format(job.module_name, job.name))
525 + for config in configs:
526 + config['job_name'] = re.sub(r'\s+', '_', config['job_name'])
527 + config['override_name'] = re.sub(r'\s+', '_', config.pop('name'))
528
581 - try:
582 - ok = job.check()
583 - except Exception as error:
584 - self.log.warning("{0}[{1}] : unhandled exception on check : {2}, skipping the job".format(
585 - job.module_name, job.name, error))
586 - continue
529 + job = Job(src.service(), mod_name, config)
530
588 - if not ok:
589 - self.log.info('{0}[{1}] : check failed'.format(job.module_name, job.name))
590 - if job.autodetection_retry() > 0:
591 - self.log.info('{0}[{1}] : will recheck every {2} second(s)'.format(
592 - job.module_name, job.name, job.autodetection_retry()))
593 - self.auto_detection_jobs.append(job)
594 - continue
531 + was_previously_active = job_statuses and job_statuses.has(job.module_name, job.real_name)
532 + if was_previously_active and job.autodetection_retry == 0:
533 + self.log.debug('{0}[{1}] was previously active, applying recovering settings'.format(
534 + job.module_name, job.real_name))
535 + job.checks = 11
536 + job.autodetection_retry = 30
537
596 - self.log.info('{0}[{1}] : check successful'.format(job.module_name, job.name))
538 + jobs.append(job)
539
598 - job.post_check(int(self.min_update_every))
540 + return jobs
541
600 - if not job.create():
601 - self.log.info('{0}[{1}] : create failed'.format(job.module_name, job.name))
542 + def setup(self):
543 + if not self.load_config():
544 + return False
545
603 - self.checked_jobs[job.module_name].append(check_name)
604 - prepared.append(job)
546 + if not self.config['enabled']:
547 + self.log.info('disabled in the configuration file')
548 + return False
549
606 - return prepared
550 + statuses = self.load_job_statuses()
551
608 - def serve(self):
609 - gc_run = self.config['gc_run']
610 - gc_interval = self.config['gc_interval']
552 + self.jobs = self.create_jobs(statuses)
553 + if not self.jobs:
554 + self.log.info('no jobs to run')
555 + return False
556 +
557 + if not IS_ATTY:
558 + abs_path = os.path.join(DIRS.var_lib, self.jobs_status_dump_name)
559 + self.saver = CachedFileSaver(abs_path)
560 + return True
561
612 - while True:
613 - self.runs += 1
562 + def start_jobs(self, *jobs):
563 + for job in jobs:
564 + if job.status not in (JOB_STATUS_INIT, JOB_STATUS_RECOVERING):
565 + continue
566
615 - # threads: main + heartbeat
616 - if threading.active_count() <= 2 and not self.auto_detection_jobs:
617 - return
567 + if job.actual_name in self.started_jobs[job.module_name]:
568 + self.log.info('{0}[{1}] : already served by another job, skipping it'.format(
569 + job.module_name, job.real_name))
570 + job.status = JOB_STATUS_DROPPED
571 + continue
572
619 - time.sleep(1)
573 + if not job.is_inited():
574 + try:
575 + job.init()
576 + except Exception as error:
577 + self.log.warning("{0}[{1}] : unhandled exception on init : {2}, skipping the job",
578 + job.module_name, job.real_name, repr(error))
579 + job.status = JOB_STATUS_DROPPED
580 + continue
581
621 - if gc_run and self.runs % gc_interval == 0:
622 - v = gc.collect()
623 - self.log.debug('GC collection run result: {0}'.format(v))
582 + try:
583 + ok = job.check()
584 + except Exception as error:
585 + self.log.warning("{0}[{1}] : unhandled exception on check : {2}, skipping the job",
586 + job.module_name, job.real_name, repr(error))
587 + job.status = JOB_STATUS_DROPPED
588 + continue
589 + if not ok:
590 + self.log.info('{0}[{1}] : check failed'.format(job.module_name, job.real_name))
591 + job.status = JOB_STATUS_RECOVERING if job.need_to_recheck() else JOB_STATUS_DROPPED
592 + continue
593 + self.log.info('{0}[{1}] : check success'.format(job.module_name, job.real_name))
594 +
595 + try:
596 + job.create()
597 + except Exception as error:
598 + self.log.error("{0}[{1}] : unhandled exception on create : {2}, skipping the job",
599 + job.module_name, job.real_name, repr(error))
600 + job.status = JOB_STATUS_DROPPED
601 + continue
602
625 - self.auto_detection_jobs = [job for job in self.auto_detection_jobs if not self.retry_job(job)]
603 + self.started_jobs[job.module_name] = job.actual_name
604 + job.status = JOB_STATUS_ACTIVE
605 + job.start()
606
627 - def retry_job(self, job):
628 - stop_retrying = True
629 - retry_later = False
607 + @staticmethod
608 + def keep_alive():
609 + if not IS_ATTY:
610 + safe_print('\n')
611 +
612 + def garbage_collection(self):
613 + if self.config['gc_run'] and self.runs % self.config['gc_interval'] == 0:
614 + v = gc.collect()
615 + self.log.debug('GC collection run result: {0}'.format(v))
616 +
617 + def restart_recovering_jobs(self):
618 + for job in self.jobs:
619 + if job.status != JOB_STATUS_RECOVERING:
620 + continue
621 + if self.runs % job.autodetection_retry != 0:
622 + continue
623 + self.start_jobs(job)
624
631 - if self.runs % job.autodetection_retry() != 0:
632 - return retry_later
625 + def cleanup_jobs(self):
626 + self.jobs = [j for j in self.jobs if j.status != JOB_STATUS_DROPPED]
627
634 - check_name = job.override_name or job.name
635 - if check_name in self.checked_jobs[job.module_name]:
636 - self.log.info("{0}[{1}]: already served by another job, give up on retrying".format(
637 - job.module_name, job.name))
638 - return stop_retrying
628 + def have_alive_jobs(self):
629 + return next(
630 + (True for job in self.jobs if job.status in (JOB_STATUS_RECOVERING, JOB_STATUS_ACTIVE)),
631 + False,
632 + )
633
634 + def save_job_statuses(self):
635 + if self.saver is None:
636 + return
637 + if self.runs % 10 != 0:
638 + return
639 + dump = JobsStatuses().from_jobs(self.jobs).dump()
640 try:
641 - ok = job.check()
641 + self.saver.save(dump)
642 except Exception as error:
643 - self.log.warning("{0}[{1}] : unhandled exception on recheck : {2}, give up on retrying".format(
644 - job.module_name, job.name, error))
645 - return stop_retrying
643 + self.log.error("error on saving jobs statuses dump : {0}".format(repr(error)))
644
647 - if not ok:
648 - self.log.info('{0}[{1}] : recheck failed, will retry in {2} second(s)'.format(
649 - job.module_name, job.name, job.autodetection_retry()))
650 - return retry_later
651 - self.log.info('{0}[{1}] : recheck successful'.format(job.module_name, job.name))
645 + def serve_once(self):
646 + if not self.have_alive_jobs():
647 + self.log.info('no jobs to serve')
648 + return False
649 +
650 + time.sleep(1)
651 + self.runs += 1
652
653 - if not job.create():
654 - return stop_retrying
653 + self.keep_alive()
654 + self.garbage_collection()
655 + self.cleanup_jobs()
656 + self.restart_recovering_jobs()
657 + self.save_job_statuses()
658 + return True
659
656 - job.post_check(int(self.min_update_every))
657 - self.checked_jobs[job.module_name].append(check_name)
658 - JobRunner(job).start()
660 + def serve(self):
661 + while self.serve_once():
662 + pass
663
660 - return stop_retrying
664 + def run(self):
665 + self.start_jobs(*self.jobs)
666 + self.serve()
667
668
663 -def parse_cmd():
669 +def parse_command_line():
670 opts = sys.argv[:][1:]
671 +
672 debug = False
673 trace = False
674 update_every = 1
675 modules_to_run = list()
676
670 - v = next((opt for opt in opts if opt.isdigit() and int(opt) >= 1), None)
671 - if v:
672 - update_every = v
673 - opts.remove(v)
677 + def find_first_positive_int(values):
678 + return next((v for v in values if v.isdigit() and int(v) >= 1), None)
679 +
680 + u = find_first_positive_int(opts)
681 + if u is not None:
682 + update_every = int(u)
683 + opts.remove(u)
684 if 'debug' in opts:
685 debug = True
686 opts.remove('debug')
@@ -680,54 +690,80 @@ def parse_cmd():
690 if opts:
691 modules_to_run = list(opts)
692
683 - return collections.namedtuple(
693 + cmd = collections.namedtuple(
694 'CMD',
695 [
696 'update_every',
697 'debug',
698 'trace',
699 'modules_to_run',
690 - ],
691 - )(
700 + ])
701 + return cmd(
702 update_every,
703 debug,
704 trace,
695 - modules_to_run,
705 + modules_to_run
706 )
707
708
709 +def guess_module(modules, *names):
710 + def guess(n):
711 + found = None
712 + for i, _ in enumerate(n):
713 + cur = [x for x in modules if x.startswith(name[:i + 1])]
714 + if not cur:
715 + return found
716 + found = cur
717 + return found
718 +
719 + guessed = list()
720 + for name in names:
721 + name = name.lower()
722 + m = guess(name)
723 + if m:
724 + guessed.extend(m)
725 + return sorted(set(guessed))
726 +
727 +
728 +def disable():
729 + if not IS_ATTY:
730 + safe_print('DISABLE')
731 + exit(0)
732 +
733 +
734 def main():
700 - cmd = parse_cmd()
701 - logger = PythonDLogger()
735 + cmd = parse_command_line()
736 + log = PythonDLogger()
737
738 if cmd.debug:
704 - logger.logger.severity = 'DEBUG'
739 + log.logger.severity = 'DEBUG'
740 if cmd.trace:
706 - logger.log_traceback = True
741 + log.log_traceback = True
742
708 - logger.info('using python v{0}'.format(PY_VERSION[0]))
743 + log.info('using python v{0}'.format(PY_VERSION[0]))
744
710 - unknown_modules = set(cmd.modules_to_run) - set(AVAILABLE_MODULES)
711 - if unknown_modules:
712 - logger.error('unknown modules : {0}'.format(sorted(list(unknown_modules))))
713 - safe_print('DISABLE')
745 + unknown = set(cmd.modules_to_run) - set(AVAILABLE_MODULES)
746 + if unknown:
747 + log.error('unknown modules : {0}'.format(sorted(list(unknown))))
748 + guessed = guess_module(AVAILABLE_MODULES, *cmd.modules_to_run)
749 + if guessed:
750 + log.info('probably you meant : \n{0}'.format(pprint.pformat(guessed, width=1)))
751 return
752
716 - plugin = Plugin(
717 - cmd.update_every,
753 + p = Plugin(
754 cmd.modules_to_run or AVAILABLE_MODULES,
755 + cmd.update_every,
756 )
757
721 - HeartBeat(1).start()
722 -
723 - if not plugin.setup():
724 - safe_print('DISABLE')
725 - return
726 -
727 - plugin.run()
728 - logger.info('exiting from main...')
729 - plugin.shutdown()
758 + try:
759 + if not p.setup():
760 + return
761 + p.run()
762 + except KeyboardInterrupt:
763 + pass
764 + log.info('exiting from main...')
765
766
732 -if __name__ == '__main__':
767 +if __name__ == "__main__":
768 main()
769 + disable()