old python.d.plugin replaced with the new one
lgz committed
Oct 13, 2017 at 00:23 UTC
41d9f2b13022ca16cdaf320eac540a1b149a6da3
1 file changed
+316
-544
plugins.d/python.d.plugin
+316
-544
@@ -1,599 +1,371 @@
1
+
2
#!/usr/bin/env bash
2
-'''':; exec "$(command -v python || command -v python3 || command -v python2 || echo "ERROR python IS NOT AVAILABLE IN THIS SYSTEM")" "$0" "$@" # '''
3
-# -*- coding: utf-8 -*-
3
+'''':; exec "$(command -v python || command -v python3 || command -v python2 ||
4
+echo "ERROR python IS NOT AVAILABLE IN THIS SYSTEM")" "$0" "$@" # '''
5
5
-# Description: netdata python modules supervisor
6
+# -*- coding: utf-8 -*-
7
+# Description:
8
# Author: Pawel Krupa (paulfantom)
9
+# Author: Ilya Mashchenko (l2isbad)
10
11
import os
12
import sys
10
-import time
13
import threading
14
+
15
from re import sub
16
+from sys import version_info, argv
17
+from time import sleep
18
+
19
+try:
20
+ from collections import OrderedDict
21
+except ImportError:
22
+ from third_party.ordereddict import OrderedDict
23
+
24
+try:
25
+ from time import monotonic as time
26
+except ImportError:
27
+ from time import time
28
+
29
+PY_VERSION = version_info[:2]
30
+PLUGIN_CONFIG_DIR = os.getenv('NETDATA_CONFIG_DIR', os.path.dirname(__file__) + '/../../../../etc/netdata') + '/'
31
+CHARTS_PY_DIR = os.path.abspath(os.getenv('NETDATA_PLUGINS_DIR', os.path.dirname(__file__)) + '/../python.d') + '/'
32
+CHARTS_PY_CONFIG_DIR = PLUGIN_CONFIG_DIR + 'python.d/'
33
+PYTHON_MODULES_DIR = CHARTS_PY_DIR + 'python_modules'
34
+
35
+sys.path.append(PYTHON_MODULES_DIR)
36
+
37
+from bases.loaders import ModuleAndConfigLoader
38
+from bases.loggers import PythonDLogger
39
+from bases.collection import setdefault_values, run_and_exit
40
14
-# -----------------------------------------------------------------------------
15
-# globals & environment setup
16
-# https://github.com/firehol/netdata/wiki/External-Plugins#environment-variables
17
-MODULE_EXTENSION = ".chart.py"
41
BASE_CONFIG = {'update_every': os.getenv('NETDATA_UPDATE_EVERY', 1),
19
- 'priority': 90000,
20
- 'retries': 10}
42
+ 'retries': 10,
43
+ 'priority': 60000,
44
+ 'autodetection_retry': 0,
45
+ 'name': str()}
46
22
-MODULES_DIR = os.path.abspath(os.getenv('NETDATA_PLUGINS_DIR',
23
- os.path.dirname(__file__)) + "/../python.d") + "/"
47
25
-CONFIG_DIR = os.getenv('NETDATA_CONFIG_DIR',
26
- os.path.dirname(__file__) + "/../../../../etc/netdata")
48
+MODULE_EXTENSION = '.chart.py'
49
+OBSOLETE_MODULES = ['apache_cache', 'gunicorn_log', 'nginx_log']
50
28
-# directories should end with '/'
29
-if CONFIG_DIR[-1] != "/":
30
- CONFIG_DIR += "/"
31
-sys.path.append(MODULES_DIR + "python_modules")
51
+JOB_INITIALIZE_MSG = 'Job initialization {status}. Module: "{module_name}", job: "{job_name}"{error}.'
52
+JOB_CREATE_CHARTS_MSG = 'CREATE {status}. Number of charts: {number_of_charts}.'
53
+MODULE_LOAD_MSG = 'Module load {status}. Module: "{module_name}", initialized jobs: {jobs_number}.'
54
+UNHANDLED_EXCEPTION = 'Unhandled exception in {method_name}(). Error: {error}.'
55
33
-PROGRAM = os.path.basename(__file__).replace(".plugin", "")
34
-DEBUG_FLAG = False
35
-TRACE_FLAG = False
36
-OVERRIDE_UPDATE_EVERY = False
56
38
-# -----------------------------------------------------------------------------
39
-# custom, third party and version specific python modules management
40
-import msg
57
+def module_ok(m):
58
+ return m.endswith(MODULE_EXTENSION) and m[:-len(MODULE_EXTENSION)] not in OBSOLETE_MODULES
59
42
-try:
43
- assert sys.version_info >= (3, 1)
44
- import importlib.machinery
45
- PY_VERSION = 3
46
- # change this hack below if we want PY_VERSION to be used in modules
47
- # import builtins
48
- # builtins.PY_VERSION = 3
49
- msg.info('Using python v3')
50
-except (AssertionError, ImportError):
51
- try:
52
- import imp
53
-
54
- # change this hack below if we want PY_VERSION to be used in modules
55
- # import __builtin__
56
- # __builtin__.PY_VERSION = 2
57
- PY_VERSION = 2
58
- msg.info('Using python v2')
59
- except ImportError:
60
- msg.fatal('Cannot start. No importlib.machinery on python3 or lack of imp on python2')
61
-# try:
62
-# import yaml
63
-# except ImportError:
64
-# msg.fatal('Cannot find yaml library')
65
-try:
66
- if PY_VERSION == 3:
67
- import pyyaml3 as yaml
68
- else:
69
- import pyyaml2 as yaml
70
-except ImportError:
71
- msg.fatal('Cannot find yaml library')
60
73
-try:
74
- from collections import OrderedDict
75
- ORDERED = True
76
- DICT = OrderedDict
77
- msg.info('YAML output is ordered')
78
-except ImportError:
79
- try:
80
- from ordereddict import OrderedDict
81
- ORDERED = True
82
- DICT = OrderedDict
83
- msg.info('YAML output is ordered')
84
- except ImportError:
85
- ORDERED = False
86
- DICT = dict
87
- msg.info('YAML output is unordered')
88
-if ORDERED:
89
- def ordered_load(stream, Loader=yaml.Loader, object_pairs_hook=OrderedDict):
90
- class OrderedLoader(Loader):
91
- pass
92
-
93
- def construct_mapping(loader, node):
94
- loader.flatten_mapping(node)
95
- return object_pairs_hook(loader.construct_pairs(node))
96
- OrderedLoader.add_constructor(
97
- yaml.resolver.BaseResolver.DEFAULT_MAPPING_TAG,
98
- construct_mapping)
99
- return yaml.load(stream, OrderedLoader)
100
-
101
-
102
-class PythonCharts(object):
103
- """
104
- Main class used to control every python module.
105
- """
106
-
107
- def __init__(self,
108
- modules=None,
109
- modules_path='../python.d/',
110
- modules_configs='../conf.d/',
111
- modules_disabled=None,
112
- modules_enabled=None,
113
- default_run=None):
61
+ALL_MODULES = [m for m in sorted(os.listdir(CHARTS_PY_DIR)) if module_ok(m)]
62
+
63
+
64
+def parse_cmd():
65
+ debug = 'debug' in argv[1:]
66
+ override_update_every = next((arg for arg in argv[1:] if arg.isdigit() and int(arg) > 1), False)
67
+ modules = [''.join([m, MODULE_EXTENSION]) for m in argv[1:] if ''.join([m, MODULE_EXTENSION]) in ALL_MODULES]
68
+ return debug, override_update_every, modules or ALL_MODULES
69
+
70
+
71
+def multi_job_check(config):
72
+ return next((True for key in config if isinstance(config[key], dict)), False)
73
+
74
+
75
+class Job(object):
76
+ def __init__(self, initialized_job, module_name, internal_name):
77
"""
115
- :param modules: list
116
- :param modules_path: str
117
- :param modules_configs: str
118
- :param modules_disabled: list
119
- :param modules_enabled: list
120
- :param default_run: bool
78
+ :param initialized_job: <Class Service>
79
+ :param module_name: <str>
80
+ :param internal_name: <str>
81
"""
82
+ self.job = initialized_job
83
+ self.internal_name = internal_name # key in Modules.jobs
84
+ self.module_name = module_name # needed in Plugin.delete_job
85
+ self.recheck_every = self.job.configuration.pop('autodetection_retry')
86
+ self.checked = False # needed in Plugin.check_job
87
+ self.created = False # needed in Plugin.create_job_charts
88
+ if OVERRIDE_UPDATE_EVERY:
89
+ self.job.update_every = int(OVERRIDE_UPDATE_EVERY)
90
123
- if modules is None:
124
- modules = []
125
- if modules_disabled is None:
126
- modules_disabled = []
91
+ def __getattr__(self, item):
92
+ return getattr(self.job, item)
93
128
- self.first_run = True
129
- # set configuration directory
130
- self.configs = modules_configs
94
+ def __repr__(self):
95
+ return self.job.__repr__()
96
132
- # load modules
133
- loaded_modules = self._load_modules(modules_path, modules, modules_disabled, modules_enabled, default_run)
97
+ def is_dead(self):
98
+ return bool(self.ident) and not self.is_alive()
99
135
- # load configuration files
136
- configured_modules = self._load_configs(loaded_modules)
100
+ def not_launched(self):
101
+ return not bool(self.ident)
102
138
- # good economy and prosperity:
139
- self.jobs = self._create_jobs(configured_modules) # type <list>
103
+ def is_autodetect(self):
104
+ return self.recheck_every
105
141
- # enable timetable override like `python.d.plugin mysql debug 1`
142
- if DEBUG_FLAG and OVERRIDE_UPDATE_EVERY:
143
- for job in self.jobs:
144
- job.create_timetable(BASE_CONFIG['update_every'])
106
146
- @staticmethod
147
- def _import_module(path, name=None):
107
+class Module(object):
108
+ def __init__(self, service, config):
109
"""
149
- Try to import module using only its path.
150
- :param path: str
151
- :param name: str
152
- :return: object
110
+ :param service: <Module>
111
+ :param config: <dict>
112
"""
113
+ self.service = service
114
+ self.name = service.__name__
115
+ self.config = self.build_jobs_configurations(config)
116
+ self.jobs = OrderedDict()
117
155
- if name is None:
156
- name = path.split('/')[-1]
157
- if name[-len(MODULE_EXTENSION):] != MODULE_EXTENSION:
158
- return None
159
- name = name[:-len(MODULE_EXTENSION)]
160
- try:
161
- if PY_VERSION == 3:
162
- return importlib.machinery.SourceFileLoader(name, path).load_module()
163
- else:
164
- return imp.load_source(name, path)
165
- except Exception as e:
166
- msg.error("Problem loading", name, str(e))
167
- return None
118
+ self.initialize_jobs()
119
169
- def _load_modules(self, path, modules, disabled, enabled, default_run):
170
- """
171
- Load modules from 'modules' list or dynamically every file from 'path' (only .chart.py files)
172
- :param path: str
173
- :param modules: list
174
- :param disabled: list
175
- :return: list
176
- """
120
+ def __repr__(self):
121
+ return '<Class Module "{name}">'.format(name=self.name)
122
178
- # check if plugin directory exists
179
- if not os.path.isdir(path):
180
- msg.fatal("cannot find charts directory ", path)
123
+ def __iter__(self):
124
+ return iter(OrderedDict(self.jobs).values())
125
182
- # load modules
183
- loaded = []
184
- if len(modules) > 0:
185
- for m in modules:
186
- if m in disabled:
187
- continue
188
- mod = self._import_module(path + m + MODULE_EXTENSION)
189
- if mod is not None:
190
- loaded.append(mod)
191
- else: # exit if plugin is not found
192
- msg.fatal('no modules found.')
193
- else:
194
- # scan directory specified in path and load all modules from there
195
- if default_run is False:
196
- names = [module for module in os.listdir(path) if module[:-9] in enabled]
197
- else:
198
- names = os.listdir(path)
199
- for mod in names:
200
- if mod.replace(MODULE_EXTENSION, "") in disabled:
201
- msg.error(mod + ": disabled module ", mod.replace(MODULE_EXTENSION, ""))
202
- continue
203
- m = self._import_module(path + mod)
204
- if m is not None:
205
- msg.debug(mod + ": loading module '" + path + mod + "'")
206
- loaded.append(m)
207
- return loaded
126
+ def __getitem__(self, item):
127
+ return self.jobs[item]
128
209
- def _load_configs(self, modules):
210
- """
211
- Append configuration in list named `config` to every module.
212
- For multi-job modules `config` list is created in _parse_config,
213
- otherwise it is created here based on BASE_CONFIG prototype with None as identifier.
214
- :param modules: list
215
- :return: list
216
- """
217
- for mod in modules:
218
- configfile = self.configs + mod.__name__ + ".conf"
219
- if os.path.isfile(configfile):
220
- msg.debug(mod.__name__ + ": loading module configuration: '" + configfile + "'")
221
- try:
222
- if not hasattr(mod, 'config'):
223
- mod.config = {}
224
- setattr(mod,
225
- 'config',
226
- self._parse_config(mod, read_config(configfile)))
227
- except Exception as e:
228
- msg.error(mod.__name__ + ": cannot parse configuration file '" + configfile + "':", str(e))
229
- else:
230
- msg.error(mod.__name__ + ": configuration file '" + configfile + "' not found. Using defaults.")
231
- # set config if not found
232
- if not hasattr(mod, 'config'):
233
- msg.debug(mod.__name__ + ": setting configuration for only one job")
234
- mod.config = {None: {}}
235
- for var in BASE_CONFIG:
236
- try:
237
- mod.config[None][var] = getattr(mod, var)
238
- except AttributeError:
239
- mod.config[None][var] = BASE_CONFIG[var]
240
- return modules
129
+ def __delitem__(self, key):
130
+ del self.jobs[key]
131
242
- @staticmethod
243
- def _parse_config(module, config):
244
- """
245
- Parse configuration file or extract configuration from module file.
246
- Example of returned dictionary:
247
- config = {'name': {
248
- 'update_every': 2,
249
- 'retries': 3,
250
- 'priority': 30000
251
- 'other_val': 123}}
252
- :param module: object
253
- :param config: dict
254
- :return: dict
255
- """
256
- if config is None:
257
- config = {}
258
- # get default values
259
- defaults = {}
260
- msg.debug(module.__name__ + ": reading configuration")
261
- for key in BASE_CONFIG:
262
- try:
263
- # get defaults from module config
264
- defaults[key] = int(config.pop(key))
265
- except (KeyError, ValueError):
266
- try:
267
- # get defaults from module source code
268
- defaults[key] = getattr(module, key)
269
- except (KeyError, ValueError, AttributeError):
270
- # if above failed, get defaults from global dict
271
- defaults[key] = BASE_CONFIG[key]
272
-
273
- # check if there are dict in config dict
274
- many_jobs = False
275
- for name in config:
276
- if isinstance(config[name], DICT):
277
- many_jobs = True
278
- break
279
-
280
- # assign variables needed by supervisor to every job configuration
281
- if many_jobs:
282
- for name in config:
283
- for key in defaults:
284
- if key not in config[name]:
285
- config[name][key] = defaults[key]
286
- # if only one job is needed, values doesn't have to be in dict (in YAML)
287
- else:
288
- config = {None: config.copy()}
289
- config[None].update(defaults)
132
+ def __len__(self):
133
+ return len(self.jobs)
134
291
- # return dictionary of jobs where every job has BASE_CONFIG variables
292
- return config
135
+ def __bool__(self):
136
+ return bool(len(self.jobs))
137
294
- @staticmethod
295
- def _create_jobs(modules):
138
+ def __nonzero__(self):
139
+ return self.__bool__()
140
+
141
+ def build_jobs_configurations(self, config):
142
"""
297
- Create jobs based on module.config dictionary and module.Service class definition.
298
- :param modules: list
299
- :return: list
143
+ :param config: <dict>
144
+ :return:
145
"""
301
- jobs = []
302
- for module in modules:
303
- for name in module.config:
304
- # register a new job
305
- conf = module.config[name]
306
- try:
307
- job = module.Service(configuration=conf, name=name)
308
- except Exception as e:
309
- msg.error(module.__name__ +
310
- ("/" + str(name) if name is not None else "") +
311
- ": cannot start job: '" +
312
- str(e))
146
+ if not config:
147
+ return {self.name: dict(BASE_CONFIG)}
148
+ if not multi_job_check(config):
149
+ return {self.name: setdefault_values(config, BASE_CONFIG)}
150
+ else:
151
+ jobs, jobs_base_config = OrderedDict(), dict()
152
+ for attr in BASE_CONFIG:
153
+ jobs_base_config[attr] = config.pop(attr, BASE_CONFIG[attr])
154
+ for job_name in config:
155
+ if not isinstance(config[job_name], dict):
156
continue
314
- else:
315
- # set chart_name (needed to plot run time graphs)
316
- job.chart_name = module.__name__
317
- if name is not None:
318
- job.chart_name += "_" + name
319
- jobs.append(job)
320
- msg.debug(module.__name__ + ("/" + str(name) if name is not None else "") + ": job added")
157
322
- return [j for j in jobs if j is not None]
158
+ defaulted = setdefault_values(config[job_name], base_dict=jobs_base_config)
159
+ config[job_name]['name'] = sub(r'\s+', '_', config[job_name]['name'])
160
+ job_name = sub(r'\s+', '_', job_name)
161
+ job_internal_name = '_'.join([self.name, job_name])
162
+ jobs[job_internal_name] = defaulted
163
+ return jobs
164
+
165
+ def initialize_jobs(self):
166
+ for job_internal_name in self.config:
167
+ job_name = job_internal_name[len(self.name) + 1:]
168
+ job_override_name = self.config[job_internal_name].pop('name')
169
+ self.config[job_internal_name]['module_name'] = self.name
170
+ self.config[job_internal_name]['job_name'] = job_name
171
+ self.config[job_internal_name]['override_name'] = job_override_name
172
+ try:
173
+ initialized_job = self.service.Service(configuration=self.config[job_internal_name])
174
+ except Exception as error:
175
+ Logger.error(JOB_INITIALIZE_MSG.format(status='FAILED',
176
+ module_name=self.name,
177
+ job_name=job_name,
178
+ error=', error: {0}'.format(error)))
179
+ continue
180
+ else:
181
+ Logger.debug(JOB_INITIALIZE_MSG.format(status='SUCCESS',
182
+ module_name=self.name,
183
+ job_name=job_name,
184
+ error=''))
185
+ self.jobs[job_internal_name] = Job(initialized_job=initialized_job,
186
+ module_name=self.name,
187
+ internal_name=job_internal_name)
188
+ del self.config
189
+ del self.service
190
+
191
+
192
+class Plugin(object):
193
+ def __init__(self):
194
+ self.loader = ModuleAndConfigLoader()
195
+ self.modules = OrderedDict()
196
+ self.sleep_time = 1
197
+ self.number_of_runs = 0
198
+ self.config, error = self.loader.load_config_from_file(PLUGIN_CONFIG_DIR + 'python.d.conf')
199
+ if error:
200
+ run_and_exit(Logger.error)(error)
201
+
202
+ if not self.config.get('enabled', True):
203
+ run_and_exit(Logger.info)('DISABLED in configuration file.')
204
+
205
+ self.load_and_initialize_modules()
206
+ if not self.modules:
207
+ run_and_exit(Logger.info)('No modules to run. Exit...')
208
+
209
+ def __iter__(self):
210
+ return iter(dict(self.modules).values())
211
+
212
+ @property
213
+ def jobs(self):
214
+ return (job for mod in self for job in mod)
215
+
216
+ @property
217
+ def dead_jobs(self):
218
+ return (job for job in self.jobs if job.is_dead())
219
+
220
+ @property
221
+ def autodetect_jobs(self):
222
+ return [job for job in self.jobs if job.not_launched()]
223
+
224
+ def enabled_modules(self):
225
+ for mod in MODULES_TO_RUN:
226
+ mod_name = mod[:-len(MODULE_EXTENSION)]
227
+ mod_path = CHARTS_PY_DIR + mod
228
+ conf_path = ''.join([CHARTS_PY_CONFIG_DIR, mod_name, '.conf'])
229
+
230
+ if DEBUG:
231
+ yield mod, mod_name, mod_path, conf_path
232
+ else:
233
+ if all([self.config.get('default_run', True),
234
+ self.config.get(mod_name, True)]):
235
+ yield mod, mod_name, mod_path, conf_path
236
+
237
+ elif all([not self.config.get('default_run'),
238
+ self.config.get(mod_name)]):
239
+ yield mod, mod_name, mod_path, conf_path
240
+
241
+ def load_and_initialize_modules(self):
242
+ for mod, mod_name, mod_path, conf_path in self.enabled_modules():
243
+ loaded_module, error = self.loader.load_module_from_file(mod_name, mod_path)
244
+ if error:
245
+ Logger.error('Module load failed. Module: {name}, error: {error}.'.format(name=mod,
246
+ error=error))
247
+ continue
248
+ loaded_config, error = self.loader.load_config_from_file(conf_path)
249
+ if error:
250
+ Logger.error('Config load failed. Module {module}, error: {error}.'.format(module=mod_name,
251
+ error=error))
252
+ initialized_module = Module(service=loaded_module, config=loaded_config)
253
+ Logger.debug(MODULE_LOAD_MSG.format(status='OK' if initialized_module else 'FAILED',
254
+ module_name=initialized_module.name,
255
+ jobs_number=len(initialized_module)))
256
+ if initialized_module:
257
+ self.modules[initialized_module.name] = initialized_module
258
324
- def _stop(self, job, reason=None):
259
+ @staticmethod
260
+ def check_job(job):
261
"""
326
- Stop specified job and remove it from self.jobs list
327
- Also notifies user about job failure if DEBUG_FLAG is set
328
- :param job: object
329
- :param reason: str
262
+ :param job: <Job>
263
+ :return:
264
"""
331
- prefix = job.__module__
332
- if job.name is not None and len(job.name) != 0:
333
- prefix += "/" + job.name
265
try:
335
- msg.error("DISABLED:", prefix)
336
- self.jobs.remove(job)
337
- except Exception as e:
338
- msg.debug("This shouldn't happen. NO " + prefix + " IN LIST:" + str(self.jobs) + " ERROR: " + str(e))
339
-
340
- # TODO remove section below and remove `reason`.
341
- prefix += ": "
342
- if reason is None:
343
- return
344
- elif reason[:3] == "no ":
345
- msg.error(prefix +
346
- "does not seem to have " +
347
- reason[3:] +
348
- "() function. Disabling it.")
349
- elif reason[:7] == "failed ":
350
- msg.error(prefix +
351
- reason[7:] +
352
- "() function reports failure.")
353
- elif reason[:13] == "configuration":
354
- msg.error(prefix +
355
- "configuration file '" +
356
- self.configs +
357
- job.__module__ +
358
- ".conf' not found. Using defaults.")
359
- elif reason[:11] == "misbehaving":
360
- msg.error(prefix + "is " + reason)
361
-
362
- def check(self):
363
- """
364
- Tries to execute check() on every job.
365
- This cannot fail thus it is catching every exception
366
- If job.check() fails job is stopped
367
- """
368
- i = 0
369
- overridden = []
370
- msg.debug("all job objects", str(self.jobs))
371
- while i < len(self.jobs):
372
- job = self.jobs[i]
373
- try:
374
- if not job.check():
375
- msg.error(job.chart_name, "check() failed - disabling job")
376
- self._stop(job)
377
- else:
378
- msg.info("CHECKED OK:", job.chart_name)
379
- i += 1
380
- try:
381
- if job.override_name is not None:
382
- new_name = job.__module__ + '_' + sub(r'\s+', '_', job.override_name)
383
- if new_name in overridden:
384
- msg.info("DROPPED:", job.name, ", job '" + job.override_name +
385
- "' is already served by another job.")
386
- self._stop(job)
387
- i -= 1
388
- else:
389
- job.name = job.override_name
390
- msg.info("RENAMED:", new_name, ", from " + job.chart_name)
391
- job.chart_name = new_name
392
- overridden.append(job.chart_name)
393
- except Exception:
394
- pass
395
- except AttributeError as e:
396
- self._stop(job)
397
- msg.error(job.chart_name, "cannot find check() function or it thrown unhandled exception.")
398
- msg.debug(str(e))
399
- except (UnboundLocalError, Exception) as e:
400
- msg.error(job.chart_name, str(e))
401
- self._stop(job)
402
- msg.debug("overridden job names:", str(overridden))
403
- msg.debug("all remaining job objects:", str(self.jobs))
404
-
405
- def create(self):
406
- """
407
- Tries to execute create() on every job.
408
- This cannot fail thus it is catching every exception.
409
- If job.create() fails job is stopped.
410
- This is also creating job run time chart.
411
- """
412
- i = 0
413
- while i < len(self.jobs):
414
- job = self.jobs[i]
415
- try:
416
- if not job.create():
417
- msg.error(job.chart_name, "create function failed.")
418
- self._stop(job)
419
- else:
420
- chart = job.chart_name
421
- sys.stdout.write(
422
- "CHART netdata.plugin_pythond_" +
423
- chart +
424
- " '' 'Execution time for " +
425
- chart +
426
- " plugin' 'milliseconds / run' python.d netdata.plugin_python area 145000 " +
427
- str(job.timetable['freq']) +
428
- '\n')
429
- sys.stdout.write("DIMENSION run_time 'run time' absolute 1 1\n\n")
430
- msg.debug("created charts for", job.chart_name)
431
- # sys.stdout.flush()
432
- i += 1
433
- except AttributeError:
434
- msg.error(job.chart_name, "cannot find create() function or it thrown unhandled exception.")
435
- self._stop(job)
436
- except (UnboundLocalError, Exception) as e:
437
- msg.error(job.chart_name, str(e))
438
- self._stop(job)
439
-
440
- def update(self):
441
- """
442
- Creates and supervises every job thread.
443
- This will stay forever and ever and ever forever and ever it'll be the one...
444
- """
445
- for job in self.jobs:
446
- job.start()
447
-
448
- while True:
449
- if threading.active_count() <= 1:
450
- msg.fatal("no more jobs")
451
- time.sleep(1)
452
-
453
-
454
-def read_config(path):
455
- """
456
- Read YAML configuration from specified file
457
- :param path: str
458
- :return: dict
459
- """
460
- try:
461
- with open(path, 'r') as stream:
462
- if ORDERED:
463
- config = ordered_load(stream, yaml.SafeLoader)
464
- else:
465
- config = yaml.load(stream)
466
- except (OSError, IOError) as error:
467
- msg.error(str(path), 'reading error:', str(error))
468
- return None
469
- except yaml.YAMLError as error:
470
- msg.error(str(path), "is malformed:", str(error))
471
- return None
472
- return config
473
-
474
-
475
-def parse_cmdline(directory, *commands):
476
- """
477
- Parse parameters from command line.
478
- :param directory: str
479
- :param commands: list of str
480
- :return: dict
481
- """
482
- global DEBUG_FLAG, TRACE_FLAG
483
- global OVERRIDE_UPDATE_EVERY
484
- global BASE_CONFIG
485
-
486
- changed_update = False
487
- mods = []
488
- for cmd in commands[1:]:
489
- if cmd == "check":
490
- pass
491
- elif cmd == "debug" or cmd == "all":
492
- DEBUG_FLAG = True
493
- # redirect stderr to stdout?
494
- elif cmd == "trace" or cmd == "all":
495
- TRACE_FLAG = True
496
- elif os.path.isfile(directory + cmd + ".chart.py") or os.path.isfile(directory + cmd):
497
- # DEBUG_FLAG = True
498
- mods.append(cmd.replace(".chart.py", ""))
266
+ check_ok = job.check()
267
+ except Exception as error:
268
+ job.error(UNHANDLED_EXCEPTION.format(method_name='check',
269
+ error=error))
270
+ return False
271
else:
500
- try:
501
- BASE_CONFIG['update_every'] = int(cmd)
502
- changed_update = True
503
- except ValueError:
504
- pass
505
- if changed_update and DEBUG_FLAG:
506
- OVERRIDE_UPDATE_EVERY = True
507
- msg.debug(PROGRAM, "overriding update interval to", str(BASE_CONFIG['update_every']))
508
-
509
- msg.debug("started from", commands[0], "with options:", *commands[1:])
510
-
511
- return mods
512
-
513
-
514
-# if __name__ == '__main__':
515
-def run():
516
- """
517
- Main program.
518
- """
519
- global DEBUG_FLAG, TRACE_FLAG, BASE_CONFIG
520
-
521
- # read configuration file
522
- disabled = ['nginx_log', 'gunicorn_log', 'apache_cache']
523
- enabled = list()
524
- default_run = True
525
- configfile = CONFIG_DIR + "python.d.conf"
526
- msg.PROGRAM = PROGRAM
527
- msg.info("reading configuration file:", configfile)
528
- log_throttle = 200
529
- log_interval = 3600
530
-
531
- conf = read_config(configfile)
532
- if conf is not None:
533
- try:
534
- # exit the whole plugin when 'enabled: no' is set in 'python.d.conf'
535
- if conf['enabled'] is False:
536
- msg.fatal('disabled in configuration file.\n')
537
- except (KeyError, TypeError):
538
- pass
539
-
540
- try:
541
- for param in BASE_CONFIG:
542
- BASE_CONFIG[param] = conf[param]
543
- except (KeyError, TypeError):
544
- pass # use default update_every from NETDATA_UPDATE_EVERY
272
+ log = job.info if check_ok else job.error
273
+ log('CHECK {status}'.format(status='OK' if check_ok else 'FAILED'))
274
+ if job.is_autodetect():
275
+ if not check_ok:
276
+ job.error('Next CHECK in {0} seconds.'.format(job.recheck_every))
277
+ return check_ok
278
279
+ @staticmethod
280
+ def create_job_charts(job):
281
+ """
282
+ :param job: <Job>
283
+ :return:
284
+ """
285
try:
547
- DEBUG_FLAG = conf['debug']
548
- except (KeyError, TypeError):
549
- pass
286
+ create_ok = job.create()
287
+ except Exception as error:
288
+ job.error(UNHANDLED_EXCEPTION.format(method_name='create',
289
+ error=error))
290
+ return False
291
+ else:
292
+ job.debug(JOB_CREATE_CHARTS_MSG.format(status='OK' if create_ok else 'FAILED',
293
+ number_of_charts=len(job.charts)))
294
+ return create_ok
295
551
- try:
552
- TRACE_FLAG = conf['trace']
553
- except (KeyError, TypeError):
554
- pass
296
+ def delete_job(self, job):
297
+ """
298
+ :param job: <Job>
299
+ :return:
300
+ """
301
+ del self.modules[job.module_name][job.internal_name]
302
556
- try:
557
- log_throttle = conf['logs_per_interval']
558
- except (KeyError, TypeError):
559
- pass
303
+ def run_check(self):
304
+ checked = list()
305
+ for job in self.jobs:
306
+ if job.name in checked:
307
+ Logger.info('DROPPED: {job_name}. Already served by another job.'.format(job_name=job.internal_name))
308
+ self.delete_job(job)
309
+ continue
310
+ ok = self.check_job(job)
311
+ if ok:
312
+ checked.append(job.name)
313
+ job.checked = True
314
+ continue
315
+ if not job.is_autodetect():
316
+ self.delete_job(job)
317
561
- try:
562
- log_interval = conf['log_interval']
563
- except (KeyError, TypeError):
564
- pass
318
+ def run_create(self):
319
+ for job in self.jobs:
320
+ if not job.checked:
321
+ continue
322
+ ok = self.create_job_charts(job)
323
+ if ok:
324
+ job.created = True
325
+ continue
326
+ self.delete_job(job)
327
566
- default_run = True if ('default_run' not in conf or conf.get('default_run')) else False
328
+ def start(self):
329
+ self.run_check()
330
+ self.run_create()
331
+ for job in self.jobs:
332
+ if job.created:
333
+ job.start()
334
568
- for k, v in conf.items():
569
- if k in ("update_every", "debug", "enabled", "default_run"):
570
- continue
571
- if default_run:
572
- if v is False:
573
- disabled.append(k)
574
- else:
575
- if v is True:
576
- enabled.append(k)
577
- # parse passed command line arguments
578
- modules = parse_cmdline(MODULES_DIR, *sys.argv)
579
- msg.DEBUG_FLAG = DEBUG_FLAG
580
- msg.TRACE_FLAG = TRACE_FLAG
581
- msg.LOG_THROTTLE = log_throttle
582
- msg.LOG_INTERVAL = log_interval
583
- msg.LOG_COUNTER = 0
584
- msg.LOG_NEXT_CHECK = 0
585
- msg.info("MODULES_DIR='" + MODULES_DIR +
586
- "', CONFIG_DIR='" + CONFIG_DIR +
587
- "', UPDATE_EVERY=" + str(BASE_CONFIG['update_every']) +
588
- ", ONLY_MODULES=" + str(modules))
589
-
590
- # run plugins
591
- charts = PythonCharts(modules, MODULES_DIR, CONFIG_DIR + "python.d/", disabled, enabled, default_run)
592
- charts.check()
593
- charts.create()
594
- charts.update()
595
- msg.fatal("finished")
335
+ while True:
336
+ if threading.active_count() <= 1 and not self.autodetect_jobs:
337
+ run_and_exit(Logger.info)('FINISHED')
338
+
339
+ sleep(self.sleep_time)
340
+ self.cleanup()
341
+ self.autodetect_retry()
342
+
343
+ def cleanup(self):
344
+ for job in self.dead_jobs:
345
+ self.delete_job(job)
346
+ for mod in self:
347
+ if not mod:
348
+ del self.modules[mod.name]
349
+
350
+ def autodetect_retry(self):
351
+ self.number_of_runs += self.sleep_time
352
+ for job in self.autodetect_jobs:
353
+ if self.number_of_runs % job.recheck_every == 0:
354
+ checked = self.check_job(job)
355
+ if checked:
356
+ created = self.create_job_charts(job)
357
+ if not created:
358
+ self.delete_job(job)
359
+ continue
360
+ job.start()
361
362
363
if __name__ == '__main__':
599
- run()
364
+ DEBUG, OVERRIDE_UPDATE_EVERY, MODULES_TO_RUN = parse_cmd()
365
+ Logger = PythonDLogger()
366
+ if DEBUG:
367
+ Logger.logger.severity = 'DEBUG'
368
+ Logger.info('Using python {version}'.format(version=PY_VERSION[0]))
369
+
370
+ plugin = Plugin()
371
+ plugin.start()