| 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 | } |