@cryptotaxi247 / netdata-1 / commits / 584d7d522

Switch to uv threads (#20250)

* Use nd_thread Switch to uv_thread_create * Remove NETDATA_THREAD_OPTION_JOINABLE option * Remove unused and commented-out code

Stelios Fragkakis committed May 12, 2025 at 19:26 UTC 584d7d52215f386a933efbd7b9f48ec41c33775d
56 files changed +161 -264
src/collectors/cgroups.plugin/cgroup-discovery.c
+3 -2
@@ -1287,7 +1287,7 @@ static inline void discovery_find_all_cgroups() {
1287 netdata_log_debug(D_CGROUP, "done searching for cgroups");
1288 }
1289
1290 -void cgroup_discovery_worker(void *ptr)
1290 +void *cgroup_discovery_worker(void *ptr)
1291 {
1292 UNUSED(ptr);
1293 uv_thread_set_name_np("P[cgroupsdisc]");
@@ -1311,7 +1311,7 @@ void cgroup_discovery_worker(void *ptr)
1311 NULL,
1312 SIMPLE_PATTERN_EXACT, true);
1313
1314 - service_register(SERVICE_THREAD_TYPE_LIBUV, NULL, NULL, NULL, false);
1314 + service_register(NULL, NULL, NULL);
1315
1316 netdata_cgroup_ebpf_initialize_shm();
1317
@@ -1342,4 +1342,5 @@ void cgroup_discovery_worker(void *ptr)
1342 worker_unregister();
1343 service_exits();
1344 __atomic_store_n(&discovery_thread.exited,1,__ATOMIC_RELAXED);
1345 + return NULL;
1346 }
src/collectors/cgroups.plugin/cgroup-internals.h
+2 -2
@@ -261,7 +261,7 @@ struct cgroup {
261 };
262
263 struct discovery_thread {
264 - uv_thread_t thread;
264 + ND_THREAD *thread;
265 uv_mutex_t mutex;
266 uv_cond_t cond_var;
267 int exited;
@@ -274,7 +274,7 @@ extern char cgroup_chart_id_prefix[];
274 extern char services_chart_id_prefix[];
275 extern uv_mutex_t cgroup_root_mutex;
276
277 -void cgroup_discovery_worker(void *ptr);
277 +void *cgroup_discovery_worker(void *ptr);
278
279 extern bool is_inside_k8s;
280 extern long system_page_size;
src/collectors/cgroups.plugin/sys_fs_cgroup.c
+4 -3
@@ -1398,9 +1398,10 @@ void *cgroups_main(void *ptr) {
1398 goto exit;
1399 }
1400
1401 - int error = uv_thread_create(&discovery_thread.thread, cgroup_discovery_worker, NULL);
1402 - if (error) {
1403 - collector_error("CGROUP: cannot create thread worker. uv_thread_create(): %s", uv_strerror(error));
1401 + discovery_thread.thread = nd_thread_create("CGDISCOVER", NETDATA_THREAD_OPTION_DEFAULT, cgroup_discovery_worker, NULL);
1402 +
1403 + if (!discovery_thread.thread) {
1404 + collector_error("CGROUP: cannot create thread worker");
1405 goto exit;
1406 }
1407
src/collectors/debugfs.plugin/module-libsensors.c
+1 -1
@@ -1321,7 +1321,7 @@ int do_module_libsensors(int update_every, const char *name __maybe_unused) {
1321 if(!libsensors) {
1322 libsensors_update_every = update_every;
1323 libsensors_running = true;
1324 - libsensors = nd_thread_create("LIBSENSORS", NETDATA_THREAD_OPTION_JOINABLE, libsensors_thread, NULL);
1324 + libsensors = nd_thread_create("LIBSENSORS", NETDATA_THREAD_OPTION_DEFAULT, libsensors_thread, NULL);
1325 }
1326
1327 return libsensors && libsensors_running ? 0 : 1;
src/collectors/diskspace.plugin/plugin_diskspace.c
+1 -1
@@ -873,7 +873,7 @@ void *diskspace_main(void *ptr) {
873
874 diskspace_slow_thread = nd_thread_create(
875 "P[diskspace slow]",
876 - NETDATA_THREAD_OPTION_JOINABLE,
876 + NETDATA_THREAD_OPTION_DEFAULT,
877 diskspace_slow_worker,
878 &slow_worker_data);
879
src/collectors/ebpf.plugin/ebpf.c
+3 -3
@@ -4244,7 +4244,7 @@ static void ebpf_initialize_data_sharing()
4244 switch (integration_with_collectors) {
4245 case NETDATA_EBPF_INTEGRATION_SOCKET: {
4246 socket_ipc =
4247 - nd_thread_create("ebpf_socket_ipc", NETDATA_THREAD_OPTION_JOINABLE, ebpf_socket_thread_ipc, NULL);
4247 + nd_thread_create("ebpf_socket_ipc", NETDATA_THREAD_OPTION_DEFAULT, ebpf_socket_thread_ipc, NULL);
4248 break;
4249 }
4250 case NETDATA_EBPF_INTEGRATION_SHM:
@@ -4389,7 +4389,7 @@ int main(int argc, char **argv)
4389 cgroup_integration_thread.start_routine = ebpf_cgroup_integration;
4390
4391 cgroup_integration_thread.thread =
4392 - nd_thread_create(cgroup_integration_thread.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_cgroup_integration, NULL);
4392 + nd_thread_create(cgroup_integration_thread.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_cgroup_integration, NULL);
4393
4394 ebpf_initialize_data_sharing();
4395
@@ -4407,7 +4407,7 @@ int main(int argc, char **argv)
4407 if (em->functions.apps_routine && (em->apps_charts || em->cgroup_charts)) {
4408 collect_pids |= 1 << i;
4409 }
4410 - st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_JOINABLE, st->start_routine, em);
4410 + st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_DEFAULT, st->start_routine, em);
4411 } else {
4412 em->lifetime = EBPF_DEFAULT_LIFETIME;
4413 }
src/collectors/ebpf.plugin/ebpf_cachestat.c
+1 -1
@@ -1760,7 +1760,7 @@ void *ebpf_cachestat_thread(void *ptr)
1760 pthread_mutex_unlock(&lock);
1761
1762 ebpf_read_cachestat.thread =
1763 - nd_thread_create(ebpf_read_cachestat.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_cachestat_thread, em);
1763 + nd_thread_create(ebpf_read_cachestat.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_cachestat_thread, em);
1764
1765 cachestat_collector(em);
1766
src/collectors/ebpf.plugin/ebpf_dcstat.c
+1 -1
@@ -1532,7 +1532,7 @@ void *ebpf_dcstat_thread(void *ptr)
1532 pthread_mutex_unlock(&lock);
1533
1534 ebpf_read_dcstat.thread =
1535 - nd_thread_create(ebpf_read_dcstat.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_dcstat_thread, em);
1535 + nd_thread_create(ebpf_read_dcstat.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_dcstat_thread, em);
1536
1537 dcstat_collector(em);
1538
src/collectors/ebpf.plugin/ebpf_fd.c
+1 -1
@@ -1558,7 +1558,7 @@ void *ebpf_fd_thread(void *ptr)
1558
1559 pthread_mutex_unlock(&lock);
1560
1561 - ebpf_read_fd.thread = nd_thread_create(ebpf_read_fd.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_fd_thread, em);
1561 + ebpf_read_fd.thread = nd_thread_create(ebpf_read_fd.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_fd_thread, em);
1562
1563 fd_collector(em);
1564
src/collectors/ebpf.plugin/ebpf_functions.c
+1 -1
@@ -31,7 +31,7 @@ static int ebpf_function_start_thread(ebpf_module_t *em, int period)
31 netdata_log_info("Starting thread %s with lifetime = %d", em->info.thread_name, period);
32 #endif
33
34 - st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_JOINABLE, st->start_routine, em);
34 + st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_DEFAULT, st->start_routine, em);
35 return st->thread ? 0 : 1;
36 }
37
src/collectors/ebpf.plugin/ebpf_shm.c
+1 -1
@@ -1395,7 +1395,7 @@ void *ebpf_shm_thread(void *ptr)
1395 pthread_mutex_unlock(&lock);
1396
1397 ebpf_read_shm.thread =
1398 - nd_thread_create(ebpf_read_shm.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_shm_thread, em);
1398 + nd_thread_create(ebpf_read_shm.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_shm_thread, em);
1399
1400 shm_collector(em);
1401
src/collectors/ebpf.plugin/ebpf_socket.c
+1 -1
@@ -3051,7 +3051,7 @@ void *ebpf_socket_thread(void *ptr)
3051 NETDATA_MAX_SOCKET_VECTOR);
3052
3053 ebpf_read_socket.thread =
3054 - nd_thread_create(ebpf_read_socket.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_socket_thread, em);
3054 + nd_thread_create(ebpf_read_socket.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_socket_thread, em);
3055
3056 pthread_mutex_lock(&lock);
3057 ebpf_socket_create_global_charts(em);
src/collectors/ebpf.plugin/ebpf_swap.c
+1 -1
@@ -1221,7 +1221,7 @@ void *ebpf_swap_thread(void *ptr)
1221 pthread_mutex_unlock(&lock);
1222
1223 ebpf_read_swap.thread =
1224 - nd_thread_create(ebpf_read_swap.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_swap_thread, em);
1224 + nd_thread_create(ebpf_read_swap.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_swap_thread, em);
1225
1226 swap_collector(em);
1227
src/collectors/ebpf.plugin/ebpf_vfs.c
+1 -1
@@ -2958,7 +2958,7 @@ void *ebpf_vfs_thread(void *ptr)
2958 pthread_mutex_unlock(&lock);
2959
2960 ebpf_read_vfs.thread =
2961 - nd_thread_create(ebpf_read_vfs.name, NETDATA_THREAD_OPTION_JOINABLE, ebpf_read_vfs_thread, em);
2961 + nd_thread_create(ebpf_read_vfs.name, NETDATA_THREAD_OPTION_DEFAULT, ebpf_read_vfs_thread, em);
2962
2963 vfs_collector(em);
2964
src/collectors/freeipmi.plugin/freeipmi_plugin.c
+3 -3
@@ -1982,9 +1982,9 @@ int main (int argc, char **argv) {
1982 },
1983 };
1984
1985 - nd_thread_create("IPMI[sensors]", NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE, netdata_ipmi_collection_thread, &sensors_data);
1986 - if(netdata_do_sel)
1987 - nd_thread_create("IPMI[sel]", NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE, netdata_ipmi_collection_thread, &sel_data);
1985 + nd_thread_create("IPMI[sensors]", NETDATA_THREAD_OPTION_DONT_LOG, netdata_ipmi_collection_thread, &sensors_data);
1986 + if (netdata_do_sel)
1987 + nd_thread_create("IPMI[sel]", NETDATA_THREAD_OPTION_DONT_LOG, netdata_ipmi_collection_thread, &sel_data);
1988
1989 // ------------------------------------------------------------------------
1990 // the main loop
src/collectors/proc.plugin/plugin_proc.c
+1 -1
@@ -217,7 +217,7 @@ void *proc_main(void *ptr)
217
218 if (inicfg_get_boolean(&netdata_config, "plugin:proc", "/proc/net/dev", CONFIG_BOOLEAN_YES)) {
219 netdata_log_debug(D_SYSTEM, "Starting thread %s.", THREAD_NETDEV_NAME);
220 - netdev_thread = nd_thread_create(THREAD_NETDEV_NAME, NETDATA_THREAD_OPTION_JOINABLE, netdev_main, NULL);
220 + netdev_thread = nd_thread_create(THREAD_NETDEV_NAME, NETDATA_THREAD_OPTION_DEFAULT, netdev_main, NULL);
221 }
222
223 inicfg_get_boolean(&netdata_config, "plugin:proc", "/proc/pagetypeinfo", CONFIG_BOOLEAN_NO);
src/collectors/profile.plugin/plugin_profile.cc
+1 -1
@@ -217,7 +217,7 @@ extern "C" void *profile_main(void *ptr) {
217 char Tag[NETDATA_THREAD_TAG_MAX + 1];
218
219 snprintfz(Tag, NETDATA_THREAD_TAG_MAX, "PROFILER[%zu]", Idx);
220 - Threads[Idx] = nd_thread_create(Tag, NETDATA_THREAD_OPTION_JOINABLE,
220 + Threads[Idx] = nd_thread_create(Tag, NETDATA_THREAD_OPTION_DEFAULT,
221 subprofile_main, static_cast<void *>(&Profilers[Idx]));
222 }
223
src/collectors/statsd.plugin/statsd.c
+1 -1
@@ -2613,7 +2613,7 @@ void *statsd_main(void *ptr) {
2613 char tag[NETDATA_THREAD_TAG_MAX + 1];
2614 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "STATSD_IN[%d]", i + 1);
2615 spinlock_init(&statsd.collection_threads_status[i].spinlock);
2616 - statsd.collection_threads_status[i].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
2616 + statsd.collection_threads_status[i].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT,
2617 statsd_collector_thread, &statsd.collection_threads_status[i]);
2618 }
2619
src/collectors/systemd-journal.plugin/systemd-main.c
+1 -1
@@ -67,7 +67,7 @@ int main(int argc __maybe_unused, char **argv __maybe_unused)
67 // ------------------------------------------------------------------------
68 // watcher thread
69
70 - nd_thread_create("SDWATCH", NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE, nd_journal_watcher_main, NULL);
70 + nd_thread_create("SDWATCH", NETDATA_THREAD_OPTION_DONT_LOG, nd_journal_watcher_main, NULL);
71
72 // ------------------------------------------------------------------------
73 // the event loop for functions
src/daemon/config/netdata-conf-db.c
+1 -1
@@ -290,7 +290,7 @@ void netdata_conf_dbengine_init(const char *hostname) {
290 if(parallel_initialization) {
291 char tag[NETDATA_THREAD_TAG_MAX + 1];
292 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "DBENGINIT[%zu]", tier);
293 - tiers_init[tier].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE, dbengine_tier_init, &tiers_init[tier]);
293 + tiers_init[tier].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, dbengine_tier_init, &tiers_init[tier]);
294 }
295 else
296 dbengine_tier_init(&tiers_init[tier]);
src/daemon/daemon-service.c
+9 -52
@@ -4,16 +4,11 @@
4
5 typedef struct service_thread {
6 pid_t tid;
7 - SERVICE_THREAD_TYPE type;
7 SERVICE_TYPE services;
8 char name[ND_THREAD_TAG_MAX + 1];
10 - bool stop_immediately;
9 bool cancelled;
10
13 - union {
14 - ND_THREAD *netdata_thread;
15 - uv_thread_t uv_thread;
16 - };
11 + ND_THREAD *netdata_thread;
12
13 force_quit_t force_quit_callback;
14 request_quit_t request_quit_callback;
@@ -27,7 +22,8 @@ struct service_globals {
22 .pid_judy = NULL,
23 };
24
30 -SERVICE_THREAD *service_register(SERVICE_THREAD_TYPE thread_type, request_quit_t request_quit_callback, force_quit_t force_quit_callback, void *data, bool update __maybe_unused) {
25 +SERVICE_THREAD *service_register(request_quit_t request_quit_callback, force_quit_t force_quit_callback, void *data)
26 +{
27 SERVICE_THREAD *sth = NULL;
28 pid_t tid = gettid_cached();
29
@@ -36,23 +32,11 @@ SERVICE_THREAD *service_register(SERVICE_THREAD_TYPE thread_type, request_quit_t
32 if(!*PValue) {
33 sth = callocz(1, sizeof(SERVICE_THREAD));
34 sth->tid = tid;
39 - sth->type = thread_type;
35 sth->request_quit_callback = request_quit_callback;
36 sth->force_quit_callback = force_quit_callback;
37 sth->data = data;
38 *PValue = sth;
44 -
45 - switch(thread_type) {
46 - default:
47 - case SERVICE_THREAD_TYPE_NETDATA:
48 - sth->netdata_thread = nd_thread_self();
49 - break;
50 -
51 - case SERVICE_THREAD_TYPE_EVENT_LOOP:
52 - case SERVICE_THREAD_TYPE_LIBUV:
53 - sth->uv_thread = uv_thread_self();
54 - break;
55 - }
39 + sth->netdata_thread = nd_thread_self();
40
41 const char *name = nd_thread_tag();
42 if(!name) name = "";
@@ -82,15 +66,11 @@ bool service_running(SERVICE_TYPE service) {
66 static __thread SERVICE_THREAD *sth = NULL;
67
68 if(unlikely(!sth))
85 - sth = service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, false);
69 + sth = service_register(NULL, NULL, NULL);
70
71 sth->services |= service;
72
89 - bool cancelled = false;
90 - if (sth->type == SERVICE_THREAD_TYPE_NETDATA)
91 - cancelled = nd_thread_signaled_to_cancel();
92 -
93 - return !sth->stop_immediately && !exit_initiated_get() && !cancelled;
73 + return !nd_thread_signaled_to_cancel() && !exit_initiated_get();
74 }
75
76 void service_signal_exit(SERVICE_TYPE service) {
@@ -103,19 +83,8 @@ void service_signal_exit(SERVICE_TYPE service) {
83 SERVICE_THREAD *sth = *PValue;
84
85 if((sth->services & service)) {
106 - sth->stop_immediately = true;
107 -
108 - switch(sth->type) {
109 - default:
110 - case SERVICE_THREAD_TYPE_NETDATA:
111 - nd_thread_signal_cancel(sth->netdata_thread);
112 - break;
113 -
114 - case SERVICE_THREAD_TYPE_EVENT_LOOP:
115 - case SERVICE_THREAD_TYPE_LIBUV:
116 - break;
117 - }
118 -
86 + nd_thread_signal_cancel(sth->netdata_thread);
87 + nd_log_daemon(NDLP_DEBUG, "SERVICE: Signal to stop : %s", sth->name);
88 if(sth->request_quit_callback) {
89 spinlock_unlock(&service_globals.lock);
90 sth->request_quit_callback(sth->data);
@@ -123,7 +92,6 @@ void service_signal_exit(SERVICE_TYPE service) {
92 }
93 }
94 }
126 -
95 spinlock_unlock(&service_globals.lock);
96 }
97
@@ -180,18 +148,7 @@ bool service_wait_exit(SERVICE_TYPE service, usec_t timeout_ut) {
148 SERVICE_THREAD *sth = *PValue;
149 if(sth->services & service && sth->tid != gettid_cached() && !sth->cancelled) {
150 sth->cancelled = true;
183 -
184 - switch(sth->type) {
185 - default:
186 - case SERVICE_THREAD_TYPE_NETDATA:
187 - nd_thread_signal_cancel(sth->netdata_thread);
188 - break;
189 -
190 - case SERVICE_THREAD_TYPE_EVENT_LOOP:
191 - case SERVICE_THREAD_TYPE_LIBUV:
192 - break;
193 - }
194 -
151 + nd_thread_signal_cancel(sth->netdata_thread);
152 if(running)
153 buffer_strcat(thread_list, ", ");
154
src/daemon/daemon-service.h
+1 -7
@@ -24,18 +24,12 @@ typedef enum {
24 SERVICE_SYSTEMD = (1 << 15),
25 } SERVICE_TYPE;
26
27 -typedef enum {
28 - SERVICE_THREAD_TYPE_NETDATA,
29 - SERVICE_THREAD_TYPE_LIBUV,
30 - SERVICE_THREAD_TYPE_EVENT_LOOP,
31 -} SERVICE_THREAD_TYPE;
32 -
27 typedef void (*force_quit_t)(void *data);
28 typedef void (*request_quit_t)(void *data);
29
30 void service_exits(void);
31 bool service_running(SERVICE_TYPE service);
38 -struct service_thread *service_register(SERVICE_THREAD_TYPE thread_type, request_quit_t request_quit_callback, force_quit_t force_quit_callback, void *data, bool update __maybe_unused);
32 +struct service_thread *service_register(request_quit_t request_quit_callback, force_quit_t force_quit_callback, void *data);
33
34 void service_signal_exit(SERVICE_TYPE service);
35 bool service_wait_exit(SERVICE_TYPE service, usec_t timeout_ut);
src/daemon/daemon-shutdown-watcher.c
+1 -1
@@ -195,7 +195,7 @@ void watcher_thread_start() {
195 completion_init(&shutdown_begin_completion);
196 completion_init(&shutdown_end_completion);
197
198 - watcher_thread = nd_thread_create("EXIT_WATCHER", NETDATA_THREAD_OPTION_JOINABLE, watcher_main, NULL);
198 + watcher_thread = nd_thread_create("EXIT_WATCHER", NETDATA_THREAD_OPTION_DEFAULT, watcher_main, NULL);
199 }
200
201 void watcher_thread_stop() {
src/daemon/daemon-shutdown.c
+1 -1
@@ -264,7 +264,7 @@ static void netdata_cleanup_and_exit(EXIT_REASON reason, bool abnormal, bool exi
264
265 ND_THREAD *th[nd_profile.storage_tiers];
266 for (size_t tier = 0; tier < nd_profile.storage_tiers; tier++)
267 - th[tier] = nd_thread_create("rrdeng-exit", NETDATA_THREAD_OPTION_JOINABLE, rrdeng_exit_background, multidb_ctx[tier]);
267 + th[tier] = nd_thread_create("rrdeng-exit", NETDATA_THREAD_OPTION_DEFAULT, rrdeng_exit_background, multidb_ctx[tier]);
268
269 // flush anything remaining again - just in case
270 rrdeng_flush_everything_and_wait(true, true, false);
src/daemon/daemon-systemd-watcher.c
+1 -1
@@ -138,7 +138,7 @@ finish:
138 void *systemd_watcher_thread(void *arg) {
139 struct netdata_static_thread *static_thread = arg;
140
141 - service_register(SERVICE_THREAD_TYPE_NETDATA, NULL, NULL, NULL, false);
141 + service_register(NULL, NULL, NULL);
142
143 listen_for_systemd_dbus_events();
144
src/daemon/dyncfg/dyncfg-unittest.c
+1 -1
@@ -580,7 +580,7 @@ int dyncfg_unittest(void) {
580 // ------------------------------------------------------------------------
581 // create the thread for testing async communication
582
583 - ND_THREAD *thread = nd_thread_create("unittest", NETDATA_THREAD_OPTION_JOINABLE, dyncfg_unittest_thread_action, NULL);
583 + ND_THREAD *thread = nd_thread_create("unittest", NETDATA_THREAD_OPTION_DEFAULT, dyncfg_unittest_thread_action, NULL);
584
585 // ------------------------------------------------------------------------
586 // single
src/daemon/libuv_workers.c
-18
@@ -109,24 +109,6 @@ void register_libuv_worker_jobs() {
109 register_libuv_worker_jobs_internal();
110 }
111
112 -// utils
113 -#define MAX_THREAD_CREATE_RETRIES (10)
114 -#define MAX_THREAD_CREATE_WAIT_MS (1000)
115 -
116 -int create_uv_thread(uv_thread_t *thread, uv_thread_cb thread_func, void *arg, int *retries)
117 -{
118 - int err;
119 -
120 - do {
121 - err = uv_thread_create(thread, thread_func, arg);
122 - if (err == 0)
123 - break;
124 - sleep_usec(MAX_THREAD_CREATE_WAIT_MS * USEC_PER_MS);
125 - } while (err == UV_EAGAIN && ++(*retries) < MAX_THREAD_CREATE_RETRIES);
126 -
127 - return err;
128 -}
129 -
112 void libuv_close_callback(uv_handle_t *handle, void *data __maybe_unused)
113 {
114 // Only close handles that aren't already closing
src/daemon/libuv_workers.h
-1
@@ -87,7 +87,6 @@ enum event_loop_job {
87 };
88
89 void register_libuv_worker_jobs();
90 -int create_uv_thread(uv_thread_t *thread, uv_thread_cb thread_func, void *arg, int *retries);
90 void libuv_close_callback(uv_handle_t *handle, void *data __maybe_unused);
91
92 #endif //NETDATA_EVENT_LOOP_H
src/daemon/main.c
+2 -2
@@ -1069,7 +1069,7 @@ int netdata_main(int argc, char **argv) {
1069
1070 if(st->enabled) {
1071 netdata_log_debug(D_SYSTEM, "Starting thread %s.", st->name);
1072 - st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_JOINABLE, st->start_routine, st);
1072 + st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_DEFAULT, st->start_routine, st);
1073 }
1074 else
1075 netdata_log_debug(D_SYSTEM, "Not starting thread %s.", st->name);
@@ -1120,7 +1120,7 @@ int netdata_main(int argc, char **argv) {
1120 struct netdata_static_thread *st = &static_threads[i];
1121 st->enabled = 1;
1122 netdata_log_debug(D_SYSTEM, "Starting thread %s.", st->name);
1123 - st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_JOINABLE, st->start_routine, st);
1123 + st->thread = nd_thread_create(st->name, NETDATA_THREAD_OPTION_DEFAULT, st->start_routine, st);
1124 }
1125 }
1126 }
src/daemon/winsvc.cc
+1 -1
@@ -141,7 +141,7 @@ static void WINAPI ServiceControlHandler(DWORD controlCode)
141 netdata_service_log("Creating cleanup thread...");
142 char tag[NETDATA_THREAD_TAG_MAX + 1];
143 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "CLEANUP");
144 - cleanup_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE, call_netdata_cleanup, &controlCode);
144 + cleanup_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, call_netdata_cleanup, &controlCode);
145
146 // Signal the stop request
147 netdata_service_log("Signalling the cleanup thread...");
src/database/engine/cache.c
+4 -4
@@ -2085,7 +2085,7 @@ PGC *pgc_create(const char *name,
2085 // last create the eviction thread
2086 {
2087 completion_init(&cache->evictor.completion);
2088 - cache->evictor.thread = nd_thread_create(name, NETDATA_THREAD_OPTION_JOINABLE, pgc_evict_thread, cache);
2088 + cache->evictor.thread = nd_thread_create(name, NETDATA_THREAD_OPTION_DEFAULT, pgc_evict_thread, cache);
2089 }
2090
2091 return cache;
@@ -2872,7 +2872,7 @@ void unittest_stress_test(void) {
2872
2873 pthread_t service_thread;
2874 nd_thread_create(&service_thread, "SERVICE",
2875 - NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG,
2875 + NETDATA_THREAD_OPTION_DONT_LOG,
2876 unittest_stress_test_service, NULL);
2877
2878 pthread_t collect_threads[pgc_uts.collect_threads];
@@ -2882,7 +2882,7 @@ void unittest_stress_test(void) {
2882 char buffer[100 + 1];
2883 snprintfz(buffer, sizeof(buffer) - 1, "COLLECT_%zu", i);
2884 nd_thread_create(&collect_threads[i], buffer,
2885 - NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG,
2885 + NETDATA_THREAD_OPTION_DONT_LOG,
2886 unittest_stress_test_collector, &collect_thread_ids[i]);
2887 }
2888
@@ -2895,7 +2895,7 @@ void unittest_stress_test(void) {
2895 snprintfz(buffer, sizeof(buffer) - 1, "QUERY_%zu", i);
2896 initstate_r(1, pgc_uts.rand_statebufs, 1024, &pgc_uts.random_data[i]);
2897 nd_thread_create(&queries_threads[i], buffer,
2898 - NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG,
2898 + NETDATA_THREAD_OPTION_DONT_LOG,
2899 unittest_stress_test_queries, &query_thread_ids[i]);
2900 }
2901
src/database/engine/mrg-unittest.c
+1 -1
@@ -173,7 +173,7 @@ int mrg_unittest(void) {
173 for(size_t i = 0; i < threads ; i++) {
174 char buf[15 + 1];
175 snprintfz(buf, sizeof(buf) - 1, "TH[%zu]", i);
176 - th[i] = nd_thread_create(buf, NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG, mrg_stress, &t);
176 + th[i] = nd_thread_create(buf, NETDATA_THREAD_OPTION_DONT_LOG, mrg_stress, &t);
177 }
178
179 sleep_usec(run_for_secs * USEC_PER_SEC);
src/database/engine/rrdengine.c
+19 -16
@@ -31,7 +31,7 @@ static inline void worker_dispatch_extent_read(struct rrdeng_cmd cmd, bool from_
31 static inline void worker_dispatch_query_prep(struct rrdeng_cmd cmd, bool from_worker);
32
33 struct rrdeng_main {
34 - uv_thread_t thread;
34 + ND_THREAD *thread;
35 uv_loop_t loop;
36 uv_async_t async;
37 uv_timer_t timer;
@@ -1339,7 +1339,7 @@ static void *flush_dirty_pages_of_section_tp_worker(struct rrdengine_instance *c
1339
1340 struct mrg_load_thread {
1341 int max_threads;
1342 - uv_thread_t thread;
1342 + ND_THREAD *thread;
1343 uv_sem_t *sem;
1344 int tier;
1345 struct rrdengine_datafile *datafile;
@@ -1350,7 +1350,7 @@ struct mrg_load_thread {
1350 size_t max_running_threads = 0;
1351 size_t running_threads = 0;
1352
1353 -void journalfile_v2_populate_retention_to_mrg_worker(void *arg)
1353 +void *journalfile_v2_populate_retention_to_mrg_worker(void *arg)
1354 {
1355 struct mrg_load_thread *mlt = arg;
1356 uv_sem_wait(mlt->sem);
@@ -1374,6 +1374,7 @@ void journalfile_v2_populate_retention_to_mrg_worker(void *arg)
1374
1375 // Signal completion - this needs to be last
1376 __atomic_store_n(&mlt->finished, true, __ATOMIC_RELEASE);
1377 + return NULL;
1378 }
1379
1380 static void after_populate_mrg(struct rrdengine_instance *ctx __maybe_unused, void *data __maybe_unused, struct completion *completion __maybe_unused, uv_work_t* req __maybe_unused, int status __maybe_unused) {
@@ -1444,7 +1445,7 @@ static void *populate_mrg_tp_worker(
1445 if (__atomic_load_n(&mlt[index].finished, __ATOMIC_RELAXED) &&
1446 __atomic_load_n(&mlt[index].tier, __ATOMIC_ACQUIRE) == tier) {
1447
1447 - rc = uv_thread_join(&(mlt[index].thread));
1448 + rc = nd_thread_join(mlt[index].thread);
1449 if (rc)
1450 nd_log_daemon(NDLP_WARNING, "Failed to join thread, rc = %d", rc);
1451
@@ -1479,11 +1480,10 @@ static void *populate_mrg_tp_worker(
1480 __atomic_store_n(&mlt[thread_index].tier, tier, __ATOMIC_RELAXED);
1481 mlt[thread_index].datafile = datafile;
1482
1482 - rc = uv_thread_create(&mlt[thread_index].thread,
1483 - journalfile_v2_populate_retention_to_mrg_worker,
1484 - &mlt[thread_index]);
1483 + mlt[thread_index].thread = nd_thread_create("MRGLOAD", NETDATA_THREAD_OPTION_DEFAULT, journalfile_v2_populate_retention_to_mrg_worker,
1484 + &mlt[thread_index]);
1485
1486 - if (rc) {
1486 + if (!mlt[thread_index].thread) {
1487 nd_log_daemon(NDLP_WARNING, "Failed to create thread, rc = %d", rc);
1488 __atomic_store_n(&mlt[thread_index].busy, false, __ATOMIC_RELEASE);
1489 spinlock_unlock(&datafile->populate_mrg.spinlock);
@@ -1504,7 +1504,7 @@ static void *populate_mrg_tp_worker(
1504
1505 if (__atomic_load_n(&mlt[index].finished, __ATOMIC_RELAXED)) {
1506 // Thread is finished, join it
1507 - rc = uv_thread_join(&(mlt[index].thread));
1507 + rc = nd_thread_join((mlt[index].thread));
1508 if (rc)
1509 nd_log_daemon(NDLP_WARNING, "Failed to join thread, rc = %d", rc);
1510
@@ -1958,11 +1958,13 @@ bool rrdeng_dbengine_spawn(struct rrdengine_instance *ctx __maybe_unused) {
1958 dbengine_initialize_structures();
1959
1960 int retries = 0;
1961 - int create_uv_thread_rc = create_uv_thread(&rrdeng_main.thread, dbengine_event_loop, &rrdeng_main, &retries);
1962 - if (create_uv_thread_rc)
1963 - nd_log_daemon(NDLP_ERR, "Failed to create DBENGINE thread, error %s, after %d retries", uv_err_name(create_uv_thread_rc), retries);
1961 +// int create_uv_thread_rc = create_uv_thread(&rrdeng_main.thread, dbengine_event_loop, &rrdeng_main, &retries);
1962 + rrdeng_main.thread = nd_thread_create("DBEV", NETDATA_THREAD_OPTION_DEFAULT, dbengine_event_loop, &rrdeng_main);
1963 +
1964 +// if (!rrdeng_main.thread)
1965 +// nd_log_daemon(NDLP_ERR, "Failed to create DBENGINE thread, error %s, after %d retries", uv_err_name(create_uv_thread_rc), retries);
1966
1965 - fatal_assert(0 == create_uv_thread_rc);
1967 + fatal_assert(0 != rrdeng_main.thread);
1968
1969 if (retries)
1970 nd_log_daemon(NDLP_WARNING, "DBENGINE thread was created after %d attempts", retries);
@@ -2036,10 +2038,10 @@ void rrdeng_calculate_tier_disk_space_percentage(void)
2038 (!__atomic_load_n(&(ctx)->atomic.migration_to_v2_running, __ATOMIC_RELAXED) && \
2039 !__atomic_load_n(&(ctx)->atomic.now_deleting_files, __ATOMIC_RELAXED))
2040
2039 -void dbengine_event_loop(void* arg) {
2041 +void *dbengine_event_loop(void* arg) {
2042 sanity_check();
2043 uv_thread_set_name_np("DBENGINE");
2042 - service_register(SERVICE_THREAD_TYPE_EVENT_LOOP, NULL, NULL, NULL, true);
2044 + service_register(NULL, NULL, NULL);
2045
2046 worker_register("DBENGINE");
2047
@@ -2265,13 +2267,14 @@ void dbengine_event_loop(void* arg) {
2267 nd_log(NDLS_DAEMON, NDLP_DEBUG, "Shutting down dbengine thread");
2268 (void) uv_loop_close(&main->loop);
2269 worker_unregister();
2270 + return NULL;
2271 }
2272
2273 void dbengine_shutdown()
2274 {
2275 rrdeng_enq_cmd(NULL, RRDENG_OPCODE_SHUTDOWN_EVLOOP, NULL, NULL, STORAGE_PRIORITY_INTERNAL_DBENGINE, NULL, NULL);
2276
2274 - int rc = uv_thread_join(&rrdeng_main.thread);
2277 + int rc = nd_thread_join(rrdeng_main.thread);
2278 if (rc)
2279 nd_log_daemon(NDLP_ERR, "DBENGINE: Failed to join thread, error %s", uv_err_name(rc));
2280 else
src/database/engine/rrdengine.h
+1 -1
@@ -465,7 +465,7 @@ bool rrdeng_ctx_tier_cap_exceeded(struct rrdengine_instance *ctx);
465 int init_rrd_files(struct rrdengine_instance *ctx);
466 void finalize_rrd_files(struct rrdengine_instance *ctx);
467 bool rrdeng_dbengine_spawn(struct rrdengine_instance *ctx);
468 -void dbengine_event_loop(void *arg);
468 +void *dbengine_event_loop(void *arg);
469
470 typedef void (*enqueue_callback_t)(struct rrdeng_cmd *cmd);
471 typedef void (*dequeue_callback_t)(struct rrdeng_cmd *cmd);
src/database/sqlite/sqlite_aclk.c
+8 -14
@@ -37,7 +37,7 @@ static void create_node_instance_result_job(const char *machine_guid, const char
37 }
38
39 struct aclk_sync_config_s {
40 - uv_thread_t thread;
40 + ND_THREAD *thread;
41 uv_loop_t loop;
42 uv_timer_t timer_req;
43 uv_async_t async;
@@ -617,14 +617,14 @@ static void free_query_list(Pvoid_t JudyL)
617 (!shutdown_requested || config->aclk_queries_running || config->alert_push_running || \
618 config->aclk_batch_job_is_running)
619
620 -static void aclk_synchronization_event_loop(void *arg)
620 +static void *aclk_synchronization_event_loop(void *arg)
621 {
622 struct aclk_sync_config_s *config = arg;
623 uv_thread_set_name_np("ACLKSYNC");
624 config->ar = aral_by_size_acquire(sizeof(struct aclk_database_cmd));
625 worker_register("ACLKSYNC");
626
627 - service_register(SERVICE_THREAD_TYPE_EVENT_LOOP, NULL, NULL, NULL, true);
627 + service_register(NULL, NULL, NULL);
628
629 worker_register_job_name(ACLK_DATABASE_NOOP, "noop");
630 worker_register_job_name(ACLK_DATABASE_NODE_STATE, "node state");
@@ -955,6 +955,7 @@ static void aclk_synchronization_event_loop(void *arg)
955 worker_unregister();
956 service_exits();
957 netdata_log_info("ACLK SYNC: Shutting down ACLK synchronization event loop");
958 + return NULL;
959 }
960
961 static void aclk_initialize_event_loop(void)
@@ -962,17 +963,10 @@ static void aclk_initialize_event_loop(void)
963 memset(&aclk_sync_config, 0, sizeof(aclk_sync_config));
964 completion_init(&aclk_sync_config.start_stop_complete);
965
965 - int retries = 0;
966 - int create_uv_thread_rc = create_uv_thread(&aclk_sync_config.thread, aclk_synchronization_event_loop, &aclk_sync_config, &retries);
967 - if (create_uv_thread_rc)
968 - nd_log_daemon(NDLP_ERR, "Failed to create ACLK synchronization thread, error %s, after %d retries", uv_err_name(create_uv_thread_rc), retries);
966 + aclk_sync_config.thread = nd_thread_create("ACLKSYNC", NETDATA_THREAD_OPTION_DEFAULT, aclk_synchronization_event_loop, &aclk_sync_config);
967 + fatal_assert(NULL != aclk_sync_config.thread);
968
970 - fatal_assert(0 == create_uv_thread_rc);
971 -
972 - if (retries)
973 - nd_log_daemon(NDLP_WARNING, "ACLK synchronization thread was created after %d attempts", retries);
969 completion_wait_for(&aclk_sync_config.start_stop_complete);
975 -
970 // Keep completion, just reset it for next use during shutdown
971 completion_reset(&aclk_sync_config.start_stop_complete);
972 }
@@ -1074,9 +1068,9 @@ void aclk_synchronization_shutdown(void)
1068 completion_wait_for(&aclk_sync_config.start_stop_complete);
1069
1070 completion_destroy(&aclk_sync_config.start_stop_complete);
1077 - int rc = uv_thread_join(&aclk_sync_config.thread);
1071 + int rc = nd_thread_join(aclk_sync_config.thread);
1072 if (rc)
1079 - nd_log_daemon(NDLP_ERR, "ACLK: Failed to join synchronization thread, error %s", uv_err_name(rc));
1073 + nd_log_daemon(NDLP_ERR, "ACLK: Failed to join synchronization thread");
1074 else
1075 nd_log_daemon(NDLP_INFO, "ACLK: synchronization thread shutdown completed");
1076 }
src/database/sqlite/sqlite_metadata.c
+18 -26
@@ -223,7 +223,7 @@ struct metadata_cmd {
223 };
224
225 struct meta_config_s {
226 - uv_thread_t thread;
226 + ND_THREAD *thread;
227 uv_loop_t loop;
228 uv_async_t async;
229 uv_timer_t timer_req;
@@ -1772,7 +1772,7 @@ struct work_payload {
1772 };
1773
1774 struct host_context_load_thread {
1775 - uv_thread_t thread;
1775 + ND_THREAD *thread;
1776 RRDHOST *host;
1777 sqlite3 *db_meta_thread;
1778 sqlite3 *db_context_thread;
@@ -1784,13 +1784,13 @@ __thread sqlite3 *db_meta_thread = NULL;
1784 __thread sqlite3 *db_context_thread = NULL;
1785 __thread bool main_context_thread = false;
1786
1787 -static void restore_host_context(void *arg)
1787 +static void *restore_host_context(void *arg)
1788 {
1789 struct host_context_load_thread *hclt = arg;
1790 RRDHOST *host = hclt->host;
1791
1792 if (!host)
1793 - return;
1793 + return NULL;
1794
1795 if (!db_meta_thread) {
1796 if (hclt->db_meta_thread) {
@@ -1837,6 +1837,7 @@ static void restore_host_context(void *arg)
1837 }
1838
1839 __atomic_store_n(&hclt->finished, true, __ATOMIC_RELEASE);
1840 + return NULL;
1841 }
1842
1843 // Callback after scan of hosts is done
@@ -1866,7 +1867,7 @@ static bool cleanup_finished_threads(struct host_context_load_thread *hclt, size
1867 if (__atomic_load_n(&(hclt[index].finished), __ATOMIC_RELAXED) ||
1868 (wait && __atomic_load_n(&(hclt[index].busy), __ATOMIC_ACQUIRE))) {
1869
1869 - int rc = uv_thread_join(&(hclt[index].thread));
1870 + int rc = nd_thread_join(hclt[index].thread);
1871 if (rc)
1872 nd_log_daemon(NDLP_WARNING, "Failed to join thread, rc = %d", rc);
1873 __atomic_store_n(&(hclt[index].busy), false, __ATOMIC_RELEASE);
@@ -1933,8 +1934,8 @@ static void ctx_hosts_load(uv_work_t *req)
1934 if (thread_found) {
1935 __atomic_store_n(&hclt[thread_index].busy, true, __ATOMIC_RELAXED);
1936 hclt[thread_index].host = host;
1936 - rc = uv_thread_create(&hclt[thread_index].thread, restore_host_context, &hclt[thread_index]);
1937 - async_exec += (rc == 0);
1937 + hclt[thread_index].thread = nd_thread_create("CTXLOAD", NETDATA_THREAD_OPTION_DEFAULT, restore_host_context, &hclt[thread_index]);
1938 + async_exec += (hclt[thread_index].thread != NULL);
1939 // if it failed, mark the thread slot as free
1940 if (rc)
1941 __atomic_store_n(&hclt[thread_index].busy, false, __ATOMIC_RELAXED);
@@ -2153,8 +2154,6 @@ static void metadata_scan_host(RRDHOST *host, BUFFER *work_buffer, bool is_worke
2154 SQLITE_FINALIZE(ml_load_stmt);
2155 SQLITE_FINALIZE(store_dimension);
2156 SQLITE_FINALIZE(store_chart);
2156 -
2157 - return;
2157 }
2158
2159
@@ -2481,10 +2480,11 @@ static void start_metadata_hosts(uv_work_t *req)
2480 #define MAX_SHUTDOWN_TIMEOUT_SECONDS (10)
2481 #define SHUTDOWN_SLEEP_INTERVAL_MS (100)
2482
2484 -static void metadata_event_loop(void *arg)
2483 +static void *metadata_event_loop(void *arg)
2484 {
2485 struct meta_config_s *config = arg;
2486 uv_thread_set_name_np(EVENT_LOOP_NAME);
2487 + service_register(NULL, NULL, NULL);
2488 worker_register(EVENT_LOOP_NAME);
2489
2490 config->ar = aral_by_size_acquire(sizeof(struct metadata_cmd));
@@ -2714,6 +2714,8 @@ static void metadata_event_loop(void *arg)
2714 worker_unregister();
2715
2716 completion_mark_complete(&config->start_stop_complete);
2717 +
2718 + return NULL;
2719 }
2720
2721 void metadata_sync_shutdown(void)
@@ -2728,9 +2730,10 @@ void metadata_sync_shutdown(void)
2730 nd_log_daemon(NDLP_DEBUG, "METADATA: Waiting for shutdown ACK");
2731 completion_wait_for(&meta_config.start_stop_complete);
2732 completion_destroy(&meta_config.start_stop_complete);
2731 - int rc = uv_thread_join(&meta_config.thread);
2733 +
2734 + int rc = nd_thread_join(meta_config.thread);
2735 if (rc)
2733 - nd_log_daemon(NDLP_ERR, "METADATA: Failed to join synchronization thread, error %s", uv_err_name(rc));
2736 + nd_log_daemon(NDLP_ERR, "METADATA: Failed to join synchronization thread");
2737 else
2738 nd_log_daemon(NDLP_INFO, "METADATA: synchronization thread shutdown completed");
2739 }
@@ -2743,15 +2746,8 @@ void metadata_sync_init(void)
2746 memset(&meta_config, 0, sizeof(meta_config));
2747 completion_init(&meta_config.start_stop_complete);
2748
2746 - int retries = 0;
2747 - int create_uv_thread_rc = create_uv_thread(&meta_config.thread, metadata_event_loop, &meta_config, &retries);
2748 - if (create_uv_thread_rc)
2749 - nd_log_daemon(NDLP_ERR, "Failed to create SQLite metadata sync thread, error %s, after %d retries", uv_err_name(create_uv_thread_rc), retries);
2750 -
2751 - fatal_assert(0 == create_uv_thread_rc);
2752 -
2753 - if (retries)
2754 - nd_log_daemon(NDLP_WARNING, "SQLite metadata sync thread was created after %d attempts", retries);
2749 + meta_config.thread = nd_thread_create("METASYNC", NETDATA_THREAD_OPTION_DEFAULT, metadata_event_loop, &meta_config);
2750 + fatal_assert(NULL != meta_config.thread);
2751
2752 // Wait for initialization
2753 completion_wait_for(&meta_config.start_stop_complete);
@@ -2991,11 +2987,7 @@ static void *metadata_unittest_threads(void)
2987 for (int i = 0; i < threads_to_create; i++) {
2988 char buf[100 + 1];
2989 snprintf(buf, sizeof(buf) - 1, "META[%d]", i);
2994 - threads[i] = nd_thread_create(
2995 - buf,
2996 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
2997 - unittest_queue_metadata,
2998 - &tu);
2990 + threads[i] = nd_thread_create(buf, NETDATA_THREAD_OPTION_DONT_LOG, unittest_queue_metadata, &tu);
2991 }
2992 (void) uv_async_send(&meta_config.async);
2993 sleep_usec(seconds_to_run * USEC_PER_SEC);
src/exporting/init_connectors.c
+1 -1
@@ -105,7 +105,7 @@ int init_connectors(struct engine *engine)
105 char threadname[ND_THREAD_TAG_MAX + 1];
106 snprintfz(threadname, ND_THREAD_TAG_MAX, "%s[%zu]", instance->config.thread_tag, instance->index);
107
108 - instance->thread = nd_thread_create(threadname, NETDATA_THREAD_OPTION_JOINABLE, instance->worker, instance);
108 + instance->thread = nd_thread_create(threadname, NETDATA_THREAD_OPTION_DEFAULT, instance->worker, instance);
109 if (!instance->thread) {
110 netdata_log_error("EXPORTING: cannot create thread worker for instance %s", instance->config.name);
111 instance->exited = 1;
src/libnetdata/aral/aral.c
+1 -5
@@ -1532,11 +1532,7 @@ int aral_stress_test(size_t threads, size_t elements, size_t seconds) {
1532 for(size_t i = 0; i < threads ; i++) {
1533 char tag[ND_THREAD_TAG_MAX + 1];
1534 snprintfz(tag, ND_THREAD_TAG_MAX, "TH[%zu]", i);
1535 - thread_ptrs[i] = nd_thread_create(
1536 - tag,
1537 - NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG,
1538 - aral_test_thread,
1539 - &auc);
1535 + thread_ptrs[i] = nd_thread_create(tag, NETDATA_THREAD_OPTION_DONT_LOG, aral_test_thread, &auc);
1536 }
1537
1538 size_t malloc_done = 0;
src/libnetdata/dictionary/dictionary-unittest.c
+4 -16
@@ -690,11 +690,7 @@ static int dictionary_unittest_threads() {
690
691 char buf[100 + 1];
692 snprintf(buf, 100, "dict%d", i);
693 - tu[i].thread = nd_thread_create(
694 - buf,
695 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
696 - unittest_dict_thread,
697 - &tu[i]);
693 + tu[i].thread = nd_thread_create(buf, NETDATA_THREAD_OPTION_DONT_LOG, unittest_dict_thread, &tu[i]);
694 }
695
696 sleep_usec(seconds_to_run * USEC_PER_SEC);
@@ -871,17 +867,9 @@ static int dictionary_unittest_view_threads() {
867 ND_THREAD *master_thread, *view_thread;
868 tv.join = 0;
869
874 - master_thread = nd_thread_create(
875 - "master",
876 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
877 - unittest_dict_master_thread,
878 - &tv);
879 -
880 - view_thread = nd_thread_create(
881 - "view",
882 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
883 - unittest_dict_view_thread,
884 - &tv);
870 + master_thread = nd_thread_create("master", NETDATA_THREAD_OPTION_DONT_LOG, unittest_dict_master_thread, &tv);
871 +
872 + view_thread = nd_thread_create("view", NETDATA_THREAD_OPTION_DONT_LOG, unittest_dict_view_thread, &tv);
873
874 sleep_usec(seconds_to_run * USEC_PER_SEC);
875
src/libnetdata/functions_evloop/functions_evloop.c
+2 -2
@@ -362,12 +362,12 @@ struct functions_evloop_globals *functions_evloop_init(size_t worker_threads, co
362
363 char tag_buffer[NETDATA_THREAD_TAG_MAX + 1];
364 snprintfz(tag_buffer, NETDATA_THREAD_TAG_MAX, "%s_READER", wg->tag);
365 - wg->reader_thread = nd_thread_create(tag_buffer, NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
365 + wg->reader_thread = nd_thread_create(tag_buffer, NETDATA_THREAD_OPTION_DONT_LOG,
366 rrd_functions_worker_globals_reader_main, wg);
367
368 for(size_t i = 0; i < wg->workers ; i++) {
369 snprintfz(tag_buffer, NETDATA_THREAD_TAG_MAX, "%s_WORK[%zu]", wg->tag, i+1);
370 - wg->worker_threads[i] = nd_thread_create(tag_buffer, NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
370 + wg->worker_threads[i] = nd_thread_create(tag_buffer, NETDATA_THREAD_OPTION_DONT_LOG,
371 rrd_functions_worker_globals_worker_main, wg);
372 }
373
src/libnetdata/local-sockets/local-sockets.h
+1 -1
@@ -1671,7 +1671,7 @@ static inline void local_sockets_namespaces(LS_STATE *ls) {
1671 workers_data[last_thread].inode = inode;
1672 workers[last_thread] = nd_thread_create(
1673 "local-sockets-worker",
1674 - NETDATA_THREAD_OPTION_JOINABLE,
1674 + NETDATA_THREAD_OPTION_DEFAULT,
1675 local_sockets_get_namespace_sockets_worker,
1676 &workers_data[last_thread]);
1677
src/libnetdata/locks/benchmark-rw.c
+4 -10
@@ -346,11 +346,8 @@ int rwlocks_stress_test(void) {
346 };
347
348 snprintf(thr_name, sizeof(thr_name), "pthread_rw%d", i);
349 - pthread_contexts[i].thread = nd_thread_create(
350 - thr_name,
351 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
352 - benchmark_thread,
353 - &pthread_contexts[i]);
349 + pthread_contexts[i].thread =
350 + nd_thread_create(thr_name, NETDATA_THREAD_OPTION_DONT_LOG, benchmark_thread, &pthread_contexts[i]);
351
352 // Initialize spinlock contexts
353 spinlock_contexts[i] = (thread_context_t){
@@ -362,11 +359,8 @@ int rwlocks_stress_test(void) {
359 };
360
361 snprintf(thr_name, sizeof(thr_name), "spin_rw%d", i);
365 - spinlock_contexts[i].thread = nd_thread_create(
366 - thr_name,
367 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
368 - benchmark_thread,
369 - &spinlock_contexts[i]);
362 + spinlock_contexts[i].thread =
363 + nd_thread_create(thr_name, NETDATA_THREAD_OPTION_DONT_LOG, benchmark_thread, &spinlock_contexts[i]);
364 }
365
366 // Run all configurations
src/libnetdata/locks/benchmark.c
+1 -1
@@ -374,7 +374,7 @@ int locks_stress_test(void) {
374 snprintf(thr_name, sizeof(thr_name), "%s%d", lock_names[type], i);
375 threads[type][i].thread = nd_thread_create(
376 thr_name,
377 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
377 + NETDATA_THREAD_OPTION_DONT_LOG,
378 benchmark_thread,
379 &threads[type][i]);
380 }
src/libnetdata/locks/waitq.c
+1 -4
@@ -247,10 +247,7 @@ static int unittest_stress(void) {
247 char thread_name[32];
248 snprintf(thread_name, sizeof(thread_name), "STRESS%d-%d", prio, t);
249 threads[thread_idx] = nd_thread_create(
250 - thread_name,
251 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
252 - stress_thread,
253 - &thread_args[thread_idx]);
250 + thread_name, NETDATA_THREAD_OPTION_DONT_LOG, stress_thread, &thread_args[thread_idx]);
251 thread_idx++;
252 }
253 }
src/libnetdata/spawn_server/log-forwarder.c
+1 -1
@@ -71,7 +71,7 @@ LOG_FORWARDER *log_forwarder_start(void) {
71 nd_log(NDLS_COLLECTORS, NDLP_ERR, "Log forwarder: Failed to set non-blocking mode");
72
73 lf->running = true;
74 - lf->thread = nd_thread_create("log-fw", NETDATA_THREAD_OPTION_JOINABLE, log_forwarder_thread_func, lf);
74 + lf->thread = nd_thread_create("log-fw", NETDATA_THREAD_OPTION_DEFAULT, log_forwarder_thread_func, lf);
75
76 nd_log(NDLS_COLLECTORS, NDLP_INFO, "Log forwarder: created thread pointer: %p", lf->thread);
77
src/libnetdata/string/string.c
+1 -1
@@ -876,7 +876,7 @@ int string_unittest(size_t entries) {
876 for (int i = 0; i < threads_to_create; i++) {
877 char buf[100 + 1];
878 snprintf(buf, 100, "string%d", i);
879 - threads[i] = nd_thread_create(buf, NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE, string_thread, &tu);
879 + threads[i] = nd_thread_create(buf, NETDATA_THREAD_OPTION_DONT_LOG, string_thread, &tu);
880 }
881 sleep_usec(seconds_to_run * USEC_PER_SEC);
882
src/libnetdata/threads/threads.c
+31 -29
@@ -22,7 +22,7 @@ struct nd_thread {
22 void *ret; // the return value of start routine
23 void *(*start_routine) (void *);
24 NETDATA_THREAD_OPTIONS options;
25 - pthread_t thread;
25 + uv_thread_t thread;
26 bool cancel_atomic;
27
28 #ifdef NETDATA_INTERNAL_CHECKS
@@ -339,15 +339,13 @@ static void nd_thread_exit(ND_THREAD *nti) {
339 }
340 spinlock_unlock(&threads_globals.running.spinlock);
341
342 - //if (nd_thread_status_check(nti, NETDATA_THREAD_OPTION_JOINABLE) != NETDATA_THREAD_OPTION_JOINABLE) {
342 spinlock_lock(&threads_globals.exited.spinlock);
343 DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(threads_globals.exited.list, nti, prev, next);
344 nti->list = ND_THREAD_LIST_EXITED;
345 spinlock_unlock(&threads_globals.exited.spinlock);
347 - //}
346 }
347
350 -static void *nd_thread_starting_point(void *ptr) {
348 +static void nd_thread_starting_point(void *ptr) {
349 ND_THREAD *nti = _nd_thread_info = (ND_THREAD *)ptr;
350 nd_thread_status_set(nti, NETDATA_THREAD_STATUS_STARTED);
351
@@ -357,12 +355,6 @@ static void *nd_thread_starting_point(void *ptr) {
355 if(nd_thread_status_check(nti, NETDATA_THREAD_OPTION_DONT_LOG_STARTUP) != NETDATA_THREAD_OPTION_DONT_LOG_STARTUP)
356 nd_log(NDLS_DAEMON, NDLP_DEBUG, "thread created with task id %d", gettid_cached());
357
360 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
361 - nd_log(NDLS_DAEMON, NDLP_WARNING, "cannot set pthread cancel type to DEFERRED.");
362 -
363 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
364 - nd_log(NDLS_DAEMON, NDLP_WARNING, "cannot set pthread cancel state to ENABLE.");
365 -
358 signals_block_all_except_deadly();
359
360 spinlock_lock(&threads_globals.running.spinlock);
@@ -374,7 +366,6 @@ static void *nd_thread_starting_point(void *ptr) {
366 nti->ret = nti->start_routine(nti->arg);
367
368 nd_thread_exit(nti);
377 - return nti;
369 }
370
371 ND_THREAD *nd_thread_self(void) {
@@ -385,6 +376,26 @@ bool nd_thread_is_me(ND_THREAD *nti) {
376 return nti && nti->thread == pthread_self();
377 }
378
379 +
380 +// utils
381 +#define MAX_THREAD_CREATE_RETRIES (10)
382 +#define MAX_THREAD_CREATE_WAIT_MS (1000)
383 +
384 +static int create_uv_thread(uv_thread_t *thread, uv_thread_cb thread_func, void *arg, int *retries)
385 +{
386 + int err;
387 +
388 + do {
389 + err = uv_thread_create(thread, thread_func, arg);
390 + if (err == 0)
391 + break;
392 +
393 + sleep_usec(MAX_THREAD_CREATE_WAIT_MS * USEC_PER_MS);
394 + } while (err == UV_EAGAIN && ++(*retries) < MAX_THREAD_CREATE_RETRIES);
395 +
396 + return err;
397 +}
398 +
399 ND_THREAD *nd_thread_create(const char *tag, NETDATA_THREAD_OPTIONS options, void *(*start_routine)(void *), void *arg)
400 {
401 ND_THREAD *nti = callocz(1, sizeof(*nti));
@@ -394,18 +405,18 @@ ND_THREAD *nd_thread_create(const char *tag, NETDATA_THREAD_OPTIONS options, voi
405 nti->options = (options & NETDATA_THREAD_OPTIONS_ALL);
406 strncpyz(nti->tag, tag, ND_THREAD_TAG_MAX);
407
397 - if ((options & NETDATA_THREAD_OPTION_JOINABLE) == 0)
398 - nd_log_daemon(NDLP_INFO, "WARNING: Creating detached thread '%s'", tag);
399 -
400 - int ret = pthread_create(&nti->thread, &threads_globals.attr, nd_thread_starting_point, nti);
408 + int retries = 0;
409 + int ret = create_uv_thread(&nti->thread, nd_thread_starting_point, nti, &retries);
410 if(ret != 0) {
411 nd_log(NDLS_DAEMON, NDLP_ERR,
403 - "failed to create new thread for %s. pthread_create() failed with code %d",
412 + "failed to create new thread for %s. uv_thread_create() failed with code %d",
413 tag, ret);
414
415 freez(nti);
416 return NULL;
417 }
418 + if (retries)
419 + nd_log_daemon(NDLP_WARNING, "nd_thread_create required %d attempts", retries);
420
421 return nti;
422 }
@@ -451,22 +462,16 @@ int nd_thread_join(ND_THREAD *nti) {
462 return 0;
463 }
464
454 - int ret = 0;
455 - bool joinable = nd_thread_status_check(nti, NETDATA_THREAD_OPTION_JOINABLE);
456 - if (joinable)
457 - ret = pthread_join(nti->thread, NULL);
458 - if(ret != 0) {
465 + int ret;
466 + if((ret = uv_thread_join(&nti->thread))) {
467 // we can't join the thread
468
469 nd_log(NDLS_DAEMON, NDLP_WARNING,
462 - "cannot join thread. pthread_join() failed with code %d. (tag=%s)",
470 + "cannot join thread. uv_thread_join() failed with code %d. (tag=%s)",
471 ret, nti->tag);
472 }
473 else {
474 // we successfully joined the thread
467 - if (joinable)
468 - nd_log(NDLS_DAEMON, NDLP_DEBUG, "Joining thread '%s', tid %d", nti->tag, nti->tid);
469 -
475 nd_thread_status_set(nti, NETDATA_THREAD_STATUS_JOINED);
476
477 spinlock_lock(&threads_globals.running.spinlock);
@@ -483,10 +488,7 @@ int nd_thread_join(ND_THREAD *nti) {
488 }
489 spinlock_unlock(&threads_globals.exited.spinlock);
490
486 - if (joinable)
487 - freez(nti);
488 - else
489 - nti->thread = 0;
491 + freez(nti);
492 }
493
494 return ret;
src/libnetdata/threads/threads.h
+6 -7
@@ -7,15 +7,14 @@
7
8 typedef enum __attribute__((packed)) {
9 NETDATA_THREAD_OPTION_DEFAULT = 0 << 0,
10 - NETDATA_THREAD_OPTION_JOINABLE = 1 << 0,
11 - NETDATA_THREAD_OPTION_DONT_LOG_STARTUP = 1 << 1,
12 - NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP = 1 << 2,
13 - NETDATA_THREAD_STATUS_STARTED = 1 << 3,
14 - NETDATA_THREAD_STATUS_FINISHED = 1 << 4,
15 - NETDATA_THREAD_STATUS_JOINED = 1 << 5,
10 + NETDATA_THREAD_OPTION_DONT_LOG_STARTUP = 1 << 0,
11 + NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP = 1 << 1,
12 + NETDATA_THREAD_STATUS_STARTED = 1 << 2,
13 + NETDATA_THREAD_STATUS_FINISHED = 1 << 3,
14 + NETDATA_THREAD_STATUS_JOINED = 1 << 4,
15 } NETDATA_THREAD_OPTIONS;
16
18 -#define NETDATA_THREAD_OPTIONS_ALL (NETDATA_THREAD_OPTION_JOINABLE | NETDATA_THREAD_OPTION_DONT_LOG_STARTUP | NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP)
17 +#define NETDATA_THREAD_OPTIONS_ALL (NETDATA_THREAD_OPTION_DONT_LOG_STARTUP | NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP)
18 #define NETDATA_THREAD_OPTION_DONT_LOG (NETDATA_THREAD_OPTION_DONT_LOG_STARTUP | NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP)
19
20 #define netdata_thread_cleanup_push(func, arg) pthread_cleanup_push(func, arg)
src/libnetdata/uuid/uuidmap.c
+1 -1
@@ -382,7 +382,7 @@ static int uuidmap_concurrent_unittest(void) {
382 snprintf(thread_name, sizeof(thread_name), "UUID-TEST-%d", i);
383 threads[i] = nd_thread_create(
384 thread_name,
385 - NETDATA_THREAD_OPTION_DONT_LOG | NETDATA_THREAD_OPTION_JOINABLE,
385 + NETDATA_THREAD_OPTION_DONT_LOG,
386 concurrent_test_thread,
387 &stats[i]);
388 }
src/ml/ml_public.cc
+2 -4
@@ -443,14 +443,12 @@ void ml_start_threads() {
443 char tag[NETDATA_THREAD_TAG_MAX + 1];
444
445 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "%s", "PREDICT");
446 - Cfg.detection_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
447 - ml_detect_main, NULL);
446 + Cfg.detection_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, ml_detect_main, NULL);
447
448 for (size_t idx = 0; idx != Cfg.num_worker_threads; idx++) {
449 ml_worker_t *worker = &Cfg.workers[idx];
450 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "TRAIN[%zu]", worker->id);
452 - worker->nd_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
453 - ml_train_main, worker);
451 + worker->nd_thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, ml_train_main, worker);
452 }
453 }
454
src/plugins.d/plugins_d.c
+1 -1
@@ -400,7 +400,7 @@ void *pluginsd_main(void *ptr) {
400 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "PD[%s]", pluginname);
401
402 // spawn a new thread for it
403 - cd->unsafe.thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
403 + cd->unsafe.thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT,
404 pluginsd_worker_thread, cd);
405 }
406 }
src/streaming/stream-connector.c
+1 -1
@@ -684,7 +684,7 @@ bool stream_connector_init(struct sender_state *s) {
684 snprintfz(tag, NETDATA_THREAD_TAG_MAX, THREAD_TAG_STREAM_SENDER "-CN" "[%d]",
685 sc->id);
686
687 - sc->thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE, stream_connector_thread, sc);
687 + sc->thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, stream_connector_thread, sc);
688 if (!sc->thread)
689 nd_log_daemon(NDLP_ERR,
690 "STREAM CONNECT '%s': failed to create new thread for client.",
src/streaming/stream-replication-sender.c
+1 -1
@@ -1719,7 +1719,7 @@ void *replication_thread_main(void *ptr) {
1719 char tag[NETDATA_THREAD_TAG_MAX + 1];
1720 snprintfz(tag, NETDATA_THREAD_TAG_MAX, "REPLAY[%zu]", i + 2);
1721 __atomic_add_fetch(&replication_buffers_allocated, sizeof(ND_THREAD *), __ATOMIC_RELAXED);
1722 - replication_globals.main_thread.threads_ptrs[i] = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
1722 + replication_globals.main_thread.threads_ptrs[i] = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT,
1723 replication_worker_thread, NULL);
1724 }
1725 }
src/streaming/stream-thread.c
+1 -1
@@ -728,7 +728,7 @@ static struct stream_thread * stream_thread_assign_and_start(RRDHOST *host) {
728 char tag[NETDATA_THREAD_TAG_MAX + 1];
729 snprintfz(tag, NETDATA_THREAD_TAG_MAX, THREAD_TAG_STREAM "[%zu]", sth->id);
730
731 - sth->thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE, stream_thread, sth);
731 + sth->thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, stream_thread, sth);
732 if (!sth->thread)
733 nd_log(NDLS_DAEMON, NDLP_ERR, "STREAM THREAD[%zu]: failed to create new thread for client.", sth->id);
734 }
src/web/api/queries/backfill.c
+1 -1
@@ -232,7 +232,7 @@ void *backfill_thread(void *ptr) {
232 for(size_t t = 0; t < threads - 1 ;t++) {
233 char tag[15];
234 snprintfz(tag, sizeof(tag), "BACKFILL[%zu]", t + 1);
235 - th[t] = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE, backfill_worker_thread, NULL);
235 + th[t] = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT, backfill_worker_thread, NULL);
236 }
237
238 backfill_worker_thread((void *)0x01);
src/web/server/static/static-threaded.c
+1 -1
@@ -384,7 +384,7 @@ void *socket_listen_main_static_threaded(void *ptr) {
384 char tag[50 + 1];
385 snprintfz(tag, sizeof(tag) - 1, "WEB[%d]", i+1);
386
387 - static_workers_private_data[i].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_JOINABLE,
387 + static_workers_private_data[i].thread = nd_thread_create(tag, NETDATA_THREAD_OPTION_DEFAULT,
388 socket_listen_main_static_threaded_worker,
389 (void *)&static_workers_private_data[i]);
390 }