replace base.py content (now it's uses only to import all FrameworkServices for backward compatibility with old version)
lgz committed
Oct 13, 2017 at 00:21 UTC
d8cef6af549782a6f372d436fc24e174e50e555f
1 file changed
+8
-1123
python.d/python_modules/base.py
+8
-1123
@@ -1,1124 +1,9 @@
1
# -*- coding: utf-8 -*-
2
-# Description: netdata python modules framework
3
-# Author: Pawel Krupa (paulfantom)
4
-
5
-# Remember:
6
-# ALL CODE NEEDS TO BE COMPATIBLE WITH Python > 2.7 and Python > 3.1
7
-# Follow PEP8 as much as it is possible
8
-# "check" and "create" CANNOT be blocking.
9
-# "update" CAN be blocking
10
-# "update" function needs to be fast, so follow:
11
-# https://wiki.python.org/moin/PythonSpeed/PerformanceTips
12
-# basically:
13
-# - use local variables wherever it is possible
14
-# - avoid dots in expressions that are executed many times
15
-# - use "join()" instead of "+"
16
-# - use "import" only at the beginning
17
-#
18
-# using ".encode()" in one thread can block other threads as well (only in python2)
19
-
20
-import os
21
-import re
22
-import socket
23
-import time
24
-import threading
25
-
26
-import urllib3
27
-
28
-from glob import glob
29
-from subprocess import Popen, PIPE
30
-from sys import exc_info
31
-
32
-try:
33
- import MySQLdb
34
- PY_MYSQL = True
35
-except ImportError:
36
- try:
37
- import pymysql as MySQLdb
38
- PY_MYSQL = True
39
- except ImportError:
40
- PY_MYSQL = False
41
-
42
-import msg
43
-
44
-
45
-PATH = os.getenv('PATH', '/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin').split(':')
46
-try:
47
- urllib3.disable_warnings()
48
-except AttributeError:
49
- msg.error('urllib3: warnings were not disabled')
50
-
51
-
52
-# class BaseService(threading.Thread):
53
-class SimpleService(threading.Thread):
54
- """
55
- Prototype of Service class.
56
- Implemented basic functionality to run jobs by `python.d.plugin`
57
- """
58
- def __init__(self, configuration=None, name=None):
59
- """
60
- This needs to be initialized in child classes
61
- :param configuration: dict
62
- :param name: str
63
- """
64
- threading.Thread.__init__(self)
65
- self._data_stream = ""
66
- self.daemon = True
67
- self.retries = 0
68
- self.retries_left = 0
69
- self.priority = 140000
70
- self.update_every = 1
71
- self.name = name
72
- self.override_name = None
73
- self.chart_name = ""
74
- self._dimensions = []
75
- self._charts = []
76
- self.__chart_set = False
77
- self.__first_run = True
78
- self.order = []
79
- self.definitions = {}
80
- self._data_from_check = dict()
81
- if configuration is None:
82
- self.error("BaseService: no configuration parameters supplied. Cannot create Service.")
83
- raise RuntimeError
84
- else:
85
- self._extract_base_config(configuration)
86
- self.timetable = {}
87
- self.create_timetable()
88
-
89
- # --- BASIC SERVICE CONFIGURATION ---
90
-
91
- def _extract_base_config(self, config):
92
- """
93
- Get basic parameters to run service
94
- Minimum config:
95
- config = {'update_every':1,
96
- 'priority':100000,
97
- 'retries':0}
98
- :param config: dict
99
- """
100
- pop = config.pop
101
- try:
102
- self.override_name = pop('name')
103
- except KeyError:
104
- pass
105
- self.update_every = int(pop('update_every'))
106
- self.priority = int(pop('priority'))
107
- self.retries = int(pop('retries'))
108
- self.retries_left = self.retries
109
- self.configuration = config
110
-
111
- def create_timetable(self, freq=None):
112
- """
113
- Create service timetable.
114
- `freq` is optional
115
- Example:
116
- timetable = {'last': 1466370091.3767564,
117
- 'next': 1466370092,
118
- 'freq': 1}
119
- :param freq: int
120
- """
121
- if freq is None:
122
- freq = self.update_every
123
- now = time.time()
124
- self.timetable = {'last': now,
125
- 'next': now - (now % freq) + freq,
126
- 'freq': freq}
127
-
128
- # --- THREAD CONFIGURATION ---
129
-
130
- def _run_once(self):
131
- """
132
- Executes self.update(interval) and draws run time chart.
133
- Return value presents exit status of update()
134
- :return: boolean
135
- """
136
- t_start = float(time.time())
137
- chart_name = self.chart_name
138
-
139
- since_last = int((t_start - self.timetable['last']) * 1000000)
140
- if self.__first_run:
141
- since_last = 0
142
-
143
- if not self.update(since_last):
144
- self.error("update function failed.")
145
- return False
146
-
147
- # draw performance graph
148
- run_time = int((time.time() - t_start) * 1000)
149
- print("BEGIN netdata.plugin_pythond_%s %s\nSET run_time = %s\nEND\n" %
150
- (self.chart_name, str(since_last), str(run_time)))
151
-
152
- self.debug(chart_name, "updated in", str(run_time), "ms")
153
- self.timetable['last'] = t_start
154
- self.__first_run = False
155
- return True
156
-
157
- def run(self):
158
- """
159
- Runs job in thread. Handles retries.
160
- Exits when job failed or timed out.
161
- :return: None
162
- """
163
- step = float(self.timetable['freq'])
164
- penalty = 0
165
- self.timetable['last'] = float(time.time() - step)
166
- self.debug("starting data collection - update frequency:", str(step), " retries allowed:", str(self.retries))
167
- while True: # run forever, unless something is wrong
168
- now = float(time.time())
169
- next = self.timetable['next'] = now - (now % step) + step + penalty
170
-
171
- # it is important to do this in a loop
172
- # sleep() is interruptable
173
- while now < next:
174
- self.debug("sleeping for", str(next - now), "secs to reach frequency of",
175
- str(step), "secs, now:", str(now), " next:", str(next), " penalty:", str(penalty))
176
- time.sleep(next - now)
177
- now = float(time.time())
178
-
179
- # do the job
180
- try:
181
- status = self._run_once()
182
- except Exception:
183
- status = False
184
-
185
- if status:
186
- # it is good
187
- self.retries_left = self.retries
188
- penalty = 0
189
- else:
190
- # it failed
191
- self.retries_left -= 1
192
- if self.retries_left <= 0:
193
- if penalty == 0:
194
- penalty = float(self.retries * step) / 2
195
- else:
196
- penalty *= 1.5
197
-
198
- if penalty > 600:
199
- penalty = 600
200
-
201
- self.retries_left = self.retries
202
- self.alert("failed to collect data for " + str(self.retries) +
203
- " times - increasing penalty to " + str(penalty) + " sec and trying again")
204
-
205
- else:
206
- self.error("failed to collect data - " + str(self.retries_left)
207
- + " retries left - penalty: " + str(penalty) + " sec")
208
-
209
- # --- CHART ---
210
-
211
- @staticmethod
212
- def _format(*args):
213
- """
214
- Escape and convert passed arguments.
215
- :param args: anything
216
- :return: list
217
- """
218
- params = []
219
- append = params.append
220
- for p in args:
221
- if p is None:
222
- append(p)
223
- continue
224
- if type(p) is not str:
225
- p = str(p)
226
- if ' ' in p:
227
- p = "'" + p + "'"
228
- append(p)
229
- return params
230
-
231
- def _line(self, instruction, *params):
232
- """
233
- Converts *params to string and joins them with one space between every one.
234
- Result is appended to self._data_stream
235
- :param params: str/int/float
236
- """
237
- tmp = list(map((lambda x: "''" if x is None or len(x) == 0 else x), params))
238
- self._data_stream += "%s %s\n" % (instruction, str(" ".join(tmp)))
239
-
240
- def chart(self, type_id, name="", title="", units="", family="",
241
- category="", chart_type="line", priority="", update_every=""):
242
- """
243
- Defines a new chart.
244
- :param type_id: str
245
- :param name: str
246
- :param title: str
247
- :param units: str
248
- :param family: str
249
- :param category: str
250
- :param chart_type: str
251
- :param priority: int/str
252
- :param update_every: int/str
253
- """
254
- self._charts.append(type_id)
255
-
256
- p = self._format(type_id, name, title, units, family, category, chart_type, priority, update_every)
257
- self._line("CHART", *p)
258
-
259
- def dimension(self, id, name=None, algorithm="absolute", multiplier=1, divisor=1, hidden=False):
260
- """
261
- Defines a new dimension for the chart
262
- :param id: str
263
- :param name: str
264
- :param algorithm: str
265
- :param multiplier: int/str
266
- :param divisor: int/str
267
- :param hidden: boolean
268
- :return:
269
- """
270
- try:
271
- int(multiplier)
272
- except TypeError:
273
- self.error("malformed dimension: multiplier is not a number:", multiplier)
274
- multiplier = 1
275
- try:
276
- int(divisor)
277
- except TypeError:
278
- self.error("malformed dimension: divisor is not a number:", divisor)
279
- divisor = 1
280
- if name is None:
281
- name = id
282
- if algorithm not in ("absolute", "incremental", "percentage-of-absolute-row", "percentage-of-incremental-row"):
283
- algorithm = "absolute"
284
-
285
- self._dimensions.append(str(id))
286
- if hidden:
287
- p = self._format(id, name, algorithm, multiplier, divisor, "hidden")
288
- else:
289
- p = self._format(id, name, algorithm, multiplier, divisor)
290
-
291
- self._line("DIMENSION", *p)
292
-
293
- def begin(self, type_id, microseconds=0):
294
- """
295
- Begin data set
296
- :param type_id: str
297
- :param microseconds: int
298
- :return: boolean
299
- """
300
- if type_id not in self._charts:
301
- self.error("wrong chart type_id:", type_id)
302
- return False
303
- try:
304
- int(microseconds)
305
- except TypeError:
306
- self.error("malformed begin statement: microseconds are not a number:", microseconds)
307
- microseconds = ""
308
-
309
- self._line("BEGIN", type_id, str(microseconds))
310
- return True
311
-
312
- def set(self, id, value):
313
- """
314
- Set value to dimension
315
- :param id: str
316
- :param value: int/float
317
- :return: boolean
318
- """
319
- if id not in self._dimensions:
320
- self.error("wrong dimension id:", id, "Available dimensions are:", *self._dimensions)
321
- return False
322
- try:
323
- value = str(int(value))
324
- except TypeError:
325
- self.error("cannot set non-numeric value:", str(value))
326
- return False
327
- self._line("SET", id, "=", str(value))
328
- self.__chart_set = True
329
- return True
330
-
331
- def end(self):
332
- if self.__chart_set:
333
- self._line("END")
334
- self.__chart_set = False
335
- else:
336
- pos = self._data_stream.rfind("BEGIN")
337
- self._data_stream = self._data_stream[:pos]
338
-
339
- def commit(self):
340
- """
341
- Upload new data to netdata.
342
- """
343
- try:
344
- print(self._data_stream)
345
- except Exception as e:
346
- msg.fatal('cannot send data to netdata:', str(e))
347
- self._data_stream = ""
348
-
349
- # --- ERROR HANDLING ---
350
-
351
- def error(self, *params):
352
- """
353
- Show error message on stderr
354
- """
355
- msg.error(self.chart_name, *params)
356
-
357
- def alert(self, *params):
358
- """
359
- Show error message on stderr
360
- """
361
- msg.alert(self.chart_name, *params)
362
-
363
- def debug(self, *params):
364
- """
365
- Show debug message on stderr
366
- """
367
- msg.debug(self.chart_name, *params)
368
-
369
- def info(self, *params):
370
- """
371
- Show information message on stderr
372
- """
373
- msg.info(self.chart_name, *params)
374
-
375
- # --- MAIN METHODS ---
376
-
377
- def _get_data(self):
378
- """
379
- Get some data
380
- :return: dict
381
- """
382
- return {}
383
-
384
- def check(self):
385
- """
386
- check() prototype
387
- :return: boolean
388
- """
389
- self.debug("Module", str(self.__module__), "doesn't implement check() function. Using default.")
390
- data = self._get_data()
391
-
392
- if data is None:
393
- self.debug("failed to receive data during check().")
394
- return False
395
-
396
- if len(data) == 0:
397
- self.debug("empty data during check().")
398
- return False
399
-
400
- self.debug("successfully received data during check(): '" + str(data) + "'")
401
- return True
402
-
403
- def create(self):
404
- """
405
- Create charts
406
- :return: boolean
407
- """
408
- data = self._data_from_check or self._get_data()
409
- if data is None:
410
- self.debug("failed to receive data during create().")
411
- return False
412
-
413
- idx = 0
414
- for name in self.order:
415
- options = self.definitions[name]['options'] + [self.priority + idx, self.update_every]
416
- self.chart(self.chart_name + "." + name, *options)
417
- # check if server has this datapoint
418
- for line in self.definitions[name]['lines']:
419
- if line[0] in data:
420
- self.dimension(*line)
421
- idx += 1
422
-
423
- self.commit()
424
- return True
425
-
426
- def update(self, interval):
427
- """
428
- Update charts
429
- :param interval: int
430
- :return: boolean
431
- """
432
- data = self._get_data()
433
- if data is None:
434
- self.debug("failed to receive data during update().")
435
- return False
436
-
437
- updated = False
438
- for chart in self.order:
439
- if self.begin(self.chart_name + "." + chart, interval):
440
- updated = True
441
- for dim in self.definitions[chart]['lines']:
442
- try:
443
- self.set(dim[0], data[dim[0]])
444
- except KeyError:
445
- pass
446
- self.end()
447
-
448
- self.commit()
449
- if not updated:
450
- self.error("no charts to update")
451
-
452
- return updated
453
-
454
- @staticmethod
455
- def find_binary(binary):
456
- try:
457
- if isinstance(binary, str):
458
- binary = os.path.basename(binary)
459
- return next(('/'.join([p, binary]) for p in PATH
460
- if os.path.isfile('/'.join([p, binary]))
461
- and os.access('/'.join([p, binary]), os.X_OK)))
462
- return None
463
- except StopIteration:
464
- return None
465
-
466
- def _add_new_dimension(self, dimension_id, chart_name, dimension=None, algorithm='incremental',
467
- multiplier=1, divisor=1, priority=65000):
468
- """
469
- :param dimension_id:
470
- :param chart_name:
471
- :param dimension:
472
- :param algorithm:
473
- :param multiplier:
474
- :param divisor:
475
- :param priority:
476
- :return:
477
- """
478
- if not all([dimension_id not in self._dimensions,
479
- chart_name in self.order,
480
- chart_name in self.definitions]):
481
- return
482
- self._dimensions.append(dimension_id)
483
- dimension_list = list(map(str, [dimension_id,
484
- dimension if dimension else dimension_id,
485
- algorithm,
486
- multiplier,
487
- divisor]))
488
- self.definitions[chart_name]['lines'].append(dimension_list)
489
- add_to_name = self.override_name or self.name
490
- job_name = ('_'.join([self.__module__, re.sub('\s+', '_', add_to_name)])
491
- if add_to_name != 'None' else self.__module__)
492
- chart = 'CHART {0}.{1} '.format(job_name, chart_name)
493
- options = '"" "{0}" {1} "{2}" {3} {4} '.format(*self.definitions[chart_name]['options'][1:6])
494
- other = '{0} {1}\n'.format(priority, self.update_every)
495
- new_dimension = "DIMENSION {0}\n".format(' '.join(dimension_list))
496
- print(chart + options + other + new_dimension)
497
-
498
-
499
-class UrlService(SimpleService):
500
- def __init__(self, configuration=None, name=None):
501
- SimpleService.__init__(self, configuration=configuration, name=name)
502
- self.url = self.configuration.get('url')
503
- self.user = self.configuration.get('user')
504
- self.password = self.configuration.get('pass')
505
- self.proxy_user = self.configuration.get('proxy_user')
506
- self.proxy_password = self.configuration.get('proxy_pass')
507
- self.proxy_url = self.configuration.get('proxy_url')
508
- self.header = self.configuration.get('header')
509
- self._manager = None
510
-
511
- def __make_headers(self, **header_kw):
512
- user = header_kw.get('user') or self.user
513
- password = header_kw.get('pass') or self.password
514
- proxy_user = header_kw.get('proxy_user') or self.proxy_user
515
- proxy_password = header_kw.get('proxy_pass') or self.proxy_password
516
- custom_header = header_kw.get('header') or self.header
517
- header_params = dict(keep_alive=True)
518
- proxy_header_params = dict()
519
- if user and password:
520
- header_params['basic_auth'] = '{user}:{password}'.format(user=user,
521
- password=password)
522
- if proxy_user and proxy_password:
523
- proxy_header_params['proxy_basic_auth'] = '{user}:{password}'.format(user=proxy_user,
524
- password=proxy_password)
525
- try:
526
- header, proxy_header = urllib3.make_headers(**header_params), urllib3.make_headers(**proxy_header_params)
527
- except TypeError as error:
528
- self.error('build_header() error: {error}'.format(error=error))
529
- return None, None
530
- else:
531
- header.update(custom_header or dict())
532
- return header, proxy_header
533
-
534
- def _build_manager(self, **header_kw):
535
- header, proxy_header = self.__make_headers(**header_kw)
536
- if header is None or proxy_header is None:
537
- return None
538
- proxy_url = header_kw.get('proxy_url') or self.proxy_url
539
- if proxy_url:
540
- manager = urllib3.ProxyManager
541
- params = dict(proxy_url=proxy_url, headers=header, proxy_headers=proxy_header)
542
- else:
543
- manager = urllib3.PoolManager
544
- params = dict(headers=header)
545
- try:
546
- return manager(**params)
547
- except (urllib3.exceptions.ProxySchemeUnknown, TypeError) as error:
548
- self.error('build_manager() error:', str(error))
549
- return None
550
-
551
- def _get_raw_data(self, url=None, manager=None):
552
- """
553
- Get raw data from http request
554
- :return: str
555
- """
556
- try:
557
- url = url or self.url
558
- manager = manager or self._manager
559
- # TODO: timeout, retries and method hardcoded..
560
- response = manager.request(method='GET',
561
- url=url,
562
- timeout=1,
563
- retries=1,
564
- headers=manager.headers)
565
- except (urllib3.exceptions.HTTPError, TypeError, AttributeError) as error:
566
- self.error('Url: {url}. Error: {error}'.format(url=url, error=error))
567
- return None
568
- if response.status == 200:
569
- return response.data.decode()
570
- self.debug('Url: {url}. Http response status code: {code}'.format(url=url, code=response.status))
571
- return None
572
-
573
- def check(self):
574
- """
575
- Format configuration data and try to connect to server
576
- :return: boolean
577
- """
578
- if not (self.url and isinstance(self.url, str)):
579
- self.error('URL is not defined or type is not <str>')
580
- return False
581
-
582
- self._manager = self._build_manager()
583
- if not self._manager:
584
- return False
585
-
586
- try:
587
- data = self._get_data()
588
- except Exception as error:
589
- self.error('_get_data() failed. Url: {url}. Error: {error}'.format(url=self.url, error=error))
590
- return False
591
-
592
- if isinstance(data, dict) and data:
593
- self._data_from_check = data
594
- return True
595
- self.error('_get_data() returned no data or type is not <dict>')
596
- return False
597
-
598
-
599
-class SocketService(SimpleService):
600
- def __init__(self, configuration=None, name=None):
601
- self._sock = None
602
- self._keep_alive = False
603
- self.host = "localhost"
604
- self.port = None
605
- self.unix_socket = None
606
- self.request = ""
607
- self.__socket_config = None
608
- self.__empty_request = "".encode()
609
- SimpleService.__init__(self, configuration=configuration, name=name)
610
-
611
- def _socketerror(self, message=None):
612
- if self.unix_socket is not None:
613
- self.error("unix socket '" + self.unix_socket + "':", message)
614
- else:
615
- if self.__socket_config is not None:
616
- af, socktype, proto, canonname, sa = self.__socket_config
617
- self.error("socket to '" + str(sa[0]) + "' port " + str(sa[1]) + ":", message)
618
- else:
619
- self.error("unknown socket:", message)
620
-
621
- def _connect2socket(self, res=None):
622
- """
623
- Connect to a socket, passing the result of getaddrinfo()
624
- :return: boolean
625
- """
626
- if res is None:
627
- res = self.__socket_config
628
- if res is None:
629
- self.error("Cannot create socket to 'None':")
630
- return False
631
-
632
- af, socktype, proto, canonname, sa = res
633
- try:
634
- self.debug("creating socket to '" + str(sa[0]) + "', port " + str(sa[1]))
635
- self._sock = socket.socket(af, socktype, proto)
636
- except socket.error as e:
637
- self.error("Failed to create socket to '" + str(sa[0]) + "', port " + str(sa[1]) + ":", str(e))
638
- self._sock = None
639
- self.__socket_config = None
640
- return False
641
-
642
- try:
643
- self.debug("connecting socket to '" + str(sa[0]) + "', port " + str(sa[1]))
644
- self._sock.connect(sa)
645
- except socket.error as e:
646
- self.error("Failed to connect to '" + str(sa[0]) + "', port " + str(sa[1]) + ":", str(e))
647
- self._disconnect()
648
- self.__socket_config = None
649
- return False
650
-
651
- self.debug("connected to '" + str(sa[0]) + "', port " + str(sa[1]))
652
- self.__socket_config = res
653
- return True
654
-
655
- def _connect2unixsocket(self):
656
- """
657
- Connect to a unix socket, given its filename
658
- :return: boolean
659
- """
660
- if self.unix_socket is None:
661
- self.error("cannot connect to unix socket 'None'")
662
- return False
663
-
664
- try:
665
- self.debug("attempting DGRAM unix socket '" + str(self.unix_socket) + "'")
666
- self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
667
- self._sock.connect(self.unix_socket)
668
- self.debug("connected DGRAM unix socket '" + str(self.unix_socket) + "'")
669
- return True
670
- except socket.error as e:
671
- self.debug("Failed to connect DGRAM unix socket '" + str(self.unix_socket) + "':", str(e))
672
-
673
- try:
674
- self.debug("attempting STREAM unix socket '" + str(self.unix_socket) + "'")
675
- self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
676
- self._sock.connect(self.unix_socket)
677
- self.debug("connected STREAM unix socket '" + str(self.unix_socket) + "'")
678
- return True
679
- except socket.error as e:
680
- self.debug("Failed to connect STREAM unix socket '" + str(self.unix_socket) + "':", str(e))
681
- self.error("Failed to connect to unix socket '" + str(self.unix_socket) + "':", str(e))
682
- self._sock = None
683
- return False
684
-
685
- def _connect(self):
686
- """
687
- Recreate socket and connect to it since sockets cannot be reused after closing
688
- Available configurations are IPv6, IPv4 or UNIX socket
689
- :return:
690
- """
691
- try:
692
- if self.unix_socket is not None:
693
- self._connect2unixsocket()
694
-
695
- else:
696
- if self.__socket_config is not None:
697
- self._connect2socket()
698
- else:
699
- for res in socket.getaddrinfo(self.host, self.port, socket.AF_UNSPEC, socket.SOCK_STREAM):
700
- if self._connect2socket(res): break
701
-
702
- except Exception as e:
703
- self._sock = None
704
- self.__socket_config = None
705
-
706
- if self._sock is not None:
707
- self._sock.setblocking(0)
708
- self._sock.settimeout(5)
709
- self.debug("set socket timeout to: " + str(self._sock.gettimeout()))
710
-
711
- def _disconnect(self):
712
- """
713
- Close socket connection
714
- :return:
715
- """
716
- if self._sock is not None:
717
- try:
718
- self.debug("closing socket")
719
- self._sock.shutdown(2) # 0 - read, 1 - write, 2 - all
720
- self._sock.close()
721
- except Exception:
722
- pass
723
- self._sock = None
724
-
725
- def _send(self):
726
- """
727
- Send request.
728
- :return: boolean
729
- """
730
- # Send request if it is needed
731
- if self.request != self.__empty_request:
732
- try:
733
- self.debug("sending request:", str(self.request))
734
- self._sock.send(self.request)
735
- except Exception as e:
736
- self._socketerror("error sending request:" + str(e))
737
- self._disconnect()
738
- return False
739
- return True
740
-
741
- def _receive(self):
742
- """
743
- Receive data from socket
744
- :return: str
745
- """
746
- data = ""
747
- while True:
748
- self.debug("receiving response")
749
- try:
750
- buf = self._sock.recv(4096)
751
- except Exception as e:
752
- self._socketerror("failed to receive response:" + str(e))
753
- self._disconnect()
754
- break
755
-
756
- if buf is None or len(buf) == 0: # handle server disconnect
757
- if data == "":
758
- self._socketerror("unexpectedly disconnected")
759
- else:
760
- self.debug("server closed the connection")
761
- self._disconnect()
762
- break
763
-
764
- self.debug("received data:", str(buf))
765
- data += buf.decode('utf-8', 'ignore')
766
- if self._check_raw_data(data):
767
- break
768
-
769
- self.debug("final response:", str(data))
770
- return data
771
-
772
- def _get_raw_data(self):
773
- """
774
- Get raw data with low-level "socket" module.
775
- :return: str
776
- """
777
- if self._sock is None:
778
- self._connect()
779
- if self._sock is None:
780
- return None
781
-
782
- # Send request if it is needed
783
- if not self._send():
784
- return None
785
-
786
- data = self._receive()
787
-
788
- if not self._keep_alive:
789
- self._disconnect()
790
-
791
- return data
792
-
793
- def _check_raw_data(self, data):
794
- """
795
- Check if all data has been gathered from socket
796
- :param data: str
797
- :return: boolean
798
- """
799
- return True
800
-
801
- def _parse_config(self):
802
- """
803
- Parse configuration data
804
- :return: boolean
805
- """
806
- if self.name is None or self.name == str(None):
807
- self.name = ""
808
- else:
809
- self.name = str(self.name)
810
-
811
- try:
812
- self.unix_socket = str(self.configuration['socket'])
813
- except (KeyError, TypeError):
814
- self.debug("No unix socket specified. Trying TCP/IP socket.")
815
- self.unix_socket = None
816
- try:
817
- self.host = str(self.configuration['host'])
818
- except (KeyError, TypeError):
819
- self.debug("No host specified. Using: '" + self.host + "'")
820
- try:
821
- self.port = int(self.configuration['port'])
822
- except (KeyError, TypeError):
823
- self.debug("No port specified. Using: '" + str(self.port) + "'")
824
-
825
- try:
826
- self.request = str(self.configuration['request'])
827
- except (KeyError, TypeError):
828
- self.debug("No request specified. Using: '" + str(self.request) + "'")
829
-
830
- self.request = self.request.encode()
831
-
832
- def check(self):
833
- self._parse_config()
834
- return SimpleService.check(self)
835
-
836
-
837
-class LogService(SimpleService):
838
- def __init__(self, configuration=None, name=None):
839
- SimpleService.__init__(self, configuration=configuration, name=name)
840
- self.log_path = self.configuration.get('path')
841
- self.__glob_path = self.log_path
842
- self._last_position = 0
843
- self.retries = 100000 # basically always retry
844
- self.__re_find = dict(current=0, run=0, maximum=60)
845
-
846
- def _get_raw_data(self):
847
- """
848
- Get log lines since last poll
849
- :return: list
850
- """
851
- lines = list()
852
- try:
853
- if self.__re_find['current'] == self.__re_find['run']:
854
- self._find_recent_log_file()
855
- size = os.path.getsize(self.log_path)
856
- if size == self._last_position:
857
- self.__re_find['current'] += 1
858
- return list() # return empty list if nothing has changed
859
- elif size < self._last_position:
860
- self._last_position = 0 # read from beginning if file has shrunk
861
-
862
- with open(self.log_path) as fp:
863
- fp.seek(self._last_position)
864
- for line in fp:
865
- lines.append(line)
866
- self._last_position = fp.tell()
867
- self.__re_find['current'] = 0
868
- except (OSError, IOError) as error:
869
- self.__re_find['current'] += 1
870
- self.error(str(error))
871
-
872
- return lines or None
873
-
874
- def _find_recent_log_file(self):
875
- """
876
- :return:
877
- """
878
- self.__re_find['run'] = self.__re_find['maximum']
879
- self.__re_find['current'] = 0
880
- self.__glob_path = self.__glob_path or self.log_path # workaround for modules w/o config files
881
- path_list = glob(self.__glob_path)
882
- if path_list:
883
- self.log_path = max(path_list)
884
- return True
885
- return False
886
-
887
- def check(self):
888
- """
889
- Parse basic configuration and check if log file exists
890
- :return: boolean
891
- """
892
- if not self.log_path:
893
- self.error("No path to log specified")
894
- return None
895
-
896
- if all([self._find_recent_log_file(),
897
- os.access(self.log_path, os.R_OK),
898
- os.path.isfile(self.log_path)]):
899
- return True
900
- self.error("Cannot access %s" % self.log_path)
901
- return False
902
-
903
- def create(self):
904
- # set cursor at last byte of log file
905
- self._last_position = os.path.getsize(self.log_path)
906
- status = SimpleService.create(self)
907
- # self._last_position = 0
908
- return status
909
-
910
-
911
-class ExecutableService(SimpleService):
912
-
913
- def __init__(self, configuration=None, name=None):
914
- SimpleService.__init__(self, configuration=configuration, name=name)
915
- self.command = None
916
-
917
- def _get_raw_data(self, stderr=False):
918
- """
919
- Get raw data from executed command
920
- :return: <list>
921
- """
922
- try:
923
- p = Popen(self.command, stdout=PIPE, stderr=PIPE)
924
- except Exception as error:
925
- self.error("Executing command", " ".join(self.command), "resulted in error:", str(error))
926
- return None
927
- data = list()
928
- std = p.stderr if stderr else p.stdout
929
- for line in std.readlines():
930
- data.append(line.decode())
931
-
932
- return data or None
933
-
934
- def check(self):
935
- """
936
- Parse basic configuration, check if command is whitelisted and is returning values
937
- :return: <boolean>
938
- """
939
- # Preference: 1. "command" from configuration file 2. "command" from plugin (if specified)
940
- if 'command' in self.configuration:
941
- self.command = self.configuration['command']
942
-
943
- # "command" must be: 1.not None 2. type <str>
944
- if not (self.command and isinstance(self.command, str)):
945
- self.error('Command is not defined or command type is not <str>')
946
- return False
947
-
948
- # Split "command" into: 1. command <str> 2. options <list>
949
- command, opts = self.command.split()[0], self.command.split()[1:]
950
-
951
- # Check for "bad" symbols in options. No pipes, redirects etc. TODO: what is missing?
952
- bad_opts = set(''.join(opts)) & set(['&', '|', ';', '>', '<'])
953
- if bad_opts:
954
- self.error("Bad command argument(s): %s" % bad_opts)
955
- return False
956
-
957
- # Find absolute path ('echo' => '/bin/echo')
958
- if '/' not in command:
959
- command = self.find_binary(command)
960
- if not command:
961
- self.error('Can\'t locate "%s" binary in PATH(%s)' % (self.command, PATH))
962
- return False
963
- # Check if binary exist and executable
964
- else:
965
- if not (os.path.isfile(command) and os.access(command, os.X_OK)):
966
- self.error('"%s" is not a file or not executable' % command)
967
- return False
968
-
969
- self.command = [command] + opts if opts else [command]
970
-
971
- try:
972
- data = self._get_data()
973
- except Exception as error:
974
- self.error('_get_data() failed. Command: %s. Error: %s' % (self.command, error))
975
- return False
976
-
977
- if isinstance(data, dict) and data:
978
- # We need this for create() method. No reason to execute get_data() again if result is not empty dict()
979
- self._data_from_check = data
980
- return True
981
- else:
982
- self.error("Command", str(self.command), "returned no data")
983
- return False
984
-
985
-
986
-class MySQLService(SimpleService):
987
-
988
- def __init__(self, configuration=None, name=None):
989
- SimpleService.__init__(self, configuration=configuration, name=name)
990
- self.__connection = None
991
- self.__conn_properties = dict()
992
- self.extra_conn_properties = dict()
993
- self.__queries = self.configuration.get('queries', dict())
994
- self.queries = dict()
995
-
996
- def __connect(self):
997
- try:
998
- connection = MySQLdb.connect(connect_timeout=self.update_every, **self.__conn_properties)
999
- except (MySQLdb.MySQLError, TypeError, AttributeError) as error:
1000
- return None, str(error)
1001
- else:
1002
- return connection, None
1003
-
1004
- def check(self):
1005
- def get_connection_properties(conf, extra_conf):
1006
- properties = dict()
1007
- if conf.get('user'):
1008
- properties['user'] = conf['user']
1009
- if conf.get('pass'):
1010
- properties['passwd'] = conf['pass']
1011
- if conf.get('socket'):
1012
- properties['unix_socket'] = conf['socket']
1013
- elif conf.get('host'):
1014
- properties['host'] = conf['host']
1015
- properties['port'] = int(conf.get('port', 3306))
1016
- elif conf.get('my.cnf'):
1017
- if MySQLdb.__name__ == 'pymysql':
1018
- self.error('"my.cnf" parsing is not working for pymysql')
1019
- else:
1020
- properties['read_default_file'] = conf['my.cnf']
1021
- if isinstance(extra_conf, dict) and extra_conf:
1022
- properties.update(extra_conf)
1023
-
1024
- return properties or None
1025
-
1026
- def is_valid_queries_dict(raw_queries, log_error):
1027
- """
1028
- :param raw_queries: dict:
1029
- :param log_error: function:
1030
- :return: dict or None
1031
-
1032
- raw_queries is valid when: type <dict> and not empty after is_valid_query(for all queries)
1033
- """
1034
- def is_valid_query(query):
1035
- return all([isinstance(query, str),
1036
- query.startswith(('SELECT', 'select', 'SHOW', 'show'))])
1037
-
1038
- if hasattr(raw_queries, 'keys') and raw_queries:
1039
- valid_queries = dict([(n, q) for n, q in raw_queries.items() if is_valid_query(q)])
1040
- bad_queries = set(raw_queries) - set(valid_queries)
1041
-
1042
- if bad_queries:
1043
- log_error('Removed query(s): %s' % bad_queries)
1044
- return valid_queries
1045
- else:
1046
- log_error('Unsupported "queries" format. Must be not empty <dict>')
1047
- return None
1048
-
1049
- if not PY_MYSQL:
1050
- self.error('MySQLdb or PyMySQL module is needed to use mysql.chart.py plugin')
1051
- return False
1052
-
1053
- # Preference: 1. "queries" from the configuration file 2. "queries" from the module
1054
- self.queries = self.__queries or self.queries
1055
- # Check if "self.queries" exist, not empty and all queries are in valid format
1056
- self.queries = is_valid_queries_dict(self.queries, self.error)
1057
- if not self.queries:
1058
- return None
1059
-
1060
- # Get connection properties
1061
- self.__conn_properties = get_connection_properties(self.configuration, self.extra_conn_properties)
1062
- if not self.__conn_properties:
1063
- self.error('Connection properties are missing')
1064
- return False
1065
-
1066
- # Create connection to the database
1067
- self.__connection, error = self.__connect()
1068
- if error:
1069
- self.error('Can\'t establish connection to MySQL: %s' % error)
1070
- return False
1071
-
1072
- try:
1073
- data = self._get_data()
1074
- except Exception as error:
1075
- self.error('_get_data() failed. Error: %s' % error)
1076
- return False
1077
-
1078
- if isinstance(data, dict) and data:
1079
- # We need this for create() method
1080
- self._data_from_check = data
1081
- return True
1082
- else:
1083
- self.error("_get_data() returned no data or type is not <dict>")
1084
- return False
1085
-
1086
- def _get_raw_data(self, description=None):
1087
- """
1088
- Get raw data from MySQL server
1089
- :return: dict: fetchall() or (fetchall(), description)
1090
- """
1091
-
1092
- if not self.__connection:
1093
- self.__connection, error = self.__connect()
1094
- if error:
1095
- return None
1096
-
1097
- raw_data = dict()
1098
- queries = dict(self.queries)
1099
- try:
1100
- with self.__connection as cursor:
1101
- for name, query in queries.items():
1102
- try:
1103
- cursor.execute(query)
1104
- except (MySQLdb.ProgrammingError, MySQLdb.OperationalError) as error:
1105
- if self.__is_error_critical(err_class=exc_info()[0], err_text=str(error)):
1106
- raise RuntimeError
1107
- self.error('Removed query: %s[%s]. Error: %s'
1108
- % (name, query, error))
1109
- self.queries.pop(name)
1110
- continue
1111
- else:
1112
- raw_data[name] = (cursor.fetchall(), cursor.description) if description else cursor.fetchall()
1113
- self.__connection.commit()
1114
- except (MySQLdb.MySQLError, RuntimeError, TypeError, AttributeError):
1115
- self.__connection.close()
1116
- self.__connection = None
1117
- return None
1118
- else:
1119
- return raw_data or None
1120
-
1121
- @staticmethod
1122
- def __is_error_critical(err_class, err_text):
1123
- return err_class == MySQLdb.OperationalError and all(['denied' not in err_text,
1124
- 'Unknown column' not in err_text])
2
+# Description: backward compatibility with old version
3
+
4
+from bases.FrameworkServices.SimpleService import SimpleService
5
+from bases.FrameworkServices.UrlService import UrlService
6
+from bases.FrameworkServices.SocketService import SocketService
7
+from bases.FrameworkServices.LogService import LogService
8
+from bases.FrameworkServices.ExecutableService import ExecutableService
9
+from bases.FrameworkServices.MySQLService import MySQLService