master
c 271 lines 9.2 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #define PULSE_INTERNALS 1
4 #include "daemon/common.h"
5
6 #define WORKER_JOB_DAEMON 0
7 #define WORKER_JOB_SQLITE3 1
8 #define WORKER_JOB_HTTP_API 2
9 #define WORKER_JOB_QUERIES 3
10 #define WORKER_JOB_INGESTION 4
11 #define WORKER_JOB_DBENGINE 5
12 #define WORKER_JOB_STRINGS 6
13 #define WORKER_JOB_DICTIONARIES 7
14 #define WORKER_JOB_ML 8
15 #define WORKER_JOB_GORILLA 9
16 #define WORKER_JOB_HEARTBEAT 10
17 #define WORKER_JOB_WORKERS 11
18 #define WORKER_JOB_MALLOC_TRACE 12
19 #define WORKER_JOB_REGISTRY 13
20 #define WORKER_JOB_ARAL 14
21 #define WORKER_JOB_NETWORK 15
22 #define WORKER_JOB_PARENTS 16
23 #define WORKER_JOB_MEMORY_EXTENDED 17
24
25 #if WORKER_UTILIZATION_MAX_JOB_TYPES < 17
26 #error "WORKER_UTILIZATION_MAX_JOB_TYPES has to be at least 14"
27 #endif
28
29 bool pulse_enabled = true;
30 bool pulse_extended_enabled = false;
31
32 static void pulse_register_workers(void) {
33 worker_register("PULSE");
34
35 worker_register_job_name(WORKER_JOB_DAEMON, "daemon");
36 worker_register_job_name(WORKER_JOB_SQLITE3, "sqlite3");
37 worker_register_job_name(WORKER_JOB_HTTP_API, "http-api");
38 worker_register_job_name(WORKER_JOB_QUERIES, "queries");
39 worker_register_job_name(WORKER_JOB_INGESTION, "ingestion");
40 worker_register_job_name(WORKER_JOB_DBENGINE, "dbengine");
41 worker_register_job_name(WORKER_JOB_STRINGS, "strings");
42 worker_register_job_name(WORKER_JOB_DICTIONARIES, "dictionaries");
43 worker_register_job_name(WORKER_JOB_ML, "ML");
44 worker_register_job_name(WORKER_JOB_GORILLA, "gorilla");
45 worker_register_job_name(WORKER_JOB_HEARTBEAT, "heartbeat");
46 worker_register_job_name(WORKER_JOB_WORKERS, "workers");
47 worker_register_job_name(WORKER_JOB_MALLOC_TRACE, "malloc trace");
48 worker_register_job_name(WORKER_JOB_REGISTRY, "registry");
49 worker_register_job_name(WORKER_JOB_ARAL, "aral");
50 worker_register_job_name(WORKER_JOB_NETWORK, "network");
51 worker_register_job_name(WORKER_JOB_PARENTS, "parents");
52 worker_register_job_name(WORKER_JOB_MEMORY_EXTENDED, "memory extended");
53 }
54
55 void pulse_thread_main(void *ptr) {
56 struct netdata_static_thread *static_thread = ptr;
57 pulse_register_workers();
58
59 int update_every =
60 (int)inicfg_get_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", localhost->rrd_update_every);
61 if (update_every < localhost->rrd_update_every) {
62 update_every = localhost->rrd_update_every;
63 inicfg_set_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", update_every);
64 }
65
66 pulse_aral_init();
67 aclk_time_histogram_init();
68
69 usec_t step = update_every * USEC_PER_SEC;
70 heartbeat_t hb;
71 heartbeat_init(&hb, USEC_PER_SEC);
72 usec_t real_step = USEC_PER_SEC;
73
74 // keep the randomness at zero
75 // to make sure we are not close to any other thread
76 hb.randomness = 0;
77
78 while (service_running(SERVICE_COLLECTORS)) {
79 worker_is_idle();
80 heartbeat_next(&hb);
81 if (real_step < step) {
82 real_step += USEC_PER_SEC;
83 continue;
84 }
85 real_step = USEC_PER_SEC;
86
87 worker_is_busy(WORKER_JOB_INGESTION);
88 pulse_ingestion_do(pulse_extended_enabled);
89
90 worker_is_busy(WORKER_JOB_HTTP_API);
91 pulse_web_do(pulse_extended_enabled);
92
93 worker_is_busy(WORKER_JOB_QUERIES);
94 pulse_queries_do(pulse_extended_enabled);
95
96 worker_is_busy(WORKER_JOB_NETWORK);
97 pulse_network_do(pulse_extended_enabled);
98
99 worker_is_busy(WORKER_JOB_ML);
100 pulse_ml_do(pulse_extended_enabled);
101
102 worker_is_busy(WORKER_JOB_GORILLA);
103 pulse_gorilla_do(pulse_extended_enabled);
104
105 worker_is_busy(WORKER_JOB_HEARTBEAT);
106 pulse_heartbeat_do(pulse_extended_enabled);
107
108 #ifdef ENABLE_DBENGINE
109 if(dbengine_enabled) {
110 worker_is_busy(WORKER_JOB_DBENGINE);
111 pulse_dbengine_do(pulse_extended_enabled);
112 dbengine_retention_statistics(pulse_extended_enabled);
113 }
114 #endif
115
116 worker_is_busy(WORKER_JOB_REGISTRY);
117 registry_statistics();
118
119 worker_is_busy(WORKER_JOB_STRINGS);
120 pulse_string_do(pulse_extended_enabled);
121
122 #ifdef DICT_WITH_STATS
123 worker_is_busy(WORKER_JOB_DICTIONARIES);
124 pulse_dictionary_do(pulse_extended_enabled);
125 #endif
126
127 worker_is_busy(WORKER_JOB_ARAL);
128 pulse_aral_do(pulse_extended_enabled);
129
130 worker_is_busy(WORKER_JOB_PARENTS);
131 pulse_parents_do(pulse_extended_enabled);
132
133 // keep this last to have access to the memory counters
134 // exposed by everyone else
135 worker_is_busy(WORKER_JOB_DAEMON);
136 pulse_daemon_do(pulse_extended_enabled);
137 }
138
139 static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
140 worker_unregister();
141 static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
142 }
143
144 // ---------------------------------------------------------------------------------------------------------------------
145 // pulse sqlite3 thread
146
147 void pulse_thread_sqlite3_main(void *ptr) {
148 struct netdata_static_thread *static_thread = ptr;
149 pulse_register_workers();
150
151 int update_every =
152 (int)inicfg_get_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", localhost->rrd_update_every);
153 if (update_every < localhost->rrd_update_every) {
154 update_every = localhost->rrd_update_every;
155 inicfg_set_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", update_every);
156 }
157
158 usec_t step = update_every * USEC_PER_SEC;
159 heartbeat_t hb;
160 heartbeat_init(&hb, USEC_PER_SEC);
161 usec_t real_step = USEC_PER_SEC;
162
163 // keep the randomness at zero
164 // to make sure we are not close to any other thread
165 hb.randomness = 0;
166
167 while (service_running(SERVICE_COLLECTORS)) {
168 worker_is_idle();
169 heartbeat_next(&hb);
170 if (real_step < step) {
171 real_step += USEC_PER_SEC;
172 continue;
173 }
174 real_step = USEC_PER_SEC;
175
176 worker_is_busy(WORKER_JOB_SQLITE3);
177 pulse_sqlite3_do(pulse_extended_enabled);
178 }
179
180 static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
181 worker_unregister();
182 static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
183 }
184
185 // ---------------------------------------------------------------------------------------------------------------------
186 // pulse workers thread
187
188 void pulse_thread_workers_main(void *ptr) {
189 struct netdata_static_thread *static_thread = ptr;
190 pulse_register_workers();
191
192 int update_every =
193 (int)inicfg_get_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", localhost->rrd_update_every);
194 if (update_every < localhost->rrd_update_every) {
195 update_every = localhost->rrd_update_every;
196 inicfg_set_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", update_every);
197 }
198
199 usec_t step = update_every * USEC_PER_SEC;
200 heartbeat_t hb;
201 heartbeat_init(&hb, USEC_PER_SEC);
202 usec_t real_step = USEC_PER_SEC;
203
204 // keep the randomness at zero
205 // to make sure we are not close to any other thread
206 hb.randomness = 0;
207
208 while (service_running(SERVICE_COLLECTORS)) {
209 worker_is_idle();
210 heartbeat_next(&hb);
211 if (real_step < step) {
212 real_step += USEC_PER_SEC;
213 continue;
214 }
215 real_step = USEC_PER_SEC;
216
217 worker_is_busy(WORKER_JOB_WORKERS);
218 pulse_workers_do(pulse_extended_enabled);
219 }
220
221 static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
222 pulse_workers_cleanup();
223 worker_unregister();
224 static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
225 }
226
227 // ---------------------------------------------------------------------------------------------------------------------
228 // pulse workers thread
229
230 void pulse_thread_memory_extended_main(void *ptr) {
231 struct netdata_static_thread *static_thread = ptr;
232 pulse_register_workers();
233
234 int update_every =
235 (int)inicfg_get_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", localhost->rrd_update_every);
236 if (update_every < localhost->rrd_update_every) {
237 update_every = localhost->rrd_update_every;
238 inicfg_set_duration_seconds(&netdata_config, CONFIG_SECTION_PULSE, "update every", update_every);
239 }
240
241 usec_t step = update_every * USEC_PER_SEC;
242 heartbeat_t hb;
243 heartbeat_init(&hb, USEC_PER_SEC);
244 usec_t real_step = USEC_PER_SEC;
245
246 // keep the randomness at zero
247 // to make sure we are not close to any other thread
248 hb.randomness = 0;
249
250 while (service_running(SERVICE_COLLECTORS)) {
251 worker_is_idle();
252 heartbeat_next(&hb);
253 if (real_step < step) {
254 real_step += USEC_PER_SEC;
255 continue;
256 }
257 real_step = USEC_PER_SEC;
258
259 #ifdef NETDATA_TRACE_ALLOCATIONS
260 worker_is_busy(WORKER_JOB_MALLOC_TRACE);
261 pulse_trace_allocations_do(pulse_extended_enabled);
262 #endif
263
264 worker_is_busy(WORKER_JOB_MEMORY_EXTENDED);
265 pulse_daemon_memory_system_do(pulse_extended_enabled);
266 }
267
268 static_thread->enabled = NETDATA_MAIN_THREAD_EXITING;
269 worker_unregister();
270 static_thread->enabled = NETDATA_MAIN_THREAD_EXITED;
271 }