move new SimpleService to a separate file
lgz committed
Oct 13, 2017 at 00:19 UTC
6282847e7d3b263c1829c7523d71a49ab502f871
1 file changed
+235
python.d/python_modules/bases/FrameworkServices/SimpleService.py
new
+235
@@ -0,0 +1,235 @@
1
+# -*- coding: utf-8 -*-
2
+# Description:
3
+# Author: Pawel Krupa (paulfantom)
4
+# Author: Ilya Mashchenko (l2isbad)
5
+
6
+from threading import Thread
7
+
8
+try:
9
+ from time import sleep, monotonic as time
10
+except ImportError:
11
+ from time import sleep, time
12
+
13
+from bases.charts import Charts, create_runtime_chart
14
+from bases.collection import OldVersionCompatibility, UsefulFuncs, safe_print
15
+from bases.loggers import PythonDLimitedLogger
16
+
17
+CHART_OBSOLETE_PENALTY = 10
18
+
19
+START_MSG = 'STARTED. Update frequency: {freq}, retries: {retries}.'
20
+UPDATE_MSG = 'UPDATE {status}. Elapsed time: {elapsed}, retries left: {retries}.'
21
+STOP_MSG = 'STOPPED after {retries_max} data collection failures in a row.'
22
+SLEEP_MSG = 'SLEEPING for {sleep_time} to reach frequency of {freq} sec.'
23
+
24
+RUNTIME_CHART_UPDATE = 'BEGIN netdata.runtime_{job_name} {since_last}\n' \
25
+ 'SET run_time = {elapsed}\n' \
26
+ 'END\n'
27
+
28
+
29
+class RuntimeCounters:
30
+ def __init__(self, configuration):
31
+ """
32
+ :param configuration: <dict>
33
+ """
34
+ self.FREQ = int(configuration.pop('update_every', 1))
35
+ self.START_RUN = 0
36
+ self.NEXT_RUN = 0
37
+ self.PREV_UPDATE = 0
38
+ self.SINCE_UPDATE = 0
39
+ self.ELAPSED = 0
40
+ self.RETRIES = 0
41
+ self.RETRIES_MAX = configuration.pop('retries', 10)
42
+
43
+ def is_sleep_time(self):
44
+ return self.START_RUN < self.NEXT_RUN
45
+
46
+
47
+class SimpleService(Thread, PythonDLimitedLogger, OldVersionCompatibility, object):
48
+ """
49
+ Prototype of Service class.
50
+ Implemented basic functionality to run jobs by `python.d.plugin`
51
+ """
52
+ def __init__(self, configuration, name=''):
53
+ """
54
+ :param configuration: <dict>
55
+ :param name: <str>
56
+ """
57
+ Thread.__init__(self)
58
+ self.daemon = True
59
+ PythonDLimitedLogger.__init__(self)
60
+ OldVersionCompatibility.__init__(self)
61
+ self.configuration = configuration
62
+ self.module_name = configuration.pop('module_name')
63
+ self.job_name = configuration.pop('job_name')
64
+ self.override_name = configuration.pop('override_name')
65
+ self.fake_name = None
66
+ self.order = list()
67
+ self.definitions = dict()
68
+ self._runtime_counters = RuntimeCounters(configuration=configuration)
69
+ self.charts = Charts(job_name=self.actual_name,
70
+ priority=configuration.pop('priority', 60000),
71
+ update_every=self._runtime_counters.FREQ)
72
+ self.functions = UsefulFuncs()
73
+
74
+ def __repr__(self):
75
+ return '<{cls_bases}: {name}>'.format(cls_bases=', '.join(c.__name__ for c in self.__class__.__bases__),
76
+ name=self.name)
77
+
78
+ @property
79
+ def name(self):
80
+ if self.job_name:
81
+ return '_'.join([self.module_name, self.override_name or self.job_name])
82
+ return self.module_name
83
+
84
+ @property
85
+ def actual_name(self):
86
+ return self.fake_name or self.name
87
+
88
+ @property
89
+ def update_every(self):
90
+ return self._runtime_counters.FREQ
91
+
92
+ @update_every.setter
93
+ def update_every(self, value):
94
+ """
95
+ :param value: <int>
96
+ :return:
97
+ """
98
+ self._runtime_counters.FREQ = value
99
+
100
+ def check(self):
101
+ """
102
+ check() prototype
103
+ :return: boolean
104
+ """
105
+ try:
106
+ data = self._get_data()
107
+ except Exception as error:
108
+ self.debug('CHECK {{error: {error}}}'.format(error=error))
109
+ else:
110
+ if data and isinstance(data, dict):
111
+ return True
112
+ self.debug('CHECK returned no data')
113
+ return False
114
+
115
+ @create_runtime_chart
116
+ def create(self):
117
+ for chart_name in self.order:
118
+ chart_config = self.definitions.get(chart_name)
119
+ if not chart_config:
120
+ self.debug('{chart_name} not in definitions'.format(chart_name=chart_name))
121
+ continue
122
+
123
+ chart_params = ([chart_name] + chart_config['options'])
124
+ ok = self.charts.add_chart(params=chart_params)
125
+ if not ok:
126
+ self.debug('"{chart}" chart no added'.format(chart=chart_name))
127
+ continue
128
+
129
+ for line in chart_config['lines']:
130
+ self.charts[chart_name].add_dimension(line)
131
+
132
+ del self.order
133
+ del self.definitions
134
+
135
+ if self.charts.empty():
136
+ return None
137
+
138
+ for chart in self.charts:
139
+ safe_print(chart.create())
140
+ return True
141
+
142
+ def run(self):
143
+ """
144
+ Runs job in thread. Handles retries.
145
+ Exits when job failed or timed out.
146
+ :return: None
147
+ """
148
+ job = self._runtime_counters
149
+ self.debug(START_MSG.format(freq=job.FREQ, retries=job.RETRIES_MAX - job.RETRIES))
150
+
151
+ while True:
152
+ job.START_RUN = time()
153
+
154
+ job.NEXT_RUN = job.START_RUN - (job.START_RUN % job.FREQ) + job.FREQ
155
+
156
+ self.sleep_until_next_run()
157
+
158
+ if job.PREV_UPDATE:
159
+ job.SINCE_UPDATE = int((job.START_RUN - job.PREV_UPDATE) * 1e6)
160
+
161
+ try:
162
+ updated = self.update(interval=job.SINCE_UPDATE)
163
+ except Exception as error:
164
+ print(error)
165
+ updated = False
166
+
167
+ if not updated:
168
+ if not self.manage_retries():
169
+ return
170
+ else:
171
+ current_time = time()
172
+ job.PREV_UPDATE = current_time
173
+ job.ELAPSED = int((current_time - job.START_RUN) * 1e3)
174
+ job.RETRIES = 0
175
+ self.safe_print(RUNTIME_CHART_UPDATE.format(job_name=self.name,
176
+ since_last=job.SINCE_UPDATE,
177
+ elapsed=job.ELAPSED))
178
+ self.debug(UPDATE_MSG.format(status='OK' if updated else 'FAILED',
179
+ elapsed=job.ELAPSED if updated else '-',
180
+ retries=job.RETRIES_MAX - job.RETRIES))
181
+
182
+ def update(self, interval):
183
+ """
184
+ :param interval: <int>
185
+ :return:
186
+ """
187
+ data = self.get_data()
188
+ if not data:
189
+ return None
190
+ charts_updated = False
191
+
192
+ for chart in self.charts.penalty_exceeded(penalty_max=CHART_OBSOLETE_PENALTY):
193
+ safe_print(chart.obsolete())
194
+ del self.charts[chart.params['id']]
195
+
196
+ for chart in self.charts:
197
+ dimension_updated = str()
198
+ for dimension in chart.dimensions:
199
+ try:
200
+ value = int(data[dimension.params['id']])
201
+ except (KeyError, TypeError):
202
+ continue
203
+ dimension_updated += dimension.set(value)
204
+
205
+ if dimension_updated:
206
+ charts_updated = True
207
+ safe_print(''.join([chart.begin(since_last=interval),
208
+ dimension_updated, 'END\n']))
209
+ else:
210
+ chart.penalty += 1
211
+ return charts_updated
212
+
213
+ def manage_retries(self):
214
+ self._runtime_counters.RETRIES += 1
215
+ if self._runtime_counters.RETRIES >= self._runtime_counters.RETRIES_MAX:
216
+ self.error(STOP_MSG.format(retries_max=self._runtime_counters.RETRIES_MAX))
217
+ return False
218
+ return True
219
+
220
+ def sleep_until_next_run(self):
221
+ job = self._runtime_counters
222
+
223
+ # sleep() is interruptable
224
+ while job.is_sleep_time():
225
+ sleep_time = job.NEXT_RUN - job.START_RUN
226
+ self.debug(SLEEP_MSG.format(sleep_time=sleep_time,
227
+ freq=job.FREQ))
228
+ sleep(sleep_time)
229
+ job.START_RUN = time()
230
+
231
+ def get_data(self):
232
+ return self._get_data()
233
+
234
+ def _get_data(self):
235
+ raise NotImplemented