@cryptotaxi247 / netdata-1 / commits / 42a58e4c9

abstracted pthreads into netdata_threads

Costa Tsaousis (ktsaou) committed Dec 28, 2017 at 04:01 UTC 42a58e4c9d0f29fec025db69cc1c14fbe660f031
30 files changed +374 -314
CMakeLists.txt
+1 -1
@@ -64,7 +64,7 @@ set(NETDATA_LINUX_FILES
64 src/sys_kernel_mm_ksm.c
65 src/sys_fs_cgroup.c
66 src/sys_fs_btrfs.c
67 - )
67 + src/threads.c src/threads.h)
68
69 set(NETDATA_COMMON_FILES
70 src/adaptive_resortable_list.c
src/Makefile.am
+42 -17
@@ -128,6 +128,8 @@ netdata_SOURCES = \
128 statsd.h \
129 storage_number.c \
130 storage_number.h \
131 + threads.c \
132 + threads.h \
133 unit_test.c \
134 unit_test.h \
135 url.c \
@@ -224,14 +226,22 @@ netdata_LDADD = \
226
227 apps_plugin_SOURCES = \
228 apps_plugin.c \
227 - avl.c avl.h \
228 - clocks.c clocks.h \
229 - common.c common.h \
229 + avl.c \
230 + avl.h \
231 + clocks.c \
232 + clocks.h \
233 + common.c \
234 + common.h \
235 inlined.h \
231 - locks.c locks.h \
236 + locks.c \
237 + locks.h \
238 log.c log.h \
233 - procfile.c procfile.h \
234 - web_buffer.c web_buffer.h \
239 + procfile.c \
240 + procfile.h \
241 + threads.c \
242 + threads.h \
243 + web_buffer.c \
244 + web_buffer.h \
245 $(NULL)
246
247 if FREEBSD
@@ -247,12 +257,18 @@ apps_plugin_LDADD = \
257
258 freeipmi_plugin_SOURCES = \
259 freeipmi_plugin.c \
250 - clocks.c clocks.h \
251 - common.c common.h \
260 + clocks.c \
261 + clocks.h \
262 + common.c \
263 + common.h \
264 inlined.h \
253 - locks.c locks.h \
265 + locks.c \
266 + locks.h \
267 log.c log.h \
255 - procfile.c procfile.h \
268 + procfile.c \
269 + procfile.h \
270 + threads.c \
271 + threads.h \
272 $(NULL)
273
274 freeipmi_plugin_LDADD = \
@@ -261,14 +277,23 @@ freeipmi_plugin_LDADD = \
277
278 cgroup_network_SOURCES = \
279 cgroup-network.c \
264 - clocks.c clocks.h \
265 - common.c common.h \
280 + clocks.c \
281 + clocks.h \
282 + common.c \
283 + common.h \
284 inlined.h \
267 - locks.c locks.h \
268 - log.c log.h \
269 - procfile.c procfile.h \
270 - popen.c popen.h \
271 - signals.c signals.h \
285 + locks.c \
286 + locks.h \
287 + log.c \
288 + log.h \
289 + procfile.c \
290 + procfile.h \
291 + popen.c \
292 + popen.h \
293 + signals.c \
294 + signals.h \
295 + threads.c \
296 + threads.h \
297 $(NULL)
298
299 cgroup_network_LDADD = \
src/backends.c
+2 -5
@@ -508,15 +508,13 @@ static void backends_main_cleanup(void *ptr) {
508 }
509
510 void *backends_main(void *ptr) {
511 - netdata_thread_welcome("BACKEND");
512 -
511 int default_port = 0;
512 int sock = -1;
513 BUFFER *b = buffer_create(1), *response = buffer_create(1);
514 int (*backend_request_formatter)(BUFFER *, const char *, RRDHOST *, const char *, RRDSET *, RRDDIM *, time_t, time_t, uint32_t) = NULL;
515 int (*backend_response_checker)(BUFFER *) = NULL;
516
519 - pthread_cleanup_push(backends_main_cleanup, ptr);
517 + netdata_thread_cleanup_push(backends_main_cleanup, ptr);
518
519 // ------------------------------------------------------------------------
520 // collect configuration options
@@ -913,7 +911,6 @@ cleanup:
911 buffer_free(b);
912 buffer_free(response);
913
916 - pthread_cleanup_pop(1);
917 - pthread_exit(NULL);
914 + netdata_thread_cleanup_pop(1);
915 return NULL;
916 }
src/common.c
-16
@@ -1116,22 +1116,6 @@ int fd_is_valid(int fd) {
1116 return fcntl(fd, F_GETFD) != -1 || errno != EBADF;
1117 }
1118
1119 -pid_t gettid(void) {
1120 -#ifdef __FreeBSD__
1121 - return (pid_t)pthread_getthreadid_np();
1122 -#elif defined(__APPLE__)
1123 -#if (defined __MAC_OS_X_VERSION_MIN_REQUIRED && __MAC_OS_X_VERSION_MIN_REQUIRED >= 1060)
1124 - uint64_t curthreadid;
1125 - pthread_threadid_np(NULL, &curthreadid);
1126 - return (pid_t)curthreadid;
1127 -#else /* __MAC_OS_X_VERSION_MIN_REQUIRED */
1128 - return (pid_t)pthread_self;
1129 -#endif /* __MAC_OS_X_VERSION_MIN_REQUIRED */
1130 -#else /* __APPLE__*/
1131 - return (pid_t)syscall(SYS_gettid);
1132 -#endif /* __FreeBSD__, __APPLE__*/
1133 -}
1134 -
1119 char *fgets_trim_len(char *buf, size_t buf_size, FILE *fp, size_t *len) {
1120 char *s = fgets(buf, (int)buf_size, fp);
1121 if (!s) return NULL;
src/common.h
+1 -2
@@ -172,6 +172,7 @@
172
173 #include "clocks.h"
174 #include "log.h"
175 +#include "threads.h"
176 #include "locks.h"
177 #include "simple_pattern.h"
178 #include "avl.h"
@@ -294,8 +295,6 @@ extern int fd_is_valid(int fd);
295
296 extern int enable_ksm;
297
297 -extern pid_t gettid(void);
298 -
298 extern int sleep_usec(usec_t usec);
299
300 extern char *fgets_trim_len(char *buf, size_t buf_size, FILE *fp, size_t *len);
src/health.c
+2 -5
@@ -348,10 +348,8 @@ static void health_main_cleanup(void *ptr) {
348 }
349
350 void *health_main(void *ptr) {
351 - netdata_thread_welcome("HEALTH");
352 -
351 BUFFER *wb = buffer_create(100);
354 - pthread_cleanup_push(health_main_cleanup, ptr);
352 + netdata_thread_cleanup_push(health_main_cleanup, ptr);
353
354 int min_run_every = (int)config_get_number(CONFIG_SECTION_HEALTH, "run at least every seconds", 10);
355 if(min_run_every < 1) min_run_every = 1;
@@ -736,7 +734,6 @@ void *health_main(void *ptr) {
734
735 buffer_free(wb);
736
739 - pthread_cleanup_pop(1);
740 - pthread_exit(NULL);
737 + netdata_thread_cleanup_pop(1);
738 return NULL;
739 }
src/locks.c
-24
@@ -1,29 +1,5 @@
1 #include "common.h"
2
3 -// ----------------------------------------------------------------------------
4 -// threads initialization
5 -
6 -static __thread char *netdata_thread_tag_name = NULL;
7 -
8 -const char *netdata_thread_tag(void) {
9 - return ((netdata_thread_tag_name && *netdata_thread_tag_name)?netdata_thread_tag_name:"unknown");
10 -}
11 -
12 -void netdata_thread_welcome_nolog(char *tag) {
13 - netdata_thread_tag_name = tag;
14 -
15 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
16 - error("%s: cannot set pthread cancel type to DEFERRED.", netdata_thread_tag());
17 -
18 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
19 - error("%s: cannot set pthread cancel state to ENABLE.", netdata_thread_tag());
20 -}
21 -
22 -void netdata_thread_welcome(char *tag) {
23 - netdata_thread_welcome_nolog(tag);
24 - info("%s: thread created with task id %d", netdata_thread_tag(), gettid());
25 -}
26 -
3 // ----------------------------------------------------------------------------
4 // automatic thread cancelability management, based on locks
5
src/locks.h
-4
@@ -7,10 +7,6 @@ typedef pthread_mutex_t netdata_mutex_t;
7 typedef pthread_rwlock_t netdata_rwlock_t;
8 #define NETDATA_RWLOCK_INITIALIZER PTHREAD_RWLOCK_INITIALIZER
9
10 -extern void netdata_thread_welcome(char *tag);
11 -extern void netdata_thread_welcome_nolog(char *tag);
12 -extern const char *netdata_thread_tag(void);
13 -
10 extern int __netdata_mutex_init(netdata_mutex_t *mutex);
11 extern int __netdata_mutex_lock(netdata_mutex_t *mutex);
12 extern int __netdata_mutex_trylock(netdata_mutex_t *mutex);
src/main.c
+22 -66
@@ -36,37 +36,37 @@ struct netdata_static_thread static_threads[] = {
36 #ifdef INTERNAL_PLUGIN_NFACCT
37 // nfacct requires root access
38 // so, we build it as an external plugin with setuid to root
39 - {"nfacct", CONFIG_SECTION_PLUGINS, "nfacct", 1, NULL, NULL, nfacct_main},
39 + {"PLUGIN_NFACCT", CONFIG_SECTION_PLUGINS, "nfacct", 1, NULL, NULL, nfacct_main},
40 #endif
41
42 #ifdef NETDATA_INTERNAL_CHECKS
43 // debugging plugin
44 - {"check", CONFIG_SECTION_PLUGINS, "checks", 0, NULL, NULL, checks_main},
44 + {"PLUGIN_CHECK", CONFIG_SECTION_PLUGINS, "checks", 0, NULL, NULL, checks_main},
45 #endif
46
47 #if defined(__FreeBSD__)
48 // FreeBSD internal plugins
49 - {"freebsd", CONFIG_SECTION_PLUGINS, "freebsd", 1, NULL, NULL, freebsd_main},
49 + {"PLUGIN_FREEBSD", CONFIG_SECTION_PLUGINS, "freebsd", 1, NULL, NULL, freebsd_main},
50 #elif defined(__APPLE__)
51 // macOS internal plugins
52 - {"macos", CONFIG_SECTION_PLUGINS, "macos", 1, NULL, NULL, macos_main},
52 + {"PLUGIN_MACOS", CONFIG_SECTION_PLUGINS, "macos", 1, NULL, NULL, macos_main},
53 #else
54 // linux internal plugins
55 - {"proc", CONFIG_SECTION_PLUGINS, "proc", 1, NULL, NULL, proc_main},
56 - {"diskspace", CONFIG_SECTION_PLUGINS, "diskspace", 1, NULL, NULL, proc_diskspace_main},
57 - {"cgroups", CONFIG_SECTION_PLUGINS, "cgroups", 1, NULL, NULL, cgroups_main},
58 - {"tc", CONFIG_SECTION_PLUGINS, "tc", 1, NULL, NULL, tc_main},
55 + {"PLUGIN_PROC", CONFIG_SECTION_PLUGINS, "proc", 1, NULL, NULL, proc_main},
56 + {"PLUGIN_DISKSPACE", CONFIG_SECTION_PLUGINS, "diskspace", 1, NULL, NULL, proc_diskspace_main},
57 + {"PLUGIN_CGROUP", CONFIG_SECTION_PLUGINS, "cgroups", 1, NULL, NULL, cgroups_main},
58 + {"PLUGIN_TC", CONFIG_SECTION_PLUGINS, "tc", 1, NULL, NULL, tc_main},
59 #endif /* __FreeBSD__, __APPLE__*/
60
61 // common plugins for all systems
62 - {"idlejitter", CONFIG_SECTION_PLUGINS, "idlejitter", 1, NULL, NULL, cpuidlejitter_main},
63 - {"backends", NULL, NULL, 1, NULL, NULL, backends_main},
64 - {"health", NULL, NULL, 1, NULL, NULL, health_main},
65 - {"plugins.d", NULL, NULL, 1, NULL, NULL, pluginsd_main},
66 - {"web", NULL, NULL, 1, NULL, NULL, socket_listen_main_multi_threaded},
67 - {"web-single-threaded", NULL, NULL, 0, NULL, NULL, socket_listen_main_single_threaded},
68 - {"push-metrics", NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
69 - {"statsd", NULL, NULL, 1, NULL, NULL, statsd_main},
62 + {"PLUGIN_IDLEJITTER", CONFIG_SECTION_PLUGINS, "idlejitter", 1, NULL, NULL, cpuidlejitter_main},
63 + {"BACKENDS", NULL, NULL, 1, NULL, NULL, backends_main},
64 + {"HEALTH", NULL, NULL, 1, NULL, NULL, health_main},
65 + {"PLUGINSD", NULL, NULL, 1, NULL, NULL, pluginsd_main},
66 + {"WEB_SERVER", NULL, NULL, 1, NULL, NULL, socket_listen_main_multi_threaded},
67 + {"WEB_SERVER_SINGLE_THREADED", NULL, NULL, 0, NULL, NULL, socket_listen_main_single_threaded},
68 + {"STREAMING", NULL, NULL, 0, NULL, NULL, rrdpush_sender_thread},
69 + {"STATSD", NULL, NULL, 1, NULL, NULL, statsd_main},
70
71 {NULL, NULL, NULL, 0, NULL, NULL, NULL}
72 };
@@ -184,22 +184,10 @@ void cancel_main_threads() {
184 for (i = 0; static_threads[i].name != NULL ; i++) {
185 if(static_threads[i].enabled) {
186 info("EXIT: Stopping master thread: %s", static_threads[i].name);
187 - int ret;
188 - if((ret = pthread_cancel(*static_threads[i].thread)) != 0)
189 - error("EXIT: pthread_cancel() failed with code %d.", ret);
190 - //else
191 - // info("MAIN: thread %s cancelled", static_threads[i].name);
192 -
187 + netdata_thread_cancel(*static_threads[i].thread);
188 static_threads[i].enabled = 0;
189 }
190 }
196 -
197 - // if, for any reason there is any child exited
198 - // catch it here
199 - info("EXIT: waiting for any unfinished child processes");
200 - siginfo_t info;
201 - waitid(P_PID, 0, &info, WEXITED|WNOHANG);
202 - info("EXIT: all threads/childs stopped.");
191 }
192
193 struct option_def option_definitions[] = {
@@ -606,8 +594,6 @@ int main(int argc, char **argv) {
594 int i;
595 int config_loaded = 0;
596 int dont_fork = 0;
609 - size_t wanted_stacksize = 0, stacksize = 0;
610 - pthread_attr_t attr;
597
598 // set the name for logging
599 program_name = "netdata";
@@ -932,21 +918,8 @@ int main(int argc, char **argv) {
918 // setup the signals we want to use
919 signals_init();
920
935 -
936 - // --------------------------------------------------------------------
937 - // get the required stack size of the threads of netdata
938 -
939 - i = pthread_attr_init(&attr);
940 - if(i != 0)
941 - fatal("pthread_attr_init() failed with code %d.", i);
942 -
943 - i = pthread_attr_getstacksize(&attr, &stacksize);
944 - if(i != 0)
945 - fatal("pthread_attr_getstacksize() failed with code %d.", i);
946 - else
947 - debug(D_OPTIONS, "initial pthread stack size is %zu bytes", stacksize);
948 -
949 - wanted_stacksize = (size_t)config_get_number(CONFIG_SECTION_GLOBAL, "pthread stack size", (long)stacksize);
921 + // setup threads configs
922 + netdata_threads_init();
923
924
925 // --------------------------------------------------------------------
@@ -1011,18 +984,7 @@ int main(int argc, char **argv) {
984 web_files_uid();
985 web_files_gid();
986
1014 -
1015 - // ------------------------------------------------------------------------
1016 - // set default pthread stack size - after we have forked
1017 -
1018 - if(stacksize < wanted_stacksize) {
1019 - i = pthread_attr_setstacksize(&attr, wanted_stacksize);
1020 - if(i != 0)
1021 - fatal("pthread_attr_setstacksize() to %zu bytes, failed with code %d.", wanted_stacksize, i);
1022 - else
1023 - debug(D_SYSTEM, "Successfully set pthread stacksize to %zu bytes", wanted_stacksize);
1024 - }
1025 -
987 + netdata_threads_init_after_fork();
988
989 // ------------------------------------------------------------------------
990 // initialize rrd, registry, health, rrdpush, etc.
@@ -1045,15 +1007,9 @@ int main(int argc, char **argv) {
1007 struct netdata_static_thread *st = &static_threads[i];
1008
1009 if(st->enabled) {
1048 - st->thread = mallocz(sizeof(pthread_t));
1049 -
1010 + st->thread = mallocz(sizeof(netdata_thread_t));
1011 debug(D_SYSTEM, "Starting thread %s.", st->name);
1051 -
1052 - if(pthread_create(st->thread, &attr, st->start_routine, st))
1053 - error("failed to create new thread for %s.", st->name);
1054 -
1055 - else if(pthread_detach(*st->thread))
1056 - error("Cannot request detach of newly created %s thread.", st->name);
1012 + netdata_thread_create(st->thread, st->name, NETDATA_THREAD_OPTION_DEFAULT, st->start_routine, st);
1013 }
1014 else debug(D_SYSTEM, "Not starting thread %s.", st->name);
1015 }
src/main.h
+1 -1
@@ -24,7 +24,7 @@ struct netdata_static_thread {
24
25 volatile sig_atomic_t enabled;
26
27 - pthread_t *thread;
27 + netdata_thread_t *thread;
28
29 void (*init_routine) (void);
30 void *(*start_routine) (void *);
src/plugin_checks.c
+2 -4
@@ -12,8 +12,7 @@ static void checks_main_cleanup(void *ptr) {
12 }
13
14 void *checks_main(void *ptr) {
15 - netdata_thread_welcome("CHECKS");
16 - pthread_cleanup_push(checks_main_cleanup, ptr);
15 + netdata_thread_cleanup_push(checks_main_cleanup, ptr);
16
17 usec_t usec = 0, susec = localhost->rrd_update_every * USEC_PER_SEC, loop_usec = 0, total_susec = 0;
18 struct timeval now, last, loop;
@@ -121,8 +120,7 @@ void *checks_main(void *ptr) {
120 rrdset_done(check3);
121 }
122
124 - pthread_cleanup_pop(1);
125 - pthread_exit(NULL);
123 + netdata_thread_cleanup_pop(1);
124 return NULL;
125 }
126
src/plugin_freebsd.c
+2 -4
@@ -76,8 +76,7 @@ static void freebsd_main_cleanup(void *ptr) {
76 }
77
78 void *freebsd_main(void *ptr) {
79 - netdata_thread_welcome("FREEBSD");
80 - pthread_cleanup_push(freebsd_main_cleanup, ptr);
79 + netdata_thread_cleanup_push(freebsd_main_cleanup, ptr);
80
81 int vdo_cpu_netdata = config_get_boolean("plugin:freebsd", "netdata server resources", 1);
82
@@ -169,7 +168,6 @@ void *freebsd_main(void *ptr) {
168 }
169 }
170
172 - pthread_cleanup_pop(1);
173 - pthread_exit(NULL);
171 + netdata_thread_cleanup_pop(1);
172 return NULL;
173 }
src/plugin_idlejitter.c
+2 -4
@@ -12,8 +12,7 @@ static void cpuidlejitter_main_cleanup(void *ptr) {
12 }
13
14 void *cpuidlejitter_main(void *ptr) {
15 - netdata_thread_welcome("IDLEJITTER");
16 - pthread_cleanup_push(cpuidlejitter_main_cleanup, ptr);
15 + netdata_thread_cleanup_push(cpuidlejitter_main_cleanup, ptr);
16
17 usec_t sleep_ut = config_get_number("plugin:idlejitter", "loop time in ms", CPU_IDLEJITTER_SLEEP_TIME_MS) * USEC_PER_MS;
18 if(sleep_ut <= 0) {
@@ -85,8 +84,7 @@ void *cpuidlejitter_main(void *ptr) {
84 }
85 }
86
88 - pthread_cleanup_pop(1);
89 - pthread_exit(NULL);
87 + netdata_thread_cleanup_pop(1);
88 return NULL;
89 }
90
src/plugin_macos.c
+2 -4
@@ -10,8 +10,7 @@ static void macos_main_cleanup(void *ptr) {
10 }
11
12 void *macos_main(void *ptr) {
13 - netdata_thread_welcome("MACOS");
14 - pthread_cleanup_push(macos_main_cleanup, ptr);
13 + netdata_thread_cleanup_push(macos_main_cleanup, ptr);
14
15 // when ZERO, attempt to do it
16 int vdo_cpu_netdata = !config_get_boolean("plugin:macos", "netdata server resources", 1);
@@ -63,7 +62,6 @@ void *macos_main(void *ptr) {
62 }
63 }
64
66 - pthread_cleanup_pop(1);
67 - pthread_exit(NULL);
65 + netdata_thread_cleanup_pop(1);
66 return NULL;
67 }
src/plugin_nfacct.c
+15 -18
@@ -751,17 +751,25 @@ static void nfacct_send_metrics() {
751
752 // ----------------------------------------------------------------------------
753
754 -void *nfacct_main(void *ptr) {
754 +static void nfacct_main_cleanup(void *ptr) {
755 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
756 + if(static_thread->enabled) {
757 + static_thread->enabled = 0;
758
757 - info("NETFILTER thread created with task id %d", gettid());
759 + info("%s: cleaning up...", netdata_thread_tag());
760
759 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
760 - error("NETFILTER: Cannot set pthread cancel type to DEFERRED.");
761 +#ifdef DO_NFACCT
762 + nfacct_cleanup();
763 +#endif
764
762 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
763 - error("NETFILTER: Cannot set pthread cancel state to ENABLE.");
765 +#ifdef DO_NFSTAT
766 + nfstat_cleanup();
767 +#endif
768 + }
769 +}
770
771 +void *nfacct_main(void *ptr) {
772 + netdata_thread_cleanup_push(nfacct_main_cleanup, ptr);
773
774 int update_every = (int)config_get_number("plugin:netfilter", "update every", localhost->rrd_update_every);
775 if(update_every < localhost->rrd_update_every)
@@ -805,18 +813,7 @@ void *nfacct_main(void *ptr) {
813 #endif
814 }
815
808 - info("NETFILTER thread exiting");
809 -
810 -#ifdef DO_NFACCT
811 - nfacct_cleanup();
812 -#endif
813 -
814 -#ifdef DO_NFSTAT
815 - nfstat_cleanup();
816 -#endif
817 -
818 - static_thread->enabled = 0;
819 - pthread_exit(NULL);
816 + netdata_thread_cleanup_pop(1);
817 return NULL;
818 }
819
src/plugin_proc.c
+2 -4
@@ -74,8 +74,7 @@ static void proc_main_cleanup(void *ptr) {
74 }
75
76 void *proc_main(void *ptr) {
77 - netdata_thread_welcome("PROC");
78 - pthread_cleanup_push(proc_main_cleanup, ptr);
77 + netdata_thread_cleanup_push(proc_main_cleanup, ptr);
78
79 int vdo_cpu_netdata = config_get_boolean("plugin:proc", "netdata server resources", 1);
80
@@ -163,8 +162,7 @@ void *proc_main(void *ptr) {
162 }
163 }
164
166 - pthread_cleanup_pop(1);
167 - pthread_exit(NULL);
165 + netdata_thread_cleanup_pop(1);
166 return NULL;
167 }
168
src/plugin_proc_diskspace.c
+2 -4
@@ -338,8 +338,7 @@ static void diskspace_main_cleanup(void *ptr) {
338 }
339
340 void *proc_diskspace_main(void *ptr) {
341 - netdata_thread_welcome("DISKSPACE");
342 - pthread_cleanup_push(diskspace_main_cleanup, ptr);
341 + netdata_thread_cleanup_push(diskspace_main_cleanup, ptr);
342
343 int vdo_cpu_netdata = config_get_boolean("plugin:proc", "netdata server resources", 1);
344
@@ -456,7 +455,6 @@ void *proc_diskspace_main(void *ptr) {
455 }
456 }
457
459 - pthread_cleanup_pop(1);
460 - pthread_exit(NULL);
458 + netdata_thread_cleanup_pop(1);
459 return NULL;
460 }
src/plugin_tc.c
+2 -4
@@ -853,8 +853,7 @@ static void tc_main_cleanup(void *ptr) {
853 }
854
855 void *tc_main(void *ptr) {
856 - netdata_thread_welcome("TC");
857 - pthread_cleanup_push(tc_main_cleanup, ptr);
856 + netdata_thread_cleanup_push(tc_main_cleanup, ptr);
857
858 struct rusage thread;
859
@@ -1158,7 +1157,6 @@ void *tc_main(void *ptr) {
1157 }
1158
1159 cleanup:
1161 - pthread_cleanup_pop(1);
1162 - pthread_exit(NULL);
1160 + netdata_thread_cleanup_pop(1);
1161 return NULL;
1162 }
src/plugins_d.c
+13 -27
@@ -498,8 +498,7 @@ static void pluginsd_worker_thread_cleanup(void *arg) {
498 }
499
500 void *pluginsd_worker_thread(void *arg) {
501 - netdata_thread_welcome_nolog("PLUGINSD_WORKER");
502 - pthread_cleanup_push(pluginsd_worker_thread_cleanup, arg);
501 + netdata_thread_cleanup_push(pluginsd_worker_thread_cleanup, arg);
502
503 struct plugind *cd = (struct plugind *)arg;
504
@@ -513,15 +512,11 @@ void *pluginsd_worker_thread(void *arg) {
512 break;
513 }
514
516 - info("PLUGINSD: '%s' running on pid %d", cd->fullfilename, cd->pid);
517 -
515 + info("%s on tid %d: connected to '%s' running on pid %d", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid);
516 count = pluginsd_process(localhost, cd, fp, 0);
519 - error("PLUGINSD: plugin '%s' disconnected.", cd->fullfilename);
520 -
517 + error("%s on tid %d: '%s' (pid %d) disconnected after %zu successful data collections (ENDs).", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, count);
518 killpid(cd->pid, SIGTERM);
519
523 - info("PLUGINSD: '%s' on pid %d stopped after %zu successful data collections (ENDs).", cd->fullfilename, cd->pid, count);
524 -
520 // get the return code
521 int code = mypclose(fp, cd->pid);
522
@@ -530,18 +525,18 @@ void *pluginsd_worker_thread(void *arg) {
525
526 if(likely(!cd->successful_collections)) {
527 // nothing collected - disable it
533 - error("PLUGINSD: '%s' exited with error code %d. Disabling it.", cd->fullfilename, code);
528 + error("%s on tid %d: '%s' (pid %d) exited with error code %d. Disabling it.", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, code);
529 cd->enabled = 0;
530 }
531 else {
532 // we have collected something
533
534 if(likely(cd->serial_failures <= 10)) {
540 - error("PLUGINSD: '%s' exited with error code %d, but has given useful output in the past (%zu times). %s", cd->fullfilename, code, cd->successful_collections, cd->enabled?"Waiting a bit before starting it again.":"Will not start it again - it is disabled.");
535 + error("%s on tid %d: '%s' (pid %d) exited with error code %d, but has given useful output in the past (%zu times). %s", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, code, cd->successful_collections, cd->enabled?"Waiting a bit before starting it again.":"Will not start it again - it is disabled.");
536 sleep((unsigned int) (cd->update_every * 10));
537 }
538 else {
544 - error("PLUGINSD: '%s' exited with error code %d, but has given useful output in the past (%zu times). We tried %zu times to restart it, but it failed to generate data. Disabling it.", cd->fullfilename, code, cd->successful_collections, cd->serial_failures);
539 + error("%s on tid %d: '%s' (pid %d) exited with error code %d, but has given useful output in the past (%zu times). We tried %zu times to restart it, but it failed to generate data. Disabling it.", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, code, cd->successful_collections, cd->serial_failures);
540 cd->enabled = 0;
541 }
542 }
@@ -553,11 +548,11 @@ void *pluginsd_worker_thread(void *arg) {
548 // we have collected nothing so far
549
550 if(likely(cd->serial_failures <= 10)) {
556 - error("PLUGINSD: '%s' (pid %d) does not generate useful output but it reports success (exits with 0). %s.", cd->fullfilename, cd->pid, cd->enabled?"Waiting a bit before starting it again.":"Will not start it again - it is disabled.");
551 + error("%s on tid %d: '%s' (pid %d) does not generate useful output but it reports success (exits with 0). %s.", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, cd->enabled?"Waiting a bit before starting it again.":"Will not start it again - it is now disabled.");
552 sleep((unsigned int) (cd->update_every * 10));
553 }
554 else {
560 - error("PLUGINSD: '%s' (pid %d) does not generate useful output, although it reports success (exits with 0), but we have tried %zu times to collect something. Disabling it.", cd->fullfilename, cd->pid, cd->serial_failures);
555 + error("%s on tid %d: '%s' (pid %d) does not generate useful output, although it reports success (exits with 0), but we have tried %zu times to collect something. Disabling it.", netdata_thread_tag(), gettid(), cd->fullfilename, cd->pid, cd->serial_failures);
556 cd->enabled = 0;
557 }
558 }
@@ -569,8 +564,7 @@ void *pluginsd_worker_thread(void *arg) {
564 if(unlikely(!cd->enabled)) break;
565 }
566
572 - pthread_cleanup_pop(1);
573 - pthread_exit(NULL);
567 + netdata_thread_cleanup_pop(1);
568 return NULL;
569 }
570
@@ -585,9 +579,7 @@ static void pluginsd_main_cleanup(void *data) {
579 for (cd = pluginsd_root; cd; cd = cd->next) {
580 if (cd->enabled && !cd->obsolete) {
581 info("PLUGINSD: Stopping plugin thread: %s", cd->id);
588 - int ret;
589 - if ((ret = pthread_cancel(cd->thread)) != 0)
590 - error("PLUGINSD: pthread_cancel() failed with code %d.", ret);
582 + netdata_thread_cancel(cd->thread);
583 }
584 }
585
@@ -596,8 +588,7 @@ static void pluginsd_main_cleanup(void *data) {
588 }
589
590 void *pluginsd_main(void *ptr) {
599 - netdata_thread_welcome("PLUGINSD");
600 - pthread_cleanup_push(pluginsd_main_cleanup, ptr);
591 + netdata_thread_cleanup_push(pluginsd_main_cleanup, ptr);
592
593 int automatic_run = config_get_boolean(CONFIG_SECTION_PLUGINS, "enable running new plugins", 1);
594 int scan_frequency = (int) config_get_number(CONFIG_SECTION_PLUGINS, "check for new plugins every", 60);
@@ -684,11 +675,7 @@ void *pluginsd_main(void *ptr) {
675
676 if(cd->enabled) {
677 // spawn a new thread for it
687 - if(unlikely(pthread_create(&cd->thread, NULL, pluginsd_worker_thread, cd) != 0))
688 - error("PLUGINSD: failed to create new thread for plugin '%s'.", cd->filename);
689 -
690 - else if(unlikely(pthread_detach(cd->thread) != 0))
691 - error("PLUGINSD: Cannot request detach of newly created thread for plugin '%s'.", cd->filename);
678 + netdata_thread_create(&cd->thread, "PLUGINSD_COLLECTOR", NETDATA_THREAD_OPTION_DEFAULT, pluginsd_worker_thread, cd);
679 }
680 }
681 }
@@ -699,7 +686,6 @@ void *pluginsd_main(void *ptr) {
686 sleep((unsigned int) scan_frequency);
687 }
688
702 - pthread_cleanup_pop(1);
703 - pthread_exit(NULL);
689 + netdata_thread_cleanup_pop(1);
690 return NULL;
691 }
src/plugins_d.h
+1 -1
@@ -27,7 +27,7 @@ struct plugind {
27 char cmd[PLUGINSD_CMD_MAX+1]; // the command that it executes
28
29 volatile pid_t pid;
30 - pthread_t thread;
30 + netdata_thread_t thread;
31
32 size_t successful_collections; // the number of times we have seen
33 // values collected from this plugin
src/rrd.h
+1 -1
@@ -437,7 +437,7 @@ struct rrdhost {
437 // the following are state information for the threading
438 // streaming metrics from this netdata to an upstream netdata
439 volatile int rrdpush_sender_spawn:1; // 1 when the sender thread has been spawn
440 - pthread_t rrdpush_sender_thread; // the sender thread
440 + netdata_thread_t rrdpush_sender_thread; // the sender thread
441
442 volatile int rrdpush_sender_connected:1; // 1 when the sender is ready to push metrics
443 int rrdpush_sender_socket; // the fd of the socket to the remote host, or -1
src/rrdpush.c
+12 -28
@@ -290,7 +290,7 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
290 rrdpush_buffer_lock(host);
291 rrdhost_wrlock(host);
292
293 - pthread_t thr = 0;
293 + netdata_thread_t thr = 0;
294
295 if(host->rrdpush_sender_spawn) {
296 info("STREAM %s [send]: signaling sending thread to stop...", host->hostname);
@@ -303,9 +303,7 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
303 thr = host->rrdpush_sender_thread;
304
305 // signal it to cancel
306 - int ret = pthread_cancel(host->rrdpush_sender_thread);
307 - if(ret != 0)
308 - error("STREAM %s [send]: pthread_cancel() returned error.", host->hostname);
306 + netdata_thread_cancel(host->rrdpush_sender_thread);
307 }
308
309 rrdhost_unlock(host);
@@ -313,12 +311,8 @@ void rrdpush_sender_thread_stop(RRDHOST *host) {
311
312 if(thr != 0) {
313 info("STREAM %s [send]: waiting for the sending thread to stop...", host->hostname);
316 -
314 void *result;
318 - int ret = pthread_join(thr, &result);
319 - if(ret != 0)
320 - error("STREAM %s [send]: pthread_join() returned error.", host->hostname);
321 -
315 + netdata_thread_join(thr, &result);
316 info("STREAM %s [send]: sending thread has exited.", host->hostname);
317 }
318 }
@@ -436,8 +430,7 @@ static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
430
431 if(!host->rrdpush_sender_join) {
432 info("STREAM %s [send]: sending thread detaches itself.", host->hostname);
439 - if(pthread_detach(pthread_self()))
440 - error("STREAM %s [send]: pthread_detach() failed.", host->hostname);
433 + netdata_thread_detach(netdata_thread_self());
434 }
435
436 host->rrdpush_sender_spawn = 0;
@@ -449,13 +442,10 @@ static void rrdpush_sender_thread_cleanup_callback(void *ptr) {
442 }
443
444 void *rrdpush_sender_thread(void *ptr) {
452 - netdata_thread_welcome_nolog("STREAM_SEND");
453 -
445 RRDHOST *host = (RRDHOST *)ptr;
446
447 if(!host->rrdpush_send_enabled || !host->rrdpush_send_destination || !*host->rrdpush_send_destination || !host->rrdpush_send_api_key || !*host->rrdpush_send_api_key) {
448 error("STREAM %s [send]: thread created (task id %d), but host has streaming disabled.", host->hostname, gettid());
458 - pthread_exit(NULL);
449 return NULL;
450 }
451
@@ -489,11 +479,11 @@ void *rrdpush_sender_thread(void *ptr) {
479
480 size_t not_connected_loops = 0;
481
492 - pthread_cleanup_push(rrdpush_sender_thread_cleanup_callback, host);
482 + netdata_thread_cleanup_push(rrdpush_sender_thread_cleanup_callback, host);
483
484 for(; host->rrdpush_send_enabled && !netdata_exit ;) {
485 // check for outstanding cancellation requests
496 - pthread_testcancel();
486 + netdata_thread_testcancel();
487
488 // if we don't have socket open, lets wait a bit
489 if(unlikely(host->rrdpush_sender_socket == -1)) {
@@ -680,8 +670,7 @@ void *rrdpush_sender_thread(void *ptr) {
670 }
671 }
672
683 - pthread_cleanup_pop(1);
684 - pthread_exit(NULL);
673 + netdata_thread_cleanup_pop(1);
674 return NULL;
675 }
676
@@ -892,8 +881,7 @@ static void rrdpush_receiver_thread_cleanup(void *ptr) {
881 }
882
883 static void *rrdpush_receiver_thread(void *ptr) {
895 - netdata_thread_welcome_nolog("STREAM_RECEIVE");
896 - pthread_cleanup_push(rrdpush_receiver_thread_cleanup, ptr);
884 + netdata_thread_cleanup_push(rrdpush_receiver_thread_cleanup, ptr);
885
886 struct rrdpush_thread *rpt = (struct rrdpush_thread *)ptr;
887 info("STREAM %s [%s]:%s: receive thread created (task id %d)", rpt->hostname, rpt->client_ip, rpt->client_port, gettid());
@@ -912,8 +900,7 @@ static void *rrdpush_receiver_thread(void *ptr) {
900 , rpt->client_port
901 );
902
915 - pthread_cleanup_pop(1);
916 - pthread_exit(NULL);
903 + netdata_thread_cleanup_pop(1);
904 return NULL;
905 }
906
@@ -921,7 +908,7 @@ static void rrdpush_sender_thread_spawn(RRDHOST *host) {
908 rrdhost_wrlock(host);
909
910 if(!host->rrdpush_sender_spawn) {
924 - if(pthread_create(&host->rrdpush_sender_thread, NULL, rrdpush_sender_thread, (void *) host))
911 + if(netdata_thread_create(&host->rrdpush_sender_thread, "STREAM_SENDER", NETDATA_THREAD_OPTION_JOINABLE, rrdpush_sender_thread, (void *) host))
912 error("STREAM %s [send]: failed to create new thread for client.", host->hostname);
913 else
914 host->rrdpush_sender_spawn = 1;
@@ -1055,16 +1042,13 @@ int rrdpush_receiver_thread_spawn(RRDHOST *host, struct web_client *w, char *url
1042 rpt->client_ip = strdupz(w->client_ip);
1043 rpt->client_port = strdupz(w->client_port);
1044 rpt->update_every = update_every;
1058 - pthread_t thread;
1045 + netdata_thread_t thread;
1046
1047 debug(D_SYSTEM, "STREAM [receive from [%s]:%s]: starting receiving thread.", w->client_ip, w->client_port);
1048
1062 - if(pthread_create(&thread, NULL, rrdpush_receiver_thread, (void *)rpt))
1049 + if(netdata_thread_create(&thread, "STREAM_RECEIVER", NETDATA_THREAD_OPTION_DEFAULT, rrdpush_receiver_thread, (void *)rpt))
1050 error("STREAM [receive from [%s]:%s]: failed to create new thread for client.", w->client_ip, w->client_port);
1051
1065 - else if(pthread_detach(thread))
1066 - error("STREAM [receive from [%s]:%s]: cannot request detach newly created thread.", w->client_ip, w->client_port);
1067 -
1052 // prevent the caller from closing the streaming socket
1053 if(w->ifd == w->ofd)
1054 w->ifd = w->ofd = -1;
src/socket.c
+2 -2
@@ -1158,7 +1158,7 @@ void poll_events(LISTEN_SOCKETS *sockets
1158
1159 int timeout = -1; // wait forever
1160
1161 - pthread_cleanup_push(poll_events_cleanup, &p);
1161 + netdata_thread_cleanup_push(poll_events_cleanup, &p);
1162
1163 for(;;) {
1164 if(unlikely(netdata_exit)) break;
@@ -1293,6 +1293,6 @@ void poll_events(LISTEN_SOCKETS *sockets
1293 }
1294 }
1295
1296 - pthread_cleanup_pop(1);
1296 + netdata_thread_cleanup_pop(1);
1297 debug(D_POLLFD, "POLLFD: LISTENER: cleanup completed");
1298 }
src/statsd.c
+11 -20
@@ -254,7 +254,7 @@ static struct statsd {
254 char *histogram_percentile_str;
255
256 int threads;
257 - pthread_t *collection_threads;
257 + netdata_thread_t *collection_threads;
258
259 LISTEN_SOCKETS sockets;
260 } statsd = {
@@ -897,13 +897,11 @@ void statsd_collector_thread_cleanup(void *data) {
897 }
898
899 void *statsd_collector_thread(void *ptr) {
900 - netdata_thread_welcome_nolog("STATSD_COLLECTOR");
901 -
900 int id = *((int *)ptr);
901 info("STATSD collector thread No %d created with task id %d", id + 1, gettid());
902
903 struct statsd_udp *d = callocz(sizeof(struct statsd_udp), 1);
906 - pthread_cleanup_push(statsd_collector_thread_cleanup, d);
904 + netdata_thread_cleanup_push(statsd_collector_thread_cleanup, d);
905
906 #ifdef HAVE_RECVMMSG
907 d->type = STATSD_SOCKET_DATA_TYPE_UDP;
@@ -929,8 +927,7 @@ void *statsd_collector_thread(void *ptr) {
927 , (void *)d
928 );
929
932 - pthread_cleanup_pop(1);
933 - pthread_exit(NULL);
930 + netdata_thread_cleanup_pop(1);
931 return NULL;
932 }
933
@@ -1982,7 +1979,7 @@ static void statsd_main_cleanup(void *data) {
1979 int i;
1980 for (i = 0; i < statsd.threads; i++) {
1981 info("STATSD: stopping data collection thread %d...", i + 1);
1985 - pthread_cancel(statsd.collection_threads[i]);
1982 + netdata_thread_cancel(statsd.collection_threads[i]);
1983 }
1984 }
1985
@@ -1994,8 +1991,7 @@ static void statsd_main_cleanup(void *data) {
1991 }
1992
1993 void *statsd_main(void *ptr) {
1997 - netdata_thread_welcome("STATSD");
1998 - pthread_cleanup_push(statsd_main_cleanup, ptr);
1994 + netdata_thread_cleanup_push(statsd_main_cleanup, ptr);
1995
1996 // ----------------------------------------------------------------------------------------------------------------
1997 // statsd configuration
@@ -2084,18 +2080,13 @@ void *statsd_main(void *ptr) {
2080 statsd_listen_sockets_setup();
2081 if(!statsd.sockets.opened) {
2082 error("STATSD: No statsd sockets to listen to. statsd will be disabled.");
2087 - pthread_exit(NULL);
2083 + goto cleanup;
2084 }
2085
2090 - statsd.collection_threads = callocz((size_t)statsd.threads, sizeof(pthread_t));
2086 + statsd.collection_threads = callocz((size_t)statsd.threads, sizeof(netdata_thread_t));
2087 int i;
2092 - for(i = 0; i < statsd.threads ;i++) {
2093 - if(pthread_create(&statsd.collection_threads[i], NULL, statsd_collector_thread, &i))
2094 - error("STATSD: failed to create child thread.");
2095 -
2096 - else if(pthread_detach(statsd.collection_threads[i]))
2097 - error("STATSD: cannot request detach of child thread.");
2098 - }
2088 + for(i = 0; i < statsd.threads ;i++)
2089 + netdata_thread_create(&statsd.collection_threads[i], "STATSD_COLLECTOR", NETDATA_THREAD_OPTION_DEFAULT, statsd_collector_thread, &i);
2090
2091 // ----------------------------------------------------------------------------------------------------------------
2092 // statsd monitoring charts
@@ -2286,7 +2277,7 @@ void *statsd_main(void *ptr) {
2277 break;
2278 }
2279
2289 - pthread_cleanup_pop(1);
2290 - pthread_exit(NULL);
2280 +cleanup:
2281 + netdata_thread_cleanup_pop(1);
2282 return NULL;
2283 }
src/sys_fs_cgroup.c
+2 -4
@@ -2684,8 +2684,7 @@ static void cgroup_main_cleanup(void *ptr) {
2684 }
2685
2686 void *cgroups_main(void *ptr) {
2687 - netdata_thread_welcome("CGROUP");
2688 - pthread_cleanup_push(cgroup_main_cleanup, ptr);
2687 + netdata_thread_cleanup_push(cgroup_main_cleanup, ptr);
2688
2689 struct rusage thread;
2690
@@ -2753,7 +2752,6 @@ void *cgroups_main(void *ptr) {
2752 }
2753 }
2754
2756 - pthread_cleanup_pop(1);
2757 - pthread_exit(NULL);
2755 + netdata_thread_cleanup_pop(1);
2756 return NULL;
2757 }
src/threads.c new
+178
@@ -0,0 +1,178 @@
1 +#include "common.h"
2 +
3 +static size_t stacksize = 0, wanted_stacksize = 0;
4 +static pthread_attr_t attr;
5 +
6 +// ----------------------------------------------------------------------------
7 +// per thread data
8 +
9 +typedef struct {
10 + void *arg;
11 + pthread_t *thread;
12 + const char *tag;
13 + void *(*start_routine) (void *);
14 + NETDATA_THREAD_OPTIONS options;
15 +} NETDATA_THREAD;
16 +
17 +static __thread NETDATA_THREAD *netdata_thread = NULL;
18 +
19 +const char *netdata_thread_tag(void) {
20 + return ((netdata_thread && netdata_thread->tag && *netdata_thread->tag)?netdata_thread->tag:"unknown");
21 +}
22 +
23 +// ----------------------------------------------------------------------------
24 +// compatibility library functions
25 +
26 +pid_t gettid(void) {
27 +#ifdef __FreeBSD__
28 +
29 + return (pid_t)pthread_getthreadid_np();
30 +
31 +#elif defined(__APPLE__)
32 +
33 + #if (defined __MAC_OS_X_VERSION_MIN_REQUIRED && __MAC_OS_X_VERSION_MIN_REQUIRED >= 1060)
34 + uint64_t curthreadid;
35 + pthread_threadid_np(NULL, &curthreadid);
36 + return (pid_t)curthreadid;
37 + #else /* __MAC_OS_X_VERSION_MIN_REQUIRED */
38 + return (pid_t)pthread_self;
39 + #endif /* __MAC_OS_X_VERSION_MIN_REQUIRED */
40 +
41 +#else /* __APPLE__*/
42 +
43 + return (pid_t)syscall(SYS_gettid);
44 +
45 +#endif /* __FreeBSD__, __APPLE__*/
46 +}
47 +
48 +// ----------------------------------------------------------------------------
49 +// early initialization
50 +
51 +void netdata_threads_init(void) {
52 + int i;
53 +
54 + // --------------------------------------------------------------------
55 + // get the required stack size of the threads of netdata
56 +
57 + i = pthread_attr_init(&attr);
58 + if(i != 0)
59 + fatal("pthread_attr_init() failed with code %d.", i);
60 +
61 + i = pthread_attr_getstacksize(&attr, &stacksize);
62 + if(i != 0)
63 + fatal("pthread_attr_getstacksize() failed with code %d.", i);
64 + else
65 + debug(D_OPTIONS, "initial pthread stack size is %zu bytes", stacksize);
66 +
67 + wanted_stacksize = (size_t)config_get_number(CONFIG_SECTION_GLOBAL, "pthread stack size", (long)stacksize);
68 +}
69 +
70 +// ----------------------------------------------------------------------------
71 +// late initialization
72 +
73 +void netdata_threads_init_after_fork(void) {
74 + int i;
75 +
76 + // ------------------------------------------------------------------------
77 + // set default pthread stack size
78 +
79 + if(stacksize < wanted_stacksize && wanted_stacksize > 0) {
80 + i = pthread_attr_setstacksize(&attr, wanted_stacksize);
81 + if(i != 0)
82 + fatal("pthread_attr_setstacksize() to %zu bytes, failed with code %d.", wanted_stacksize, i);
83 + else
84 + debug(D_SYSTEM, "Successfully set pthread stacksize to %zu bytes", wanted_stacksize);
85 + }
86 +}
87 +
88 +
89 +// ----------------------------------------------------------------------------
90 +// netdata_thread_create
91 +
92 +static void thread_cleanup(void *ptr) {
93 + if(netdata_thread != ptr) {
94 + NETDATA_THREAD *info = (NETDATA_THREAD *)ptr;
95 + error("THREADS: internal error - thread local variable does not match the one passed to this function. Expected thread '%s', passed thread '%s'", netdata_thread->tag, info->tag);
96 + }
97 +
98 + if(!(netdata_thread->options & NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP))
99 + info("%s: thread with task id %d finished", netdata_thread_tag(), gettid());
100 +
101 + freez((void *)netdata_thread->tag);
102 + freez(netdata_thread);
103 + netdata_thread->tag = NULL;
104 + netdata_thread = NULL;
105 +}
106 +
107 +static void *thread_start(void *ptr) {
108 + netdata_thread = (NETDATA_THREAD *)ptr;
109 +
110 + if(!(netdata_thread->options & NETDATA_THREAD_OPTION_DONT_LOG_STARTUP))
111 + info("%s: thread created with task id %d", netdata_thread_tag(), gettid());
112 +
113 + if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
114 + error("%s: cannot set pthread cancel type to DEFERRED.", netdata_thread_tag());
115 +
116 + if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
117 + error("%s: cannot set pthread cancel state to ENABLE.", netdata_thread_tag());
118 +
119 + void *ret = NULL;
120 + pthread_cleanup_push(thread_cleanup, ptr);
121 + ret = netdata_thread->start_routine(netdata_thread->arg);
122 + pthread_cleanup_pop(1);
123 +
124 + return ret;
125 +}
126 +
127 +int netdata_thread_create(netdata_thread_t *thread, const char *tag, NETDATA_THREAD_OPTIONS options, void *(*start_routine) (void *), void *arg) {
128 + NETDATA_THREAD *info = mallocz(sizeof(NETDATA_THREAD));
129 + info->arg = arg;
130 + info->thread = thread;
131 + info->tag = strdupz(tag);
132 + info->start_routine = start_routine;
133 + info->options = options;
134 +
135 + int ret = pthread_create(thread, &attr, thread_start, info);
136 + if(ret != 0)
137 + error("%s: failed to create new thread for %s. pthread_create() failed with code %d", netdata_thread_tag(), tag, ret);
138 +
139 + else {
140 + if (!(options & NETDATA_THREAD_OPTION_JOINABLE)) {
141 + int ret2 = pthread_detach(*thread);
142 + if (ret2 != 0)
143 + error("%s: cannot request detach of newly created %s thread. pthread_detach() failed with code %d", netdata_thread_tag(), tag, ret2);
144 + }
145 + }
146 +
147 + return ret;
148 +}
149 +
150 +// ----------------------------------------------------------------------------
151 +// netdata_thread_cancel
152 +
153 +int netdata_thread_cancel(netdata_thread_t thread) {
154 + int ret = pthread_cancel(thread);
155 + if(ret != 0)
156 + error("%s: cannot cancel thread. pthread_cancel() failed with code %d.", netdata_thread_tag(), ret);
157 +
158 + return ret;
159 +}
160 +
161 +// ----------------------------------------------------------------------------
162 +// netdata_thread_join
163 +
164 +int netdata_thread_join(netdata_thread_t thread, void **retval) {
165 + int ret = pthread_join(thread, retval);
166 + if(ret != 0)
167 + error("%s: cannot join thread. pthread_join() failed with code %d.", netdata_thread_tag(), ret);
168 +
169 + return ret;
170 +}
171 +
172 +int netdata_thread_detach(pthread_t thread) {
173 + int ret = pthread_detach(thread);
174 + if(ret != 0)
175 + error("%s: cannot detach thread. pthread_detach() failed with code %d.", netdata_thread_tag(), ret);
176 +
177 + return ret;
178 +}
src/threads.h new
+32
@@ -0,0 +1,32 @@
1 +#ifndef NETDATA_THREADS_H
2 +#define NETDATA_THREADS_H
3 +
4 +extern pid_t gettid(void);
5 +
6 +typedef enum {
7 + NETDATA_THREAD_OPTION_DEFAULT = 0 << 0,
8 + NETDATA_THREAD_OPTION_JOINABLE = 1 << 0,
9 + NETDATA_THREAD_OPTION_DONT_LOG_STARTUP = 1 << 1,
10 + NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP = 1 << 2,
11 + NETDATA_THREAD_OPTION_DONT_LOG = NETDATA_THREAD_OPTION_DONT_LOG_STARTUP|NETDATA_THREAD_OPTION_DONT_LOG_CLEANUP,
12 +} NETDATA_THREAD_OPTIONS;
13 +
14 +#define netdata_thread_cleanup_push(func, arg) pthread_cleanup_push(func, arg)
15 +#define netdata_thread_cleanup_pop(execute) pthread_cleanup_pop(execute)
16 +
17 +typedef pthread_t netdata_thread_t;
18 +
19 +extern const char *netdata_thread_tag(void);
20 +
21 +extern void netdata_threads_init(void);
22 +extern void netdata_threads_init_after_fork(void);
23 +
24 +extern int netdata_thread_create(netdata_thread_t *thread, const char *tag, NETDATA_THREAD_OPTIONS options, void *(*start_routine) (void *), void *arg);
25 +extern int netdata_thread_cancel(netdata_thread_t thread);
26 +extern int netdata_thread_join(netdata_thread_t thread, void **retval);
27 +extern int netdata_thread_detach(pthread_t thread);
28 +
29 +#define netdata_thread_self pthread_self
30 +#define netdata_thread_testcancel pthread_testcancel
31 +
32 +#endif //NETDATA_THREADS_H
src/web_client.c
+2 -4
@@ -1766,8 +1766,7 @@ static void web_client_main_cleanup(void *ptr) {
1766 }
1767
1768 void *web_client_main(void *ptr) {
1769 - netdata_thread_welcome_nolog("WEB_CLIENT");
1770 - pthread_cleanup_push(web_client_main_cleanup, ptr);
1769 + netdata_thread_cleanup_push(web_client_main_cleanup, ptr);
1770
1771 struct web_client *w = ptr;
1772 struct pollfd fds[2], *ifd, *ofd;
@@ -1898,7 +1897,6 @@ void *web_client_main(void *ptr) {
1897 w->ifd = -1;
1898 w->ofd = -1;
1899
1901 - pthread_cleanup_pop(1);
1902 - pthread_exit(NULL);
1900 + netdata_thread_cleanup_pop(1);
1901 return NULL;
1902 }
src/web_client.h
+1 -1
@@ -155,7 +155,7 @@ struct web_client {
155 size_t stats_received_bytes;
156 size_t stats_sent_bytes;
157
158 - pthread_t thread; // the thread servicing this client
158 + netdata_thread_t thread; // the thread servicing this client
159
160 struct web_client *prev;
161 struct web_client *next;
src/web_server.c
+19 -39
@@ -91,8 +91,6 @@ static inline void cleanup_web_clients(void) {
91 for (w = web_clients; w;) {
92 if (web_client_check_obsolete(w)) {
93 debug(D_WEB_CLIENT, "%llu: Removing client.", w->id);
94 - // pthread_cancel(w->thread);
95 - // pthread_join(w->thread, NULL);
94 w = web_client_free(w);
95 #ifdef NETDATA_INTERNAL_CHECKS
96 log_allocations();
@@ -104,8 +102,8 @@ static inline void cleanup_web_clients(void) {
102
103 // 1. it accepts new incoming requests on our port
104 // 2. creates a new web_client for each connection received
107 -// 3. spawns a new pthread to serve the client (this is optimal for keep-alive clients)
108 -// 4. cleans up old web_clients that their pthreads have been exited
105 +// 3. spawns a new netdata_thread to serve the client (this is optimal for keep-alive clients)
106 +// 4. cleans up old web_clients that their netdata_threads have been exited
107
108 #define CLEANUP_EVERY_EVENTS 100
109
@@ -131,20 +129,14 @@ static void socket_listen_main_multi_threaded_cleanup(void *data) {
129 for(w = web_clients; w ; w = w->next) {
130 if(!web_client_check_obsolete(w)) {
131 WEB_CLIENT_IS_OBSOLETE(w);
134 -
132 info("LISTENER: Stopping web client %s, id %llu", w->client_ip, w->id);
136 - int ret;
137 - if ((ret = pthread_cancel(w->thread)) != 0)
138 - error("LISTENER: pthread_cancel() failed with code %d, id %llu.", ret, w->id);
139 - //else
140 - // info("LISTENER: web client thread %s cancelled, id %llu", w->client_ip, w->id);
133 + netdata_thread_cancel(w->thread);
134 }
135 }
136 }
137
138 void *socket_listen_main_multi_threaded(void *ptr) {
146 - netdata_thread_welcome("WEBSERVER_MULTITHREADED");
147 - pthread_cleanup_push(socket_listen_main_multi_threaded_cleanup, ptr);
139 + netdata_thread_cleanup_push(socket_listen_main_multi_threaded_cleanup, ptr);
140
141 web_server_mode = WEB_SERVER_MODE_MULTI_THREADED;
142
@@ -201,14 +193,8 @@ void *socket_listen_main_multi_threaded(void *ptr) {
193 else
194 web_client_set_tcp(w);
195
204 - if(pthread_create(&w->thread, NULL, web_client_main, w) != 0) {
205 - error("%llu: failed to create new thread for web client.", w->id);
206 - WEB_CLIENT_IS_OBSOLETE(w);
207 - }
208 - else if(pthread_detach(w->thread) != 0) {
209 - error("%llu: Cannot request detach of newly created web client thread.", w->id);
196 + if(netdata_thread_create(&w->thread, "WEB_CLIENT", NETDATA_THREAD_OPTION_DEFAULT, web_client_main, w) != 0)
197 WEB_CLIENT_IS_OBSOLETE(w);
211 - }
198 }
199 }
200
@@ -220,8 +206,7 @@ void *socket_listen_main_multi_threaded(void *ptr) {
206 }
207 }
208
223 - pthread_cleanup_pop(1);
224 - pthread_exit(NULL);
209 + netdata_thread_cleanup_pop(1);
210 return NULL;
211 }
212
@@ -285,8 +270,7 @@ static void socket_listen_main_single_threaded_cleanup(void *data) {
270 }
271
272 void *socket_listen_main_single_threaded(void *ptr) {
288 - netdata_thread_welcome("WEBSERVER_SINGLETHREADED");
289 - pthread_cleanup_push(socket_listen_main_single_threaded_cleanup, ptr);
273 + netdata_thread_cleanup_push(socket_listen_main_single_threaded_cleanup, ptr);
274 web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
275
276 struct web_client *w;
@@ -400,8 +384,7 @@ void *socket_listen_main_single_threaded(void *ptr) {
384 }
385 }
386
403 - pthread_cleanup_pop(1);
404 - pthread_exit(NULL);
387 + netdata_thread_cleanup_pop(1);
388 return NULL;
389 }
390
@@ -510,18 +493,19 @@ static int web_server_snd_callback(int fd, int socktype, void *data, short int *
493 return 0;
494 }
495
513 -void *socket_listen_main_single_threaded(void *ptr) {
496 +static void socket_listen_main_single_threaded_cleanup(void *ptr) {
497 struct netdata_static_thread *static_thread = (struct netdata_static_thread *)ptr;
498 + if(static_thread->enabled) {
499 + static_thread->enabled = 0;
500
516 - web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
517 -
518 - info("Single-threaded WEB SERVER thread created with task id %d", gettid());
519 -
520 - if(pthread_setcanceltype(PTHREAD_CANCEL_DEFERRED, NULL) != 0)
521 - error("Cannot set pthread cancel type to DEFERRED.");
501 + info("%s: cleaning up...", netdata_thread_tag());
502 + listen_sockets_close(&api_sockets);
503 + }
504 +}
505
523 - if(pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL) != 0)
524 - error("Cannot set pthread cancel state to ENABLE.");
506 +void *socket_listen_main_single_threaded(void *ptr) {
507 + netdata_thread_cleanup_push(socket_listen_main_single_threaded_cleanup, ptr);
508 + web_server_mode = WEB_SERVER_MODE_SINGLE_THREADED;
509
510 if(!api_sockets.opened)
511 fatal("LISTENER: no listen sockets available.");
@@ -535,11 +519,7 @@ void *socket_listen_main_single_threaded(void *ptr) {
519 , NULL
520 );
521
538 - debug(D_WEB_CLIENT, "LISTENER: exit!");
539 - listen_sockets_close(&api_sockets);
540 -
541 - static_thread->enabled = 0;
542 - pthread_exit(NULL);
522 + netdata_thread_cleanup_pop(1);
523 return NULL;
524 }
525 #endif